Outgestion Pipeline
Design of the outgestion pipeline for pushing data to external partner systems.
Overview
The outgestion pipeline is the counterpart to the ingestion pipeline. While ingestion brings data into the registry from external systems, outgestion pushes registry data changes outward to partner systems, external registries, and downstream consumers.
The pipeline is fully asynchronous, built on Celery workers with Redis as the message broker. When a change request is approved in the registry, the outgestion pipeline captures the change, transforms it into the partner's expected schema using Jinja2 templates, and publishes it via the WebSub protocol.
Key design goals
Event-driven delivery -- registry changes are automatically propagated to subscribed external systems without polling.
Schema flexibility -- each partner system can receive data in its own format via configurable Jinja2 templates.
Asynchronous processing -- decouples the registry approval workflow from external delivery, so partner system availability does not block registry operations.
Horizontal scalability -- scales out by adding Celery workers.
Traceability -- every outgoing record is traceable from the originating change request through transformation to final publication.
Fault tolerance -- every stage supports retry up to a configurable
worker_max_attemptsthreshold before marking as failed.
Outgestion topics
An outgestion topic defines a delivery channel for a specific register and data model combination. Topics use the WebSub protocol to register with a hub and publish content to subscribers.
Each topic is uniquely constrained by (register_id, data_model_id). When a registry change matches a topic's register and data model, the outgestion pipeline picks it up for processing.
Topic lifecycle
Creation -- an administrator creates a topic via the staff portal API, specifying the register, data model, and WebSub topic URL.
Registration -- the topic registration beat producer detects the new topic (status:
PENDING) and dispatches a worker to register it with the WebSub hub.Active -- once registered (
PROCESSED), the topic is ready to receive and publish data.Deactivation -- topics can be toggled inactive, halting publication without deleting configuration.
Topic states
PENDING
Awaiting WebSub hub registration
PROCESSING
Registration request in flight
PROCESSED
Successfully registered with hub
FAILED
Registration failed after max attempts
Template-based transformation
Outgestion templates convert internal registry data into the external schema expected by partner systems. Templates are Jinja2 files stored in MinIO object storage.
Each template is bound to a unique (register_id, data_model_id) combination, matching the corresponding outgestion topic. When the transformation worker processes an outgoing record, it:
Retrieves the raw change payload from
outgoing_raw_data_payloads.Looks up the
OutgoingTemplatematching the record's register and data model.Fetches the Jinja2 template file from MinIO via
TemplateHelper.Optionally expands the data using JSON-LD expansion (for linked data interoperability).
Renders the template with the expanded data to produce the transformed payload.
Stores the result in
outgoing_transformed_data_payloads.
This approach allows the same registry data to be delivered in different schemas to different partners -- each partner's topic has its own template.
Event publishing flow
The outgestion pipeline processes data through three asynchronous stages, each driven by a Celery beat producer that polls for PENDING records and dispatches workers.
Stage 1: Topic registration
Beat producer: outgest_topic_register_beat_producer Worker: outgest_topic_register_worker
Registers newly created outgestion topics with the WebSub hub. This is a one-time operation per topic.
Beat producer queries
outgoing_topicsfor rows withwebsub_register_status = PENDING.Updates status to
PROCESSINGand dispatches a worker per topic.Worker calls
WebsubHelper.register_topic()which POSTshub.mode=registerto the WebSub hub.On success: status becomes
PROCESSED. On failure: retries up toworker_max_attempts, thenFAILED.
Stage 2: Data transformation
Beat producer: outgest_data_transformation_beat_producer Worker: outgest_data_transformation_worker
Transforms raw outgoing data into the partner's expected format.
Beat producer queries
outgoing_raw_datafor rows withtransformation_status = PENDING.Updates status to
PROCESSINGand dispatches a worker per record.Worker fetches the raw payload, locates the matching outgoing template, and renders via Jinja2.
Transformed output is stored in
outgoing_transformed_data_payloads.On success:
transformation_status = PROCESSEDandpublish_statusis set toPENDING, triggering Stage 3.
Stage 3: Data publishing
Beat producer: outgest_data_publish_beat_producer Worker: outgest_data_publish_worker
Publishes transformed data to the WebSub hub for delivery to subscribers.
Beat producer queries
outgoing_raw_datafor rows withpublish_status = PENDING.Updates status to
PROCESSINGand dispatches a worker per record.Worker fetches the transformed payload and calls
WebsubHelper.publish()which POSTs the content to the WebSub hub withhub.mode=publish.The hub then distributes the content to all subscribers of that topic.
On success:
publish_status = PROCESSED. On failure: retries up toworker_max_attempts, thenFAILED.
Flow diagram
Configuration models
outgoing_topics
Defines WebSub topic endpoints for each register and data model combination.
topic_id
UUID (PK)
Unique identifier
register_id
String (indexed)
Target register
data_model_id
String (indexed)
Target data model
websub_topic
String
WebSub topic URL
description
String
Human-readable description
is_active
Boolean
Whether topic accepts new data
websub_register_status
String
Registration state (PENDING, PROCESSING, PROCESSED, FAILED)
websub_register_datetime
DateTime
Last registration attempt timestamp
websub_register_number_of_attempts
Integer
Retry counter
websub_register_latest_error_code
String
Last error message
Unique constraint: (data_model_id, register_id)
outgoing_templates
Defines Jinja2 transformation templates for each register and data model combination.
template_id
UUID (PK)
Unique identifier
data_model_id
String (indexed)
Target data model
register_id
String (indexed)
Target register
template_file_id
String
Reference to Jinja2 template file in MinIO
Unique constraint: (data_model_id, register_id)
outgoing_raw_data
Tracks each outgoing record through the transformation and publishing stages.
outgest_id
String (PK)
Unique outgestion identifier
change_request_id
String (indexed)
Originating change request
internal_record_id
String (indexed)
Registry record identifier
register_id
String (indexed)
Register the record belongs to
data_model_id
String (indexed)
Data model used
topic_id
String (indexed)
Target outgestion topic
changed_by
String
User who made the change
changed_at
DateTime
When the change was made
approved_by
String
User who approved the change
approved_at
DateTime
When the change was approved
changed_by_partner_id
String (indexed)
Partner that originated the change (if any)
transformation_status
String (indexed)
Transformation stage state
publish_status
String (indexed)
Publishing stage state
outgoing_raw_data_payloads
Stores the original registry data before transformation.
change_request_id
String (PK)
Links to the change request
raw_data_json
JSONB
Raw data in JSON format
raw_data_xml
Text
Raw data in XML format
outgoing_transformed_data_payloads
Stores the template-transformed data ready for publishing.
change_request_id
String (PK)
Links to the change request
transformed_data_json
JSONB
Transformed data in JSON format
transformed_data_xml
Text
Transformed data in XML format
Error handling
Every pipeline stage tracks retry state independently:
Attempt counter -- incremented on each try.
Error code -- the latest error message is stored for debugging.
Status rollback -- on failure below max attempts, status resets to
PENDINGfor the beat producer to re-dispatch.Terminal failure -- after
worker_max_attempts(default: 5) the status moves toFAILEDfor manual inspection.
Each stage uses database transactions with rollback on error, ensuring no partial state is committed.
API endpoints
The outgestion configuration is managed through the staff portal API under the /outgestion-config prefix.
Topic management
POST /create_topic
outgestTopic:create
Create new outgestion topic
POST /get_all_topics
outgestTopic:view
List topics (paginated)
POST /get_topic
outgestTopic:view
Get single topic details
POST /update_topic
outgestTopic:edit
Update topic configuration
POST /toggle_topicstatus
outgestTopic:edit
Activate or deactivate topic
POST /re_register_topic
outgestTopic:edit
Re-trigger WebSub registration
POST /delete_topic
outgestTopic:delete
Delete inactive topic
Template management
POST /create_template
outgestTemplate:create
Create new transformation template
POST /get_template
outgestTemplate:view
Get single template details
POST /get_all_templates
outgestTemplate:view
List templates (paginated)
POST /update_template
outgestTemplate:edit
Update template configuration
POST /delete_template
outgestTemplate:delete
Delete template
Last updated
Was this helpful?