Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
923e5ed
feat(api): delegate durable lifecycle to storage admission
DecisionNerd Aug 16, 2026
feff3f5
test(bindings): cover filesystem admission parity
DecisionNerd Aug 16, 2026
19326f3
test(api): certify concurrent project admission
DecisionNerd Aug 16, 2026
c95f8bb
feat(storage): gate project lifecycle admission (#780)
DecisionNerd Aug 16, 2026
d4a50e8
feat(storage): route durable publishers through admission (#780)
DecisionNerd Aug 16, 2026
8f909f9
docs(storage): freeze lifecycle admission contract (#780)
DecisionNerd Aug 16, 2026
94bfb29
fix(storage): preserve writer concurrency under admission
DecisionNerd Aug 16, 2026
8c43083
fix(storage): anchor lifecycle admission traversal (#780)
DecisionNerd Aug 16, 2026
a8965c8
fix(storage): preserve admitted publication authority (#780)
DecisionNerd Aug 17, 2026
745ca22
fix(api): preserve ephemeral lifecycle mode (#780)
DecisionNerd Aug 17, 2026
993f0bd
style(bindings): format admission smoke coverage
DecisionNerd Aug 17, 2026
47713d3
fix(storage): reconcile admitted publication outcomes (#780)
DecisionNerd Aug 17, 2026
17fd144
test(cli): map filesystem admission in Bazel (#780)
DecisionNerd Aug 17, 2026
a6d70c4
fix(storage): align platform admission checks (#780)
DecisionNerd Aug 17, 2026
e96121a
test(bdd): align absent project lifecycle contract (#780)
DecisionNerd Aug 17, 2026
4f4cc8e
test(bdd): share admissible absent-root fixture (#780)
DecisionNerd Aug 17, 2026
c1c632d
fix(storage): release inherited checkpoint reads (#780)
DecisionNerd Aug 17, 2026
7ac7483
fix(storage): close lifecycle review races (#780)
DecisionNerd Aug 17, 2026
a80cb4a
fix(storage): preserve admitted commit ordering (#780)
DecisionNerd Aug 17, 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
5 changes: 1 addition & 4 deletions crates/graphforge-api/src/algorithm_runs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -519,10 +519,7 @@ fn publish(
.collect(),
participants,
};
let receipt = match graphforge_storage::stage_project_generation(
graph.resolved_generation.container_root(),
&request,
)? {
let receipt = match graph.stage_project_generation(&request)? {
ProjectStageOutcome::AlreadyPublished(receipt) => receipt,
ProjectStageOutcome::Staged(staged) => staged
.validate(
Expand Down
5 changes: 1 addition & 4 deletions crates/graphforge-api/src/belief_projection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1313,10 +1313,7 @@ fn publish_attachment(
.collect(),
participants,
};
let receipt = match graphforge_storage::stage_project_generation(
graph.resolved_generation.container_root(),
&request,
)? {
let receipt = match graph.stage_project_generation(&request)? {
ProjectStageOutcome::AlreadyPublished(receipt) => receipt,
ProjectStageOutcome::Staged(staged) => staged
.validate(
Expand Down
47 changes: 25 additions & 22 deletions crates/graphforge-api/src/capabilities.rs
Original file line number Diff line number Diff line change
Expand Up @@ -228,28 +228,27 @@ impl GraphForge {
capabilities,
participants,
};
let generation_uuid =
match graphforge_storage::stage_project_generation(root, &publication)? {
ProjectStageOutcome::AlreadyPublished(receipt) => receipt.generation_uuid,
ProjectStageOutcome::Staged(staged) => {
let expected_parent = parent.generation_uuid();
staged
.validate(
|_| Ok(()),
|actual_parent, _| {
if actual_parent.generation_uuid() != expected_parent {
return Err(GfError::Validation(
"project generation changed before capability publication"
.into(),
));
}
Ok(())
},
)?
.publish()?
.generation_uuid
}
};
let generation_uuid = match self.stage_project_generation(&publication)? {
ProjectStageOutcome::AlreadyPublished(receipt) => receipt.generation_uuid,
ProjectStageOutcome::Staged(staged) => {
let expected_parent = parent.generation_uuid();
staged
.validate(
|_| Ok(()),
|actual_parent, _| {
if actual_parent.generation_uuid() != expected_parent {
return Err(GfError::Validation(
"project generation changed before capability publication"
.into(),
));
}
Ok(())
},
)?
.publish()?
.generation_uuid
}
};
*self
.current_generation_uuid
.lock()
Expand Down Expand Up @@ -396,6 +395,10 @@ mod tests {
#[test]
fn capability_enable_is_atomic_and_idempotent() {
let graph = GraphForge::new(None).unwrap();
assert_eq!(
graph.lifecycle_mode,
graphforge_storage::filesystem_admission::ProjectLifecycleMode::Ephemeral
);
let request = EnableCapabilityRequest {
context: WriteContext {
operation_uuid: OperationId(Uuid::now_v7()),
Expand Down
16 changes: 15 additions & 1 deletion crates/graphforge-api/src/checkpoint_graph_diff.rs
Original file line number Diff line number Diff line change
Expand Up @@ -38,15 +38,29 @@ pub(crate) struct LogicalGraphRecords {
/// relationship types and endpoints, and user properties. Storage surrogate
/// IDs, timestamps, Parquet row groups, archive paths, and batch boundaries do
/// not participate.
#[cfg(test)]
pub(crate) fn extract_logical_graph_records(
generation: &ResolvedProjectGeneration,
cancellation: Option<&CancellationToken>,
) -> Result<LogicalGraphRecords, GfError> {
extract_logical_graph_records_with_mode(
generation,
cancellation,
graphforge_storage::filesystem_admission::ProjectLifecycleMode::Durable,
)
}

pub(crate) fn extract_logical_graph_records_with_mode(
generation: &ResolvedProjectGeneration,
cancellation: Option<&CancellationToken>,
lifecycle_mode: graphforge_storage::filesystem_admission::ProjectLifecycleMode,
) -> Result<LogicalGraphRecords, GfError> {
checkpoint(cancellation)?;
let graph = GraphForge::open_resolved_with_mode(
let graph = GraphForge::open_resolved_with_lifecycle_mode(
generation.container_root().to_path_buf(),
generation.clone(),
true,
lifecycle_mode,
)?;

let nodes = streamed_logical_rows(
Expand Down
87 changes: 74 additions & 13 deletions crates/graphforge-api/src/checkpoints.rs
Original file line number Diff line number Diff line change
Expand Up @@ -638,14 +638,15 @@ impl CheckpointView {
impl GraphForge {
/// Create a durable named checkpoint.
pub fn checkpoint(&self, request: CheckpointRequest) -> Result<ExecutionResult, GfError> {
let receipt = graphforge_storage::create_checkpoint(
let receipt = graphforge_storage::create_checkpoint_with_mode(
self.resolved_generation.container_root(),
&graphforge_storage::CheckpointCreateRequest {
operation_uuid: request.idempotency_key.0,
name: request.name,
description: request.description,
actor_uuid: request.actor_uuid,
},
self.lifecycle_mode,
)?;
Ok(receipt_result(&receipt))
}
Expand All @@ -657,7 +658,10 @@ impl GraphForge {
) -> Result<ExecutionResult, GfError> {
let ListCheckpointsRequest { page } = request;
cancellation(&page)?;
let rows = graphforge_storage::list_checkpoints(self.resolved_generation.container_root())?;
let rows = graphforge_storage::list_checkpoints_with_mode(
self.resolved_generation.container_root(),
self.lifecycle_mode,
)?;
let snapshot = checkpoint_list_snapshot(&rows);
let binding = request_binding("checkpoint-list", 0, 0);
let cursors = rows
Expand Down Expand Up @@ -688,9 +692,10 @@ impl GraphForge {
request: ShowCheckpointRequest,
) -> Result<ExecutionResult, GfError> {
let ShowCheckpointRequest { name } = request;
let (checkpoint, _) = graphforge_storage::open_checkpoint_generation(
let (checkpoint, _) = graphforge_storage::open_checkpoint_generation_with_mode(
self.resolved_generation.container_root(),
&name,
self.lifecycle_mode,
)?;
checkpoint_rows(std::slice::from_ref(&checkpoint), None)
}
Expand All @@ -713,14 +718,16 @@ impl GraphForge {

/// Open an immutable view pinned to the named checkpoint generation.
pub fn open_checkpoint(&self, name: &str) -> Result<CheckpointView, GfError> {
let (checkpoint, generation) = graphforge_storage::open_checkpoint_generation(
let (checkpoint, generation) = graphforge_storage::open_checkpoint_generation_with_mode(
self.resolved_generation.container_root(),
name,
self.lifecycle_mode,
)?;
let graph = Self::open_resolved_with_mode(
let graph = Self::open_resolved_with_lifecycle_mode(
self.resolved_generation.container_root().to_owned(),
generation,
true,
self.lifecycle_mode,
)?;
Ok(CheckpointView { checkpoint, graph })
}
Expand All @@ -730,13 +737,14 @@ impl GraphForge {
&self,
request: DeleteCheckpointRequest,
) -> Result<ExecutionResult, GfError> {
let receipt = graphforge_storage::delete_checkpoint(
let receipt = graphforge_storage::delete_checkpoint_with_mode(
self.resolved_generation.container_root(),
&graphforge_storage::CheckpointDeleteRequest {
operation_uuid: request.idempotency_key.0,
name: request.name,
actor_uuid: request.actor_uuid,
},
self.lifecycle_mode,
)?;
Ok(receipt_result(&receipt))
}
Expand All @@ -751,11 +759,12 @@ impl GraphForge {
}
let container_root = self.resolved_generation.container_root().to_path_buf();
let clock = self.clock.lock().expect("clock lock poisoned").clone();
let lifecycle_mode = self.lifecycle_mode;
let write_options = self.write_options.clone();
let resource_policy = self.resource_policy.clone();
let select_clock = Arc::clone(&clock);
let prepared = std::cell::RefCell::new(None);
let (receipt, resolved) = graphforge_storage::revert_checkpoint(
let (receipt, resolved) = graphforge_storage::revert_checkpoint_with_mode(
&container_root,
&graphforge_storage::CheckpointRevertRequest {
operation_uuid: request.idempotency_key.0,
Expand All @@ -765,7 +774,7 @@ impl GraphForge {
},
move || select_clock(),
|generation| {
validate_revert_source(generation)?;
validate_revert_source(generation, lifecycle_mode)?;
prepared.replace(Some(GraphForge::open_resolved_with_options(
container_root.clone(),
generation.clone(),
Expand All @@ -778,12 +787,14 @@ impl GraphForge {
)?));
Ok(())
},
self.lifecycle_mode,
)?;
let result = receipt_result(&receipt);

let mut reopened = prepared
.into_inner()
.expect("successful revert validation prepares the replacement facade");
reopened.lifecycle_mode = lifecycle_mode;
reopened.resolved_generation = resolved;
*reopened
.current_generation_uuid
Expand Down Expand Up @@ -838,16 +849,19 @@ impl GraphForge {
let to = resolve(&to)?;
match detail {
CheckpointDiffDetail::Summary => summary_diff(&from, &to, scope, binding, &page),
CheckpointDiffDetail::Records => record_diff(&from, &to, scope, binding, &page),
CheckpointDiffDetail::Records => {
record_diff(&from, &to, scope, binding, &page, self.lifecycle_mode)
}
}
}

fn resolve_selector(&self, selector: &CheckpointSelector) -> Result<DiffEndpoint, GfError> {
let (checkpoint_uuid, generation) = match selector {
CheckpointSelector::Named(name) => {
let (row, generation) = graphforge_storage::open_checkpoint_generation(
let (row, generation) = graphforge_storage::open_checkpoint_generation_with_mode(
self.resolved_generation.container_root(),
name,
self.lifecycle_mode,
)?;
(row.checkpoint_uuid, generation)
}
Expand Down Expand Up @@ -945,9 +959,10 @@ fn record_diff(
scope: CheckpointDiffScope,
binding: Uuid,
page: &PageRequest,
lifecycle_mode: graphforge_storage::filesystem_admission::ProjectLifecycleMode,
) -> Result<ExecutionResult, GfError> {
let left = logical_records(&from.generation, scope, page)?;
let right = logical_records(&to.generation, scope, page)?;
let left = logical_records(&from.generation, scope, page, lifecycle_mode)?;
let right = logical_records(&to.generation, scope, page, lifecycle_mode)?;
let mut keys = left.keys().chain(right.keys()).cloned().collect::<Vec<_>>();
keys.sort();
keys.dedup();
Expand Down Expand Up @@ -1023,6 +1038,7 @@ fn logical_records(
generation: &graphforge_storage::ResolvedProjectGeneration,
scope: CheckpointDiffScope,
page: &PageRequest,
lifecycle_mode: graphforge_storage::filesystem_admission::ProjectLifecycleMode,
) -> Result<LogicalRecords, GfError> {
let adapters = record_adapters()?;
let mut out = BTreeMap::new();
Expand All @@ -1035,9 +1051,10 @@ fn logical_records(
if descriptor.capability_id == "graph"
&& matches!(descriptor.record_family_id.as_str(), "snapshot" | "files")
{
let records = crate::checkpoint_graph_diff::extract_logical_graph_records(
let records = crate::checkpoint_graph_diff::extract_logical_graph_records_with_mode(
generation,
page.cancellation.as_ref(),
lifecycle_mode,
)?;
for (family, records) in [("nodes", records.nodes), ("edges", records.edges)] {
for record in records {
Expand Down Expand Up @@ -1307,6 +1324,7 @@ type Inventory =

fn validate_revert_source(
generation: &graphforge_storage::ResolvedProjectGeneration,
lifecycle_mode: graphforge_storage::filesystem_admission::ProjectLifecycleMode,
) -> Result<(), GfError> {
generation.validate_complete_participant_inventory()?;
let _workspace = crate::hydrate_graph_workspace(generation, true)?;
Expand All @@ -1330,6 +1348,7 @@ fn validate_revert_source(
generation,
CheckpointDiffScope::All,
&PageRequest::default(),
lifecycle_mode,
)?;
// Run each domain owner's decoder as well as the generic checkpoint adapters.
// These readers enforce each ledger's schema and ledger-local invariants.
Expand Down Expand Up @@ -3358,6 +3377,48 @@ mod tests {
);
}

#[test]
fn in_memory_checkpoint_lifecycle_remains_ephemeral() {
let mut graph = GraphForge::new(None).unwrap();
graph.execute("CREATE (:Person {name: 'before'})").unwrap();
graph
.checkpoint(CheckpointRequest {
name: "Ephemeral".into(),
description: None,
idempotency_key: operation(230),
actor_uuid: None,
})
.unwrap();
graph.execute("CREATE (:Person {name: 'after'})").unwrap();
graph
.revert_to_checkpoint(RevertCheckpointRequest {
name: "Ephemeral".into(),
reason: "restore ephemeral checkpoint".into(),
idempotency_key: operation(232),
actor_uuid: None,
})
.unwrap();
assert_eq!(
graph.lifecycle_mode,
graphforge_storage::filesystem_admission::ProjectLifecycleMode::Ephemeral
);
assert_eq!(
graph
.list_checkpoints(ListCheckpointsRequest::default())
.unwrap()
.stats
.rows_produced,
1
);
graph
.delete_checkpoint(DeleteCheckpointRequest {
name: "Ephemeral".into(),
idempotency_key: operation(233),
actor_uuid: None,
})
.unwrap();
}

#[test]
fn list_and_summary_diff_are_arrow_ordered_and_page_bound() {
let directory = tempdir().unwrap();
Expand Down
Loading