How to Debug Enricher Failures in Flowsint: A Complete Troubleshooting Guide
To debug enricher failures in Flowsint, trace errors through the five-phase pipeline (async_init → preprocess → scan → postprocess → execute) by checking logs for the sketch ID, validating parameters against vault secrets in flowsint-core, and inspecting the scan method where third-party API calls occur.
Enrichers in the reconurge/flowsint repository orchestrate third-party data collection through a standardized lifecycle defined in the abstract Enricher base class. When an enrichment run fails, the error typically surfaces in one of five distinct execution phases, each with specific logging signatures and distinct debugging strategies located in flowsint-core/src/flowsint_core/core/enricher_base.py.
Understanding the Enricher Lifecycle
Flowsint’s enrichers follow a strict execution order defined in EnricherBase.execute. Understanding this sequence is essential for pinpointing where failures originate.
Phase 1: Parameter Resolution with async_init()
The async_init() method resolves vault secrets and validates enricher parameters through EnricherBase.resolve_params. Failures here typically involve missing vault keys or validation errors raising InvalidEnricherParams.
According to the source code in enricher_base.py (lines 41-54), initialization errors surface immediately with clear vault key references.
Phase 2: Input Validation in preprocess()
The preprocess() method validates and normalizes raw input values using the Pydantic InputType. Invalid items are silently skipped, with warnings emitted via Logger.warn if no items survive validation (lines 4-11 in the base implementation).
Phase 3: Third-Party API Calls in scan()
The concrete scan() method implements the actual third-party integration (e.g., WHOXY, Etherscan, Maigret). This is where external API timeouts, authentication failures, and rate limiting most commonly occur.
Phase 4: Data Transformation in postprocess()
The optional postprocess() method transforms raw results before persistence. Logic errors here are rare but can involve mismatched data structures.
Phase 5: Execution Orchestration and Graph Writes in execute()
The execute() method ties all steps together, logging start/finish events and flushing pending Neo4j batch operations via self._graph_service.flush(). Unhandled exceptions from any previous phase are caught and logged with Logger.error (lines 39-44).
Common Failure Points and Diagnostic Logs
Each pipeline phase writes distinct log signatures that include the sketch ID, enabling precise log filtering for specific investigations.
-
Parameter Resolution Failures: Look for
InvalidEnricherParamsexceptions inasync_initoutput. Verify vault secrets exist inflowsint-coreVaultService or that defaults are provided in the configuration. -
Preprocessing Warnings: Check for
Logger.warnmessages indicating filtered items. If no items survive validation, the enricher emits a specific warning (lines 4-11). -
Scanning Exceptions: API failures appear as
Logger.errorentries caught in theexecutewrapper (lines 39-44). These include the full enricher name and exception traceback. -
Neo4j Batch Failures: Graph write errors surface during the flush operation at the end of
execute. Check Neo4j service logs or container logs for transaction failures.
Step-by-Step Debugging Procedure
-
Enable Detailed Logging: Flowsint writes logs to stdout/stderr. When running locally via
docker compose up, view real-time logs with:docker compose logs -f -
Identify the Failing Enricher: The error message includes the enricher name (
Enricher {self.name()} errored). List all registered enrichers using the registry utility:from flowsint_enrichers import ENRICHER_REGISTRY print([e["name"] for e in ENRICHER_REGISTRY.list()])The registry implementation resides in
flowsint-enrichers/src/flowsint_enrichers/registry.py. -
Check Parameter Resolution: Confirm vault secrets exist or provide defaults. Manually inspect resolved parameters after initialization:
enricher = MyEnricher(sketch_id="demo", scan_id="run1", params={}) await enricher.async_init() print(enricher.get_params()) # shows resolved secrets -
Inspect Preprocess Output: Add temporary debug prints inside
preprocessto identify filtered items:class MyEnricher(Enricher): def preprocess(self, values): cleaned = super().preprocess(values) print("DEBUG: validated", cleaned) # temporary inspection return cleaned -
Wrap the Scan Call: Surround third-party API calls with explicit error handling:
async def scan(self, data): try: response = await external_api.call(data) return response except Exception as exc: Logger.error(self.sketch_id, {"message": f"API failure: {exc}"}) raise -
Review Neo4j Batch Flushing: Force immediate graph writes to surface persistence errors earlier:
self.graph_service.flush() -
Run the Enricher in Isolation: Use the test suite in
flowsint-enrichers/tests/for fast feedback. Run specific enricher tests with:pytest -k MyEnricher flowsint-enrichers/tests/test_vault_integration.py
Practical Debugging Code Examples
Inspecting Resolved Parameters from Vault Secrets
This example demonstrates how to manually initialize an enricher and inspect resolved vault secrets before execution:
from flowsint_enrichers import ENRICHER_REGISTRY
# Example: debugging a domain-to-history enricher
enricher_cls = ENRICHER_REGISTRY.get_enricher(
"domain_to_history", sketch_id="demo", scan_id="run1"
)
await enricher_cls.async_init()
print("Resolved params:", enricher_cls.get_params())
Adding Debug Logging to Custom Enrichers
Inject granular logging into a custom enricher to trace data flow through each phase:
from flowsint_core.core.enricher_base import Enricher
from flowsint_core.core.logger import Logger
class DebugEnricher(Enricher):
"""Enricher that shows every step in the logs."""
InputType = Domain
OutputType = Domain
@classmethod
def name(cls):
return "debug_enricher"
@classmethod
def category(cls):
return "Domain"
@classmethod
def key(cls):
return "domain"
async def scan(self, data):
Logger.info(self.sketch_id, {"message": f"Scanning {len(data)} items"})
results = []
for item in data:
try:
# Replace with actual API call
results.append({"domain": item.domain, "status": "ok"})
except Exception as e:
Logger.error(self.sketch_id, {"message": f"API error: {e}"})
raise
return results
Running an Enricher in Isolation
Execute an enricher outside the UI to eliminate orchestration variables:
from flowsint_enrichers import ENRICHER_REGISTRY
enricher = ENRICHER_REGISTRY.get_enricher(
"website_to_links", sketch_id="demo", scan_id="run1"
)
# Raw input – list of website URLs
inputs = ["https://example.com", "https://invalid"]
results = await enricher.execute(inputs)
print("Enricher output:", results)
Key Source Files for Deep Debugging
-
flowsint-core/src/flowsint_core/core/enricher_base.py: Contains the core lifecycle implementation (async_init,execute, error handling) and theLoggerintegration. -
flowsint-enrichers/src/flowsint_enrichers/registry.py: ImplementsENRICHER_REGISTRYfor enricher discovery and instantiation. -
flowsint-enrichers/tests/test_vault_integration.py: Demonstrates successful and failed vault secret handling patterns. -
flowsint-enrichers/src/flowsint_enrichers/website/to_links.py: Real-world reference implementation showing properscanmethod structure. -
flowsint-core/src/flowsint_core/core/vault.py: Defines the Vault protocol supplying secrets duringasync_init.
Summary
- Trace the five-phase pipeline: Parameter resolution → preprocessing → scanning → postprocessing → graph writes.
- Filter logs by sketch ID: Every log message includes the sketch ID for precise correlation.
- Validate vault secrets early: Check
async_init(lines 41-54) forInvalidEnricherParamsbefore investigating downstream errors. - Isolate the enricher: Run standalone tests or manual scripts to eliminate UI and orchestration complexity.
- Force graph flushes: Call
self.graph_service.flush()to surface Neo4j write errors immediately rather than at pipeline end.
Frequently Asked Questions
How do I identify which enricher is failing in Flowsint?
The failing enricher name appears in the error message as Enricher {self.name()} errored. You can also list all available enrichers by importing ENRICHER_REGISTRY from flowsint_enrichers and calling .list(), which returns metadata for all registered enrichers including their names and categories.
What does a missing vault secret error look like in the logs?
Vault secret failures raise InvalidEnricherParams during the async_init phase (lines 41-54 in enricher_base.py). The error message includes the missing key name. Check flowsint-core VaultService to confirm the secret exists, or ensure your enricher provides a default value in the parameter schema.
Why are my input items being silently skipped during preprocessing?
Items that fail Pydantic InputType validation are filtered out in preprocess() without throwing exceptions. If all items are filtered, the method emits a Logger.warn message (lines 4-11). Add temporary debug prints inside your enricher's preprocess method to inspect the raw input versus validated output.
How can I debug Neo4j batch write failures?
Neo4j operations are batched and flushed at the end of execute(). To surface these errors earlier, manually call self.graph_service.flush() inside your scan or postprocess method. If writes fail, the error propagates to the execute-level exception handler (lines 39-44) and logs with the sketch ID. Check Neo4j container logs alongside the enricher logs for transaction-level details.
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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →