diff --git a/node-graph/libraries/core-types/src/record.rs b/node-graph/libraries/core-types/src/record.rs index 25ca2c7939..ef55cd924c 100644 --- a/node-graph/libraries/core-types/src/record.rs +++ b/node-graph/libraries/core-types/src/record.rs @@ -1644,12 +1644,6 @@ impl RecordSource { } } -/// # Safety -/// `rec` must be a live record of `layout`. -pub unsafe fn copy_record_bytes(layout: &Layout, rec: Rec<'_>) -> Box<[u8]> { - unsafe { std::slice::from_raw_parts(rec.ptr(), layout.size) }.into() -} - /// A node's own output frame, minted by its caller from the caller's frame /// space: the one closing surface for every exit. Writes land through it, /// [`Self::lift`] and [`Self::finish`] serve the record, and it carries the diff --git a/node-graph/nodes/gcore/src/memo.rs b/node-graph/nodes/gcore/src/memo.rs index 6558fa1d51..3ce26695cc 100644 --- a/node-graph/nodes/gcore/src/memo.rs +++ b/node-graph/nodes/gcore/src/memo.rs @@ -1,32 +1,33 @@ -use core_types::arena::ArenaCell; use core_types::context::{Ctx, CtxSnapshot, DeriveCtx, ExtractAll, ModifyIndex}; -use core_types::frame_table::{FrameTable, Lookup}; use core_types::gpoll::{Finality, GPoll}; use core_types::graphene_hash::CacheHash; -use core_types::record::{FrameClaim, LevelStatus, MaterializedSpan, Served, copy_record_bytes}; +use core_types::record::{FrameClaim, LevelStatus, MaterializedSpan, OwnedRecord, Promotion, Served}; use core_types::registry::cache_key; use std::sync::Arc; use std::sync::Mutex; -/// The memo entry: the level lives in the persistent region, which no -/// evaluation resets, so a hit copies its lane's bytes and the parked -/// references they carry stay live without a re-park. +/// The memo entry: the deep copies survive a persistent flush, while the span +/// serves lanes straight out of the persistent region (generation-guarded), so +/// a hit before the next flush allocates nothing. #[derive(Debug)] pub struct MemoLevel { key: u64, /// The persistent region the level was promoted into, resolvable only /// while that region's epoch is live. - span: MaterializedSpan, + span: Option, + lanes: Vec, finality: Finality, } /// Helps speed up repeated renders in a computationally-heavy part of the node graph. /// -/// Promotes the last record (a scalar wire) or the last whole level (a leveled -/// wire) that flowed through this node into the persistent region and serves -/// it on subsequent renders if the context has not changed. A leveled wire's -/// cache key normalizes the addressed lane away, so per-lane pulls share one -/// materialization of the content instead of re-evaluating it per lane. +/// Stores a deep copy of the last record (a scalar wire) or the last whole +/// level (a leveled wire) that flowed through this node and replays it on +/// subsequent renders if the context has not changed. The owned copies survive +/// a persistent flush, so this is the memo for content whose recomputation is +/// expensive. A leveled wire's cache key normalizes the addressed lane away, +/// so per-lane pulls share one materialization of the content instead of +/// re-evaluating it per lane. #[node_macro::node(category("General"), path(graphene_core::memo))] fn memoize<'e, 'l>( ctx: impl Ctx + CacheHash + DeriveCtx + ExtractArena<'e> + ModifyIndex + Copy, @@ -50,8 +51,124 @@ fn memoize<'e, 'l>( } false => cache_key(&ctx), }; - // The region the level is promoted into, which the executor flushes only - // between evaluations, so bytes copied out of it stay readable for this one. + let mut slot = slot; + let promotion = Promotion::new(ctx.arena(), slot.frames().bounds(), ctx.scope().persistent()); + let persistent = ctx.scope().persistent(); + let finalized = |value: Served<'e>, finality: &Finality| match finality { + Finality::AllFinal => GPoll::Final(value), + Finality::Partial => GPoll::Partial(value), + }; + // The claim is this node's output frame: a hit fills it from the cached + // bytes, and every valueless exit drops it with the frame still claimed. + let serve = |entry: &MemoLevel, mut slot: FrameClaim<'e, 'l>| { + if lane >= entry.lanes.len() { + // The cached level ends here; the past-end signal serves drains. + return GPoll::Error(Box::new(core_types::gpoll::GraphError::past_end())); + } + if let Some(span) = entry.span + && let Some(src) = span.lane(persistent, lane, content.layout()) + { + // SAFETY: the span resolved in generation, so the lane is live and + // immutable at the layout it was promoted under. + unsafe { slot.fill_copy(src) }; + // SAFETY: the copy images a complete record of this layout. + return finalized(unsafe { slot.finish_served() }, &entry.finality); + } + match entry.lanes[lane].replay_into(&mut slot, ctx.arena()) { + // SAFETY: the replay completes the record in the frame. + Some(()) => finalized(unsafe { slot.finish_served() }, &entry.finality), + None => GPoll::arena_exhausted(), + } + }; + if let Some(entry) = cache.lock().unwrap().as_ref() + && entry.key == key + { + return serve(entry, slot); + } + if leveled { + return match content.materialize_level(&ctx, ctx.arena()) { + LevelStatus::Batch(batch, finality) => { + let layout = content.layout(); + // SAFETY: the batch came from this edge, so it carries the edge's layout. + let lanes: Vec = (0..batch.len()).map(|index| unsafe { OwnedRecord::copy_out(layout, batch.get(index).rec()) }).collect(); + let entry = MemoLevel { + key, + // SAFETY: as above. + span: unsafe { MaterializedSpan::to_persistent(&batch, &promotion) }, + lanes, + finality, + }; + let result = serve(&entry, slot); + *cache.lock().unwrap() = Some(entry); + result + } + LevelStatus::Pending => GPoll::Pending, + LevelStatus::Error(error) => GPoll::Error(Box::new(error)), + }; + } + // The output layout is the content's, so the claim is the content's frame. + let result = content.serve(&ctx, slot); + let publishable = match &result { + GPoll::Final(served) => Some((served.record(), Finality::AllFinal)), + GPoll::Partial(served) => Some((served.record(), Finality::Partial)), + GPoll::Pending | GPoll::Fallback(_) | GPoll::Error(_) => None, + }; + if let Some((value, finality)) = publishable { + let layout = content.layout(); + // SAFETY: the value came from this edge, so it carries the edge's + // layout, and one record of it is a batch of one lane. + let batch = unsafe { core_types::node::RecordBatch::new(layout.rec(value).ptr(), 1, layout) }; + // SAFETY: as above. + let copy = unsafe { OwnedRecord::copy_out(layout, layout.rec(value)) }; + *cache.lock().unwrap() = Some(MemoLevel { + key, + // SAFETY: as above. + span: unsafe { MaterializedSpan::to_persistent(&batch, &promotion) }, + lanes: vec![copy], + finality, + }); + } + result +} + +/// The span memo's entry: the span resolves only while the persistent epoch +/// that published it is live, so a flush costs a re-publish and nothing else. +#[derive(Debug)] +pub struct SpanLevel { + key: u64, + span: MaterializedSpan, + finality: Finality, +} + +/// The cheap memo the compiler inserts at context boundaries: a published +/// level lives in the persistent region and serves cross-evaluation hits by +/// byte copy until the next flush, which costs a re-publish rather than a deep +/// copy. +#[node_macro::node(category(""), path(graphene_core::memo))] +fn frame_memo<'e, 'l>( + ctx: impl Ctx + CacheHash + DeriveCtx + ExtractArena<'e> + ModifyIndex + Copy, + #[data] cache: Arc>>, + content: impl Node>, + slot: FrameClaim<'e, 'l>, +) -> GPoll> { + // A scalar wire's value may depend on the consuming lane (index readers), + // so only a leveled wire, whose level covers every lane by construction, + // keys with the lane normalized away. + let leveled = content.layout().depth > 0; + let lane = match leveled { + true => ctx.index() as usize, + false => 0, + }; + let key = match leveled { + true => { + let mut keyed = *ctx; + keyed.set_index(0); + cache_key(&keyed) + } + false => cache_key(&ctx), + }; + let mut slot = slot; + let promotion = Promotion::new(ctx.arena(), slot.frames().bounds(), ctx.scope().persistent()); let persistent = ctx.scope().persistent(); let finalized = |value: Served<'e>, finality: Finality| match finality { Finality::AllFinal => GPoll::Final(value), @@ -60,15 +177,15 @@ fn memoize<'e, 'l>( // The claim is this node's output frame: a hit fills it from the published // bytes, and every valueless exit drops it with the frame still claimed. let serve = |src: *const u8, finality: Finality, mut slot: FrameClaim<'e, 'l>| { - // SAFETY: the source images a complete record of this layout whose - // parked payloads outlive the evaluation. + // SAFETY: the source images a complete record of this layout, and the + // span resolved in generation, so its parked payloads are live. unsafe { slot.fill_copy(src) }; // SAFETY: the copy images a complete record of this layout. finalized(unsafe { slot.finish_served() }, finality) }; let past_end = || GPoll::Error(Box::new(core_types::gpoll::GraphError::past_end())); let entry = cache.lock().unwrap().as_ref().filter(|entry| entry.key == key).map(|entry| (entry.span, entry.finality)); - // A span that no longer resolves was flushed; the miss below re-promotes it. + // A span that no longer resolves was flushed; the miss below re-publishes. if let Some((span, finality)) = entry && let Some(published) = span.batch(persistent, content.layout()) { @@ -82,8 +199,8 @@ fn memoize<'e, 'l>( return match content.materialize_level(&ctx, ctx.arena()) { LevelStatus::Batch(batch, finality) => { // SAFETY: the batch came from this edge, so it carries the edge's layout. - let span = unsafe { MaterializedSpan::promote(&batch, persistent) }; - *cache.lock().unwrap() = span.map(|span| MemoLevel { key, span, finality }); + let span = unsafe { MaterializedSpan::to_persistent(&batch, &promotion) }; + *cache.lock().unwrap() = span.map(|span| SpanLevel { key, span, finality }); match lane < batch.len() { // The publishing evaluation reads the resident batch, not the copy. true => serve(batch.get(lane).rec().ptr(), finality, slot), @@ -105,63 +222,14 @@ fn memoize<'e, 'l>( let layout = content.layout(); // SAFETY: the value came from this edge, so it carries the edge's // layout, and one record of it is a batch of one lane. - let span = unsafe { MaterializedSpan::promote(&core_types::node::RecordBatch::new(layout.rec(value).ptr(), 1, layout), persistent) }; - *cache.lock().unwrap() = span.map(|span| MemoLevel { key, span, finality }); + let batch = unsafe { core_types::node::RecordBatch::new(layout.rec(value).ptr(), 1, layout) }; + // SAFETY: as above. + let span = unsafe { MaterializedSpan::to_persistent(&batch, &promotion) }; + *cache.lock().unwrap() = span.map(|span| SpanLevel { key, span, finality }); } result } -#[node_macro::node(category(""), path(graphene_core::memo))] -fn frame_memo<'e, 'l>( - ctx: impl Ctx + CacheHash + ExtractArena<'e>, - #[data] cell: ArenaCell, 32>>, - content: impl Node>, - frame: FrameClaim<'e, 'l>, -) -> GPoll> { - let arena = ctx.arena(); - let table = match cell.load(arena) { - Some(table) => table, - None => match arena.alloc(FrameTable::new()) { - Some((table, weak)) => { - cell.store(weak); - table - } - None => return content.serve(&ctx, frame), - }, - }; - // SAFETY: published bytes are same-frame copies of this edge's records, - // so they carry the edge's layout with live parked references, and the - // claim is that layout's frame. - let revive = |mut frame: FrameClaim<'e, 'l>, bytes: &Box<[u8]>| unsafe { - frame.fill_copy(bytes.as_ptr()); - frame.finish_served() - }; - match table.lookup(cache_key(ctx)) { - Lookup::Hit(Finality::AllFinal, bytes) => GPoll::Final(revive(frame, bytes)), - Lookup::Hit(Finality::Partial, bytes) => GPoll::Partial(revive(frame, bytes)), - Lookup::Vacant(slot) => match content.serve(&ctx, frame) { - GPoll::Final(served) => { - // SAFETY: the value came from this edge, so it carries the edge's layout. - // The serve's own frame answers this pull; the publish feeds later ones. - let bytes = unsafe { copy_record_bytes(content.layout(), content.layout().rec(served.record())) }; - slot.publish(bytes, Finality::AllFinal); - GPoll::Final(served) - } - GPoll::Partial(served) => { - // SAFETY: as above. - let bytes = unsafe { copy_record_bytes(content.layout(), content.layout().rec(served.record())) }; - slot.publish(bytes, Finality::Partial); - GPoll::Partial(served) - } - unpublishable => { - slot.release(); - unpublishable - } - }, - Lookup::Full => content.serve(&ctx, frame), - } -} - type MonitorValue = Arc>>; /// The Monitor node is used by the editor to access the data flowing through @@ -298,8 +366,9 @@ mod tests { core_types::record::register_deep_element_clone::(deep, deep_repark); let arena = Arena::new(4096).unwrap(); + let persistent = Arena::new(4096).unwrap(); let generations = []; - let scope = scope_fixture(&generations, &arena); + let scope = scope_fixture(&generations, &arena).with_persistent(&persistent); let ctx = ContextImpl::root(&scope); let layout = element_layout::(); @@ -310,7 +379,39 @@ mod tests { assert_eq!( memoized.eval(&ctx, &frames), GPoll::Final(Payload("deep".to_string(), 2)), - "the hit replays through both halves of the deep glue" + "promoting across regions runs both halves of the deep glue" + ); + } + + #[test] + fn a_promote_within_one_region_shares_instead_of_cloning() { + let frames = core_types::record::test_frames(1 << 16); + #[derive(Clone, Debug, PartialEq, dyn_any::DynAny)] + struct Shared(String, u32); + unsafe fn deep(ptr: *const u8) -> Box { + let value = unsafe { core_types::record::borrow_element::(core_types::record::Rec::new(ptr)) }; + Box::new(Shared(value.0.clone(), value.1 + 1)) + } + unsafe fn deep_repark(value: &(dyn std::any::Any + Send + Sync), dst: *mut u8, arena: &Arena) -> Option<()> { + let value = value.downcast_ref::().expect("an element replays at its own type"); + unsafe { core_types::record::write_element(dst, Shared(value.0.clone(), value.1 + 1), arena) } + } + core_types::record::register_deep_element_clone::(deep, deep_repark); + + let arena = Arena::new(4096).unwrap(); + let generations = []; + let scope = scope_fixture(&generations, &arena); + let ctx = ContextImpl::root(&scope); + + let layout = element_layout::(); + let memoized = MemoizeNode::new(lifted::(Shared("shared".to_string(), 0)), &layout); + let memoized = core_types::record::RecordExtract::::new(memoized, &layout); + + assert_eq!(memoized.eval(&ctx, &frames), GPoll::Final(Shared("shared".to_string(), 0)), "the miss serves the live value"); + assert_eq!( + memoized.eval(&ctx, &frames), + GPoll::Final(Shared("shared".to_string(), 0)), + "a payload already living as long as the span is shared, so no glue runs" ); } @@ -411,7 +512,29 @@ mod tests { } #[test] - fn a_flush_invalidates_every_persistent_span() { + fn a_flush_invalidates_every_span_memo_entry() { + let frames = core_types::record::test_frames(1 << 16); + let arena = Arena::new(4096).unwrap(); + let mut persistent = Arena::new(4096).unwrap(); + let generations = []; + + let layout = element_layout::(); + let memoized = FrameMemoNode::new(counting(), &layout); + let memoized = core_types::record::RecordExtract::::new(memoized, &layout); + let eval = |persistent: &Arena| { + let scope = scope_fixture(&generations, &arena).with_persistent(persistent); + memoized.eval(&ContextImpl::root(&scope), &frames) + }; + + assert_eq!(eval(&persistent), GPoll::Final(1)); + assert_eq!(eval(&persistent), GPoll::Final(1), "the published level serves the hit"); + persistent.reset(); + assert_eq!(eval(&persistent), GPoll::Final(2), "the flush invalidates the span"); + assert_eq!(eval(&persistent), GPoll::Final(2), "the miss re-published the level"); + } + + #[test] + fn the_owned_tier_survives_a_flush_without_recomputing() { let frames = core_types::record::test_frames(1 << 16); let arena = Arena::new(4096).unwrap(); let mut persistent = Arena::new(4096).unwrap(); @@ -420,16 +543,17 @@ mod tests { let layout = element_layout::(); let memoized = MemoizeNode::new(counting(), &layout); let memoized = core_types::record::RecordExtract::::new(memoized, &layout); - let eval = |persistent: &Arena| { - let scope = scope_fixture(&generations, &arena).with_persistent(persistent); + let eval = |persistent: &Arena, generations: &[(SourceId, u64)]| { + let scope = scope_fixture(generations, &arena).with_persistent(persistent); memoized.eval(&ContextImpl::root(&scope), &frames) }; - assert_eq!(eval(&persistent), GPoll::Final(1)); - assert_eq!(eval(&persistent), GPoll::Final(1), "the promoted level serves the hit"); + assert_eq!(eval(&persistent, &generations), GPoll::Final(1)); + assert_eq!(eval(&persistent, &generations), GPoll::Final(1), "the promoted level serves the hit"); persistent.reset(); - assert_eq!(eval(&persistent), GPoll::Final(2), "the flush invalidates the span"); - assert_eq!(eval(&persistent), GPoll::Final(2), "the miss re-promoted the level"); + assert_eq!(eval(&persistent, &generations), GPoll::Final(1), "the owned copies replay across the flush"); + let bumped = [(7 as SourceId, 3)]; + assert_eq!(eval(&persistent, &bumped), GPoll::Final(2), "only a key change recomputes"); } #[test] @@ -441,7 +565,7 @@ mod tests { let generations = []; let layout = element_layout::(); - let memoized = MemoizeNode::new(counting(), &layout); + let memoized = FrameMemoNode::new(counting(), &layout); let memoized = core_types::record::RecordExtract::::new(memoized, &layout); let eval = |persistent: &Arena| { let scope = scope_fixture(&generations, &arena).with_persistent(persistent); @@ -462,12 +586,30 @@ mod tests { let scope = scope_fixture(&generations, &arena).with_persistent(&persistent); let ctx = ContextImpl::root(&scope); + let layout = element_layout::(); + let memoized = FrameMemoNode::new(counting(), &layout); + let memoized = core_types::record::RecordExtract::::new(memoized, &layout); + + assert_eq!(memoized.eval(&ctx, &frames), GPoll::Final(1)); + assert_eq!(memoized.eval(&ctx, &frames), GPoll::Final(2), "an unpublished level recomputes"); + assert!(persistent.exhausted(), "the refused promote marks the region for a flush"); + } + + #[test] + fn a_refused_promote_leaves_the_owned_tier_serving() { + let frames = core_types::record::test_frames(1 << 16); + let arena = Arena::new(4096).unwrap(); + let persistent = Arena::new(0).unwrap(); + let generations = []; + let scope = scope_fixture(&generations, &arena).with_persistent(&persistent); + let ctx = ContextImpl::root(&scope); + let layout = element_layout::(); let memoized = MemoizeNode::new(counting(), &layout); let memoized = core_types::record::RecordExtract::::new(memoized, &layout); assert_eq!(memoized.eval(&ctx, &frames), GPoll::Final(1)); - assert_eq!(memoized.eval(&ctx, &frames), GPoll::Final(2), "an unpromoted level recomputes"); + assert_eq!(memoized.eval(&ctx, &frames), GPoll::Final(1), "the owned tier answers where the promote was refused"); assert!(persistent.exhausted(), "the refused promote marks the region for a flush"); } @@ -514,7 +656,7 @@ mod tests { let memo = FrameMemoNode::new(lifted::("lent out".to_string()), &layout); let GPoll::Final(first) = core_types::record::serve_edge(&memo, &ctx, &frames) else { - panic!("the miss must fill the frame table"); + panic!("the miss must publish the record"); }; let GPoll::Final(second) = core_types::record::serve_edge(&memo, &ctx, &frames) else { panic!("the hit must revive the published record"); @@ -524,4 +666,32 @@ mod tests { assert_eq!(first, "lent out"); assert!(std::ptr::eq(first, second), "the hit shares the parked payload"); } + + #[test] + fn a_span_memo_hit_crosses_evaluations_on_one_payload() { + let frames = core_types::record::test_frames(1 << 16); + let mut arena = Arena::new(4096).unwrap(); + let persistent = Arena::new(4096).unwrap(); + let generations = []; + + let layout = element_layout::(); + let memo = FrameMemoNode::new(lifted::("published".to_string()), &layout); + let served_at = |arena: &Arena| { + let scope = scope_fixture(&generations, arena).with_persistent(&persistent); + let ctx = ContextImpl::root(&scope); + let GPoll::Final(value) = core_types::record::serve_edge(&memo, &ctx, &frames) else { + panic!("the span memo must serve a final record"); + }; + let element: &String = unsafe { core_types::record::borrow_element(layout.rec(&value)) }; + assert_eq!(element, "published"); + std::ptr::from_ref(element) + }; + + served_at(&arena); + let first = served_at(&arena); + arena.reset(); + let second = served_at(&arena); + assert_eq!(first, second, "the hit names the published payload rather than re-parking it"); + assert!(persistent.contains(first.cast::()), "the payload the hits share lives in the persistent region"); + } }