Actions processing
Ocean integrations can execute integration actions, which write data or trigger operations in the third-party tool, such as dispatching a GitHub workflow or creating a Jira issue.
This page describes how Ocean's actions processor picks up action runs from Port and executes them. The set of available actions is specific to each integration.
By default, actions run with the integration's own credentials. To run them as the user who triggered the workflow, see identity propagation.
Execution flow
- Trigger - A user, automation, or workflow runs an integration action in Port. Port creates a pending action run.
- Claim - The integration's actions processor polls Port for pending runs and claims a batch of them. A claimed run is hidden from other integration instances for a limited time, so the same run is not executed twice.
- Queue - Ocean places each claimed run in a queue, as described in queues and partitioning.
- Execute - A worker takes the run from its queue and executes the action against the third-party API.
- Report - Ocean reports the run's logs, status label, and final result (success or failure) back to Port, where you can follow it on the run's page.
Some actions hand work off to a long-running process in the third-party tool, such as a CI pipeline. In that case the run stays in progress until the integration receives a live event reporting that the process finished.
Queues and partitioning
Ocean uses two kinds of queues:
- Global queue - Runs that can execute in parallel.
- Partition queues - Runs that must execute one after another. Each action can define a partition key, for example the target repository. Runs with the same partition key execute sequentially, in order, to avoid conflicting changes.
A pool of workers processes runs from all queues in a round-robin order, so a busy partition does not block the others.
If the same run arrives twice while it is still being processed, Ocean skips the duplicate.
Flow control
The actions processor protects both the integration and the third-party API from overload:
- Buffer limit - Ocean stops claiming new runs when the total number of queued runs reaches
runsBufferHighWatermark, and resumes once workers free up space. - Per-action limit - A single action cannot fill more than
maxRunsBufferUtilPctPerActionpercent of the buffer, so a burst of runs for one action does not starve the others. - Rate limiting - Before executing a run, Ocean checks whether the third-party API is close to its rate limit. If it is, Ocean pauses execution until the limit resets, backing off for up to 10 seconds at a time.
Status labels
While a run is in progress, the integration can attach a short status label, such as Dispatching workflow or Workflow running, which Port displays next to the run.
A failed run always ends with a label. If the integration does not provide a specific one, Ocean sets the label to Execution failed.
What an action looks like in code
An integration implements each action as an executor, a class that extends Ocean's AbstractExecutor. The examples below are adapted from the Jira, GitHub, and Azure DevOps integrations.
Declare the action
The integration declares its actions and their inputs in its .port/spec.json file. Port uses this declaration to show the action and its inputs in the workflow builder:
{
"actionsProcessingEnabled": true,
"actions": [
{
"name": "add_comment",
"description": "Add a comment to a Jira issue",
"inputs": [
{ "name": "issueKey", "title": "Issue key", "type": "string", "required": true },
{ "name": "comment", "title": "Comment", "type": "string", "required": true }
]
}
]
}
Implement the executor
The executor's ACTION_NAME matches the action's name in the spec. The inputs sent from Port are available in run.execution_properties:
from port_ocean.context.ocean import ocean
from port_ocean.core.handlers.actions.abstract_executor import AbstractExecutor
from port_ocean.core.models import IntegrationRun, WorkflowNodeRun
class AddCommentExecutor(AbstractExecutor):
ACTION_NAME = "add_comment"
WEBHOOK_PROCESSOR_CLASS = None
async def is_close_to_rate_limit(self, run: IntegrationRun) -> bool:
return False
async def get_remaining_seconds_until_rate_limit(self, run: IntegrationRun) -> float:
return 0.0
async def execute(self, run: IntegrationRun) -> None:
issue_key = run.execution_properties["issueKey"]
# Log to the run in Port, and set its status label
await ocean.port_client.post_run_log(
run,
f"Adding comment to Jira issue {issue_key}",
status_label="Adding comment",
)
comment = await self.client.add_comment(
issue_key, run.execution_properties["comment"]
)
# Expose an output to the next nodes in the workflow
if isinstance(run, WorkflowNodeRun):
run.output = {"issueKey": issue_key, "commentId": comment.id}
# Mark the run as successful
await ocean.port_client.report_run_completed(
run,
success=True,
message=f"Added comment to {issue_key}",
status_label="Comment added",
)
If execute raises an exception, Ocean marks the run as failed and reports the error message to Port. To set a specific status label for the failure, raise an ActionExecutionError:
from port_ocean.exceptions.execution_manager import ActionExecutionError
class AddCommentError(ActionExecutionError):
DEFAULT_STATUS_LABEL = "Comment failed"
Executors are registered once, when the integration starts:
ocean.register_action_executor(AddCommentExecutor())
Partition keys and rate limits
The GitHub integration runs actions on the same repository one after another, and pauses execution when its GitHub API quota is running low:
class AbstractGithubExecutor(AbstractExecutor):
async def _get_partition_key(self, run: IntegrationRun) -> str | None:
# Runs with the same key execute sequentially. None means "run in parallel".
org = run.execution_properties.get("org")
repo = run.execution_properties.get("repo")
if not isinstance(org, str) or not isinstance(repo, str):
return None
return f"{org}/{repo}"
async def is_close_to_rate_limit(self, run: IntegrationRun) -> bool:
info = (await self._get_rest_client(run)).get_rate_limit_status()
return bool(info) and info.remaining < MIN_REMAINING_RATE_LIMIT_FOR_ACTIONS
async def get_remaining_seconds_until_rate_limit(self, run: IntegrationRun) -> float:
info = (await self._get_rest_client(run)).get_rate_limit_status()
return info.seconds_until_reset if info else 0.0
Long-running actions
When an action starts a long-running process, the executor marks the run as started and returns. A live events processor reports the result once the third-party tool sends a webhook about it.
The Azure DevOps trigger_pipeline action links the pipeline run to the Port run using an external ID:
class TriggerPipelineExecutor(AbstractExecutor):
ACTION_NAME = "trigger_pipeline"
# The processor that reports completion, and the endpoint it listens on
WEBHOOK_PROCESSOR_CLASS = PipelineRunActionWebhookProcessor
WEBHOOK_PATH = "/webhook"
async def execute(self, run: IntegrationRun) -> None:
pipeline_run = await self.client.run_pipeline(project_id, pipeline_id, options)
external_id = build_external_id(project_id, pipeline_id, pipeline_run["id"])
link = pipeline_run["_links"]["web"]["href"]
# The run stays in progress in Port, with a link to the pipeline run
await ocean.port_client.update_run_started(run, link, external_id)
class PipelineRunActionWebhookProcessor(AbstractWebhookProcessor):
@classmethod
def get_processor_type(cls) -> WebhookProcessorType:
# An ACTION processor updates runs, not catalog entities
return WebhookProcessorType.ACTION
async def handle_event(self, payload, resource_config) -> WebhookEventRawResults:
run = payload["resource"]["run"]
external_id = build_external_id(project_id, pipeline_id, run["id"])
# Find the Port run that started this pipeline run, and complete it
port_run = await ocean.port_client.find_run_by_external_id(external_id)
if port_run and ocean.port_client.is_run_in_progress(port_run):
await ocean.port_client.report_run_completed(
port_run, run["result"] == "succeeded", f"Pipeline run {run['result']}"
)
return WebhookEventRawResults(updated_raw_results=[], deleted_raw_results=[])
Configuration
Action processing is disabled by default. To enable it, set actionsProcessor.enabled to true (OCEAN__ACTIONS_PROCESSOR__ENABLED environment variable). The integration's event listener must also support actions.
actionsProcessor:
enabled: true
workersCount: 3
runsBufferHighWatermark: 300
pollCheckIntervalSeconds: 10
visibilityTimeoutMs: 60000
maxRunsBufferUtilPctPerAction: 30
| Parameter | Environment variable | Description | Default |
|---|---|---|---|
enabled | OCEAN__ACTIONS_PROCESSOR__ENABLED | Whether to start the actions processor and poll Port for pending runs. | false |
workersCount | OCEAN__ACTIONS_PROCESSOR__WORKERS_COUNT | The number of concurrent workers. Tune according to the CPU and memory allocated to the integration. | 3 |
runsBufferHighWatermark | OCEAN__ACTIONS_PROCESSOR__RUNS_BUFFER_HIGH_WATERMARK | The maximum number of queued runs before Ocean stops claiming new ones (1-1000). | 300 |
pollCheckIntervalSeconds | OCEAN__ACTIONS_PROCESSOR__POLL_CHECK_INTERVAL_SECONDS | The time in seconds between polls for pending runs. | 10 |
visibilityTimeoutMs | OCEAN__ACTIONS_PROCESSOR__VISIBILITY_TIMEOUT_MS | The time in milliseconds a claimed run stays hidden from other consumers before it can be claimed again (1-600000). | 60000 |
maxRunsBufferUtilPctPerAction | OCEAN__ACTIONS_PROCESSOR__MAX_RUNS_BUFFER_UTIL_PCT_PER_ACTION | The maximum percentage of the buffer a single action can use before Ocean stops claiming its runs (1-100). | 30 |
On shutdown, Ocean stops claiming new runs and waits up to maxWaitSecondsBeforeShutdown seconds (OCEAN__MAX_WAIT_SECONDS_BEFORE_SHUTDOWN, default 5) for workers to finish their current runs.