kestra-io / kestra-io/plugin-oci

[Plugin] OCI — Data Integration (Workspaces & Pipelines)

Open
#9 0 comments 0 reactions 0 assignees View on GitHub
area/plugin
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

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.