diff --git a/.github/workflows/reference-compose.yml b/.github/workflows/reference-compose.yml index 8f320166..7c973b3e 100644 --- a/.github/workflows/reference-compose.yml +++ b/.github/workflows/reference-compose.yml @@ -38,7 +38,8 @@ permissions: concurrency: group: cell-reference-compose-${{ github.event.pull_request.number || github.ref }} - cancel-in-progress: true + # Finish frozen comparisons and retain their evidence before testing a new head. + cancel-in-progress: false jobs: smoke: @@ -129,9 +130,14 @@ jobs: if-no-files-found: error retention-days: 7 - routing: + routing_measurements: runs-on: ubuntu-24.04 timeout-minutes: 90 + strategy: + # Each mode retains all four pairs. Both jobs must pass the routing gate. + fail-fast: false + matrix: + mode: [leased, object_only] services: rustfs: image: ghcr.io/rustfs/rustfs:1.0.0-glibc@sha256:bffcab0c9d647aab0055d1c69d340b202d0909966b385932d4ead1aeb7602858 @@ -173,11 +179,23 @@ jobs: run: | python3 -B -m unittest discover -s crates/cellule-peer-http/qualification -p 'test_*.py' python3 crates/cellule-peer-http/qualification/routing.py \ - --baseline "$ROUTING_BASELINE" --state "$RUNNER_TEMP/cell-routing" + --baseline "$ROUTING_BASELINE" --state "$RUNNER_TEMP/cell-routing" \ + --mode "${{ matrix.mode }}" - uses: actions/upload-artifact@ea165f8d65b6e75b540449e92b4886f43607fa02 # v4 if: always() with: - name: cell-routing-${{ github.run_id }}-${{ github.run_attempt }} + name: cell-routing-${{ matrix.mode }}-${{ github.run_id }}-${{ github.run_attempt }} path: ${{ runner.temp }}/cell-routing/evidence/ if-no-files-found: error retention-days: 7 + + routing: + if: ${{ always() }} + needs: routing_measurements + runs-on: ubuntu-24.04 + timeout-minutes: 5 + steps: + - name: Require all original routing modes + env: + MEASUREMENTS_RESULT: ${{ needs.routing_measurements.result }} + run: test "$MEASUREMENTS_RESULT" = success diff --git a/crates/cellule-ltx/docs/publication.md b/crates/cellule-ltx/docs/publication.md index 3b6efa87..1a8c95ac 100644 --- a/crates/cellule-ltx/docs/publication.md +++ b/crates/cellule-ltx/docs/publication.md @@ -22,6 +22,7 @@ sequenceDiagram | `CellReplica::prepare` | Verifies cuts and writes immutable root dependencies. | | `prepare_bundle` | Selects this Cell's exact rows from a shared bundle. | | `prepare_compaction` | Rewrites representation without changing logical state. | +| `prepare_after_compaction` | Appends to a private compaction while retaining its original authority predecessor. | | Runtime CAS | Names the authoritative owner and exact root. | | `Db::prune_captured` | Removes only the successfully published batch. | @@ -29,3 +30,10 @@ A failed CAS leaves unreachable content, never an acknowledged state. Provider retry pins and rechecks the selected capture bytes, so path replacement cannot change an in-flight proposal. The host owns request admission, deadlines, and reconciliation after ambiguous results. + +A representation-only compaction can remain private while its successor append +uploads. `prepare_after_compaction` verifies that the compaction preserves the +predecessor's position, commit sequence, Cell and incarnation. The runtime selects +the append's schema and can choose the final root with one CAS +against the original authority record. Every immutable dependency still finishes +uploading before the successor proposal is returned. diff --git a/crates/cellule-ltx/src/replica/prepare.rs b/crates/cellule-ltx/src/replica/prepare.rs index a900b3bd..02288307 100644 --- a/crates/cellule-ltx/src/replica/prepare.rs +++ b/crates/cellule-ltx/src/replica/prepare.rs @@ -7,6 +7,43 @@ use super::*; impl CellReplica { + /// Appends captured cuts to a private representation-only compaction. + /// + /// The successor retains the compaction's original predecessor for one + /// authority CAS. Both proposals remain immutable and fully uploaded; + /// this method neither publishes the compaction nor grants ownership. + pub async fn prepare_after_compaction( + &self, + compacted: &PreparedRoot, + cuts: &CaptureBatch, + commit_sequence: u64, + schema: u32, + ) -> Result { + let predecessor = compacted + .predecessor + .ok_or(LtxError::InvalidState("compaction has no predecessor"))?; + let root = compacted.root(); + if root.cell != self.cell + || root.incarnation != self.incarnation + || root.cell != predecessor.cell + || root.incarnation != predecessor.incarnation + || root.position != predecessor.position + || root.commit_sequence != predecessor.commit_sequence + { + return Err(LtxError::InvalidState( + "append requires a representation-only compaction", + )); + } + let mut successor = self + .prepare(Some(&root), cuts, commit_sequence, schema) + .await?; + // The private compaction authenticates identical logical state. The + // final root replaces that original state directly, after every new + // dependency has uploaded through the normal preparation path. + successor.predecessor = Some(predecessor); + Ok(successor) + } + /// Verifies and uploads a new immutable root without changing authority. pub async fn prepare( &self, diff --git a/crates/cellule-ltx/tests/cell/roots/compaction.rs b/crates/cellule-ltx/tests/cell/roots/compaction.rs index 4aa8e589..70cd4c9b 100644 --- a/crates/cellule-ltx/tests/cell/roots/compaction.rs +++ b/crates/cellule-ltx/tests/cell/roots/compaction.rs @@ -1,6 +1,90 @@ //! Scheduled compaction, large-frame streaming, and corruption refusal. use super::*; +use cellule_ltx::LtxError; + +#[tokio::test] +async fn private_compaction_append_retains_original_predecessor_and_exact_root() { + let directory = tempfile::tempdir().unwrap(); + let mut writer = Db::open(&directory.path().join("writer.sqlite"), Limits::default()).unwrap(); + let store = Store::new(Arc::new(InMemory::new())); + let cell = replica(store.clone(), [231; 32], [232; 16]); + writer + .transaction(|tx| { + tx.execute_batch("CREATE TABLE counter(value); INSERT INTO counter VALUES(0)") + }) + .unwrap(); + let initial_cuts = writer.capture().unwrap(); + let initial = cell.prepare(None, &initial_cuts, 1, 1).await.unwrap(); + writer + .transaction(|tx| { + tx.execute("UPDATE counter SET value=1", [])?; + Ok(()) + }) + .unwrap(); + let cuts = writer.capture().unwrap(); + let base = cell + .prepare(Some(&initial.root()), &cuts, 2, 1) + .await + .unwrap(); + let count = base.verified().segment_count(); + let compacted = cell + .prepare_compaction(&base.root(), 0..count, 1, directory.path()) + .await + .unwrap(); + writer + .transaction(|tx| { + tx.execute("UPDATE counter SET value=2", [])?; + Ok(()) + }) + .unwrap(); + let cuts = writer.capture().unwrap(); + for invalid in [&initial, &base] { + assert!(matches!( + cell.prepare_after_compaction(invalid, &cuts, 3, 1).await, + Err(LtxError::InvalidState(_)) + )); + } + let migrated = cell + .prepare_after_compaction(&compacted, &cuts, 3, 2) + .await + .unwrap(); + assert_eq!(migrated.verified().schema(), 2); + assert_eq!(migrated.predecessor(), Some(base.root())); + let foreign = replica(store, [230; 32], [232; 16]); + assert!(matches!( + foreign + .prepare_after_compaction(&compacted, &cuts, 3, 1) + .await, + Err(LtxError::InvalidState(_)) + )); + + let standard = cell + .prepare(Some(&compacted.root()), &cuts, 3, 1) + .await + .unwrap(); + let combined = cell + .prepare_after_compaction(&compacted, &cuts, 3, 1) + .await + .unwrap(); + assert_eq!(combined.predecessor(), Some(base.root())); + assert_eq!(standard.predecessor(), Some(compacted.root())); + assert_eq!( + combined.root(), + standard.root(), + "publication composition must preserve the exact immutable format" + ); + let restored = directory.path().join("restored.sqlite"); + combined.verified().restore(&restored).await.unwrap(); + let connection = cellule_ltx::rusqlite::Connection::open(restored).unwrap(); + assert_eq!( + connection + .query_row("SELECT value FROM counter", [], |row| row.get::<_, i64>(0)) + .unwrap(), + 2 + ); + writer.close().unwrap(); +} #[tokio::test(flavor = "multi_thread")] async fn range_compaction_reads_and_reserves_only_the_selected_data() { diff --git a/crates/cellule-peer-http/docs/routing.md b/crates/cellule-peer-http/docs/routing.md index 90a59087..99a21b3f 100644 --- a/crates/cellule-peer-http/docs/routing.md +++ b/crates/cellule-peer-http/docs/routing.md @@ -266,6 +266,12 @@ receiver latency. Cold and fresh-client lanes still include discovery and Describe. Manual workflow runs accept `routing_baseline` and `routing_only` to repeat a specific comparison without repeating Compose scaling. +CI measures leased and object-only routing in separate isolated-provider jobs. +Each mode retains four adjacent baseline/candidate pairs, all lanes, raw samples, +physical read and hop checks, and exact recovery. The final `routing` check +requires both mode jobs to pass. `routing.py --mode leased` or +`--mode object_only` runs one complete mode; omitting `--mode` runs both. + The larger CI sample sizes and adjacent pairs retain the original 10% limits. Serial full-profile comparisons still failed calibration after increasing the sample sizes, including paced local p99 at 1.13–1.15. Adjacent windows control diff --git a/crates/cellule-peer-http/qualification/routing.py b/crates/cellule-peer-http/qualification/routing.py index a5a8686b..f8a095c0 100644 --- a/crates/cellule-peer-http/qualification/routing.py +++ b/crates/cellule-peer-http/qualification/routing.py @@ -218,15 +218,18 @@ def parse_measurement(version, mode, index, evidence): return rows -def compare(rows): +def compare(rows, modes=None): comparisons = [] failures = [] - if {row["mode"] for row in rows} != set(SELECTORS): + modes = tuple(SELECTORS) if modes is None else tuple(modes) + if not modes or not set(modes) <= set(SELECTORS): + raise RuntimeError("Invalid routing modes") + if {row["mode"] for row in rows} != set(modes): raise RuntimeError("Incomplete routing modes") safe_unleased_baseline = any(row["reads"] > 0 for row in rows if row["mode"] == "object_only" and row["version"] == "baseline" and row["lane"] == "local_query") - for mode in SELECTORS: + for mode in modes: keys = sorted({(row["lane"], row["concurrency"]) for row in rows if row["mode"] == mode}) for lane, concurrency in keys: selected = {version: [row for row in rows if row["mode"] == mode @@ -261,7 +264,10 @@ def main(): parser = argparse.ArgumentParser(description=__doc__) parser.add_argument("--baseline", required=True) parser.add_argument("--state", type=Path, required=True) + parser.add_argument("--mode", choices=SELECTORS, + help="Measure one complete mode; CI requires both mode jobs") args = parser.parse_args() + modes = tuple(SELECTORS) if args.mode is None else (args.mode,) state = args.state.resolve() state.mkdir() evidence = state / "evidence" @@ -270,7 +276,7 @@ def main(): revisions = {"baseline": output("git", "rev-parse", "--verify", f"{args.baseline}^{{commit}}"), "candidate": output("git", "rev-parse", "HEAD")} manifest = {"order": ORDER, "pairs": PAIRS, "queries": QUERIES, "commands_per_lane": COMMANDS, - "paced_bursts": PACED_BURSTS, + "paced_bursts": PACED_BURSTS, "modes": modes, "selectors": SELECTORS, "host": platform.uname()._asdict(), "test_only_transplant": [str(PEER / path) for path in ( "src/performance_tests.rs", "src/lib.rs", "src/tests.rs", "Cargo.toml")] + ["Cargo.lock (peer test dependencies only)"], @@ -286,13 +292,13 @@ def main(): write_json(evidence / "manifest.json", manifest) write_json(evidence / "manifest.json", manifest) rows = [] - for mode in SELECTORS: + for mode in modes: for pair in range(len(PAIRS)): binaries = {version: Path(manifest[version]["binary"]) for version in revisions} if any(digest(binary) != manifest[version]["binary_sha256"] for version, binary in binaries.items()): raise RuntimeError("Frozen binary changed") rows.extend(measure_pair(mode, pair, binaries, evidence)) - summary = compare(rows) + summary = compare(rows, modes) write_json(evidence / "comparison.json", summary) if summary["failures"]: raise RuntimeError(f"Performance gate failed: {summary['failures']}") diff --git a/crates/cellule-peer-http/qualification/test_routing.py b/crates/cellule-peer-http/qualification/test_routing.py index 2de53ad7..52e6eac9 100644 --- a/crates/cellule-peer-http/qualification/test_routing.py +++ b/crates/cellule-peer-http/qualification/test_routing.py @@ -32,6 +32,33 @@ def test_missing_repeat_is_rejected(self): with self.assertRaisesRegex(RuntimeError, "Incomplete repeats"): routing.compare(rows) + def test_single_mode_requires_explicit_selection(self): + rows = [row for row in measurements() if row["mode"] == "leased"] + with self.assertRaisesRegex(RuntimeError, "Incomplete routing modes"): + routing.compare(rows) + + def test_each_partition_retains_repeats_and_regression_gates(self): + for mode in routing.SELECTORS: + with self.subTest(mode=mode): + rows = [row for row in measurements() if row["mode"] == mode] + self.assertEqual(routing.compare(rows, (mode,))["failures"], []) + for metric, value in (("throughput", 899), ("p95_ms", 2.21), ("p99_ms", 3.31)): + regressed = copy.deepcopy(rows) + for row in regressed: + if row["version"] == "candidate": + row[metric] = value + self.assertEqual(routing.compare(regressed, (mode,))["failures"], + [f"{mode}/local_query/c16"]) + with self.assertRaisesRegex(RuntimeError, "Incomplete repeats"): + routing.compare(rows[:-1], (mode,)) + + def test_partition_rejects_unexpected_or_invalid_modes(self): + with self.assertRaisesRegex(RuntimeError, "Incomplete routing modes"): + routing.compare(measurements(), ("leased",)) + for modes in ((), ("unknown",)): + with self.assertRaisesRegex(RuntimeError, "Invalid routing modes"): + routing.compare(measurements(), modes) + def test_cold_forwarded_discovery_regression_fails(self): rows = measurements() for row in rows: diff --git a/crates/cellule-runtime/docs/deployment.md b/crates/cellule-runtime/docs/deployment.md index 4764e27f..e5a1e046 100644 --- a/crates/cellule-runtime/docs/deployment.md +++ b/crates/cellule-runtime/docs/deployment.md @@ -129,6 +129,10 @@ flowchart LR The receiving node resolves the exact owner from `control.json` and its signed live advertisement. It dispatches locally when it owns the Cell, otherwise it forwards once. +`CellClient::runtime_with_peer` forwards without a local dispatcher lookup when +the runtime has no active or activating Cells. The shared resource ledger proves +that miss; the destination still checks its own authority and admission. + A stale endpoint retry is allowed only when the first attempt definitely did not start. An ambiguous mutation returns evidence for `Resolve`. diff --git a/crates/cellule-runtime/docs/runtime.md b/crates/cellule-runtime/docs/runtime.md index 4f5d4392..c07d7450 100644 --- a/crates/cellule-runtime/docs/runtime.md +++ b/crates/cellule-runtime/docs/runtime.md @@ -325,6 +325,7 @@ The actor never reruns a handler after SQLite may have started it. `Resolve` rea **Bounded upload concurrency** - Root preparation also overlaps independent content-addressed uploads. The LTX body and index, changed and initial directory nodes, and root metadata use bounded concurrency under the runtime's shared I/O permits. +- When foreground compaction clears the segment-debt bound, its successor append retains the original authority predecessor. One fenced CAS selects the final root after all dependencies upload; the unchanged intermediate root stays private. Compaction cascades and quiet-period publication keep their existing bounds. - Initial directory construction retains at most eight encoded nodes awaiting upload. - The proposal remains private until every dependency upload completes, so authority cannot observe a partial root. diff --git a/crates/cellule-runtime/src/cell/actor/acquire.rs b/crates/cellule-runtime/src/cell/actor/acquire.rs index 16c13bc3..2faedc58 100644 --- a/crates/cellule-runtime/src/cell/actor/acquire.rs +++ b/crates/cellule-runtime/src/cell/actor/acquire.rs @@ -13,6 +13,17 @@ impl CellRuntime { /// the peer transport remains responsible for resolving remote authority. pub(crate) async fn has_local_owner(&self, cell: CellId) -> crate::Result { self.ensure_running()?; + if self.inner.pool.active_cells() == 0 { + // Activation reserves a Cell before enqueueing it and holds that + // charge through worker teardown. Zero therefore proves a local + // miss without waking the dispatcher. A concurrent activation can + // change the destination hint; the peer still verifies ownership. + self.ensure_running()?; + if self.inner.sender.is_closed() { + return Err(Error::RuntimeClosed); + } + return Ok(false); + } let (reply, response) = oneshot::channel(); self.inner .sender diff --git a/crates/cellule-runtime/src/cell/actor/tests.rs b/crates/cellule-runtime/src/cell/actor/tests.rs index ec064a2c..ab03ed0c 100644 --- a/crates/cellule-runtime/src/cell/actor/tests.rs +++ b/crates/cellule-runtime/src/cell/actor/tests.rs @@ -1,5 +1,107 @@ use super::{CellRuntimeStats, bounded_u32}; +#[tokio::test] +async fn empty_runtime_miss_does_not_wait_for_the_dispatcher() { + use super::*; + use futures_util::FutureExt; + + let pool = SqlWorkerPool::new(1, 1).unwrap(); + let runtime = CellRuntime::new(pool, 1 << 20, SessionId::from_bytes([62; 16])).unwrap(); + // On this current-thread executor the actor has not been polled. An empty + // runtime must forward without first waking that actor for a negative lookup. + assert!(matches!( + runtime + .has_local_owner(CellId::from_bytes([62; 32])) + .now_or_never(), + Some(Ok(false)) + )); + runtime.shutdown().await.unwrap(); +} + +#[tokio::test] +async fn empty_runtime_miss_preserves_lease_and_closed_dispatcher_errors() { + use super::*; + + let fenced = CellRuntime::new_with_replica_host_requiring_node_lease( + SqlWorkerPool::new(1, 1).unwrap(), + 1 << 20, + SessionId::from_bytes([64; 16]), + cellule_ltx::Host::default(), + ) + .unwrap(); + let cell = CellId::from_bytes([64; 32]); + assert!(matches!( + fenced.has_local_owner(cell).await, + Err(Error::Fenced) + )); + fenced.shutdown().await.unwrap(); + + let mut runtime = CellRuntime::new( + SqlWorkerPool::new(1, 1).unwrap(), + 1 << 20, + SessionId::from_bytes([65; 16]), + ) + .unwrap(); + let (sender, receiver) = mpsc::channel(1); + drop(receiver); + let original = std::mem::replace( + &mut Arc::get_mut(&mut runtime.inner).unwrap().sender, + sender, + ); + assert!(matches!( + runtime.has_local_owner(cell).await, + Err(Error::RuntimeClosed) + )); + Arc::get_mut(&mut runtime.inner).unwrap().sender = original; + runtime.shutdown().await.unwrap(); +} + +#[tokio::test] +async fn activation_reservation_prevents_the_empty_runtime_shortcut() { + use super::*; + use futures_util::FutureExt; + + let pool = SqlWorkerPool::new(1, 1).unwrap(); + let mut runtime = + CellRuntime::new(pool.clone(), 1 << 20, SessionId::from_bytes([63; 16])).unwrap(); + let cell = CellId::from_bytes([63; 32]); + let (sender, mut receiver) = mpsc::channel(1); + let original = std::mem::replace( + &mut Arc::get_mut(&mut runtime.inner).unwrap().sender, + sender, + ); + let reservation = pool.reserve_activation().unwrap(); + { + let lookup = runtime.has_local_owner(cell); + tokio::pin!(lookup); + assert!(lookup.as_mut().now_or_never().is_none()); + match receiver.try_recv().unwrap() { + Message::Lookup { + cell: requested, + require_resident, + reply, + } => { + assert_eq!(requested, cell); + assert!(!require_resident); + assert!(reply.send(None).is_ok()); + } + _ => panic!("activation in progress must consult actor admission"), + } + assert!(!lookup.await.unwrap()); + } + drop(reservation); + assert!(matches!( + runtime.has_local_owner(cell).now_or_never(), + Some(Ok(false)) + )); + Arc::get_mut(&mut runtime.inner).unwrap().sender = original; + runtime.shutdown().await.unwrap(); + assert!(matches!( + runtime.has_local_owner(cell).await, + Err(Error::RuntimeClosed) + )); +} + #[tokio::test] async fn shutdown_reports_deactivation_failure_completed_before_it_started() { assert!(matches!( diff --git a/crates/cellule-runtime/src/client/mod.rs b/crates/cellule-runtime/src/client/mod.rs index 6b929422..4b3fbad6 100644 --- a/crates/cellule-runtime/src/client/mod.rs +++ b/crates/cellule-runtime/src/client/mod.rs @@ -798,7 +798,8 @@ impl CellClient { /// A lease-fenced resident actor supplies a local route without metadata. /// Unleased local actors check fresh catalog and authority. A local miss /// delegates without those reads; the peer round trip resolves the remote - /// owner and verifies enrollment. This does not acquire an idle Cell. + /// owner and verifies enrollment. A runtime with no active or activating + /// Cells also avoids a dispatcher lookup. This does not acquire an idle Cell. #[must_use] pub fn runtime_with_peer( registry: Arc, diff --git a/crates/cellule-runtime/src/publication/mod.rs b/crates/cellule-runtime/src/publication/mod.rs index b4e75578..5cfb331b 100644 --- a/crates/cellule-runtime/src/publication/mod.rs +++ b/crates/cellule-runtime/src/publication/mod.rs @@ -48,6 +48,20 @@ pub struct CellPublisher { telemetry: crate::fleet::telemetry::CellTelemetryHandle, } +enum AppendBase { + Published(Option), + Compacted(Box), +} + +impl AppendBase { + fn published(&self) -> Option { + match self { + Self::Published(root) => *root, + Self::Compacted(prepared) => prepared.predecessor(), + } + } +} + impl CellPublisher { /// Creates an owner-bound publisher with a private compaction scratch directory. /// @@ -325,7 +339,12 @@ impl CellPublisher { return Err(Error::Control("bootstrap control already has a root")); } let result = self - .prepare_cuts(None, cuts, 0, self.observed.value().schema) + .prepare_cuts( + &AppendBase::Published(None), + cuts, + 0, + self.observed.value().schema, + ) .await; self.record_publication_cost(); let prepared = result?; @@ -372,22 +391,30 @@ impl CellPublisher { ) -> Result { let base = self.compact_before_append(cuts.segments.len()).await?; let prepared = match self - .prepare_cuts(base.as_ref(), cuts, commit_sequence, schema) + .prepare_cuts(&base, cuts, commit_sequence, schema) .await { Ok(prepared) => prepared, Err(Error::Ltx(error)) if error.is_cell_graph_limit() => { - let Some(root) = base else { + let Some(root) = base.published() else { return Err(error.into()); }; let Some(compacted) = self.force_full_compaction(&root).await? else { return Err(error.into()); }; - self.prepare_cuts(Some(&compacted), cuts, commit_sequence, schema) - .await? + self.prepare_cuts( + &AppendBase::Published(Some(compacted)), + cuts, + commit_sequence, + schema, + ) + .await? } Err(error) => return Err(error), }; + if matches!(base, AppendBase::Compacted(_)) { + self.appends_since_compaction_check = 0; + } self.note_append(&prepared); Ok(prepared) } @@ -402,12 +429,9 @@ impl CellPublisher { ); } - async fn compact_before_append( - &mut self, - incoming_segments: usize, - ) -> Result> { + async fn compact_before_append(&mut self, incoming_segments: usize) -> Result { let Some(mut base) = self.observed.value().ltx_root() else { - return Ok(None); + return Ok(AppendBase::Published(None)); }; let segment_count = match self.segment_count { Some(count) => count, @@ -422,7 +446,7 @@ impl CellPublisher { let debt_limit = COMPACTION_DEBT_SEGMENTS.min(segment_limit); let under_pressure = projected >= debt_limit; if !under_pressure { - return Ok(Some(base)); + return Ok(AppendBase::Published(Some(base))); } tracing::debug!( segments = segment_count, @@ -439,16 +463,18 @@ impl CellPublisher { base = compacted; } self.appends_since_compaction_check = 0; - return Ok(Some(base)); + return Ok(AppendBase::Published(Some(base))); }; - let next_due_ms = self.observed.value().next_due_ms; - base = self.publish_prepared(&prepared, next_due_ms).await?; let compacted_segments = prepared.verified().segment_count(); - self.segment_count = Some(compacted_segments); if compacted_segments.saturating_add(incoming_segments) < debt_limit { - self.appends_since_compaction_check = 0; - return Ok(Some(base)); + // Keep the unchanged-state proposal private. Its append can + // replace the original authority root with one fenced CAS. + // Cascades that still exceed the bound publish as before. + return Ok(AppendBase::Compacted(Box::new(prepared))); } + let next_due_ms = self.observed.value().next_due_ms; + base = self.publish_prepared(&prepared, next_due_ms).await?; + self.segment_count = Some(compacted_segments); } Err(Error::Control( "Cell compaction cascade exceeded level limit", @@ -526,7 +552,7 @@ impl CellPublisher { async fn prepare_cuts( &mut self, - base: Option<&cellule_ltx::RootRef>, + base: &AppendBase, cuts: &cellule_ltx::CaptureBatch, commit_sequence: u64, schema: u32, @@ -534,7 +560,20 @@ impl CellPublisher { let mut backoff = Backoff::default(); loop { let replica = self.replica.clone(); - let attempt = replica.prepare(base, cuts, commit_sequence, schema); + let attempt = async { + match base { + AppendBase::Published(root) => { + replica + .prepare(root.as_ref(), cuts, commit_sequence, schema) + .await + } + AppendBase::Compacted(prepared) => { + replica + .prepare_after_compaction(prepared, cuts, commit_sequence, schema) + .await + } + } + }; tokio::pin!(attempt); let result = loop { tokio::select! { diff --git a/crates/cellule-runtime/src/publication/tests.rs b/crates/cellule-runtime/src/publication/tests.rs index b9d943fb..8855b1f3 100644 --- a/crates/cellule-runtime/src/publication/tests.rs +++ b/crates/cellule-runtime/src/publication/tests.rs @@ -16,6 +16,15 @@ use crate::identity::{CellId, Digest, SessionId}; #[tokio::test] async fn quiet_compaction_publishes_exact_root_after_eight_appends() { + verify_compaction_append(1).await; +} + +#[tokio::test] +async fn schema_migration_combines_foreground_compaction_without_an_intermediate_cas() { + verify_compaction_append(2).await; +} + +async fn verify_compaction_append(schema: u32) { let directory = tempfile::tempdir().unwrap(); let mut database = Db::open(&directory.path().join("cell.sqlite"), Limits::default()).unwrap(); let cell = CellId::from_bytes([41; 32]); @@ -140,8 +149,47 @@ async fn quiet_compaction_publishes_exact_root_after_eight_appends() { }) .unwrap(); let cuts = database.capture_deferred().unwrap(); - let prepared = publisher.prepare_append(&cuts, 38, 1).await.unwrap(); - publisher.publish_prepared(&prepared, None).await.unwrap(); + let before_revision = publisher.control().value().revision; + assert!(matches!( + publisher.prepare_append(&cuts, 37, schema).await, + Err(Error::Ltx(cellule_ltx::LtxError::InvalidState(_))) + )); + assert_eq!(publisher.control().value().revision, before_revision); + assert_eq!( + publisher + .authority + .load(cell) + .await + .unwrap() + .unwrap() + .value(), + publisher.control().value(), + "a rejected successor must leave its compaction private" + ); + let prepared = publisher.prepare_append(&cuts, 38, schema).await.unwrap(); + assert_eq!(publisher.control().value().ltx_root(), Some(at_ceiling)); + assert_eq!( + publisher + .authority + .load(cell) + .await + .unwrap() + .unwrap() + .value() + .ltx_root(), + Some(at_ceiling) + ); + assert_eq!(prepared.predecessor(), Some(at_ceiling)); + if schema == 1 { + publisher.publish_prepared(&prepared, None).await.unwrap(); + } else { + publisher + .publish_migration(&prepared, None, Digest::from_bytes([46; 32]), schema) + .await + .unwrap(); + } + assert_eq!(publisher.control().value().schema, schema); + assert_eq!(publisher.control().value().revision, before_revision + 1); let forced = publisher.control().value().ltx_root().unwrap(); assert!(replica.open_root(&forced).await.unwrap().segment_count() < 32); database.close().unwrap();