Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
97 commits
Select commit Hold shift + click to select a range
9138a51
feat(fleet): add durable movement foundations and bounded reconciler
forhappy Oct 1, 2026
801bc4f
feat(fleet): run real three-node overload and controller reconstruction
forhappy Oct 1, 2026
cd77e41
fix(fleet): fence unaccepted dispatches atomically before retry
forhappy Oct 1, 2026
cddb5c9
Merge remote-tracking branch 'origin/main' into codex/fleet-operation…
forhappy Oct 1, 2026
734650d
fix(docs): compile website examples with target dependency artifacts
forhappy Oct 1, 2026
ecec878
feat(fleet): resume lost releases under a replacement controller
forhappy Oct 1, 2026
b61768f
fix: use atomic try_update across supported Rust toolchains
forhappy Oct 1, 2026
7be792b
feat(fleet): cordon maintenance nodes before settled evacuation
forhappy Oct 1, 2026
b8ec303
feat(runtime): quiesce Cell foreground work for maintenance
forhappy Oct 1, 2026
5d98eb0
feat(fleet): release quiesced maintenance Cells canonically
forhappy Oct 1, 2026
df55295
feat(fleet): plan busy maintenance from configured Cell envelopes
forhappy Oct 1, 2026
869f0e4
Merge remote-tracking branch 'origin/main' into codex/fleet-operation…
forhappy Oct 1, 2026
e7aa692
docs(fleet): record busy demand and merged-source evidence
forhappy Oct 1, 2026
5fc341b
test(fleet): cover Effect and Activity maintenance recovery
forhappy Oct 1, 2026
135492e
feat(fleet): fence boot startup with durable enrollment
forhappy Oct 1, 2026
e472441
fix(runtime): join read replica closure across retained clones
forhappy Oct 1, 2026
c67c4e5
feat(fleet): pin exact reader sources before enrollment
forhappy Oct 1, 2026
9c99d3a
feat(fleet): own durable reader enrollment and retirement
forhappy Oct 1, 2026
29780a3
feat(fleet): confirm every follower fence before maintenance close
forhappy Oct 1, 2026
6092203
feat(fleet): own requested maintenance rotation through host drain
forhappy Oct 2, 2026
cc7aa30
feat(fleet): prepare exact follower enrollment and retain retirement …
forhappy Oct 2, 2026
7375109
feat(fleet): journal managed followers before native enrollment
forhappy Oct 2, 2026
5872122
feat(fleet): bind planning to the complete durable enrollment roster
forhappy Oct 2, 2026
8cff2c3
feat(fleet): observe retained reader enrollment through paused I/O
forhappy Oct 2, 2026
489c5a9
feat(fleet): expose bounded managed follower enrollment progress
forhappy Oct 2, 2026
353904e
docs(fleet): self-contained control plane design and implementation p…
forhappy Oct 2, 2026
9faedb9
test(fleet): retain reader fixture admission and live boot headroom
forhappy Oct 2, 2026
bb08c90
feat(fleet): expose retained durability supervisor progress
forhappy Oct 2, 2026
ac739bc
feat(fleet): retain request-bound native snapshots
forhappy Oct 2, 2026
c40761e
feat(fleet): drive real count convergence with native observations
forhappy Oct 2, 2026
aa12ad9
feat(fleet): refresh running boot intent before lease renewal
forhappy Oct 2, 2026
284dd8c
test(qualification): capture canonical roots after bounded publicatio…
forhappy Oct 2, 2026
9f7d3c6
feat(fleet): expose canonical reader lifetime observations
forhappy Oct 2, 2026
4dc71df
fix(fleet): bind reader continuation to complete native states
forhappy Oct 2, 2026
d5eebff
perf(node): avoid repeated verification during canonical decoding
forhappy Oct 2, 2026
8698183
fix(readers): reconcile retained enrollment work periodically
forhappy Oct 2, 2026
0132b60
fix(readers): preserve reconciliation across lease boundaries
forhappy Oct 2, 2026
6726aae
fix(host): retain node task joins across shutdown retries
forhappy Oct 2, 2026
4cc8ed3
Merge main into fleet operations foundations
forhappy Oct 2, 2026
5d350d1
fix(host): retain complete drain attempts across caller cancellation
forhappy Oct 2, 2026
9ed469c
fix(fleet): distinguish source release refusal from unknown close
forhappy Oct 2, 2026
1bd98d8
test(runtime): qualify Cron dispatch across maintenance handoff
forhappy Oct 2, 2026
a2806db
fix(fleet): inspect and adopt actual clean-release successors
forhappy Oct 2, 2026
238b816
fix(fleet): avoid competing with oldest-first pressure eviction
forhappy Oct 2, 2026
4d493ac
fix(host): bind native closing to checked boot withdrawal
forhappy Oct 2, 2026
7541774
fix(qualification): defer evidence writes beyond arrival clocks
forhappy Oct 2, 2026
b36f763
feat(host): evacuate managed readers with checked replacements
forhappy Oct 2, 2026
d8e09d0
feat(host): collect and recheck full native fleet inventories
forhappy Oct 2, 2026
57559a5
Merge main into fleet operations and preserve qualification clocks
forhappy Oct 2, 2026
3721ee6
feat(host): traverse and recheck foreign follower authority
forhappy Oct 3, 2026
7b6df30
feat(host): verify managed follower evacuation against live replacements
forhappy Oct 3, 2026
8a43d80
feat(host): verify cross-node fleet role coverage
forhappy Oct 3, 2026
af40830
feat(host): retain role coverage in reconciler observations
forhappy Oct 3, 2026
db33914
feat(runtime): retire canonically recovered follower ensembles
forhappy Oct 3, 2026
6fba324
feat(host): publish recovered follower enrollment closure
forhappy Oct 3, 2026
a9cee96
feat(host): publish failed boot retirement with process evidence
forhappy Oct 3, 2026
92313e8
feat(host): close failed receiver reader enrollments
forhappy Oct 3, 2026
72c5d3a
feat(host): persist and revalidate reader replacement evidence
forhappy Oct 3, 2026
2a15523
feat(host): persist follower replacement evidence and relocate minion
forhappy Oct 3, 2026
8e4976e
feat(runtime): expose complete pinned recovery manifest inventory
forhappy Oct 3, 2026
de01024
feat(host): confirm original process closure before recovery
forhappy Oct 3, 2026
8659f49
feat(runtime): retain original owner history before departure
forhappy Oct 3, 2026
b0552ee
feat(runtime): verify complete bounded catalog traversals
forhappy Oct 3, 2026
e532964
feat(fleet): retain complete original failed-boot writers
forhappy Oct 3, 2026
6a9443f
feat(runtime): retain canonical acquisition metadata before admission
forhappy Oct 3, 2026
1898f05
feat(runtime): verify exact root prefixes before fleet serving
forhappy Oct 3, 2026
a2edbe5
feat(fleet): verify original sealed recovery suffix before serving
forhappy Oct 3, 2026
4f646b5
fix(publication): overlap durable lineage with native uploads
forhappy Oct 3, 2026
704fa11
fix(ltx): compose final predecessor before lineage retention
forhappy Oct 3, 2026
37e3184
feat(fleet): reload complete committed original writer inventory
forhappy Oct 3, 2026
a3de63c
feat(fleet): collect complete original boot suffix inputs
forhappy Oct 3, 2026
1e7d34d
merge: preserve exact receipt root barrier with latest main
forhappy Oct 3, 2026
c78ab9d
docs(fleet): record suffix checkpoint and remaining delivery
forhappy Oct 3, 2026
842e71e
feat(fleet): verify all original writers against native successors
forhappy Oct 3, 2026
087da5f
merge: retain native derivation and verified descriptor reuse
forhappy Oct 3, 2026
8f0aa1f
docs(fleet): document native successor collection and repair cookbook…
forhappy Oct 3, 2026
5e19875
merge: preserve original root barrier and add final drained restore p…
forhappy Oct 3, 2026
c5a572f
docs(fleet): record successor verification and remaining delivery
forhappy Oct 3, 2026
831b83a
feat(fleet): retain original successor proofs in planner observations
forhappy Oct 3, 2026
f655139
test(fleet): qualify inherited recovery and stabilize lease fixtures
forhappy Oct 3, 2026
4d27858
fix(ci): satisfy Rust 1.99 error-size and slice lints
forhappy Oct 3, 2026
fe378b0
feat(fleet): retain fresh failed-boot closure observations
forhappy Oct 3, 2026
7ebc0c1
fix(fleet): pin process confirmation to retired boot evidence
forhappy Oct 3, 2026
942ca99
feat(host): retain fresh role evacuation proofs in fleet observations
forhappy Oct 3, 2026
0e4049e
Merge main while preserving minion and Axum CI verification
forhappy Oct 3, 2026
f86dd6c
Merge main and combine bounded receipt roots with exact restore evidence
forhappy Oct 3, 2026
9037a04
docs(fleet): record qualified role observation and remaining delivery
forhappy Oct 3, 2026
b1390b2
Observe original accepted fleet executor work
forhappy Oct 4, 2026
617819a
Fence source enrollment after maintenance closure
forhappy Oct 4, 2026
8c49a81
fleet: retain original maintenance enrollments
forhappy Oct 4, 2026
0eabc9a
fleet: match every original maintenance policy obligation
forhappy Oct 4, 2026
1b65e0a
fleet: discover every native maintenance donor policy
forhappy Oct 4, 2026
2a3726b
fleet: confirm original enrollment nonexecution
forhappy Oct 4, 2026
bd381d6
fleet: join exact original reader requests
forhappy Oct 4, 2026
f0a15a6
fleet: retain exact final reader roots
forhappy Oct 4, 2026
85f731e
Merge main and preserve fleet maintenance admission
forhappy Oct 4, 2026
9110095
fix(runtime): preserve replica rotation across retries
forhappy Oct 4, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
The table of contents is too big for display.
Diff view
Diff view
  •  
  •  
  •  
19 changes: 13 additions & 6 deletions .github/workflows/rust.yml
Original file line number Diff line number Diff line change
Expand Up @@ -68,7 +68,7 @@ jobs:
AWS_DEFAULT_REGION: us-east-1
run: |
aws --endpoint-url http://127.0.0.1:9000 s3api create-bucket --bucket cellule-ci
- name: Run balanced three-process storefront smoke
- name: Run balanced three-process storefront smoke and replica retry
env:
AWS_ACCESS_KEY_ID: cellule-ci
AWS_SECRET_ACCESS_KEY: cellule-integration-fixture
Expand All @@ -78,11 +78,18 @@ jobs:
CELLULE_TEST_PREFIX: reference-smoke
CELLULE_PERF_ITERATIONS: "1"
run: |
smoke_test=process_performance::reference_balanced_three_process_fleet_end_to_end_performance
cargo test -p cellule-app --test integration "$smoke_test" \
--locked -- --ignored --exact --list | grep -Fx "$smoke_test: test"
cargo test -p cellule-app --test integration "$smoke_test" \
--locked -- --ignored --exact --nocapture
for smoke_test in \
process_performance::reference_balanced_three_process_fleet_end_to_end_performance \
process_performance::reference_replica_retry_three_process_fleet; do
cargo test -p cellule-app --test integration "$smoke_test" \
--locked -- --ignored --exact --list | grep -Fx "$smoke_test: test"
cargo test -p cellule-app --test integration "$smoke_test" \
--locked -- --ignored --exact --nocapture
done
- name: Test fleet journal and public controller models
# Each case runs concurrent native runtimes and SQL workers. Bound whole
# fixtures on the hosted runner; retain each case's deadlines and races.
run: cargo test -p cellule-host --example fleet_operations --all-features --locked -- --test-threads=2
- name: Verify Axum HTTP publication, retries and cold recovery against RustFS
env:
AWS_ACCESS_KEY_ID: cellule-ci
Expand Down
6 changes: 6 additions & 0 deletions Cargo.lock

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

1 change: 1 addition & 0 deletions cookbook/Cargo.lock

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

2 changes: 1 addition & 1 deletion cookbook/apps/file-vault/src/model.rs
Original file line number Diff line number Diff line change
Expand Up @@ -62,7 +62,7 @@ impl<'de> Deserialize<'de> for Etag {
));
}
let mut bytes = [0; 32];
for (index, pair) in value.as_bytes().chunks_exact(2).enumerate() {
for (index, pair) in value.as_bytes().as_chunks::<2>().0.iter().enumerate() {
let digit = |b: u8| {
if b.is_ascii_digit() {
b - b'0'
Expand Down
2 changes: 1 addition & 1 deletion cookbook/apps/release-pipeline/src/model.rs
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ impl std::str::FromStr for ReleaseId {
return Err(Error::Identity("release identity requires 32 hex digits"));
}
let mut bytes = [0; 16];
for (position, pair) in text.as_bytes().chunks_exact(2).enumerate() {
for (position, pair) in text.as_bytes().as_chunks::<2>().0.iter().enumerate() {
bytes[position] = digit(pair[0])? * 16 + digit(pair[1])?;
}
Self::from_bytes(bytes)
Expand Down
2 changes: 1 addition & 1 deletion cookbook/apps/settings/src/model.rs
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ impl<'de> Deserialize<'de> for Version {
));
}
let mut bytes = [0; 28];
for (index, pair) in value.as_bytes().chunks_exact(2).enumerate() {
for (index, pair) in value.as_bytes().as_chunks::<2>().0.iter().enumerate() {
let digit = |b: u8| {
if b.is_ascii_digit() {
b - b'0'
Expand Down
2 changes: 1 addition & 1 deletion cookbook/apps/tenant-workspace/src/model.rs
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@ pub(crate) fn unhex<const N: usize>(text: &str) -> Result<[u8; N], Error> {
}
};
let mut bytes = [0; N];
for (i, pair) in text.as_bytes().chunks_exact(2).enumerate() {
for (i, pair) in text.as_bytes().as_chunks::<2>().0.iter().enumerate() {
bytes[i] = digit(pair[0]) * 16 + digit(pair[1]);
}
Ok(bytes)
Expand Down
2 changes: 1 addition & 1 deletion cookbook/apps/webhook-delivery/src/ingress/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -109,7 +109,7 @@ fn delivery_key(value: &str) -> Result<[u8; 32]> {
return Err("delivery key must be 64 lowercase hex digits".into());
}
let mut key = [0; 32];
for (index, pair) in value.as_bytes().chunks_exact(2).enumerate() {
for (index, pair) in value.as_bytes().as_chunks::<2>().0.iter().enumerate() {
key[index] = u8::from_str_radix(std::str::from_utf8(pair)?, 16)?;
}
Ok(key)
Expand Down
16 changes: 16 additions & 0 deletions crates/cellule-app/docs/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -143,6 +143,9 @@ sequenceDiagram
- The host wires `CellClient::with_read_replicas`, `ReplicaReadRouter`, and an
authenticated `ReplicaPeerClient`.
- Selection and retries share one five-second deadline.
- Replica retries preserve each reader's original placement position, so a
fallback does not reset round-robin ties for the next query. Outstanding
attempts still take priority when choosing the least busy reader.
- Authors receive typed capabilities, not storage or transport handles.

## Verification map
Expand All @@ -165,6 +168,19 @@ cargo test -p cellule-app --locked
- RustFS, constrained containers, and continuous-traffic rollout require the
separate qualification environment.

The provider-backed `process_performance::reference_replica_retry_three_process_fleet`
case injects one initial replica refusal across three independent hosts. It
checks the same 12 successful reads, balanced receiver counts, automatic
recruitment/refresh, receipt readback, policy eviction and joined shutdown as the
ordinary process smoke. With an isolated RustFS bucket, prefix and explicit
credentials configured, run:

```sh
cargo test -p cellule-app --test integration --locked \
process_performance::reference_replica_retry_three_process_fleet \
-- --ignored --exact --nocapture
```

See [PERFORMANCE.md](../PERFORMANCE.md), [AGENTS.md](../AGENTS.md), and the
[framework quickstart](../../../docs/quickstart.md).

Expand Down
43 changes: 43 additions & 0 deletions crates/cellule-app/performance/2026-09-29-write-capacity.md
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,49 @@ volume, object prefix, and evidence directory. The workflow uploads raw
samples, logs, binary digests, and provider details even if a repeat fails.
Its output requires review before a capacity claim.

## Arrival evidence collection

The driver retains each scheduled outcome in its existing bounded sample
vector during the arrival window and accepted-work drain. It writes and flushes
the sample TSV after capturing the original end clocks. Synchronous evidence
storage therefore cannot delay the next scheduled arrival or extend the measured
drain. The offered rates, concurrency, ten-second arrival window, two-second
drain allowance, and classification of late arrivals remain unchanged.

PR #37's follower-proof run `37060783753`, repeat 2, recorded 58/60 successful
actions at the first uniform point. Arrivals 38 and 39 were `scheduler_late`
and were never dispatched; every dispatched action succeeded. The driver had
no CPU throttling. The previous loop synchronously flushed completed samples
before scheduling each next arrival, exposing the arrival clock to evidence
storage stalls. The traces do not identify the individual stalled syscall;
deferring these writes removes that known blocking path. This change alone
does not establish a qualified capacity result: fresh provider repeats remain
required.

## Canonical root capture after client work

Follower-proof responses can precede object publication. After joining client
work and checking every receipt ledger, the driver observes canonical roots
under one two-second deadline for the complete original serving roster. It
pins owner session and endpoint, epoch, incarnation, code and schema. Changed,
missing, unreadable or expired authority cannot pass. The barrier only reads
authority; it does not rotate epochs or force publication.

`capacity-root-barrier-3.tsv` records actual duration, complete read passes and
Linux boot-clock bounds. The duration includes clock reads, which are measured
separately for the centisecond clock comparison. `capacity-roots-3.tsv` retains
actual canonical roots and adds each Cell's minimum acknowledged sequence.
The independent verifier derives these minima from all arrival records and
retains the existing root coverage and original owner assertions. Scaling
stages emit the corresponding `entity-root-barrier-N.tsv` evidence.

This post-load observation is separate from response and arrival latency.
The ten-second arrival windows, two-second accepted-work drain allowance,
offered rates, concurrency, readback, overload and follower-proof gates remain.
Each barrier must finish within two seconds; a stalled publisher still fails.
Earlier artifacts retain their historical verifier contract and cannot provide
the new barrier evidence. Later telemetry cannot repair a stale root capture.

## First isolated object-proof result

[CI run 36649534205](https://github.com/crabbuild/cellule/actions/runs/36649534205)
Expand Down
34 changes: 30 additions & 4 deletions crates/cellule-app/qualification/entities.py
Original file line number Diff line number Diff line change
Expand Up @@ -330,6 +330,32 @@ def verify_root_coverage(roots: list[dict], positions: dict[int, list[int]],
assert int(row["root_sequence"]) >= max(positions[entity]), "published root does not cover writes"


def verify_root_barrier(control: Path, prefix: str, nodes: int, roots: list[dict],
positions: dict[int, list[int]], windows: list[dict]) -> dict:
metadata, = rows(control / f"{prefix}-root-barrier-{nodes}.tsv")
cells = nodes * CELLS_PER_NODE
assert int(metadata["nodes"]) == nodes and int(metadata["cells"]) == cells, "root barrier roster changed"
assert int(metadata["limit_us"]) == CAPACITY_DRAIN_GRACE_US, "root barrier budget changed"
elapsed_us, reads = int(metadata["elapsed_us"]), int(metadata["reads"])
assert 0 <= elapsed_us < CAPACITY_DRAIN_GRACE_US, "root barrier exceeded original budget"
assert reads >= cells and reads % cells == 0, "root barrier did not traverse the complete roster"
started, ended = int(metadata["started_boot_ms"]), int(metadata["ended_boot_ms"])
last_window = max(window["ended_boot_ms"] for window in windows if window["nodes"] == nodes)
assert started >= last_window and ended >= started, "root barrier preceded accepted client work"
# /proc/uptime has 10ms resolution. Bound the difference by that bucket
# plus the measured clock reads enclosed by the driver's elapsed timer.
clock_read_us = int(metadata["clock_read_us"])
assert 0 <= clock_read_us <= elapsed_us, "invalid root barrier clock-read duration"
assert abs((ended - started) * 1000 - elapsed_us) <= 10_000 + clock_read_us, "root barrier clocks disagree"
assert [int(row["entity"]) for row in roots] == list(range(cells))
for row in roots:
entity = int(row["entity"])
assert int(row["minimum_sequence"]) == max(positions[entity]), "root barrier omitted an acknowledged sequence"
return dict(nodes=nodes, cells=cells, limit_us=CAPACITY_DRAIN_GRACE_US,
elapsed_us=elapsed_us, reads=reads, started_boot_ms=started, ended_boot_ms=ended,
clock_read_us=clock_read_us)


def verify_follower_roots(control: Path, roots: list[dict], positions: dict[int, list[int]],
identity: dict[int, tuple], cells: int) -> dict:
# A follower proof may precede object publication. Retain the serving
Expand Down Expand Up @@ -405,7 +431,7 @@ def verify_entities(control: Path, capacity: bool = False, follower: bool = Fals
ingress = list(map(int, (control / f"{evidence_prefix}-ingress-{stage}.txt").read_text().split()))
assert len(ingress) == stage and min(ingress) > 0 and max(ingress) - min(ingress) <= 1
assert len({value[0] for value in identity.values()}) == stages[-1] * CELLS_PER_NODE, "entity targets collapsed"
windows, positions, root_recovery = [], {}, None
windows, positions, root_barriers, root_recovery = [], {}, [], None
for nodes in stages:
if capacity:
windows.extend(verify_capacity_windows(control, positions))
Expand All @@ -414,10 +440,10 @@ def verify_entities(control: Path, capacity: bool = False, follower: bool = Fals
for rate, concurrency in POINTS:
windows.append(verify_window(control, nodes, shape, rate, concurrency, len(windows), positions))
roots = rows(control / f"{evidence_prefix}-roots-{nodes}.tsv")
verify_root_coverage(roots, positions, identity, nodes * CELLS_PER_NODE)
root_barriers.append(verify_root_barrier(control, evidence_prefix, nodes, roots, positions, windows))
if follower:
root_recovery = verify_follower_roots(control, roots, positions, identity, nodes * CELLS_PER_NODE)
else:
verify_root_coverage(roots, positions, identity, nodes * CELLS_PER_NODE)
for node in range(nodes):
assert sum(window["acknowledged_writes_by_node"][node] for window in windows if window["nodes"] == nodes) > 0
resources = {}
Expand Down Expand Up @@ -503,7 +529,7 @@ def verify_entities(control: Path, capacity: bool = False, follower: bool = Fals
node["published_roots_per_second"] for node in fully_served["node_durability"].values()),
)
extra["capacity_curves"] = capacity_curves
return dict(integrity_verified=True, windows=windows, resources=resources,
return dict(integrity_verified=True, windows=windows, resources=resources, root_barriers=root_barriers,
verified_cells=len(positions), acknowledged_writes=sum(map(len, positions.values())),
raw_sha256={path.name: hashlib.sha256(path.read_bytes()).hexdigest()
for path in sorted(control.glob("*.tsv"))}, **extra)
75 changes: 74 additions & 1 deletion crates/cellule-app/qualification/test_entities.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@
import unittest
from unittest.mock import patch

from entities import destination, verify_capacity_windows, verify_follower_proof, verify_follower_roots, verify_object_operations, verify_root_coverage, verify_timing_evidence, verify_window
from entities import destination, verify_capacity_windows, verify_follower_proof, verify_follower_roots, verify_object_operations, verify_root_barrier, verify_root_coverage, verify_timing_evidence, verify_window


class EntityWindowEvidence(unittest.TestCase):
Expand Down Expand Up @@ -228,6 +228,79 @@ def test_rate_mislabeled_as_fully_served_is_rejected(self):
verify_capacity_windows(self.root, {})


class RootBarrierEvidence(unittest.TestCase):
def setUp(self):
temporary = tempfile.TemporaryDirectory()
self.addCleanup(temporary.cleanup)
self.control = Path(temporary.name)
self.path = self.control / "capacity-root-barrier-3.tsv"
self.metadata = dict(nodes="3", cells="12", limit_us="2000000", elapsed_us="10000",
reads="24", started_boot_ms="110000", ended_boot_ms="110010", clock_read_us="0")
self.roots = [dict(entity=str(entity), minimum_sequence="5") for entity in range(12)]
self.positions = {entity: [2, 3, 5] for entity in range(12)}
self.windows = [dict(nodes=3, ended_boot_ms=110000)]

def verify(self):
self.path.write_text("\t".join(self.metadata) + "\n" + "\t".join(self.metadata.values()) + "\n")
return verify_root_barrier(self.control, "capacity", 3, self.roots, self.positions, self.windows)

def test_complete_original_roster_and_budget_are_required(self):
self.assertEqual(self.verify()["reads"], 24)

def test_missing_barrier_is_rejected(self):
with self.assertRaises(FileNotFoundError):
verify_root_barrier(self.control, "capacity", 3, self.roots, self.positions, self.windows)

def test_late_and_extended_barriers_are_rejected(self):
for field, value, message in [("elapsed_us", "2000000", "exceeded original budget"),
("elapsed_us", "2000001", "exceeded original budget"),
("elapsed_us", "-1", "exceeded original budget"),
("limit_us", "2000001", "budget changed")]:
with self.subTest(field=field, value=value):
original = self.metadata[field]
self.metadata[field] = value
with self.assertRaisesRegex(AssertionError, message):
self.verify()
self.metadata[field] = original

def test_incomplete_roster_and_read_passes_are_rejected(self):
for field, value, message in [("cells", "11", "roster changed"),
("nodes", "2", "roster changed"),
("reads", "11", "complete roster"),
("reads", "13", "complete roster")]:
with self.subTest(field=field, value=value):
original = self.metadata[field]
self.metadata[field] = value
with self.assertRaisesRegex(AssertionError, message):
self.verify()
self.metadata[field] = original

def test_barrier_cannot_precede_client_work_or_forge_elapsed_time(self):
for field, value, message in [("started_boot_ms", "109999", "preceded accepted"),
("ended_boot_ms", "109999", "preceded accepted"),
("elapsed_us", "21001", "clocks disagree")]:
with self.subTest(field=field, value=value):
original = self.metadata[field]
self.metadata[field] = value
with self.assertRaisesRegex(AssertionError, message):
self.verify()
self.metadata[field] = original

def test_each_minimum_is_derived_from_all_independent_acknowledgements(self):
for minimum in ("3", "6"):
self.roots[-1]["minimum_sequence"] = minimum
with self.assertRaisesRegex(AssertionError, "omitted an acknowledged sequence"):
self.verify()

def test_clock_read_duration_is_measured_inside_the_original_budget(self):
for duration in ("-1", "10001"):
self.metadata["clock_read_us"] = duration
with self.assertRaisesRegex(AssertionError, "clock-read duration"):
self.verify()
self.metadata.update(clock_read_us="100", elapsed_us="20100")
self.assertEqual(self.verify()["clock_read_us"], 100)


class FollowerRootDrainEvidence(unittest.TestCase):
def setUp(self):
temporary = tempfile.TemporaryDirectory()
Expand Down
1 change: 1 addition & 0 deletions crates/cellule-app/tests/entities/process.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ use tokio::net::TcpListener;
mod drain;
mod driver;
mod observation;
mod root_capture;

async fn endpoints(sync: &Path, count: usize) -> Vec<SocketAddr> {
let mut endpoints = Vec::new();
Expand Down
Loading
Loading