Phase 2b: Blob RPC (Stat/Get/Put/LoadManifest) #7
@@ -386,6 +386,7 @@ mod tests {
|
||||
bind_rpc_lan: None,
|
||||
bind_rpc_tailscale: None,
|
||||
tls: None,
|
||||
blob_store_root: None,
|
||||
};
|
||||
let g = ClusterGossip::bootstrap(&cfg, "solo").await.unwrap();
|
||||
let id = g.self_chitchat_id().await;
|
||||
@@ -404,6 +405,7 @@ mod tests {
|
||||
bind_rpc_lan: None,
|
||||
bind_rpc_tailscale: None,
|
||||
tls: None,
|
||||
blob_store_root: None,
|
||||
};
|
||||
let err = ClusterGossip::bootstrap(&cfg, "")
|
||||
.await
|
||||
@@ -422,6 +424,7 @@ mod tests {
|
||||
bind_rpc_lan: None,
|
||||
bind_rpc_tailscale: None,
|
||||
tls: None,
|
||||
blob_store_root: None,
|
||||
};
|
||||
// ClusterConfig::validate rejects this first — that's what we want:
|
||||
// the daemon should refuse to bootstrap gossip on a malformed config.
|
||||
@@ -448,6 +451,7 @@ mod tests {
|
||||
bind_rpc_lan: None,
|
||||
bind_rpc_tailscale: None,
|
||||
tls: None,
|
||||
blob_store_root: None,
|
||||
};
|
||||
// Node B: uses A as seed.
|
||||
let cfg_b = ClusterConfig {
|
||||
@@ -463,6 +467,7 @@ mod tests {
|
||||
bind_rpc_lan: None,
|
||||
bind_rpc_tailscale: None,
|
||||
tls: None,
|
||||
blob_store_root: None,
|
||||
};
|
||||
|
||||
let gossip_a = ClusterGossip::bootstrap(&cfg_a, "a").await.unwrap();
|
||||
@@ -534,6 +539,7 @@ mod tests {
|
||||
bind_rpc_lan: None,
|
||||
bind_rpc_tailscale: None,
|
||||
tls: None,
|
||||
blob_store_root: None,
|
||||
};
|
||||
let g = ClusterGossip::bootstrap(&cfg, "solo").await.unwrap();
|
||||
// Solo cluster — peers() must never include self.
|
||||
|
||||
+506
-10
@@ -25,6 +25,7 @@
|
||||
//! containing this node's local view of the cluster (its own
|
||||
//! name+zone, plus every peer it currently knows about via gossip).
|
||||
|
||||
use crate::cluster::blob::{BlobId, BlobManifest, BlobStat, BlobStore};
|
||||
use crate::cluster::gossip::{ClusterGossip, PeerView};
|
||||
use anyhow::{bail, Context, Result};
|
||||
use quinn::{Connection, ConnectionError};
|
||||
@@ -32,9 +33,11 @@ use serde::{Deserialize, Serialize};
|
||||
use std::sync::Arc;
|
||||
|
||||
/// Cap on a single request or response, including the method tag.
|
||||
/// Matches `MAX_MESSAGE_BYTES` in [`super::transport`] so both layers
|
||||
/// bound memory the same way.
|
||||
pub const MAX_MESSAGE_BYTES: usize = 16 * 1024;
|
||||
/// 16 MiB is generous enough to hold one 4 MiB blob chunk with plenty
|
||||
/// of framing headroom; multi-chunk / whole-blob transfers still fit
|
||||
/// well under that ceiling for anything up to a few MB. Streaming
|
||||
/// (many-GB) put/get lands in Phase 2c.
|
||||
pub const MAX_MESSAGE_BYTES: usize = 16 * 1024 * 1024;
|
||||
|
||||
/// RPC method tag byte.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
@@ -42,6 +45,18 @@ pub const MAX_MESSAGE_BYTES: usize = 16 * 1024;
|
||||
pub enum Method {
|
||||
Ping = 0x01,
|
||||
PeerStatus = 0x02,
|
||||
/// `payload`: 32-byte `BlobId`. Reply: JSON `BlobStat` or
|
||||
/// single-byte [`ErrorCode::NotFound`].
|
||||
BlobStat = 0x03,
|
||||
/// `payload`: 32-byte `BlobId`. Reply: raw blob bytes or
|
||||
/// single-byte [`ErrorCode::NotFound`].
|
||||
BlobGet = 0x04,
|
||||
/// `payload`: raw blob bytes. Reply: 32-byte `BlobId` of the stored
|
||||
/// content.
|
||||
BlobPut = 0x05,
|
||||
/// `payload`: 32-byte `BlobId`. Reply: JSON `BlobManifest` or
|
||||
/// single-byte [`ErrorCode::NotFound`].
|
||||
BlobLoadManifest = 0x06,
|
||||
}
|
||||
|
||||
impl Method {
|
||||
@@ -51,6 +66,10 @@ impl Method {
|
||||
match b {
|
||||
0x01 => Some(Method::Ping),
|
||||
0x02 => Some(Method::PeerStatus),
|
||||
0x03 => Some(Method::BlobStat),
|
||||
0x04 => Some(Method::BlobGet),
|
||||
0x05 => Some(Method::BlobPut),
|
||||
0x06 => Some(Method::BlobLoadManifest),
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
@@ -70,6 +89,15 @@ pub enum ErrorCode {
|
||||
EmptyRequest = 0xf0,
|
||||
UnknownMethod = 0xf1,
|
||||
HandlerFailure = 0xf2,
|
||||
/// Blob (or manifest) not present in the local store.
|
||||
NotFound = 0xf3,
|
||||
/// Request payload was structurally wrong (e.g. wrong length for a
|
||||
/// 32-byte hash).
|
||||
InvalidRequest = 0xf4,
|
||||
/// Server-side subsystem required for this method wasn't
|
||||
/// configured (e.g. Blob RPCs called on a node with no local
|
||||
/// blob store).
|
||||
NotConfigured = 0xf5,
|
||||
}
|
||||
|
||||
impl ErrorCode {
|
||||
@@ -81,6 +109,9 @@ impl ErrorCode {
|
||||
ErrorCode::EmptyRequest => "empty request",
|
||||
ErrorCode::UnknownMethod => "unknown method",
|
||||
ErrorCode::HandlerFailure => "handler failure",
|
||||
ErrorCode::NotFound => "not found",
|
||||
ErrorCode::InvalidRequest => "invalid request",
|
||||
ErrorCode::NotConfigured => "server subsystem not configured",
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -99,36 +130,64 @@ pub struct PeerStatusReply {
|
||||
}
|
||||
|
||||
/// The concrete RPC handler used by the daemon. Holds Arc references
|
||||
/// to the state a request might need to read (currently just the
|
||||
/// gossip service; Phase 2+ adds the blob store, metadata index, etc).
|
||||
/// to the state a request might need to read: the gossip service
|
||||
/// (always), and optionally a local blob store (for Blob* methods).
|
||||
///
|
||||
/// Not `Clone` on its own — wrap in `Arc<RpcRouter>` so a single
|
||||
/// instance backs the accept loop plus any explicit dispatch calls.
|
||||
pub struct RpcRouter {
|
||||
gossip: Arc<ClusterGossip>,
|
||||
blob_store: Option<Arc<BlobStore>>,
|
||||
local_name: String,
|
||||
local_zone: String,
|
||||
}
|
||||
|
||||
/// Result of dispatching a request: either a real reply (`Ok`) or a
|
||||
/// well-known error code the wire layer surfaces as a single byte.
|
||||
/// Extracted so `handle` stays free of `Vec<u8>` shell games.
|
||||
enum HandlerOutcome {
|
||||
Reply(Vec<u8>),
|
||||
Error(ErrorCode),
|
||||
}
|
||||
|
||||
impl RpcRouter {
|
||||
pub fn new(gossip: Arc<ClusterGossip>, local_name: String, local_zone: String) -> Self {
|
||||
Self {
|
||||
gossip,
|
||||
blob_store: None,
|
||||
local_name,
|
||||
local_zone,
|
||||
}
|
||||
}
|
||||
|
||||
/// Attach a local blob store. Enables the `Blob*` methods; nodes
|
||||
/// without a store return [`ErrorCode::NotConfigured`] for those.
|
||||
pub fn with_blob_store(mut self, store: Arc<BlobStore>) -> Self {
|
||||
self.blob_store = Some(store);
|
||||
self
|
||||
}
|
||||
|
||||
pub fn blob_store(&self) -> Option<&Arc<BlobStore>> {
|
||||
self.blob_store.as_ref()
|
||||
}
|
||||
|
||||
/// Dispatch a single request. Called by [`serve_connection`] for
|
||||
/// every accepted bidi stream. Test code may call it directly to
|
||||
/// bypass the transport.
|
||||
pub async fn handle(&self, method: Method, payload: &[u8]) -> Result<Vec<u8>> {
|
||||
match self.handle_outcome(method, payload).await? {
|
||||
HandlerOutcome::Reply(bytes) => Ok(bytes),
|
||||
HandlerOutcome::Error(code) => Ok(vec![code.as_byte()]),
|
||||
}
|
||||
}
|
||||
|
||||
async fn handle_outcome(&self, method: Method, payload: &[u8]) -> Result<HandlerOutcome> {
|
||||
match method {
|
||||
Method::Ping => {
|
||||
let mut reply = Vec::with_capacity(5 + payload.len());
|
||||
reply.extend_from_slice(b"pong:");
|
||||
reply.extend_from_slice(payload);
|
||||
Ok(reply)
|
||||
Ok(HandlerOutcome::Reply(reply))
|
||||
}
|
||||
Method::PeerStatus => {
|
||||
let peers = self.gossip.peers().await;
|
||||
@@ -146,11 +205,93 @@ impl RpcRouter {
|
||||
MAX_MESSAGE_BYTES
|
||||
);
|
||||
}
|
||||
Ok(json)
|
||||
Ok(HandlerOutcome::Reply(json))
|
||||
}
|
||||
Method::BlobStat => {
|
||||
let store = match &self.blob_store {
|
||||
Some(s) => s,
|
||||
None => return Ok(HandlerOutcome::Error(ErrorCode::NotConfigured)),
|
||||
};
|
||||
let id = match decode_blob_id(payload) {
|
||||
Some(id) => id,
|
||||
None => return Ok(HandlerOutcome::Error(ErrorCode::InvalidRequest)),
|
||||
};
|
||||
match store.stat(&id).await? {
|
||||
Some(stat) => {
|
||||
let json = serde_json::to_vec(&stat)
|
||||
.context("encoding BlobStat as JSON")?;
|
||||
Ok(HandlerOutcome::Reply(json))
|
||||
}
|
||||
None => Ok(HandlerOutcome::Error(ErrorCode::NotFound)),
|
||||
}
|
||||
}
|
||||
Method::BlobGet => {
|
||||
let store = match &self.blob_store {
|
||||
Some(s) => s,
|
||||
None => return Ok(HandlerOutcome::Error(ErrorCode::NotConfigured)),
|
||||
};
|
||||
let id = match decode_blob_id(payload) {
|
||||
Some(id) => id,
|
||||
None => return Ok(HandlerOutcome::Error(ErrorCode::InvalidRequest)),
|
||||
};
|
||||
match store.get_bytes(&id).await? {
|
||||
Some(bytes) => {
|
||||
if bytes.len() > MAX_MESSAGE_BYTES {
|
||||
bail!(
|
||||
"blob {} at {} bytes exceeds RPC cap {}; use streaming variant (Phase 2c)",
|
||||
id.to_hex(),
|
||||
bytes.len(),
|
||||
MAX_MESSAGE_BYTES
|
||||
);
|
||||
}
|
||||
Ok(HandlerOutcome::Reply(bytes))
|
||||
}
|
||||
None => Ok(HandlerOutcome::Error(ErrorCode::NotFound)),
|
||||
}
|
||||
}
|
||||
Method::BlobPut => {
|
||||
let store = match &self.blob_store {
|
||||
Some(s) => s,
|
||||
None => return Ok(HandlerOutcome::Error(ErrorCode::NotConfigured)),
|
||||
};
|
||||
// Empty payload is a legitimate empty-blob put — falls
|
||||
// through to store.put_bytes(&[]) which yields the
|
||||
// hash of the empty byte sequence.
|
||||
let id = store.put_bytes(payload).await?;
|
||||
Ok(HandlerOutcome::Reply(id.as_bytes().to_vec()))
|
||||
}
|
||||
Method::BlobLoadManifest => {
|
||||
let store = match &self.blob_store {
|
||||
Some(s) => s,
|
||||
None => return Ok(HandlerOutcome::Error(ErrorCode::NotConfigured)),
|
||||
};
|
||||
let id = match decode_blob_id(payload) {
|
||||
Some(id) => id,
|
||||
None => return Ok(HandlerOutcome::Error(ErrorCode::InvalidRequest)),
|
||||
};
|
||||
match store.load_manifest(&id).await? {
|
||||
Some(manifest) => {
|
||||
let json = serde_json::to_vec(&manifest)
|
||||
.context("encoding BlobManifest as JSON")?;
|
||||
Ok(HandlerOutcome::Reply(json))
|
||||
}
|
||||
None => Ok(HandlerOutcome::Error(ErrorCode::NotFound)),
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Parse a payload as a 32-byte BlobId. Returns `None` for any other
|
||||
/// length so the caller can surface [`ErrorCode::InvalidRequest`].
|
||||
fn decode_blob_id(payload: &[u8]) -> Option<BlobId> {
|
||||
if payload.len() != 32 {
|
||||
return None;
|
||||
}
|
||||
let mut buf = [0u8; 32];
|
||||
buf.copy_from_slice(payload);
|
||||
Some(BlobId::from_bytes(buf))
|
||||
}
|
||||
|
||||
/// Server-side: loop accepting bidi streams on `conn`, dispatch to
|
||||
/// `router`, write the reply. Returns cleanly when the peer closes the
|
||||
@@ -262,15 +403,106 @@ pub async fn call_peer_status(conn: &Connection) -> Result<PeerStatusReply> {
|
||||
/// Recognise a single-byte reply as one of our error codes. Returns
|
||||
/// `None` for any other single-byte value (which is a valid reply,
|
||||
/// just an unusually short one).
|
||||
fn decode_error(b: u8) -> Option<ErrorCode> {
|
||||
pub fn decode_error(b: u8) -> Option<ErrorCode> {
|
||||
match b {
|
||||
0xf0 => Some(ErrorCode::EmptyRequest),
|
||||
0xf1 => Some(ErrorCode::UnknownMethod),
|
||||
0xf2 => Some(ErrorCode::HandlerFailure),
|
||||
0xf3 => Some(ErrorCode::NotFound),
|
||||
0xf4 => Some(ErrorCode::InvalidRequest),
|
||||
0xf5 => Some(ErrorCode::NotConfigured),
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
|
||||
// ── Blob RPC client helpers ───────────────────────────────────────────
|
||||
|
||||
/// Ask the peer for a blob's size + chunk count. `Ok(None)` when the
|
||||
/// peer replies [`ErrorCode::NotFound`]; other error codes surface as
|
||||
/// `Err`.
|
||||
pub async fn call_blob_stat(
|
||||
conn: &Connection,
|
||||
id: &BlobId,
|
||||
) -> Result<Option<BlobStat>> {
|
||||
let reply = rpc_call(conn, Method::BlobStat, id.as_bytes()).await?;
|
||||
if reply.len() == 1 {
|
||||
match decode_error(reply[0]) {
|
||||
Some(ErrorCode::NotFound) => return Ok(None),
|
||||
Some(code) => bail!("peer replied with error: {}", code.describe()),
|
||||
None => {} // single-byte JSON like "1" is technically possible; fall through
|
||||
}
|
||||
}
|
||||
let stat = serde_json::from_slice(&reply).context("decoding BlobStat JSON")?;
|
||||
Ok(Some(stat))
|
||||
}
|
||||
|
||||
/// Fetch a blob's raw bytes. `Ok(None)` when the peer replies
|
||||
/// [`ErrorCode::NotFound`]. Callers that need many-GB transfers should
|
||||
/// use the streaming variant (Phase 2c); this helper caps at
|
||||
/// [`MAX_MESSAGE_BYTES`].
|
||||
pub async fn call_blob_get(
|
||||
conn: &Connection,
|
||||
id: &BlobId,
|
||||
) -> Result<Option<Vec<u8>>> {
|
||||
let reply = rpc_call(conn, Method::BlobGet, id.as_bytes()).await?;
|
||||
if reply.len() == 1 {
|
||||
match decode_error(reply[0]) {
|
||||
Some(ErrorCode::NotFound) => return Ok(None),
|
||||
Some(code) => bail!("peer replied with error: {}", code.describe()),
|
||||
None => {} // single-byte blob is legitimate content
|
||||
}
|
||||
}
|
||||
Ok(Some(reply))
|
||||
}
|
||||
|
||||
/// Upload a blob to the peer's store. Returns the BlobId the peer
|
||||
/// assigned; must match the local hash of `bytes` (proves the peer
|
||||
/// stored what we sent).
|
||||
pub async fn call_blob_put(conn: &Connection, bytes: &[u8]) -> Result<BlobId> {
|
||||
let reply = rpc_call(conn, Method::BlobPut, bytes).await?;
|
||||
if reply.len() == 1 {
|
||||
if let Some(code) = decode_error(reply[0]) {
|
||||
bail!("peer replied with error: {}", code.describe());
|
||||
}
|
||||
}
|
||||
if reply.len() != 32 {
|
||||
bail!(
|
||||
"expected 32-byte BlobId in reply, got {} bytes",
|
||||
reply.len()
|
||||
);
|
||||
}
|
||||
let mut buf = [0u8; 32];
|
||||
buf.copy_from_slice(&reply);
|
||||
let assigned = BlobId::from_bytes(buf);
|
||||
let expected = BlobId::from_bytes(blake3::hash(bytes).into());
|
||||
if assigned != expected {
|
||||
bail!(
|
||||
"peer returned BlobId {} but content hashes to {}; corruption or protocol drift",
|
||||
assigned.to_hex(),
|
||||
expected.to_hex()
|
||||
);
|
||||
}
|
||||
Ok(assigned)
|
||||
}
|
||||
|
||||
/// Fetch a blob's manifest (chunk-list + size). `Ok(None)` on NotFound.
|
||||
pub async fn call_blob_load_manifest(
|
||||
conn: &Connection,
|
||||
id: &BlobId,
|
||||
) -> Result<Option<BlobManifest>> {
|
||||
let reply = rpc_call(conn, Method::BlobLoadManifest, id.as_bytes()).await?;
|
||||
if reply.len() == 1 {
|
||||
match decode_error(reply[0]) {
|
||||
Some(ErrorCode::NotFound) => return Ok(None),
|
||||
Some(code) => bail!("peer replied with error: {}", code.describe()),
|
||||
None => {}
|
||||
}
|
||||
}
|
||||
let manifest =
|
||||
serde_json::from_slice(&reply).context("decoding BlobManifest JSON")?;
|
||||
Ok(Some(manifest))
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
@@ -304,12 +536,40 @@ mod tests {
|
||||
fn method_round_trips_byte_encoding() {
|
||||
assert_eq!(Method::Ping.as_byte(), 0x01);
|
||||
assert_eq!(Method::PeerStatus.as_byte(), 0x02);
|
||||
assert_eq!(Method::from_byte(0x01), Some(Method::Ping));
|
||||
assert_eq!(Method::from_byte(0x02), Some(Method::PeerStatus));
|
||||
assert_eq!(Method::BlobStat.as_byte(), 0x03);
|
||||
assert_eq!(Method::BlobGet.as_byte(), 0x04);
|
||||
assert_eq!(Method::BlobPut.as_byte(), 0x05);
|
||||
assert_eq!(Method::BlobLoadManifest.as_byte(), 0x06);
|
||||
for m in [
|
||||
Method::Ping,
|
||||
Method::PeerStatus,
|
||||
Method::BlobStat,
|
||||
Method::BlobGet,
|
||||
Method::BlobPut,
|
||||
Method::BlobLoadManifest,
|
||||
] {
|
||||
assert_eq!(Method::from_byte(m.as_byte()), Some(m));
|
||||
}
|
||||
assert_eq!(Method::from_byte(0x00), None);
|
||||
assert_eq!(Method::from_byte(0xff), None);
|
||||
}
|
||||
|
||||
fn open_blob_store() -> (tempfile::TempDir, Arc<BlobStore>) {
|
||||
let tmp = tempfile::TempDir::new().unwrap();
|
||||
let store = Arc::new(BlobStore::open(tmp.path().to_path_buf()).unwrap());
|
||||
(tmp, store)
|
||||
}
|
||||
|
||||
async fn router_with_blobs(name: &str, port: u16) -> (tempfile::TempDir, Arc<RpcRouter>) {
|
||||
let gossip = bootstrap_gossip(name, port).await;
|
||||
let (tmp, store) = open_blob_store();
|
||||
let router = Arc::new(
|
||||
RpcRouter::new(gossip, name.into(), "fabric-10g".into())
|
||||
.with_blob_store(store),
|
||||
);
|
||||
(tmp, router)
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn dispatch_returns_pong_for_ping() {
|
||||
let gossip = bootstrap_gossip("solo", next_port()).await;
|
||||
@@ -506,5 +766,241 @@ mod tests {
|
||||
assert_eq!(ErrorCode::EmptyRequest.describe(), "empty request");
|
||||
assert_eq!(ErrorCode::UnknownMethod.describe(), "unknown method");
|
||||
assert_eq!(ErrorCode::HandlerFailure.describe(), "handler failure");
|
||||
assert_eq!(ErrorCode::NotFound.describe(), "not found");
|
||||
assert_eq!(ErrorCode::InvalidRequest.describe(), "invalid request");
|
||||
assert_eq!(
|
||||
ErrorCode::NotConfigured.describe(),
|
||||
"server subsystem not configured"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn decode_error_covers_all_known_codes() {
|
||||
for code in [
|
||||
ErrorCode::EmptyRequest,
|
||||
ErrorCode::UnknownMethod,
|
||||
ErrorCode::HandlerFailure,
|
||||
ErrorCode::NotFound,
|
||||
ErrorCode::InvalidRequest,
|
||||
ErrorCode::NotConfigured,
|
||||
] {
|
||||
assert_eq!(decode_error(code.as_byte()), Some(code));
|
||||
}
|
||||
assert_eq!(decode_error(0x00), None);
|
||||
assert_eq!(decode_error(0xff), None);
|
||||
}
|
||||
|
||||
// ── Phase 2b: Blob RPC ─────────────────────────────────────────────
|
||||
|
||||
#[tokio::test]
|
||||
async fn blob_rpcs_return_not_configured_without_store() {
|
||||
// Router built via `new` alone (no `.with_blob_store`) must
|
||||
// refuse Blob* methods with a well-known error code.
|
||||
let gossip = bootstrap_gossip("solo", next_port()).await;
|
||||
let router = RpcRouter::new(gossip, "solo".into(), "z".into());
|
||||
for method in [
|
||||
Method::BlobStat,
|
||||
Method::BlobGet,
|
||||
Method::BlobPut,
|
||||
Method::BlobLoadManifest,
|
||||
] {
|
||||
let mut req = vec![method.as_byte()];
|
||||
req.extend_from_slice(&[0u8; 32]);
|
||||
let reply = dispatch(&router, &req).await;
|
||||
assert_eq!(
|
||||
reply,
|
||||
vec![ErrorCode::NotConfigured.as_byte()],
|
||||
"method {method:?} should be NotConfigured without a store"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn blob_stat_returns_not_found_for_missing() {
|
||||
let (_tmp, router) = router_with_blobs("solo", next_port()).await;
|
||||
let missing = BlobId::from_bytes([0u8; 32]);
|
||||
let mut req = vec![Method::BlobStat.as_byte()];
|
||||
req.extend_from_slice(missing.as_bytes());
|
||||
let reply = dispatch(&router, &req).await;
|
||||
assert_eq!(reply, vec![ErrorCode::NotFound.as_byte()]);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn blob_stat_returns_json_for_existing() {
|
||||
let (_tmp, router) = router_with_blobs("solo", next_port()).await;
|
||||
let id = router
|
||||
.blob_store()
|
||||
.unwrap()
|
||||
.put_bytes(b"tiny content")
|
||||
.await
|
||||
.unwrap();
|
||||
let mut req = vec![Method::BlobStat.as_byte()];
|
||||
req.extend_from_slice(id.as_bytes());
|
||||
let reply = dispatch(&router, &req).await;
|
||||
let stat: BlobStat = serde_json::from_slice(&reply).unwrap();
|
||||
assert_eq!(stat.total_size, b"tiny content".len() as u64);
|
||||
assert_eq!(stat.chunk_count, 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn blob_stat_returns_invalid_request_for_bad_length() {
|
||||
let (_tmp, router) = router_with_blobs("solo", next_port()).await;
|
||||
// Payload is only 5 bytes; a valid BlobId is 32.
|
||||
let req = vec![Method::BlobStat.as_byte(), 1, 2, 3, 4, 5];
|
||||
let reply = dispatch(&router, &req).await;
|
||||
assert_eq!(reply, vec![ErrorCode::InvalidRequest.as_byte()]);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn blob_get_returns_content_bytes() {
|
||||
let (_tmp, router) = router_with_blobs("solo", next_port()).await;
|
||||
let payload: &[u8] = b"contents to retrieve";
|
||||
let id = router.blob_store().unwrap().put_bytes(payload).await.unwrap();
|
||||
let mut req = vec![Method::BlobGet.as_byte()];
|
||||
req.extend_from_slice(id.as_bytes());
|
||||
let reply = dispatch(&router, &req).await;
|
||||
assert_eq!(reply, payload);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn blob_put_stores_bytes_and_returns_hash() {
|
||||
let (_tmp, router) = router_with_blobs("solo", next_port()).await;
|
||||
let payload = b"put via rpc";
|
||||
let mut req = vec![Method::BlobPut.as_byte()];
|
||||
req.extend_from_slice(payload);
|
||||
let reply = dispatch(&router, &req).await;
|
||||
assert_eq!(reply.len(), 32);
|
||||
let mut id_bytes = [0u8; 32];
|
||||
id_bytes.copy_from_slice(&reply);
|
||||
let assigned = BlobId::from_bytes(id_bytes);
|
||||
let expected = BlobId::from_bytes(blake3::hash(payload).into());
|
||||
assert_eq!(assigned, expected);
|
||||
// Round-trip: the bytes are now readable via the store.
|
||||
let round = router
|
||||
.blob_store()
|
||||
.unwrap()
|
||||
.get_bytes(&assigned)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(round.as_deref(), Some(payload.as_slice()));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn blob_load_manifest_returns_json_for_existing() {
|
||||
let (_tmp, router) = router_with_blobs("solo", next_port()).await;
|
||||
let data = vec![0x77u8; 4 * 1024 * 1024 + 1]; // 2 chunks
|
||||
let id = router.blob_store().unwrap().put_bytes(&data).await.unwrap();
|
||||
let mut req = vec![Method::BlobLoadManifest.as_byte()];
|
||||
req.extend_from_slice(id.as_bytes());
|
||||
let reply = dispatch(&router, &req).await;
|
||||
let manifest: BlobManifest = serde_json::from_slice(&reply).unwrap();
|
||||
assert_eq!(manifest.blob_id, id);
|
||||
assert_eq!(manifest.total_size, data.len() as u64);
|
||||
assert_eq!(manifest.chunks.len(), 2);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn blob_load_manifest_returns_not_found_for_missing() {
|
||||
let (_tmp, router) = router_with_blobs("solo", next_port()).await;
|
||||
let missing = BlobId::from_bytes([0u8; 32]);
|
||||
let mut req = vec![Method::BlobLoadManifest.as_byte()];
|
||||
req.extend_from_slice(missing.as_bytes());
|
||||
let reply = dispatch(&router, &req).await;
|
||||
assert_eq!(reply, vec![ErrorCode::NotFound.as_byte()]);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn end_to_end_blob_put_stat_get_over_real_quic() {
|
||||
// The full loop: B → A over real QUIC + mTLS.
|
||||
// 1. B puts a blob on A (BlobPut).
|
||||
// 2. B queries stat + fetches it back (BlobStat + BlobGet).
|
||||
// 3. B asks for the manifest (BlobLoadManifest).
|
||||
// Every step goes through the actual wire, no shortcuts.
|
||||
let (id_a, id_b) = NodeIdentity::generate_test_pair("a", "b").unwrap();
|
||||
let (_tmp, router) = router_with_blobs("a", next_port()).await;
|
||||
|
||||
let server = QuicServer::bind(loopback(0), id_a).unwrap();
|
||||
let server_addr = server.local_addr().unwrap();
|
||||
let router_srv = router.clone();
|
||||
let accept_task = tokio::spawn(async move {
|
||||
if let Some(Ok(conn)) = server.accept().await {
|
||||
let _ = serve_connection(conn, router_srv).await;
|
||||
}
|
||||
});
|
||||
|
||||
let client = QuicClient::new(loopback(0), id_b).unwrap();
|
||||
let conn = client.connect(server_addr, "a").await.unwrap();
|
||||
|
||||
let payload: &[u8] = b"cross-node payload";
|
||||
let assigned = call_blob_put(&conn, payload).await.unwrap();
|
||||
let expected = BlobId::from_bytes(blake3::hash(payload).into());
|
||||
assert_eq!(assigned, expected);
|
||||
|
||||
let stat = call_blob_stat(&conn, &assigned).await.unwrap().unwrap();
|
||||
assert_eq!(stat.total_size, payload.len() as u64);
|
||||
assert_eq!(stat.chunk_count, 1);
|
||||
|
||||
let round = call_blob_get(&conn, &assigned).await.unwrap().unwrap();
|
||||
assert_eq!(round, payload);
|
||||
|
||||
let manifest = call_blob_load_manifest(&conn, &assigned)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(manifest.blob_id, assigned);
|
||||
assert_eq!(manifest.chunks.len(), 1);
|
||||
|
||||
// NotFound path also works over the wire.
|
||||
let ghost = BlobId::from_bytes([0u8; 32]);
|
||||
assert!(call_blob_stat(&conn, &ghost).await.unwrap().is_none());
|
||||
assert!(call_blob_get(&conn, &ghost).await.unwrap().is_none());
|
||||
assert!(call_blob_load_manifest(&conn, &ghost)
|
||||
.await
|
||||
.unwrap()
|
||||
.is_none());
|
||||
|
||||
conn.close(quinn::VarInt::from_u32(0), b"done");
|
||||
client.shutdown().await;
|
||||
tokio::time::sleep(Duration::from_millis(50)).await;
|
||||
accept_task.abort();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn end_to_end_multi_chunk_blob_over_real_quic() {
|
||||
// 6 MB blob → 2 chunks. Round-trips whole via BlobPut/Get.
|
||||
// Also confirms the RPC layer's MAX_MESSAGE_BYTES bump from
|
||||
// 16 KiB to 16 MiB actually took effect end-to-end.
|
||||
let (id_a, id_b) = NodeIdentity::generate_test_pair("a", "b").unwrap();
|
||||
let (_tmp, router) = router_with_blobs("a", next_port()).await;
|
||||
|
||||
let server = QuicServer::bind(loopback(0), id_a).unwrap();
|
||||
let server_addr = server.local_addr().unwrap();
|
||||
let router_srv = router.clone();
|
||||
let accept_task = tokio::spawn(async move {
|
||||
if let Some(Ok(conn)) = server.accept().await {
|
||||
let _ = serve_connection(conn, router_srv).await;
|
||||
}
|
||||
});
|
||||
|
||||
let client = QuicClient::new(loopback(0), id_b).unwrap();
|
||||
let conn = client.connect(server_addr, "a").await.unwrap();
|
||||
|
||||
let payload: Vec<u8> = (0..6 * 1024 * 1024)
|
||||
.map(|i| (i % 251) as u8)
|
||||
.collect();
|
||||
let assigned = call_blob_put(&conn, &payload).await.unwrap();
|
||||
let manifest = call_blob_load_manifest(&conn, &assigned)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(manifest.chunks.len(), 2, "6 MB should split into 2 chunks");
|
||||
let round = call_blob_get(&conn, &assigned).await.unwrap().unwrap();
|
||||
assert_eq!(round.len(), payload.len());
|
||||
assert_eq!(round, payload);
|
||||
|
||||
conn.close(quinn::VarInt::from_u32(0), b"done");
|
||||
client.shutdown().await;
|
||||
tokio::time::sleep(Duration::from_millis(50)).await;
|
||||
accept_task.abort();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -12,6 +12,7 @@
|
||||
//! When `[cluster]` is absent from the config the daemon runs exactly
|
||||
//! as it did pre-v2: no gossip, no RPC, no metric ticker.
|
||||
|
||||
use crate::cluster::blob::BlobStore;
|
||||
use crate::cluster::gossip::ClusterGossip;
|
||||
use crate::cluster::rpc::{serve_connection, RpcRouter};
|
||||
use crate::cluster::transport::{NodeIdentity, QuicServer};
|
||||
@@ -39,6 +40,10 @@ pub struct ClusterServices {
|
||||
/// Live gossip service. `Arc` so read-only consumers on other
|
||||
/// subsystems can hold references without blocking shutdown.
|
||||
pub gossip: Arc<ClusterGossip>,
|
||||
/// Local blob store, when `blob_store_root` was provided at start
|
||||
/// time. `None` means this node runs gossip-only (Blob RPCs return
|
||||
/// [`crate::cluster::rpc::ErrorCode::NotConfigured`]).
|
||||
pub blob_store: Option<Arc<BlobStore>>,
|
||||
/// QUIC accept-loop task. `None` when `[cluster.tls]` was absent
|
||||
/// and RPC therefore didn't come up.
|
||||
accept_task: Option<JoinHandle<()>>,
|
||||
@@ -70,6 +75,7 @@ impl ClusterServices {
|
||||
local_name: String,
|
||||
hot_dir: PathBuf,
|
||||
hot_max_bytes: u64,
|
||||
blob_store_root: Option<PathBuf>,
|
||||
) -> Result<Self> {
|
||||
let gossip = Arc::new(
|
||||
ClusterGossip::bootstrap(cluster, &local_name)
|
||||
@@ -80,6 +86,23 @@ impl ClusterServices {
|
||||
// Publish the static config value once. Used-bytes updates every tick.
|
||||
gossip.set_hot_max(hot_max_bytes).await;
|
||||
|
||||
// Open the local blob store if a root path was supplied. Kept
|
||||
// outside the TLS branch: a node can serve blobs to callers
|
||||
// without RPC (via in-process API) or over RPC (once TLS is
|
||||
// configured too).
|
||||
let blob_store: Option<Arc<BlobStore>> = match blob_store_root {
|
||||
Some(root) => {
|
||||
let store = BlobStore::open(root.clone())
|
||||
.with_context(|| format!("opening blob store at {}", root.display()))?;
|
||||
tracing::info!("blob store opened at {}", root.display());
|
||||
Some(Arc::new(store))
|
||||
}
|
||||
None => {
|
||||
tracing::info!("no blob store configured; Blob RPCs will return NotConfigured");
|
||||
None
|
||||
}
|
||||
};
|
||||
|
||||
// RPC server + accept loop — only when TLS material is configured.
|
||||
let accept_task = match &cluster.tls {
|
||||
Some(tls) => {
|
||||
@@ -92,11 +115,15 @@ impl ClusterServices {
|
||||
.context("no RPC bind address (need bind_lan or bind_tailscale)")?;
|
||||
let server = QuicServer::bind(bind, identity)
|
||||
.context("binding QUIC RPC server")?;
|
||||
let router = Arc::new(RpcRouter::new(
|
||||
let mut router = RpcRouter::new(
|
||||
gossip.clone(),
|
||||
local_name.clone(),
|
||||
cluster.zone.clone(),
|
||||
));
|
||||
);
|
||||
if let Some(store) = &blob_store {
|
||||
router = router.with_blob_store(store.clone());
|
||||
}
|
||||
let router = Arc::new(router);
|
||||
tracing::info!("cluster RPC server listening on {}", bind);
|
||||
Some(tokio::spawn(async move {
|
||||
accept_forever(server, router).await;
|
||||
@@ -124,11 +151,17 @@ impl ClusterServices {
|
||||
|
||||
Ok(Self {
|
||||
gossip,
|
||||
blob_store,
|
||||
accept_task,
|
||||
metric_task,
|
||||
})
|
||||
}
|
||||
|
||||
/// Whether a local blob store was configured at start time.
|
||||
pub fn blob_store_enabled(&self) -> bool {
|
||||
self.blob_store.is_some()
|
||||
}
|
||||
|
||||
/// Whether the QUIC RPC server is running. `false` when `[cluster.tls]`
|
||||
/// was absent at start time.
|
||||
pub fn rpc_enabled(&self) -> bool {
|
||||
@@ -223,9 +256,14 @@ mod tests {
|
||||
|
||||
/// Dedicated port range for services tests (44000+) so we don't
|
||||
/// collide with gossip (41000+), transport (42000+), rpc (43000+).
|
||||
///
|
||||
/// Bump by 2 so each returned port `p` also implicitly reserves
|
||||
/// `p + 1` — RPC binds to `bind_lan.port + 1` (see
|
||||
/// `ClusterConfig::rpc_lan`), so returning consecutive ports would
|
||||
/// have one test's RPC step on the next test's gossip.
|
||||
static NEXT_PORT: AtomicU16 = AtomicU16::new(44001);
|
||||
fn next_port() -> u16 {
|
||||
NEXT_PORT.fetch_add(1, Ordering::Relaxed)
|
||||
NEXT_PORT.fetch_add(2, Ordering::Relaxed)
|
||||
}
|
||||
|
||||
fn loopback(port: u16) -> SocketAddr {
|
||||
@@ -262,10 +300,12 @@ mod tests {
|
||||
"test-node".into(),
|
||||
tmp.path().to_path_buf(),
|
||||
1_000_000,
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(!svc.rpc_enabled(), "no [cluster.tls] → RPC disabled");
|
||||
assert!(!svc.blob_store_enabled(), "no blob root → blob store disabled");
|
||||
// Gossip must still be functional.
|
||||
let self_id = svc.gossip.self_chitchat_id().await;
|
||||
assert_eq!(self_id.node_id.as_ref(), "test-node");
|
||||
@@ -304,7 +344,8 @@ mod tests {
|
||||
std::fs::create_dir_all(&hot_dir).unwrap();
|
||||
std::fs::write(hot_dir.join("blob"), vec![0u8; 4096]).unwrap();
|
||||
|
||||
let svc = ClusterServices::start(&cfg_a, "a".into(), hot_dir.clone(), 1_000_000)
|
||||
let svc =
|
||||
ClusterServices::start(&cfg_a, "a".into(), hot_dir.clone(), 1_000_000, None)
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(svc.rpc_enabled());
|
||||
@@ -352,6 +393,7 @@ mod tests {
|
||||
"publisher".into(),
|
||||
tmp.path().to_path_buf(),
|
||||
10 * 1024 * 1024,
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
@@ -391,11 +433,11 @@ mod tests {
|
||||
};
|
||||
let tmp = tempfile::TempDir::new().unwrap();
|
||||
let svc_a =
|
||||
ClusterServices::start(&cfg_a, "a".into(), tmp.path().to_path_buf(), 1_000_000)
|
||||
ClusterServices::start(&cfg_a, "a".into(), tmp.path().to_path_buf(), 1_000_000, None)
|
||||
.await
|
||||
.unwrap();
|
||||
let svc_b =
|
||||
ClusterServices::start(&cfg_b, "b".into(), tmp.path().to_path_buf(), 1_000_000)
|
||||
ClusterServices::start(&cfg_b, "b".into(), tmp.path().to_path_buf(), 1_000_000, None)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
@@ -415,4 +457,79 @@ mod tests {
|
||||
svc_a.shutdown();
|
||||
svc_b.shutdown();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn services_with_blob_store_serves_blob_rpc_end_to_end() {
|
||||
// Prove that config → services → RPC path plumbs BlobStore
|
||||
// through correctly. Node A runs with a blob store; B dials it
|
||||
// and puts + gets a blob over real QUIC + mTLS.
|
||||
use crate::cluster::rpc::{call_blob_get, call_blob_put};
|
||||
use crate::cluster::transport::{FleetCa, NodeIdentity, QuicClient};
|
||||
use crate::config::ClusterTlsConfig;
|
||||
|
||||
let tmp = tempfile::TempDir::new().unwrap();
|
||||
let ca_dir = tmp.path().join("ca");
|
||||
let a_tls_dir = tmp.path().join("a-tls");
|
||||
let b_tls_dir = tmp.path().join("b-tls");
|
||||
let blob_root = tmp.path().join("a-blobs");
|
||||
|
||||
let ca = FleetCa::generate("test CA").unwrap();
|
||||
ca.save(&ca_dir).unwrap();
|
||||
ca.sign_leaf_to_pem("a", &a_tls_dir).unwrap();
|
||||
ca.sign_leaf_to_pem("b", &b_tls_dir).unwrap();
|
||||
|
||||
let port_a = next_port();
|
||||
let cfg_a = ClusterConfig {
|
||||
zone: "fabric-10g".into(),
|
||||
bind_lan: Some(loopback(port_a)),
|
||||
tls: Some(ClusterTlsConfig {
|
||||
ca_cert: a_tls_dir.join("ca.crt"),
|
||||
node_cert: a_tls_dir.join("node.crt"),
|
||||
node_key: a_tls_dir.join("node.key"),
|
||||
}),
|
||||
blob_store_root: Some(blob_root.clone()),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let hot_dir = tmp.path().join("a-hot");
|
||||
std::fs::create_dir_all(&hot_dir).unwrap();
|
||||
let svc = ClusterServices::start(
|
||||
&cfg_a,
|
||||
"a".into(),
|
||||
hot_dir,
|
||||
1_000_000,
|
||||
Some(blob_root.clone()),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(svc.rpc_enabled());
|
||||
assert!(svc.blob_store_enabled());
|
||||
|
||||
let id_b = NodeIdentity::from_pem_dir(&b_tls_dir).unwrap();
|
||||
let client = QuicClient::new(loopback(0), id_b).unwrap();
|
||||
tokio::time::sleep(Duration::from_millis(50)).await;
|
||||
let rpc_addr = cfg_a.rpc_lan().unwrap();
|
||||
let conn = client.connect(rpc_addr, "a").await.unwrap();
|
||||
|
||||
let payload: &[u8] = b"cross-node blob over rpc";
|
||||
let id = call_blob_put(&conn, payload).await.unwrap();
|
||||
let round = call_blob_get(&conn, &id).await.unwrap().unwrap();
|
||||
assert_eq!(round, payload);
|
||||
|
||||
// Independently confirm the bytes are physically on disk in A's
|
||||
// store — proves the RPC didn't just echo back but actually
|
||||
// wrote through to BlobStore.
|
||||
let local_round = svc
|
||||
.blob_store
|
||||
.as_ref()
|
||||
.unwrap()
|
||||
.get_bytes(&id)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(local_round.as_deref(), Some(payload));
|
||||
|
||||
conn.close(quinn::VarInt::from_u32(0), b"done");
|
||||
client.shutdown().await;
|
||||
svc.shutdown();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -137,6 +137,12 @@ pub struct ClusterConfig {
|
||||
/// used; absent means "gossip only, no RPC" for now.
|
||||
#[serde(default)]
|
||||
pub tls: Option<ClusterTlsConfig>,
|
||||
/// Optional local blob-store root (Phase 2). When set, the daemon
|
||||
/// opens a content-addressed store at this path and serves it via
|
||||
/// the Blob* RPC methods. Absent means the node participates in
|
||||
/// gossip + peer-status but returns `NotConfigured` for Blob RPCs.
|
||||
#[serde(default)]
|
||||
pub blob_store_root: Option<PathBuf>,
|
||||
}
|
||||
|
||||
/// Compute the default RPC address for a gossip address: same IP, port + 1.
|
||||
@@ -395,6 +401,7 @@ tailscale_addr = "100.64.1.5:7701"
|
||||
bind_rpc_lan: None,
|
||||
bind_rpc_tailscale: None,
|
||||
tls: None,
|
||||
blob_store_root: None,
|
||||
};
|
||||
let err = cluster.validate().unwrap_err().to_string();
|
||||
assert!(
|
||||
@@ -426,6 +433,7 @@ tailscale_addr = "100.64.1.5:7701"
|
||||
bind_rpc_lan: None,
|
||||
bind_rpc_tailscale: None,
|
||||
tls: None,
|
||||
blob_store_root: None,
|
||||
};
|
||||
let err = cluster.validate().unwrap_err().to_string();
|
||||
assert!(
|
||||
|
||||
@@ -35,17 +35,20 @@ pub async fn run(cfg: Config, mut manifest: Manifest) -> Result<()> {
|
||||
Some(cluster_cfg) => {
|
||||
let hot_dir = cfg.hot.path.clone();
|
||||
let hot_max_bytes = cfg.hot.max_gb.saturating_mul(1024 * 1024 * 1024);
|
||||
let blob_root = cluster_cfg.blob_store_root.clone();
|
||||
match ClusterServices::start(
|
||||
cluster_cfg,
|
||||
cfg.node.name.clone(),
|
||||
hot_dir,
|
||||
hot_max_bytes,
|
||||
blob_root,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(svc) => {
|
||||
tracing::info!(
|
||||
rpc_enabled = svc.rpc_enabled(),
|
||||
blob_store_enabled = svc.blob_store_enabled(),
|
||||
zone = %cluster_cfg.zone,
|
||||
"cluster services online"
|
||||
);
|
||||
|
||||
Reference in New Issue
Block a user