Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
46 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
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 .github/workflows/rust.yml
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,8 @@ jobs:
--locked -- --ignored --exact --list | grep -Fx "$smoke_test: test"
cargo test -p cellule-app --test integration "$smoke_test" \
--locked -- --ignored --exact --nocapture
- name: Test fleet journal and public controller models
run: cargo test -p cellule-host --example fleet_operations --all-features --locked
- name: Test local LTX without replica
run: cargo test -p cellule-ltx --no-default-features --locked
- name: Test replica feature
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.

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
31 changes: 29 additions & 2 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_entities(control: Path, capacity: bool = False, follower: bool = False) -> dict:
assert not follower or capacity
stages = (3,) if capacity else STAGES
Expand All @@ -347,7 +373,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 = [], {}
windows, positions, root_barriers = [], {}, []
for nodes in stages:
if capacity:
windows.extend(verify_capacity_windows(control, positions))
Expand All @@ -357,6 +383,7 @@ def verify_entities(control: Path, capacity: bool = False, follower: bool = Fals
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))
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 @@ -441,7 +468,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_object_operations, verify_root_coverage, verify_timing_evidence, verify_window
from entities import destination, verify_capacity_windows, verify_follower_proof, 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 ObjectOperationEvidence(unittest.TestCase):
def test_missing_provider_operation_is_rejected(self):
with tempfile.TemporaryDirectory() as root:
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 @@ -19,6 +19,7 @@ use tokio::net::TcpListener;

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