How to Implement Custom Enrichers in Flowsint: A Complete Developer Guide

To implement a custom enricher in Flowsint, subclass the Enricher base class from flowsint-core, apply the @flowsint_enricher decorator, declare InputType and OutputType, and implement the scan() and postprocess() methods.

Flowsint's enrichment architecture relies on modular units called enrichers that transform input entities into graph-connected output entities. Learning how to implement custom enrichers in Flowsint allows you to extend the platform with proprietary intelligence sources and external tools. This guide references the actual reconurge/flowsint source code, from the abstract base class in flowsint-core to production-ready implementations in the flowsint-enrichers package.

Understand the Enricher Base Class

Every custom enricher inherits from the abstract base class defined in flowsint-core/src/flowsint_core/core/enricher_base.py. The Enricher class enforces a strict contract through two core type attributes and two methods you must implement.

The base class provides:

  • InputType and OutputType — Pydantic model class attributes that declare what the enricher consumes and produces.
  • Automatic schema generation — input_schema() and output_schema() produce JSON schemas for API validation.
  • Execution pipeline — execute() orchestrates async_init() → preprocess() → scan() → postprocess() → graph flush.
  • Graph helpers — create_node() and create_relationship() let you write Neo4j entities directly from Pydantic objects.

You only need to implement scan() for the intelligence-gathering logic and postprocess() for graph construction.

Steps to Implement Custom Enrichers in Flowsint

Choose Input and Output Types

Select Pydantic models from the flowsint-types package. For example, use Domain as InputType and Ip as OutputType for a DNS resolution enricher. Set these as base models, not wrapped in List.

Create the Enricher Class

Place your file under flowsint-enrichers/src/flowsint_enrichers/<input>/. A typical path is flowsint-enrichers/src/flowsint_enrichers/domain/to_ip.py.

Import the required components and apply the registry decorator:

from flowsint_core.core.enricher_base import Enricher
from flowsint_enrichers.registry import flowsint_enricher
from flowsint_types import Domain, Ip

Define a class that inherits from Enricher, attach @flowsint_enricher, and implement the three metadata class methods: name(), category(), and key().

Implement the scan() Method

The scan() method receives List[InputType] and must return List[OutputType]. This is where you execute DNS lookups, call APIs, or invoke external tools.

async def scan(self, data: List[InputType]) -> List[OutputType]:
    results: List[OutputType] = []
    for domain in data:
        ip_address = socket.gethostbyname(domain.domain)
        results.append(Ip(address=ip_address))
    return results

Implement the postprocess() Method

The postprocess() method receives the output from scan() plus the original input data. Use it to persist nodes and relationships to the graph:

def postprocess(self, results: List[OutputType], input_data: List[InputType]) -> List[OutputType]:
    for domain, ip in zip(input_data, results):
        self.create_node(domain)
        self.create_node(ip)
        self.create_relationship(domain, ip, "RESOLVES_TO")
    return results

Add Optional Parameters

If your enricher needs configuration, define get_params_schema() as a class method. Access resolved values via self.params and vault secrets via self.get_secret() inside scan().

@classmethod
def get_params_schema(cls) -> List[Dict[str, Any]]:
    return [
        {"name": "mode", "type": "select", "required": True, "default": "passive",
         "options": [{"label": "Passive", "value": "passive"}]},
        {"name": "PDCP_API_KEY", "type": "vaultSecret", "description": "ProjectDiscovery key"}
    ]

Complete Example: Domain-to-IP Resolver

The following production-ready example, located at flowsint-enrichers/src/flowsint_enrichers/domain/to_ip.py, resolves domains to IP addresses using Python's socket library:

import socket
from typing import List
from flowsint_enrichers.registry import flowsint_enricher
from flowsint_core.core.enricher_base import Enricher
from flowsint_core.core.logger import Logger
from flowsint_types import Domain, Ip

@flowsint_enricher
class DomainToIpEnricher(Enricher):
    """Resolves domain names to their IP addresses using DNS."""
    InputType = Domain
    OutputType = Ip

    @classmethod
    def name(cls) -> str:
        return "domain_to_ip"

    @classmethod
    def category(cls) -> str:
        return "Domain"

    @classmethod
    def key(cls) -> str:
        return "domain"

    async def scan(self, data: List[InputType]) -> List[OutputType]:
        results: List[OutputType] = []
        for domain in data:
            try:
                ip_address = socket.gethostbyname(domain.domain)
                results.append(Ip(address=ip_address))
                Logger.info(self.sketch_id, {"message": f"Resolved {domain.domain} -> {ip_address}"})
            except socket.gaierror as e:
                Logger.warn(self.sketch_id, {"message": f"Could not resolve {domain.domain}: {e}"})
        return results

    def postprocess(self, results: List[OutputType], input_data: List[InputType]) -> List[OutputType]:
        for domain, ip in zip(input_data, results):
            self.create_node(domain)
            self.create_node(ip)
            self.create_relationship(domain, ip, "RESOLVES_TO")
        return results

InputType = DomainToIpEnricher.InputType
OutputType = DomainToIpEnricher.OutputType

Notice the exported aliases at the bottom. These make type imports convenient for downstream consumers.

Integrate External Tools

Many real-world enrichers invoke external tools rather than standard libraries. The SubfinderTool at tools/network/subfinder.py demonstrates this pattern. Inside scan(), instantiate the tool and process its output:

from tools.network.subfinder import SubfinderTool

async def scan(self, data: List[InputType]) -> List[OutputType]:
    results: List[OutputType] = []
    subfinder = SubfinderTool()
    for domain in data:
        subdomains = subfinder.launch(domain.domain)
        for sub in subdomains:
            results.append(Domain(domain=sub, root=False))
    return results

Keep tool imports inside the method or at module scope depending on startup cost, but always handle exceptions to prevent one failure from killing the entire batch.

Handle Vault Secrets and Parameters

The base class resolves parameters during async_init(), which runs before scan(). Never attempt to read vault secrets before this phase. Inside scan(), retrieve sensitive values safely:

async def scan(self, data: List[InputType]) -> List[OutputType]:
    mode = self.params.get("mode")
    api_key = self.get_secret("PDCP_API_KEY")
    # Pass api_key to your tool or API client

    return results

The strict Pydantic model built by resolve_params() uses extra="forbid". If you send undeclared fields in the request, validation fails immediately.

Testing and Auto-Discovery

After you implement custom enrichers in Flowsint, write tests in flowsint-enrichers/tests/ to validate metadata, type mappings, and scan logic. The existing test suite flowsint-enrichers/tests/test_domain_to_ip.py covers these exact checks.

Once your file lives under flowsint-enrichers/src/flowsint_enrichers/ and carries the @flowsint_enricher decorator, the API auto-discovers it via load_all_enrichers() on startup. Restart the server to register the new enricher.

Common Pitfalls When You Implement Custom Enrichers in Flowsint

  • Undefined type attributes. The base class raises NotImplementedError if InputType or OutputType are missing. Always assign them as base Pydantic models, never List[Domain].
  • Accessing params too early. Vault secrets resolve only after async_init(). Use self.get_secret() inside scan(), not during class construction.
  • Missing decorator. Without @flowsint_enricher, the registry ignores the class entirely.
  • Forgetting graph helpers. If postprocess() does not call self.create_node() and self.create_relationship(), no data persists to Neo4j.

Summary

  • Subclass Enricher from flowsint-core/src/flowsint_core/core/enricher_base.py and apply @flowsint_enricher.
  • Declare InputType and OutputType using Pydantic models from flowsint-types.
  • Implement scan() to perform enrichment logic and return List[OutputType].
  • Implement postprocess() to call self.create_node() and self.create_relationship().
  • Define get_params_schema() when you need configurable options or vault secrets.
  • Export InputType and OutputType aliases at the bottom of the file.
  • Restart the API to trigger auto-discovery via load_all_enrichers().

Frequently Asked Questions

What file path should I use for a new enricher?

Place the file under flowsint-enrichers/src/flowsint_enrichers/<input_category>/. For example, a domain resolver belongs at flowsint-enrichers/src/flowsint_enrichers/domain/to_ip.py. The directory structure maps to the input type and keeps the registry organized.

Why is my enricher not showing up in the API?

The class must inherit from Enricher, carry the @flowsint_enricher decorator, and reside inside the flowsint-enrichers/src/flowsint_enrichers/ tree. If both conditions are met, restart the API server so load_all_enrichers() can scan and register the module.

How do I access vault secrets inside an enricher?

Define the secret in get_params_schema() with "type": "vaultSecret". Then call self.get_secret("YOUR_KEY_NAME") inside scan(). The base class fetches the actual value during async_init(), so accessing it earlier returns an unresolved placeholder.

What is the difference between scan() and postprocess()?

scan() performs the actual intelligence gathering and returns a list of output entities. postprocess() receives those results alongside the original inputs and is responsible for persisting graph nodes and relationships via self.create_node() and self.create_relationship().

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 →