Resolves the 142 findings from tidal/docs/reviews/CODE_REVIEW_m0-m10.md across the engine, server, net, and CLI surfaces: - WAL/session-journal durability, checkpoint format, and crash-recovery hardening - Replication shipper/receiver, tenant isolation, and migration paths - Cluster scatter-gather, router, standalone server + health/offload endpoints - tidalctl refactored into command modules with JSON output and WAL-state tooling - Cohort, governance, signal-ledger, and vector-registry correctness fixes - Expanded UAT/integration/durability test coverage across all milestones
634 lines
18 KiB
Rust
634 lines
18 KiB
Rust
// Integration-test exemptions (same posture as the tidaldb integration tests):
|
|
// unwrap on known-good fixtures, short-lived read guards, and loop-counter
|
|
// casts in throughput math are idiomatic here.
|
|
#![allow(
|
|
clippy::unwrap_used,
|
|
clippy::significant_drop_tightening,
|
|
clippy::cast_precision_loss
|
|
)]
|
|
//! gRPC transport integration tests (tier-2 distributed testing).
|
|
//!
|
|
//! Proves the full path: signal write → WAL batch encode → gRPC ship →
|
|
//! gRPC receive → `apply_payload` → `SignalLedger` replication.
|
|
//!
|
|
//! Uses `GrpcTransport` on localhost instead of in-process crossbeam channels.
|
|
//! All `TidalDb` instances run in the same test process — this validates the
|
|
//! transport layer's serialization, delivery, and idempotency guarantees.
|
|
//!
|
|
//! **Not covered here (requires tier-3 multi-process harness):**
|
|
//! - Spawning separate `tidal-server cluster` OS processes
|
|
//! - iptables/pfctl network partition injection
|
|
//! - HLC clock skew simulation across process boundaries
|
|
//! - Rolling upgrade with mixed binary versions
|
|
//! - HTTP runbook endpoint verification against real cluster
|
|
|
|
use std::{
|
|
collections::HashMap,
|
|
net::SocketAddr,
|
|
thread,
|
|
time::{Duration, Instant},
|
|
};
|
|
|
|
use tidal_net::{GrpcTransport, config::GrpcTransportConfig};
|
|
use tidaldb::{
|
|
TidalDb,
|
|
db::config::{NodeConfig, NodeRole},
|
|
replication::{
|
|
WalSegmentId,
|
|
receiver::apply_payload,
|
|
shard::{RegionId, ShardId},
|
|
transport::{Transport, WalSegmentPayload},
|
|
},
|
|
schema::{DecaySpec, EntityId, EntityKind, SchemaBuilder, Timestamp, Window},
|
|
signals::{NoopWalWriter, SignalLedger},
|
|
wal::format::batch::{EventRecord, encode_batch},
|
|
};
|
|
|
|
// ── Helpers ────────────────────────────────────────────────────────────────
|
|
|
|
fn free_addr() -> SocketAddr {
|
|
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
|
|
listener.local_addr().unwrap()
|
|
}
|
|
|
|
fn m8_schema() -> tidaldb::schema::Schema {
|
|
let mut builder = SchemaBuilder::new();
|
|
let _ = builder
|
|
.signal(
|
|
"view",
|
|
EntityKind::Item,
|
|
DecaySpec::Exponential {
|
|
half_life: Duration::from_secs(7 * 24 * 3600),
|
|
},
|
|
)
|
|
.windows(&[Window::OneHour, Window::TwentyFourHours])
|
|
.velocity(false)
|
|
.add();
|
|
let _ = builder
|
|
.signal(
|
|
"like",
|
|
EntityKind::Item,
|
|
DecaySpec::Exponential {
|
|
half_life: Duration::from_secs(24 * 3600),
|
|
},
|
|
)
|
|
.windows(&[Window::OneHour])
|
|
.velocity(false)
|
|
.add();
|
|
builder.build().unwrap()
|
|
}
|
|
|
|
/// A node in the gRPC cluster.
|
|
struct GrpcNode {
|
|
db: TidalDb,
|
|
transport: GrpcTransport,
|
|
}
|
|
|
|
/// Build a leader + follower pair connected via `GrpcTransport`.
|
|
fn build_pair() -> (GrpcNode, GrpcNode, HashMap<String, u8>) {
|
|
let schema = m8_schema();
|
|
let addr0 = free_addr();
|
|
let addr1 = free_addr();
|
|
|
|
let t0 = GrpcTransport::new(GrpcTransportConfig {
|
|
local_shard: ShardId(0),
|
|
listen_addr: addr0,
|
|
peers: HashMap::from([(ShardId(1), addr1)]),
|
|
insecure: true,
|
|
..Default::default()
|
|
})
|
|
.unwrap();
|
|
|
|
let t1 = GrpcTransport::new(GrpcTransportConfig {
|
|
local_shard: ShardId(1),
|
|
listen_addr: addr1,
|
|
peers: HashMap::from([(ShardId(0), addr0)]),
|
|
insecure: true,
|
|
..Default::default()
|
|
})
|
|
.unwrap();
|
|
|
|
thread::sleep(Duration::from_millis(200));
|
|
|
|
let db0 = TidalDb::builder()
|
|
.ephemeral()
|
|
.with_schema(schema.clone())
|
|
.with_cluster(NodeConfig {
|
|
role: NodeRole::Single,
|
|
shard_id: ShardId(0),
|
|
peer_shards: vec![ShardId(1)],
|
|
..NodeConfig::default()
|
|
})
|
|
.open()
|
|
.unwrap();
|
|
|
|
let db1 = TidalDb::builder()
|
|
.ephemeral()
|
|
.with_schema(schema.clone())
|
|
.with_cluster(NodeConfig {
|
|
role: NodeRole::Single,
|
|
shard_id: ShardId(1),
|
|
peer_shards: vec![ShardId(0)],
|
|
..NodeConfig::default()
|
|
})
|
|
.open()
|
|
.unwrap();
|
|
|
|
let scratch = SignalLedger::new(schema, Box::new(NoopWalWriter));
|
|
let sig_ids: HashMap<String, u8> = ["view", "like"]
|
|
.iter()
|
|
.filter_map(|name| {
|
|
scratch
|
|
.resolve_signal_type(name)
|
|
.ok()
|
|
.map(|id| (name.to_string(), id.as_u16() as u8))
|
|
})
|
|
.collect();
|
|
|
|
(
|
|
GrpcNode {
|
|
db: db0,
|
|
transport: t0,
|
|
},
|
|
GrpcNode {
|
|
db: db1,
|
|
transport: t1,
|
|
},
|
|
sig_ids,
|
|
)
|
|
}
|
|
|
|
/// Write a signal to a node, encode a WAL batch, and ship via gRPC.
|
|
fn write_and_ship(
|
|
node: &GrpcNode,
|
|
signal_type: &str,
|
|
entity_id: EntityId,
|
|
weight: f64,
|
|
seqno: u64,
|
|
sig_ids: &HashMap<String, u8>,
|
|
target_shard: ShardId,
|
|
) {
|
|
let ts = Timestamp::now();
|
|
node.db.signal(signal_type, entity_id, weight, ts).unwrap();
|
|
|
|
let type_id = *sig_ids.get(signal_type).unwrap();
|
|
let events = [EventRecord::signal(
|
|
entity_id.as_u64(),
|
|
type_id,
|
|
weight as f32,
|
|
ts.as_nanos(),
|
|
)];
|
|
let bytes = encode_batch(&events, seqno, ts.as_nanos()).unwrap();
|
|
|
|
let payload = WalSegmentPayload {
|
|
id: WalSegmentId::new(RegionId::SINGLE, ShardId(0), seqno),
|
|
bytes,
|
|
event_count: 1,
|
|
leader_last_seq: seqno,
|
|
};
|
|
node.transport
|
|
.send_segment(target_shard, payload)
|
|
.expect("gRPC ship failed");
|
|
}
|
|
|
|
/// Receive a payload from the transport and apply it to the follower's ledger.
|
|
fn recv_and_apply(node: &GrpcNode) {
|
|
let payload = node
|
|
.transport
|
|
.recv_segment()
|
|
.expect("recv_segment returned None");
|
|
let ledger = node.db.ledger().unwrap().clone();
|
|
let rep_state = node.db.replication_state().clone();
|
|
apply_payload(&payload.bytes, ShardId(0), &ledger, &rep_state, None)
|
|
.expect("apply_payload failed");
|
|
}
|
|
|
|
// ── UAT Tests ──────────────────────────────────────────────────────────────
|
|
|
|
/// Step 1: Cross-region signal replication over gRPC.
|
|
///
|
|
/// Write signals on leader, ship via gRPC, receive + apply on follower.
|
|
/// Verify decay scores match (6 decimal places).
|
|
#[test]
|
|
fn uat_step1_grpc_replication() {
|
|
let (leader, follower, sig_ids) = build_pair();
|
|
|
|
for i in 1..=10u64 {
|
|
write_and_ship(
|
|
&leader,
|
|
"view",
|
|
EntityId::new(i),
|
|
1.0,
|
|
i,
|
|
&sig_ids,
|
|
ShardId(1),
|
|
);
|
|
recv_and_apply(&follower);
|
|
}
|
|
|
|
for i in 1..=10u64 {
|
|
let eid = EntityId::new(i);
|
|
let l = leader
|
|
.db
|
|
.read_decay_score(eid, "view", 0)
|
|
.unwrap()
|
|
.unwrap_or(0.0);
|
|
let f = follower
|
|
.db
|
|
.read_decay_score(eid, "view", 0)
|
|
.unwrap()
|
|
.unwrap_or(0.0);
|
|
assert!(
|
|
(l - f).abs() < 1e-6,
|
|
"entity {i}: leader={l} vs follower={f}"
|
|
);
|
|
}
|
|
}
|
|
|
|
/// Step 2: Idempotent replay — same batch shipped twice, applied once.
|
|
#[test]
|
|
fn uat_step2_idempotent_replay() {
|
|
let (leader, follower, sig_ids) = build_pair();
|
|
|
|
let ts = Timestamp::now();
|
|
let type_id = *sig_ids.get("view").unwrap();
|
|
let events = [EventRecord::signal(42, type_id, 1.0, ts.as_nanos())];
|
|
let bytes = encode_batch(&events, 1, ts.as_nanos()).unwrap();
|
|
|
|
// Ship same batch twice.
|
|
for _ in 0..2 {
|
|
let payload = WalSegmentPayload {
|
|
id: WalSegmentId::new(RegionId::SINGLE, ShardId(0), 1),
|
|
bytes: bytes.clone(),
|
|
event_count: 1,
|
|
leader_last_seq: 1,
|
|
};
|
|
leader.transport.send_segment(ShardId(1), payload).unwrap();
|
|
}
|
|
|
|
// Receive and apply both.
|
|
recv_and_apply(&follower);
|
|
recv_and_apply(&follower);
|
|
|
|
let score = follower
|
|
.db
|
|
.read_decay_score(EntityId::new(42), "view", 0)
|
|
.unwrap()
|
|
.unwrap_or(0.0);
|
|
|
|
// Should be ~1.0 (applied once), not ~2.0.
|
|
assert!(score > 0.5 && score < 1.5, "expected ~1.0, got {score}");
|
|
}
|
|
|
|
/// Step 3: Mixed signal types replicate correctly.
|
|
#[test]
|
|
fn uat_step3_mixed_signals() {
|
|
let (leader, follower, sig_ids) = build_pair();
|
|
|
|
// Views.
|
|
for i in 1..=5u64 {
|
|
write_and_ship(
|
|
&leader,
|
|
"view",
|
|
EntityId::new(i),
|
|
1.0,
|
|
i,
|
|
&sig_ids,
|
|
ShardId(1),
|
|
);
|
|
recv_and_apply(&follower);
|
|
}
|
|
// Likes.
|
|
for i in 1..=5u64 {
|
|
write_and_ship(
|
|
&leader,
|
|
"like",
|
|
EntityId::new(i),
|
|
2.0,
|
|
i + 5,
|
|
&sig_ids,
|
|
ShardId(1),
|
|
);
|
|
recv_and_apply(&follower);
|
|
}
|
|
|
|
for i in 1..=5u64 {
|
|
let eid = EntityId::new(i);
|
|
let view = follower
|
|
.db
|
|
.read_decay_score(eid, "view", 0)
|
|
.unwrap()
|
|
.unwrap_or(0.0);
|
|
let like = follower
|
|
.db
|
|
.read_decay_score(eid, "like", 0)
|
|
.unwrap()
|
|
.unwrap_or(0.0);
|
|
assert!(view > 0.0, "entity {i} view missing");
|
|
assert!(like > 0.0, "entity {i} like missing");
|
|
}
|
|
}
|
|
|
|
/// Step 4: Performance — 100 signals replicate within 2s over gRPC.
|
|
#[test]
|
|
fn perf_replication_latency() {
|
|
let (leader, follower, sig_ids) = build_pair();
|
|
|
|
let count = 100u64;
|
|
let start = Instant::now();
|
|
|
|
for i in 1..=count {
|
|
write_and_ship(
|
|
&leader,
|
|
"view",
|
|
EntityId::new(i),
|
|
1.0,
|
|
i,
|
|
&sig_ids,
|
|
ShardId(1),
|
|
);
|
|
}
|
|
for _ in 1..=count {
|
|
recv_and_apply(&follower);
|
|
}
|
|
|
|
let elapsed = start.elapsed();
|
|
assert!(
|
|
elapsed < Duration::from_secs(2),
|
|
"{count} signals took {elapsed:?}, exceeds 2s"
|
|
);
|
|
eprintln!(
|
|
"perf_replication_latency: {count} signals in {elapsed:?} ({:.0} signals/sec)",
|
|
count as f64 / elapsed.as_secs_f64()
|
|
);
|
|
}
|
|
|
|
/// Step 5: Three-node replication — leader ships to 2 followers.
|
|
#[test]
|
|
fn uat_step5_three_node_replication() {
|
|
let schema = m8_schema();
|
|
let addrs: Vec<SocketAddr> = (0..3).map(|_| free_addr()).collect();
|
|
|
|
let transports: Vec<GrpcTransport> = (0..3u16)
|
|
.map(|i| {
|
|
let mut peers = HashMap::new();
|
|
for j in 0..3u16 {
|
|
if i != j {
|
|
peers.insert(ShardId(j), addrs[j as usize]);
|
|
}
|
|
}
|
|
GrpcTransport::new(GrpcTransportConfig {
|
|
local_shard: ShardId(i),
|
|
listen_addr: addrs[i as usize],
|
|
peers,
|
|
insecure: true,
|
|
..Default::default()
|
|
})
|
|
.unwrap()
|
|
})
|
|
.collect();
|
|
|
|
thread::sleep(Duration::from_millis(200));
|
|
|
|
let dbs: Vec<TidalDb> = (0..3u16)
|
|
.map(|i| {
|
|
TidalDb::builder()
|
|
.ephemeral()
|
|
.with_schema(schema.clone())
|
|
.with_cluster(NodeConfig {
|
|
role: NodeRole::Single,
|
|
shard_id: ShardId(i),
|
|
peer_shards: (0..3u16).filter(|&j| j != i).map(ShardId).collect(),
|
|
..NodeConfig::default()
|
|
})
|
|
.open()
|
|
.unwrap()
|
|
})
|
|
.collect();
|
|
|
|
let scratch = SignalLedger::new(schema, Box::new(NoopWalWriter));
|
|
let view_id = scratch.resolve_signal_type("view").unwrap().as_u16() as u8;
|
|
|
|
// Write 5 signals on leader (node 0), ship to both followers.
|
|
for seq in 1..=5u64 {
|
|
let ts = Timestamp::now();
|
|
let eid = EntityId::new(seq);
|
|
dbs[0].signal("view", eid, 1.0, ts).unwrap();
|
|
|
|
let events = [EventRecord::signal(seq, view_id, 1.0, ts.as_nanos())];
|
|
let bytes = encode_batch(&events, seq, ts.as_nanos()).unwrap();
|
|
|
|
// Ship to follower 1 and 2.
|
|
for target in [ShardId(1), ShardId(2)] {
|
|
let payload = WalSegmentPayload {
|
|
id: WalSegmentId::new(RegionId::SINGLE, ShardId(0), seq),
|
|
bytes: bytes.clone(),
|
|
event_count: 1,
|
|
leader_last_seq: seq,
|
|
};
|
|
transports[0].send_segment(target, payload).unwrap();
|
|
}
|
|
}
|
|
|
|
// Receive and apply on both followers.
|
|
for follower_idx in [1, 2] {
|
|
for _ in 1..=5u64 {
|
|
let payload = transports[follower_idx]
|
|
.recv_segment()
|
|
.expect("recv failed");
|
|
let ledger = dbs[follower_idx].ledger().unwrap().clone();
|
|
let rep = dbs[follower_idx].replication_state().clone();
|
|
apply_payload(&payload.bytes, ShardId(0), &ledger, &rep, None).unwrap();
|
|
}
|
|
}
|
|
|
|
// Verify all 3 nodes agree.
|
|
for i in 1..=5u64 {
|
|
let eid = EntityId::new(i);
|
|
let scores: Vec<f64> = (0..3)
|
|
.map(|n| {
|
|
dbs[n]
|
|
.read_decay_score(eid, "view", 0)
|
|
.unwrap()
|
|
.unwrap_or(0.0)
|
|
})
|
|
.collect();
|
|
assert!(
|
|
(scores[0] - scores[1]).abs() < 1e-6 && (scores[0] - scores[2]).abs() < 1e-6,
|
|
"entity {i}: scores={scores:?}"
|
|
);
|
|
}
|
|
}
|
|
|
|
/// Basic seed-and-converge: write on leader, ship via gRPC, verify follower matches.
|
|
#[test]
|
|
fn seed_and_converge() {
|
|
let (leader, follower, sig_ids) = build_pair();
|
|
|
|
// Seed data.
|
|
for i in 1..=5u64 {
|
|
write_and_ship(
|
|
&leader,
|
|
"view",
|
|
EntityId::new(i),
|
|
1.0,
|
|
i,
|
|
&sig_ids,
|
|
ShardId(1),
|
|
);
|
|
recv_and_apply(&follower);
|
|
}
|
|
|
|
// Verify convergence.
|
|
for i in 1..=5u64 {
|
|
let eid = EntityId::new(i);
|
|
let l = leader
|
|
.db
|
|
.read_decay_score(eid, "view", 0)
|
|
.unwrap()
|
|
.unwrap_or(0.0);
|
|
let f = follower
|
|
.db
|
|
.read_decay_score(eid, "view", 0)
|
|
.unwrap()
|
|
.unwrap_or(0.0);
|
|
assert!(
|
|
(l - f).abs() < 1e-6,
|
|
"entity {i}: leader={l} vs follower={f}"
|
|
);
|
|
}
|
|
}
|
|
|
|
/// Simulated partition: stop shipping to follower, continue writing on leader,
|
|
/// then resume shipping and verify the follower catches up.
|
|
///
|
|
/// This tests the "heal after partition" scenario at the gRPC transport layer.
|
|
/// The WAL is durable on the leader; the follower replays missed segments.
|
|
#[test]
|
|
fn partition_heal_convergence() {
|
|
let (leader, follower, sig_ids) = build_pair();
|
|
|
|
// Phase 1: Normal replication — ship 5 signals.
|
|
for i in 1..=5u64 {
|
|
write_and_ship(
|
|
&leader,
|
|
"view",
|
|
EntityId::new(i),
|
|
1.0,
|
|
i,
|
|
&sig_ids,
|
|
ShardId(1),
|
|
);
|
|
recv_and_apply(&follower);
|
|
}
|
|
|
|
// Phase 2: "Partition" — write 5 more on leader, DON'T ship to follower.
|
|
// We still encode and store them for later replay.
|
|
let mut missed_payloads = Vec::new();
|
|
for i in 6..=10u64 {
|
|
let ts = Timestamp::now();
|
|
leader.db.signal("view", EntityId::new(i), 1.0, ts).unwrap();
|
|
|
|
let type_id = *sig_ids.get("view").unwrap();
|
|
let events = [EventRecord::signal(i, type_id, 1.0, ts.as_nanos())];
|
|
let bytes = encode_batch(&events, i, ts.as_nanos()).unwrap();
|
|
missed_payloads.push(WalSegmentPayload {
|
|
id: WalSegmentId::new(RegionId::SINGLE, ShardId(0), i),
|
|
bytes,
|
|
event_count: 1,
|
|
leader_last_seq: i,
|
|
});
|
|
}
|
|
|
|
// Verify follower is missing entities 6-10.
|
|
for i in 6..=10u64 {
|
|
let score = follower
|
|
.db
|
|
.read_decay_score(EntityId::new(i), "view", 0)
|
|
.unwrap();
|
|
assert!(
|
|
score.is_none() || score.unwrap() == 0.0,
|
|
"entity {i} should not exist on follower during partition"
|
|
);
|
|
}
|
|
|
|
// Phase 3: "Heal" — ship the missed segments via gRPC.
|
|
for payload in missed_payloads {
|
|
leader
|
|
.transport
|
|
.send_segment(ShardId(1), payload)
|
|
.expect("gRPC ship failed");
|
|
recv_and_apply(&follower);
|
|
}
|
|
|
|
// Verify all 10 entities converge.
|
|
for i in 1..=10u64 {
|
|
let eid = EntityId::new(i);
|
|
let l = leader
|
|
.db
|
|
.read_decay_score(eid, "view", 0)
|
|
.unwrap()
|
|
.unwrap_or(0.0);
|
|
let f = follower
|
|
.db
|
|
.read_decay_score(eid, "view", 0)
|
|
.unwrap()
|
|
.unwrap_or(0.0);
|
|
assert!(
|
|
(l - f).abs() < 1e-6,
|
|
"entity {i} after heal: leader={l} vs follower={f}"
|
|
);
|
|
}
|
|
}
|
|
|
|
/// Follower receives data during partition via a different leader, then
|
|
/// the original leader's batches arrive — idempotency ensures no duplication.
|
|
/// This approximates the "degraded query during partition" UAT step.
|
|
#[test]
|
|
fn degraded_follower_still_serves_old_data() {
|
|
let (leader, follower, sig_ids) = build_pair();
|
|
|
|
// Ship some initial data.
|
|
for i in 1..=5u64 {
|
|
write_and_ship(
|
|
&leader,
|
|
"view",
|
|
EntityId::new(i),
|
|
1.0,
|
|
i,
|
|
&sig_ids,
|
|
ShardId(1),
|
|
);
|
|
recv_and_apply(&follower);
|
|
}
|
|
|
|
// "Partition" — write more on leader without shipping.
|
|
let ts = Timestamp::now();
|
|
leader
|
|
.db
|
|
.signal("view", EntityId::new(100), 5.0, ts)
|
|
.unwrap();
|
|
|
|
// Follower should still serve the original 5 entities (degraded but available).
|
|
for i in 1..=5u64 {
|
|
let score = follower
|
|
.db
|
|
.read_decay_score(EntityId::new(i), "view", 0)
|
|
.unwrap();
|
|
assert!(
|
|
score.is_some(),
|
|
"entity {i} should be readable on follower during partition"
|
|
);
|
|
}
|
|
|
|
// Entity 100 should NOT be on follower.
|
|
let missing = follower
|
|
.db
|
|
.read_decay_score(EntityId::new(100), "view", 0)
|
|
.unwrap();
|
|
assert!(
|
|
missing.is_none() || missing.unwrap() == 0.0,
|
|
"entity 100 should not be on follower during partition"
|
|
);
|
|
}
|