Verified against code 2026-08-26 (demo branch). The produce path is fully synchronous — one request is read, executed, and answered before the next read. There are no writer goroutines or reply channels (that was an earlier design; see server-design-decisions.md for why it was dropped).
flowchart TD
CLIENT[TCP Client] -->|ProduceRequest| SERVER[Server.handleConn]
SERVER -->|ReadFrame| DECODE[protobuf decode]
DECODE --> DISPATCH[Server.dispatch]
DISPATCH -->|ProduceRequest| HANDLE[Server.handleProduce]
HANDLE -->|Broker.Produce| BROKER[Broker]
BROKER -->|topic lookup| TOPIC[Topic.Produce]
TOPIC -->|partitionForKey| ROUTE[FNV-1a key hash mod numPartitions]
ROUTE -->|direct call| PART[partition.Append]
PART -->|mu.Lock| OFFSET[assign offset]
OFFSET --> SEG[active segment]
SEG -->|Encode| LOG[.log file]
SEG -->|entry at interval| IDX[.index file]
SEG -->|entry at interval| TIDX[.timeindex file]
SEG -->|messageCount++| NEXT{size >= segmentSize?}
NEXT -->|no| RESULT[return partition, offset]
NEXT -->|yes| ROTATE[Partition.rotate]
ROTATE -->|Close current| CLOSED[segment sealed + synced]
ROTATE -->|Create new| NEWSEG[new active segment]
NEWSEG --> RESULT
RESULT --> HANDLE2[Server.handleProduce]
HANDLE2 -->|ProduceResponse| RESP[protobuf encode]
RESP -->|WriteFrame| CLIENT
Example trace — TCP client produces to topic "orders", key="user-1":
1. Client sends: ProduceRequest{topic:"orders", key:"user-1", value:"hello"}
2. Server reads frame (4-byte length + protobuf)
3. Server dispatches to handleProduce
4. Broker.Produce("orders", "user-1", "hello")
5. Topic.Produce: FNV-1a("user-1") % 3 = partition 1
6. Direct call to partitions[1].Append() — synchronous, no channels
7. Partition assigns offset=7, writes to segment
8. Result returns up the call stack
9. Server returns: ProduceResponse{partition:1, offset:7}
10. Server writes frame to TCP connection
flowchart TD
CLIENT[TCP Client] -->|ConsumeRequest| SERVER[Server.handleConn]
SERVER -->|ReadFrame| DECODE[protobuf decode]
DECODE --> DISPATCH[Server.dispatch]
DISPATCH -->|ConsumeRequest| HANDLE[Server.handleConsume]
HANDLE -->|Broker.Consume| BROKER[Broker]
BROKER -->|topic lookup| TOPIC[Topic.Consume]
TOPIC -->|partition range check| PART[partition.Read]
PART -->|mu.Lock| FIND[findSegment]
FIND -->|active or closed| SEG[segment]
SEG -->|on-demand open| OPEN[.log + .index]
OPEN -->|index lookup| POS[byte position]
POS -->|seek + read| DECODE2[DecodeMessage]
DECODE2 --> CRC{CRC valid?}
CRC -->|yes| RESULT[return ConsumeResult]
CRC -->|no| ERR[return error]
RESULT -->|ConsumeResponse| RESP[protobuf encode]
RESP -->|WriteFrame| CLIENT
ERR -->|ErrorResponse| RESP
Example trace — TCP client consumes from topic "orders", partition 1, offset 7:
1. Client sends: ConsumeRequest{topic:"orders", partition:1, offset:7}
2. Server reads frame
3. Server dispatches to handleConsume
4. Broker.Consume("orders", 1, 7)
5. Topic.Consume: range check OK
6. partition.Read(7): lock, findSegment, read from disk
7. Decode message, verify CRC
8. Server returns: ConsumeResponse{offset:7, key:"user-1", value:"hello"}
9. Server writes frame to TCP connection
flowchart TD
ACCEPT[net.Listener.Accept] -->|new connection| CONN[handleConn goroutine]
CONN --> LOOP[read loop]
LOOP -->|ReadFrame| REQ[decode request]
REQ --> DISPATCH[dispatch to handler]
DISPATCH --> CALL[call Broker method]
CALL --> RESP[encode response]
RESP -->|WriteFrame| LOOP
LOOP -->|EOF or error| DONE[goroutine exits]
DONE -->|wg.Done| WAIT[sync.WaitGroup]
The wire protocol carries six request types. The three metadata requests exist because clients cannot discover topics or partitions from the data path alone — offsets are per-partition, so partition count is essential routing information.
flowchart TD
subgraph Create
C1[CreateTopicRequest] --> C2[Broker.CreateTopic] --> C3[Topic with N partitions]
end
subgraph Discover
L1[ListTopicsRequest] --> L2[Broker.Topics] --> L3["[]string topic names"]
P1[ListPartitionsRequest] --> P2[Broker.NumPartitions] --> P3[int partition count]
end
ListPartitions was added for the demo CLI: mqclient consume --all uses it
to discover the partition range, drain every partition, and merge records by
timestamp (best-effort — ordering is per-partition/per-key only).
flowchart TD
P[Producer] -->|Topic.Produce| PA[Partition.Append]
PA -->|assign offset| SEG[active segment]
SEG -->|Encode| LOG[.log file]
SEG -->|entry at interval| IDX[.index file]
SEG -->|entry at interval| TIDX[.timeindex file]
SEG -->|messageCount++| NEXT{size >= segmentSize?}
NEXT -->|no| DONE[return bytes written]
NEXT -->|yes| ROTATE[Partition.rotate]
ROTATE -->|Close current| CLOSED[segment sealed + synced]
ROTATE -->|Create new| NEWSEG[new active segment]
NEWSEG --> DONE
Durability note: no fsync runs on the append path. Data lands in the OS page cache and is synced only when a segment is closed (rotation) or the partition closes. A process crash loses nothing (page cache survives); an OS crash loses the un-synced tail — the fdatasync-window tradeoff documented in the Kafka study. Recovery truncates torn records via CRC validation.
Example trace — partition-0, baseOffset=0, indexInterval=3 (via WithIndexInterval(3); the default is 1, i.e. index every append):
NextOffset = baseOffset + messageCount = 0 + 0 = 0
Append(key="user-1", value="hello"):
1. offset = NextOffset() = 0 + 0 = 0
2. msg.Offset = 0
3. Encode -> 39 bytes (28 header + 6 key + 5 value)
4. Write .log at position 0
5. messageCount=1, size=39
6. 1 % 3 != 0 -> skip index write
Append(key="user-2", value="world"):
1. offset = NextOffset() = 0 + 1 = 1
2. Encode -> 39 bytes (28 + 6 key + 5 value)
3. Write .log at position 39
4. messageCount=2, size=78
5. 2 % 3 != 0 -> skip index write
Append(key="user-3", value="!"):
1. offset = NextOffset() = 0 + 2 = 2
2. Encode -> 35 bytes (28 + 6 key + 1 value)
3. Write .log at position 78
4. messageCount=3, size=113
5. 3 % 3 == 0 -> write .index: offset=2 -> bytePos=78
6. 3 % 3 == 0 -> write .timeindex: offset=2 -> timestamp=1700000002
flowchart TD
C[Consumer] -->|ConsumeRecord| PR[Partition.Read]
PR -->|findSegment| FS{offset in active segment?}
FS -->|yes| ACTIVE[active segment]
FS -->|no| BINARY[binary search closed list]
BINARY --> CLOSED[opened closed segment]
ACTIVE --> READ[Segment.Read]
CLOSED --> READ
READ -->|on-demand open| OPEN[.log + .index]
READ -->|index lookup| POS[byte position]
READ -->|seek + read| DECODE[DecodeMessage]
DECODE --> CRC{CRC valid?}
CRC -->|yes| RET[return ConsumeRecord]
CRC -->|no| ERR[return ErrCRCMismatch]
Example trace — reading offset 2 (from append trace above):
Read(offset=2):
1. findSegment: 2 >= baseOffset(0) -> active segment
2. Segment.Read(2):
a. Index lookup: binary search .index for offset 2
- Entry found: offset=2, bytePos=78
b. Seek .log to position 78
c. Read 28 bytes (header)
d. Parse: offset=2, size=7, crc=0xABCD, timestamp=1700000002, keylen=6
e. Read 7 bytes (key + value)
f. Decode: key="user-3", value="!"
g. Verify CRC -> valid
3. Return ConsumeRecord{Offset:2, Timestamp:1700000002, Key:"user-3", Value:"!"}
Example trace — reading offset 1 (not indexed):
Read(offset=1):
1. findSegment: active segment
2. Segment.Read(1):
a. Index lookup: binary search .index
- Nearest entry: offset=2, bytePos=78
- But we want offset 1, which is BEFORE offset 2
- Index returns bytePos=0 (start of segment)
b. Seek .log to position 0
c. Read header at 0: offset=0 -> skip (not target)
d. Read header at 39: offset=1 -> found
e. Read body at 39+28=67: key="user-2", value="world"
f. Verify CRC -> valid
3. Return ConsumeRecord{Offset:1, Timestamp:1700000001, Key:"user-2", Value:"world"}
flowchart TD
START[Recover] -->|os.ReadDir| SEGMENTS[segment list]
SEGMENTS -->|all but last| SEAL[RecoverSegment + Release]
SEGMENTS -->|last| OPEN_ACTIVE[recoverActiveSegment]
OPEN_ACTIVE -->|scan .log| SCAN[scanSegment]
SCAN -->|valid messages| VALID[count + size]
VALID -->|truncate .index| TRUNC_IDX[remove stale entries]
VALID -->|truncate .timeindex| TRUNC_TIDX[remove stale entries]
TRUNC_IDX --> READY[partition ready]
TRUNC_TIDX --> READY
SCAN -->|empty file| EMPTY[messageCount = 0]
EMPTY --> READY
Example trace — recovery after crash:
Partition directory:
00000000000000000000.log (3 messages, 99 bytes)
00000000000000000000.index (1 entry: offset=2, bytePos=68)
00000000000000000000.timeindex (1 entry: offset=2, timestamp=1700000002)
00000000000000000099.log (partial message, 15 bytes - crash mid-write)
00000000000000000099.index (0 entries)
00000000000000000099.timeindex (0 entries)
Recover:
1. os.ReadDir: [00000000000000000000, 00000000000000000099] (already sorted)
2. RecoverSegment 00000000000000000000, then Release (free fds)
3. recoverActiveSegment: 00000000000000000099
4. scanSegment:
- Read header at 0: only 15 bytes available (need 28)
- Partial message -> stop
- validCount = 0, validSize = 0
5. Truncate .index to 0 entries
6. Truncate .timeindex to 0 entries
7. Reopen index and timeindex
8. Partition ready: baseOffset=99, messageCount=0
9. NextOffset = 99 + 0 = 99 (next message gets offset 99)
flowchart TD
A[Segment.Append] -->|write .log| WLOG{write ok?}
WLOG -->|no| ERR1[return error: write to segment]
WLOG -->|yes| WIDX[write .index]
WIDX -->|write ok?| WIDX_OK{ok?}
WIDX_OK -->|no| ERR2[return error: write index]
WIDX_OK -->|yes| WTIDX[write .timeindex]
WTIDX -->|write ok?| WTIDX_OK{ok?}
WTIDX_OK -->|no| ERR3[return error: write time index]
WTIDX_OK -->|yes| OK[return bytes written]
Note: syncSegment (fdatasync with EINTR retry, 3 attempts) is NOT called during Append. It runs only in Close — on rotation and on partition close.
Example trace — successful append:
Append(key="k", value="v"):
1. Encode -> 34 bytes
2. Write .log: 34 bytes at position 99 -> ok
3. Write .index: offset=0, bytePos=99 -> ok (interval=1: every append)
4. Write .timeindex: offset=0, timestamp -> ok
5. Return 34 bytes written; size and messageCount updated
Example trace — write failure:
Append(key="k", value="v"):
1. Write .log: 34 bytes at position 99 -> ok
2. Write .index: offset=0, bytePos=99 -> DISK FULL
3. Return error: "write index: write: no space left on device"
4. .log has 34 bytes, .index has 0 entries (inconsistent)
5. Next Append will fail (offset collision) or recovery will fix