Skip to content

Latest commit

 

History

History
315 lines (274 loc) · 11.2 KB

File metadata and controls

315 lines (274 loc) · 11.2 KB

Control Flow Diagrams

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).

Full Stack: Produce

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
Loading

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

Full Stack: Consume

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
Loading

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

Server Connection Lifecycle

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]
Loading

Metadata Flows: Discovery

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
Loading

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).

Normal Flow: Append (partition level)

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
Loading

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

Normal Flow: Read (partition level)

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]
Loading

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"}

Failure Mode: Recovery on Startup (partition level)

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
Loading

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)

Failure Mode: Write Errors (partition level)

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]
Loading

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