Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
10 changes: 8 additions & 2 deletions crates/ember-cluster/src/topology.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
33 changes: 14 additions & 19 deletions crates/ember-server/src/cluster.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ pub struct ClusterCoordinator {
gossip: Mutex<GossipEngine>,
migration: Mutex<MigrationManager>,
local_id: NodeId,
gossip_port_offset: u16,
/// bound UDP socket for gossip, set after spawn_gossip
udp_socket: Mutex<Option<Arc<UdpSocket>>>,
}
Expand All @@ -51,20 +52,18 @@ impl ClusterCoordinator {
) -> (Self, mpsc::Receiver<GossipEvent>) {
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
Expand All @@ -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),
};

Expand Down Expand Up @@ -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())
Expand Down Expand Up @@ -412,17 +412,8 @@ impl ClusterCoordinator {
bind_addr: SocketAddr,
mut event_rx: mpsc::Receiver<GossipEvent>,
) {
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),
Expand Down Expand Up @@ -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);
}
}
Expand Down
129 changes: 129 additions & 0 deletions tests/integration/src/cli.rs
Original file line number Diff line number Diff line change
@@ -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::<u64>().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"
);
}
Loading