SignalDB Flight Communication Design¶
1. Introduction¶
This document outlines the design for Apache Arrow Flight as the primary communication mechanism in SignalDB, both for inter-service communication and external client access. The design leverages the performance benefits of Arrow Flight while maintaining compatibility with the current architecture.
Current Implementation Status: ✅ Complete - Full Flight communication with WAL integration is implemented and production-ready. All integration tests passing.
2. Background¶
2.1 Current Architecture¶
SignalDB currently uses:
- Apache Arrow Flight as the primary inter-service communication mechanism ✅ Implemented
- OTLP data received via gRPC and HTTP at the Acceptor
- Direct Flight communication between Acceptor and Writer
- Direct Flight communication between Router and Querier
- Object storage integration for Parquet persistence
Current architecture with WAL integration:
External Clients
│
▼ (OTLP/gRPC)
┌─────────────┐ ┌──────┐ Flight ┌─────────────┐ ┌─────────────┐
│ Acceptor │───▶│ WAL │────────▶│ Writer │───▶│ Storage │
└─────────────┘ └──────┘ └─────────────┘ └─────────────┘
(OTLP) (Disk) (Flight) (Parquet)
│
▼
┌─────────────┐ HTTP ┌─────────────┐ Flight ┌─────────────┐
│ Clients │◀─────────│ Router │───────▶│ Querier │
└─────────────┘ (Tempo) └─────────────┘ (DataFusion) └─────────────┘
What's Working (✅ Complete):
- OTLP clients send telemetry data to the Acceptor
- Acceptor writes data to WAL for durability, then converts OTLP to Arrow format
- Acceptor forwards data to Writer via Flight (with Storage capability routing)
- Writer receives Arrow data and persists to Parquet storage
- Writer marks WAL entries as processed after successful storage
- Router exposes HTTP endpoints (Tempo API) and forwards queries via Flight
- Querier executes DataFusion queries against Parquet storage
- All services discover each other via catalog-based service registry
2.2 Apache Arrow Flight¶
Apache Arrow Flight is a high-performance client-server framework designed for efficient transfer of large datasets over network interfaces.
Key benefits include:
- Native Arrow format transfer (no serialization/deserialization overhead)
- High throughput, low latency data transfer
- Streaming capabilities
- Built on gRPC with authentication and encryption support
3. Design Goals¶
- ✅ Achieved: Improve performance of data transfer between components
- ✅ Achieved: Provide a high-performance query interface for external clients
- ✅ Achieved: Maintain logical separation of components while supporting monolithic deployment
- ✅ Achieved: Eliminate the need for a separate message bus
- ✅ Achieved: Support both in-process and networked communication with the same code
- ✅ Achieved: Implement WAL-based durability with automatic recovery
- ✅ Achieved: Provide capability-based service discovery and routing
4. Flight Integration Design¶
4.1 Current Implementation¶
The current architecture uses Flight as the primary data transfer mechanism:
External Clients
│
▼
┌─────────────┐ ┌─────────────┐ ┌─────────────┐
│ Acceptor │───▶│ Writer │───▶│ Storage │
└─────────────┘ └─────────────┘ └─────────────┘
│
▼
┌─────────────┐
│ Clients │◀─────────────────────│ Router │───▶│ Querier │
└─────────────┘ └─────────────┘ └─────────────┘
All data-intensive communication between components uses Flight.
4.2 Component Flight Services ✅ Implemented¶
Each component implements a Flight service:
- Acceptor: no Flight server; acts as a Flight client forwarding data to the Writer
- IcebergWriterFlightService: Receives data from Acceptor and writes to Iceberg tables
- QuerierFlightService: Executes queries against storage and returns results
- SignalDBFlightService (Router): Exposes HTTP API and forwards requests to Querier via Flight. It also writes uploaded eval results (
POST /api/v1/evals/results) to a Writer as logs batches, through the samecommon::flight::forward::forward_batch_to_writerDoPutthe Acceptor uses; unlike the Acceptor it keeps no WAL, so the upload is durable once the Writer acks - CompactorFlightService: Admin-only
DoActioninterface for compaction management
4.3 External Flight Interface¶
The Router exposes Flight capabilities via HTTP endpoints, providing:
- Query execution via Tempo-compatible API
- Trace retrieval and search functionality
- Administrative operations
4.4 Supported Flight RPC Methods¶
The following table shows the Flight RPC methods supported by each service:
| Method | Router | Querier | Writer | Compactor | Description |
|---|---|---|---|---|---|
Handshake |
✅ | ✅ | ✅ | ❌ | Protocol version exchange |
ListFlights |
✅ | ✅ | Empty stream | ❌ | List available query types |
GetFlightInfo |
✅ | ❌ | ❌ | ❌ | Get metadata for a query |
GetSchema |
✅ | ✅ | ❌ | ❌ | Get schema for a query type |
DoGet |
✅ | ✅ | ❌ | ❌ | Execute query and stream results |
DoPut |
❌ | ❌ | ✅ | ❌ | Write data to storage |
DoExchange |
❌ | ❌ | ❌ | ❌ | Not implemented |
DoAction |
❌ | ❌ | ❌ | ✅ | Admin commands (compactor only) |
ListActions |
Empty stream | Empty stream | Empty stream | ✅ | List admin commands |
Legend: ✅ implemented, ❌ returns unimplemented, "Empty stream" succeeds but yields nothing (a no-op, not an error).
Note: The Router is the primary client-facing Flight interface. Clients typically connect to the Router for all query operations. The Compactor's Flight service (src/compactor/src/flight.rs) is an admin interface only: do_action supports compact_now, compact_status, and compact_dry_run; every other RPC returns unimplemented.
Ticket and Command Grammar¶
There are two layers with different grammars -- the Router's descriptor commands and the Querier's do_get tickets:
Router (src/router/src/endpoints/flight.rs) -- get_flight_info, get_schema, and list_flights recognize these FlightDescriptor cmd values:
| Command | Description | Notes |
|---|---|---|
traces |
Trace/span schema and metadata | do_get currently returns an empty stream (placeholder) |
trace_by_id?id={id} |
Single-trace schema and metadata | do_get currently returns an empty stream (placeholder) |
logs |
Log schema and metadata | do_get currently returns an empty stream (placeholder) |
metrics |
Metric schema and metadata | do_get currently returns an empty stream (placeholder) |
Any other do_get ticket (including find_trace:..., search_traces:..., and raw SQL) is proxied verbatim to a Querier discovered via the QueryExecution capability, with request metadata forwarded.
Querier (parse_ticket in src/querier/src/flight.rs) -- do_get tickets use this grammar:
| Ticket | Description |
|---|---|
find_trace:{tenant_slug}:{dataset_slug}:{trace_id}[:{start}:{end}] |
Single trace lookup; the optional trailing segments are unix-second time hints (either may be empty) that prune the scanned range. Routers only append them when a hint is present, so the 3-part form remains valid. A missing trace yields a Flight not_found status, not an empty stream |
search_traces:{tenant_slug}:{dataset_slug}:{params_json} |
Trace search (SearchQueryParams as JSON; unknown fields are ignored on deserialization) |
trace_tags:{tenant_slug}:{dataset_slug}:{params_json} |
Trace tag-name discovery (TraceTagsParams as JSON: nanosecond start/end, optional scope). Returns keys observed in the window (resource/span) plus the fixed intrinsics, grouped by scope (#1073) |
trace_tag_values:{tenant_slug}:{dataset_slug}:{tag}:{params_json} |
Distinct values of one (unscoped) trace tag in the window (TraceTagValuesParams as JSON: nanosecond start/end). status/kind return their static enum; an unknown tag returns an empty list, never an error (#1073) |
query_logs:{tenant_slug}:{dataset_slug}:{params_json} |
LogQL log query (LogQueryParams as JSON: LogQL string, nanosecond start/end, limit, direction). Returns the projected log columns ordered by timestamp |
query_logs_labels:{tenant_slug}:{dataset_slug}:{start}:{end} |
Log label names in the nanosecond window |
query_logs_label_values:{tenant_slug}:{dataset_slug}:{label}:{start}:{end} |
Distinct values of one log label in the window |
query_logs_series:{tenant_slug}:{dataset_slug}:{params_json} |
Series (label sets) matching a stream selector (LogSeriesParams as JSON) |
query_logs_detected_fields:{tenant_slug}:{dataset_slug}:{params_json} |
Attribute-field discovery: sampled keys with inferred type and approximate cardinality (DetectedFieldsParams as JSON) |
query_metric:{tenant_slug}:{dataset_slug}:{params_json} |
LogQL metric query (MetricQueryParams as JSON: LogQL string, nanosecond start/end, step). Returns a matrix bucketed by date_bin(step) |
query_promql:{tenant_slug}:{dataset_slug}:{params_json} |
PromQL query (PromQlQueryParams as JSON: PromQL string, nanosecond start/end, step). Returns a matrix over the metrics tables |
query_ir:{tenant_slug}:{dataset_slug}:{params_json} |
Native Query IR (IrQueryParams as JSON: a versioned IR document plus the server-stamped now_ns for deterministic relative-time resolution). The querier validates and lowers the single-signal IR to a DataFusion plan; returns the declared rows/series/table envelope, or a one-row flamegraph_json/graph_json batch for the flamegraph/graph envelopes |
| anything else | Treated as a raw SQL query executed via DataFusion |
Whichever ticket a query arrives on, it executes in a session built by
querier::session_config_from from [querier.datafusion]. Those options are
not merely a tuning surface: split_file_groups_by_statistics is what allows an
ordered scan over files that attest the table's
declared sort order to skip a sort it
would otherwise have to perform. Whether that ordering survives from the scan to
the physical plan depends on the options and optimizer rules actually in force,
so the function is public and the ordering tests plan against it rather than
against a session of their own — a rule that quietly dropped the ordering would
break no result, it would only make queries slow again. The same section also
carries the query's scan shape (batch_size, target_partitions,
sort_spill_reservation_mb) — the shared
common::datafusion_runtime::ScanShape the compactor applies too, so a
sort's unspillable per-batch reservation stays inside whatever memory pool
the querier is running under (#1359).
The standalone querier binary additionally serves Tempo's tempopb.Querier
gRPC protocol on the same port as Flight (see the
Tempo API reference);
that protocol does not use tickets.
There is no trace_by_id?id=... ticket form at the Querier; that command exists only in the Router's metadata path. The Tempo HTTP endpoints bypass the Router's Flight commands entirely and send find_trace:/search_traces:/SQL tickets straight to the Querier.
Self-Monitoring Anti-Loop Guard¶
When self-monitoring is enabled, both Flight handlers apply the anti-loop
guard from src/common/src/self_monitoring/suppress.rs: requests that touch
the reserved _system tenant are processed with OpenTelemetry export
suppressed, so handling SignalDB's own telemetry does not generate more of
it. The Writer's do_put suppresses per batch based on the tenant in the
Flight metadata (its background WAL loop does the same per WAL entry), and
the Querier's do_get resolves the tenant — authenticated caller, the
op:{tenant_slug}:... ticket segment, or the x-tenant-id header for raw
SQL — before creating its processing span. The suppression marker is a tokio
task-local: it crosses neither the Flight hop between services nor
tokio::spawn, which is why each handler carries its own call site.
Boundary Spans (RPC Semantic Conventions)¶
Every Flight handler roots its request in a semconv RPC SERVER span built by
the factories in src/common/src/self_monitoring/spans.rs — the single
sanctioned construction path for boundary spans. Flight has no semantic
convention of its own, so it is modeled as plain gRPC: spans are named by the
fully-qualified logical method, disambiguated by a low-cardinality detail
segment where one exists —
arrow.flight.protocol.FlightService/DoGet query_ir (Querier, ticket verb),
…/DoPut (Writer), …/DoAction flush (Writer, action type),
…/DoAction compact_dry_run (Compactor, action type) —
and carry rpc.system.name = grpc, rpc.method, and the string
rpc.response.status_code. Status mapping follows the RPC semconv asymmetry:
a server span is marked failed only for server-fault gRPC codes (UNKNOWN,
DEADLINE_EXCEEDED, UNIMPLEMENTED, INTERNAL, UNAVAILABLE, DATA_LOSS);
codes like NOT_FOUND are the caller's problem and leave the span status
unset. Raw-SQL tickets have no op: prefix, so their first :-segment is
query text — the Querier only appends the verb when it matches a short
lowercase identifier, keeping SQL out of span names.
Trace Context Propagation¶
W3C trace context (traceparent/tracestate) is propagated across SignalDB
service boundaries by src/common/src/flight/trace_context.rs, so a single
distributed trace can span the acceptor, writer, router, and querier. Every
function routes through the global OpenTelemetry text-map propagator, which is
a no-op unless self-monitoring is enabled. A parent must be adopted before
its span is first entered; span links, by contrast, may be added at any time.
Four carriers move the context, matching how each path already exchanges metadata:
| Carrier | Path | Direction |
|---|---|---|
JSON app_metadata on the first FlightData message |
Acceptor → Writer do_put |
inject / extract |
| gRPC request metadata headers | Router → Querier do_get |
inject / extract |
| HTTP request headers | external caller → Router query APIs | extract (server side) |
| Span links | WAL batch fan-in (background processor) | link |
The same do_put app_metadata JSON also carries ingest_id, a content
fingerprint of the batch: an xxh3-128 hash of tenant, dataset, WAL operation
and the Arrow IPC bytes, formatted as a uuid (retry_dedup::batch_fingerprint
in the acceptor). The handler stores it in the acceptor WAL entry's metadata
next to the routing fields, so the hot path and the WAL retry consumer forward
the entry under the same id. Two entries holding byte-identical batches share
an ingest_id, which is what lets the writer drop a client's resend (an OTLP
exporter retrying after its own timeout) even when it lands in a different
acceptor WAL entry, at another acceptor replica, or after an acceptor restart.
Acceptor WAL entries written before the field existed carry no stored id, and
the retry consumer forwards those under their WAL entry id.
The acceptor picks the destination writer for a do_put by rendezvous
(highest random weight) hashing on ingest_id instead of round-robin, so every
copy of a batch reaches the same writer as long as the writer set is
unchanged; see InMemoryFlightTransport::get_client_for_capability_keyed. That
writer remembers each id for [writer].ingest_dedup_window (default 1h) and
rebuilds the cache from its own WAL at startup; a repeat is acked and its WAL
entries marked processed, so it is never committed. ingest_id is absent for
acceptors that predate the field, which fall back to non-deduped behavior; a
present-but-unparseable id is rejected with invalid_argument.
Each acceptor also keeps a per-process cache of the fingerprints it flushed in
the last [acceptor].retry_dedup_window (default 5m). It is a cheap first
line: a resend that returns to the same acceptor is acked and retired without
a Flight round trip or a second writer WAL append. 0s disables that cache
only; the writer's dedup still applies. Both caches live in
common::ingest_dedup.
Write path. At do_put the Writer records the active span's context into
the WAL entry metadata alongside the routing fields. Because the background
WalProcessor commits a batch that fans in entries from many independent
ingest requests, its span cannot adopt a single parent — it reads each entry's
stored context and adds one span link per distinct ingest trace, keeping
every source trace reachable from the batch span instead of leaving it a
detached root.
Read path. An http_trace_context_middleware at the Router's HTTP boundary
roots each request in a server span whose parent is the caller-supplied
traceparent, so an external client that propagates trace context sees
SignalDB's query trace join theirs. Downstream #[instrument] handler spans
become children of that span, and each Router → Querier Flight call runs
inside a semconv RPC CLIENT span (do_get_client_span) whose context is
injected into the request metadata — so the querier's SERVER span is the
client span's child and the trace reads SERVER → CLIENT → SERVER. The
middleware mirrors the anti-loop guard above: _system tenant requests bypass
the span so self-monitoring queries are not re-instrumented and re-ingested.
Response direction. The same middleware returns the server span's context
to the caller on every response — Server-Timing: traceparent;desc="..."
plus the W3C traceresponse header, formatted by
trace_context::format_traceparent (which refuses invalid all-zero contexts,
so nothing is emitted when self-monitoring is disabled). See
Trace Context on HTTP Responses for the
caller-facing contract.
Error Recording on Query Spans¶
A failing query is only useful in a trace if the reason survives. By the time a
do_get error reaches the caller it has been flattened into a transport
Status that the Router strips down to a bare HTTP code, so the querier records
the cause where it still exists: the whole do_get body runs inside a single
error boundary that, on any Err, calls
common::self_monitoring::record_span_exception to attach an OpenTelemetry
exception event (exception.message) and an error status to the
…FlightService/DoGet server span. Because the boundary wraps the entire request, every
failure path — ticket parsing, cross-tenant rejection, query execution, and
result conversion — is captured, not just execution errors. The helper is a
no-op when self-monitoring is disabled (Span::current() is the disabled span),
so this costs nothing on the hot path.
5. Implementation Details¶
5.1 Current Data Flow ✅ Working¶
Trace Ingestion Flow:¶
sequenceDiagram
participant C as OTLP client
participant A as Acceptor
participant W as Writer (Storage capability)
participant O as Object store
C->>A: OTLP traces (gRPC)
A->>A: convert to Arrow (otlp_traces_to_arrow)
A->>A: append to Acceptor WAL + flush
A->>W: Flight DoPut (Arrow batches)
W->>W: transform v1 to physical-v4, append to Writer WAL
W-->>A: confirm
A->>A: mark WAL entry processed
A-->>C: acknowledge
Note over W,O: asynchronous, WalProcessor 5s loop
W->>O: Iceberg commit (Parquet files)
- Acceptor receives OTLP trace data via gRPC
- Acceptor converts OTLP to Arrow format using
otlp_traces_to_arrow - Acceptor appends the batch to its WAL and flushes for durability
- Acceptor uses Flight
DoPutto send Arrow data to Writer (Storage capability) - Writer transforms to the physical-v4 intermediate storage shape (function name says "v2" for historical reasons; it resolves a fixed
"physical-v4"literal, same as the logs transform's"physical-v3"— notSCHEMA_DEFINITIONS.current_trace_version(), nowphysical-v5, the typed attribute layout: this transform's only job is bridging the wire's v1 shape to the last pre-typed version, and the typed-container splitting that carries a batch from there to whatever the table's actual current schema is happens generically afterward, inIcebergTableWriter::append_batches_with_marker) and appends to its own WAL — the WAL of the batch's own tenant/dataset/signal, one instance per combination, so a poisoned segment or a slow flush on the append path does not block another tenant'sdo_put(the background drain over those WALs commits groups concurrently, so a tenant whose Iceberg round trip is slow no longer delays another tenant's commit in the same cycle). This transform resolves a materialization plan once per schema version (compiled-schema-materializer) rather than dispatching per field per batch — see theflight-schemasskill. - Writer confirms after its WAL flush (it does not block the confirm on the Iceberg commit); Acceptor marks its WAL entry as processed
- Writer's
WalProcessorasynchronously commits WAL entries to Iceberg (Parquet in the object store), coalescing pending entries per(tenant, dataset, table)— a group commits when[writer].commit_intervalelapses or its rows reach[writer].max_uncommitted_rows. This caps the Iceberg snapshot / catalog-metadata write rate independent of ingest rate.
Steps 5 and 7 route the same batch at different times — once to pick its WAL, once to pick its Iceberg table — so both call the writer's single routing::route. It trims the metadata tenant/dataset ids, substitutes the deployment default for a blank or absent one, and validates the result; the commit path falls back to the ids of the WAL the entry lives in, which are what routing already returned on ingest. Before this was shared (#1319), a padded or empty metadata tenant landed in one WAL and committed under a different Iceberg tenant, so a row's destination depended on whether it was committed live or replayed after a restart. Signal-to-table mapping lives in the same function: one table per signal, except metrics, which honour the metadata's target_table and otherwise fall back to metrics (the wide table; a pre-cutover deployment falls back to the legacy metrics_gauge).
Because the commit is asynchronous, ingested data is queryable only once committed (bounded by commit_interval). A caller needing read-your-writes forces an immediate commit with the Writer Flight do_action("flush") (advertised via list_actions). The action is tenant-scoped: the scope is taken from the request's x-tenant-id (required) and x-dataset-id (optional) gRPC metadata — the same tenant identity the ingest path carries — and it force-commits only that tenant's (optionally that dataset's) pending groups. A request without x-tenant-id is rejected, so a caller can neither flush every tenant nor a tenant it names only in the payload. Tests use common::testing::flush_storage_writers(transport, tenant, dataset) for a deterministic barrier.
Query Flow:¶
sequenceDiagram
participant C as HTTP client
participant R as Router
participant Q as Querier
participant O as Object store
C->>R: Tempo API query (HTTP)
R->>Q: Flight do_get ticket
Q->>O: DataFusion scan (Parquet)
Q-->>R: Arrow RecordBatch stream
R-->>C: JSON response
- Client sends HTTP query to Router
- Router forwards query to Querier via Flight
- Querier executes query using DataFusion against Parquet files
- Results streamed back to client via Flight → HTTP
Parquet footer caching. Step 3 is dominated by opening candidate Parquet
files rather than by reading them: pruning (partition, statistics, bloom)
already reduces a point lookup to a handful of bytes, while every candidate
file's footer still has to be fetched and decoded, one object-store round-trip
each. The querier's session therefore installs
common::parquet_metadata_cache::CacheParquetMetadata, a physical optimizer
rule that makes every Parquet scan read footers through the shared
RuntimeEnv's metadata cache. DataFusion only wires that cache up on its own
ListingTable path, so scans from the Iceberg table provider would otherwise
re-read every footer on every query. The cache is sized by
[querier].parquet_metadata_cache_mb (0 disables it), lives on the one
RuntimeEnv shared by every per-request session, and is safe because Iceberg
data files are immutable — DataFusion still revalidates each hit against the
current object metadata.
5.2 Schema Design ✅ Implemented¶
Flight schemas are defined in src/common/src/flight/schema.rs with conversions for:
- OTLP traces → Arrow schema
- OTLP metrics → Arrow schema
- OTLP logs → Arrow schema
A resource without service.name is stored as unknown
(common::flight::conversion::UNKNOWN_SERVICE_NAME) on every signal — the
acceptor's OTLP conversion substitutes it for traces and logs (their v1
batches carry a service_name column), and the writer's metrics transforms
substitute it when they re-derive service_name from resource_json —
because the Iceberg service_name column is non-nullable and a missing
attribute must not dead-letter the batch.
Metric values that JSON cannot carry — NaN (Prometheus's staleness marker,
0/0 rates) and ±Inf — travel in the wire data_json as the strings "NaN",
"+Inf", "-Inf" (common::flight::conversion::f64_to_json /
json_to_f64), never as null; the writer maps them back to the same
non-finite doubles, and a data point with no value at all lands as NaN. This
keeps the non-nullable metrics.value column satisfiable for a gauge/sum
row, so one such point can no longer make the writer reject a whole
batch and pin its WAL entry forever (#1061). Histogram explicit_bounds keep
a +Inf bound for the same reason.
The reverse direction, OTLP → Prometheus series
(conversion_prometheus::from_otel), downsamples exponential histograms to
classic _bucket/_count/_sum series: bucket i of scale s gets upper
bound base^(i+1) with base = 2^(2^-s) (negative buckets -base^i), and
the zero bucket folds into the lowest bound (#748).
A related isolation applies one step later, at the Iceberg commit itself:
IcebergTableWriter::append_batches_with_marker prepares (transforms and
coerces) every WAL entry in a commit group independently and returns a
CommitOutcome { committed, rejected } rather than failing the whole group
on the first entry it cannot shape into the table. A rejected entry retires
immediately via Wal::dead_letter_rejected; its healthy same-cycle
neighbours still commit together in one all-or-nothing Iceberg transaction.
One poison entry can therefore no longer block its neighbours from
committing, or make a caller's do_action("flush") fail once the poison
entry is retired (W2, writer review 2026-08-23).
The writer's per-do_put "Received data" line is DEBUG; per-request
handler lines on the acceptor are DEBUG too — steady-state INFO shows WAL
batch commits, not individual requests.
Dictionary-Safe Encode/Decode¶
RecordBatch encoding goes through common::flight::batches_to_compressed_flight_data,
which forwards any dictionary batches IpcDataGenerator::encode produces
ahead of each data batch. Decoding on the receiving side goes through
common::flight::decode::flight_data_vec_to_batches — a wrapper around
arrow_flight::decode::FlightRecordBatchStream for call sites that buffer a
Vec<FlightData> before decoding — rather than
arrow_flight::utils::flight_data_to_batches, whose dictionaries_by_id map
is always empty and so cannot decode a stream containing dictionary batches.
No column in SignalDB's own schemas is dictionary-encoded yet; this only
removes the transport-level blocker for adopting one in the future.
5.3 Service Discovery Integration ✅ Implemented¶
Components discover each other via:
- Catalog-based service registry with PostgreSQL/SQLite backend
- ServiceBootstrap pattern for automatic registration on startup
- Capability-based routing (TraceIngestion, Storage, QueryExecution, Routing)
- Heartbeat monitoring with automatic TTL-based cleanup
- Flight endpoint discovery with connection pooling
Timeouts¶
Two independent bounds, deliberately not the same value:
| Bound | Covers | Default |
|---|---|---|
connect_timeout |
Establishing the channel | 30s |
request_timeout |
Each request on that channel | querier.query_timeout + 30s grace |
tonic's Endpoint::timeout() is a per-request deadline, not a connect
timeout — Endpoint::connect_timeout() is the latter. Conflating them makes
the caller abort before the callee's own timeout can fire, which replaces a
diagnosable DeadlineExceeded (HTTP 504) with an opaque Cancelled. The
request deadline is therefore derived from the configured query budget, so
raising querier.query_timeout raises the caller's patience with it.
6. WAL Integration ✅ Implemented¶
Write-Ahead Log provides durability and crash recovery capabilities:
┌─────────────┐ ┌─────────────┐ ┌─────────────┐
│ Acceptor │───▶│ WAL │───▶│ Writer │
└─────────────┘ └─────────────┘ └─────────────┘
Implemented WAL Features:¶
- ✅ Durability: Write incoming data to WAL before acknowledgment
- ✅ Recovery: Automatic replay of unprocessed entries on restart
- ✅ Batching: Efficient batch processing with configurable flush policies
- ✅ Entry Tracking: WAL entries marked as processed after successful storage
- ✅ Configurable Storage: Persistent WAL directories with segment rotation
Future WAL Enhancements:¶
- Compression: WAL segment compression for storage efficiency
- Replication: WAL replication for high availability
- Retention Policies: Automatic cleanup of old WAL segments
6.2 Enhanced Buffering¶
For handling backpressure and improving performance:
- In-memory buffering in Writer before Parquet persistence
- Configurable flush policies (size, time, or count-based)
- Real-time query support for buffered data
6.3 Multi-Writer Replication¶
For high availability:
- Hash-based data distribution across multiple Writers
- Replication factor configuration
- Automatic failover handling
7. Monolithic Binary Implementation ✅ Current¶
The current monolithic binary (cargo run --bin signaldb) starts all services in a single process:
- Services communicate via Flight using localhost endpoints
- Automatic service discovery via catalog
- Single configuration file for all components
8. Client SDK Integration¶
Flight communication enables:
- High-performance data transfer
- Streaming query results
- Native Arrow format support
- gRPC-based transport with authentication
9. Performance Benefits ✅ Achieved¶
Current implementation provides:
- Zero-copy data transfer: Arrow format maintained throughout pipeline
- Streaming capabilities: Large query results can be streamed
- Protocol efficiency: gRPC transport with minimal overhead
- Schema evolution: Arrow schema support for versioning
10. Deployment Modes¶
10.1 Monolithic Mode ✅ Current¶
- All services in single process
- Flight communication via localhost
- Simplified deployment and configuration
10.2 Microservices Mode ✅ Supported¶
- Services deployed independently
- Flight communication via network
- Service discovery via the shared catalog database
- Individual scaling and failure isolation
11. Conclusion¶
✅ Phase 2 Complete: SignalDB's Arrow Flight implementation with WAL integration is production-ready, providing:
Achieved Goals:
- High-performance Flight-based inter-service communication
- WAL-based durability with crash recovery
- Catalog-based service discovery with capability routing
- Complete elimination of message bus dependencies
- Support for both monolithic and distributed deployments
- Integration test coverage in
tests-integration/
Performance Benefits:
- Zero-copy data transfer via Arrow Flight
- Efficient service discovery with connection pooling
- Durability guarantees through WAL persistence
- Streaming query capabilities with DataFusion
Production Readiness:
- Robust error handling and retry logic
- Automatic service registration and health monitoring
- Configurable WAL and storage options
- Comprehensive logging and debugging capabilities
The Flight-based architecture with WAL integration provides a solid, production-ready foundation for observability data processing at scale.
The writer's
do_putv1→storage transformation resolves materialized-label allowlists per tenant (a tenant schema override replaces the global set).