kestra-io / kestra-io/plugin-oci
[Plugin] OCI — Data Integration (Workspaces & Pipelines)
- Dominant language
- Java
- Stars
- 0
- Forks
- 0
- Avg merge
- 18h 3m
- Merged PRs (30d)
- 1
Description
## Summary
The OCI Data Integration sub-plugin for `plugin-oci` enables Kestra flows to trigger and monitor OCI DI tasks, pipelines, and data loader tasks — bridging Kestra's orchestration layer with Oracle's managed ETL/ELT service. This lets teams sequence OCI Data Integration workloads alongside other pipeline steps (database startup, file landing, notification) without leaving the Kestra flow.
## Motivation
Organizations using OCI Data Integration to move data between Oracle data sources face an orchestration gap: OCI DI provides its own scheduling, but it cannot conditionally trigger based on upstream events, chain with non-OCI steps, or integrate with centralized alerting. A Kestra plugin task that submits a DI task run and waits for completion fills that gap, making OCI DI a composable step in any flow.
## Context
Part of the OCI Plugin Suite EPIC: https://github.com/kestra-io/plugin-oci/issues/2
Reference: similar async-run-and-wait pattern in `plugin-dbt` and `plugin-databricks`.
## API Reference
- **Official docs**: https://docs.oracle.com/en-us/iaas/api/#/en/dataintegration/latest/
- **Authentication**: Config-file (`~/.oci/config`), instance principal, or `SimpleAuthenticationDetailsProvider`
- **Base URL pattern**: `https://dataintegration.{region}.oci.oraclecloud.com/20200430/`
- **SDK**: OCI Java SDK v3.87.0 via BOM
## Gradle Dependencies
Add to `build.gradle`:
```groovy
// OCI Java SDK BOM
implementation platform("com.oracle.oci.sdk:oci-java-sdk-bom:3.87.0")
// Data Integration
implementation "com.oracle.oci.sdk:oci-java-sdk-dataintegration"
implementation "com.oracle.oci.sdk:oci-java-sdk-common"
```
## Plugin Structure
- **Repository**: `plugin-oci`
- **Namespace**: `io.kestra.plugin.oci.dataintegration`
- **Sub-plugins**: `dataintegration` (workspaces, task runs, pipeline runs)
## Suggested Tasks
1. `CreateTaskRun` — submit a DI task run in a workspace, optionally wait for completion
2. `GetTaskRun` — fetch task run status and emit lifecycle state as output
3. `ListTaskRuns` — list task runs in a workspace with optional state filter
4. `CreatePipelineRun` — submit a DI pipeline run, optionally wait for completion
5. `GetPipelineRun` — fetch pipeline run details
6. `ListWorkspaces` — list DI workspaces in a compartment
7. Add `TaskRunCompletedTrigger` — poll until a task run completes
8. Write unit + integration tests
## YAML Examples
### Example 1 — Run a DI task and wait for it to succeed
```yaml
id: run_oci_di_task
namespace: company.data
tasks:
- id: run_task
type: io.kestra.plugin.oci.dataintegration.CreateTaskRun
region: eu-frankfurt-1
tenancyOcid: "{{ secret('OCI_TENANCY_OCID') }}"
userId: "{{ secret('OCI_USER_OCID') }}"
fingerprint: "{{ secret('OCI_FINGERPRINT') }}"
privateKey: "{{ secret('OCI_PRIVATE_KEY') }}"
workspaceId: "{{ secret('OCI_DI_WORKSPACE_OCID') }}"
applicationKey: "{{ secret('OCI_DI_APP_KEY') }}"
taskKey: "{{ secret('OCI_DI_TASK_KEY') }}"
taskRunName: kestra-run-{{ execution.id }}
waitForCompletion: true
- id: log_result
type: io.kestra.plugin.core.log.Log
message: "Task run {{ outputs.run_task.taskRunKey }} — status: {{ outputs.run_task.status }}"
```
### Example 2 — Submit a pipeline run and list all task runs
```yaml
id: run_oci_di_pipeline
namespace: company.data
tasks:
- id: pipeline_run
type: io.kestra.plugin.oci.dataintegration.CreatePipelineRun
region: eu-frankfurt-1
tenancyOcid: "{{ secret('OCI_TENANCY_OCID') }}"
userId: "{{ secret('OCI_USER_OCID') }}"
fingerprint: "{{ secret('OCI_FINGERPRINT') }}"
privateKey: "{{ secret('OCI_PRIVATE_KEY') }}"
workspaceId: "{{ secret('OCI_DI_WORKSPACE_OCID') }}"
applicationKey: "{{ secret('OCI_DI_APP_KEY') }}"
pipelineKey: "{{ secret('OCI_DI_PIPELINE_KEY') }}"
waitForCompletion: true
- id: list_task_runs
type: io.kestra.plugin.oci.dataintegration.ListTaskRuns
region: eu-frankfurt-1
tenancyOcid: "{{ secret('OCI_TENANCY_OCID') }}"
userId: "{{ secret('OCI_USER_OCID') }}"
fingerprint: "{{ secret('OCI_FINGERPRINT') }}"
privateKey: "{{ secret('OCI_PRIVATE_KEY') }}"
workspaceId: "{{ secret('OCI_DI_WORKSPACE_OCID') }}"
applicationKey: "{{ secret('OCI_DI_APP_KEY') }}"
- id: log_count
type: io.kestra.plugin.core.log.Log
message: "Pipeline spawned {{ outputs.list_task_runs.count }} task runs"
```
### Example 3 — React when a DI task run completes
```yaml
id: on_di_task_run_complete
namespace: company.data
triggers:
- id: task_run_watcher
type: io.kestra.plugin.oci.dataintegration.TaskRunCompletedTrigger
region: eu-frankfurt-1
tenancyOcid: "{{ secret('OCI_TENANCY_OCID') }}"
userId: "{{ secret('OCI_USER_OCID') }}"
fingerprint: "{{ secret('OCI_FINGERPRINT') }}"
privateKey: "{{ secret('OCI_PRIVATE_KEY') }}"
workspaceId: "{{ secret('OCI_DI_WORKSPACE_OCID') }}"
taskRunKey: "{{ secret('OCI_DI_TASK_RUN_KEY') }}"
interval: PT2M
tasks:
- id: handle
type: io.kestra.plugin.core.log.Log
message: "DI task run {{ trigger.taskRunKey }} completed — status: {{ trigger.status }}"
```
## Acceptance Criteria
- [ ] `CreateTaskRun`, `GetTaskRun`, `ListTaskRuns`, `CreatePipelineRun`, `GetPipelineRun`, `ListWorkspaces` tasks implemented
- [ ] `CreateTaskRun` and `CreatePipelineRun` support `waitForCompletion: true`
- [ ] `TaskRunCompletedTrigger` polling trigger implemented
- [ ] All `Property` fields support Kestra expression language
- [ ] Unit + integration tests pass (`./gradlew test`)
- [ ] `package-info.java` with `@PluginSubGroup(category = PluginSubGroup.PluginCategory.CLOUD)`
- [ ] Build passes with `./gradlew build`
Contributor guide
No contributing guide indexed for this repository
Research direction
Start with build.gradle for the OCI Java SDK dependencies, then compare the async-run-and-wait patterns in plugin-dbt and plugin-databricks. Review the planned task and trigger APIs, add package-info.java with the required cloud subgroup annotation, and run ./gradlew test and ./gradlew build; done means all listed tasks, polling behavior, expression fields, and tests pass.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- java
- Domain
- cloud
- Issue type
- Feature
- Difficulty
- 5/5
- Estimated time
- Over a week
- Activity status
- Active
- Clarity
- Mostly clear
- Newbie friendliness
- 35/100