diff --git a/README.md b/README.md index 4f7e4d9a..37156925 100644 --- a/README.md +++ b/README.md @@ -235,7 +235,7 @@ contributions welcome — see [CONTRIBUTING.md](CONTRIBUTING.md). | 4 | clustering (raft, gossip, slots, migration) | ✅ complete | | 5 | developer experience (observability, CLI, clients) | 🚧 in progress | -**current**: 85 commands, 861 tests, ~18k lines of code (excluding tests) +**current**: 85 commands, 886 tests, ~18k lines of code (excluding tests) ## security diff --git a/crates/ember-cluster/src/topology.rs b/crates/ember-cluster/src/topology.rs index d63d799b..a2d4d1de 100644 --- a/crates/ember-cluster/src/topology.rs +++ b/crates/ember-cluster/src/topology.rs @@ -141,9 +141,15 @@ pub struct ClusterNode { } impl ClusterNode { - /// Creates a new primary node. + /// Creates a new primary node with the default bus port offset (10000). pub fn new_primary(id: NodeId, addr: SocketAddr) -> Self { - let cluster_bus_addr = SocketAddr::new(addr.ip(), addr.port() + 10000); + Self::new_primary_with_offset(id, addr, 10000) + } + + /// Creates a new primary node with a custom bus port offset. + pub fn new_primary_with_offset(id: NodeId, addr: SocketAddr, bus_port_offset: u16) -> Self { + let cluster_bus_addr = + SocketAddr::new(addr.ip(), addr.port().wrapping_add(bus_port_offset)); Self { id, addr, diff --git a/crates/ember-server/src/cluster.rs b/crates/ember-server/src/cluster.rs index 2ee97a43..b2d66687 100644 --- a/crates/ember-server/src/cluster.rs +++ b/crates/ember-server/src/cluster.rs @@ -26,6 +26,7 @@ pub struct ClusterCoordinator { gossip: Mutex, migration: Mutex, local_id: NodeId, + gossip_port_offset: u16, /// bound UDP socket for gossip, set after spawn_gossip udp_socket: Mutex>>, } @@ -51,20 +52,18 @@ impl ClusterCoordinator { ) -> (Self, mpsc::Receiver) { let (event_tx, event_rx) = mpsc::channel(256); - let gossip_addr = SocketAddr::new( - bind_addr.ip(), - bind_addr.port() + gossip_config.gossip_port_offset, - ); + let port_offset = gossip_config.gossip_port_offset; + let gossip_addr = SocketAddr::new(bind_addr.ip(), bind_addr.port() + port_offset); let gossip = GossipEngine::new(local_id, gossip_addr, gossip_config, event_tx); let state = if bootstrap { - let mut node = ClusterNode::new_primary(local_id, bind_addr); + let mut node = ClusterNode::new_primary_with_offset(local_id, bind_addr, port_offset); node.set_myself(); ClusterState::single_node(node) } else { let mut cs = ClusterState::new(local_id); - let mut node = ClusterNode::new_primary(local_id, bind_addr); + let mut node = ClusterNode::new_primary_with_offset(local_id, bind_addr, port_offset); node.set_myself(); cs.add_node(node); cs @@ -75,6 +74,7 @@ impl ClusterCoordinator { gossip: Mutex::new(gossip), migration: Mutex::new(MigrationManager::new()), local_id, + gossip_port_offset: port_offset, udp_socket: Mutex::new(None), }; @@ -166,7 +166,7 @@ impl ClusterCoordinator { // add to cluster state as well let mut state = self.state.write().await; - let node = ClusterNode::new_primary(new_id, addr); + let node = ClusterNode::new_primary_with_offset(new_id, addr, self.gossip_port_offset); state.add_node(node); Frame::Simple("OK".into()) @@ -412,17 +412,8 @@ impl ClusterCoordinator { bind_addr: SocketAddr, mut event_rx: mpsc::Receiver, ) { - let gossip_addr = { - let gossip = self.gossip.lock().await; - // gossip engine was initialized with the gossip address - // we need to bind to the same port - drop(gossip); - - // compute gossip address from bind_addr + offset - // the offset was baked into the GossipEngine, but we need - // to derive it for the UDP bind - SocketAddr::new(bind_addr.ip(), bind_addr.port() + 10000) - }; + let gossip_addr = + SocketAddr::new(bind_addr.ip(), bind_addr.port() + self.gossip_port_offset); let socket = match UdpSocket::bind(gossip_addr).await { Ok(s) => Arc::new(s), @@ -495,7 +486,11 @@ impl ClusterCoordinator { GossipEvent::MemberJoined(id, addr) => { info!("cluster: node {} joined at {}", id, addr); if !state.nodes.contains_key(&id) { - let node = ClusterNode::new_primary(id, addr); + let node = ClusterNode::new_primary_with_offset( + id, + addr, + coordinator.gossip_port_offset, + ); state.add_node(node); } } diff --git a/tests/integration/src/cli.rs b/tests/integration/src/cli.rs new file mode 100644 index 00000000..ac96275e --- /dev/null +++ b/tests/integration/src/cli.rs @@ -0,0 +1,129 @@ +//! Integration tests for the emberkv-cli binary. + +use crate::helpers::{run_cli, ServerOptions, TestServer}; + +/// Helper — start a plain (non-cluster) server. +fn plain_server() -> TestServer { + TestServer::start() +} + +/// Helper — start a cluster-enabled server. +fn cluster_server() -> TestServer { + TestServer::start_with(ServerOptions { + cluster_enabled: true, + cluster_bootstrap: true, + ..Default::default() + }) +} + +// -- one-shot mode -- + +#[tokio::test] +async fn cli_oneshot_ping() { + let server = plain_server(); + let output = run_cli(server.port, &["PING"]); + + assert!(output.status.success(), "exit code: {:?}", output.status); + let stdout = String::from_utf8_lossy(&output.stdout); + assert!( + stdout.contains("PONG"), + "expected PONG in stdout, got: {stdout}" + ); +} + +#[tokio::test] +async fn cli_oneshot_set_get() { + let server = plain_server(); + + let set_out = run_cli(server.port, &["SET", "clikey", "clival"]); + assert!(set_out.status.success()); + let stdout = String::from_utf8_lossy(&set_out.stdout); + assert!(stdout.contains("OK"), "expected OK, got: {stdout}"); + + let get_out = run_cli(server.port, &["GET", "clikey"]); + assert!(get_out.status.success()); + let stdout = String::from_utf8_lossy(&get_out.stdout); + assert!( + stdout.contains("clival"), + "expected clival in stdout, got: {stdout}" + ); +} + +#[tokio::test] +async fn cli_oneshot_error() { + let server = plain_server(); + let output = run_cli(server.port, &["NOTAREALCOMMAND"]); + + let stdout = String::from_utf8_lossy(&output.stdout); + let stderr = String::from_utf8_lossy(&output.stderr); + let combined = format!("{stdout}{stderr}"); + assert!( + combined.contains("unknown command") + || combined.contains("error") + || combined.contains("ERR"), + "expected error output, got stdout={stdout} stderr={stderr}" + ); +} + +// -- cluster subcommands -- + +#[tokio::test] +async fn cli_cluster_info() { + let server = cluster_server(); + let output = run_cli(server.port, &["cluster", "info"]); + + assert!(output.status.success(), "exit code: {:?}", output.status); + let stdout = String::from_utf8_lossy(&output.stdout); + assert!( + stdout.contains("cluster_state") || stdout.contains("cluster_slots"), + "expected cluster info in stdout, got: {stdout}" + ); +} + +#[tokio::test] +async fn cli_cluster_keyslot() { + let server = cluster_server(); + let output = run_cli(server.port, &["cluster", "keyslot", "foo"]); + + assert!(output.status.success(), "exit code: {:?}", output.status); + let stdout = String::from_utf8_lossy(&output.stdout); + // output should contain the integer slot number + let trimmed = stdout.trim(); + assert!( + trimmed.contains("integer") || trimmed.parse::().is_ok(), + "expected integer slot output, got: {trimmed}" + ); +} + +// -- benchmark -- + +#[tokio::test] +async fn cli_benchmark_smoke() { + let server = plain_server(); + let output = run_cli(server.port, &["benchmark", "-n", "100", "-c", "2", "-q"]); + + assert!( + output.status.success(), + "benchmark failed: {:?}\nstderr: {}", + output.status, + String::from_utf8_lossy(&output.stderr) + ); + let stdout = String::from_utf8_lossy(&output.stdout); + assert!( + stdout.contains("rps"), + "expected 'rps' in benchmark output, got: {stdout}" + ); +} + +// -- error handling -- + +#[tokio::test] +async fn cli_connection_refused() { + // connect to a port that nothing is listening on + let output = run_cli(1, &["PING"]); + + assert!( + !output.status.success(), + "expected non-zero exit code for refused connection" + ); +} diff --git a/tests/integration/src/cluster.rs b/tests/integration/src/cluster.rs new file mode 100644 index 00000000..c028ab61 --- /dev/null +++ b/tests/integration/src/cluster.rs @@ -0,0 +1,307 @@ +//! Integration tests for cluster command dispatch. + +use ember_protocol::Frame; + +use crate::helpers::{ServerOptions, TestServer}; + +/// Starts a single-node cluster server (all 16384 slots assigned). +fn cluster_server() -> TestServer { + TestServer::start_with(ServerOptions { + cluster_enabled: true, + cluster_bootstrap: true, + ..Default::default() + }) +} + +/// Starts a cluster server without bootstrapping (no slots assigned). +fn cluster_server_empty() -> TestServer { + TestServer::start_with(ServerOptions { + cluster_enabled: true, + ..Default::default() + }) +} + +// -- query commands -- + +#[tokio::test] +async fn cluster_info() { + let server = cluster_server(); + let mut c = server.connect().await; + + let resp = c.cmd(&["CLUSTER", "INFO"]).await; + match resp { + Frame::Bulk(data) => { + let text = String::from_utf8_lossy(&data); + assert!( + text.contains("cluster_state:ok"), + "expected cluster_state:ok, got: {text}" + ); + assert!( + text.contains("cluster_slots_assigned:16384"), + "expected 16384 slots assigned, got: {text}" + ); + } + other => panic!("expected Bulk, got {other:?}"), + } +} + +#[tokio::test] +async fn cluster_nodes() { + let server = cluster_server(); + let mut c = server.connect().await; + + let resp = c.cmd(&["CLUSTER", "NODES"]).await; + match resp { + Frame::Bulk(data) => { + let text = String::from_utf8_lossy(&data); + assert!( + text.contains("myself"), + "node entry should contain 'myself'" + ); + assert!( + text.contains("master"), + "node entry should contain 'master'" + ); + } + other => panic!("expected Bulk, got {other:?}"), + } +} + +#[tokio::test] +async fn cluster_slots() { + let server = cluster_server(); + let mut c = server.connect().await; + + let resp = c.cmd(&["CLUSTER", "SLOTS"]).await; + match resp { + Frame::Array(ref entries) => { + assert!( + !entries.is_empty(), + "bootstrapped cluster should have slot ranges" + ); + } + other => panic!("expected Array, got {other:?}"), + } +} + +#[tokio::test] +async fn cluster_myid() { + let server = cluster_server(); + let mut c = server.connect().await; + + let id = c + .get_bulk(&["CLUSTER", "MYID"]) + .await + .expect("MYID should return a bulk string"); + + // UUID v4 format: 8-4-4-4-12 hex digits + assert_eq!(id.len(), 36, "node ID should be a 36-char UUID, got: {id}"); + assert_eq!(id.matches('-').count(), 4, "UUID should have 4 hyphens"); +} + +#[tokio::test] +async fn cluster_keyslot() { + let server = cluster_server(); + let mut c = server.connect().await; + + let slot = c.get_int(&["CLUSTER", "KEYSLOT", "foo"]).await; + assert!(slot >= 0 && slot < 16384, "slot out of range: {slot}"); + + // same key should always hash to the same slot + let slot2 = c.get_int(&["CLUSTER", "KEYSLOT", "foo"]).await; + assert_eq!(slot, slot2); +} + +// -- state-changing commands -- + +#[tokio::test] +async fn cluster_addslots() { + let server = cluster_server_empty(); + let mut c = server.connect().await; + + // empty cluster should have 0 slots + let info = c.cmd(&["CLUSTER", "INFO"]).await; + let text = match &info { + Frame::Bulk(data) => String::from_utf8_lossy(data).to_string(), + other => panic!("expected Bulk, got {other:?}"), + }; + assert!(text.contains("cluster_slots_assigned:0"), "got: {text}"); + + // add a few slots + c.ok(&["CLUSTER", "ADDSLOTS", "0", "1", "2"]).await; + + let info = c.cmd(&["CLUSTER", "INFO"]).await; + let text = match &info { + Frame::Bulk(data) => String::from_utf8_lossy(data).to_string(), + other => panic!("expected Bulk, got {other:?}"), + }; + assert!(text.contains("cluster_slots_assigned:3"), "got: {text}"); +} + +#[tokio::test] +async fn cluster_delslots() { + let server = cluster_server_empty(); + let mut c = server.connect().await; + + c.ok(&["CLUSTER", "ADDSLOTS", "100", "101", "102"]).await; + c.ok(&["CLUSTER", "DELSLOTS", "101"]).await; + + let info = c.cmd(&["CLUSTER", "INFO"]).await; + let text = match &info { + Frame::Bulk(data) => String::from_utf8_lossy(data).to_string(), + other => panic!("expected Bulk, got {other:?}"), + }; + assert!(text.contains("cluster_slots_assigned:2"), "got: {text}"); +} + +#[tokio::test] +async fn cluster_forget_unknown() { + let server = cluster_server(); + let mut c = server.connect().await; + + let err = c + .err(&["CLUSTER", "FORGET", "00000000-0000-0000-0000-000000000000"]) + .await; + assert!( + err.contains("Unknown node") || err.contains("unknown"), + "expected unknown node error, got: {err}" + ); +} + +// -- data routing -- + +#[tokio::test] +async fn cluster_set_get_owned_slot() { + let server = cluster_server(); + let mut c = server.connect().await; + + // bootstrapped node owns all slots, so data commands should work + c.ok(&["SET", "hello", "world"]).await; + let val = c.get_bulk(&["GET", "hello"]).await; + assert_eq!(val, Some("world".into())); +} + +#[tokio::test] +async fn cluster_no_slots_returns_error() { + let server = cluster_server_empty(); + let mut c = server.connect().await; + + let err = c.err(&["SET", "key", "val"]).await; + assert!( + err.contains("CLUSTERDOWN"), + "expected CLUSTERDOWN error, got: {err}" + ); +} + +// -- slot queries -- + +#[tokio::test] +async fn cluster_countkeysinslot() { + let server = cluster_server(); + let mut c = server.connect().await; + + // figure out which slot "mykey" hashes to + let slot = c.get_int(&["CLUSTER", "KEYSLOT", "mykey"]).await; + + c.ok(&["SET", "mykey", "value"]).await; + let count = c + .get_int(&["CLUSTER", "COUNTKEYSINSLOT", &slot.to_string()]) + .await; + assert_eq!(count, 1); +} + +#[tokio::test] +async fn cluster_getkeysinslot() { + let server = cluster_server(); + let mut c = server.connect().await; + + let slot = c.get_int(&["CLUSTER", "KEYSLOT", "slotkey"]).await; + c.ok(&["SET", "slotkey", "val"]).await; + + let resp = c + .cmd(&["CLUSTER", "GETKEYSINSLOT", &slot.to_string(), "10"]) + .await; + match resp { + Frame::Array(keys) => { + assert_eq!(keys.len(), 1); + assert!(matches!(&keys[0], Frame::Bulk(b) if b == &b"slotkey"[..])); + } + other => panic!("expected Array, got {other:?}"), + } +} + +// -- stubs / error responses -- + +#[tokio::test] +async fn cluster_replicate_stub() { + let server = cluster_server(); + let mut c = server.connect().await; + + let err = c + .err(&[ + "CLUSTER", + "REPLICATE", + "00000000-0000-0000-0000-000000000000", + ]) + .await; + assert!( + err.contains("not yet supported"), + "expected stub error, got: {err}" + ); +} + +#[tokio::test] +async fn cluster_failover_stub() { + let server = cluster_server(); + let mut c = server.connect().await; + + let err = c.err(&["CLUSTER", "FAILOVER"]).await; + assert!( + err.contains("not yet supported"), + "expected stub error, got: {err}" + ); +} + +#[tokio::test] +async fn cluster_migrate_stub() { + let server = cluster_server(); + let mut c = server.connect().await; + + let err = c + .err(&["MIGRATE", "127.0.0.1", "6380", "key", "0", "1000"]) + .await; + assert!( + err.contains("not yet implemented"), + "expected stub error, got: {err}" + ); +} + +#[tokio::test] +async fn cluster_setslot_importing() { + let server = cluster_server(); + let mut c = server.connect().await; + + // importing from self should fail + let my_id = c + .get_bulk(&["CLUSTER", "MYID"]) + .await + .expect("MYID should return a string"); + + let err = c + .err(&["CLUSTER", "SETSLOT", "0", "IMPORTING", &my_id]) + .await; + assert!( + err.contains("can't import from myself") || err.contains("ERR"), + "expected error for self-import, got: {err}" + ); +} + +#[tokio::test] +async fn cluster_meet_invalid() { + let server = cluster_server(); + let mut c = server.connect().await; + + // invalid port + let err = c.err(&["CLUSTER", "MEET", "127.0.0.1", "99999"]).await; + assert!(!err.is_empty(), "expected error for invalid meet address"); +} diff --git a/tests/integration/src/helpers.rs b/tests/integration/src/helpers.rs index ca3cab55..702200e6 100644 --- a/tests/integration/src/helpers.rs +++ b/tests/integration/src/helpers.rs @@ -27,6 +27,10 @@ pub struct ServerOptions { /// Use an existing path without taking ownership. /// If both `data_dir` and `data_dir_path` are set, `data_dir_path` wins. pub data_dir_path: Option, + /// Start with cluster support enabled. + pub cluster_enabled: bool, + /// Bootstrap as a single-node cluster owning all 16384 slots. + pub cluster_bootstrap: bool, } impl TestServer { @@ -39,10 +43,10 @@ impl TestServer { /// Starts a new ember-server with custom options. pub fn start_with(opts: ServerOptions) -> Self { - let port = find_free_port(); - let binary = server_binary(); + let port = find_free_port(); + let mut cmd = Command::new(&binary); cmd.arg("--port").arg(port.to_string()); cmd.arg("--host").arg("127.0.0.1"); @@ -54,6 +58,17 @@ impl TestServer { cmd.arg("--requirepass").arg(pass); } + if opts.cluster_enabled { + cmd.arg("--cluster-enabled"); + // use a small offset so the gossip port stays in valid u16 range. + // random test ports are often >55000, and the default +10000 would + // overflow past 65535. + cmd.arg("--cluster-port-offset").arg("1"); + } + if opts.cluster_bootstrap { + cmd.arg("--cluster-bootstrap"); + } + let data_dir = if opts.appendonly { cmd.arg("--appendonly"); cmd.arg("--appendfsync").arg("always"); @@ -220,32 +235,42 @@ fn find_free_port() -> u16 { listener.local_addr().unwrap().port() } -/// Locates the ember-server binary in the cargo target directory. -fn server_binary() -> PathBuf { - // cargo sets OUT_DIR for build scripts, but for integration tests - // we can find the binary relative to the test binary itself +/// Locates a binary in the cargo target directory. +fn find_binary(name: &str) -> PathBuf { let mut path = std::env::current_exe().unwrap(); // test binary is in target/debug/deps/ — go up to target/debug/ path.pop(); if path.ends_with("deps") { path.pop(); } - path.push("ember-server"); + path.push(name); if !path.exists() { - // try release - let mut release = std::env::current_exe().unwrap(); - release.pop(); - if release.ends_with("deps") { - release.pop(); - } - release.push("ember-server"); - if release.exists() { - return release; - } panic!( - "ember-server binary not found. run `cargo build` first.\nlooked at: {}", + "{name} binary not found. run `cargo build` first.\nlooked at: {}", path.display() ); } path } + +/// Locates the ember-server binary in the cargo target directory. +fn server_binary() -> PathBuf { + find_binary("ember-server") +} + +/// Locates the ember-cli binary in the cargo target directory. +pub fn cli_binary() -> PathBuf { + find_binary("ember-cli") +} + +/// Runs the CLI binary with the given args and captures output. +pub fn run_cli(port: u16, args: &[&str]) -> std::process::Output { + Command::new(cli_binary()) + .arg("--port") + .arg(port.to_string()) + .arg("--host") + .arg("127.0.0.1") + .args(args) + .output() + .expect("failed to run emberkv-cli") +} diff --git a/tests/integration/src/main.rs b/tests/integration/src/main.rs index e9af83ca..7ee5389e 100644 --- a/tests/integration/src/main.rs +++ b/tests/integration/src/main.rs @@ -2,6 +2,8 @@ mod helpers; mod auth; mod basic_operations; +mod cli; +mod cluster; mod data_types; mod persistence; mod pubsub;