Implement live events
Live events keep the catalog up to date between resyncs. In code they are webhook processors: HTTP POST handlers that authenticate the request, decide which kinds are affected, and return updated or deleted raw objects.
For how Ocean queues, routes, and retries live events at runtime, see live events processing.
Processor interface
Subclass AbstractWebhookProcessor and implement:
| Method | Purpose |
|---|---|
should_process_event | Return whether this processor handles the event. |
get_matching_kinds | Return the kinds the event affects. |
authenticate | Verify the request (signature, token, and similar). |
validate_payload | Reject malformed payloads. |
handle_event | Fetch current data if needed and return raw results. |
handle_event must return WebhookEventRawResults with:
updated_raw_results- raw objects to upsert via mapping.deleted_raw_results- raw objects that identify deletes.
Prefer fetching the latest resource from the third-party API instead of registering the webhook payload as-is. Webhook bodies are often partial and can diverge from resync shapes.
Example
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,
)
from client import MyClient
from kinds import ObjectKind
class ProjectWebhookProcessor(AbstractWebhookProcessor):
async def should_process_event(self, event: WebhookEvent) -> bool:
return event.payload.get("eventType", "").startswith("project.")
async def get_matching_kinds(self, event: WebhookEvent) -> list[str]:
return [ObjectKind.PROJECT]
async def authenticate(
self, payload: EventPayload, headers: EventHeaders
) -> bool:
return True
async def validate_payload(self, payload: EventPayload) -> bool:
return "project" in payload
async def handle_event(
self, payload: EventPayload, resource_config: ResourceConfig
) -> WebhookEventRawResults:
client = MyClient()
project_id = payload["project"]["id"]
if payload.get("eventType") == "project.deleted":
return WebhookEventRawResults(
updated_raw_results=[],
deleted_raw_results=[payload["project"]],
)
project = await client.get_project(project_id)
if not project:
return WebhookEventRawResults(
updated_raw_results=[],
deleted_raw_results=[payload["project"]],
)
return WebhookEventRawResults(
updated_raw_results=[project],
deleted_raw_results=[],
)
Register processors
Register processors on a fixed path under the integration's /integration prefix (for example /webhook):
from port_ocean.context.ocean import ocean
from webhook_processors.project_webhook_processor import ProjectWebhookProcessor
ocean.add_webhook_processor("/webhook", ProjectWebhookProcessor)
Register the matching webhook URL in the third-party system from @ocean.on_start(), using your public base URL. See advanced configuration for base URL and path prefix settings.
Ocean processes live events only for kinds that exist in the integration's mapping. If you remove a kind from the mapping, live events for it are ignored.