Conversation
Offsets are per-partition, so consumers need partition count metadata to discover a topic's full key space. Adds list_partitions to the Request and Response envelopes, MESSAGE_TYPE_LIST_PARTITIONS, and the server handler backed by the existing Broker.NumPartitions.
Happy path (3-partition topic reports 3) and unknown-topic error path.
Command-line client with create, produce, consume, topics, and partitions subcommands over the length-prefixed protobuf protocol. consume --all discovers partition count via the new ListPartitions RPC, drains every partition, and merges by timestamp best-effort; ordering remains per-partition/per-key only (documented in usage). 5s dial timeout so bad addresses fail fast.
Covers produce/consume round trip with key affinity, end-of-log termination for sequential fetch, empty-partition behavior, the full consume --all wiring (flags -> partition discovery -> drain -> merge), and server-error surfacing. Adds Server.SetListener so external packages can run the server on an ephemeral port they inspect.
Builds mqclient, starts docker compose (--build so image matches the checkout), creates a unique per-run 3-partition topic (the named volume persists topics between runs), produces keyed messages showing FNV-1a partition assignment, consumes merged across partitions, verifies bytes at rest with mqvalidate, and cleans up. Compatible with asciinema rec.
README leads with the demo GIF (real PTY recording via asciinema 2.4, GIF via agg 1.9.0): keyed produce with FNV-1a partition assignment, merged consume across partitions, mqvalidate on-disk check. Ordering guarantee stated the way Kafka states it: per-partition/per-key only, cross-partition merge best-effort. docs/demo.cast kept for terminal playback (asciinema play). Replaces the synthetic cast generator, which appended stderr after stdout and produced out-of-order frames. demo.sh gains a real-request readiness probe (a bare TCP connect can succeed against the old container during compose recreate, RST-ing the next request) and moves the MQ helper above first use.
module message-queue -> github.com/ga11221/message-queue across go.mod, all imports, proto go_package, and regenerated bindings. Required before the repo goes public; mechanical rename, no behavior change. Full suite green.
Produce path is synchronous (direct call stack, no writer goroutine or reply channels — that was the pre-implementation design; see server-design-decisions.md). syncSegment removed from the Append diagram: fdatasync runs only on Close/rotation, not per append. Index writes labeled by interval with the default (1) made explicit. New metadata-flows section covers CreateTopic/ListTopics and the ListPartitions RPC added for the demo CLI. Verified against code 2026-08-26.
Layer-responsibility table (server/broker/topic/partition/segment/ message with explicit does-NOT-know boundaries), the big-picture Mermaid diagram from TCP client down to the three-file segment format, produce-to-bytes and consume-from-bytes data flows, the five key invariants (per-partition offsets, key affinity, immutable closed segments, no-fsync append path, synchronous request lifecycle), and an explicit not-yet list with deferral reasons.
Export boundary corrected (Broker is the entry point, not Topic). Coordination table: partition writes use the partition mutex, not a single-writer assumption (contention measured ~12%, negligible). fdatasync claims removed from the append path (sync happens only at segment Close/rotation). Rotation is size-based only (1GB); the '1 hour elapsed' trigger was never implemented. Index sparsity clarified: default interval is 1 (dense), with the measured syscall cost noted and the Phase 2.5 in-memory-index target referenced. The 'Batch + Single Writer' section relabeled as the pre-implementation plan; current design is synchronous, batching deferred to Phase 2.5 with the measured knee cited. Kafka batching material retained as design grounding. Verified against code 2026-08-26.
Post-process asciicast timestamps to insert 0.5-2.0s pauses between section headers, produce output lines, and consume/validate results. The original GIF compressed all produce+consume into ~40ms — too fast for a human viewer. New GIF: 29s total, ~20 frames. Docker idle (13.5s) stays compressed by the cast's idle_time_limit (2s). Re-rendered with agg v1.9.0 (agg-x86_64-unknown-linux-musl).
Add .env environment headers to the top of each .txt file so they are self-contained (reader sees hardware/kernel/go version before the numbers). Generate CSV exports (baseline.csv, sweep.csv) — 140 and 25 rows respectively — importable to spreadsheet tools. New bench-results/README.md with summary tables: baseline encode/decode/ append/read/lookup medians and the batch-size sweep knee curve with the 16-64 target annotated. Reproduction instructions included.
Header: reframe from 'exactly-once via 2PC' (unimplemented) to 'Phase 1 complete; Phase 2 in progress'. Architecture diagram: replace stale '[buffer channel] --> single writer' with accurate synchronous call stack (links to docs/phase1/architecture.md). Design decisions: 'Single writer per partition' reworded to 'per-partition locking' with the measured contention cited (~12%, negligible vs 86% syscall cost). Phase table: 'Not started' updated to 'Design complete, implementation next'. Add CI badge. Drop the '~15 seconds' GIF caption (timing changed with paced re-render).
partition.Read() and partition.NextOffset() were taking write locks despite being pure read operations. This serialized all reads on a partition unnecessarily. Upgraded partition.mu from Mutex to RWMutex and use RLock for Read/NextOffset, allowing concurrent reads.
Documents the four-layer locking hierarchy (Server -> Broker -> Topic -> Partition -> Segment), parallelism model, shutdown coordination, and known limitations.
Five ADRs covering key design decisions: - ADR-001: RWMutex for partition reads (b3b4d2f) - ADR-002: Synchronous architecture, no writer goroutine (7a295e2) - ADR-003: Sparse index with interval=1 (04d95e8) - ADR-004: fdatasync over fsync (fbd3b8b) - ADR-005: 28-byte fixed-size message header (b3475d3) Also fixes stale architecture diagram in README and adds concurrency model link to docs table.
- ADR-006: Length-prefixed protobuf wire protocol (12aff04) - ADR-007: FNV-1a hash for key-based partition routing (d7d0f4b) - ADR-008: Release closed segment file handles after recovery (1af0d20) - ADR-009: Fuzz testing with committed seed corpus (6b89d1a) - ADR-010: Immutable closed segments (chmod 0444) (70f8b6b)
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Phase 1 complete: single-node log engine + broker + TCP server
Merges the demo branch into main, completing Phase 1 of the Kafka-style message queue.
Features
Benchmarking
Documentation
Bug fix
Test coverage
77%+ across all packages. Race-detected. Fuzz-tested (seed corpus committed).
Module path
Renamed to github.com/ga11221/message-queue for public visibility.