-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathdoc.go
More file actions
93 lines (93 loc) · 3.72 KB
/
Copy pathdoc.go
File metadata and controls
93 lines (93 loc) · 3.72 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
// Package log implements an append-only, segment-based message storage engine
// inspired by Apache Kafka's architecture.
//
// The core insight from the Kafka paper is that append-only sequential I/O is
// 10-100x faster than random writes. Every other optimization (zero-copy,
// batching, OS page cache) is secondary to this one choice.
//
// # Entry Point
//
// Broker is the entry point for external packages. Topic, partition, and
// segment are unexported implementation details. External packages interact
// with the Broker via New, CreateTopic, Produce, Consume, and Close.
//
// Topic constructors (createTopic, recoverTopic) are unexported — only the
// Broker calls them. This ensures all topic lifecycle is managed through
// the Broker, which coordinates map access and topic locking.
//
// # Storage Layout
//
// A partition is a directory of segment files. Each segment stores a contiguous
// range of messages in an append-only .log file, backed by sparse index files
// for fast offset and timestamp lookups.
//
// partition-0/
// 00000000000000000000.log # messages
// 00000000000000000000.index # offset -> byte position
// 00000000000000000000.timeindex # offset -> timestamp
// 00000000000000150000.log
// 00000000000000150000.index
// 00000000000000150000.timeindex
//
// # Message Format
//
// Each message is encoded as a fixed 28-byte header followed by variable-length
// key and value:
//
// | offset (8B) | size (4B) | crc (4B) | timestamp (8B) | keylen (4B) | key | value |
//
// The CRC-32 covers timestamp, key, and value for bit-rot detection. Offset and
// timestamp are broker-assigned — the producer only supplies key and value.
//
// # Concurrency Model
//
// Two-layer locking with clean separation of concerns:
//
// Layer 1 — Broker map (Broker.mu, sync.RWMutex):
// - Protects the topics map only. Held briefly for map lookup/insert.
// - No disk I/O under this lock.
//
// Layer 2 — Per-topic (Topic.mu, sync.RWMutex):
// - Protects the topic during operations (Produce, Consume, etc.) and
// lifecycle (close).
// - Produce/Consume acquire mu.RLock — multiple operations can proceed
// concurrently.
// - close acquires mu.Lock — waits for in-flight operations to complete
// before closing file handles.
//
// This ensures a topic cannot be closed while Produce/Consume is in flight,
// without holding the broker lock during disk I/O.
//
// Parallelism comes from multiple partitions (single writer per partition),
// not from splitting writes within one. This avoids lock-free complexity
// (gap problems on crash, pwrite requirements).
//
// # Dependency Flow
//
// Dependencies flow one direction: Broker -> Topic -> partition -> segment -> file.
// Each layer owns a single concern:
//
// - Broker (exported): topic management, map coordination, request routing
// - Topic (unexported): naming, partition routing (key hash, round-robin)
// - partition (unexported): offset assignment, segment rotation, recovery
// - segment (unexported): byte-level storage, CRC validation, index management
//
// No layer depends on a higher layer. This makes each layer independently
// testable and replaceable.
//
// # Directory Ownership
//
// Only two os.MkdirAll calls total:
//
// - createTopic creates data/<topic>/
// - createPartition creates data/<topic>/<partition>/
//
// Everything below partition operates on files within an existing directory.
//
// # Recovery
//
// On startup, each segment is scanned forward. Valid messages are counted, then
// the .log file, .index, and .timeindex are truncated to remove any partial or
// corrupted trailing data. Segments are immutable after Close() by Go-level
// convention — file handles are nilled to enforce this.
package log