Skip to content
26 changes: 22 additions & 4 deletions .github/workflows/reference-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
8 changes: 8 additions & 0 deletions crates/cellule-ltx/docs/publication.md
Original file line number Diff line number Diff line change
Expand Up @@ -22,10 +22,18 @@ 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. |

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.
37 changes: 37 additions & 0 deletions crates/cellule-ltx/src/replica/prepare.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<PreparedRoot> {
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,
Expand Down
84 changes: 84 additions & 0 deletions crates/cellule-ltx/tests/cell/roots/compaction.rs
Original file line number Diff line number Diff line change
@@ -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() {
Expand Down
6 changes: 6 additions & 0 deletions crates/cellule-peer-http/docs/routing.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
18 changes: 12 additions & 6 deletions crates/cellule-peer-http/qualification/routing.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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"
Expand All @@ -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)"],
Expand All @@ -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']}")
Expand Down
27 changes: 27 additions & 0 deletions crates/cellule-peer-http/qualification/test_routing.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
4 changes: 4 additions & 0 deletions crates/cellule-runtime/docs/deployment.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`.

<a id="replica-reads"></a>
Expand Down
1 change: 1 addition & 0 deletions crates/cellule-runtime/docs/runtime.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down
11 changes: 11 additions & 0 deletions crates/cellule-runtime/src/cell/actor/acquire.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<bool> {
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
Expand Down
Loading
Loading