How Celery Tasks Handle Async Enricher Processing in Flowsint

Flowsint uses Celery's distributed task queue to execute enrichment jobs asynchronously, bridging synchronous Celery workers with async enricher implementations via asyncio.run() while persisting results to a Scan table indexed by the task UUID.

Flowsint leverages Celery to offload compute-intensive enrichment operations from the web request cycle. This architecture, implemented in the reconurge/flowsint repository, allows the platform to execute both Python-based and template-driven enrichers asynchronously while maintaining full traceability through a dedicated database model.

The Celery Task Architecture

Flowsint defines two primary Celery tasks in flowsint-core/src/flowsint_core/tasks/enricher.py to handle different enrichment paradigms. Both tasks follow a similar execution pattern but differ in how they instantiate the enricher logic.

Python-Based Enrichers

The run_enricher task executes enrichers that implement the Enricher interface defined in flowsint_core.core.enricher_base. When invoked, the task performs the following steps:

  1. Creates a new Scan database row using a UUID that matches the Celery task ID (self.request.id)
  2. Optionally instantiates a Vault service when an owner_id is provided
  3. Retrieves the enricher class from the global ENRICHER_REGISTRY
  4. Passes JSON-serialized input objects to the enricher via asyncio.run(enricher.execute(...))
  5. Updates the Scan status to COMPLETED and stores JSON-serializable results in scan.details on success
  6. Rolls back the transaction, marks the Scan as FAILED, and updates the Celery state to FAILURE on error

Reference: Task definition and execution flow—lines 26-73 of enricher.py.

Template-Driven Enrichers

The run_template_enricher task handles enrichers defined as YAML templates stored in the database. This task mirrors the Python-based flow but instantiates a TemplateEnricher class instead:

  • Loads the template configuration from the database
  • Creates a Scan record with the Celery task UUID
  • Resolves vault credentials if owner_id is present
  • Executes the template logic via asyncio.run(enricher.execute(...))
  • Persists results to the same Scan table structure

Reference: Template task implementation—lines 95-152 of enricher.py.

Bridging Synchronous Workers with Async Enrichers

Celery workers operate on synchronous Python code, while Flowsint enrichers implement asynchronous execution patterns for I/O-bound operations like HTTP requests.

The bridge occurs through the asyncio.run() method:


# From flowsint_core/tasks/enricher.py (lines 69-70, 45-46)

results = asyncio.run(enricher.execute(serialized_objects))

This pattern starts a temporary event loop within the synchronous Celery worker, executes the async execute method to completion, and returns the results to the worker process. The Enricher base class requires all implementations to define async def execute(self, objects) according to the interface in enricher_base.py.

Results are converted to JSON-serializable format using to_json_serializable() before storage, ensuring compatibility with Celery's JSON serializer and PostgreSQL JSON columns.

Celery Configuration and Time Limits

The Celery application is configured in flowsint-core/src/flowsint_core/core/celery.py:

celery = Celery(
    "flowsint",
    broker=settings.CELERY_BROKER_URL,
    backend=settings.CELERY_RESULT_BACKEND,
    include=[
        "flowsint_core.tasks.event",
        "flowsint_core.tasks.enricher",
        "flowsint_core.tasks.flow",
    ],
)

celery.conf.update(
    task_serializer="json",
    accept_content=["json"],
    result_serializer="json",
    timezone="UTC",
    enable_utc=True,
    task_track_started=True,
    task_time_limit=3600,          # 1 hour per task

    worker_max_tasks_per_child=1000,
    worker_prefetch_multiplier=4,
)

Key configuration parameters:

  • Broker and Backend: URLs are sourced from settings.CELERY_BROKER_URL and settings.CELERY_RESULT_BACKEND in the config module
  • Task Discovery: The include list ensures both enrichment tasks are auto-discovered by workers
  • Time Limits: task_time_limit=3600 prevents runaway enrichers from consuming resources indefinitely
  • Prefetch Control: worker_prefetch_multiplier=4 limits task pre-fetching, improving fairness when many long-running enrichers are queued

Error Handling and State Management

Both tasks implement comprehensive error handling through try/except/finally blocks:

  1. Database Transactions: On any exception, session.rollback() is invoked and the Scan record updates with status = EventLevel.FAILED and the error message
  2. Celery State: self.update_state(state=states.FAILURE) notifies the Celery monitor of task failure
  3. Audit Logging: Errors are printed and logged via the internal Logger for debugging and compliance purposes

This dual-layer state tracking ensures the frontend can poll either Celery's result backend or the Scan table to determine task status.

End-to-End Implementation Examples

Enqueuing Python Enrichers

To trigger a Python-based enrichment job:

from flowsint_core.tasks.enricher import run_enricher

# objects is a list of dicts from the API payload

task = run_enricher.delay(
    enricher_name="github_user_lookup",
    serialized_objects=objects,
    sketch_id="a1b2c3d4-e5f6-7g8h-9i0j-k1l2m3n4o5p6",
    owner_id="d4e5f6a7-b8c9-0d1e-f2g3-h4i5j6k7l8m9",
)

scan_id = task.id  # UUID shared between Celery and Scan table

Enqueuing Template Enrichers

For YAML template-based enrichments:

from flowsint_core.tasks.enricher import run_template_enricher

task = run_template_enricher.delay(
    template_name="whois_lookup",
    serialized_objects=objects,
    sketch_id=None,  # optional

    owner_id="123e4567-e89b-12d3-a456-426614174000",
)

scan_id = task.id

Retrieving Results

Since the Celery task ID matches the Scan table UUID, clients can poll for results:

from flowsint_core.core.models import Scan
from flowsint_core.core.postgre_db import SessionLocal

def get_scan_result(scan_id: str):
    with SessionLocal() as db:
        scan = db.query(Scan).filter(Scan.id == scan_id).first()
        if not scan:
            raise ValueError("Scan not found")
        return {
            "status": scan.status,
            "details": scan.details,
            "error": scan.error,
        }

The status field reflects Celery task states: PENDING, COMPLETED, or FAILED.

Summary

  • Flowsint uses two Celery tasks—run_enricher and run_template_enricher—to process enrichment jobs asynchronously outside the request-response cycle
  • The asyncio.run() bridge in flowsint_core/tasks/enricher.py allows synchronous Celery workers to execute async enricher implementations
  • Task status and results are persisted to a Scan table using the Celery task UUID as the primary key, enabling reliable polling mechanisms
  • Configuration in flowsint_core/core/celery.py enforces 1-hour time limits and JSON serialization for task safety
  • Comprehensive error handling updates both the database transaction state and Celery's internal failure state when enrichers encounter exceptions

Frequently Asked Questions

Why does Flowsint use asyncio.run() instead of native async Celery tasks?

Celery workers historically operate in synchronous contexts. Rather than implementing complex async worker pools, Flowsint uses asyncio.run() as a bridge to execute async enricher code within standard synchronous Celery tasks. This pattern (found in lines 69-70 of run_enricher) creates a temporary event loop, runs the async I/O operations to completion, and returns control to the Celery worker, simplifying deployment and compatibility.

How does Flowsint track the status of long-running enrichment tasks?

Flowsint implements a dual-tracking system using the Scan model in flowsint_core/core/models.py. When a task starts, it creates a Scan row with an ID matching self.request.id (the Celery task UUID). The task updates this record with COMPLETED or FAILED status upon completion. Clients can query the PostgreSQL Scan table directly using the task ID, avoiding the need to maintain Celery result backend connections for status checks.

What happens when an enricher fails during execution?

Both tasks wrap execution in try/except/finally blocks. On exception, the code calls session.rollback() to revert database changes, updates the Scan record with status=EventLevel.FAILED and the error message, and invokes self.update_state(state=states.FAILURE) to notify Celery's monitoring system. This ensures atomicity between the database state and Celery's internal task state while preserving error details for debugging.

Where is the Celery broker configuration defined in Flowsint?

The broker URL and result backend are defined in flowsint_core/core/config.py as CELERY_BROKER_URL and CELERY_RESULT_BACKEND. These settings are consumed by the Celery app instantiation in flowsint_core/core/celery.py. The configuration supports standard Celery broker protocols (Redis, RabbitMQ, etc.) via environment variables, allowing flexible deployment across different infrastructure environments.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →