From f2411a7757f1900b3366986cdfb27926dc540121 Mon Sep 17 00:00:00 2001 From: OceanLi Date: Thu, 10 Sep 2026 17:59:04 -0400 Subject: [PATCH 1/2] feat(stream): bound the batch member and observation lists stream.batch refused an empty ops list and named the writer it would take for a batch that writes nothing, and did not refuse a large one. In atomic mode the whole list runs inside one BEGIN IMMEDIATE, so the statement count inside a single writer hold was whatever the caller sent, ceilinged only by the 8 MiB frame. The verb now admits at most 1000 members and at most 100 observed entries, refusing with invalid_input that names both the cap and the count sent. The check runs at admission, before any member is parsed into an action, any note plan is prepared, or a writer is requested, and it applies to the list in both modes so one input is refused the same way whatever the mode. Six acceptance arms: a list at the cap commits, a list at the cap with the observation list at its own cap commits, one member over refuses naming 1001 and 1000 with the stream head unchanged as the control that the refusal preceded the writer, one observation over refuses on its own number, per-member mode refuses the same list the same way, and a list of individually-invalid members still refuses on the cap, which is what proves the check runs before members are interpreted. Help names both caps. The ignored cap_measure_writer_hold arm measures the hold with the store's own open-transaction registry rather than wall time around the dispatch, so preparation before the writer is not counted as hold. Implements ADR-174 Amendment 7. --- crates/khive-pack-kg/src/handler_defs.rs | 2 +- crates/khive-pack-kg/src/handlers/stream.rs | 26 +++ .../src/handlers/stream_batch_cap_tests.rs | 203 ++++++++++++++++++ .../src/handlers/stream_tests.rs | 3 + 4 files changed, 233 insertions(+), 1 deletion(-) create mode 100644 crates/khive-pack-kg/src/handlers/stream_batch_cap_tests.rs diff --git a/crates/khive-pack-kg/src/handler_defs.rs b/crates/khive-pack-kg/src/handler_defs.rs index d50052370..566ee338d 100644 --- a/crates/khive-pack-kg/src/handler_defs.rs +++ b/crates/khive-pack-kg/src/handler_defs.rs @@ -56,7 +56,7 @@ pub(crate) static KG_HANDLERS: [HandlerDef; 24] = [ }, HandlerDef { name: "stream.batch", - description: "Run append and keyed write members in one request. Atomic mode (the default whenever fence is present) runs the list in one writer transaction: appends to one stream take consecutive numbers, each expected_seq is checked immediately before its own append inside that transaction, and any member refusal refuses the whole batch with details.member set and nothing written. Per-member mode runs each member in its own transaction in list order: a refusal is returned as that member's value with domain_disposition not_committed and its siblings stand; numbers on one stream increase with list position but another writer's append may fall between them. Reads are not members. Result {results: [...], committed: true} in list order.", + description: "Run append and keyed write members in one request. Atomic mode (the default whenever fence is present) runs the list in one writer transaction: appends to one stream take consecutive numbers, each expected_seq is checked immediately before its own append inside that transaction, and any member refusal refuses the whole batch with details.member set and nothing written. Per-member mode runs each member in its own transaction in list order: a refusal is returned as that member's value with domain_disposition not_committed and its siblings stand; numbers on one stream increase with list position but another writer's append may fall between them. Reads are not members. A call sending more than 1000 members, or more than 100 observed entries, is refused with invalid_input naming the cap and the count sent, before any member is parsed or any writer is requested; both caps apply in both modes. Result {results: [...], committed: true} in list order.", visibility: Visibility::Verb, category: VerbCategory::Commissive, params: &[ diff --git a/crates/khive-pack-kg/src/handlers/stream.rs b/crates/khive-pack-kg/src/handlers/stream.rs index 21937d3aa..999016343 100644 --- a/crates/khive-pack-kg/src/handlers/stream.rs +++ b/crates/khive-pack-kg/src/handlers/stream.rs @@ -13,6 +13,16 @@ use uuid::Uuid; use super::common::{canonical_note_kind, deser}; use crate::KgPack; +/// The most members one `stream.batch` call admits (ADR-174 Amendment 7). In +/// atomic mode the whole list runs inside one writer transaction, so this is +/// the bound on how long one caller may hold the single writer. +pub(crate) const MAX_BATCH_MEMBERS: usize = 1000; + +/// The most `observed` entries one `stream.batch` call admits (ADR-174 +/// Amendment 7). Each entry is a read the batch performs while holding the +/// writer only to decide whether to proceed. +pub(crate) const MAX_BATCH_OBSERVED: usize = 100; + #[derive(Deserialize)] #[serde(deny_unknown_fields)] struct AppendParams { @@ -298,6 +308,22 @@ impl KgPack { .into(), )); } + // Admission bounds, before any member is parsed into an action, any + // note plan is prepared, or a writer is requested (ADR-174 A7.3). + if p.ops.len() > MAX_BATCH_MEMBERS { + return Err(RuntimeError::InvalidInput(format!( + "stream.batch admits at most {MAX_BATCH_MEMBERS} members; this call sent {}: the whole list runs inside one writer transaction", + p.ops.len() + ))); + } + if let Some(entries) = observed.as_ref().and_then(Value::as_array) { + if entries.len() > MAX_BATCH_OBSERVED { + return Err(RuntimeError::InvalidInput(format!( + "stream.batch admits at most {MAX_BATCH_OBSERVED} observed entries; this call sent {}: an observation is a read taken while holding the writer", + entries.len() + ))); + } + } let fence = fence .map(|value| batch_fence(value, registry)) .transpose()?; diff --git a/crates/khive-pack-kg/src/handlers/stream_batch_cap_tests.rs b/crates/khive-pack-kg/src/handlers/stream_batch_cap_tests.rs new file mode 100644 index 000000000..ab733b8cb --- /dev/null +++ b/crates/khive-pack-kg/src/handlers/stream_batch_cap_tests.rs @@ -0,0 +1,203 @@ +//! ADR-174 A7: the batch admits a bounded list. +use super::*; + +use crate::handlers::stream::{MAX_BATCH_MEMBERS, MAX_BATCH_OBSERVED}; + +fn appends(count: usize) -> Vec { + (0..count) + .map(|n| json!({"op":"append","stream":"cap/a","record":n})) + .collect() +} + +fn invalid_input(error: RuntimeError) -> String { + match error { + RuntimeError::InvalidInput(message) => message, + other => panic!("invalid_input, got {other:?}"), + } +} + +#[tokio::test] +async fn cap_arm1_a_list_at_the_member_cap_commits() { + let (_, reg) = surface(); + let result = reg + .dispatch( + "stream.batch", + json!({"atomic":true,"ops":appends(MAX_BATCH_MEMBERS)}), + ) + .await + .unwrap(); + assert_eq!(result["committed"], true); + assert_eq!( + result["results"].as_array().unwrap().len(), + MAX_BATCH_MEMBERS + ); + assert_eq!(result["results"][MAX_BATCH_MEMBERS - 1]["seq"], 1000); + assert_eq!(heads(®, &["cap/a"]).await, vec![1000]); +} + +#[tokio::test] +async fn cap_arm1_b_the_member_cap_and_the_observation_cap_hold_together() { + let (_, reg) = surface(); + let mut observed = Vec::with_capacity(MAX_BATCH_OBSERVED); + for n in 0..MAX_BATCH_OBSERVED { + let key = format!("cap-key-{n}"); + lease(®, &key).await; + observed.push(json!({"key":key,"kind":"head","version":1})); + } + let result = reg + .dispatch( + "stream.batch", + json!({"atomic":true,"observed":observed,"ops":appends(MAX_BATCH_MEMBERS)}), + ) + .await + .unwrap(); + assert_eq!(result["committed"], true); + assert_eq!(heads(®, &["cap/a"]).await, vec![1000]); +} + +#[tokio::test] +async fn cap_arm2_one_member_over_the_cap_refuses_naming_the_cap_and_the_count() { + let (rt, reg) = surface(); + let before = population(&rt).await; + let error = reg + .dispatch( + "stream.batch", + json!({"atomic":true,"ops":appends(MAX_BATCH_MEMBERS + 1)}), + ) + .await + .unwrap_err(); + let message = invalid_input(error); + assert!(message.contains("1001"), "{message}"); + assert!(message.contains("1000"), "{message}"); + // The refusal preceded the writer: the head never moved and no row landed. + assert_eq!(heads(®, &["cap/a"]).await, vec![0]); + assert_eq!(population(&rt).await, before); +} + +#[tokio::test] +async fn cap_arm3_one_observation_over_the_cap_refuses_on_its_own_number() { + let (rt, reg) = surface(); + let before = population(&rt).await; + let observed: Vec = (0..MAX_BATCH_OBSERVED + 1) + .map(|n| json!({"key":format!("cap-key-{n}"),"kind":"head","version":1})) + .collect(); + let error = reg + .dispatch( + "stream.batch", + json!({"atomic":true,"observed":observed,"ops":appends(3)}), + ) + .await + .unwrap_err(); + let message = invalid_input(error); + assert!(message.contains("101"), "{message}"); + assert!(message.contains("100"), "{message}"); + assert!(message.contains("observed"), "{message}"); + assert_eq!(population(&rt).await, before); +} + +#[tokio::test] +async fn cap_arm4_per_member_mode_refuses_the_same_list_the_same_way() { + let (rt, reg) = surface(); + let before = population(&rt).await; + let error = reg + .dispatch( + "stream.batch", + json!({"atomic":false,"ops":appends(MAX_BATCH_MEMBERS + 1)}), + ) + .await + .unwrap_err(); + let message = invalid_input(error); + assert!(message.contains("1001"), "{message}"); + assert!(message.contains("1000"), "{message}"); + // No member ran, so there are no per-member values to inspect and nothing + // in the store: the mode that could have absorbed the list still refuses it. + assert_eq!(heads(®, &["cap/a"]).await, vec![0]); + assert_eq!(population(&rt).await, before); +} + +#[tokio::test] +async fn cap_arm5_the_cap_is_read_before_any_member_is_interpreted() { + let (_, reg) = surface(); + // Every member here is independently invalid: no op string at all. If the + // cap were checked after members were parsed, this would refuse on member 0. + let ops: Vec = (0..MAX_BATCH_MEMBERS + 1) + .map(|n| json!({"not_an_op":n})) + .collect(); + let message = invalid_input( + reg.dispatch("stream.batch", json!({"atomic":true,"ops":ops})) + .await + .unwrap_err(), + ); + assert!(message.contains("1001"), "{message}"); + assert!(!message.contains("member 0"), "{message}"); +} + +#[tokio::test] +async fn cap_arm6_help_names_both_caps() { + let (_, reg) = surface(); + let help = reg + .dispatch("stream.batch", json!({"help":true})) + .await + .unwrap(); + let text = help["description"].as_str().unwrap(); + for required in ["1000 members", "100 observed", "invalid_input", "both modes"] { + assert!(text.contains(required), "missing {required}: {help}"); + } +} + +/// Writer hold time for one atomic batch, measured with the store's own +/// open-transaction registry rather than with wall time around the dispatch, +/// so preparation before the writer is not counted as hold. +/// +/// Ignored: it is a measurement, not an assertion. Run it as +/// `KHIVE_CAP_BENCH_MEMBERS=1000 cargo test -p khive-pack-kg --lib +/// cap_measure_writer_hold -- --ignored --nocapture`. The count above the cap +/// is measured with the cap constant raised, which is the only way to observe +/// the hold this amendment exists to prevent. +#[tokio::test] +#[ignore] +async fn cap_measure_writer_hold() { + use khive_runtime::RuntimeConfig; + use std::time::{Duration, Instant}; + + let members: usize = std::env::var("KHIVE_CAP_BENCH_MEMBERS") + .ok() + .and_then(|v| v.parse().ok()) + .unwrap_or(MAX_BATCH_MEMBERS); + let root = std::env::temp_dir().join(format!("khive-cap-bench-{}", uuid::Uuid::new_v4())); + std::fs::create_dir_all(&root).unwrap(); + let rt = KhiveRuntime::new(RuntimeConfig { + db_path: Some(root.join("khive.db")), + ..Default::default() + }) + .unwrap(); + let mut builder = VerbRegistryBuilder::new(); + builder.register(crate::KgPack::new(rt.clone())); + let reg = builder.build().unwrap(); + + let ops = appends(members); + let peak_micros = std::sync::Arc::new(std::sync::atomic::AtomicU64::new(0)); + let sampler = peak_micros.clone(); + let watcher = tokio::spawn(async move { + loop { + if let Some((_, age, _)) = khive_storage::tx_registry::oldest() { + sampler.fetch_max(age.as_micros() as u64, std::sync::atomic::Ordering::Relaxed); + } + tokio::time::sleep(Duration::from_millis(1)).await; + } + }); + + let started = Instant::now(); + let result = reg + .dispatch("stream.batch", json!({"atomic":true,"ops":ops})) + .await; + let dispatch = started.elapsed(); + watcher.abort(); + let peak = Duration::from_micros(peak_micros.load(std::sync::atomic::Ordering::Relaxed)); + let committed = result.as_ref().map(|v| v["committed"].clone()); + println!( + "cap-bench members={members} dispatch={dispatch:?} peak_writer_hold={peak:?} committed={committed:?}" + ); + let _ = std::fs::remove_dir_all(&root); + result.unwrap(); +} diff --git a/crates/khive-pack-kg/src/handlers/stream_tests.rs b/crates/khive-pack-kg/src/handlers/stream_tests.rs index 1cd77f9be..0b2f2f862 100644 --- a/crates/khive-pack-kg/src/handlers/stream_tests.rs +++ b/crates/khive-pack-kg/src/handlers/stream_tests.rs @@ -1507,3 +1507,6 @@ mod expiry_tests; #[path = "stream_observed_id_tests.rs"] mod observed_id_tests; + +#[path = "stream_batch_cap_tests.rs"] +mod batch_cap_tests; From 8247c8d57c4047a99be4a7d9a5008609f0db0fe2 Mon Sep 17 00:00:00 2001 From: OceanLi Date: Thu, 10 Sep 2026 18:06:42 -0400 Subject: [PATCH 2/2] style: rustfmt the help arm's required-string list --- .../khive-pack-kg/src/handlers/stream_batch_cap_tests.rs | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/crates/khive-pack-kg/src/handlers/stream_batch_cap_tests.rs b/crates/khive-pack-kg/src/handlers/stream_batch_cap_tests.rs index ab733b8cb..f1fa86860 100644 --- a/crates/khive-pack-kg/src/handlers/stream_batch_cap_tests.rs +++ b/crates/khive-pack-kg/src/handlers/stream_batch_cap_tests.rs @@ -140,7 +140,12 @@ async fn cap_arm6_help_names_both_caps() { .await .unwrap(); let text = help["description"].as_str().unwrap(); - for required in ["1000 members", "100 observed", "invalid_input", "both modes"] { + for required in [ + "1000 members", + "100 observed", + "invalid_input", + "both modes", + ] { assert!(text.contains(required), "missing {required}: {help}"); } }