> For the complete documentation index, see llms.txt.
Skip to main content

Check out Port for yourself ➜ 

Live events processing

Live events keep the catalog up to date between resyncs. This page describes what happens inside an Ocean integration when a third-party tool sends it an event, usually as a webhook.

Which events are supported, and how to register the webhook in the third-party tool, is specific to each integration and documented on its page.

Live events flow​

  1. Receive - The third-party tool sends an HTTP POST request to one of the integration's live events endpoints. Ocean exposes these endpoints under the /integration path prefix, for example /integration/webhook.
  2. Queue - Ocean adds the event to a queue and responds to the sender immediately, so slow processing never causes the third-party tool to time out.
  3. Route - Ocean passes the event to every live events processor registered for that endpoint. Each processor decides whether it handles the event, and which kinds the event affects.
  4. Authenticate and validate - The processor verifies that the request is legitimate, for example by checking a webhook signature, and that the payload has the expected structure. Events that fail either check are rejected.
  5. Fetch - The processor usually fetches the latest version of the resource from the third-party API, since webhook payloads often contain only partial data.
  6. Send to Port - Ocean sends the updated and deleted resources to Port, which applies your mapping and updates the catalog.
Kinds must be in the mapping

Ocean processes a live event only for kinds that exist in the integration's mapping. If you remove a kind from the mapping, live events for it are ignored.

What a live events processor looks like​

Each step of the flow above maps to a method of a live events processor. The following processor is adapted from the Jira integration, and handles Jira project events:

webhook_processors/project_webhook_processor.py
from port_ocean.core.handlers.port_app_config.models import ResourceConfig
from port_ocean.core.handlers.webhook.abstract_webhook_processor import (
AbstractWebhookProcessor,
)
from port_ocean.core.handlers.webhook.webhook_event import (
EventHeaders,
EventPayload,
WebhookEvent,
WebhookEventRawResults,
)


class ProjectWebhookProcessor(AbstractWebhookProcessor):
# Route: does this processor handle the event?
async def should_process_event(self, event: WebhookEvent) -> bool:
return event.payload.get("webhookEvent", "").startswith("project_")

# Route: which kinds in the mapping does the event affect?
async def get_matching_kinds(self, event: WebhookEvent) -> list[str]:
return ["project"]

# Authenticate and validate: reject illegitimate or malformed requests.
async def authenticate(self, payload: EventPayload, headers: EventHeaders) -> bool:
return True

async def validate_payload(self, payload: EventPayload) -> bool:
return "project" in payload

# Fetch: get the latest version of the resource from the Jira API.
async def handle_event(
self, payload: EventPayload, resource_config: ResourceConfig
) -> WebhookEventRawResults:
project_key = payload["project"]["key"]

if payload["webhookEvent"] == "project_soft_deleted":
return WebhookEventRawResults(
updated_raw_results=[],
deleted_raw_results=[payload["project"]],
)

project = await get_or_create_jira_client().get_single_project(project_key)
return WebhookEventRawResults(
updated_raw_results=[project],
deleted_raw_results=[],
)

Ocean then sends the updated_raw_results and deleted_raw_results to Port, which maps them using the mapping of each matching kind.

Processors are registered in main.py. Several processors can share the same endpoint, and Ocean passes each event to all of them:

main.py
ocean.add_webhook_processor("/webhook", IssueWebhookProcessor)
ocean.add_webhook_processor("/webhook", ProjectWebhookProcessor)

Integrations that support webhook signatures verify them in authenticate. For example, the GitHub integration compares the x-hub-signature-256 header with an HMAC SHA-256 of the payload, computed with the webhook secret.

Queuing and workers​

Each live events endpoint has its own queue. Ocean processes events from each queue using a pool of workers, configured with the eventWorkersCount parameter (OCEAN__EVENT_WORKERS_COUNT environment variable). The default is 1, which processes events from each endpoint one at a time.

Increasing the number of workers processes events in parallel, at the cost of more concurrent calls to the third-party API.

When live events at scale is enabled, events are consumed from a Redis stream instead of the in-memory queue. To process events in parallel in this mode, scale the number of consumer pods.

Retries and timeouts​

If processing an event fails with a retryable error, Ocean retries it with exponential backoff. By default, each event is retried up to 3 times, starting with a 1-second delay and capped at 30 seconds between attempts. Each integration decides which errors are retryable, and can override these values.

A processor overrides the retry behavior with class attributes and the should_retry method. By default, only errors of type RetryableError are retried:

import httpx


class ProjectWebhookProcessor(AbstractWebhookProcessor):
max_retries = 5
initial_retry_delay_seconds = 2.0
max_retry_delay_seconds = 60.0

def should_retry(self, error: Exception) -> bool:
# Keep the default behavior, and also retry when Jira returns a 5xx error
if super().should_retry(error):
return True
return (
isinstance(error, httpx.HTTPStatusError)
and error.response.status_code >= 500
)

Each event has a processing time limit, including retries, set by the maxEventProcessingSeconds parameter (OCEAN__MAX_EVENT_PROCESSING_SECONDS environment variable, default 90). Events that exceed it fail.

When the integration shuts down, Ocean cancels in-flight events and waits up to maxWaitSecondsBeforeShutdown seconds (OCEAN__MAX_WAIT_SECONDS_BEFORE_SHUTDOWN, default 5) for them to stop.

Live events and resyncs​

Live events are processed independently of resyncs. They keep running while a full resync is in progress, and a resync never waits for live events to finish. For how Port reconciles entities updated by both, see sync interactions.