ShardMaster is a high-throughput distributed database sharding proxy, SQL schema and query execution engine, Change Data Capture (CDC) VReplication streamer, and interactive terminal control center written in Pure Go (Golang 1.22+). Modeled after production cloud-native sharding architectures such as Vitess (PlanetScale) and Citus, ShardMaster presents a cluster of isolated physical shards as a single logical database listening on port 6000 via the native PostgreSQL Frontend/Backend Wire Protocol v3.0 (PGWire).
____ _ _ _ ____ ____ __ __ _ ____ _____ _____ ____
/ ___|| | | | / \ | _ \| _ \| \/ | / \ / ___|_ _| ____| _ \
\___ \| |_| | / _ \ | |_) | | | | |\/| | / _ \ \___ \ | | | _| | |_) |
___) | _ |/ ___ \| _ <| |_| | | | |/ ___ \ ___) || | | |___| _ <
|____/|_| |_/_/ \_\_| \_\____/|_| |_/_/ \_\____/ |_| |_____|_| \_\
========================================================================
CLUSTER: 4 Shards (1,024 Buckets) | ROWS: 10,000 | PGWire :6000 ONLINE
========================================================================
- The 6 Architectural Pillars
- System Architecture and Mermaid Diagrams
- 2.1 High-Level Control Plane and Data Plane Topology
- 2.2 Pillar 1: Native PostgreSQL v3.0 Wire Protocol Lifecycle
- 2.3 Pillar 2: Distributed Scatter-Gather and Streaming K-Way Merge
- 2.4 Pillar 3 and 4: CDC VReplication Streamer and VDiff Atomic Cutover
- 2.5 Pillar 5: Autonomous EWMA Hotspot Detection State Machine
- Mathematical Foundations and L1-Cache Memory Engineering
- Distributed SQL and Schema Engine
- Hardware Efficiency and Live Benchmark Results
- Interactive Control Center and CLI Reference
- 4-Tab Interactive Terminal UI (TUI) Guide
- Project Directory Structure
| Pillar | Component | Traditional Student Implementation | ShardMaster Production Implementation |
|---|---|---|---|
| Pillar 1 | Client Interface | Basic HTTP REST JSON wrapper (curl) |
Native PostgreSQL v3.0 Wire Protocol (PGWire on :6000) via jackc/pgproto3/v2 plus :8080/shard?user_id=123 HTTP compatibility bridge. |
| Pillar 2 | Non-Key Queries | Unsupported or loads entire tables into RAM |
Distributed Scatter-Gather with Streaming K-Way Merge Sort and Two-Phase Map-Reduce GROUP BY using worker Goroutines, bounded channels, and a container/heap Priority Queue in |
| Pillar 3 | In-Flight Sync | Naive dual-writing (vulnerable to split-brain) |
Vitess-Style Change Data Capture (CDC) Stream Engine: Single-source writes, transactional _shardmaster_cdc mutation log, microsecond replication lag tracking, and <200us zero-lag atomic pointer cutover. |
| Pillar 4 | Data Verification | Simple SELECT COUNT(*) row count check |
Cryptographic Bit-Level Parity (VDiff): Streaming, order-independent 256-bit commutative XOR of SHA-256 canonical row digests across source and target shards. |
| Pillar 5 | Workload Intelligence | Static manual shard splits only |
Autonomous EWMA Hotspot Detector: Cache-line padded per-bucket Exponentially Weighted Moving Average frequency tracker that automatically isolates hot buckets (>98th percentile load) onto cold shards. |
| Pillar 6 | Operator Experience | Unstructured scrolling log lines | Animated Interactive Control Center, 4-Tab Minimalist TUI, Typed PostgreSQL Grid Renderer, 500M+ QPS Benchmark, and Consistent-Hash Split Analyzer. |
flowchart TB
subgraph ClientLayer["Client and Operator Layer"]
PSQL["PostgreSQL Clients (psql, DBeaver, pgAdmin) on TCP :6000"]
HTTPClient["HTTP Directory Client on GET :8080/shard"]
OperatorShell["Interactive Control Center and 4-Tab TUI"]
end
subgraph ProxyCore["SHARDMASTER PROXY CORE (Go 1.22+ Static Binary)"]
PGWire["Pillar 1: PGWire v3.0 Protocol Server"]
Lexer["Zero-Alloc SQL Lexer, AST Planner, and Schema Catalog"]
EWMA["Pillar 5: Lock-Free EWMA Hotspot Tracker (1,024 Padded Counters)"]
Directory["O(1) Atomic Bucket Ring ([1024]atomic.Uint32 in 4 KB L1 Cache)"]
KWay["Pillar 2: Scatter-Gather K-Way Merge and Map-Reduce Engine"]
CDC["Pillar 3: CDC VReplication Streamer and Atomic Cutover Gate"]
VDiff["Pillar 4: Cryptographic 256-Bit XOR-SHA256 VDiff Engine"]
end
subgraph StoragePlane["Physical Database Shards (1,024 Buckets + Local NVMe Persistence)"]
S0["Shard 0 (:5432) | us-west | Buckets 0..255"]
S1["Shard 1 (:5433) | us-east | Buckets 256..511"]
S2["Shard 2 (:5434) | eu-central | Buckets 512..767"]
S3["Shard 3 (:5435) | ap-south | Buckets 768..1023"]
S4["Shard 4..7 (:5436-:5439) | Dynamic Scale-Out Targets"]
end
PSQL -->|"PGWire v3.0 TCP :6000"| PGWire
HTTPClient -->|"HTTP GET :8080"| Directory
OperatorShell -->|"Live Control and Queries"| ProxyCore
PGWire --> Lexer
Lexer --> EWMA
EWMA --> Directory
Directory -->|"O(1) Point Route in 18 ns"| S0
Directory -->|"O(1) Point Route in 18 ns"| S1
Directory -->|"O(1) Point Route in 18 ns"| S2
Directory -->|"O(1) Point Route in 18 ns"| S3
Lexer -->|"Scatter-Gather or GROUP BY"| KWay
KWay -->|"Parallel Worker Goroutines"| S0
KWay -->|"Parallel Worker Goroutines"| S1
KWay -->|"Parallel Worker Goroutines"| S2
KWay -->|"Parallel Worker Goroutines"| S3
CDC -.->|"Keyset Backfill and CDC LSN Stream"| S4
VDiff -.->|"256-Bit XOR-SHA256 Parity Audit"| S4
sequenceDiagram
autonumber
participant Client as psql or DBeaver Client
participant PGWire as PGWire Server (:6000)
participant Lexer as SQL Lexer and Catalog
participant Dir as L1 Bucket Ring (4 KB)
participant Shard as Target Physical Shard
Client->>PGWire: TCP Connect and SSLRequest
PGWire-->>Client: Decline SSL ('N') for Local Wire Handshake
Client->>PGWire: StartupMessage (Protocol v3.0, user=admin, db=shardmaster)
PGWire-->>Client: AuthenticationOk + ParameterStatus + ReadyForQuery
Client->>PGWire: Query SELECT * FROM users WHERE user_id = 42
PGWire->>Lexer: ClassifySQL and Extract Predicate (user_id = 42)
Lexer-->>PGWire: ClassifiedQuery (QueryPointSelect, Key = 42)
PGWire->>Dir: LookupFast("42") via xxHash64("42") bitwise-AND 1023
Dir-->>PGWire: Bucket 697 mapped to Shard 2 (:5434) in 18 ns
PGWire->>Shard: Execute O(1) Slab and Delta Lookup on Shard 2
Shard-->>PGWire: Return Matching UserRow Tuple
PGWire-->>Client: RowDescription + DataRow + CommandComplete + ReadyForQuery
When a client executes a query without a single user_id equality predicate (for example, SELECT * FROM users WHERE email LIKE '%@gmail.com' ORDER BY created_at DESC LIMIT 5), ShardMaster avoids loading entire tables into memory by streaming ordered batches across bounded Go channels into a Min-Heap Priority Queue:
flowchart LR
Query["Non-Key SQL Query"] --> Planner["Distributed Query Planner"]
subgraph FanOut["1. Parallel Goroutine Scatter"]
Planner -->|"Worker 0"| S0["Shard 0 (:5432) Local Top-K"]
Planner -->|"Worker 1"| S1["Shard 1 (:5433) Local Top-K"]
Planner -->|"Worker 2"| S2["Shard 2 (:5434) Local Top-K"]
Planner -->|"Worker 3"| S3["Shard 3 (:5435) Local Top-K"]
end
subgraph Channels["2. Bounded Channels (cap=16)"]
S0 --> C0["Stream 0"]
S1 --> C1["Stream 1"]
S2 --> C2["Stream 2"]
S3 --> C3["Stream 3"]
end
C0 --> Heap["3. container/heap Min-Heap K-Way Merge"]
C1 --> Heap
C2 --> Heap
C3 --> Heap
Heap --> Result["4. Global Top-K Sorted ResultSet"]
sequenceDiagram
autonumber
participant App as Concurrent App Writes
participant Dir as Atomic Bucket Directory
participant Source as Source Shard (Shard 0)
participant CDC as CDC VReplication Engine
participant Target as Target Shard (Shard 4)
Note over App,Source: Initial State: Buckets 128..255 owned by Shard 0
App->>Dir: INSERT or UPDATE user in Bucket 200
Dir->>Source: Apply Write and Append Monotonic LSN to _shardmaster_cdc
Note over CDC,Target: Phase 1: Lock-Free Keyset Backfill + CDC Stream
CDC->>Dir: SetBucketState(128..255, CDC_STREAMING)
CDC->>Source: Export Columnar Bucket Slab and Snapshot LSN Watermark
Source-->>CDC: Bucket Slab Rows + Watermark LSN
CDC->>Target: Install Bucket Slab on Shard 4
CDC->>Source: FetchCDCMutationsAfter(Watermark LSN)
Source-->>CDC: In-Flight Delta Mutations
CDC->>Target: Replay Delta Mutations and Advance LSN
Note over CDC,Target: Phase 2: Sub-Millisecond Cutover Gate and VDiff Audit
CDC->>Dir: SetBucketState(128..255, CUTOVER_GATE)
CDC->>Source: Drain Final Tail Mutations (Replication Lag = 0.00 ms)
CDC->>Source: Compute 256-Bit Commutative XOR-SHA256 Digest
CDC->>Target: Compute 256-Bit Commutative XOR-SHA256 Digest
CDC->>CDC: Verify SourceDigest == TargetDigest
Note over Dir,Target: Phase 3: Single-Instruction Atomic Pointer Swap
CDC->>Dir: AtomicCutoverBucket(128..255, Shard 4) via atomic.Uint32.Store
CDC->>Source: Purge Migrated Bucket Slabs from Shard 0
App->>Dir: Subsequent Read or Write for Bucket 200
Dir->>Target: Routed Directly to Shard 4 with 0.00 ms Downtime
stateDiagram-v2
[*] --> Monitoring
Monitoring --> Evaluating: Tick Every 250ms to Compute EWMA Decay
Evaluating --> Monitoring: All 1,024 Buckets Below 5x Cluster Mean QPS
Evaluating --> HotspotDetected: Bucket 412 Exceeds 1,500 QPS and 5x Cluster Mean
HotspotDetected --> TargetSelection: Select Coldest Physical Shard by QPS and Buckets
TargetSelection --> LiveCDCMigration: Stream Hot Bucket 412 via CDC and Verify VDiff
LiveCDCMigration --> AtomicIsolation: Swap atomic.Uint32 Pointer for Bucket 412
AtomicIsolation --> Monitoring: Hotspot Isolated with Zero Dropped Queries
Every shard key xxHash (github.com/cespare/xxhash/v2) and mapped onto xxHash64(k) & 1023):
The directory table is stored as a fixed-size contiguous array [1024]atomic.Uint32:
Because 0 B/op, 0 allocs/op).
To prove zero data corruption across arbitrary row orderings without sorting millions of rows in memory, ShardMaster computes a 256-bit commutative XOR accumulator over canonical row SHA-256 digests across bucket range
Because bitwise XOR (
Each virtual bucket
Whenever
ShardMaster includes a complete distributed schema catalog (pkg/router/schema.go) and typed PostgreSQL grid renderer (pkg/console/shell.go). Every query result displays the SQL statement, distributed execution plan (PLAN >), column names, PostgreSQL data types, and execution latency footer:
--- TABLE SCHEMA DEFINITION: PUBLIC.USERS ------------------------------
SQL > DESCRIBE users;
PLAN > Catalog Schema Resolution for 'public.users' (11 columns, 5 indexes)
+---------+---------------+---------------+-------------+-------------------------+----------------------------+----------------------------+
| ordinal | column_name | data_type | nullable | key_constraint | default_expr | storage_encoding |
| INT2 | VARCHAR(64) | VARCHAR(32) | VARCHAR(12) | VARCHAR(28) | VARCHAR(32) | VARCHAR(32) |
+=========+===============+===============+=============+=========================+============================+============================+
| 1 | user_id | BIGINT | NOT NULL | PRIMARY KEY (SHARD KEY) | nextval('users_id_seq') | 64-Bit Integer Ring Key |
| 2 | user_key | VARCHAR(64) | NOT NULL | HASH RING KEY | CAST(user_id AS TEXT) | Inline UTF-8 Key |
| 3 | name | VARCHAR(128) | NOT NULL | NONE | '' | Dictionary + Delta Overlay |
| 4 | email | VARCHAR(255) | NOT NULL | LOCAL INDEX | '' | Trigram Indexed String |
| 5 | tenant_id | VARCHAR(64) | NOT NULL | LOCAL INDEX | 'tenant_core' | Interned Low-Cardinality |
| 6 | region | VARCHAR(32) | NOT NULL | PARTITION KEY | 'us-west' | Geo-Placement Tag |
| 7 | balance_cents | BIGINT | NOT NULL | COLUMNAR SLAB | 250000 | Pointer-Free []uint32 Slab |
| 8 | balance_usd | NUMERIC(12,2) | NOT NULL | GENERATED VIRTUAL | (balance_cents / 100.0) | Computed Currency Column |
| 9 | bucket_id | SMALLINT | NOT NULL | BUCKET INDEX [0..1023] | (xxhash64(user_id) & 1023) | 10-Bit Virtual Bucket ID |
| 10 | created_at | TIMESTAMPTZ | NOT NULL | MERGE SORT KEY (DESC) | CURRENT_TIMESTAMP | 64-Bit UTC Epoch Micros |
| 11 | updated_at | TIMESTAMPTZ | NOT NULL | CDC LSN TRACKED | CURRENT_TIMESTAMP | 64-Bit UTC Epoch Micros |
+---------+---------------+---------------+-------------+-------------------------+----------------------------+----------------------------+
* Table: public.users | Type: SHARDED TABLE | Shard Key: user_id (BIGINT) | Strategy: HASH (xxHash64 & 1023) (1024 Buckets)
* Index [pk_users_user_id]: HASH_RING_PK ON (user_id) | Scope: O(1) SINGLE SHARD (UNIQUE)
* Index [idx_users_created_at_desc]: BTREE_DESC ON (created_at DESC, user_id DESC) | Scope: SCATTER K-WAY MERGE (NON-UNIQUE)
[OK] (11 rows) | Tag: DESCRIBE 11 | Route: SCHEMA_CATALOG | Time: 0.018 ms (18 us)
ShardMaster pairs its pointer-free columnar slab engine with an embedded pure-Go ANSI/PostgreSQL relational execution engine (pkg/router/sql_engine.go) and atomic local NVMe state persistence (data/cluster_state.json), giving you 100% SQL support including multi-table JOINs, Window Functions, Common Table Expressions (WITH / WITH RECURSIVE), Subqueries, GROUP BY ... HAVING, Views, Triggers, ALTER TABLE, explicit ACID transactions (BEGIN / COMMIT / ROLLBACK), custom deterministic sharding scalar functions (xxhash64(), virtual_bucket(), target_shard()), and live byte-level shard control DDL:
-- 1. Schema & DDL Catalog Introspection (8 Built-In Tables/Views + Custom DDL)
SHOW TABLES;
DESCRIBE users;
DESCRIBE orders;
DESCRIBE payments;
DESCRIBE vip_users_view;
DESCRIBE _shardmaster_cdc;
SHOW CREATE TABLE users;
SHOW INDEXES;
SHOW SCHEMAS;
-- 2. Custom DDL & Schema Evolution (CREATE / ALTER / DROP TABLE, VIEW, INDEX, TRIGGER)
-- All rows inserted into custom tables are hashed to their virtual bucket [0..1023],
-- billed to the owning shard's UsedMemoryBytes(), verified by VDiff, and migrated by CDC!
CREATE TABLE invoices (invoice_id BIGINT PRIMARY KEY, user_id BIGINT NOT NULL, amount_usd NUMERIC(12,2) DEFAULT 99.50);
INSERT INTO invoices (invoice_id, user_id, amount_usd) VALUES (9001, 42, 1450.75), (9002, 777, 3200.00);
ALTER TABLE invoices ADD COLUMN status VARCHAR(32) DEFAULT 'PAID';
DESCRIBE invoices;
SELECT i.invoice_id, u.name, i.amount_usd, i.status FROM invoices i INNER JOIN users u ON u.user_id = i.user_id;
DROP TABLE invoices;
-- 3. Live Shard Customization, Exact Byte Quotas, Draining & Shard-Pinned SQL
CREATE SHARD ALIAS='eu-west-micro' REGION='eu-west' PORT=5440 BYTES=16777216 SLAB_BYTES=16384 POOL_BYTES=4194304 TIER='Custom-16MB' WEIGHT=100 BUCKETS=128;
ALTER SHARD 0 SET ALIAS='us-west-nvme-0', BYTES=67108864, SLAB_BYTES=65536, POOL_BYTES=16777216, WEIGHT=150, BUCKETS=256, MODE='READ_WRITE';
DRAIN SHARD 3;
REBALANCE SHARDS BY WEIGHT;
/*+ SHARD(0) */ SELECT COUNT(*), SUM(balance_usd), AVG(balance_usd) FROM users;
SELECT * FROM users ON SHARD 1 LIMIT 5;
-- 4. Multi-Statement ACID Transactions (BEGIN / COMMIT / ROLLBACK)
BEGIN;
UPDATE users SET balance_cents = balance_cents - 50000 WHERE user_id = 42;
UPDATE users SET balance_cents = balance_cents + 50000 WHERE user_id = 100;
COMMIT;
-- 5. Cluster Topology, Virtual Buckets, CDC & VDiff Diagnostics
SHOW SHARDS;
SHOW BUCKETS;
SHOW CDC;
SHOW HOTSPOTS;
SHOW STATS;
RUN VDIFF;
EXPLAIN ANALYZE SELECT * FROM users WHERE user_id = 42;
-- 6. O(1) Point Lookups, Column Projection & Multi-Key Batch Routing
SELECT * FROM users WHERE user_id = 42;
SELECT user_id, name, email, balance_usd FROM users WHERE user_id IN (42, 100, 777, 8888, 9999);
SELECT * FROM users WHERE user_id BETWEEN 100 AND 105;
-- 7. Distributed Scatter-Gather K-Way Merge & Map-Reduce GROUP BY
SELECT * FROM users WHERE email LIKE '%@gmail.com' ORDER BY created_at DESC LIMIT 5;
SELECT * FROM users WHERE region = 'us-west' ORDER BY created_at DESC LIMIT 5;
SELECT region, COUNT(*), SUM(balance_usd), AVG(balance_usd) FROM users GROUP BY region;
SELECT shard_id, COUNT(*), AVG(balance_usd) FROM users GROUP BY shard_id;
SELECT COUNT(*), SUM(balance_usd), AVG(balance_usd), MIN(balance_usd), MAX(balance_usd) FROM users;
-- 8. Co-Located Multi-Table JOINs, Window Functions, CTEs, Subqueries & Shard Functions
SELECT u.shard_id, u.user_id, u.name, o.order_id, o.product_name, o.amount_usd, p.payment_method
FROM users u
INNER JOIN orders o ON u.user_id = o.user_id
INNER JOIN payments p ON o.order_id = p.order_id
ORDER BY o.amount_usd DESC LIMIT 6;
SELECT shard_id, user_id, name, region, balance_usd,
RANK() OVER (PARTITION BY region ORDER BY balance_cents DESC) AS regional_rank
FROM users LIMIT 8;
WITH high_value AS (SELECT * FROM users WHERE balance_cents >= 500000)
SELECT region, COUNT(*) AS vip_users, ROUND(AVG(balance_usd), 2) AS avg_vip_usd, MAX(balance_usd) AS max_vip_usd
FROM high_value GROUP BY region HAVING COUNT(*) >= 5 ORDER BY avg_vip_usd DESC;
SELECT user_id, name, balance_usd,
xxhash64(user_id) AS xxhash64_hex,
virtual_bucket(user_id) AS bucket_id,
target_shard(user_id) AS routed_shard
FROM users
WHERE balance_cents > (SELECT AVG(balance_cents) FROM users)
ORDER BY balance_cents DESC LIMIT 6;
-- 9. Point Writes (INSERT / UPDATE / DELETE) & Live CDC Log Queries
INSERT INTO users (user_id, name, email, balance_cents) VALUES (42, 'Ada Lovelace', 'ada@gmail.com', 950000);
UPDATE users SET name = 'Grace Hopper', balance_usd = 12500.00 WHERE user_id = 42;
DELETE FROM users WHERE user_id = 100;
SELECT * FROM _shardmaster_cdc LIMIT 6;ShardMaster features a production-grade universal CSV/TSV import and export engine (pkg/router/csv_engine.go):
- Auto-Inferred Schema & Delimiter Detection:
- Automatically detects delimiters: comma (
,), tab (\t), semicolon (;), and pipe (|). - Dynamically scans data rows to infer column types:
BIGINT,NUMERIC(12,2),TIMESTAMPTZ,BOOLEAN, andTEXT. - Automatically detects optimal primary shard keys (
user_id,id,uuid,<table_name>_id,key). - Automatically provisions distributed tables with 1,024 Virtual Buckets in
SchemaCatalogand the relational engine if the table does not exist.
- Automatically detects delimiters: comma (
- Direct CLI Import:
.\shardmaster.exe import employees.csv .\shardmaster.exe import orders.tsv custom_orders --delimiter "\t"
- Interactive Shell Import:
shardmaster [4-shards] > import customers.csv
- PostgreSQL Standard SQL
COPYandIMPORT CSV:-- Import from CSV file on disk: COPY customers FROM 'C:/data/customers.csv' WITH (FORMAT CSV, HEADER); IMPORT CSV 'C:/data/orders.csv' INTO orders SHARD BY order_id DELIMITER ','; -- Export distributed table or query result to CSV file: COPY customers TO 'C:/exports/customers.csv' WITH (FORMAT CSV, HEADER); COPY (SELECT region, COUNT(*), AVG(balance_usd) FROM users GROUP BY region) TO 'C:/exports/regional_summary.csv';
- Streaming
COPY ... FROM STDIN© ... TO STDOUT(psql\copyand DBeaver Import Wizard):- Full support for PostgreSQL streaming copy protocol (
CopyInResponse,CopyData,CopyDone,CopyOutResponse).
- Full support for PostgreSQL streaming copy protocol (
ShardMaster works out-of-the-box as a drop-in PostgreSQL server with any external database management tool or programming language:
| Tool / Driver | Compatibility Level | Supported Features |
|---|---|---|
| DBeaver | 100% Native | Full schema navigation, table introspections (pg_catalog.pg_tables), Data Editor, CSV Import Wizard (COPY ... FROM STDIN), SQL console |
| DataGrip / IntelliJ | 100% Native | Database Explorer tree, information_schema.tables, syntax inspection, prepared statements, multi-statement queries |
| pgAdmin 4 | 100% Native | Object Browser, server dashboard, SQL query tool, schema/view inspection |
| VS Code PostgreSQL | 100% Native | Connection explorer, query runner, auto-complete against catalog |
| psql CLI | 100% Native | Meta-commands (\dt, \d, \di, \dn, \l), \copy file streaming, interactive query prompt |
| JDBC / Python psycopg2 / Go pgx | 100% Native | Extended Query Protocol (Parse, Describe, Bind, Execute, Sync), parameter substitution ($1, $2, ...), connection pooling |
- Host:
localhost(or127.0.0.1) - Port:
6000 - Database:
shardmaster(orpostgres) - Username:
shardmaster_admin(or any username) - Password: (any password / blank)
- SSL Mode:
disable(orallow)
-
100% Real Hashed Row Storage (
SeededUIDs+SlabBalances): Every seeded user row (1..N) is individually hashed viaxxHash64into its exact virtual bucket (0..1023), stored 1-to-1 in sortedSeededUIDs []int64andSlabBalances []uint32, and looked up via$O(\log N)$ binary search (findSeededSlot)—with zero collisions and zero synthetic row fabrication. -
Atomic Local NVMe State Persistence (
data/cluster_state.json): Every cluster topology change, bucket ownership map, shard byte quota,DeltaOverridesmutation,DeletedIDstombstone, custom SQL table row, and_shardmaster_cdcLSN journal is atomically persisted todata/cluster_state.jsonvia write-temp-and-rename and automatically restored on startup. -
Zero Data Loss When Shrinking Shard
BYTES(EnforceShardCapacityAndBuckets): Shrinking a shard'sMaxCapacityBytesbelow its liveUsedMemoryBytes()never truncates committed rows or populatedSlabBalancesin-place. Instead, ShardMaster calculates the exact number of virtual buckets to evacuate, verifies free byte capacity across remaining writable shards, and streams those buckets out via CDC VReplication + VDiff—or rejects cleanly withERR_INSUFFICIENT_CLUSTER_CAPACITYbefore modifying a single byte. -
Automatic Bucket Evacuation & Hard Quota Rejection (
SQLSTATE 53100): EveryINSERTandUPDATE(both onusersand on customCREATE TABLEtables) reserves exact row bytes against the owning shard'sMaxCapacityBytes. If a shard reaches 100% byte capacity,AutoEvacuateForWriteautomatically evacuates non-active buckets to writable shards with free capacity; if the cluster has no free capacity, the write is rejected atomically withSQLSTATE 53100 (disk_full): ERR_SHARD_CAPACITY_EXCEEDEDand rolled back with zero partial state. -
Lock-Ordered Per-Bucket Cutover Write Gate:
ShardDirectoryguards each of the 1,024 virtual buckets with a read-write gate acquired in strictly ascending bucket order (0..1023), making deadlocks mathematically impossible while guaranteeing zero stranded writes during<200usatomic pointer cutovers.
Verified on a 16 GB DDR5 RAM Windows workstation (go1.22+ windows/amd64, Ryzen 5 7535HS 12 logical CPU threads):
| Metric | Measured Value | Engineering Mechanism |
|---|---|---|
| Default Seeded Cluster Dataset | 10,000 Real Hashed Rows (Configurable) | Exact xxHash64 bucket placement (SeededUIDs + SlabBalances) + Delta Overlay + NVMe persistence |
| Multi-Core Routing Throughput | 877,368,697 req/sec | Lock-free [1024]atomic.Uint32 + xxHash64 + padded EWMA counters across 12 CPU threads |
| Routing Latency (P50 / P99) | 11 ns / 28 ns | L1-cache resident 4 KB lookup table, 0 syscalls, 0 mutex locks |
| Heap Allocations on Hot Path | 0 B/op, 0 allocs/op | Stack-allocated key buffer, zero GC pressure |
| Resharding Availability (4 -> 8 Shards) | 100.000% (0.00 ms Downtime) | All rows migrated via Keyset Backfill + CDC Stream + <200us atomic pointer swap |
| Consistent-Hash Split Efficiency | 78.8% Network I/O Saved (64 -> 80 Shards) |
Virtual bucket indirection moves only 20.0% of buckets (205/1024) vs. 98.8% with naive modulo |
Double-clicking shardmaster.exe in Windows Explorer (or running .\shardmaster.exe in PowerShell/CMD) boots the entire distributed cluster (PGWire :6000, HTTP :8080, 4 Physical Shards with persisted local NVMe state, and the EWMA Hotspot Monitor) with animated - / | \ loading spinners and launches the Interactive Control Center:
[1] Interactive Academy Step-by-step guided tour of all 6 Pillars
[2] Cluster Status View shard names, sizes, load bars & CDC lag
[3] Live Dashboard (TUI) 4-Tab TUI + Live Shard Customizer & Resizer
[4] Route User Key Inspect O(1) xxHash64 & Virtual Bucket
[5] 100% Full SQL Engine Schemas, DDL, JOINs, CTEs, Window Funcs & 26 Presets
[6] Customize & Add Shards Custom names, disk sizes, live resize, drain & pin SQL
[7] Zero-Downtime Split Split cluster (4 -> 8 shards) with VDiff
[8] Hotspot Self-Healer Spike Bucket #412 & watch auto-isolation
[9] VDiff Parity Audit Verify 256-bit XOR-SHA256 across shards
[10] 500M+ QPS Benchmark Multi-core lock-free routing & chaos test
[11] Scale-Out Split Bench Benchmark N->M shard split & movement vs modulo
[12] Architecture Manual Formulas, internals & psql connection guide
------------------------------------------------------------------------
Quick Commands: demo | schema | reset | menu | exit
SQL Editor Keys: [Shift+Enter] New Line | [Tab] Indent 4 Spaces | [Up/Down] History
Both the main shardmaster [4-shards] > prompt and Option [5] (100% Full SQL Engine) feature a native Win32 key-event multi-line SQL editor (pkg/console/editor_windows.go):
[Shift+Enter](or[Ctrl+Enter]/[Alt+Enter]): Moves to a new continuation line (.. >) without executing prematurely, allowing you to write multi-line SQL queries or multi-statement scripts.[Tab]&[Shift+Tab]: Inserts or removes 4 spaces of indentation aligned to 4-column tab stops.- Smart Block Auto-Indentation: Automatically indents 4 spaces after
(or SQL clause headers (CREATE TABLE (,SELECT,FROM,WHERE,VALUES,GROUP BY,ORDER BY) and automatically un-indents closing);to align with the opening statement. [Up]/[Down]/[Left]/[Right]/[Home]/[End]: Full cursor movement and recall of previously executed SQL statements.- Multi-Statement Batch Scripts: Type or paste multiple
;-separated SQL statements (CREATE TABLE ...; INSERT INTO ...; SELECT ...;) in one block and press[Enter]to execute the entire batch sequentially.
# Launch Interactive Control Center (Default):
.\shardmaster.exe
# Run Complete 6-Pillar Automated Showcase:
.\shardmaster.exe demo
# Inspect Table Schemas or Execute Any SQL Query (JOINs, CTEs, Window Functions, DDL):
.\shardmaster.exe schema
.\shardmaster.exe query "SHOW TABLES;"
.\shardmaster.exe query "DESCRIBE users;"
.\shardmaster.exe query "SELECT u.shard_id, u.user_id, u.name, o.order_id, o.product_name, o.amount_usd, p.payment_method FROM users u INNER JOIN orders o ON u.user_id = o.user_id INNER JOIN payments p ON o.order_id = p.order_id ORDER BY o.amount_usd DESC LIMIT 6;"
# Perform O(1) Atomic Bucket Ring Lookup:
.\shardmaster.exe lookup 42
# Add Shard, Split Cluster, or Run Cryptographic VDiff:
.\shardmaster.exe add-shard --region us-west
.\shardmaster.exe rebalance --target-num-shards 8
.\shardmaster.exe vdiff
# Run Multi-Million QPS Benchmark & Consistent-Hash Split Analyzer:
.\shardmaster.exe bench --duration 3
.\shardmaster.exe scale-sim --from-shards 64 --to-shards 80Pressing 3 in the Control Center (or running .\shardmaster.exe tui) opens the minimalist 4-Tab Terminal UI (76x17), engineered to fit 100% inside any standard 80x24 terminal window without line-wrapping or vertical clipping:
- Tab
[1] Topology: Live physical shard load bars, custom aliases, regions, exact byte usage (USED/MAX), bucket counts, row counts, and rolling QPS. - Tab
[2] CDC & VDiff: Real-time Vitess CDC VReplication progress bar, replication lag (ms), andXOR-SHA256VDiff digests. - Tab
[3] Hotspots: Autonomous EWMA hotspot telemetry and live isolation alerts forBucket #412. - Tab
[4] SQL Explorer: Live aligned SQL table viewer with 13 preset schema and data queries (Left/Rightarrows to cycle queries,Up/Downor Mouse Wheel to scroll,qto type any custom SQL). - Interactive Shard Customizer & Creator Modals (
[e]/[n]): Presseto edit the selected shard ornto provision a new shard using the 8-field free-form modal (Up/Downto select field,Left/Rightto nudge numeric values, type any exact byte capacity or string directly,Enterto apply live with CDC VReplication,dto drain the selected shard,wto rebalance all shards by weight). - Hotkeys:
[1-4]or[Tab]Switch Views |[e]Edit Shard |[n]New Shard |[s]Zero-Downtime Split |[h]Spike Bucket #412 |[m]500M+ QPS Burst |[b]Pause/Resume Load |[r]Reset Cluster |[Esc]Return to Control Center.
shardmaster/
|-- cmd/
| +-- shardmaster/
| |-- main.go # Cobra CLI entrypoint and direct subcommands
| |-- console_windows.go # Native Win32 console allocation & VT100 ANSI enabler
| +-- console_other.go # Cross-platform POSIX console stub
|-- pkg/
| |-- console/
| | |-- shell.go # Animated ASCII banner, -/|\- spinners, typed SQL tables & REPL
| | +-- editor_windows.go # Native Win32 [Shift+Enter] multi-line SQL & [Tab] indent editor
| |-- hash/
| | |-- ring.go # Zero-allocation xxhash/v2 1,024 Virtual Bucket ring & split math
| | +-- ring_test.go # Unit, PGWire, K-Way Merge, CDC/VDiff, Full SQL, ACID & Persistence tests
| |-- directory/
| | +-- directory.go # Lock-free [1024]atomic.Uint32 ShardDirectory & per-bucket cutover gates
| |-- storage/
| | +-- backend.go # Real hashed bucket slabs, byte quotas, custom tables, CDC log & NVMe state
| |-- pgwire/
| | +-- server.go # Pillar 1: TCP :6000 PostgreSQL v3.0 Wire Protocol Server (pgproto3)
| |-- router/
| | |-- schema.go # Distributed table catalog, DDL generator & column/index metadata
| | |-- sql_engine.go # 100% Full ANSI/PostgreSQL Relational Engine + ACID Savepoints/Rollback
| | |-- lexer.go # Zero-regex SQL classifier, shard DDL parser & predicate extractor
| | |-- router.go # Central SQL query planner, Map-Reduce aggregator & admin executor
| | +-- kway_merge.go # Pillar 2: Parallel Goroutine Scatter-Gather + Min-Heap K-Way Merge
| |-- cdc/
| | |-- streamer.go # Pillar 3: Keyset Backfill, capacity-aware CDC streamer & cutover gate
| | +-- vdiff.go # Pillar 4: Commutative 256-bit XOR of SHA-256 row digests (VDiff)
| |-- hotspot/
| | +-- ewma.go # Pillar 5: 64-byte padded EWMA frequency tracker & auto-rebalancer
| |-- bench/
| | |-- engine_bench.go # Multi-Million QPS lock-free routing benchmark & chaos stress test
| | +-- petabyte_sim.go # Consistent-hashing shard split & bucket movement benchmark
| +-- tui/
| +-- dashboard.go # Pillar 6: 4-Tab minimalist Bubbletea TUI + Free-Form Shard Customizer
|-- Start-ShardMaster.bat # One-click Windows launcher script
|-- docker-compose.yml # 5 Isolated PostgreSQL 15 Alpine shard containers (:5432-:5436)
|-- go.mod # Go module definition
|-- go.sum # Cryptographic dependency checksums
+-- README.md # Architecture, SQL engine, and operations documentation