Verified Telemetry Object Protocol Engine
Replay-safe, manifest-driven telemetry transfer from Kafka, files, and syslog spools into object storage.
VTOP Engine moves telemetry into object storage while protecting the source commit point.
It ingests telemetry from:
- Kafka topics
- append-only log files
- syslog spool files
The longer-term cluster direction is a VTOP-owned native Rust broker and control plane. Kafka remains supported as an optional edge/source adapter, not as VTOP's coordinator or correctness dependency. See the native broker architecture.
For every batch, VTOP:
- reads records from a source
- forms a batch
- compresses the batch
- computes a checksum
- creates a manifest
- uploads the object and manifest
- verifies the uploaded data
- commits source progress only after verification succeeds
Important
VTOP does not commit Kafka offsets, file byte offsets, or syslog spool offsets until the destination object has been verified.
Note
VTOP Engine is currently a prototype / reference implementation of a proposed protocol and proposed method.
This repository is intended to support candidate-invention disclosure work. It is not patented or patent-pending.
- Overview
- Status
- Why VTOP exists
- Core rule
- How it works
- State machine
- Supported source modes
- Format detection
- Architecture
- Quick start
- Build and test
- CLI usage
- Docker lab
- Example manifest
- Metrics
- Upload and verification
- Replay and recovery
- Known limitations
- Roadmap
- Documentation
- License
Most log-to-object-storage pipelines can move bytes into a bucket.
The harder problem is knowing when it is safe to advance the source position.
Did the object land intact,
and is it safe to commit the source offset now?
VTOP addresses this with a source-agnostic safety model:
| Capability | Purpose |
|---|---|
| Manifest-bound transfer | Binds source progress, object URI, checksum, format, compression, and verification state. |
| Verify before commit | Prevents source progress from advancing before destination verification. |
| Replay-safe state store | Allows recovery without silently losing unverified data. |
| Explicit state machine | Makes unsafe transitions visible and testable. |
| Pluggable sources and backends | Applies the same safety model to Kafka, files, syslog spools, and object storage backends. |
SOURCE_COMMITTED is forbidden until VERIFIED is true.
A source progress marker is never committed until the batch completes this lifecycle:
DISCOVERED
→ BATCHING
→ SEALED
→ COMPRESSED
→ CHECKSUMMED
→ OBJECT_UPLOADED
→ MANIFEST_UPLOADED
→ VERIFIED
→ SOURCE_COMMITTED
Failure can happen from any state:
ANY_STATE
→ FAILED
→ REPLAY_REQUIRED
→ BATCHING
Caution
Transitions such as SEALED → SOURCE_COMMITTED or OBJECT_UPLOADED → SOURCE_COMMITTED are invalid.
At a high level, VTOP separates source progress from destination durability.
Source records
→ batch
→ compressed object
→ checksum
→ manifest
→ upload object
→ upload manifest
→ verify object and manifest
→ commit source progress
The manifest acts as the transfer evidence record.
It links:
- source type
- source name
- source progress marker
- object URI
- object checksum
- compression type
- detected format
- batch metadata
- reproducible manifest self-hash and optional keyed-BLAKE3 authentication
The state machine is the enforcement point for safe progress.
Only this final transition is valid:
VERIFIED → SOURCE_COMMITTED
Invalid examples:
SEALED → SOURCE_COMMITTED
COMPRESSED → SOURCE_COMMITTED
CHECKSUMMED → SOURCE_COMMITTED
OBJECT_UPLOADED → SOURCE_COMMITTED
MANIFEST_UPLOADED → SOURCE_COMMITTED
Relevant implementation:
| Source mode | Progress marker | Behavior |
|---|---|---|
| Kafka | topic, partition, offset range | Uses a Kafka consumer with auto-commit disabled. Each batch contains records from one topic and one partition. Offsets are committed only after verification. |
| File | path, inode, byte range | Reads append-only files line by line. Partial trailing lines are not committed. Replay resumes from the last safe byte offset. |
| Syslog spool | spool ID, byte range | Treats rsyslog or syslog-ng spool files as append-only inputs. External collectors own syslog delivery; VTOP owns batching, upload, verification, replay, and commit safety. |
VTOP is not fixed to one telemetry format.
When a stream does not explicitly define a format in examples/streams.yaml, the engine detects the format per batch.
Supported detected formats:
| Format | Example output extension |
|---|---|
| CEF | .cef.gz |
| JSON | .json.gz |
| JSON Lines | .jsonl.gz |
| Syslog | .syslog.gz |
| Plain text | .txt.gz |
A single engine can process different formats across different streams.
For example:
source A → CEF
source B → JSON Lines
source C → syslog
source D → mixed batches
The detected format is recorded in the manifest.
Relevant implementation:
The engine reads from a source, writes a compressed object plus a manifest to object storage, verifies what it wrote, and only then advances the source commit point. Verification failure means the source is never committed, so the data stays replayable.
flowchart LR
S["Kafka · Files · Syslog spool"]
E(["VTOP engine"])
O[("S3 / MinIO")]
D[("State store")]
S -->|"read batch"| E
E -->|"1 · upload object + manifest"| O
O -->|"2 · verify size + checksum"| E
E ==>|"3 · commit progress<br/>ONLY after VERIFIED"| S
E <-.->|"batch state · replay ledger"| D
classDef src fill:#eef4ff,stroke:#4a72b8,stroke-width:1px,color:#12263f
classDef eng fill:#e8f7ee,stroke:#2e8b57,stroke-width:2px,color:#12263f
classDef store fill:#fff6e5,stroke:#c07f19,stroke-width:1px,color:#12263f
class S src
class E eng
class O,D store
Steps 1 → 2 → 3 are the whole protocol: the thick arrow is the one rule that must never break — see Core rule.
VTOP is organized as a Rust workspace.
crates/
vtop-core/ protocol-independent logic:
state machine, batching, manifests,
checksums, compression, partitioning,
config, replay
vtop-log/ Kafka-independent native broker storage:
framed records, active/sealed segments,
crash recovery, sparse indexes, manifests
vtop-adapters/ source adapters:
Kafka, file, syslog spool
vtop-upload/ upload backends:
native S3, s3cmd, awscli, MinIO, mock
vtop-state/ durable SQLite / feature-gated PostgreSQL state store
vtop-cli/ vtopctl CLI and engine runtime
examples/ example config and sample streams
docs/ protocol, architecture, security, invention notes
docker/ container build files and seed scripts
tests/ integration tests
benchmarks/ benchmark and performance test support
Tip
Keep detailed internals in docs/. The README should stay focused on orientation, quick start, and key guarantees.
Full architecture documentation:
Run the full Docker lab:
docker compose up -d
docker compose logs -f vtop-engineOr build and run locally:
cargo build --release
cargo run -p vtop-cli -- discover --config examples/config.yamlcargo fmt --all --check
cargo clippy --workspace --all-targets -- -D warnings
cargo test --workspace
cargo build --releaseCI runs formatting, linting, tests, and release build on push and pull request.
Kafka is enabled by default for compatibility with existing deployments, but
it is a feature-gated adapter. A native/file/syslog-only CLI build does not
compile or link rdkafka:
cargo build -p vtop-cli --no-default-featuresWorkflow file:
.github/workflows/ci.yml
The CLI binary is vtopctl.
cargo run -p vtop-cli -- run \
--config examples/config.yaml
cargo run -p vtop-cli -- discover \
--config examples/config.yaml
cargo run -p vtop-cli -- process-once \
--source kafka \
--config examples/config.yaml
cargo run -p vtop-cli -- process-once \
--source file \
--config examples/config.yaml
cargo run -p vtop-cli -- replay \
--batch-id <batch_id> \
--config examples/config.yaml
cargo run -p vtop-cli -- status \
--config examples/config.yaml
cargo run -p vtop-cli -- list-batches \
--config examples/config.yaml \
--json
cargo run -p vtop-cli -- verify-manifest \
--manifest s3://telemetry-data/.../batch.manifest.json \
--config examples/config.yaml
# PostgreSQL only: run once with the privileged migration secret before the
# engine starts with its DML-only runtime secret.
cargo run -p vtop-cli --features postgres -- migrate \
--config examples/config.yamlCommon CLI behavior:
| Option | Purpose |
|---|---|
--json |
machine-readable output |
--log-level |
runtime log level |
| non-zero exit | command failure |
| secret-safe output | commands should not print credentials |
PostgreSQL schema changes are never run by normal engine startup. Use a
separate migration identity for vtopctl migrate, then give the runtime role
only schema USAGE and SELECT, INSERT, UPDATE on batches. See
PostgreSQL deployment.
The Docker lab provides Kafka, MinIO, seeded telemetry, and the VTOP engine.
| Service | Purpose |
|---|---|
kafka |
test Kafka broker |
kafka-ui |
browser UI at http://localhost:8080 |
minio |
S3-compatible object storage |
minio-init |
bucket initialization |
kafka-init |
seeded test events |
vtop-engine |
VTOP runtime |
rsyslog |
optional syslog collector profile |
MinIO endpoints:
API: http://localhost:9000
Console: http://localhost:9001
Bucket: telemetry-data
Start the lab:
docker compose up -d
docker compose logs -f vtop-engineStart Kafka, MinIO, and seed data:
docker compose up -d kafka minio minio-init kafka-initStart the engine:
docker compose up -d vtop-engine
docker compose logs -f vtop-engineExpected lifecycle events:
format_detected
object_uploaded
manifest_uploaded
verification_passed
source_committed
Open the MinIO console:
http://localhost:9001
Then inspect the telemetry-data bucket.
Generate test input files:
docker/seed-events.sh cef 200 > ./data/input/auth.cef.log
docker/seed-events.sh json 200 > ./data/input/app.json.log
docker/seed-events.sh syslog 200 > ./data/input/sys.syslog.log
docker/seed-events.sh mixed 500 > ./data/input/mixed.logRun the engine:
docker compose up -d vtop-engine
docker compose logs -f vtop-engineGenerate additional test data at any time:
docker/seed-events.sh <cef|json|jsonl|syslog|mixed> [count]Infrastructure-free file-flow test:
tests/integration_file_to_minio.rs
{
"protocol": "VTOP",
"version": "0.2",
"batch_id": "vtop-20260618T150000Z-app_events-p0-481000-482499-1a2b3c4d",
"tenant": "default",
"source_type": "kafka",
"source_name": "app_events",
"format": "cef",
"compression": "gzip",
"record_count": 1500,
"source_progress": {
"source_type": "kafka",
"topic": "app_events",
"partition": 0,
"start_offset": 481000,
"end_offset": 482499,
"consumer_group": "vtop-engine"
},
"object": {
"uri": "s3://telemetry-data/tenant=default/source=app/format=cef/year=2026/month=06/day=18/hour=15/vtop-....cef.gz",
"size_bytes": 924822,
"sha256": "abc123..."
},
"manifest": {
"uri": "s3://telemetry-data/.../vtop-....manifest.json",
"sha256": "def456...",
"mac": "0123abcd..."
},
"state": "manifest_uploaded",
"verification_status": "not_verified"
}Note
The manifest is written at MANIFEST_UPLOADED, before storage-side verification.
The authoritative post-verification state lives in the state store and can be queried with vtopctl status or vtopctl list-batches --json.
The manifest self-hash is reproducible and detects accidental changes, but an
attacker able to rewrite the manifest can recompute it. Set
manifest_mac_key_env to the name of an environment variable containing a
32-byte hex key to add manifest.mac, a keyed BLAKE3 authenticator. Both
embedded values are blanked for canonicalization. The key itself is never
serialized. Enabling a key intentionally rejects older unsigned manifests;
verify the backlog before cutover. Key rotation is not implemented yet.
VTOP emits structured per-batch metrics.
Example:
3 records, 114 B->80 B (1.43x, 29.8% saved) in 6 ms | 500 rec/s, 0.00 MiB/s up |
stages: compress=0ms checksum=0ms put_obj=0ms put_manifest=0ms verify=0ms commit=0ms
Each batch records:
| Metric area | Examples |
|---|---|
| Size and transfer | uncompressed bytes, compressed bytes, compression ratio, percentage saved |
| Latency | compression, checksum, object upload, manifest upload, verification, commit |
| Throughput | records/sec, uncompressed MiB/sec, effective upload MiB/sec |
vtopctl process-once --json includes the full metrics object per batch.
Prometheus metrics are exported by the engine at /metrics when VTOP_METRICS_ADDR is set. See observability/ for the optional Grafana LGTM stack (Alloy + Mimir/Loki/Tempo) and dashboards.
Relevant implementation:
VTOP supports multiple upload backends.
| Backend | Purpose | Verification level |
|---|---|---|
| native S3 | primary S3-compatible backend | service-computed SHA-256 or streamed BLAKE3 |
| AWS CLI | command-based backend | downloads and hashes stored content |
| s3cmd | command-based backend | downloads and hashes stored content |
| MinIO client | command-based backend | downloads and hashes stored content |
| LocalFS | local/air-gapped backend | streams stored files through the configured digest |
| mock | tests and local integration flow | hashes stored in-memory content |
Important
Strong verification is the default. A sidecar, ETag, or uploader-written user metadata is never accepted as proof of stored content.
The awscli, s3cmd, and minio compatibility backends are an explicit
opt-in. They require upload.command_binary to be an absolute path, verify the
tool's --version identity at startup, clear the child environment, and apply
wall-clock and captured-output limits. Add only the exact runtime variable
names the selected tool needs to upload.command_env_allowlist; values are
resolved at startup and are never serialized. Native s3_native does not
spawn an external process and does not use these settings.
VTOP recovery is designed around one rule:
Unverified data must remain replayable.
Recovery behavior:
| Crash point | Recovery action |
|---|---|
| before object upload | replay from source |
| after object upload but before verification | replay from source |
| after verification but before source commit | retry source commit |
| after source commit | batch is complete |
If verification fails:
batch → FAILED
source progress → not committed
If commit fails after verification:
batch → VERIFIED
recovery → retries source commit
Relevant tests:
tests/integration_replay.rs
tests/integration_state_recovery.rs
VTOP is currently a prototype. The following limits are known and intentional.
| Area | Current behavior | Planned direction |
|---|---|---|
| Large objects | native S3 backend uses single-part put_object |
add multipart upload |
| Large records / whole files | max_bytes is a hard per-source/per-batch ceiling; an oversized record is rejected without advancing source progress |
raise the explicit budget only when the deployment has matching memory headroom |
| Partial upload recovery | replays from source instead of resuming half-written local objects | add resumable local staging |
| Command backend verification cost | aws, s3cmd, and mc download each stored object to hash it |
prefer native S3 SHA-256 when read-back bandwidth is costly |
| Syslog timestamps | received_time_* is not yet extracted into the spool marker |
add timestamp extraction |
| Manifest integrity | self-hash plus optional keyed-BLAKE3 authentication; key rotation not implemented | add multi-key rotation and public-key signatures if required |
| Object immutability | S3 Object Lock is designed but not implemented | add Object Lock profile |
| Metrics export | Prometheus /metrics implemented (opt-in via VTOP_METRICS_ADDR); OpenTelemetry trace export not yet |
add OTLP span export |
| Kafka integration test | requires live broker and is ignored by default | add optional CI service profile |
| Binary / pre-compressed inputs | supported via the file source whole_file mode (archived verbatim, byte-exact) |
streaming for very large files |
| Local filesystem backend | available (backend: localfs, objects under local_path/<bucket>/<key>; sidecars are inventory hints only) |
— |
| Checksums | SHA-256 and BLAKE3, or disabled (size-only); strong verification defaults on and require_strong_verification: false is an explicit weak-mode opt-out |
— |
Completed:
- local filesystem upload backend (
backend: localfs) - BLAKE3 checksum strategy (and checksum-disabled mode)
- binary / pre-compressed input framing (file source
whole_filemode) - strong-verification gate (
require_strong_verification) - Prometheus metrics exporter — the engine serves
/metrics(/healthz,/readyz) behindVTOP_METRICS_ADDR; see observability/ for the optional Grafana LGTM stack and dashboards - end-to-end smoke + live-broker Kafka CI over the full compose lab
- Kafka is isolated behind the optional
kafkaCargo feature - first Kafka-independent native segment-log storage kernel
- bounded native produce/fetch wire codec and TLS-1.3 mTLS local-broker library with durable producer-epoch fencing and committed-only fetch
Planned implementation areas:
- bounded file/syslog/whole-file reads, pre-clone Kafka record checks, and
streaming local compression;
max_bytesis enforced before source progress advances (native object upload remains single-part) - native three-node metadata/control-plane prototype
- multipart upload support
- optional keyed-BLAKE3 manifest authentication via a named secret env var
- manifest MAC key rotation / optional public-key signatures
- S3 Object Lock profile
- OpenTelemetry trace export (the metrics endpoint exists; spans do not yet)
- million-file benchmark suite
| Document | Purpose |
|---|---|
| docs/README.md | documentation index |
| docs/VTOP_PROTOCOL_DRAFT.md | protocol draft and conformance profiles |
| docs/ARCHITECTURE.md | architecture and runtime flow |
| docs/SECURITY_MODEL.md | security model and normative rules |
| docs/INVENTION_DISCLOSURE_DRAFT.md | candidate-invention disclosure draft |
| docs/PRIOR_ART_SEARCH_PLAN.md | prior-art search plan |
MIT © 2026 Tamir Suliman.
