diff --git a/document/graph-storage/src/from_runtime.rs b/document/graph-storage/src/from_runtime.rs index 110f60200c..0c5937a04a 100644 --- a/document/graph-storage/src/from_runtime.rs +++ b/document/graph-storage/src/from_runtime.rs @@ -182,33 +182,40 @@ fn convert_resources(resources: &graphene_resource::ResourceRegistry, referenced if !referenced.contains(&id) { continue; } - let Some(info) = resources.info(&id) else { continue }; - - let mut entry = crate::ResourceEntry { - hash: info.hash.copied(), - hash_timestamp: TimeStamp::ORIGIN, - ..Default::default() - }; - for (position, source) in info.sources.iter().enumerate() { - let key = crate::SourceKey { - priority: crate::Priority::new(position as f64).expect("enumerate index is finite"), - peer, - }; - let body = serde_json::to_value(source).map_err(|error| ConversionError::SerializationError(error.to_string()))?; - entry.set_source( - key, - crate::SourceValue { - source: body, - timestamp: TimeStamp::ORIGIN, - }, - ); + if let Some(entry) = convert_resource_entry(resources, id, peer)? { + registry.resources.insert(id, entry); } - - registry.resources.insert(id, entry); } Ok(()) } +/// Snapshot one runtime resource into a storage entry, or `None` if the runtime has no such resource. +pub fn convert_resource_entry(resources: &graphene_resource::ResourceRegistry, id: ResourceId, peer: PeerId) -> Result, ConversionError> { + let Some(info) = resources.info(&id) else { return Ok(None) }; + + let mut entry = crate::ResourceEntry { + hash: info.hash.copied(), + hash_timestamp: TimeStamp::ORIGIN, + ..Default::default() + }; + for (position, source) in info.sources.iter().enumerate() { + let key = crate::SourceKey { + priority: crate::Priority::new(position as f64).expect("enumerate index is finite"), + peer, + }; + let body = serde_json::to_value(source).map_err(|error| ConversionError::SerializationError(error.to_string()))?; + entry.set_source( + key, + crate::SourceValue { + source: body, + timestamp: TimeStamp::ORIGIN, + }, + ); + } + + Ok(Some(entry)) +} + /// Collect the `ResourceId`s referenced by `TaggedValue::Resource` inputs anywhere in the network /// (recursively through nested networks). These are the resources the document actually uses; the /// runtime cache may hold more (history-retained orphans) that shouldn't be snapshotted into storage. @@ -279,7 +286,7 @@ fn convert_network( metadata_path, runtime_node_id: *runtime_node_id, }; - let mut node = convert_node(doc_node, location, registry, ctx)?; + let mut node = convert_node(doc_node, location, registry, ctx, true)?; node.attributes.set(node::ORIGINAL_NODE_ID, serde_json::json!(runtime_node_id.0), TimeStamp::ORIGIN); registry.node_instances.insert(global_id, node); } @@ -351,7 +358,13 @@ struct NodeLocation<'a> { runtime_node_id: RuntimeNodeId, } -fn convert_node(doc_node: &DocumentNode, location: NodeLocation<'_>, registry: &mut Registry, ctx: &mut ConversionContext<'_, M>) -> Result { +fn convert_node( + doc_node: &DocumentNode, + location: NodeLocation<'_>, + registry: &mut Registry, + ctx: &mut ConversionContext<'_, M>, + recurse: bool, +) -> Result { let NodeLocation { network_id, parent_path, @@ -383,7 +396,7 @@ fn convert_node(doc_node: &DocumentNode, locatio } else { metadata_path }; - let implementation = convert_implementation(&doc_node.implementation, &node_path, child_metadata_path, registry, ctx)?; + let implementation = convert_implementation(&doc_node.implementation, &node_path, child_metadata_path, registry, ctx, recurse)?; // Defaults match `DocumentNode::default()`; `to_runtime` rehydrates absent keys from the same defaults. let mut attributes = crate::Attributes::new(); @@ -533,6 +546,7 @@ fn convert_implementation( child_metadata_path: &[RuntimeNodeId], registry: &mut Registry, ctx: &mut ConversionContext<'_, M>, + recurse: bool, ) -> Result { Ok(match implementation { DocumentNodeImplementation::ProtoNode(identifier) => { @@ -564,10 +578,146 @@ fn convert_implementation( // `to_runtime` -> `from_runtime` round trip reproduces the same `NetworkId` (and thus the // same node-path hashes underneath it). let nested_network_id = current_node_path.owned_network_id(ctx.peer); - convert_network(nested_network, nested_network_id, Some(current_node_path), child_metadata_path, registry, ctx)?; + if recurse { + convert_network(nested_network, nested_network_id, Some(current_node_path), child_metadata_path, registry, ctx)?; + } Implementation::Network(nested_network_id) } // TODO: Support Extract in the Registry format. DocumentNodeImplementation::Extract => return Err(ConversionError::UnsupportedImplementation), }) } + +/// Derives the stable storage IDs for entities addressed by runtime network paths, walking the +/// same path-hash chain as a full conversion so scoped and whole-document conversions agree. +pub struct PathResolver { + peer: PeerId, +} + +impl PathResolver { + pub fn new(peer: PeerId) -> Self { + Self { peer } + } + + /// The storage `NetworkId` of the network at `local_path` (`ROOT_NETWORK` for the empty path). + pub fn network_id(&self, local_path: &[RuntimeNodeId]) -> NetworkId { + match self.owner_path(local_path) { + None => ROOT_NETWORK, + Some(owner) => owner.owned_network_id(self.peer), + } + } + + /// The storage `NodeId` of the node with `local_id` inside the network at `local_path`. + pub fn node_id(&self, local_path: &[RuntimeNodeId], local_id: RuntimeNodeId) -> NodeId { + let owner = self.owner_path(local_path); + child_path(owner.as_ref(), self.network_id(local_path), local_id).to_global_id(self.peer) + } + + /// The `NodePath` of the node owning the network at `local_path`, or `None` for the root network. + fn owner_path(&self, local_path: &[RuntimeNodeId]) -> Option { + let mut owner: Option = None; + for &local_id in local_path { + let network_id = match &owner { + None => ROOT_NETWORK, + Some(owner) => owner.owned_network_id(self.peer), + }; + owner = Some(child_path(owner.as_ref(), network_id, local_id)); + } + owner + } +} + +/// Converts a chosen subset of runtime entities into a scratch [`Registry`] through the same +/// encoders as a whole-document conversion, so staging can derive deltas for just those entities. +pub struct ScopedConversion<'m, M: NodeMetadataSource + ?Sized> { + ctx: ConversionContext<'m, M>, + resolver: PathResolver, +} + +impl<'m, M: NodeMetadataSource + ?Sized> ScopedConversion<'m, M> { + pub fn new(metadata: &'m M, peer: PeerId) -> Self { + Self { + ctx: ConversionContext { + declaration_ids: HashMap::new(), + declaration_bytes: HashMap::new(), + network_ids: HashMap::new(), + metadata, + peer, + }, + resolver: PathResolver::new(peer), + } + } + + /// Converts the node with `local_id` in the network at `local_path` into `registry` under its + /// stable IDs. `recurse` also converts the contents of any nested network the node implements; + /// without it only the node itself is converted, with its nested network referenced by ID. + pub fn convert_node_at(&mut self, registry: &mut Registry, local_path: &[RuntimeNodeId], local_id: RuntimeNodeId, doc_node: &DocumentNode, recurse: bool) -> Result<(), ConversionError> { + let owner_path = self.resolver.owner_path(local_path); + let location = NodeLocation { + network_id: self.resolver.network_id(local_path), + parent_path: owner_path.as_ref(), + metadata_path: local_path, + runtime_node_id: local_id, + }; + + let mut node = convert_node(doc_node, location, registry, &mut self.ctx, recurse)?; + node.attributes.set(node::ORIGINAL_NODE_ID, serde_json::json!(local_id.0), TimeStamp::ORIGIN); + registry.node_instances.insert(self.resolver.node_id(local_path, local_id), node); + Ok(()) + } + + /// Converts the entry of the network at `local_path` (its export slots and network-level + /// attributes, not its nodes) into `registry`. + pub fn convert_network_entry_at(&mut self, registry: &mut Registry, local_path: &[RuntimeNodeId], node_network: &NodeNetwork) -> Result<(), ConversionError> { + let owner_path = self.resolver.owner_path(local_path); + let network_id = self.resolver.network_id(local_path); + + let exports = node_network + .exports + .iter() + .map(|export| { + Ok(ExportSlot { + target: Some(convert_input(export, owner_path.as_ref(), network_id, self.ctx.peer)?), + timestamp: TimeStamp::ORIGIN, + }) + }) + .collect::, ConversionError>>()?; + + let mut attributes = crate::Attributes::new(); + write_ui_network_attributes(&mut attributes, self.ctx.metadata, local_path, TimeStamp::ORIGIN)?; + write_scope_injections(&mut attributes, node_network, owner_path.as_ref(), network_id, self.ctx.peer, TimeStamp::ORIGIN)?; + + registry.networks.insert(network_id, Network { exports, attributes }); + Ok(()) + } + + /// The extracted proto-node declaration bytes, for the caller to persist into its byte store. + pub fn finish(self) -> DeclarationBytes { + self.ctx.declaration_bytes + } +} + +/// The `TaggedValue::Resource` IDs referenced by a stored node's value inputs. +/// +/// Peeks at the externally-tagged serde shape (`{"Resource": }`) rather than deserializing the +/// whole `TaggedValue`, which can be arbitrarily large. `resource_ref_shape_matches_serde` asserts +/// this stays in lockstep with the real serialization. +pub fn node_value_resource_refs(node: &Node) -> impl Iterator + '_ { + node.inputs.iter().filter_map(|slot| match &slot.input { + NodeInput::Value { value, .. } => value.get("Resource").and_then(|id| serde_json::from_value(id.clone()).ok()), + _ => None, + }) +} + +#[cfg(test)] +mod resource_ref_shape { + use super::*; + + #[test] + fn resource_ref_shape_matches_serde() { + let id = ResourceId::from_hash(&ResourceHash::from(b"shape test".as_slice())); + let value = serde_json::to_value(TaggedValue::Resource(id)).expect("TaggedValue::Resource serializes"); + let parsed: Option = value.get("Resource").and_then(|inner| serde_json::from_value(inner.clone()).ok()); + assert_eq!(parsed, Some(id), "The shape peek in node_value_resource_refs must match TaggedValue's serde form"); + } +} diff --git a/document/graph-storage/src/lib.rs b/document/graph-storage/src/lib.rs index 806c4cb494..37034537eb 100644 --- a/document/graph-storage/src/lib.rs +++ b/document/graph-storage/src/lib.rs @@ -29,7 +29,7 @@ pub use resources::*; pub use session::*; #[cfg(any(feature = "conversion", test))] -pub use from_runtime::{RuntimeConversion, decode_declaration, encode_declaration}; +pub use from_runtime::{PathResolver, RuntimeConversion, ScopedConversion, convert_resource_entry, decode_declaration, encode_declaration, node_value_resource_refs}; #[cfg(any(feature = "conversion", test))] pub use metadata_source::{InputMetadataEntry, NetworkMetadataEntry, NoMetadata, NodeMetadataEntry, NodeMetadataSource, Position}; #[cfg(any(feature = "conversion", test))] diff --git a/document/graph-storage/src/session.rs b/document/graph-storage/src/session.rs index 8edfee51ca..714e7701f3 100644 --- a/document/graph-storage/src/session.rs +++ b/document/graph-storage/src/session.rs @@ -129,6 +129,12 @@ impl Session { Ok(revs) } + /// Stages ops computed outside the session (an incremental runtime projection) as hot ops, the + /// same way `stage_from_runtime` stages a diff's ops. + pub fn stage_computed_ops(&mut self, ops: Vec) -> Result, CrdtError> { + self.stage_ops(ops) + } + /// Apply each op as a hot op with a freshly-ticked timestamp, returning the staged frames in /// order. Each tick is strictly later than the last, so the final frame carries the latest /// timestamp, which is what the caller passes to `retire`.