Skip to content

Latest commit

 

History

History
194 lines (157 loc) · 15.8 KB

File metadata and controls

194 lines (157 loc) · 15.8 KB

Progress Tracker

Timeline

Phase Component ETA Status
1 Log engine + Broker + Server/RPC (single-node) 1-2 weeks COMPLETE
2 Exactly-once via 2PC (single-node) 2-3 weeks IN PROGRESS
3 Replication (ISR, leader election, replica fetch) 2-3 weeks NOT STARTED
4 Stream processing (windowing, state, fault tolerance) 2-3 weeks NOT STARTED

Phase 4 is conditional — only if Phases 1-3 complete in 7-9 weeks.

Phase 1: Single-node Log Engine

Component Status Date Notes
Message format DONE 08/18 28-byte header, CRC-32, big-endian
Index (offset -> position) DONE 08/18 Sparse, configurable interval, binary search
Time index (offset -> timestamp) DONE 08/18 Sparse, configurable interval, binary search
Segment (create, open, close, append, read) DONE 08/18 EINTR retry, fdatasync, message count tracking
Segment recovery DONE 08/18 Scan forward, truncate partial messages + stale indexes
Partition (create, open, close, append, read) DONE 08/18 Offset assignment, rotation, binary search for closed segments
Topic (create, recover, produce, consume) DONE 08/24 Static partition count, FNV-1a key routing, round-robin nil key
Producer API types DONE 08/24 pkg/producer/record.go
Consumer API types DONE 08/24 pkg/consumer/record.go
Broker (single node) DONE 08/24 Topic management, produce/consume routing, recovery
Partition mutex DONE 08/25 sync.Mutex, concurrent access safety
Server / RPC DONE 08/25 Protobuf wire protocol, goroutine-per-connection, integration tests
Fuzz testing DONE 08/25 4 targets; found+fixed DecodeMessage uint32 underflow/overflow bug
Docker packaging DONE 08/25 Multi-stage Dockerfile, compose w/ named volume, smoke-tested
External server test DONE 08/25 TestExternalServer gated by MQ_ADDR env var, -count=1 required

Phase 2: Exactly-once via 2PC (single-node)

Component Status Date Notes
Coordinator NOT STARTED Transaction orchestration
Participants NOT STARTED Log + consumer offset
Recovery NOT STARTED In-doubt resolution
Idempotency keys NOT STARTED Deduplicate retries

Phase 3: Replication

Component Status Date Notes
ISR tracking NOT STARTED In-sync replica set
Leader election NOT STARTED Preferred leader, unclean leader election
Replica fetch protocol NOT STARTED Follower pulls from leader
min.insync.replicas NOT STARTED Durability guarantee

Phase 4: Stream Processing (conditional)

Component Status Date Notes
Windowing NOT STARTED Tumbling, sliding
State management NOT STARTED RocksDB or in-memory
Fault tolerance NOT STARTED Checkpoint to log

CI / Infrastructure

Component Status Date Notes
GitHub Actions DONE 08/20 vet, test -race, coverage gate 77%
Coverage reports DONE 08/20 67 tests, 79.9% coverage
MIT License DONE 08/20
SSH commit signing DONE 08/20 Verified GSS: commits + v0.1.0 tag signed (ED25519)

Test Coverage

Package Tests Coverage Notes
message.go 6 100% encode/decode, CRC, edge cases
index.go 4 ~80% write, sparse, lookup, empty
timeindex.go 3 ~80% write, lookup, empty
segment.go 18 ~85% create, open, recovery, scan, append, read, close, missing index files
partition.go 8 ~85% create, append, read, rotation, recovery, close
topic.go 9 ~80% create, recover, produce, consume, key routing, round-robin, close
broker.go 12 83.8% create, recover, produce, consume, routing, key consistency
server.go 6 78.1% produce, consume, create topic, list topics, errors, concurrent produce
protocol.go 1 ~75% round-trip framing
Total 67 79.9% Remaining: file system error paths, listen error

Control Flow Diagrams

See docs/phase1/control-flow.md — 7 Mermaid diagrams:

  • Full stack: Produce (TCP client -> server -> broker -> partition)
  • Full stack: Consume (TCP client -> server -> broker -> partition)
  • Server connection lifecycle
  • Append (partition level)
  • Read (partition level)
  • Recovery on startup
  • Write errors

Key Learnings

Phase 1

  • Message type is unexported — external packages use ProduceRecord/ConsumeRecord
  • Offset and timestamp are broker-assigned — producer only sets Key+Value
  • CRC-32 over SHA-256 — hardware-accelerated, 1-in-4-billion error rate, sufficient for bit-rot
  • Big-endian throughout — network byte order, consistent serialization
  • Index uses message count interval — not bytes (simpler, aligns with offset semantics)
  • Segment is immutable after Close() — Go-level convention, handles nil'd
  • Recovery truncates index files — both .index and .timeindex to match valid message count
  • Partition delegates to Segment — clean boundary, Segment owns storage, Partition owns routing
  • Segment.messageCount is single source of truth — Partition derives nextOffset from it
  • On-demand file opening — closed segments nil handles, open .log/.index on read
  • EINTR retry — 3 attempts with immediate retry (no backoff)
  • 100% coverage unrealistic without VFS abstraction — 77.6% practical ceiling
  • Directory ownership: topic + partition only — segment, index, timeindex never create directories. Only two MkdirAll calls: one in CreateTopic, one in createPartition. Everything else operates on files within the partition directory.
  • Unexport internal types — partition and segment are unexported. Topic is the only exported type. Enforces layering at compile time. PartitionOption stays exported (used in CreateTopic/RecoverTopic signatures).
  • sync.Mutex over RWMutex for partitions — RWMutex has reentrancy issues: Read() holds RLock, calls nextOffset(), which also tries RLock. Deadlock when a writer is waiting. Solution: internal nextOffsetLocked() helper, all callers use it under the existing lock.
  • Import aliasing for colliding names — when a package name matches stdlib (e.g., "log"), alias both: mqlog "message-queue/pkg/log" and stdlog "log". Happens in both server and cmd/server.
  • Protobuf framing is trivial — 4-byte big-endian length prefix + protobuf bytes. Max frame size (16MB) prevents OOM from malicious clients. WriteFrame/ReadFrame are ~40 lines.
  • Partition mutex added late — realized concurrent Produce from multiple TCP connections would corrupt segment files. The server design decisions document forced this realization before coding.
  • Design decisions before code — the 6-question decision tree (connection model, wire protocol, request lifecycle, partition sync, batching, broker API) prevented rework. Each decision had clear concurrency implications.
  • 2PC before replication — single-node 2PC doesn't need replication. Move to Phase 2, replication to Phase 3. Each phase is independently testable.
  • Batching is orthogonal to 2PC — 2PC is about atomicity (commit all or abort all), not batching. Don't conflate them in the partition interface.

Fuzz Testing (08/25)

  • Fuzzer found a real bug in under 3 seconds — DecodeMessage panicked on crafted input. Hand-written tests (6 message tests, 100% coverage) never found it; the fuzzer did in 3s / ~500k execs. Coverage measures execution, not adversarial input.
  • Two distinct integer bugs in one decode path:
    1. Underflow: size - keyLen where both are uint32. If keyLen > size, the subtraction wraps to ~4 billion -> slice bounds panic at buf[keyEnd:valueEnd]. Fix: validate keyLen <= size.
    2. Overflow in the guard itself: the check len(buf) < HeaderSize+size overflows when size is near MaxUint32 (HeaderSize+size wraps small, check passes). Fix: rewrite as size > uint32(len(buf))-HeaderSize — safe because len(buf) >= HeaderSize was already checked.
  • Lesson: never add unsigned values inside a bounds check — compute against remaining capacity instead (size > remaining). The guard must be as overflow-safe as the code it protects.
  • Failing inputs become permanent regression tests — Go fuzzing writes crashers to testdata/fuzz//; these run in every normal go test, so the bug can never silently return. Committed both crasher files.
  • Nil seeds fail at build time — f.Add(nil, nil) doesn't type-match []byte params; use empty slices []byte{} instead.
  • Fuzz targets chosen by trust boundary — anything that parses untrusted bytes (DecodeMessage from disk/network, ReadFrame from TCP) gets fuzzed; internal state machines don't. 4 targets, all serialization-layer.
  • Round-trip fuzzing catches silent corruption too — FuzzMessageRoundTrip/FuzzFrameRoundTrip assert decoded == original for arbitrary key/value bytes, not just no-panic. Catches encode/decode asymmetry that fixed test vectors miss.

Docker + External Integration Testing (08/25)

  • Container keeps data off the host disk — the original motivation. Named volume mq-data mounted at /data inside the container; verified persistence across container restarts (topics survived). Multi-stage build: golang:1.26-alpine build -> alpine runtime, CGO_ENABLED=0 static binary.
  • FNV-1a routing gotcha in tests — produce response returns the actual partition; consumers must use that partition, not assume 0. First smoke-test client hardcoded partition 0 and failed with "offset out of range".
  • External integration test gated by env var — TestExternalServer runs create/produce/consume against a live server when MQ_ADDR=localhost:9092 go test -count=1 -run TestExternalServer. Skipped without MQ_ADDR so CI stays container-free.
  • Go test cache cannot see external state — cache is keyed on code inputs and env vars read via os.Getenv, but not server-side state. Restarting the container doesn't invalidate a cached PASS. Always -count=1 when testing against external infrastructure.
  • Unique topic names per run — integration-<unixnano> avoids conflicts with data persisted in the container volume across repeated runs.
  • Fuzz tests never touch the network — they call DecodeMessage/ReadFrame directly in-process; unit/integration tests spin up their own ephemeral server. Only TestExternalServer exercises a running container over TCP.

Dual Docker Daemon Gotcha (08/25)

  • Two engines were running simultaneously — Docker Desktop (WSL2 backend) AND Ubuntu's native docker.service (systemd, started at boot). Desktop took over /run/docker.sock at 11:52, so the CLI talked to Desktop while the native daemon kept running invisible in the background.
  • An orphaned broker container hijacked port 9092 — started under the native daemon before Desktop integration activated. Every integration test passed against it; its data went to native-daemon storage, so the mq-data volume in Docker Desktop stayed empty. Tests "passing" with the compose container stopped was the tell.
  • Diagnosis path that found it: held a probe connection open and ran ss -tn, which showed ESTAB to 172.19.0.2:9092 — a bridge subnet belonging to no network in docker network ls. That subnet belonged to the native daemon (br-f94ee353a59a, 172.19.0.0/16).
  • Fix: sudo systemctl disable --now docker.service docker.socket containerd + systemctl mask to prevent boot re-enable. Single engine now: Docker Desktop serves both GUI and CLI.
  • Lesson: docker ps only shows one daemon's containers — a listening port can be served by an engine your CLI cannot see. If container state and CLI output disagree, check for multiple daemons (ps aux | grep dockerd, ip -br addr | grep 172.) before trusting any test result against "the" container.

Benchmarks + Profiling (08/25)

First profile-driven look at Phase 1 performance (AMD Ryzen 7 7735HS, Go 1.26, 100-byte values). Suite: pkg/log/bench_test.go.

Benchmark Result
EncodeMessage 74ns (16B) to 735ns (1KB); ~1.4 GB/s at 1KB
DecodeMessage (incl CRC verify) 107ns (16B) to 696ns (1KB)
SegmentAppend 3.1us/op, 43 MB/s, 272 B/op
PartitionAppend 3.2us/op (+mutex, offset assign: ~negligible)
ConcurrentProduce (parallel writers) 3.6us/op — only ~12% degradation, mutex NOT the bottleneck
PartitionRead @10k depth 10.7us/op
PartitionRead @1M depth 15.5us/op (only +45% despite 100x depth)
IndexLookup @10 entries 3.2us/op (!)

Key findings from CPU/memory/block profiles:

  • 86% of append CPU is in write syscalls — every append issues THREE separate pwrites (.log, .index, .timeindex each written independently). timeIndex.write alone is 37% cumulative. This is the top optimization target: batch/buffer the three writes into one, or group-commit multiple records per syscall.
  • 98% of allocations come from Encode() — fresh buffer per record (272 B/op for a 134-byte record). Fix candidate: sync.Pool or caller-provided buffer.
  • Index lookup does a ReadAt syscall PER BINARY SEARCH PROBE — explains why looking up among just 10 entries costs 3.2us (~20 probes x syscall latency). The index file is small enough to hold in memory (or mmap). Same pattern likely in timeindex. Biggest read-path win available.
  • Mutex contention is a non-issue — parallel writers degrade throughput only ~12%. Do not add lock-free complexity yet.
  • fsync is not on the hot path by design — sync happens only on segment close/rotation, matching the fdatasync-window tradeoff documented in the Kafka study. Durability window exists but throughput doesn't pay for it per-append.
  • Read scales gracefully with depth (+45% from 10k to 1M) because binary search bounds probe count; the constant factor is syscall overhead, which finding #3 addresses.

These numbers justify the deferred optimizations with data: (1) in-memory/mmap indexes, (2) batched writes, (3) encoder buffer reuse — in that order.

Batch-Size Sweep (08/25) — Phase 1.5

BenchmarkSegmentAppendBatch writes N pre-encoded records per write() call (N in {1,4,16,64,256}) to find where syscall amortization saturates. Per-record cost on ext4 (WSL2 VHDX), 134-byte records:

Records/write ns per record Throughput
1 1296 ~100 MB/s
4 424 ~277 MB/s
16 207 ~495 MB/s
64 144 ~740 MB/s
256 127 ~1060 MB/s

Findings:

  • The knee is at roughly 16-64 records (~2-8KB per write): batching 16 records already removes ~84% of the syscall tax; going from 64 to 256 buys only ~12% more. Matches the back-of-envelope crossover estimate within 2x.
  • At batch=256 we approach medium bandwidth (~1 GB/s) — the bottleneck shifts from syscall tax to actual storage throughput, confirming the "syscall-bound vs disk-bound" taxonomy. A buffered-writes design should target batches of 32-128 records; beyond that you are paying latency for nothing.
  • Pitfall found: benchmark scratch space must be real disk. The first sweep ran against /tmp tmpfs (RAM-backed): numbers were faster but meaningless for storage workloads, and a multi-GB run exhausted the 4.9GB mount mid-benchmark (ENOSPC after 93s). benchDir now creates scratch under ~/.cache/mq-bench (ext4) with explicit cleanup.
  • Pitfall found: bound file growth inside write benchmarks. Writing one ever-growing log made late samples slower than early ones (+107% variance). The benchmark now rewinds its scratch file every 64MB; variance collapsed to a few percent.
  • Baseline was recaptured on ext4 after the tmpfs fix (append moved 3.21us -> 3.88us/op; the old tmpfs baseline was flattering). All committed results files now measure the same medium as future comparisons.