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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 0 additions & 2 deletions ci/no_clock_and_no_json.sh
Original file line number Diff line number Diff line change
Expand Up @@ -63,8 +63,6 @@ common_allow_globs=(
# other operational controls. They remain subject to the JSON/encoding/version gates.
clock_allow_globs=(
"${common_allow_globs[@]}"
--glob '!**/api/infra/rate_limit.rs' # transport-layer DoS rate limiting (permitted)
--glob '!**/api/transport/b0x.rs' # transport-layer rate limiting (permitted)
--glob '!**/jni/ble_events.rs' # BLE event buffering / runtime wakeups
--glob '!**/deterministic_state_machine/dsm_sdk/src/sdk/bluetooth_transport.rs' # BLE retries / ACK timeouts / reconnect backoff
--glob '!**/deterministic_state_machine/dsm_sdk/src/bluetooth/pairing_orchestrator.rs' # BLE handshake freshness / retry windows
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -568,24 +568,29 @@ async fn committed_at<S: RouteSeats>(
/// (§9 route chains, rule 4).
pub async fn read_cell<S: RouteSeats>(seats: &S, cell: &RoutedCell) -> CellEvidence {
let (route, namespace, key) = (cell.route(), cell.namespace(), cell.key());
let mut values = Vec::with_capacity(route.seats().len());
for seat in route.seats() {
values.push(seats.read_values(seat, namespace, key).await);
}
let mut cycles = Vec::with_capacity(values.len());
for (seat, held) in route.seats().iter().zip(&values) {
let holds_any = held.as_ref().is_some_and(|v| !v.is_empty());
cycles.push(if holds_any {
seats.close(seat).await
} else {
None
});
}
// Every seat is asked at once, and the answers stay in route order: a
// seat that does not answer costs one timeout, not one per seat.
let values: Vec<Option<Vec<Vec<u8>>>> = futures::future::join_all(
route
.seats()
.iter()
.map(|seat| seats.read_values(seat, namespace, key)),
)
.await;
let cycles: Vec<Option<u64>> =
futures::future::join_all(route.seats().iter().zip(&values).map(
|(seat, held)| async move {
if held.as_ref().is_some_and(|v| !v.is_empty()) {
seats.close(seat).await
} else {
None
}
},
))
.await;
let members = seats.members();
if cycles.iter().any(Option::is_some) {
for member in &members {
seats.sync_mirror(member).await;
}
futures::future::join_all(members.iter().map(|member| seats.sync_mirror(member))).await;
}
let leader = route.leader();
let leader_cycle = cycles.first().copied().flatten();
Expand Down
3 changes: 3 additions & 0 deletions dsm_storage_node/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -59,4 +59,7 @@ name = "storage_node"
path = "src/main.rs"

[dev-dependencies]
# Parse the node's own sources for clock reads (tests/no_clock_reads.rs).
proc-macro2 = "1"
rcgen = "0.14.8"
syn = { version = "2", features = ["full", "visit"] }
7 changes: 5 additions & 2 deletions dsm_storage_node/src/api/objects/bytecommit.rs
Original file line number Diff line number Diff line change
Expand Up @@ -44,8 +44,11 @@ use dsm::utils::text_id;
const CYCLE_HEADER: &str = "x-cycle";
const ECHO_HEADER: &str = "x-dsm-node-id";
/// Upper bound on cycles fetched from one member in one sync, so one request
/// does bounded work; a later sync continues where this one stopped.
const MAX_SYNC_CYCLES: u64 = 1024;
/// does bounded work; a later sync continues where this one stopped, since
/// every cycle is kept as it is fetched. At this bound one sync of a member
/// far behind is 128 fetches, well inside what a client waits for one
/// request (the SDK's member client waits 30 s).
const MAX_SYNC_CYCLES: u64 = 128;

pub fn create_router(state: Arc<AppState>) -> Router<()> {
Router::new()
Expand Down
17 changes: 16 additions & 1 deletion dsm_storage_node/src/set_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,17 @@ use rustls::pki_types::pem::PemObject;
use rustls::pki_types::CertificateDer;
use rustls::{ClientConfig, RootCertStore};
use std::sync::Arc;
use std::time::Duration;

/// How long a node waits to connect to a set-mate, and for one answer from
/// it. A set-mate that accepts a connection and never answers fails the
/// fetch within these bounds instead of holding the mirror sync, and with it
/// the node's one-at-a-time sync lock, open. They are transport liveness
/// bounds, as an unreachable set-mate is: a fetch that times out stores
/// nothing and orders nothing, so no protocol fact depends on them (storage
/// spec §1 rule 4).
pub const SET_MATE_CONNECT_TIMEOUT: Duration = Duration::from_secs(5);
pub const SET_MATE_REQUEST_TIMEOUT: Duration = Duration::from_secs(10);

/// A client pinned to `set_ca_pem`, the storage set's CA certificate.
pub fn pinned_set_client(set_ca_pem: &[u8]) -> anyhow::Result<Client> {
Expand All @@ -27,7 +38,11 @@ pub fn pinned_set_client(set_ca_pem: &[u8]) -> anyhow::Result<Client> {
.with_root_certificates(root_store)
.with_no_client_auth();

Ok(Client::builder().use_preconfigured_tls(config).build()?)
Ok(Client::builder()
.use_preconfigured_tls(config)
.connect_timeout(SET_MATE_CONNECT_TIMEOUT)
.timeout(SET_MATE_REQUEST_TIMEOUT)
.build()?)
}

#[cfg(test)]
Expand Down
33 changes: 33 additions & 0 deletions dsm_storage_node/tests/bytecommit_chain.rs
Original file line number Diff line number Diff line change
Expand Up @@ -345,6 +345,39 @@ async fn a_node_in_no_set_refuses_a_mirror_sync() {
assert_eq!(sync(&app(state)).await, StatusCode::CONFLICT);
}

/// A set-mate that accepts the connection and never answers fails the sync
/// within the set client's bounds, instead of holding the sync, and the
/// node's one-at-a-time sync lock, open. The node answers BAD_GATEWAY, and
/// the next sync is answered too: the lock was released.
#[tokio::test]
async fn a_set_mate_that_never_answers_fails_the_sync_and_frees_it() {
let (silent, silent_url) = listener().await;
// Accepts every connection and answers none of them.
tokio::spawn(async move {
let mut held = Vec::new();
while let Ok((socket, _)) = silent.accept().await {
held.push(socket);
}
});
let b = app(node_state(
"bc_silent_b",
"dsm-node-b",
&["dsm-node-a", "dsm-node-b"],
&[("dsm-node-a", &silent_url)],
)
.await);
let within = dsm_storage_node::set_client::SET_MATE_REQUEST_TIMEOUT * 3;
for attempt in 0..2 {
let answered = match tokio::time::timeout(within, sync(&b)).await {
Ok(status) => status,
Err(elapsed) => {
panic!("sync {attempt} hung on a set-mate that never answers: {elapsed}")
}
};
assert_eq!(answered, StatusCode::BAD_GATEWAY, "sync {attempt}");
}
}

/// A set-mate that does not answer at its configured endpoint fails the
/// sync; the node does not report the sync as done.
#[tokio::test]
Expand Down
175 changes: 175 additions & 0 deletions dsm_storage_node/tests/no_clock_reads.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,175 @@
// SPDX-License-Identifier: Apache-2.0
//! No storage-node source reads a clock or runs a timer: ordering inside a
//! node is its arrival sequence, and a cycle closes when asked, never on a
//! schedule (storage spec §1 rule 4, §14). Every file under `src/` is parsed:
//! paths and `use` trees are read whole, so a grouped or renamed import is
//! seen, and an unparsed macro body is read token by token. A name in a
//! comment or a string never counts. Bounding how long a client waits for a
//! set-mate is a `Duration` handed to the HTTP client; it reads nothing.

mod common;

use proc_macro2::{TokenStream, TokenTree};
use std::path::{Path, PathBuf};
use syn::visit::Visit;
use syn::UseTree;

/// Names only a clock read, a timestamp or a timer needs.
const CLOCK_NAMES: [&str; 4] = ["Instant", "SystemTime", "UNIX_EPOCH", "chrono"];

#[derive(Default)]
struct Clocks {
found: Vec<String>,
}

impl Clocks {
fn path(&mut self, segments: &[String]) {
for name in segments {
if CLOCK_NAMES.contains(&name.as_str()) {
self.found.push(name.clone());
}
}
if segments
.windows(2)
.any(|pair| pair[0] == "tokio" && pair[1] == "time")
{
self.found.push("tokio::time".to_string());
}
}
}

/// Every full path a `use` tree brings into scope. A glob ends at its
/// prefix, so `use tokio::*` is the path `tokio`, which brings `time` in.
fn use_paths(tree: &UseTree, prefix: Vec<String>, out: &mut Vec<Vec<String>>) {
match tree {
UseTree::Path(step) => {
let mut next = prefix;
next.push(step.ident.to_string());
use_paths(&step.tree, next, out);
}
UseTree::Name(leaf) => {
let mut full = prefix;
full.push(leaf.ident.to_string());
out.push(full);
}
UseTree::Rename(leaf) => {
let mut full = prefix;
full.push(leaf.ident.to_string());
out.push(full);
}
UseTree::Glob(_) => {
let mut full = prefix;
if full == ["tokio"] {
full.push("time".to_string());
}
out.push(full);
}
UseTree::Group(group) => {
for item in &group.items {
use_paths(item, prefix.clone(), out);
}
}
}
}

/// A macro body's identifiers and punctuation in order, with every literal
/// as an empty token so nothing joins across it.
fn tokens(stream: TokenStream, out: &mut Vec<String>) {
for tree in stream {
match tree {
TokenTree::Group(group) => tokens(group.stream(), out),
TokenTree::Ident(ident) => out.push(ident.to_string()),
TokenTree::Punct(punct) => out.push(punct.as_char().to_string()),
TokenTree::Literal(_) => out.push(String::new()),
}
}
}

impl<'ast> Visit<'ast> for Clocks {
fn visit_path(&mut self, path: &'ast syn::Path) {
let segments: Vec<String> = path.segments.iter().map(|s| s.ident.to_string()).collect();
self.path(&segments);
syn::visit::visit_path(self, path);
}

fn visit_item_use(&mut self, item: &'ast syn::ItemUse) {
let mut paths = Vec::new();
use_paths(&item.tree, Vec::new(), &mut paths);
for path in paths {
self.path(&path);
}
}

fn visit_macro(&mut self, mac: &'ast syn::Macro) {
let mut flat = Vec::new();
tokens(mac.tokens.clone(), &mut flat);
let joined: Vec<String> = flat
.split(|t| t == ":")
.filter(|run| !run.is_empty())
.flat_map(|run| run.iter().cloned())
.collect();
self.path(&joined);
syn::visit::visit_macro(self, mac);
}
}

fn sources(dir: &Path, found: &mut Vec<PathBuf>) {
for entry in common::ok_or_panic(std::fs::read_dir(dir), "list a source directory") {
let path = common::ok_or_panic(entry, "read a source directory entry").path();
if path.is_dir() {
sources(&path, found);
} else if path.extension().is_some_and(|ext| ext == "rs") {
found.push(path);
}
}
}

fn clock_reads(source: &str, context: &str) -> Vec<String> {
let file = common::ok_or_panic(syn::parse_file(source), context);
let mut clocks = Clocks::default();
clocks.visit_file(&file);
clocks.found
}

#[test]
fn a_clock_read_is_found_and_a_comment_or_string_is_not() {
let reads =
"fn f() { let t = std::time::Instant::now(); log::info!(\"{:?}\", SystemTime::now()); }";
assert_eq!(clock_reads(reads, "parse reads"), ["Instant", "SystemTime"]);
let grouped = "use tokio::{sync::Mutex, time};";
assert_eq!(
clock_reads(grouped, "parse a grouped import"),
["tokio::time"]
);
let renamed = "use std::time::Instant as Tick; fn f() { Tick::now(); }";
assert_eq!(clock_reads(renamed, "parse a renamed import"), ["Instant"]);
let glob = "use tokio::*; fn f() { time::sleep(d); }";
assert_eq!(clock_reads(glob, "parse a glob import"), ["tokio::time"]);
let in_macro = "fn f() { tokio::select! { _ = tokio::time::sleep(d) => {} } }";
assert_eq!(clock_reads(in_macro, "parse a macro body"), ["tokio::time"]);
let quiet =
"// Instant::now()\n/// SystemTime\nfn f() -> &'static str { \"tokio::time chrono\" }";
assert_eq!(clock_reads(quiet, "parse quiet"), Vec::<String>::new());
}

#[test]
fn no_node_source_reads_a_clock() {
let root = Path::new(env!("CARGO_MANIFEST_DIR")).join("src");
let mut files = Vec::new();
sources(&root, &mut files);
for expected in ["lib.rs", "main.rs", "api/transport/b0x.rs", "set_client.rs"] {
assert!(files.contains(&root.join(expected)), "{expected} not swept");
}
let mut found = Vec::new();
for file in &files {
let source = common::ok_or_panic(std::fs::read_to_string(file), "read a node source");
for name in clock_reads(&source, &file.display().to_string()) {
found.push(format!("{}: {name}", file.display()));
}
}
assert_eq!(
found,
Vec::<String>::new(),
"node sources that read a clock"
);
}
Loading
Loading