KGraph Crawler Architecture & Pipeline
Deep-dive into the Apache Tika 4.0 Pipes lifecycle, FetchEmitTuple routing, error resilience, and C4 component diagrams powering the RobOS KGraph Crawler.
Table of contents
- 1. Architectural Foundations: Apache Tika 4.0 Pipes
- 2. The
FetchEmitTupleLifecycle - 3. C4 Component Topology
- 4. Fault Tolerance & Crash Isolation
- 5. Next Steps
1. Architectural Foundations: Apache Tika 4.0 Pipes
The RobOS KGraph Crawler is built on top of Apache Tika 4.0 Pipes, a decoupled, high-performance architecture engineered for enterprise content extraction, asynchronous streaming, and fault-tolerant batch processing.
In traditional crawler setups, content discovery, data fetching, text/AST parsing, and graph indexing are tightly bound into monolithic processes. If an unexpected binary file crashes a parser or exhausts heap memory, the entire crawler dies, leaving the Knowledge Graph in an indeterminate state.
Tika 4.0 Pipes decouples these concerns into four distinct phases:
sequenceDiagram
autonumber
participant I as PipesIterator
participant F as Fetcher
participant G as tika-grpc Server
participant P as Polyglot Parsers
participant S as Schema Classifier
participant E as Emitter
participant K as Modular KGraph Store
I->>F: Emit FetchEmitTuple(fetchKey, emitKey, metadata)
F->>G: Stream bytes / payload over gRPC
G->>P: Parse AST, DDL, OpenAPI, text
P-->>G: Structured Metadata + Extracted AST
G-->>S: Return Parsed Output over gRPC
S->>S: Match existing SHACL shape or infer new package
S->>E: Hand off validated OSLC JSON-LD node
E->>E: Run W3C SHACL shape validation gate
E->>K: Write to .robos/kgraphs/<pkg>/package.jsonld
The 4 Phases of the Pipeline
- Iteration (
PipesIterator):- Traverses target resources (git workspaces, file directories, database catalogs, Kafka clusters, cloud API endpoints).
- Generates lightweight
FetchEmitTupleobjects containing routing coordinates without reading actual payload bytes.
- Fetching (
Fetcher):- Retrieves the payload stream for a given
fetchKey. - RobOS fetchers integrate with local UNIX
passGPG encrypted keystores (robos:hasCredential) to acquire connection credentials just-in-time, preventing secret leakage.
- Retrieves the payload stream for a given
- Parsing (
tika-grpc& Polyglot Parsers):- Evaluates content type (MIME detection), extracts structural syntax trees (ASTs), parses schema definitions (OpenAPI, Protobuf, SQL DDL), and analyzes metadata.
- Executes inside isolated worker processes managed by
tika-grpc.
- Emission (
Emitter):- Transforms extracted tokens into OSLC JSON-LD 1.1 entities.
- Evaluates target schema package (existing vs novel).
- Enforces W3C SHACL shape constraints.
- Persists valid nodes to
.robos/kgraphs/<pkg>/package.jsonldand updates.robos/kgraph.yaml.
2. The FetchEmitTuple Lifecycle
At the heart of the pipeline is the FetchEmitTuple. In RobOS, FetchEmitTuple is augmented with SDLC metadata:
{
"fetchKey": "file:///home/ndipiazza/source/robos/packages/order-service/openapi.yaml",
"emitKey": "urn:robos:service:order-service",
"metadata": {
"robos:repository": "github.com/nddipiazza/robos",
"robos:packageHint": "services",
"robos:fileMimeType": "application/x-yaml",
"robos:credentialUrn": "urn:robos:pass:devops/git/github",
"robos:discoveredAt": "2026-09-11T15:30:00Z"
},
"container": null,
"routingParams": {
"inferSchemaIfMissing": "true",
"validateSHACL": "true"
}
}
Tuple Lifecycle States
| State | Description | Transition Trigger |
|---|---|---|
ENQUEUED | Tuple generated by PipesIterator and enqueued into the in-memory or Kafka pipeline queue. | Iterator detects a new or updated asset. |
FETCHING | Fetcher connects to the data source and opens a payload stream. | Worker thread pulls tuple from queue. |
PARSING | Stream is dispatched to tika-grpc for AST and metadata extraction. | Payload stream successfully acquired. |
CLASSIFYING | Extracted schema is matched against existing SHACL shapes or flagged for novel package synthesis. | Tika parser returns structured metadata. |
EMITTING | OSLC JSON-LD node is validated against SHACL constraints and written to disk. | Classification completed with zero fatal schema errors. |
COMPLETED | Node committed to Git-backed package store; aggregated graph synchronized. | Disk I/O succeeds and index updated. |
DEAD_LETTER | Tuple failed after max retry attempts (e.g., auth failure, unparseable corrupt payload). | Max retries exceeded; routed to DLQ. |
3. C4 Component Topology
flowchart TB
subgraph HostSystem["RobOS Host System / Ephemeral Sandbox"]
subgraph CrawlerEngine["RobOS KGraph Crawler Harness"]
Orchestrator["Crawl Orchestrator<br/>(Task Runner / Agent Scheduler)"]
Queue["Tuple Work Queue<br/>(Non-blocking bounded buffer)"]
WorkerPool["Concurrent Worker Pool<br/>(Configurable thread count)"]
DLQ["Dead Letter Queue<br/>(~/.config/robos/crawler/dlq.json)"]
end
subgraph PipeSubsystem["RobOS Pipes Modules"]
Iterators["RobOS Custom Iterators<br/>(Workspace, DataSource, KGraph)"]
Fetchers["RobOS Custom Fetchers<br/>(Workspace, Database, Cloud)"]
Emitters["RobOS Custom Emitters<br/>(KGraph, SchemaPackage, LivingDocs)"]
end
subgraph SecurityBoundary["Security & Credential Boundary"]
PassStore["UNIX GPG pass Store<br/>(~/.password-store/)"]
PassBridge["PassCredentialBridge<br/>(GPG Decryption & Session Cache)"]
end
subgraph TikaDaemon["tika-grpc Process Isolation"]
GrpcServer["tika-grpc Server<br/>(:50051)"]
ParserPool["Tika Polyglot Parsers<br/>(AST, SQL, OpenAPI, PDF, Code)"]
end
subgraph KGraphStorage["SDLC Knowledge Graph Store"]
PackagesDir[".robos/kgraphs/<pkg>/package.jsonld"]
AggregatedDoc[".robos/knowledge-graph.jsonld"]
PkgIndex[".robos/kgraph.yaml"]
end
end
Orchestrator --> Iterators
Iterators --> Queue
Queue --> WorkerPool
WorkerPool --> Fetchers
Fetchers --> PassBridge
PassBridge --> PassStore
Fetchers --> GrpcServer
GrpcServer --> ParserPool
ParserPool --> GrpcServer
GrpcServer --> WorkerPool
WorkerPool --> Emitters
Emitters --> PackagesDir
Emitters --> AggregatedDoc
Emitters --> PkgIndex
WorkerPool -- "On Failure" --> DLQ
4. Fault Tolerance & Crash Isolation
Content extraction engines frequently process untrusted, corrupted, or adversarially crafted files:
- 500 MB malformed SQL dump files
- Circular JSON-LD or YAML references (Billion Laughs attack)
- Deeply nested ZIP or JAR archives (Zip bombs)
- Memory-leaking binary parsers
The Three-Layer Isolation Boundary
- Process Isolation (
tika-grpc):- The parser runtime executes in a separate process or Docker container from the RobOS Crawler Harness.
- If a parser encounters a fatal JVM
OutOfMemoryErroror segmentation fault, only that worker thread/process terminates. - The
tika-grpcdaemon supervisor automatically replaces the worker and returns a structured gRPC error status (DEADLINE_EXCEEDEDorRESOURCE_EXHAUSTED) to the crawler.
- Circuit Breaking & Dead-Letter Queue (DLQ):
- If a specific
FetchEmitTuplefails repeatedly, it is marked asDEAD_LETTERand persisted to~/.config/robos/crawler/dlq.json. - The crawl continues uninterrupted across the remaining thousands of files.
- If a specific
- Atomic File Emission:
- The
RobOSKGraphEmitterwrites to temporary shadow files (package.jsonld.tmp) before renaming, ensuring that a power cut or system kill signal never corrupts the underlying Git repository.
- The
5. Next Steps
- Explore Custom Iterators & Fetchers to configure data source adapters.
- Learn how the crawler synthesizes new schema packages in Emitters & Schema Inference.
- Inspect the tika-grpc Streaming Service.