# How Celery Tasks Handle Async Enricher Processing in Flowsint

> Discover how Celery tasks manage async enricher processing in Flowsint. Learn about bridging sync workers with async enrichers and persisting results to the Scan table.

- Repository: [reconurge/flowsint](https://github.com/reconurge/flowsint)
- Tags: internals
- Published: 2026-06-03

---

**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`](https://github.com/reconurge/flowsint/blob/main/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`](https://github.com/reconurge/flowsint/blob/main/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`](https://github.com/reconurge/flowsint/blob/main/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:

```python

# 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`](https://github.com/reconurge/flowsint/blob/main/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`](https://github.com/reconurge/flowsint/blob/main/flowsint-core/src/flowsint_core/core/celery.py):

```python
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:

```python
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:

```python
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:

```python
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`](https://github.com/reconurge/flowsint/blob/main/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`](https://github.com/reconurge/flowsint/blob/main/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`](https://github.com/reconurge/flowsint/blob/main/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`](https://github.com/reconurge/flowsint/blob/main/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`](https://github.com/reconurge/flowsint/blob/main/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.