Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
89 commits
Select commit Hold shift + click to select a range
c590664
perf: share Cell publication and authorize follower append windows
forhappy Oct 7, 2026
4778f27
Fix shared publication admission and queued-frame append grants
forhappy Oct 7, 2026
d351d88
Preserve Fleet mode and evaluator identity in comparison evidence
forhappy Oct 7, 2026
120e4aa
Delta shared counters without subtracting cumulative summaries
forhappy Oct 7, 2026
d7dc0d4
Model bundle selection recovery and delayed materialization
forhappy Oct 7, 2026
0924dfd
Specify both bundle upload successor branches
forhappy Oct 7, 2026
57196d2
Make unsafe bundle branch assignments explicit
forhappy Oct 7, 2026
d30e6f1
Document verified write gap and gate serving on node liveness
forhappy Oct 7, 2026
a3787a8
Reuse canonical packs for singleton publication cohorts
forhappy Oct 7, 2026
20e21b7
Model independent Cell progress across binding closure
forhappy Oct 7, 2026
75ae439
Make independent checkpoint authority assumptions explicit
forhappy Oct 7, 2026
954cacf
Model accepted follower suffix during Cell binding drain
forhappy Oct 7, 2026
417eaac
Record shared write qualification failures and proof obligations
forhappy Oct 7, 2026
f7b2080
Add verified node bundle coverage and complete Cell drain
forhappy Oct 7, 2026
7fc0793
Reserve bundle inventory before pinning Cell authority
forhappy Oct 7, 2026
dd8928e
Document PR 67 write performance reevaluation
forhappy Oct 7, 2026
075b2cd
Use WAL NORMAL for externally durable runtime activations
forhappy Oct 7, 2026
2267656
Reuse canonical packs for bounded recovery materialization
forhappy Oct 7, 2026
ad66590
Index shared coverage and checkpoint exact root cohorts
forhappy Oct 7, 2026
6e4ade6
Retain dense authenticated histories and stream recovered materializa…
forhappy Oct 8, 2026
720bd4e
Recover selected bundle prefixes with complete follower tails
forhappy Oct 8, 2026
13ab0ec
Materialize fenced bundle recovery and atomically close its catalog
forhappy Oct 8, 2026
07f2712
Reconcile fenced provisional bundle reservations
forhappy Oct 8, 2026
4957985
Select native bundle coverage with one authority CAS
forhappy Oct 8, 2026
b6e8a7e
Report matched application TPS and latency at 4957985
forhappy Oct 8, 2026
67cde41
Distinguish sampled audit status codes from verified counts
forhappy Oct 8, 2026
b185672
Avoid node coverage CAS for sparse published roots
forhappy Oct 8, 2026
e900bfa
Report matched sparse-root write measurements and failures
forhappy Oct 8, 2026
bb2f089
Coalesce follower-proven roots before saturated preparation admission
forhappy Oct 8, 2026
ddb030c
Revert "Coalesce follower-proven roots before saturated preparation a…
forhappy Oct 8, 2026
61f0136
Record matched Fleet root-delay regression and rollback
forhappy Oct 8, 2026
d4bec27
Retain complete native captures in an ordered publication feed
forhappy Oct 8, 2026
2dfbaa2
Record current Fleet TPS and incomplete comparison evidence
forhappy Oct 8, 2026
3e43601
Release selected bundle captures through the actor publisher
forhappy Oct 8, 2026
7b986b0
Record measured selected-capture TPS and failed parity gates
forhappy Oct 8, 2026
2dc1917
Retain and join bounded node bundle publication
forhappy Oct 8, 2026
d4b2785
Record managed-producer TPS regression and failed qualification
forhappy Oct 8, 2026
5193c7c
Join root coverage overtaken by native bundle selection
forhappy Oct 8, 2026
e40ecd6
Bound full-compaction future stack usage on the minimum Rust version
forhappy Oct 8, 2026
df922e7
Record measured recovery fix without claiming write throughput gain
forhappy Oct 8, 2026
9d4e632
Verify fresh bundle cohorts from one bounded origin read
forhappy Oct 8, 2026
ec127ba
Record fresh bundle read optimization and unqualified write measurements
forhappy Oct 8, 2026
6d62d41
Retire verified captures and materialize selected roots asynchronously
forhappy Oct 8, 2026
a0be6e9
Record asynchronous-root measurements and availability regression
forhappy Oct 8, 2026
4a8f55c
Preserve selected receipts across confirmed asynchronous checkpoints
forhappy Oct 9, 2026
a52fe28
Record failed checkpoint-continuity performance qualification
forhappy Oct 9, 2026
524635f
docs(perf): record honest release repeat and selection deadline failures
forhappy Oct 9, 2026
6c909a6
fix(runtime): keep the publisher free while bundle selection waits
forhappy Oct 9, 2026
d852525
fix(runtime): preserve checked gate IDs across Rust versions
forhappy Oct 9, 2026
11843f6
docs(perf): record selection-readiness measurements and merge gates
forhappy Oct 9, 2026
399e908
perf(runtime): group fresh historical bundle verification
forhappy Oct 9, 2026
2257384
Revert "perf(runtime): group fresh historical bundle verification"
forhappy Oct 9, 2026
433e303
docs(perf): record withdrawn historical-read experiment
forhappy Oct 9, 2026
addc1ce
Keep bundle receipt admission progressing through checkpoint pressure
forhappy Oct 9, 2026
6dfde17
Record receipt-pressure verification and qualify architecture compari…
forhappy Oct 9, 2026
607eca9
Preserve durable command replays during publication pressure
forhappy Oct 9, 2026
c45693a
Retain admitted mutations after a pressure outcome probe
forhappy Oct 9, 2026
9f2a39d
Record durable-replay verification and publication admission gap
forhappy Oct 9, 2026
8bde901
Observe native submission waits before follower issuance
forhappy Oct 9, 2026
e7b9396
Record publication wait behind global issuance lock
forhappy Oct 9, 2026
90f099a
Coalesce fresh historical ranges within original memory admission
forhappy Oct 9, 2026
3f45936
Record bounded history measurements and revised node target
forhappy Oct 9, 2026
29b9415
Overlap bounded fresh Cell base verification
forhappy Oct 9, 2026
cf49ef5
Deploy matching Cell population in performance fixtures
forhappy Oct 9, 2026
60c3093
Record matched 2000-Cell write and read diagnostics
forhappy Oct 9, 2026
d12eac8
Bound Cell closure callbacks to preserve node lease renewal
forhappy Oct 9, 2026
6102252
Record closure recovery measurements and preserve driver errors
forhappy Oct 9, 2026
1502cff
Overlap bounded authenticated catalog and history reads
forhappy Oct 9, 2026
98a368a
Format the catalog row decoder signature
forhappy Oct 9, 2026
7d52e88
Record catalog overlap measurements and remaining bottlenecks
forhappy Oct 9, 2026
30cd960
Reduce repeated catalog scans during bundle encoding
forhappy Oct 9, 2026
b048409
Record encoder cost measurements and remaining parity gaps
forhappy Oct 9, 2026
4a5b001
Coalesce authenticated catalog and history range reads
forhappy Oct 9, 2026
c64d3ad
Record metadata window costs and failed parity measurements
forhappy Oct 9, 2026
10370d2
docs: explain measured write publication bottleneck
forhappy Oct 9, 2026
2b6db11
Keep bundle base verification slots busy
forhappy Oct 9, 2026
fa79e54
Revert "Keep bundle base verification slots busy"
forhappy Oct 9, 2026
9b79fd2
Record failed base-slot trial and fresh write comparison
forhappy Oct 9, 2026
01de6e1
Verify publication pressure against retained write telemetry
forhappy Oct 9, 2026
c384039
Correct pinned celld bundle flush source anchor
forhappy Oct 9, 2026
6e4dc59
Decouple native admission and pipeline ordered follower appends
forhappy Oct 9, 2026
a42d344
Reuse freshly matched bundle metadata during selection
forhappy Oct 9, 2026
3cad11e
Preserve a short group commit window in the native pipeline
forhappy Oct 9, 2026
d5cec48
Credit completed follower rounds while assembling the next batch
forhappy Oct 9, 2026
3ed5aaf
Keep the crash child marker independent of the test harness line
forhappy Oct 9, 2026
6fa5278
Record native pipeline measurements and remaining qualification failures
forhappy Oct 9, 2026
f858ed3
Expose original fence causes in the SQL diagnostic service
forhappy Oct 9, 2026
cf4785c
Reuse completed command admission for publication memory
forhappy Oct 9, 2026
1f3d78d
Record publication admission recovery and write regressions
forhappy Oct 9, 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
7 changes: 7 additions & 0 deletions .github/workflows/coordination-model.yml
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,10 @@ on:
- "crates/cellule-runtime/model/**"
- "crates/cellule-runtime/src/coordination/mod.rs"
- "crates/cellule-runtime/src/coordination/**"
- "crates/cellule-runtime/src/control/**"
- "crates/cellule-runtime/src/follower/**"
- "crates/cellule-runtime/src/node/**"
- "crates/cellule-runtime/src/publication/**"
- ".github/workflows/coordination-model.yml"
schedule:
- cron: "17 3 * * 1"
Expand All @@ -32,6 +36,9 @@ jobs:
chmod +x crates/cellule-runtime/model/check.sh
crates/cellule-runtime/model/check.sh fast
crates/cellule-runtime/model/check.sh negative
crates/cellule-runtime/model/check.sh write-proofs
crates/cellule-runtime/model/check.sh bundle-coverage
crates/cellule-runtime/model/check.sh binding-drain

broad:
if: github.event_name != 'pull_request'
Expand Down
1 change: 1 addition & 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 crates/cellule-axum/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ tempfile.workspace = true
tokio = { workspace = true, features = ["macros", "net", "rt-multi-thread", "signal"] }
tokio-util = { workspace = true, features = ["rt"] }
tower = { version = "0.5", features = ["util"] }
tracing-subscriber = "0.3"

[[example]]
name = "typed-api-service"
Expand Down
4 changes: 4 additions & 0 deletions crates/cellule-axum/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -120,6 +120,10 @@ For REST routes `POST /orders` and `GET /orders/{id}`, run the
cargo run -p cellule-axum --example sql --locked
```

The SQL service writes runtime warnings, including the original cause when a
command fences its Cell, to stderr. Logging stays in the embedding application;
HTTP clients receive the same outcome-aware error envelope.

## Write a manual handler

Add the adapter alongside Axum and your Cellule application crates:
Expand Down
140 changes: 140 additions & 0 deletions crates/cellule-axum/examples/fleet/authority.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,9 @@
//! One serialized directory view for heartbeats and durability CAS callbacks.
use super::*;
use cellule_runtime::node::durability::NodeLogAuthority;
use cellule_runtime::node::durability::{
BundleCheckpoint, NodeBundleAuthority, NodeBundlePublicationAuthority,
};
use cellule_runtime::node::log::{NodeLogRetirementObservation, NodeLogRotationBarrier};
use cellule_runtime::node::{
NodeAdvertisement, NodeCapacity, NodeFailureDomain, VersionedNodeAdvertisement,
Expand Down Expand Up @@ -151,6 +154,10 @@ impl Authority {
.directory
.recruit_log(&state.observed, 1, 16 << 20, 3, clock()?)
.await?;
state.observed = self
.directory
.initialize_bundle_lane(&state.observed, 1, clock()?)
.await?;
let members = state
.observed
.advertisement()
Expand Down Expand Up @@ -199,6 +206,139 @@ impl Authority {
}
}

impl NodeBundleAuthority for Authority {
fn bind<'a>(
&'a self,
authority: &'a cellule_runtime::control::authority::CellAuthority,
observed: &'a cellule_runtime::control::authority::VersionedControl,
) -> BoxFuture<'a, Result<cellule_runtime::control::authority::VersionedControl>> {
Box::pin(async move {
let mut state = self.state.lock().await;
let current = self.current(&state, 1).await?;
let (next, pinned) = self
.directory
.bind_bundle_cell(&current, authority, observed, clock()?)
.await?;
state.observed = next;
self.lease.check()?;
Ok(pinned)
})
}

fn close<'a>(
&'a self,
authority: &'a cellule_runtime::control::authority::CellAuthority,
observed: &'a cellule_runtime::control::authority::VersionedControl,
issued: cellule_runtime::node::log::CellIssuedRange,
) -> BoxFuture<'a, Result<()>> {
Box::pin(async move {
// Runtime joins the complete issued producer prefix before entering
// this mutex, including captures whose Fleet ACK preceded selection.
let mut state = self.state.lock().await;
let mut current = self.current(&state, issued.log_epoch()).await?;
let proof = self
.directory
.load_bundle_coverage(authority, observed, cellule_ltx::Limits::default())
.await?;
current = self
.directory
.checkpoint_bundle_cell(
&current,
authority,
&proof,
cellule_ltx::Limits::default(),
clock()?,
)
.await?;
state.observed = current;
state.observed = self
.directory
.begin_bundle_close(&state.observed, proof.binding(), issued, clock()?)
.await?;
state.observed = self
.directory
.finish_bundle_close(&state.observed, proof.binding(), issued, clock()?)
.await?;
self.lease.check()
})
}
}

impl NodeBundlePublicationAuthority for Authority {
fn select<'a>(
&'a self,
captures: &'a [cellule_runtime::node::log_shipper::AssignedCapture],
lease: &'a NodeLeaseGuard,
) -> BoxFuture<'a, Result<Vec<cellule_runtime::node::bundle::BundleCoverageProof>>> {
Box::pin(async move {
let mut state = self.state.lock().await;
let current = self.current(&state, 1).await?;
let frames = captures
.iter()
.flat_map(|c| c.frames().iter().cloned())
.collect::<Vec<_>>();
let assignments = captures.iter().map(|c| c.assignment()).collect::<Vec<_>>();
let prepared = self
.directory
.prepare_node_bundle(&current, &frames, &assignments, clock()?)
.await?;
let (next, proofs) = self
.directory
.select_node_bundle(
&current,
&prepared,
lease,
cellule_ltx::Limits::default(),
clock()?,
)
.await?;
state.observed = next;
self.lease.check()?;
Ok(proofs)
})
}

fn checkpoint<'a>(&'a self, checkpoints: &'a [BundleCheckpoint]) -> BoxFuture<'a, Result<()>> {
Box::pin(async move {
let mut state = self.state.lock().await;
let current = self.current(&state, 1).await?;
let mut ready = Vec::new();
for checkpoint in checkpoints {
let root = checkpoint.root();
let observed = checkpoint
.authority()
.load(cellule_runtime::CellId::from_bytes(root.cell))
.await?
.ok_or(Error::Fenced)?;
let actual = observed.value().ltx_root().ok_or(Error::Fenced)?;
if observed.value().bundle_binding == Some(checkpoint.proof().binding())
&& actual.cell == root.cell
&& actual.incarnation == root.incarnation
&& actual.commit_sequence > root.commit_sequence
{
// The later root's original notification is retained by its
// joined publisher. This observation releases no locators.
continue;
}
ready.push((checkpoint.authority(), checkpoint.proof()));
}
state.observed = if ready.is_empty() {
current
} else {
self.directory
.checkpoint_bundle_cells(
&current,
&ready,
cellule_ltx::Limits::default(),
clock()?,
)
.await?
};
self.lease.check()
})
}
}

fn capacity(follower: Option<&cellule_runtime::follower::FollowerStore>) -> NodeCapacity {
// Memory/disk hints are fixture admission ceilings, not measured physical
// headroom. Follower retention and free budget come from the canonical store.
Expand Down
4 changes: 3 additions & 1 deletion crates/cellule-axum/examples/fleet/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -136,7 +136,9 @@ impl Config {
cellule_ltx::Limits::default(),
runtime.telemetry_handle(),
)?;
runtime.install_node_durability(application, config.build()?)?;
let durability = config.build()?;
runtime.install_node_durability(application, durability.clone())?;
durability.start_bundle_publication(enrollment.authority.clone())?;
Ok(())
}
.await;
Expand Down
88 changes: 67 additions & 21 deletions crates/cellule-axum/examples/fleet/server.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
//! Private follower process: mTLS, fresh enrollment and canonical fsynced storage.
//! Private follower process: signed grants, mTLS and canonical fsynced storage.
use super::*;
use axum::{
Extension, Router,
Expand All @@ -9,7 +9,9 @@ use axum::{
};
use bytes::Bytes;
use cellule_peer_http::PeerTlsIdentity;
use cellule_runtime::follower::FollowerStore;
use cellule_runtime::follower::{
AppendGrantIssuer, AppendGrantPeer, FollowerStore, GrantedFollowerAppend,
};
use ed25519_dalek::VerifyingKey;
use prost::Message;
use tokio::sync::{OwnedSemaphorePermit, Semaphore};
Expand Down Expand Up @@ -41,6 +43,7 @@ pub(super) async fn serve(
cellule_ltx::Limits::default(),
cellule_ltx::DiskBudget::new(1 << 30),
)?
.with_append_grant_receiver(session)
.with_telemetry(
cellule_runtime::fleet::telemetry::CellTelemetryHandle::from_sink(metrics.clone()),
);
Expand Down Expand Up @@ -165,21 +168,30 @@ async fn handle_inner(server: Server, peer: PeerTlsIdentity, encoded: Bytes) ->
if request.member != server.member.as_bytes() {
return Err(Error::PeerAuthorization("capacity log member differs"));
}
let enrolled = server
.metrics
.enrollment(
true,
let enrolled = if matches!(request.operation, 2..=4) {
Some(
server
.directory
.peer_verifier(sender, peer.certificate(), peer.public_key(), clock()?),
.metrics
.enrollment(
true,
server.directory.peer_verifier(
sender,
peer.certificate(),
peer.public_key(),
clock()?,
),
)
.await?,
)
.await;
let enrolled = enrolled?;
} else {
None
};
// Directory I/O consumed time. Preserve the accepted horizon across wall
// clock rollback, while monotonic elapsed time still expires the request.
let now = wire::request_time(started_ms, clock()?, started_at.elapsed())?;
validate_request(&request, now, "after directory verification")?;
server.lease.check()?;
let deadline = request.deadline_ms;
let mut reply = wire::Reply {
member: server.member.as_bytes().to_vec(),
request_digest: blake3::hash(&body).as_bytes().to_vec(),
Expand All @@ -190,20 +202,21 @@ async fn handle_inner(server: Server, peer: PeerTlsIdentity, encoded: Bytes) ->
if sender != leader {
return Err(Error::PeerAuthorization("capacity append sender differs"));
}
enrolled.authorize_log_append(
server.member,
request.epoch,
request.covered_through,
clock()?,
)?;
let phase = std::time::Instant::now();
let result = server
.store
.append(
leader,
request.epoch,
request.frames.into_iter().map(Bytes::from).collect(),
request.covered_through,
.append_granted(
GrantedFollowerAppend {
peer: AppendGrantPeer {
session: sender,
certificate: peer.certificate(),
public_key: peer.public_key(),
},
log_epoch: request.epoch,
grant: Digest::try_from(request.grant.as_slice())?,
frames: request.frames.into_iter().map(Bytes::from).collect(),
},
server.lease.clone(),
)
.await;
server
Expand All @@ -213,6 +226,33 @@ async fn handle_inner(server: Server, peer: PeerTlsIdentity, encoded: Bytes) ->
reply.base_sequence = receipt.base_sequence;
reply.durable_through = receipt.durable_through;
}
5 => {
if sender != leader {
return Err(Error::PeerAuthorization("capacity grant sender differs"));
}
let grant = server
.metrics
.enrollment(
true,
server.store.open_append_grant(
AppendGrantPeer {
session: sender,
certificate: peer.certificate(),
public_key: peer.public_key(),
},
request.epoch,
request.first_sequence,
AppendGrantIssuer {
directory: &server.directory,
signing_key: server.tls.signing_key(),
lease: &server.lease,
now_ms: now,
},
),
)
.await?;
reply.grant = grant.encode();
}
3 => {
if sender != leader {
return Err(Error::PeerAuthorization("capacity retire sender differs"));
Expand Down Expand Up @@ -262,7 +302,13 @@ async fn handle_inner(server: Server, peer: PeerTlsIdentity, encoded: Bytes) ->
}
_ => return Err(Error::PeerAuthorization("capacity operation differs")),
}
drop(enrolled);
server.lease.check()?;
if wire::request_time(started_ms, clock()?, started_at.elapsed())? >= deadline {
return Err(Error::PeerAuthorization(
"capacity reply exceeded request horizon",
));
}
let reply = wire::sign(
reply.encode_to_vec(),
server.tls.signing_key(),
Expand Down
Loading