fix: adjust auth protocol to datastore protocol
This commit is contained in:
parent
ebee109e95
commit
8ac45e593c
19 changed files with 1150 additions and 272 deletions
10
Cargo.lock
generated
10
Cargo.lock
generated
|
|
@ -1129,6 +1129,7 @@ dependencies = [
|
|||
"rand_core 0.6.4",
|
||||
"serde",
|
||||
"serde_json",
|
||||
"shared-types",
|
||||
"swactor",
|
||||
"tokio",
|
||||
]
|
||||
|
|
@ -4094,6 +4095,14 @@ dependencies = [
|
|||
"lazy_static",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "shared-types"
|
||||
version = "0.1.0"
|
||||
dependencies = [
|
||||
"blake3",
|
||||
"serde",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "shlex"
|
||||
version = "1.3.0"
|
||||
|
|
@ -4387,6 +4396,7 @@ dependencies = [
|
|||
"runtime-dashboard",
|
||||
"serde",
|
||||
"serde_json",
|
||||
"shared-types",
|
||||
"swactor",
|
||||
"swactor-std",
|
||||
"tempfile",
|
||||
|
|
|
|||
|
|
@ -1,5 +1,5 @@
|
|||
[workspace]
|
||||
members = [".", "crates/python", "crates/wasm", "crates/bin-runner", "crates/simulation", "crates/runtime-dashboard", "crates/distribution", "crates/std", "crates/datastore", "tests/docker", "xtask"]
|
||||
members = [".", "crates/python", "crates/wasm", "crates/bin-runner", "crates/simulation", "crates/runtime-dashboard", "crates/distribution", "crates/std", "crates/datastore", "crates/shared-types", "tests/docker", "xtask"]
|
||||
exclude = ["tools/depgraph"]
|
||||
|
||||
[package]
|
||||
|
|
|
|||
|
|
@ -145,10 +145,10 @@ SignedRequestPayload {
|
|||
}
|
||||
|
||||
DatastoreAction = enum {
|
||||
Put { path, content_hash, size_bytes, tags },
|
||||
Get { path },
|
||||
Delete { path },
|
||||
List { prefix },
|
||||
Put { name, content_hash, size_bytes, tags },
|
||||
Get { content_hash },
|
||||
Delete { content_hash },
|
||||
List { name_filter },
|
||||
}
|
||||
```
|
||||
|
||||
|
|
@ -285,12 +285,12 @@ Each protocol flow from `DATASTORE_PROTOCOL.md` §6 has a clear auth integration
|
|||
|
||||
| Protocol Flow | Auth Path 1 (Direct) | Auth Path 2 (Signed Request) |
|
||||
|---------------|----------------------|------------------------------|
|
||||
| §6.1 PUT | Connection-level ACL check | `SignedRequest { action: Put { path, content_hash, size_bytes, tags }, .. }` |
|
||||
| §6.2 GET (Local) | Connection-level ACL check | `SignedRequest { action: Get { path }, .. }` |
|
||||
| §6.3 GET (Remote) | Connection-level ACL check | `SignedRequest { action: Get { path }, .. }` → node handles remote fetch internally |
|
||||
| §6.4 DELETE | Connection-level ACL check | `SignedRequest { action: Delete { path }, .. }` |
|
||||
| §6.5 LIST (Local) | Connection-level ACL check | `SignedRequest { action: List { prefix }, .. }` |
|
||||
| §6.6 LIST (Swarm-Wide) | Connection-level ACL check | `SignedRequest { action: List { prefix }, .. }` → node handles fan-out internally |
|
||||
| §6.1 PUT | Connection-level ACL check | `SignedRequest { action: Put { name, content_hash, size_bytes, tags }, .. }` |
|
||||
| §6.2 GET (Local) | Connection-level ACL check | `SignedRequest { action: Get { content_hash }, .. }` |
|
||||
| §6.3 GET (Remote) | Connection-level ACL check | `SignedRequest { action: Get { content_hash }, .. }` → node handles remote fetch internally |
|
||||
| §6.4 DELETE | Connection-level ACL check | `SignedRequest { action: Delete { content_hash }, .. }` |
|
||||
| §6.5 LIST (Local) | Connection-level ACL check | `SignedRequest { action: List { name_filter }, .. }` |
|
||||
| §6.6 LIST (Swarm-Wide) | Connection-level ACL check | `SignedRequest { action: List { name_filter }, .. }` → node handles fan-out internally |
|
||||
|
||||
In all cases, auth is enforced *before* the request reaches the actor system. Internal inter-node communication (DHT replication, chunk transfers between cluster members) is not subject to auth checks.
|
||||
|
||||
|
|
|
|||
|
|
@ -6,6 +6,7 @@ edition = "2024"
|
|||
[dependencies]
|
||||
swactor = { path = "../..", features = ["serde", "transport"] }
|
||||
distribution = { path = "../distribution" }
|
||||
shared-types = { path = "../shared-types" }
|
||||
serde = { version = "1", features = ["derive"] }
|
||||
serde_json = "1"
|
||||
blake3 = "1"
|
||||
|
|
|
|||
170
crates/datastore/src/actors/gateway.rs
Normal file
170
crates/datastore/src/actors/gateway.rs
Normal file
|
|
@ -0,0 +1,170 @@
|
|||
//! GatewayActor — auth enforcement point for the datastore.
|
||||
//!
|
||||
//! Sits in front of the `DatastoreNode` coordinator. All external requests
|
||||
//! pass through the gateway, which checks authorization before forwarding
|
||||
//! to the internal actors.
|
||||
//!
|
||||
//! ```text
|
||||
//! External Client → GatewayActor → DatastoreNode → MetadataActor/BlobStoreActor
|
||||
//! (auth check) (dispatch) (auth-unaware)
|
||||
//! ```
|
||||
|
||||
use std::path::PathBuf;
|
||||
use std::time::{SystemTime, UNIX_EPOCH};
|
||||
|
||||
use swactor::actor::{ActorAddress, ActorInterface, Ctx};
|
||||
|
||||
use crate::auth::{AuthzEngine, AuthzResult, DatastoreAction};
|
||||
use crate::messages::{DatastoreNodeMsg, DatastoreResponse, GatewayMsg};
|
||||
|
||||
/// The auth gateway actor wrapping an `AuthzEngine`.
|
||||
pub struct GatewayActor {
|
||||
engine: AuthzEngine,
|
||||
datastore_node: ActorAddress,
|
||||
acl_path: Option<PathBuf>,
|
||||
}
|
||||
|
||||
impl GatewayActor {
|
||||
pub fn new(
|
||||
engine: AuthzEngine,
|
||||
datastore_node: ActorAddress,
|
||||
acl_path: Option<PathBuf>,
|
||||
) -> Self {
|
||||
Self {
|
||||
engine,
|
||||
datastore_node,
|
||||
acl_path,
|
||||
}
|
||||
}
|
||||
|
||||
fn now_secs() -> u64 {
|
||||
SystemTime::now()
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.unwrap_or_default()
|
||||
.as_secs()
|
||||
}
|
||||
|
||||
fn persist_acl(&self) {
|
||||
if let Some(ref path) = self.acl_path {
|
||||
let _ = self.engine.acl.save(path);
|
||||
}
|
||||
}
|
||||
|
||||
fn handle_signed_request(&mut self, ctx: &Ctx, request: crate::auth::SignedRequest, reply_to: ActorAddress) {
|
||||
let now = Self::now_secs();
|
||||
match self.engine.check_signed_request(&request, now) {
|
||||
AuthzResult::Allowed => {
|
||||
let msg = action_to_node_msg(request.payload.action, reply_to);
|
||||
let _ = ctx.send(self.datastore_node, msg);
|
||||
}
|
||||
AuthzResult::Denied(reason) => {
|
||||
let _ = ctx.send(reply_to, DatastoreResponse::Denied { reason });
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn handle_check_connection(&self, ctx: &Ctx, node_id: distribution::types::NodeId, reply_to: ActorAddress) {
|
||||
match self.engine.check_node(&node_id) {
|
||||
AuthzResult::Allowed => {
|
||||
let _ = ctx.send(reply_to, DatastoreResponse::Bool(true));
|
||||
}
|
||||
AuthzResult::Denied(reason) => {
|
||||
let _ = ctx.send(reply_to, DatastoreResponse::Denied { reason });
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn handle_grant(&mut self, ctx: &Ctx, requester: distribution::types::NodeId, key: distribution::types::NodeId, reply_to: ActorAddress) {
|
||||
match self.engine.grant(&requester, key) {
|
||||
Ok(()) => {
|
||||
self.persist_acl();
|
||||
let _ = ctx.send(reply_to, DatastoreResponse::Bool(true));
|
||||
}
|
||||
Err(reason) => {
|
||||
let _ = ctx.send(reply_to, DatastoreResponse::Denied { reason });
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn handle_revoke(&mut self, ctx: &Ctx, requester: distribution::types::NodeId, key: distribution::types::NodeId, reply_to: ActorAddress) {
|
||||
match self.engine.revoke(&requester, key) {
|
||||
Ok(()) => {
|
||||
self.persist_acl();
|
||||
let _ = ctx.send(reply_to, DatastoreResponse::Bool(true));
|
||||
}
|
||||
Err(reason) => {
|
||||
let _ = ctx.send(reply_to, DatastoreResponse::Denied { reason });
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl ActorInterface for GatewayActor {
|
||||
type Incoming = GatewayMsg;
|
||||
type Response = DatastoreResponse;
|
||||
|
||||
fn handle(&mut self, ctx: &Ctx, msg: GatewayMsg) {
|
||||
match msg {
|
||||
GatewayMsg::HandleSignedRequest { request, reply_to } => {
|
||||
self.handle_signed_request(ctx, request, reply_to);
|
||||
}
|
||||
GatewayMsg::CheckConnection { node_id, reply_to } => {
|
||||
self.handle_check_connection(ctx, node_id, reply_to);
|
||||
}
|
||||
GatewayMsg::Grant {
|
||||
requester,
|
||||
key,
|
||||
reply_to,
|
||||
} => {
|
||||
self.handle_grant(ctx, requester, key, reply_to);
|
||||
}
|
||||
GatewayMsg::Revoke {
|
||||
requester,
|
||||
key,
|
||||
reply_to,
|
||||
} => {
|
||||
self.handle_revoke(ctx, requester, key, reply_to);
|
||||
}
|
||||
GatewayMsg::NonceGcTick => {
|
||||
self.engine.gc_nonces(Self::now_secs());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Translate a `DatastoreAction` into the corresponding `DatastoreNodeMsg`.
|
||||
fn action_to_node_msg(action: DatastoreAction, reply_to: ActorAddress) -> DatastoreNodeMsg {
|
||||
match action {
|
||||
DatastoreAction::Get { content_hash } => DatastoreNodeMsg::Get {
|
||||
content_hash,
|
||||
reply_to,
|
||||
},
|
||||
DatastoreAction::Delete { content_hash } => DatastoreNodeMsg::Delete {
|
||||
content_hash,
|
||||
reply_to,
|
||||
},
|
||||
DatastoreAction::List { name_filter } => DatastoreNodeMsg::List {
|
||||
name_filter,
|
||||
all: false,
|
||||
reply_to,
|
||||
},
|
||||
DatastoreAction::Put {
|
||||
name,
|
||||
content_hash: _,
|
||||
size_bytes: _,
|
||||
tags,
|
||||
} => {
|
||||
// Put via signed request is an authorization of the operation.
|
||||
// The actual data upload happens separately. We forward as a
|
||||
// zero-data Put — the DatastoreNode will handle the metadata.
|
||||
// In the full flow, the data is uploaded separately and the
|
||||
// signed request only authorizes it.
|
||||
DatastoreNodeMsg::Put {
|
||||
data: Vec::new(),
|
||||
name,
|
||||
tags,
|
||||
reply_to,
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -1,9 +1,11 @@
|
|||
pub mod blob_store;
|
||||
pub mod datastore_node;
|
||||
pub mod gateway;
|
||||
pub mod metadata;
|
||||
pub mod transfer;
|
||||
|
||||
pub use blob_store::BlobStoreActor;
|
||||
pub use datastore_node::DatastoreNode;
|
||||
pub use gateway::GatewayActor;
|
||||
pub use metadata::MetadataActor;
|
||||
pub use transfer::TransferActor;
|
||||
|
|
|
|||
254
crates/datastore/src/auth.rs
Normal file
254
crates/datastore/src/auth.rs
Normal file
|
|
@ -0,0 +1,254 @@
|
|||
//! Authorization types and engine for the distributed datastore.
|
||||
//!
|
||||
//! Enforces binary access control (authorized or not) using ed25519 identities.
|
||||
//! Two auth paths:
|
||||
//! - **Path 1 (Direct iroh):** connection-level `check_node` against the ACL.
|
||||
//! - **Path 2 (Browser relay):** per-request `check_signed_request` with
|
||||
//! signature, timestamp, nonce, and ACL verification.
|
||||
|
||||
use std::collections::{BTreeMap, HashMap, HashSet};
|
||||
use std::io;
|
||||
use std::path::Path;
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
use distribution::crypto;
|
||||
use distribution::types::{NodeId, Signature};
|
||||
use shared_types::ContentHash;
|
||||
|
||||
// ─── DatastoreAction ────────────────────────────────────────────────────────
|
||||
|
||||
/// An action a client wants to perform on the datastore.
|
||||
///
|
||||
/// Carried inside a `SignedRequestPayload` for browser-relay auth (Auth Path 2).
|
||||
/// Aligned to match `DatastoreNodeMsg` variants — content-hash-first addressing.
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub enum DatastoreAction {
|
||||
Put {
|
||||
name: Option<String>,
|
||||
content_hash: ContentHash,
|
||||
size_bytes: u64,
|
||||
tags: BTreeMap<String, String>,
|
||||
},
|
||||
Get {
|
||||
content_hash: ContentHash,
|
||||
},
|
||||
Delete {
|
||||
content_hash: ContentHash,
|
||||
},
|
||||
List {
|
||||
name_filter: Option<String>,
|
||||
},
|
||||
}
|
||||
|
||||
// ─── SignedRequestPayload ───────────────────────────────────────────────────
|
||||
|
||||
/// The signable payload of a client request.
|
||||
///
|
||||
/// Serialized canonically (serde_json) and signed by the client's ed25519 key.
|
||||
/// Includes timestamp and nonce for replay protection.
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct SignedRequestPayload {
|
||||
pub action: DatastoreAction,
|
||||
/// Unix timestamp in seconds.
|
||||
pub timestamp: u64,
|
||||
/// 16 random bytes — prevents replay within the timestamp window.
|
||||
pub nonce: [u8; 16],
|
||||
}
|
||||
|
||||
// ─── SignedRequest ──────────────────────────────────────────────────────────
|
||||
|
||||
/// A signed request envelope for browser-relay auth (Auth Path 2).
|
||||
///
|
||||
/// The relay forwards this opaquely — it cannot forge, modify, or replay it.
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct SignedRequest {
|
||||
pub payload: SignedRequestPayload,
|
||||
/// The client's ed25519 public key.
|
||||
pub public_key: NodeId,
|
||||
/// ed25519 signature over the canonical serialization of `payload`.
|
||||
pub signature: Signature,
|
||||
}
|
||||
|
||||
// ─── AccessControlList ──────────────────────────────────────────────────────
|
||||
|
||||
/// The datastore's access control list.
|
||||
///
|
||||
/// Persisted as `acl.json` alongside the datastore's `storage_path`.
|
||||
/// The owner always has implicit full access.
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct AccessControlList {
|
||||
/// The datastore owner's public key — always has full access.
|
||||
pub owner: NodeId,
|
||||
/// Explicitly authorized client keys.
|
||||
pub authorized_keys: HashSet<NodeId>,
|
||||
}
|
||||
|
||||
impl AccessControlList {
|
||||
/// Load an ACL from disk, or create a default one with the given owner.
|
||||
pub fn load_or_create(path: &Path, owner: NodeId) -> io::Result<Self> {
|
||||
if path.exists() {
|
||||
let data = std::fs::read_to_string(path)?;
|
||||
let acl: AccessControlList = serde_json::from_str(&data)
|
||||
.map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
|
||||
Ok(acl)
|
||||
} else {
|
||||
let acl = AccessControlList {
|
||||
owner,
|
||||
authorized_keys: HashSet::new(),
|
||||
};
|
||||
acl.save(path)?;
|
||||
Ok(acl)
|
||||
}
|
||||
}
|
||||
|
||||
/// Persist the ACL to disk as JSON.
|
||||
pub fn save(&self, path: &Path) -> io::Result<()> {
|
||||
if let Some(parent) = path.parent() {
|
||||
std::fs::create_dir_all(parent)?;
|
||||
}
|
||||
let json = serde_json::to_string_pretty(self)
|
||||
.map_err(|e| io::Error::new(io::ErrorKind::Other, e))?;
|
||||
std::fs::write(path, json)
|
||||
}
|
||||
}
|
||||
|
||||
// ─── AuthzResult ────────────────────────────────────────────────────────────
|
||||
|
||||
/// The outcome of an authorization check.
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub enum AuthzResult {
|
||||
Allowed,
|
||||
Denied(DeniedReason),
|
||||
}
|
||||
|
||||
/// Why a request was denied.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub enum DeniedReason {
|
||||
/// The key is not in the ACL.
|
||||
NotAuthorized,
|
||||
/// The ed25519 signature is invalid.
|
||||
InvalidSignature,
|
||||
/// The request timestamp is outside the ±300s window.
|
||||
RequestExpired,
|
||||
/// The nonce has already been seen within the time window.
|
||||
ReplayDetected,
|
||||
}
|
||||
|
||||
// ─── Signing / Verification ─────────────────────────────────────────────────
|
||||
|
||||
/// Sign a request payload, returning a complete `SignedRequest` envelope.
|
||||
pub fn sign_request(keypair: &crypto::Keypair, payload: SignedRequestPayload) -> SignedRequest {
|
||||
let bytes = serde_json::to_vec(&payload).expect("SignedRequestPayload is always serializable");
|
||||
let signature = keypair.sign(&bytes);
|
||||
SignedRequest {
|
||||
payload,
|
||||
public_key: keypair.node_id(),
|
||||
signature,
|
||||
}
|
||||
}
|
||||
|
||||
/// Verify a `SignedRequest`'s signature against its embedded `public_key`.
|
||||
///
|
||||
/// Checks only signature validity — does NOT check timestamp, nonce, or ACL.
|
||||
pub fn verify_signed_request(request: &SignedRequest) -> bool {
|
||||
let Ok(bytes) = serde_json::to_vec(&request.payload) else {
|
||||
return false;
|
||||
};
|
||||
crypto::verify(&request.public_key, &bytes, &request.signature)
|
||||
}
|
||||
|
||||
// ─── AuthzEngine ────────────────────────────────────────────────────────────
|
||||
|
||||
/// Authorization engine — checks requests against the ACL and replay state.
|
||||
///
|
||||
/// Sits at the edge of the actor system (Auth Gate / GatewayActor) and decides
|
||||
/// whether to accept or reject external requests before they reach internal actors.
|
||||
#[derive(Debug)]
|
||||
pub struct AuthzEngine {
|
||||
pub acl: AccessControlList,
|
||||
seen_nonces: HashMap<[u8; 16], u64>,
|
||||
timestamp_window: u64,
|
||||
}
|
||||
|
||||
impl AuthzEngine {
|
||||
/// Create a new engine with the given ACL and a default 300-second window.
|
||||
pub fn new(acl: AccessControlList) -> Self {
|
||||
Self {
|
||||
acl,
|
||||
seen_nonces: HashMap::new(),
|
||||
timestamp_window: 300,
|
||||
}
|
||||
}
|
||||
|
||||
/// Check whether a `NodeId` is authorized (connection-level, Auth Path 1).
|
||||
///
|
||||
/// The owner always has implicit access. Other keys must be in `authorized_keys`.
|
||||
pub fn check_node(&self, node_id: &NodeId) -> AuthzResult {
|
||||
if *node_id == self.acl.owner || self.acl.authorized_keys.contains(node_id) {
|
||||
AuthzResult::Allowed
|
||||
} else {
|
||||
AuthzResult::Denied(DeniedReason::NotAuthorized)
|
||||
}
|
||||
}
|
||||
|
||||
/// Verify and authorize a signed request (Auth Path 2).
|
||||
///
|
||||
/// Four-step verification in strict order:
|
||||
/// 1. Signature validity
|
||||
/// 2. Timestamp freshness (±window)
|
||||
/// 3. Nonce uniqueness
|
||||
/// 4. ACL check
|
||||
pub fn check_signed_request(&mut self, request: &SignedRequest, now: u64) -> AuthzResult {
|
||||
// 1. Signature
|
||||
if !verify_signed_request(request) {
|
||||
return AuthzResult::Denied(DeniedReason::InvalidSignature);
|
||||
}
|
||||
|
||||
// 2. Timestamp freshness
|
||||
let ts = request.payload.timestamp;
|
||||
let diff = if now >= ts { now - ts } else { ts - now };
|
||||
if diff > self.timestamp_window {
|
||||
return AuthzResult::Denied(DeniedReason::RequestExpired);
|
||||
}
|
||||
|
||||
// 3. Nonce uniqueness
|
||||
if self.seen_nonces.contains_key(&request.payload.nonce) {
|
||||
return AuthzResult::Denied(DeniedReason::ReplayDetected);
|
||||
}
|
||||
self.seen_nonces.insert(request.payload.nonce, ts);
|
||||
|
||||
// 4. ACL check
|
||||
self.check_node(&request.public_key)
|
||||
}
|
||||
|
||||
/// Grant access to a `NodeId`. Owner-only, idempotent.
|
||||
pub fn grant(&mut self, requester: &NodeId, key: NodeId) -> Result<(), DeniedReason> {
|
||||
if *requester != self.acl.owner {
|
||||
return Err(DeniedReason::NotAuthorized);
|
||||
}
|
||||
self.acl.authorized_keys.insert(key);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Revoke access from a `NodeId`. Owner-only, idempotent.
|
||||
/// Revoking the owner is a no-op (owner's implicit access cannot be removed).
|
||||
pub fn revoke(&mut self, requester: &NodeId, key: NodeId) -> Result<(), DeniedReason> {
|
||||
if *requester != self.acl.owner {
|
||||
return Err(DeniedReason::NotAuthorized);
|
||||
}
|
||||
// Owner's implicit access cannot be removed.
|
||||
if key != self.acl.owner {
|
||||
self.acl.authorized_keys.remove(&key);
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Evict nonces whose timestamps fall outside the current window.
|
||||
pub fn gc_nonces(&mut self, now: u64) {
|
||||
self.seen_nonces.retain(|_nonce, ts| {
|
||||
let diff = if now >= *ts { now - *ts } else { *ts - now };
|
||||
diff <= self.timestamp_window
|
||||
});
|
||||
}
|
||||
}
|
||||
|
|
@ -3,6 +3,7 @@ pub mod messages;
|
|||
pub mod chunking;
|
||||
pub mod storage;
|
||||
pub mod actors;
|
||||
pub mod auth;
|
||||
pub mod cli;
|
||||
pub mod metrics;
|
||||
#[cfg(feature = "node")]
|
||||
|
|
|
|||
|
|
@ -16,6 +16,7 @@ use swactor::transport::NetworkMessage;
|
|||
|
||||
use distribution::types::NodeId;
|
||||
|
||||
use crate::auth::{DeniedReason, SignedRequest};
|
||||
use crate::types::{ContentHash, ObjectEntry, ObjectManifest};
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
|
|
@ -365,8 +366,41 @@ pub enum DatastoreResponse {
|
|||
NotFound,
|
||||
/// An error occurred.
|
||||
Error { reason: String },
|
||||
/// Request was denied by the auth layer.
|
||||
Denied { reason: DeniedReason },
|
||||
/// Boolean response (e.g. HasChunk).
|
||||
Bool(bool),
|
||||
/// List of chunk hashes.
|
||||
ChunkList { hashes: Vec<ContentHash> },
|
||||
}
|
||||
|
||||
// ─── GatewayMsg ────────────────────────────────────────────────────────────
|
||||
|
||||
/// Messages handled by the `GatewayActor` — the auth enforcement point.
|
||||
#[derive(Debug, Clone)]
|
||||
pub enum GatewayMsg {
|
||||
/// Auth Path 2: verify a signed request and dispatch if allowed.
|
||||
HandleSignedRequest {
|
||||
request: SignedRequest,
|
||||
reply_to: ActorAddress,
|
||||
},
|
||||
/// Auth Path 1: check whether a node is authorized for connection.
|
||||
CheckConnection {
|
||||
node_id: NodeId,
|
||||
reply_to: ActorAddress,
|
||||
},
|
||||
/// Owner-only: grant access to a key.
|
||||
Grant {
|
||||
requester: NodeId,
|
||||
key: NodeId,
|
||||
reply_to: ActorAddress,
|
||||
},
|
||||
/// Owner-only: revoke access from a key.
|
||||
Revoke {
|
||||
requester: NodeId,
|
||||
key: NodeId,
|
||||
reply_to: ActorAddress,
|
||||
},
|
||||
/// Periodic nonce garbage collection tick.
|
||||
NonceGcTick,
|
||||
}
|
||||
|
|
|
|||
|
|
@ -4,111 +4,13 @@
|
|||
//! `blake3(blob_bytes)`. Names are optional metadata, not keys.
|
||||
|
||||
use std::collections::BTreeMap;
|
||||
use std::fmt;
|
||||
use std::path::PathBuf;
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
use distribution::types::NodeId;
|
||||
|
||||
// ─── ContentHash ────────────────────────────────────────────────────────────
|
||||
|
||||
/// A blake3 content hash (32 bytes).
|
||||
///
|
||||
/// The primary identifier for blobs and the DHT key. Mirrors the `NodeId`
|
||||
/// pattern from `distribution::types` — XOR distance for DHT routing, compact
|
||||
/// Debug/Display for logging.
|
||||
#[derive(Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
|
||||
pub struct ContentHash(pub [u8; 32]);
|
||||
|
||||
impl ContentHash {
|
||||
/// Compute the blake3 hash of the given data.
|
||||
pub fn of(data: &[u8]) -> Self {
|
||||
let hash = blake3::hash(data);
|
||||
ContentHash(*hash.as_bytes())
|
||||
}
|
||||
|
||||
/// XOR distance between two content hashes (Kademlia metric).
|
||||
pub fn xor_distance(&self, other: &ContentHash) -> [u8; 32] {
|
||||
let mut out = [0u8; 32];
|
||||
for i in 0..32 {
|
||||
out[i] = self.0[i] ^ other.0[i];
|
||||
}
|
||||
out
|
||||
}
|
||||
|
||||
/// Number of leading zero bits in the XOR distance to `other`.
|
||||
/// Returns 0..=256. Used to select the k-bucket index in the metadata DHT.
|
||||
pub fn xor_leading_zeros(&self, other: &ContentHash) -> u32 {
|
||||
let dist = self.xor_distance(other);
|
||||
let mut zeros = 0u32;
|
||||
for byte in dist {
|
||||
if byte == 0 {
|
||||
zeros += 8;
|
||||
} else {
|
||||
zeros += byte.leading_zeros();
|
||||
break;
|
||||
}
|
||||
}
|
||||
zeros
|
||||
}
|
||||
|
||||
/// Parse a 64-character hex string into a ContentHash.
|
||||
/// Returns `None` if the string is not exactly 64 hex characters.
|
||||
pub fn from_hex(hex: &str) -> Option<Self> {
|
||||
if hex.len() != 64 {
|
||||
return None;
|
||||
}
|
||||
let mut bytes = [0u8; 32];
|
||||
for (i, chunk) in hex.as_bytes().chunks(2).enumerate() {
|
||||
let hi = hex_digit(chunk[0])?;
|
||||
let lo = hex_digit(chunk[1])?;
|
||||
bytes[i] = (hi << 4) | lo;
|
||||
}
|
||||
Some(ContentHash(bytes))
|
||||
}
|
||||
|
||||
/// Encode as lowercase hex string.
|
||||
pub fn to_hex(&self) -> String {
|
||||
let mut s = String::with_capacity(64);
|
||||
for b in &self.0 {
|
||||
use fmt::Write;
|
||||
write!(s, "{:02x}", b).unwrap();
|
||||
}
|
||||
s
|
||||
}
|
||||
|
||||
/// The zero hash (all zeroes). Used as a sentinel.
|
||||
pub const ZERO: ContentHash = ContentHash([0u8; 32]);
|
||||
}
|
||||
|
||||
impl fmt::Debug for ContentHash {
|
||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
write!(f, "Hash(")?;
|
||||
for b in &self.0[..4] {
|
||||
write!(f, "{:02x}", b)?;
|
||||
}
|
||||
write!(f, "\u{2026})")
|
||||
}
|
||||
}
|
||||
|
||||
impl fmt::Display for ContentHash {
|
||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
for b in &self.0[..8] {
|
||||
write!(f, "{:02x}", b)?;
|
||||
}
|
||||
write!(f, "\u{2026}")
|
||||
}
|
||||
}
|
||||
|
||||
fn hex_digit(b: u8) -> Option<u8> {
|
||||
match b {
|
||||
b'0'..=b'9' => Some(b - b'0'),
|
||||
b'a'..=b'f' => Some(b - b'a' + 10),
|
||||
b'A'..=b'F' => Some(b - b'A' + 10),
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
pub use shared_types::ContentHash;
|
||||
|
||||
// ─── ObjectEntry ────────────────────────────────────────────────────────────
|
||||
|
||||
|
|
|
|||
58
crates/datastore/tests/acl_persistence_tests.rs
Normal file
58
crates/datastore/tests/acl_persistence_tests.rs
Normal file
|
|
@ -0,0 +1,58 @@
|
|||
//! ACL file persistence tests — roundtrip save/load.
|
||||
|
||||
use std::collections::HashSet;
|
||||
|
||||
use distribution::crypto::Keypair;
|
||||
use swactor_datastore::auth::AccessControlList;
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
// 1. Save + load preserves owner and authorized_keys
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
|
||||
#[test]
|
||||
fn save_and_load_roundtrip() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let path = dir.path().join("acl.json");
|
||||
|
||||
let owner = Keypair::generate().node_id();
|
||||
let client_a = Keypair::generate().node_id();
|
||||
let client_b = Keypair::generate().node_id();
|
||||
|
||||
let mut keys = HashSet::new();
|
||||
keys.insert(client_a);
|
||||
keys.insert(client_b);
|
||||
|
||||
let acl = AccessControlList {
|
||||
owner,
|
||||
authorized_keys: keys.clone(),
|
||||
};
|
||||
acl.save(&path).unwrap();
|
||||
|
||||
let loaded = AccessControlList::load_or_create(&path, owner).unwrap();
|
||||
assert_eq!(loaded.owner, owner);
|
||||
assert_eq!(loaded.authorized_keys, keys);
|
||||
}
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
// 2. load_or_create on missing file creates default
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
|
||||
#[test]
|
||||
fn load_or_create_on_missing_file_creates_default() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let path = dir.path().join("nonexistent/acl.json");
|
||||
|
||||
let owner = Keypair::generate().node_id();
|
||||
let acl = AccessControlList::load_or_create(&path, owner).unwrap();
|
||||
|
||||
assert_eq!(acl.owner, owner);
|
||||
assert!(acl.authorized_keys.is_empty());
|
||||
|
||||
// File should now exist
|
||||
assert!(path.exists());
|
||||
|
||||
// Loading again should give same result
|
||||
let acl2 = AccessControlList::load_or_create(&path, owner).unwrap();
|
||||
assert_eq!(acl2.owner, owner);
|
||||
assert!(acl2.authorized_keys.is_empty());
|
||||
}
|
||||
291
crates/datastore/tests/auth_scenario_tests.rs
Normal file
291
crates/datastore/tests/auth_scenario_tests.rs
Normal file
|
|
@ -0,0 +1,291 @@
|
|||
//! Scenario tests for AuthzEngine — no actor system, pure auth logic.
|
||||
|
||||
use std::collections::HashSet;
|
||||
|
||||
use distribution::crypto::Keypair;
|
||||
use shared_types::ContentHash;
|
||||
use swactor_datastore::auth::{
|
||||
sign_request, AccessControlList, AuthzEngine, AuthzResult, DatastoreAction, DeniedReason,
|
||||
SignedRequestPayload,
|
||||
};
|
||||
|
||||
fn owner_engine() -> (Keypair, AuthzEngine) {
|
||||
let owner_kp = Keypair::generate();
|
||||
let acl = AccessControlList {
|
||||
owner: owner_kp.node_id(),
|
||||
authorized_keys: HashSet::new(),
|
||||
};
|
||||
(owner_kp, AuthzEngine::new(acl))
|
||||
}
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
// 1. Owner always allowed; random key denied
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
|
||||
#[test]
|
||||
fn owner_is_always_allowed() {
|
||||
let (owner_kp, engine) = owner_engine();
|
||||
assert_eq!(engine.check_node(&owner_kp.node_id()), AuthzResult::Allowed);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn unknown_key_is_denied() {
|
||||
let (_owner_kp, engine) = owner_engine();
|
||||
let stranger = Keypair::generate().node_id();
|
||||
assert_eq!(
|
||||
engine.check_node(&stranger),
|
||||
AuthzResult::Denied(DeniedReason::NotAuthorized)
|
||||
);
|
||||
}
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
// 2. Grant → access → revoke → denied (lifecycle story)
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
|
||||
#[test]
|
||||
fn grant_then_revoke_lifecycle() {
|
||||
let (owner_kp, mut engine) = owner_engine();
|
||||
let client = Keypair::generate().node_id();
|
||||
|
||||
// Initially denied
|
||||
assert_eq!(
|
||||
engine.check_node(&client),
|
||||
AuthzResult::Denied(DeniedReason::NotAuthorized)
|
||||
);
|
||||
|
||||
// Grant
|
||||
engine.grant(&owner_kp.node_id(), client).unwrap();
|
||||
assert_eq!(engine.check_node(&client), AuthzResult::Allowed);
|
||||
|
||||
// Revoke
|
||||
engine.revoke(&owner_kp.node_id(), client).unwrap();
|
||||
assert_eq!(
|
||||
engine.check_node(&client),
|
||||
AuthzResult::Denied(DeniedReason::NotAuthorized)
|
||||
);
|
||||
}
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
// 3. Only owner can grant/revoke; non-owner gets NotAuthorized
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
|
||||
#[test]
|
||||
fn non_owner_cannot_grant() {
|
||||
let (_owner_kp, mut engine) = owner_engine();
|
||||
let impostor = Keypair::generate().node_id();
|
||||
let target = Keypair::generate().node_id();
|
||||
|
||||
assert_eq!(
|
||||
engine.grant(&impostor, target),
|
||||
Err(DeniedReason::NotAuthorized)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn non_owner_cannot_revoke() {
|
||||
let (owner_kp, mut engine) = owner_engine();
|
||||
let client = Keypair::generate().node_id();
|
||||
engine.grant(&owner_kp.node_id(), client).unwrap();
|
||||
|
||||
let impostor = Keypair::generate().node_id();
|
||||
assert_eq!(
|
||||
engine.revoke(&impostor, client),
|
||||
Err(DeniedReason::NotAuthorized)
|
||||
);
|
||||
|
||||
// Client still authorized
|
||||
assert_eq!(engine.check_node(&client), AuthzResult::Allowed);
|
||||
}
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
// 4. Cannot revoke owner's implicit access
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
|
||||
#[test]
|
||||
fn revoking_owner_is_noop() {
|
||||
let (owner_kp, mut engine) = owner_engine();
|
||||
let owner_id = owner_kp.node_id();
|
||||
|
||||
// Attempt to revoke owner — should succeed (idempotent no-op) but owner remains allowed
|
||||
engine.revoke(&owner_id, owner_id).unwrap();
|
||||
assert_eq!(engine.check_node(&owner_id), AuthzResult::Allowed);
|
||||
}
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
// 5. Signed request happy path (sign → verify → allowed)
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
|
||||
#[test]
|
||||
fn signed_request_happy_path() {
|
||||
let (owner_kp, mut engine) = owner_engine();
|
||||
let now = 1_000_000u64;
|
||||
|
||||
let payload = SignedRequestPayload {
|
||||
action: DatastoreAction::Get {
|
||||
content_hash: ContentHash::of(b"hello"),
|
||||
},
|
||||
timestamp: now,
|
||||
nonce: [1; 16],
|
||||
};
|
||||
let request = sign_request(&owner_kp, payload);
|
||||
|
||||
assert_eq!(
|
||||
engine.check_signed_request(&request, now),
|
||||
AuthzResult::Allowed
|
||||
);
|
||||
}
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
// 6. Tampered signature → InvalidSignature
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
|
||||
#[test]
|
||||
fn tampered_signature_is_rejected() {
|
||||
let (owner_kp, mut engine) = owner_engine();
|
||||
let now = 1_000_000u64;
|
||||
|
||||
let payload = SignedRequestPayload {
|
||||
action: DatastoreAction::Get {
|
||||
content_hash: ContentHash::of(b"hello"),
|
||||
},
|
||||
timestamp: now,
|
||||
nonce: [2; 16],
|
||||
};
|
||||
let mut request = sign_request(&owner_kp, payload);
|
||||
// Tamper with signature
|
||||
request.signature.0[0] ^= 0xFF;
|
||||
|
||||
assert_eq!(
|
||||
engine.check_signed_request(&request, now),
|
||||
AuthzResult::Denied(DeniedReason::InvalidSignature)
|
||||
);
|
||||
}
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
// 7. Stale timestamp → RequestExpired
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
|
||||
#[test]
|
||||
fn stale_timestamp_is_rejected() {
|
||||
let (owner_kp, mut engine) = owner_engine();
|
||||
let now = 1_000_000u64;
|
||||
|
||||
let payload = SignedRequestPayload {
|
||||
action: DatastoreAction::Get {
|
||||
content_hash: ContentHash::of(b"stale"),
|
||||
},
|
||||
timestamp: now - 400, // 400s ago, outside 300s window
|
||||
nonce: [3; 16],
|
||||
};
|
||||
let request = sign_request(&owner_kp, payload);
|
||||
|
||||
assert_eq!(
|
||||
engine.check_signed_request(&request, now),
|
||||
AuthzResult::Denied(DeniedReason::RequestExpired)
|
||||
);
|
||||
}
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
// 8. Replayed nonce → ReplayDetected
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
|
||||
#[test]
|
||||
fn replayed_nonce_is_rejected() {
|
||||
let (owner_kp, mut engine) = owner_engine();
|
||||
let now = 1_000_000u64;
|
||||
|
||||
let payload = SignedRequestPayload {
|
||||
action: DatastoreAction::Get {
|
||||
content_hash: ContentHash::of(b"first"),
|
||||
},
|
||||
timestamp: now,
|
||||
nonce: [4; 16],
|
||||
};
|
||||
let request = sign_request(&owner_kp, payload);
|
||||
|
||||
// First time — allowed
|
||||
assert_eq!(
|
||||
engine.check_signed_request(&request, now),
|
||||
AuthzResult::Allowed
|
||||
);
|
||||
|
||||
// Replay — same nonce
|
||||
let payload2 = SignedRequestPayload {
|
||||
action: DatastoreAction::Get {
|
||||
content_hash: ContentHash::of(b"first"),
|
||||
},
|
||||
timestamp: now,
|
||||
nonce: [4; 16],
|
||||
};
|
||||
let request2 = sign_request(&owner_kp, payload2);
|
||||
assert_eq!(
|
||||
engine.check_signed_request(&request2, now),
|
||||
AuthzResult::Denied(DeniedReason::ReplayDetected)
|
||||
);
|
||||
}
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
// 9. Nonce GC frees old nonces for reuse
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
|
||||
#[test]
|
||||
fn nonce_gc_frees_old_nonces() {
|
||||
let (owner_kp, mut engine) = owner_engine();
|
||||
let t0 = 1_000_000u64;
|
||||
|
||||
let payload = SignedRequestPayload {
|
||||
action: DatastoreAction::Get {
|
||||
content_hash: ContentHash::of(b"gc-test"),
|
||||
},
|
||||
timestamp: t0,
|
||||
nonce: [5; 16],
|
||||
};
|
||||
let request = sign_request(&owner_kp, payload);
|
||||
assert_eq!(
|
||||
engine.check_signed_request(&request, t0),
|
||||
AuthzResult::Allowed
|
||||
);
|
||||
|
||||
// Advance time past the window and GC
|
||||
let t1 = t0 + 400;
|
||||
engine.gc_nonces(t1);
|
||||
|
||||
// Same nonce but with current timestamp — no longer flagged as replay
|
||||
let payload2 = SignedRequestPayload {
|
||||
action: DatastoreAction::Get {
|
||||
content_hash: ContentHash::of(b"gc-test"),
|
||||
},
|
||||
timestamp: t1,
|
||||
nonce: [5; 16],
|
||||
};
|
||||
let request2 = sign_request(&owner_kp, payload2);
|
||||
assert_eq!(
|
||||
engine.check_signed_request(&request2, t1),
|
||||
AuthzResult::Allowed
|
||||
);
|
||||
}
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
// 10. Unauthorized key with valid signature → NotAuthorized
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
|
||||
#[test]
|
||||
fn unauthorized_key_with_valid_signature_is_denied() {
|
||||
let (_owner_kp, mut engine) = owner_engine();
|
||||
let stranger_kp = Keypair::generate();
|
||||
let now = 1_000_000u64;
|
||||
|
||||
let payload = SignedRequestPayload {
|
||||
action: DatastoreAction::Get {
|
||||
content_hash: ContentHash::of(b"intrusion"),
|
||||
},
|
||||
timestamp: now,
|
||||
nonce: [6; 16],
|
||||
};
|
||||
let request = sign_request(&stranger_kp, payload);
|
||||
|
||||
assert_eq!(
|
||||
engine.check_signed_request(&request, now),
|
||||
AuthzResult::Denied(DeniedReason::NotAuthorized)
|
||||
);
|
||||
}
|
||||
202
crates/datastore/tests/gateway_tests.rs
Normal file
202
crates/datastore/tests/gateway_tests.rs
Normal file
|
|
@ -0,0 +1,202 @@
|
|||
//! Actor-level tests for the GatewayActor.
|
||||
//!
|
||||
//! Uses the swactor runtime to spawn GatewayActor + DatastoreNode and verify
|
||||
//! that authorized requests flow through while unauthorized ones are denied.
|
||||
|
||||
mod common;
|
||||
|
||||
use std::collections::HashSet;
|
||||
use std::time::{SystemTime, UNIX_EPOCH};
|
||||
|
||||
use common::{spawn_blob_store, spawn_metadata, test_runtime, tick_until_recv};
|
||||
|
||||
use distribution::crypto::Keypair;
|
||||
use shared_types::ContentHash;
|
||||
use swactor_datastore::actors::{DatastoreNode, GatewayActor};
|
||||
use swactor_datastore::auth::{
|
||||
sign_request, AccessControlList, AuthzEngine, DatastoreAction, DeniedReason,
|
||||
SignedRequestPayload,
|
||||
};
|
||||
use swactor_datastore::messages::{DatastoreResponse, GatewayMsg};
|
||||
use swactor_datastore::types::DatastoreConfig;
|
||||
|
||||
struct GatewayHarness {
|
||||
rt: swactor::runtime::Runtime,
|
||||
gateway: swactor::actor::ActorAddress,
|
||||
inbox: swactor::runtime::Inbox<DatastoreResponse>,
|
||||
owner_kp: Keypair,
|
||||
}
|
||||
|
||||
impl GatewayHarness {
|
||||
fn new() -> Self {
|
||||
let owner_kp = Keypair::generate();
|
||||
let rt = test_runtime();
|
||||
|
||||
let blob_store = spawn_blob_store(&rt);
|
||||
let node_id = owner_kp.node_id();
|
||||
let metadata = spawn_metadata(&rt, node_id);
|
||||
|
||||
let mut config = DatastoreConfig::default();
|
||||
config.chunk_size = 64;
|
||||
let datastore_node = rt
|
||||
.spawn(DatastoreNode::new(node_id, blob_store, metadata, config))
|
||||
.unwrap();
|
||||
|
||||
let acl = AccessControlList {
|
||||
owner: owner_kp.node_id(),
|
||||
authorized_keys: HashSet::new(),
|
||||
};
|
||||
let engine = AuthzEngine::new(acl);
|
||||
let gateway = rt
|
||||
.spawn(GatewayActor::new(engine, datastore_node, None))
|
||||
.unwrap();
|
||||
|
||||
let inbox = rt.new_inbox::<DatastoreResponse>().unwrap();
|
||||
|
||||
// Let actors initialize
|
||||
for _ in 0..3 {
|
||||
rt.tick();
|
||||
}
|
||||
|
||||
Self {
|
||||
rt,
|
||||
gateway,
|
||||
inbox,
|
||||
owner_kp,
|
||||
}
|
||||
}
|
||||
|
||||
fn reply_addr(&self) -> swactor::actor::ActorAddress {
|
||||
*self.inbox.addr()
|
||||
}
|
||||
}
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
// 1. Authorized signed GET dispatches and returns result
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
|
||||
#[test]
|
||||
fn authorized_signed_get_flows_through_to_datastore() {
|
||||
let h = GatewayHarness::new();
|
||||
let now = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_secs();
|
||||
|
||||
// PUT some data first via signed request
|
||||
let data = b"gateway test data";
|
||||
let content_hash = ContentHash::of(data);
|
||||
|
||||
// Store data by sending directly to datastore through a Put via gateway
|
||||
// (For simplicity, store via DatastoreNode first, then GET through gateway)
|
||||
// Actually, let's just do a GET for a nonexistent hash — we should get NotFound (not Denied)
|
||||
let payload = SignedRequestPayload {
|
||||
action: DatastoreAction::Get { content_hash },
|
||||
timestamp: now,
|
||||
nonce: [10; 16],
|
||||
};
|
||||
let request = sign_request(&h.owner_kp, payload);
|
||||
|
||||
h.rt.send_to(
|
||||
h.gateway,
|
||||
GatewayMsg::HandleSignedRequest {
|
||||
request,
|
||||
reply_to: h.reply_addr(),
|
||||
},
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
let resp = tick_until_recv(&h.rt, &h.inbox, 30).unwrap();
|
||||
// Should get NotFound (authorized, but object doesn't exist) — NOT Denied
|
||||
assert!(
|
||||
matches!(resp, DatastoreResponse::NotFound),
|
||||
"expected NotFound (authorized but missing), got {resp:?}"
|
||||
);
|
||||
}
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
// 2. Unauthorized signed GET returns Denied
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
|
||||
#[test]
|
||||
fn unauthorized_signed_get_returns_denied() {
|
||||
let h = GatewayHarness::new();
|
||||
let stranger_kp = Keypair::generate();
|
||||
let now = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_secs();
|
||||
|
||||
let payload = SignedRequestPayload {
|
||||
action: DatastoreAction::Get {
|
||||
content_hash: ContentHash::of(b"unauthorized"),
|
||||
},
|
||||
timestamp: now,
|
||||
nonce: [11; 16],
|
||||
};
|
||||
let request = sign_request(&stranger_kp, payload);
|
||||
|
||||
h.rt.send_to(
|
||||
h.gateway,
|
||||
GatewayMsg::HandleSignedRequest {
|
||||
request,
|
||||
reply_to: h.reply_addr(),
|
||||
},
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
let resp = tick_until_recv(&h.rt, &h.inbox, 30).unwrap();
|
||||
assert!(
|
||||
matches!(
|
||||
resp,
|
||||
DatastoreResponse::Denied {
|
||||
reason: DeniedReason::NotAuthorized
|
||||
}
|
||||
),
|
||||
"expected Denied(NotAuthorized), got {resp:?}"
|
||||
);
|
||||
}
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
// 3. Connection check allows/denies correctly
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
|
||||
#[test]
|
||||
fn check_connection_allows_owner() {
|
||||
let h = GatewayHarness::new();
|
||||
|
||||
h.rt.send_to(
|
||||
h.gateway,
|
||||
GatewayMsg::CheckConnection {
|
||||
node_id: h.owner_kp.node_id(),
|
||||
reply_to: h.reply_addr(),
|
||||
},
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
let resp = tick_until_recv(&h.rt, &h.inbox, 20).unwrap();
|
||||
assert!(
|
||||
matches!(resp, DatastoreResponse::Bool(true)),
|
||||
"expected Bool(true), got {resp:?}"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn check_connection_denies_stranger() {
|
||||
let h = GatewayHarness::new();
|
||||
let stranger = Keypair::generate().node_id();
|
||||
|
||||
h.rt.send_to(
|
||||
h.gateway,
|
||||
GatewayMsg::CheckConnection {
|
||||
node_id: stranger,
|
||||
reply_to: h.reply_addr(),
|
||||
},
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
let resp = tick_until_recv(&h.rt, &h.inbox, 20).unwrap();
|
||||
assert!(
|
||||
matches!(
|
||||
resp,
|
||||
DatastoreResponse::Denied {
|
||||
reason: DeniedReason::NotAuthorized
|
||||
}
|
||||
),
|
||||
"expected Denied(NotAuthorized), got {resp:?}"
|
||||
);
|
||||
}
|
||||
|
|
@ -10,6 +10,7 @@ iroh = ["dep:iroh", "dep:tokio"]
|
|||
|
||||
[dependencies]
|
||||
swactor = { path = "../..", features = ["serde", "transport"] }
|
||||
shared-types = { path = "../shared-types" }
|
||||
ed25519-dalek = { version = "2", features = ["rand_core"] }
|
||||
rand_core = { version = "0.6", features = ["getrandom"] }
|
||||
serde = { version = "1", features = ["derive"] }
|
||||
|
|
|
|||
|
|
@ -1,145 +0,0 @@
|
|||
use std::collections::HashSet;
|
||||
use std::fmt;
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
use crate::types::{NodeId, Signature};
|
||||
|
||||
// ─── ContentHash ────────────────────────────────────────────────────────────
|
||||
|
||||
/// A 32-byte blake3 digest used as content address for chunks and manifests.
|
||||
#[derive(Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
|
||||
pub struct ContentHash(pub [u8; 32]);
|
||||
|
||||
impl fmt::Debug for ContentHash {
|
||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
write!(f, "ContentHash(")?;
|
||||
for b in &self.0[..4] {
|
||||
write!(f, "{:02x}", b)?;
|
||||
}
|
||||
write!(f, "\u{2026})")
|
||||
}
|
||||
}
|
||||
|
||||
// ─── DatastoreAction ────────────────────────────────────────────────────────
|
||||
|
||||
/// An action a client wants to perform on the datastore.
|
||||
///
|
||||
/// Carried inside a `SignedRequestPayload` for browser-relay auth (Auth Path 2).
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub enum DatastoreAction {
|
||||
Put {
|
||||
path: String,
|
||||
content_hash: ContentHash,
|
||||
size_bytes: u64,
|
||||
tags: std::collections::BTreeMap<String, String>,
|
||||
},
|
||||
Get {
|
||||
path: String,
|
||||
},
|
||||
Delete {
|
||||
path: String,
|
||||
},
|
||||
List {
|
||||
prefix: Option<String>,
|
||||
},
|
||||
}
|
||||
|
||||
// ─── SignedRequestPayload ───────────────────────────────────────────────────
|
||||
|
||||
/// The signable payload of a client request.
|
||||
///
|
||||
/// Serialized canonically (serde_json) and signed by the client's ed25519 key.
|
||||
/// Includes timestamp and nonce for replay protection.
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct SignedRequestPayload {
|
||||
pub action: DatastoreAction,
|
||||
/// Unix timestamp in seconds.
|
||||
pub timestamp: u64,
|
||||
/// 16 random bytes — prevents replay within the timestamp window.
|
||||
pub nonce: [u8; 16],
|
||||
}
|
||||
|
||||
// ─── SignedRequest ──────────────────────────────────────────────────────────
|
||||
|
||||
/// A signed request envelope for browser-relay auth (Auth Path 2).
|
||||
///
|
||||
/// The relay forwards this opaquely — it cannot forge, modify, or replay it.
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct SignedRequest {
|
||||
pub payload: SignedRequestPayload,
|
||||
/// The client's ed25519 public key.
|
||||
pub public_key: NodeId,
|
||||
/// ed25519 signature over the canonical serialization of `payload`.
|
||||
pub signature: Signature,
|
||||
}
|
||||
|
||||
// ─── AccessControlList ──────────────────────────────────────────────────────
|
||||
|
||||
/// The datastore's access control list.
|
||||
///
|
||||
/// Persisted as `acl.json` alongside the datastore's `storage_path`.
|
||||
/// The owner always has implicit full access.
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct AccessControlList {
|
||||
/// The datastore owner's public key — always has full access.
|
||||
pub owner: NodeId,
|
||||
/// Explicitly authorized client keys.
|
||||
pub authorized_keys: HashSet<NodeId>,
|
||||
}
|
||||
|
||||
// ─── AuthzResult ────────────────────────────────────────────────────────────
|
||||
|
||||
/// The outcome of an authorization check.
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub enum AuthzResult {
|
||||
Allowed,
|
||||
Denied(DeniedReason),
|
||||
}
|
||||
|
||||
/// Why a request was denied.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub enum DeniedReason {
|
||||
/// The key is not in the ACL.
|
||||
NotAuthorized,
|
||||
/// The ed25519 signature is invalid.
|
||||
InvalidSignature,
|
||||
/// The request timestamp is outside the ±300s window.
|
||||
RequestExpired,
|
||||
/// The nonce has already been seen within the time window.
|
||||
ReplayDetected,
|
||||
}
|
||||
|
||||
// ─── AuthzEngine ────────────────────────────────────────────────────────────
|
||||
|
||||
/// Authorization engine — checks requests against the ACL and replay state.
|
||||
///
|
||||
/// Sits at the edge of the actor system (Auth Gate) and decides whether
|
||||
/// to accept or reject external requests before they reach the actors.
|
||||
#[derive(Debug)]
|
||||
pub struct AuthzEngine {
|
||||
pub acl: AccessControlList,
|
||||
// Nonce tracking and other runtime state will be added during implementation.
|
||||
}
|
||||
|
||||
impl AuthzEngine {
|
||||
/// Check whether a `NodeId` is authorized (connection-level, Auth Path 1).
|
||||
pub fn check_node(&self, _node_id: &NodeId) -> AuthzResult {
|
||||
todo!()
|
||||
}
|
||||
|
||||
/// Verify and authorize a signed request (Auth Path 2).
|
||||
pub fn check_signed_request(&self, _request: &SignedRequest) -> AuthzResult {
|
||||
todo!()
|
||||
}
|
||||
|
||||
/// Grant access to a `NodeId`. Owner-only operation.
|
||||
pub fn grant(&mut self, _requester: &NodeId, _key: NodeId) -> Result<(), DeniedReason> {
|
||||
todo!()
|
||||
}
|
||||
|
||||
/// Revoke access from a `NodeId`. Owner-only operation.
|
||||
pub fn revoke(&mut self, _requester: &NodeId, _key: NodeId) -> Result<(), DeniedReason> {
|
||||
todo!()
|
||||
}
|
||||
}
|
||||
|
|
@ -1,6 +1,5 @@
|
|||
use ed25519_dalek::{Signer, Verifier};
|
||||
|
||||
use crate::auth::{SignedRequest, SignedRequestPayload};
|
||||
use crate::types::{DirectoryEntry, DirectoryEntryPayload, NodeId, Signature};
|
||||
|
||||
// ─── Keypair ────────────────────────────────────────────────────────────────
|
||||
|
|
@ -42,13 +41,6 @@ impl Keypair {
|
|||
Signature(sig.to_bytes())
|
||||
}
|
||||
|
||||
/// Sign a request payload, returning a complete `SignedRequest` envelope.
|
||||
///
|
||||
/// Used by clients for Auth Path 2 (browser relay).
|
||||
pub fn sign_request(&self, _payload: &SignedRequestPayload) -> SignedRequest {
|
||||
todo!()
|
||||
}
|
||||
|
||||
/// Sign a directory entry payload, returning a complete `DirectoryEntry`.
|
||||
pub fn sign_directory_entry(
|
||||
&self,
|
||||
|
|
@ -82,14 +74,6 @@ pub fn verify(node_id: &NodeId, msg: &[u8], sig: &Signature) -> bool {
|
|||
vk.verify(msg, &signature).is_ok()
|
||||
}
|
||||
|
||||
/// Verify a `SignedRequest`'s signature against its embedded `public_key`.
|
||||
///
|
||||
/// Checks only signature validity — does NOT check timestamp, nonce, or ACL.
|
||||
/// Use `AuthzEngine::check_signed_request` for full verification.
|
||||
pub fn verify_signed_request(_request: &SignedRequest) -> bool {
|
||||
todo!()
|
||||
}
|
||||
|
||||
/// Verify a `DirectoryEntry`'s signature against its embedded `node_id`.
|
||||
pub fn verify_directory_entry(entry: &DirectoryEntry) -> bool {
|
||||
let payload = entry.payload();
|
||||
|
|
|
|||
|
|
@ -14,4 +14,3 @@ pub mod snapshot;
|
|||
pub mod driver;
|
||||
#[cfg(feature = "iroh")]
|
||||
pub mod iroh_driver;
|
||||
pub mod auth;
|
||||
|
|
|
|||
8
crates/shared-types/Cargo.toml
Normal file
8
crates/shared-types/Cargo.toml
Normal file
|
|
@ -0,0 +1,8 @@
|
|||
[package]
|
||||
name = "shared-types"
|
||||
version = "0.1.0"
|
||||
edition = "2024"
|
||||
|
||||
[dependencies]
|
||||
serde = { version = "1", features = ["derive"] }
|
||||
blake3 = "1"
|
||||
106
crates/shared-types/src/lib.rs
Normal file
106
crates/shared-types/src/lib.rs
Normal file
|
|
@ -0,0 +1,106 @@
|
|||
//! Shared types used across the swactor crate ecosystem.
|
||||
//!
|
||||
//! Contains `ContentHash` — the blake3-based content address used by
|
||||
//! both the datastore and distribution layers.
|
||||
|
||||
use std::fmt;
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
// ─── ContentHash ────────────────────────────────────────────────────────────
|
||||
|
||||
/// A blake3 content hash (32 bytes).
|
||||
///
|
||||
/// The primary identifier for blobs and the DHT key. XOR distance for DHT
|
||||
/// routing, compact Debug/Display for logging.
|
||||
#[derive(Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
|
||||
pub struct ContentHash(pub [u8; 32]);
|
||||
|
||||
impl ContentHash {
|
||||
/// Compute the blake3 hash of the given data.
|
||||
pub fn of(data: &[u8]) -> Self {
|
||||
let hash = blake3::hash(data);
|
||||
ContentHash(*hash.as_bytes())
|
||||
}
|
||||
|
||||
/// XOR distance between two content hashes (Kademlia metric).
|
||||
pub fn xor_distance(&self, other: &ContentHash) -> [u8; 32] {
|
||||
let mut out = [0u8; 32];
|
||||
for i in 0..32 {
|
||||
out[i] = self.0[i] ^ other.0[i];
|
||||
}
|
||||
out
|
||||
}
|
||||
|
||||
/// Number of leading zero bits in the XOR distance to `other`.
|
||||
/// Returns 0..=256. Used to select the k-bucket index in the metadata DHT.
|
||||
pub fn xor_leading_zeros(&self, other: &ContentHash) -> u32 {
|
||||
let dist = self.xor_distance(other);
|
||||
let mut zeros = 0u32;
|
||||
for byte in dist {
|
||||
if byte == 0 {
|
||||
zeros += 8;
|
||||
} else {
|
||||
zeros += byte.leading_zeros();
|
||||
break;
|
||||
}
|
||||
}
|
||||
zeros
|
||||
}
|
||||
|
||||
/// Parse a 64-character hex string into a ContentHash.
|
||||
/// Returns `None` if the string is not exactly 64 hex characters.
|
||||
pub fn from_hex(hex: &str) -> Option<Self> {
|
||||
if hex.len() != 64 {
|
||||
return None;
|
||||
}
|
||||
let mut bytes = [0u8; 32];
|
||||
for (i, chunk) in hex.as_bytes().chunks(2).enumerate() {
|
||||
let hi = hex_digit(chunk[0])?;
|
||||
let lo = hex_digit(chunk[1])?;
|
||||
bytes[i] = (hi << 4) | lo;
|
||||
}
|
||||
Some(ContentHash(bytes))
|
||||
}
|
||||
|
||||
/// Encode as lowercase hex string.
|
||||
pub fn to_hex(&self) -> String {
|
||||
let mut s = String::with_capacity(64);
|
||||
for b in &self.0 {
|
||||
use fmt::Write;
|
||||
write!(s, "{:02x}", b).unwrap();
|
||||
}
|
||||
s
|
||||
}
|
||||
|
||||
/// The zero hash (all zeroes). Used as a sentinel.
|
||||
pub const ZERO: ContentHash = ContentHash([0u8; 32]);
|
||||
}
|
||||
|
||||
impl fmt::Debug for ContentHash {
|
||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
write!(f, "Hash(")?;
|
||||
for b in &self.0[..4] {
|
||||
write!(f, "{:02x}", b)?;
|
||||
}
|
||||
write!(f, "\u{2026})")
|
||||
}
|
||||
}
|
||||
|
||||
impl fmt::Display for ContentHash {
|
||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
for b in &self.0[..8] {
|
||||
write!(f, "{:02x}", b)?;
|
||||
}
|
||||
write!(f, "\u{2026}")
|
||||
}
|
||||
}
|
||||
|
||||
fn hex_digit(b: u8) -> Option<u8> {
|
||||
match b {
|
||||
b'0'..=b'9' => Some(b - b'0'),
|
||||
b'a'..=b'f' => Some(b - b'a' + 10),
|
||||
b'A'..=b'F' => Some(b - b'A' + 10),
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
Loading…
Reference in a new issue