How Flowsint's Enricher Architecture Processes Data: The Complete Pipeline

Flowsint's enricher architecture processes data through a strict, six-stage life-cycle defined in the Enricher abstract base class, orchestrating construction, parameter resolution, input validation, asynchronous enrichment, post-processing, and graph persistence.

Flowsint (reconurge/flowsint) treats every enrichment step as a self-contained Enricher that plugs into a unified pipeline. The architecture centers on the abstract base class located in flowsint-core/src/flowsint_core/core/enricher_base.py, which automates parameter validation, type coercion, and Neo4j graph operations. Developers implementing Flowsint's enricher architecture need only declare an InputType, an OutputType, and an asynchronous scan method; the framework manages orchestration and persistence.

The Six-Stage Life-Cycle of Flowsint's Enricher Architecture

The Enricher base class defines a rigid yet extensible contract. Concrete subclasses override specific hooks while the superclass handles validation and graph interaction.

Construction and Parameter Setup

When an enricher is instantiated, its constructor receives optional identifiers (sketch_id, scan_id), a params schema, an optional vault for secret retrieval, and an optional graph service. The base class stores these values and immediately builds a strict Pydantic model from the declared params_schema via build_params_model, as implemented in flowsint-core/src/flowsint_core/core/enricher_base.py.

Asynchronous Initialization (async_init)

Before any data touches the scan method, the enricher enters the async_init phase. This stage resolves all parameters through resolve_params, including secrets fetched from the injected vault. Once resolved, parameters are validated against the pre-built Pydantic model; any mismatch raises InvalidEnricherParams. This design guarantees that Flowsint's enricher architecture processes data only after the configuration is fully resolved and type-checked.

Pre-processing (preprocess)

The preprocess stage automatically validates incoming raw values against the enricher’s declared InputType. The base class uses Pydantic’s TypeAdapter to coerce raw strings into the primary field of the input model. Invalid items are skipped, and the enricher logs a warning if no valid inputs remain. This automated filtering ensures that the downstream scan method receives a clean list of validated objects.

Scanning (scan)

Subclasses implement the core enrichment logic inside the asynchronous scan method. This method accepts a list of validated InputType objects and must return a list of dictionaries or Pydantic models representing the enriched results. Because scan is async, Flowsint's enricher architecture natively supports I/O-bound workloads such as external DNS lookups, API calls, or scraping tasks without blocking the event loop.

Post-processing (postprocess)

After scan completes, the pipeline invokes postprocess. The default implementation in the base class returns raw results unchanged, but concrete enrichers may override this hook to merge metadata, deduplicate entries, or transform the output shape before final delivery.

Graph Interaction and Persistence

Throughout the life-cycle, enrichers persist data into Neo4j through helper methods injected via the GraphService. These helpers, which delegate to self._graph_service.flush() for batching, include:

  • create_node(node_obj) – Creates a node from a Pydantic object, automatically adding meta-properties such as type, sketch_id, and timestamps.
  • create_relationship(from_obj, to_obj, rel_label) – Creates a relationship between two existing nodes.
  • log_graph_message(message) – Writes custom log entries directly into the graph.

This graph interaction layer ensures that enrichment results are not merely returned but also stored in the graph database for downstream analysis.

The Execution Pipeline That Processes Enrichment Data

The public execute method in flowsint-core/src/flowsint_core/core/enricher_base.py orchestrates the entire flow. When called, it performs the following steps in order:

  1. Calls async_init to resolve and validate all parameters.
  2. Runs preprocess on the supplied raw values to coerce and filter inputs.
  3. Awaits scan to perform the actual enrichment work.
  4. Passes results through postprocess for any final transformations.
  5. Flushes pending graph operations via self._graph_service.flush().
  6. Logs start, completion, and error messages through the central Logger.

The method ultimately returns the processed list of enriched records to the caller. Concurrently, any side effects—such as newly created nodes and relationships—are committed to Neo4j via the flushed graph service.

Implementing a Concrete Enricher

Developers extend the base class by specifying InputType, OutputType, and the scan implementation. The following example from the Flowsint source demonstrates a domain-to-IP resolver:


# example_enricher.py

from flowsint_core.core.enricher_base import Enricher
from flowsint_types import Domain, Ip

class DomainToIpEnricher(Enricher):
    """Resolve a domain name to its IP addresses."""
    InputType = Domain          # expects a `domain` field

    OutputType = Ip             # will emit objects with an `ip` field

    @classmethod
    def name(cls):
        return "domain_to_ip"

    @classmethod
    def category(cls):
        return "Network"

    @classmethod
    def key(cls):
        return "domain"

    @classmethod
    def get_params_schema(cls):
        # No external secrets needed for this simple example

        return []

    async def scan(self, values: list[Domain]) -> list[dict]:
        results = []
        for domain in values:
            # Pretend we call a DNS resolver (omitted for brevity)

            ip_addr = "93.184.216.34"
            results.append({"ip": ip_addr, "source_domain": domain.domain})
        return results

In this implementation, Flowsint's enricher architecture handles all parameter and input validation automatically. The developer only writes the DNS resolution logic inside scan, relying on the base class for type checking and graph persistence.

Running and Interacting with Enrichers

Consumers retrieve enrichers by name from the global registry and invoke execute with raw inputs. The registry abstracts away instantiation and configuration, allowing a single method call to run the full enrichment pipeline:


# usage.py

from flowsint_enrichers import ENRICHER_REGISTRY

# Retrieve the enricher by name (registry auto-discovers it)

enricher = ENRICHER_REGISTRY.get_enricher("domain_to_ip", sketch_id="demo", scan_id="run1")

# Execute the enricher on raw input data

raw_domains = ["example.com", "python.org"]
enriched = await enricher.execute(raw_domains)

print(enriched)

# => [{'ip': '93.184.216.34', 'source_domain': 'example.com'},

#     {'ip': '93.184.216.34', 'source_domain': 'python.org'}]

Inside an enricher, graph persistence is handled through concise helper calls:


# graph interaction example (inside an enricher)

async def scan(self, values):
    for domain in values:
        ip = await resolve_dns(domain.domain)
        ip_obj = Ip(ip=ip)                 # Pydantic model

        self.create_node(ip_obj)           # Persist node

        self.create_relationship(domain, ip_obj, rel_label="RESOLVES_TO")

Registration and Auto-Discovery in the Enricher Ecosystem

Flowsint enables plug-and-play enrichment through automatic registration. The @flowsint_enricher decorator, defined in flowsint-enrichers/src/flowsint_enrichers/registry.py, adds each decorated class to the global ENRICHER_REGISTRY. This mechanism supports look-ups by name via ENRICHER_REGISTRY.get_enricher(name, sketch_id, scan_id) and bulk loading through load_all_enrichers(). A real-world reference implementation can be found in flowsint-enrichers/src/flowsint_enrichers/website/to_domain.py, which demonstrates how a concrete enricher is structured and registered.

Summary

  • Flowsint's enricher architecture processes data through a rigid life-cycle defined in flowsint-core/src/flowsint_core/core/enricher_base.py.
  • The Enricher base class automates construction, parameter resolution, input validation, graph persistence, and execution orchestration.
  • Developers implement only the asynchronous scan method, plus optional postprocess overrides, to define custom enrichment logic.
  • Inputs are coerced and validated against Pydantic InputType models during preprocess, while parameters are resolved and validated during async_init.
  • Results are written to Neo4j through create_node, create_relationship, and log_graph_message, with batch flushing handled by the injected GraphService.
  • The @flowsint_enricher decorator in flowsint-enrichers/src/flowsint_enrichers/registry.py enables automatic discovery and retrieval via ENRICHER_REGISTRY.

Frequently Asked Questions

What is the role of the Enricher abstract base class in Flowsint?

The Enricher abstract base class defines the uniform life-cycle that every enrichment step follows. Located in flowsint-core/src/flowsint_core/core/enricher_base.py, it enforces methods for initialization, pre-processing, scanning, and post-processing while providing concrete helpers for graph persistence and parameter validation.

How does Flowsint validate enricher inputs and parameters?

The framework validates inputs during preprocess by coercing raw values against the declared InputType using Pydantic’s TypeAdapter. Parameters are resolved and validated during async_init against a dynamically built Pydantic model derived from get_params_schema. Failures in parameter validation raise InvalidEnricherParams.

Can enrichers in Flowsint perform asynchronous I/O operations?

Yes. The scan method is explicitly asynchronous, allowing enrichers to await external API calls, DNS resolutions, or database queries without blocking the execution pipeline. The execute entry point manages the event loop coordination automatically.

How are enrichers discovered and instantiated at runtime?

Flowsint uses the @flowsint_enricher class decorator to populate the global ENRICHER_REGISTRY at import time. Runtime code can then retrieve instances by name using ENRICHER_REGISTRY.get_enricher(name, sketch_id, scan_id) or load all available enrichers in bulk via load_all_enrichers().

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 →