Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
28 commits
Select commit Hold shift + click to select a range
2f4ce38
Reserve HTTP listeners throughout SSH publication fixtures
forhappy Oct 4, 2026
2888af5
Gate deployment and workspace access on packed storage format
forhappy Oct 4, 2026
38ec7eb
Select packed production registry and certify repository creation
forhappy Oct 4, 2026
699810d
Register and recover exact catalog initialization commands
forhappy Oct 4, 2026
6490293
fix: retire initialization recovery pins with shared receipts
forhappy Oct 4, 2026
50b9a2b
Preserve first preparation admission across cold restart
forhappy Oct 4, 2026
dcef9c8
Register exact custody intents before namespace admission
forhappy Oct 4, 2026
1a11162
Own registered staging and preparation commands through recovery
forhappy Oct 4, 2026
26bee4f
Fence restored custody against current durable Cell ownership
forhappy Oct 4, 2026
0a33a67
Reconstruct registered staging custody and wake fenced bound workers
forhappy Oct 4, 2026
29e785d
Retire expired custody originals through bounded fair maintenance
forhappy Oct 4, 2026
da0bde9
Recover production initialization after expired custody retirement
forhappy Oct 4, 2026
b768e2b
Automatically drain staging after authenticated custody retirement
forhappy Oct 4, 2026
4d67294
Bound publication admission and dispatch across repositories
forhappy Oct 4, 2026
c325983
Own resident publication recovery through eviction and shutdown
forhappy Oct 4, 2026
45ad384
Retain certified serving generations through owned read drain
forhappy Oct 4, 2026
e9ea1b3
Register exact serving acquisition and renewal through shared custody
forhappy Oct 4, 2026
dab8bd0
Observe current serving roots and drain exact pins before closure
forhappy Oct 4, 2026
2012f86
Link the production cutover draft and preserve release status
forhappy Oct 4, 2026
b21f9d7
Own serving generation renewal and exact physical drain
forhappy Oct 4, 2026
02307fd
Pool resident serving generations and join them before Cell release
forhappy Oct 4, 2026
1e32347
Fence serving construction against shutdown and join rejected owners
forhappy Oct 4, 2026
b89d929
Serve browser refs through owned immutable joint snapshots
forhappy Oct 4, 2026
b1bd813
Serve certified native bodies through shared owned pack caches
forhappy Oct 4, 2026
cc4a963
Serve browser objects and ancestry through certified snapshots
forhappy Oct 4, 2026
6ab78d6
Serve native fetch through owned certified history workspaces
forhappy Oct 4, 2026
6101d58
Replace retired hydration cursor with pinned native catalog inputs
forhappy Oct 4, 2026
d86f315
Retain staging worker admission through physical native work
forhappy Oct 4, 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
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/canopy-server/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ cellule-store = { git = "https://github.com/crabbuild/cellule.git", rev = "16106
ed25519-dalek = "2"
flate2 = "1.1"
futures-core = "0.3"
futures-util = { version = "0.3", default-features = false, features = ["std"] }
hex = "0.4"
http-body = "1"
object_store = "0.14.1"
Expand Down
4 changes: 4 additions & 0 deletions crates/canopy-server/src/admission.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,10 @@ pub(crate) struct AccountAdmission {
}

impl AccountAdmission {
pub(crate) fn available(&self) -> usize {
self.total.available_permits()
}

pub(crate) fn new(
limit: usize,
total_capacity: &'static str,
Expand Down
8 changes: 8 additions & 0 deletions crates/canopy-server/src/deployment/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,14 @@ mod recovery;
mod root;
pub use recovery::WorkerConfig;

/// Incompatible repository deployment and local cache format.
pub const STORAGE_FORMAT: &str = "canopy-pack-v1";

/// Read-only admission before local reclamation, probes or Cell activation.
pub(crate) async fn validate_service_root(store: &Store, prefix: &Path) -> Result<()> {
root::validate_service(store, prefix).await
}

/// Application-wide admission shared by nodes and offline administration.
#[derive(Clone)]
pub struct Deployment {
Expand Down
2 changes: 1 addition & 1 deletion crates/canopy-server/src/deployment/recovery.rs
Original file line number Diff line number Diff line change
Expand Up @@ -168,7 +168,7 @@ impl Deployment {
let (module, schema) = if entry.namespace() == directory::DIRECTORY {
(DirectoryModule::NAME, directory::SCHEMA)
} else if entry.namespace() == REPOSITORIES {
(RepositoryModule::NAME, include_str!("../schema.sql"))
(RepositoryModule::NAME, crate::REPOSITORY_SCHEMA)
} else {
return Err(Error::Release("unknown maintenance Cell namespace").into());
};
Expand Down
75 changes: 59 additions & 16 deletions crates/canopy-server/src/deployment/root.rs
Original file line number Diff line number Diff line change
Expand Up @@ -70,21 +70,74 @@ pub(super) struct RootClaim {
token: ETag,
}

#[derive(Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
struct RootEnvelope<P> {
format: String,
purpose: P,
}

fn encode(purpose: &RootPurpose) -> Result<Bytes> {
if matches!(purpose, RootPurpose::Backup { source, pin, .. } | RootPurpose::Restore { source, pin, .. }
if source.len() > 4096 || pin.len() > 4096)
{
return Err(Error::Backup("root reservation exceeds size limit"));
}
let body = Bytes::from(serde_json::to_vec(&RootEnvelope {
format: STORAGE_FORMAT.into(),
purpose,
})?);
if body.len() > 4096 {
return Err(Error::Backup("root reservation exceeds size limit"));
}
Ok(body)
}

fn path(root: &Path) -> Path {
root.clone().join("canopy-root-v1.json")
}

pub(super) async fn load(store: &Store, root: &Path) -> Result<Option<RootClaim>> {
match store.get_with_etag_bounded(&path(root), 4096).await {
Ok((bytes, token)) => Ok(Some(RootClaim {
purpose: serde_json::from_slice(&bytes)?,
token,
})),
Ok((bytes, token)) => {
let envelope: RootEnvelope<RootPurpose> = serde_json::from_slice(&bytes)?;
if envelope.format != STORAGE_FORMAT {
return Err(Error::Backup("unrecognized Canopy storage format"));
}
Ok(Some(RootClaim {
purpose: envelope.purpose,
token,
}))
}
Err(StorageError::NotFound { .. }) => Ok(None),
Err(error) => Err(error.into()),
}
}

pub(super) async fn validate_service(store: &Store, root: &Path) -> Result<()> {
let mut claim = load(store, root).await?;
if claim.is_none()
&& ApplicationIdentityStore::new(store.clone(), root.clone())
.load()
.await?
.is_some()
{
// A concurrent initializer reserves its marker before its identity.
claim = load(store, root).await?;
if claim.is_none() {
return Err(Error::Backup(
"destination already contains an application identity",
));
}
}
if claim.is_some_and(|claim| !claim.purpose.permits_service()) {
return Err(Error::Backup(
"backup or unfinished restore prefix cannot serve",
));
}
Ok(())
}

pub(super) async fn reserve(store: &Store, root: &Path, purpose: RootPurpose) -> Result<RootClaim> {
if load(store, root).await?.is_none()
&& ApplicationIdentityStore::new(store.clone(), root.clone())
Expand All @@ -99,10 +152,7 @@ pub(super) async fn reserve(store: &Store, root: &Path, purpose: RootPurpose) ->
"destination already contains an application identity",
));
}
let body = Bytes::from(serde_json::to_vec(&purpose)?);
if body.len() > 4096 {
return Err(Error::Backup("root reservation exceeds size limit"));
}
let body = encode(&purpose)?;
match store.create_strict_with_etag(&path(root), body).await {
Ok(token) => Ok(RootClaim { purpose, token }),
Err(error) => match load(store, root).await? {
Expand Down Expand Up @@ -131,14 +181,7 @@ impl RootClaim {
}
RootPurpose::Service => return Err(Error::Backup("service root cannot finish a copy")),
}
match store
.update(
&path(root),
Bytes::from(serde_json::to_vec(&next)?),
self.token,
)
.await
{
match store.update(&path(root), encode(&next)?, self.token).await {
Ok(_) => Ok(()),
Err(error) => match load(store, root).await? {
Some(current) if current.purpose == next => Ok(()),
Expand Down
78 changes: 78 additions & 0 deletions crates/canopy-server/src/deployment/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,84 @@ type TestResult = std::result::Result<(), Box<dyn std::error::Error>>;

mod retained_maintenance;

#[tokio::test]
async fn root_reservation_limit_includes_the_format_envelope_before_any_write() -> TestResult {
let deployment = fixture()?;
let mut purpose = root::RootPurpose::Backup {
source: String::new(),
pin: uuid::Uuid::new_v4().to_string(),
complete: false,
};
let previous_overhead = serde_json::to_vec(&purpose)?.len();
if let root::RootPurpose::Backup { source, .. } = &mut purpose {
*source = "s".repeat(4096 - previous_overhead);
}
assert_eq!(serde_json::to_vec(&purpose)?.len(), 4096);
assert!(matches!(
root::reserve(deployment.layout.store(), &deployment.prefix, purpose).await,
Err(Error::Backup("root reservation exceeds size limit"))
));
assert!(
root::load(deployment.layout.store(), &deployment.prefix)
.await?
.is_none()
);
assert!(deployment.identities.load().await?.is_none());
Ok(())
}

#[tokio::test]
async fn old_or_unknown_root_format_cannot_initialize_identity_or_release() -> TestResult {
for bytes in [
br#"{"kind":"service"}"#.as_slice(),
br#"{"purpose":{"kind":"service"}}"#,
br#"{"format":"future-format","purpose":{"kind":"service"}}"#,
br#"{"format":"canopy-pack-v1","purpose":{"kind":"service"},"extra":true}"#,
] {
let deployment = fixture()?;
let path = deployment.prefix.clone().join("canopy-root-v1.json");
let original = bytes::Bytes::copy_from_slice(bytes);
deployment
.layout
.store()
.create_strict(&path, original.clone())
.await?;
let before = deployment.layout.store().get_with_etag(&path).await?;
assert!(
root::load(deployment.layout.store(), &deployment.prefix)
.await
.is_err()
);
assert!(deployment.initialize().await.is_err());
assert!(deployment.identities.load().await?.is_none());
assert!(deployment.releases.load().await?.is_none());
assert_eq!(
deployment.layout.store().get_with_etag(&path).await?,
before
);
}
Ok(())
}

#[tokio::test]
async fn packed_root_envelope_reuses_purpose_and_is_rejected_by_old_decoder() -> TestResult {
let deployment = fixture()?;
deployment.initialize().await?;
let path = deployment.prefix.clone().join("canopy-root-v1.json");
let (bytes, _) = deployment.layout.store().get_with_etag(&path).await?;
assert_eq!(
serde_json::from_slice::<serde_json::Value>(&bytes)?,
serde_json::json!({
"format": "canopy-pack-v1", "purpose": { "kind": "service" }
})
);
// The old decoder is exactly the existing tagged RootPurpose type.
assert!(serde_json::from_slice::<root::RootPurpose>(&bytes).is_err());
deployment.initialize().await?;
deployment.require_ready().await?;
Ok(())
}

fn fixture() -> std::result::Result<Deployment, Box<dyn std::error::Error>> {
let app = CanopyApplication::compile(build_descriptor(
include_bytes!("../../../../Cargo.lock"),
Expand Down
16 changes: 16 additions & 0 deletions crates/canopy-server/src/git_cache/artifacts.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,15 +6,27 @@ use canopy_object_storage::artifact::{ArtifactKind, ArtifactStore};
struct Writer {
file: File,
_cache: Arc<GitCache>,
_owner: crate::git_objects::ReadOwner,
}
impl GitCache {
/// The isolated verifier calls this exactly once on its fresh private cache.
/// Reserve the complete pair before creating files or reading the provider.
/// No second pack copy or blob-as-artifact wrapper is involved.
#[cfg(test)]
pub(crate) async fn download_native(
self: &Arc<Self>,
store: &ArtifactStore,
descriptor: NativePackDescriptor,
) -> Result<(), MetadataError> {
self.download_native_owned(store, descriptor, Arc::new(()))
.await
}

pub(crate) async fn download_native_owned(
self: &Arc<Self>,
store: &ArtifactStore,
descriptor: NativePackDescriptor,
owner: crate::git_objects::ReadOwner,
) -> Result<(), MetadataError> {
descriptor
.validate(store.repository(), self.object_format)
Expand All @@ -23,7 +35,9 @@ impl GitCache {
return Err(MetadataError::Integrity);
}
let cache = Arc::clone(self);
let admission = Arc::clone(&owner);
tokio::task::spawn_blocking(move || {
let _owner = admission;
let size = descriptor
.pack
.size
Expand All @@ -38,6 +52,7 @@ impl GitCache {
(ArtifactKind::Index, descriptor.index, "idx"),
] {
let cache = Arc::clone(self);
let admission = Arc::clone(&owner);
let mut writer = tokio::task::spawn_blocking(move || {
let path = cache.git_dir().join(format!(
"objects/pack/pack-{}.{}",
Expand All @@ -47,6 +62,7 @@ impl GitCache {
Ok::<_, MetadataError>(Writer {
file: File::create_new(path)?,
_cache: cache,
_owner: admission,
})
})
.await??;
Expand Down
5 changes: 5 additions & 0 deletions crates/canopy-server/src/git_cache/cleanup.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ pub(super) struct Cleanup {
pub(super) path: PathBuf,
pub(super) reservation: Option<DiskReservation>,
pub(super) objects: Option<Arc<GitCache>>,
pub(super) owner: Option<crate::git_objects::ReadOwner>,
}

impl Cleanup {
Expand Down Expand Up @@ -45,6 +46,7 @@ impl Cleanup {
// Successful removal permits normal field teardown.
self.reservation.take();
self.objects.take();
self.owner.take();
}).await;
return;
}
Expand All @@ -68,6 +70,9 @@ impl Drop for Cleanup {
if let Some(reservation) = self.reservation.take() {
std::mem::forget(reservation);
}
if let Some(owner) = self.owner.take() {
std::mem::forget(owner);
}
if let Some(objects) = self.objects.take() {
std::mem::forget(objects);
}
Expand Down
9 changes: 7 additions & 2 deletions crates/canopy-server/src/git_cache/maintenance.rs
Original file line number Diff line number Diff line change
@@ -1,9 +1,13 @@
//! Immutable cache generations: never repack/delete files beneath an active reader.
use super::*;
#[cfg(test)]
use crate::git_http::{GitHttpError, GitProcess, WORKER_DEADLINE, read_bounded};
#[cfg(test)]
use std::process::Stdio;
#[cfg(test)]
use tokio::io::AsyncWriteExt;

#[cfg(test)]
fn worker_error(error: GitHttpError) -> CacheError {
io::Error::other(error).into()
}
Expand Down Expand Up @@ -247,6 +251,7 @@ impl GitCache {
}
/// Reuse only packs whose *every* object was verified and durably recorded
/// during this ingestion. Extra/unverified objects disable this optimization.
#[cfg(test)]
pub(crate) async fn retain_verified_packs(
self: &Arc<Self>,
source: Arc<Self>,
Expand Down Expand Up @@ -318,6 +323,7 @@ impl GitCache {

/// Enumerate a captured cache into a new self-contained pack. Old objects
/// and packs are untouched; dropping the last old reader reclaims them.
#[cfg(test)]
pub(crate) async fn repacked(
self: &Arc<Self>,
root: PathBuf,
Expand Down Expand Up @@ -346,7 +352,6 @@ impl GitCache {
next.reservation()?
.try_grow(reserve)
.map_err(io::Error::other)?;
let coverage = self.prepared.lock().await.clone();
let durable = self
.durable_packs
.read()
Expand Down Expand Up @@ -497,7 +502,6 @@ impl GitCache {
result.map_err(worker_error)?;
next.pack_files
.store(1, std::sync::atomic::Ordering::Relaxed);
*next.prepared.lock().await = coverage;
*next
.durable_packs
.write()
Expand All @@ -506,6 +510,7 @@ impl GitCache {
}
}

#[cfg(test)]
async fn finish<T: Send + 'static>(
process: &mut GitProcess<T>,
stderr: Vec<u8>,
Expand Down
Loading
Loading