1use std::cell::{Cell, RefCell};
8use std::collections::{BTreeMap, HashMap, HashSet, VecDeque};
9use std::fmt;
10use std::panic::{catch_unwind, resume_unwind, AssertUnwindSafe};
11use std::rc::Rc;
12
13use crate::batch::{batch as run_batch, BatchCtx};
14use crate::checkpoint::GraphCheckpointJson;
15use crate::ctx::{Ctx, DepRecord, DepTerminal, WaveData};
16use crate::dispatcher::{default_dispatcher, Dispatcher, NodeFn};
17use crate::environment::EnvironmentDrivers;
18use crate::node::{Core, GraphArena, Node, NodeOpts, Status};
19use crate::operators::Operator;
20use crate::protocol::{AnyValue, LockId, Message, PullDemand, Tier};
21use crate::versioning::{NodeVersion, NodeVersioningPolicy};
22
23#[derive(Clone, Default)]
25pub struct GraphOptions {
26 pub name: Option<String>,
28 pub profile: bool,
30 pub dispatcher: Option<Dispatcher>,
32 pub environment: EnvironmentDrivers,
34 pub versioning: Option<NodeVersioningPolicy>,
36}
37
38impl fmt::Debug for GraphOptions {
39 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
40 f.debug_struct("GraphOptions")
41 .field("name", &self.name)
42 .field("profile", &self.profile)
43 .field("dispatcher", &self.dispatcher)
44 .field("environment", &self.environment)
45 .field("versioning", &self.versioning)
46 .finish()
47 }
48}
49
50impl GraphOptions {
51 pub fn named(name: impl Into<String>) -> Self {
53 Self {
54 name: Some(name.into()),
55 profile: false,
56 dispatcher: None,
57 environment: EnvironmentDrivers::default(),
58 versioning: None,
59 }
60 }
61}
62
63#[derive(Debug, Clone, Default)]
66pub struct GraphNodeOpts {
67 pub name: Option<String>,
69 pub meta: BTreeMap<String, String>,
71 pub node: NodeOpts,
73 pub restore: Option<RestoreFactoryMeta>,
75}
76
77impl GraphNodeOpts {
78 pub fn named(name: impl Into<String>) -> Self {
80 Self {
81 name: Some(name.into()),
82 ..Self::default()
83 }
84 }
85}
86
87#[derive(Debug, Clone, Default)]
89pub struct TopologyGroupOptions {
90 pub name: Option<String>,
92}
93
94impl TopologyGroupOptions {
95 pub fn named(name: impl Into<String>) -> Self {
97 Self {
98 name: Some(name.into()),
99 }
100 }
101}
102
103#[derive(Debug, Clone, PartialEq)]
104pub struct RestoreFactoryMeta {
106 pub ref_: String,
108 pub config: Option<GraphCheckpointJson>,
110 pub config_version: Option<GraphCheckpointJson>,
112}
113
114impl RestoreFactoryMeta {
115 pub fn registry_ref(ref_: impl Into<String>) -> Self {
117 Self {
118 ref_: ref_.into(),
119 config: None,
120 config_version: None,
121 }
122 }
123
124 pub fn with_config(mut self, config: GraphCheckpointJson) -> Self {
126 self.config = Some(config);
127 self
128 }
129
130 pub fn with_config_version(mut self, config_version: GraphCheckpointJson) -> Self {
132 self.config_version = Some(config_version);
133 self
134 }
135}
136
137#[derive(Clone)]
138struct Entry {
139 core: Core,
140 id: String,
141 name: Option<String>,
142 factory: String,
143 meta: BTreeMap<String, String>,
144 restore: Option<RestoreFactoryMeta>,
145}
146
147struct Mount {
148 at: String,
149 graph: Graph,
150 topology_observer: RefCell<Option<GraphTopologyObserver>>,
151}
152
153#[derive(Clone)]
154struct TopologyObserverEntry {
155 path: Option<String>,
156 sink: Rc<dyn Fn(TopologyEvent)>,
157}
158
159struct GraphInner {
160 name: Option<String>,
161 arena: GraphArena,
162 dispatcher: Dispatcher,
163 environment: EnvironmentDrivers,
164 versioning: Option<NodeVersioningPolicy>,
165 profile_enabled: Cell<bool>,
166 entries: RefCell<Vec<Entry>>,
167 by_id: RefCell<HashMap<String, Core>>,
168 retired_ids: RefCell<HashSet<String>>,
169 mounts: RefCell<Vec<Mount>>,
170 seq: Cell<usize>,
171 synth_seq: Cell<usize>,
172 synth_ids: RefCell<HashMap<(usize, usize, u64), String>>,
173 clock: Rc<Cell<u64>>,
174 topology_observers: RefCell<Vec<Option<TopologyObserverEntry>>>,
175 topology_observer_active: Cell<usize>,
176 topology_observer_free: RefCell<Vec<usize>>,
177 topology_delivering: Cell<bool>,
178 topology_queue: RefCell<VecDeque<TopologyEvent>>,
179}
180
181#[derive(Clone)]
183pub struct Graph {
184 inner: Rc<GraphInner>,
185}
186
187struct TopologyGroupInner {
188 graph: Graph,
189 name: Option<String>,
190 members: RefCell<Vec<Core>>,
191 released: Cell<bool>,
192}
193
194#[derive(Clone)]
199pub struct TopologyGroup {
200 inner: Rc<TopologyGroupInner>,
201}
202
203impl TopologyGroup {
204 fn new(graph: Graph, opts: TopologyGroupOptions) -> Self {
205 Self {
206 inner: Rc::new(TopologyGroupInner {
207 graph,
208 name: opts.name,
209 members: RefCell::new(Vec::new()),
210 released: Cell::new(false),
211 }),
212 }
213 }
214
215 pub fn name(&self) -> Option<&str> {
217 self.inner.name.as_deref()
218 }
219
220 pub fn is_released(&self) -> bool {
222 self.inner.released.get()
223 }
224
225 pub fn add<T: 'static>(&self, node: &Node<T>) -> Node<T> {
227 self.assert_live();
228 self.inner.graph.assert_registered_core(
229 &node.erased(),
230 &format!("topology group '{}' member", self.name().unwrap_or("group")),
231 );
232 self.track(node.erased());
233 node.clone()
234 }
235
236 pub fn node<T: 'static, F: Fn(&Ctx) + 'static>(&self, deps: Vec<Core>, f: F) -> Node<T> {
238 self.node_opts(deps, f, GraphNodeOpts::default())
239 }
240
241 pub fn node_opts<T: 'static, F: Fn(&Ctx) + 'static>(
243 &self,
244 deps: Vec<Core>,
245 f: F,
246 opts: GraphNodeOpts,
247 ) -> Node<T> {
248 self.assert_live();
249 self.track_node(self.inner.graph.node_opts(deps, f, opts))
250 }
251
252 pub fn state<T: 'static>(&self, initial: T) -> Node<T> {
254 self.state_opts(initial, GraphNodeOpts::default())
255 }
256
257 pub fn state_opts<T: 'static>(&self, initial: T, opts: GraphNodeOpts) -> Node<T> {
259 self.assert_live();
260 self.track_node(self.inner.graph.state_opts(initial, opts))
261 }
262
263 pub fn state_empty<T: 'static>(&self) -> Node<T> {
265 self.state_empty_opts(GraphNodeOpts::default())
266 }
267
268 pub fn state_empty_opts<T: 'static>(&self, opts: GraphNodeOpts) -> Node<T> {
270 self.assert_live();
271 self.track_node(self.inner.graph.state_empty_opts(opts))
272 }
273
274 pub fn producer<T: 'static, F: Fn(&Ctx) + 'static>(&self, f: F) -> Node<T> {
276 self.producer_opts(f, GraphNodeOpts::default())
277 }
278
279 pub fn producer_opts<T: 'static, F: Fn(&Ctx) + 'static>(
281 &self,
282 f: F,
283 opts: GraphNodeOpts,
284 ) -> Node<T> {
285 self.assert_live();
286 self.track_node(self.inner.graph.producer_opts(f, opts))
287 }
288
289 pub fn derived<T: 'static, F: Fn(&Values<'_>) -> Option<T> + 'static>(
291 &self,
292 deps: Vec<Core>,
293 f: F,
294 ) -> Node<T> {
295 self.derived_opts(deps, f, GraphNodeOpts::default())
296 }
297
298 pub fn derived_opts<T: 'static, F: Fn(&Values<'_>) -> Option<T> + 'static>(
300 &self,
301 deps: Vec<Core>,
302 f: F,
303 opts: GraphNodeOpts,
304 ) -> Node<T> {
305 self.assert_live();
306 self.track_node(self.inner.graph.derived_opts(deps, f, opts))
307 }
308
309 pub fn effect<F: Fn(&Values<'_>) -> Option<Box<dyn FnOnce()>> + 'static>(
311 &self,
312 deps: Vec<Core>,
313 f: F,
314 ) -> Node<()> {
315 self.effect_opts(deps, f, GraphNodeOpts::default())
316 }
317
318 pub fn effect_opts<F: Fn(&Values<'_>) -> Option<Box<dyn FnOnce()>> + 'static>(
320 &self,
321 deps: Vec<Core>,
322 f: F,
323 opts: GraphNodeOpts,
324 ) -> Node<()> {
325 self.assert_live();
326 self.track_node(self.inner.graph.effect_opts(deps, f, opts))
327 }
328
329 pub fn init_node<T: 'static>(
331 &self,
332 op: Operator<T>,
333 deps: Vec<Core>,
334 opts: GraphNodeOpts,
335 ) -> Node<T> {
336 self.assert_live();
337 self.track_node(self.inner.graph.init_node(op, deps, opts))
338 }
339
340 pub fn release(&self) {
342 let reason = self
343 .inner
344 .name
345 .clone()
346 .unwrap_or_else(|| "topology group".to_owned());
347 self.release_with_reason(&reason);
348 }
349
350 pub fn release_with_reason(&self, reason: &str) {
352 if self.inner.released.get() {
353 return;
354 }
355 let members = self.inner.members.borrow().clone();
356 self.inner.graph.release_nodes(&members, reason);
357 self.inner.members.borrow_mut().clear();
358 self.inner.released.set(true);
359 }
360
361 fn track_node<T: 'static>(&self, node: Node<T>) -> Node<T> {
362 self.assert_live();
363 self.track(node.erased());
364 node
365 }
366
367 fn track(&self, core: Core) {
368 let mut members = self.inner.members.borrow_mut();
369 if !members.iter().any(|member| member.ptr_eq(&core)) {
370 members.push(core);
371 }
372 }
373
374 fn assert_live(&self) {
375 assert!(
376 !self.inner.released.get(),
377 "topology group '{}' has been released (D152)",
378 self.name().unwrap_or("group")
379 );
380 }
381}
382
383impl Graph {
384 pub fn new(opts: GraphOptions) -> Self {
386 let dispatcher = opts.dispatcher.unwrap_or_else(default_dispatcher);
387 let mut environment = opts.environment;
388 if environment.local_async_driver().is_none() {
389 environment.set_local_async_driver(dispatcher.local_async_driver());
390 }
391 if opts.profile {
392 dispatcher.set_recording(true);
393 }
394 Self {
395 inner: Rc::new(GraphInner {
396 name: opts.name,
397 arena: GraphArena::new(),
398 dispatcher,
399 environment,
400 versioning: opts.versioning,
401 profile_enabled: Cell::new(opts.profile),
402 entries: RefCell::new(Vec::new()),
403 by_id: RefCell::new(HashMap::new()),
404 retired_ids: RefCell::new(HashSet::new()),
405 mounts: RefCell::new(Vec::new()),
406 seq: Cell::new(0),
407 synth_seq: Cell::new(0),
408 synth_ids: RefCell::new(HashMap::new()),
409 clock: Rc::new(Cell::new(0)),
410 topology_observers: RefCell::new(Vec::new()),
411 topology_observer_active: Cell::new(0),
412 topology_observer_free: RefCell::new(Vec::new()),
413 topology_delivering: Cell::new(false),
414 topology_queue: RefCell::new(VecDeque::new()),
415 }),
416 }
417 }
418
419 pub fn name(&self) -> Option<&str> {
421 self.inner.name.as_deref()
422 }
423
424 pub fn topology_group(&self) -> TopologyGroup {
427 self.topology_group_opts(TopologyGroupOptions::default())
428 }
429
430 pub fn topology_group_opts(&self, opts: TopologyGroupOptions) -> TopologyGroup {
432 TopologyGroup::new(self.clone(), opts)
433 }
434
435 pub fn find(&self, id: &str) -> Option<GraphNode> {
437 self.inner
438 .by_id
439 .borrow()
440 .get(id)
441 .cloned()
442 .map(GraphNode::new)
443 }
444
445 pub fn node<T: 'static, F: Fn(&Ctx) + 'static>(&self, deps: Vec<Core>, f: F) -> Node<T> {
447 self.node_opts(deps, f, GraphNodeOpts::default())
448 }
449
450 pub fn node_opts<T: 'static, F: Fn(&Ctx) + 'static>(
452 &self,
453 deps: Vec<Core>,
454 f: F,
455 opts: GraphNodeOpts,
456 ) -> Node<T> {
457 self.node_opts_initial(deps, f, opts, None)
458 }
459
460 pub(crate) fn node_opts_initial<T: 'static, F: Fn(&Ctx) + 'static>(
461 &self,
462 deps: Vec<Core>,
463 f: F,
464 opts: GraphNodeOpts,
465 initial: Option<T>,
466 ) -> Node<T> {
467 self.assert_graph_local_deps(&deps, opts.name.as_deref().unwrap_or("node"));
468 let node = Node::derived_opts_initial_in_arena_with_dispatcher(
469 &self.inner.arena,
470 self.inner.dispatcher.clone(),
471 deps,
472 self.apply_default_node_opts(opts.node.clone()),
473 initial,
474 f,
475 );
476 self.add(node, "node", opts)
477 }
478
479 pub fn state<T: 'static>(&self, initial: T) -> Node<T> {
481 self.state_opts(initial, GraphNodeOpts::default())
482 }
483
484 pub fn state_opts<T: 'static>(&self, initial: T, opts: GraphNodeOpts) -> Node<T> {
486 let node = Node::state_opts_in_arena_with_dispatcher(
487 &self.inner.arena,
488 self.inner.dispatcher.clone(),
489 initial,
490 self.apply_default_node_opts(opts.node.clone()),
491 );
492 self.add(node, "state", opts)
493 }
494
495 pub fn state_empty<T: 'static>(&self) -> Node<T> {
497 self.state_empty_opts(GraphNodeOpts::default())
498 }
499
500 pub fn state_empty_opts<T: 'static>(&self, opts: GraphNodeOpts) -> Node<T> {
502 let node = Node::state_empty_opts_in_arena_with_dispatcher(
503 &self.inner.arena,
504 self.inner.dispatcher.clone(),
505 self.apply_default_node_opts(opts.node.clone()),
506 );
507 self.add(node, "state", opts)
508 }
509
510 pub fn producer<T: 'static, F: Fn(&Ctx) + 'static>(&self, f: F) -> Node<T> {
512 self.producer_opts(f, GraphNodeOpts::default())
513 }
514
515 pub fn producer_opts<T: 'static, F: Fn(&Ctx) + 'static>(
517 &self,
518 f: F,
519 opts: GraphNodeOpts,
520 ) -> Node<T> {
521 let node = Node::producer_opts_in_arena_with_dispatcher(
522 &self.inner.arena,
523 self.inner.dispatcher.clone(),
524 self.apply_default_node_opts(opts.node.clone()),
525 f,
526 );
527 self.add(node, "producer", opts)
528 }
529
530 pub fn derived<T: 'static, F: Fn(&Values<'_>) -> Option<T> + 'static>(
533 &self,
534 deps: Vec<Core>,
535 f: F,
536 ) -> Node<T> {
537 self.derived_opts(deps, f, GraphNodeOpts::default())
538 }
539
540 pub fn derived_opts<T: 'static, F: Fn(&Values<'_>) -> Option<T> + 'static>(
542 &self,
543 deps: Vec<Core>,
544 f: F,
545 opts: GraphNodeOpts,
546 ) -> Node<T> {
547 self.assert_graph_local_deps(&deps, opts.name.as_deref().unwrap_or("derived"));
548 let node = Node::derived_opts_in_arena_with_dispatcher(
549 &self.inner.arena,
550 self.inner.dispatcher.clone(),
551 deps,
552 self.apply_default_node_opts(opts.node.clone()),
553 move |ctx| {
554 let values = Values::new(ctx.dep_records());
555 if let Some(value) = f(&values) {
556 ctx.emit(value);
557 }
558 },
559 );
560 self.add(node, "derived", opts)
561 }
562
563 pub fn effect<F: Fn(&Values<'_>) -> Option<Box<dyn FnOnce()>> + 'static>(
565 &self,
566 deps: Vec<Core>,
567 f: F,
568 ) -> Node<()> {
569 self.effect_opts(deps, f, GraphNodeOpts::default())
570 }
571
572 pub fn effect_opts<F: Fn(&Values<'_>) -> Option<Box<dyn FnOnce()>> + 'static>(
574 &self,
575 deps: Vec<Core>,
576 f: F,
577 opts: GraphNodeOpts,
578 ) -> Node<()> {
579 self.assert_graph_local_deps(&deps, opts.name.as_deref().unwrap_or("effect"));
580 let node = Node::derived_opts_in_arena_with_dispatcher(
581 &self.inner.arena,
582 self.inner.dispatcher.clone(),
583 deps,
584 self.apply_default_node_opts(opts.node.clone()),
585 move |ctx| {
586 let values = Values::new(ctx.dep_records());
587 if let Some(cleanup) = f(&values) {
588 ctx.on_deactivation(cleanup);
589 }
590 },
591 );
592 self.add(node, "effect", opts)
593 }
594
595 pub fn batch<R>(&self, f: impl FnOnce(&BatchCtx) -> R) -> R {
597 run_batch(f)
598 }
599
600 pub fn init_node<T: 'static>(
602 &self,
603 op: Operator<T>,
604 deps: Vec<Core>,
605 opts: GraphNodeOpts,
606 ) -> Node<T> {
607 self.assert_graph_local_deps(&deps, opts.name.as_deref().unwrap_or(op.factory));
608 let node = crate::operators::init_node_in_arena_with_dispatcher(
609 op.clone(),
610 &self.inner.arena,
611 self.inner.dispatcher.clone(),
612 deps,
613 self.apply_default_node_opts(opts.node.clone()),
614 );
615 self.add(node, op.factory, opts)
616 }
617
618 pub fn mount(&self, child: Graph, at: impl Into<String>) {
620 let at = at.into();
621 assert!(!at.is_empty(), "mount path must not be empty");
622 assert!(
623 !child.contains_graph(self),
624 "mount would create a graph cycle"
625 );
626 assert!(
627 !self.inner.mounts.borrow().iter().any(|m| m.at == at),
628 "duplicate mount path '{at}'"
629 );
630 if self.inner.profile_enabled.get() {
631 child.enable_profile_recorder_recursive();
632 }
633 self.inner.mounts.borrow_mut().push(Mount {
634 at: at.clone(),
635 graph: child,
636 topology_observer: RefCell::new(None),
637 });
638 if self.has_topology_observers() {
639 self.ensure_mounted_topology_forwarders();
640 }
641 self.emit_topology_mount_changed(at);
642 }
643
644 pub fn describe(&self) -> DescribeSnapshot {
646 self.describe_opts(DescribeOpts::default())
647 }
648
649 pub fn describe_opts(&self, opts: DescribeOpts) -> DescribeSnapshot {
651 let snap = self.describe_with_prefix("");
652 if let Some(explain) = opts.explain {
653 explain_subset(snap, &explain)
654 } else {
655 snap
656 }
657 }
658
659 pub fn observe(&self) -> ObserveStream {
661 ObserveStream {
662 graph: self.clone(),
663 path: None,
664 }
665 }
666
667 pub fn observe_path(&self, path: &str) -> ObserveStream {
669 ObserveStream {
670 graph: self.clone(),
671 path: Some(path.to_owned()),
672 }
673 }
674
675 pub fn observe_topology(&self) -> TopologyStream {
679 TopologyStream {
680 graph: self.clone(),
681 path: None,
682 }
683 }
684
685 pub fn observe_topology_path(&self, path: &str) -> TopologyStream {
687 TopologyStream {
688 graph: self.clone(),
689 path: Some(path.to_owned()),
690 }
691 }
692
693 pub fn profile(&self) -> Profile {
695 let mut total_invokes = 0;
696 let mut nodes = BTreeMap::new();
697 self.collect_profile("", &mut nodes, &mut total_invokes);
698 Profile {
699 total_invokes,
700 nodes,
701 }
702 }
703
704 fn collect_profile(
705 &self,
706 prefix: &str,
707 nodes: &mut BTreeMap<String, NodeProfile>,
708 total_invokes: &mut u64,
709 ) {
710 for entry in self.inner.entries.borrow().iter() {
711 let stat = entry
712 .core
713 .handle()
714 .and_then(|handle| self.inner.dispatcher.stat_for(handle));
715 let invokes = stat.map_or(0, |s| s.invokes);
716 *total_invokes += invokes;
717 nodes.insert(
718 format!("{prefix}{}", entry.id),
719 NodeProfile {
720 invokes,
721 total_duration_ns: stat.map_or(0, |s| s.total_duration_ns),
722 last_duration_ns: stat.map_or(0, |s| s.last_duration_ns),
723 status: entry.core.status(),
724 },
725 );
726 }
727 for mount in self.inner.mounts.borrow().iter() {
728 mount
729 .graph
730 .collect_profile(&format!("{prefix}{}::", mount.at), nodes, total_invokes);
731 }
732 }
733
734 fn enable_profile_recorder_recursive(&self) {
735 self.inner.dispatcher.set_recording(true);
736 self.inner.profile_enabled.set(true);
737 for mount in self.inner.mounts.borrow().iter() {
738 mount.graph.enable_profile_recorder_recursive();
739 }
740 }
741
742 pub(crate) fn checkpoint_entries(&self) -> Vec<crate::checkpoint::CheckpointEntry> {
743 self.inner
744 .entries
745 .borrow()
746 .iter()
747 .map(|entry| crate::checkpoint::CheckpointEntry {
748 id: entry.id.clone(),
749 name: entry.name.clone(),
750 factory: entry.factory.clone(),
751 meta: entry.meta.clone(),
752 restore: entry.restore.clone(),
753 core: entry.core.clone(),
754 unregistered: false,
755 })
756 .collect()
757 }
758
759 pub(crate) fn checkpoint_mounts(&self) -> Vec<(String, Graph)> {
760 self.inner
761 .mounts
762 .borrow()
763 .iter()
764 .map(|mount| (mount.at.clone(), mount.graph.clone()))
765 .collect()
766 }
767
768 pub(crate) fn checkpoint_synthetic_id_for_core(&self, core: &Core) -> String {
769 self.synthetic_id_for_core(core, "")
770 }
771
772 pub(crate) fn contains_core(&self, core: &Core) -> bool {
773 core.same_graph_arena(&self.inner.arena)
774 && self
775 .inner
776 .entries
777 .borrow()
778 .iter()
779 .any(|entry| entry.core.ptr_eq(core))
780 }
781
782 pub(crate) fn restore_state_json_with_id(
783 &self,
784 id: String,
785 opts: GraphNodeOpts,
786 ) -> Node<GraphCheckpointJson> {
787 let node = Node::state_empty_opts_in_arena_with_dispatcher(
788 &self.inner.arena,
789 self.inner.dispatcher.clone(),
790 self.apply_default_node_opts(opts.node.clone()),
791 );
792 self.add_with_id(node, "state", id, opts)
793 }
794
795 pub(crate) fn restore_node_with_id(
796 &self,
797 id: String,
798 factory: String,
799 deps: Vec<Core>,
800 f: NodeFn,
801 opts: GraphNodeOpts,
802 ) -> Node<GraphCheckpointJson> {
803 self.assert_graph_local_deps(&deps, &id);
804 let node = Node::derived_opts_in_arena_with_dispatcher(
805 &self.inner.arena,
806 self.inner.dispatcher.clone(),
807 deps,
808 self.apply_default_node_opts(opts.node.clone()),
809 move |ctx| f(ctx),
810 );
811 self.add_with_id(node, &factory, id, opts)
812 }
813
814 pub(crate) fn mount_restored(&self, child: Graph, at: String) {
815 self.mount(child, at);
816 }
817
818 fn add<T: 'static>(&self, node: Node<T>, factory: &str, opts: GraphNodeOpts) -> Node<T> {
819 let id = opts.name.clone().unwrap_or_else(|| {
820 let next = self.inner.seq.get();
821 self.inner.seq.set(next + 1);
822 format!("{factory}#{next}")
823 });
824 self.add_with_id(node, factory, id, opts)
825 }
826
827 fn apply_default_node_opts(&self, mut opts: NodeOpts) -> NodeOpts {
828 if opts.versioning.is_none() {
829 opts.versioning = self.inner.versioning.clone();
830 }
831 opts
832 }
833
834 pub(crate) fn empty_source<T: 'static>(&self, factory: &str, opts: GraphNodeOpts) -> Node<T> {
835 let node = Node::state_empty_opts_in_arena_with_dispatcher(
836 &self.inner.arena,
837 self.inner.dispatcher.clone(),
838 self.apply_default_node_opts(opts.node.clone()),
839 );
840 self.add(node, factory, opts)
841 }
842
843 fn add_with_id<T: 'static>(
844 &self,
845 node: Node<T>,
846 factory: &str,
847 id: String,
848 opts: GraphNodeOpts,
849 ) -> Node<T> {
850 assert!(
851 node.erased().same_graph_arena(&self.inner.arena),
852 "graph node belongs to a different graph arena"
853 );
854 let core = node.erased();
855 core.set_environment(self.inner.environment.clone());
856 assert!(
857 !self.inner.by_id.borrow().contains_key(&id),
858 "duplicate graph node id '{id}'"
859 );
860 assert!(
861 !self.inner.retired_ids.borrow().contains(&id),
862 "graph node id '{id}' has been released from this graph lifecycle (D122)"
863 );
864 let entry = Entry {
865 core: core.clone(),
866 id: id.clone(),
867 name: opts.name,
868 factory: factory.to_owned(),
869 meta: opts.meta,
870 restore: opts.restore,
871 };
872 let weak_inner = Rc::downgrade(&self.inner);
873 let topology_id = id.clone();
874 core.set_topology_deps_changed_observer(Rc::new(move |old_deps, new_deps| {
875 if let Some(inner) = weak_inner.upgrade() {
876 Graph { inner }.emit_topology_deps_changed(&topology_id, old_deps, new_deps);
877 }
878 }));
879 self.inner.by_id.borrow_mut().insert(id.clone(), core);
880 self.inner.entries.borrow_mut().push(entry);
881 if let Some(n) = id
882 .rsplit_once('#')
883 .and_then(|(_, n)| n.parse::<usize>().ok())
884 {
885 self.inner.seq.set(self.inner.seq.get().max(n + 1));
886 }
887 self.emit_topology_node_registered(&id);
888 node
889 }
890
891 pub(crate) fn release_nodes(&self, nodes: &[Core], reason: &str) {
892 let mut release_entries: Vec<(String, Core)> = Vec::new();
893 for node in nodes {
894 if release_entries.iter().any(|(_, core)| core.ptr_eq(node)) {
895 continue;
896 }
897 let found = self
898 .inner
899 .entries
900 .borrow()
901 .iter()
902 .find(|entry| entry.core.ptr_eq(node))
903 .map(|entry| (entry.id.clone(), entry.core.clone()));
904 if let Some(entry) = found {
905 release_entries.push(entry);
906 }
907 }
908 if release_entries.is_empty() {
909 return;
910 }
911 let release_ids = release_entries
912 .iter()
913 .map(|(id, _)| id.clone())
914 .collect::<HashSet<_>>();
915 let release_cores = release_entries
916 .iter()
917 .map(|(_, core)| core.clone())
918 .collect::<Vec<_>>();
919 for entry in self.inner.entries.borrow().iter() {
920 if release_ids.contains(&entry.id) {
921 continue;
922 }
923 for dep in entry.core.deps() {
924 if let Some((dep_id, _)) = release_entries
925 .iter()
926 .find(|(_, released)| released.ptr_eq(&dep))
927 {
928 panic!(
929 "graph: cannot release node group for {reason}; '{}' still depends on '{}' (D122)",
930 entry.id, dep_id
931 );
932 }
933 }
934 }
935 for (id, core) in &release_entries {
936 assert!(
937 core.runtime_is_quiescent_for_release(),
938 "graph: cannot release node group for {reason}; '{id}' is not runtime-quiescent (D124)"
939 );
940 let internal_subscribers = release_cores
941 .iter()
942 .filter(|dependent| dependent.is_active())
943 .flat_map(|dependent| dependent.deps())
944 .filter(|dep| dep.ptr_eq(core))
945 .count();
946 assert!(
947 core.external_subscriber_count_for_release() <= internal_subscribers,
948 "graph: cannot release node group for {reason}; '{id}' still has live subscribers (D124)"
949 );
950 }
951 let release_events = if self.has_topology_observers() {
952 release_entries
953 .iter()
954 .filter_map(|(id, core)| {
955 self.inner
956 .entries
957 .borrow()
958 .iter()
959 .find(|entry| entry.id == *id && entry.core.ptr_eq(core))
960 .map(|entry| {
961 (
962 entry.id.clone(),
963 entry.factory.clone(),
964 self.topology_deps(&entry.core.deps()),
965 )
966 })
967 })
968 .collect::<Vec<_>>()
969 } else {
970 Vec::new()
971 };
972 {
973 let mut entries = self.inner.entries.borrow_mut();
974 entries.retain(|entry| !release_ids.contains(&entry.id));
975 }
976 {
977 let mut by_id = self.inner.by_id.borrow_mut();
978 let mut retired = self.inner.retired_ids.borrow_mut();
979 for id in &release_ids {
980 by_id.remove(id);
981 retired.insert(id.clone());
982 }
983 }
984 for core in &release_cores {
985 core.detach_graph_observer_subscribers_for_release();
986 }
987 let mut pending = release_cores.iter().rev().cloned().collect::<Vec<_>>();
988 let mut first_panic = None;
989 while !pending.is_empty() {
990 let before = pending.len();
991 let mut index = 0;
992 while index < pending.len() {
993 if pending[index].subscriber_count() != 0 {
994 index += 1;
995 continue;
996 }
997 let core = pending.remove(index);
998 let result = catch_unwind(AssertUnwindSafe(|| {
999 assert!(
1000 core.release_runtime_for_graph(),
1001 "graph: cannot release node group for {reason}; runtime became non-quiescent (D124)"
1002 );
1003 }));
1004 if let Err(panic) = result {
1005 if first_panic.is_none() {
1006 first_panic = Some(panic);
1007 }
1008 }
1009 }
1010 if pending.len() == before {
1011 break;
1012 }
1013 }
1014 if pending.is_empty() {
1015 for (id, factory, deps) in release_events {
1016 self.emit_topology_node_released(id, factory, deps);
1017 }
1018 }
1019 if let Some(panic) = first_panic {
1020 resume_unwind(panic);
1021 }
1022 assert!(
1023 pending.is_empty(),
1024 "graph: cannot release node group for {reason}; internal release order did not quiesce (D124)"
1025 );
1026 }
1027
1028 pub(crate) fn retain<T: 'static>(&self, node: &Node<T>, reason: &str) -> Box<dyn FnOnce()> {
1029 assert!(
1030 node.erased().same_graph_arena(&self.inner.arena),
1031 "graph: cannot retain node for {reason}; node belongs to a different graph"
1032 );
1033 node.subscribe(|_| {})
1034 }
1035
1036 fn assert_graph_local_deps(&self, deps: &[Core], label: &str) {
1037 for dep in deps {
1038 assert!(
1039 dep.same_graph_arena(&self.inner.arena),
1040 "{label} dep belongs to a different graph; cross-graph deps require a wire bridge"
1041 );
1042 }
1043 }
1044
1045 fn assert_registered_core(&self, core: &Core, label: &str) {
1046 assert!(
1047 core.same_graph_arena(&self.inner.arena),
1048 "{label} belongs to a different graph; cross-graph deps require a wire bridge"
1049 );
1050 assert!(
1051 self.inner
1052 .entries
1053 .borrow()
1054 .iter()
1055 .any(|entry| entry.core.ptr_eq(core)),
1056 "{label} is not a registered graph node (D152)"
1057 );
1058 }
1059
1060 fn registered_id_for_core(&self, core: &Core, prefix: &str) -> Option<String> {
1061 self.inner
1062 .entries
1063 .borrow()
1064 .iter()
1065 .find(|entry| entry.core.ptr_eq(core))
1066 .map(|entry| format!("{prefix}{}", entry.id))
1067 }
1068
1069 fn synthetic_id_for_core(&self, core: &Core, prefix: &str) -> String {
1070 let key = core.identity_key();
1071 let id = {
1072 let mut synth = self.inner.synth_ids.borrow_mut();
1073 if let Some(id) = synth.get(&key) {
1074 id.clone()
1075 } else {
1076 let next = self.inner.synth_seq.get();
1077 self.inner.synth_seq.set(next + 1);
1078 let id = format!(
1079 "~{}#{next}",
1080 core.factory().unwrap_or_else(|| "?".to_owned())
1081 );
1082 synth.insert(key, id.clone());
1083 id
1084 }
1085 };
1086 format!("{prefix}{id}")
1087 }
1088
1089 fn id_for_core(&self, core: &Core, prefix: &str) -> String {
1090 self.registered_id_for_core(core, prefix)
1091 .unwrap_or_else(|| self.synthetic_id_for_core(core, prefix))
1092 }
1093
1094 fn topology_deps(&self, deps: &[Core]) -> Vec<String> {
1095 deps.iter().map(|dep| self.id_for_core(dep, "")).collect()
1096 }
1097
1098 fn has_topology_observers(&self) -> bool {
1099 self.inner.topology_observer_active.get() > 0
1100 }
1101
1102 fn emit_topology_node_registered(&self, id: &str) {
1103 if !self.has_topology_observers() {
1104 return;
1105 }
1106 let Some(entry) = self
1107 .inner
1108 .entries
1109 .borrow()
1110 .iter()
1111 .find(|entry| entry.id == id)
1112 .cloned()
1113 else {
1114 return;
1115 };
1116 let seq = self.inner.clock.get();
1117 self.inner.clock.set(seq + 1);
1118 self.emit_topology_event(TopologyEvent {
1119 kind: TopologyEventKind::NodeRegistered,
1120 path: entry.id,
1121 deps: self.topology_deps(&entry.core.deps()),
1122 prev_deps: None,
1123 factory: Some(entry.factory),
1124 seq,
1125 });
1126 }
1127
1128 fn emit_topology_deps_changed(&self, id: &str, old_deps: &[Core], new_deps: &[Core]) {
1129 if !self.has_topology_observers() {
1130 return;
1131 }
1132 if self.find(id).is_none() {
1133 return;
1134 }
1135 let seq = self.inner.clock.get();
1136 self.inner.clock.set(seq + 1);
1137 self.emit_topology_event(TopologyEvent {
1138 kind: TopologyEventKind::DepsChanged,
1139 path: id.to_owned(),
1140 deps: self.topology_deps(new_deps),
1141 prev_deps: Some(self.topology_deps(old_deps)),
1142 factory: None,
1143 seq,
1144 });
1145 }
1146
1147 fn emit_topology_node_released(&self, path: String, factory: String, deps: Vec<String>) {
1148 if !self.has_topology_observers() {
1149 return;
1150 }
1151 let seq = self.inner.clock.get();
1152 self.inner.clock.set(seq + 1);
1153 self.emit_topology_event(TopologyEvent {
1154 kind: TopologyEventKind::NodeReleased,
1155 path,
1156 deps,
1157 prev_deps: None,
1158 factory: Some(factory),
1159 seq,
1160 });
1161 }
1162
1163 fn emit_topology_mount_changed(&self, path: String) {
1164 if !self.has_topology_observers() {
1165 return;
1166 }
1167 let seq = self.inner.clock.get();
1168 self.inner.clock.set(seq + 1);
1169 self.emit_topology_event(TopologyEvent {
1170 kind: TopologyEventKind::MountChanged,
1171 path,
1172 deps: Vec::new(),
1173 prev_deps: None,
1174 factory: Some("mount".to_owned()),
1175 seq,
1176 });
1177 }
1178
1179 fn ensure_mounted_topology_forwarders(&self) {
1180 if !self.has_topology_observers() {
1181 return;
1182 }
1183 let parent = Rc::downgrade(&self.inner);
1184 for mount in self.inner.mounts.borrow().iter() {
1185 if mount.topology_observer.borrow().is_some() {
1186 continue;
1187 }
1188 let parent = parent.clone();
1189 let mount_path = mount.at.clone();
1190 let observer = mount.graph.observe_topology().subscribe(move |event| {
1191 if let Some(inner) = parent.upgrade() {
1192 Graph { inner }.emit_mounted_topology_event(&mount_path, event);
1193 }
1194 });
1195 *mount.topology_observer.borrow_mut() = Some(observer);
1196 }
1197 }
1198
1199 fn release_mounted_topology_forwarders(&self) {
1200 for mount in self.inner.mounts.borrow().iter() {
1201 mount.topology_observer.borrow_mut().take();
1202 }
1203 }
1204
1205 fn emit_mounted_topology_event(&self, mount_path: &str, event: TopologyEvent) {
1206 if !self.has_topology_observers() {
1207 return;
1208 }
1209 let seq = self.inner.clock.get();
1210 self.inner.clock.set(seq + 1);
1211 self.emit_topology_event(TopologyEvent {
1212 kind: event.kind,
1213 path: prefix_topology_path(mount_path, &event.path),
1214 deps: event
1215 .deps
1216 .iter()
1217 .map(|dep| prefix_topology_path(mount_path, dep))
1218 .collect(),
1219 prev_deps: event.prev_deps.map(|deps| {
1220 deps.iter()
1221 .map(|dep| prefix_topology_path(mount_path, dep))
1222 .collect()
1223 }),
1224 factory: event.factory,
1225 seq,
1226 });
1227 }
1228
1229 fn emit_topology_event(&self, event: TopologyEvent) {
1230 if self.inner.topology_delivering.get() {
1231 self.inner.topology_queue.borrow_mut().push_back(event);
1232 return;
1233 }
1234 self.inner.topology_delivering.set(true);
1235 let mut current = Some(event);
1236 while let Some(event) = current {
1237 let observers = self
1238 .inner
1239 .topology_observers
1240 .borrow()
1241 .iter()
1242 .enumerate()
1243 .filter_map(|(id, entry)| entry.clone().map(|entry| (id, entry)))
1244 .collect::<Vec<_>>();
1245 for (id, observer) in observers {
1246 let still_active = self
1247 .inner
1248 .topology_observers
1249 .borrow()
1250 .get(id)
1251 .and_then(Option::as_ref)
1252 .is_some_and(|current| Rc::ptr_eq(¤t.sink, &observer.sink));
1253 if !still_active {
1254 continue;
1255 }
1256 if observer
1257 .path
1258 .as_deref()
1259 .is_some_and(|path| !topology_path_matches(&event.path, path))
1260 {
1261 continue;
1262 }
1263 let _ = catch_unwind(AssertUnwindSafe(|| (observer.sink)(event.clone())));
1264 }
1265 current = self.inner.topology_queue.borrow_mut().pop_front();
1266 }
1267 self.inner.topology_delivering.set(false);
1268 }
1269
1270 fn describe_with_prefix(&self, prefix: &str) -> DescribeSnapshot {
1271 let entries = self.inner.entries.borrow().clone();
1272 let mut nodes = Vec::with_capacity(entries.len());
1273 let mut edges = Vec::new();
1274 let mut discovered: Vec<Core> = Vec::new();
1275 for entry in &entries {
1276 let id = format!("{prefix}{}", entry.id);
1277 let deps: Vec<String> = entry
1278 .core
1279 .deps()
1280 .iter()
1281 .map(|dep| {
1282 if self.registered_id_for_core(dep, prefix).is_none()
1283 && !discovered.iter().any(|d| d.ptr_eq(dep))
1284 {
1285 discovered.push(dep.clone());
1286 }
1287 self.id_for_core(dep, prefix)
1288 })
1289 .collect();
1290 for dep_id in &deps {
1291 edges.push(DescribeEdge {
1292 from: dep_id.clone(),
1293 to: id.clone(),
1294 });
1295 }
1296 nodes.push(DescribeNode {
1297 id,
1298 name: entry.name.clone(),
1299 factory: entry.factory.clone(),
1300 status: entry.core.status(),
1301 value: entry.core.cache_any().as_ref().map(describe_value),
1302 version: entry.core.version(),
1303 deps,
1304 meta: if entry.meta.is_empty() {
1305 None
1306 } else {
1307 Some(entry.meta.clone())
1308 },
1309 });
1310 }
1311 let mut seen = HashSet::new();
1312 let mut i = 0;
1313 while i < discovered.len() {
1314 let core = discovered[i].clone();
1315 i += 1;
1316 let key = core.identity_key();
1317 if !seen.insert(key) {
1318 continue;
1319 }
1320 let id = self.synthetic_id_for_core(&core, prefix);
1321 let deps: Vec<String> = core
1322 .deps()
1323 .iter()
1324 .map(|dep| {
1325 if self.registered_id_for_core(dep, prefix).is_none()
1326 && !discovered.iter().any(|d| d.ptr_eq(dep))
1327 {
1328 discovered.push(dep.clone());
1329 }
1330 self.id_for_core(dep, prefix)
1331 })
1332 .collect();
1333 for dep_id in &deps {
1334 edges.push(DescribeEdge {
1335 from: dep_id.clone(),
1336 to: id.clone(),
1337 });
1338 }
1339 nodes.push(DescribeNode {
1340 id,
1341 name: None,
1342 factory: core.factory().unwrap_or_else(|| "?".to_owned()),
1343 status: core.status(),
1344 value: core.cache_any().as_ref().map(describe_value),
1345 version: core.version(),
1346 deps,
1347 meta: None,
1348 });
1349 }
1350 let subgraphs = self
1351 .inner
1352 .mounts
1353 .borrow()
1354 .iter()
1355 .map(|m| {
1356 let mount_path = format!("{prefix}{}", m.at);
1357 let mut child = m.graph.describe_with_prefix(&format!("{mount_path}::"));
1358 child.name = Some(mount_path);
1359 child
1360 })
1361 .collect::<Vec<_>>();
1362 DescribeSnapshot {
1363 name: self.inner.name.clone(),
1364 nodes,
1365 edges,
1366 subgraphs: if subgraphs.is_empty() {
1367 None
1368 } else {
1369 Some(subgraphs)
1370 },
1371 }
1372 }
1373
1374 fn observe_targets(
1375 &self,
1376 path: Option<&str>,
1377 sink: impl Fn(ObserveEvent) + 'static,
1378 ) -> GraphObserver {
1379 let targets = self.observe_target_entries(path);
1380 let clock = self.inner.clock.clone();
1381 let sink = Rc::new(sink);
1382 let mut observer = GraphObserver {
1383 unsubs: Vec::with_capacity(targets.len()),
1384 };
1385 for (path, core) in targets {
1386 let sink = sink.clone();
1387 let clock = clock.clone();
1388 let subscribed = catch_unwind(AssertUnwindSafe(|| {
1389 Node::<AnyValue>::from_core(core).subscribe_graph_observer(move |msg| {
1390 let seq = clock.get();
1391 clock.set(seq + 1);
1392 sink(ObserveEvent {
1393 path: path.clone(),
1394 msg: ObserveMessage::from_message(msg),
1395 tier: msg.tier(),
1396 seq,
1397 });
1398 })
1399 }));
1400 let unsub = match subscribed {
1401 Ok(unsub) => unsub,
1402 Err(payload) => {
1403 for unsub in observer.unsubs.drain(..).flatten() {
1404 unsub();
1405 }
1406 resume_unwind(payload);
1407 }
1408 };
1409 observer.unsubs.push(Some(unsub));
1410 }
1411 observer
1412 }
1413
1414 fn observe_target_entries(&self, path: Option<&str>) -> Vec<(String, Core)> {
1415 let mut all = Vec::new();
1416 self.collect_observe_entries("", &mut all);
1417 if let Some(path) = path {
1418 let prefix = format!("{path}::");
1419 all.into_iter()
1420 .filter(|(id, _)| id == path || id.starts_with(&prefix))
1421 .collect()
1422 } else {
1423 all
1424 }
1425 }
1426
1427 fn collect_observe_entries(&self, prefix: &str, out: &mut Vec<(String, Core)>) {
1428 out.extend(
1429 self.inner
1430 .entries
1431 .borrow()
1432 .iter()
1433 .map(|entry| (format!("{prefix}{}", entry.id), entry.core.clone())),
1434 );
1435 for mount in self.inner.mounts.borrow().iter() {
1436 mount
1437 .graph
1438 .collect_observe_entries(&format!("{prefix}{}::", mount.at), out);
1439 }
1440 }
1441
1442 fn contains_graph(&self, target: &Graph) -> bool {
1443 Rc::ptr_eq(&self.inner, &target.inner)
1444 || self
1445 .inner
1446 .mounts
1447 .borrow()
1448 .iter()
1449 .any(|m| m.graph.contains_graph(target))
1450 }
1451}
1452
1453pub fn graph() -> Graph {
1455 Graph::new(GraphOptions::default())
1456}
1457
1458pub fn graph_opts(opts: GraphOptions) -> Graph {
1460 Graph::new(opts)
1461}
1462
1463#[derive(Debug, Clone, Default, PartialEq, Eq)]
1464pub struct DescribeOpts {
1466 pub explain: Option<Explain>,
1468}
1469
1470#[derive(Debug, Clone, PartialEq, Eq)]
1471pub struct Explain {
1473 pub from: String,
1475 pub to: String,
1477}
1478
1479pub struct Values<'a> {
1486 records: &'a [DepRecord],
1487}
1488
1489impl<'a> Values<'a> {
1490 fn new(records: &'a [DepRecord]) -> Self {
1491 Self { records }
1492 }
1493
1494 pub fn len(&self) -> usize {
1496 self.records.len()
1497 }
1498
1499 pub fn is_empty(&self) -> bool {
1501 self.records.is_empty()
1502 }
1503
1504 pub fn prev<T: 'static>(&self, i: usize) -> Option<Rc<T>> {
1507 self.records
1508 .get(i)
1509 .and_then(|r| r.prev_data.clone())
1510 .and_then(|a| a.downcast::<T>().ok())
1511 }
1512
1513 pub fn batches<T: 'static>(&self, i: usize) -> Vec<Vec<Rc<T>>> {
1515 self.records
1516 .get(i)
1517 .map(|r| {
1518 r.wave_data
1519 .iter()
1520 .map(|wave| wave.iter().filter_map(WaveData::data::<T>).collect())
1521 .collect()
1522 })
1523 .unwrap_or_default()
1524 }
1525
1526 pub fn terminal(&self, i: usize) -> Option<&DepTerminal> {
1528 self.records.get(i).and_then(|r| r.terminal.as_ref())
1529 }
1530}
1531
1532#[derive(Clone)]
1537pub struct GraphNode {
1538 core: Core,
1539}
1540
1541impl GraphNode {
1542 fn new(core: Core) -> Self {
1543 Self { core }
1544 }
1545
1546 pub fn core(&self) -> Core {
1548 self.core.clone()
1549 }
1550
1551 pub fn status(&self) -> Status {
1553 self.core.status()
1554 }
1555
1556 pub fn version(&self) -> Option<NodeVersion> {
1558 self.core.version()
1559 }
1560
1561 pub fn cache_any(&self) -> Option<AnyValue> {
1563 self.core.cache_any()
1564 }
1565
1566 pub fn deps(&self) -> Vec<Core> {
1568 self.core.deps()
1569 }
1570
1571 pub fn subscribe(&self, sink: impl Fn(&Message<AnyValue>) + 'static) -> Box<dyn FnOnce()> {
1573 Node::<AnyValue>::from_core(self.core.clone()).subscribe(sink)
1574 }
1575
1576 pub fn down(&self, msgs: Vec<Message<AnyValue>>) {
1580 Node::<AnyValue>::from_core(self.core.clone()).down(msgs)
1581 }
1582}
1583
1584#[derive(Debug, Clone, PartialEq)]
1586pub struct DescribeSnapshot {
1587 pub name: Option<String>,
1589 pub nodes: Vec<DescribeNode>,
1591 pub edges: Vec<DescribeEdge>,
1593 pub subgraphs: Option<Vec<DescribeSnapshot>>,
1595}
1596
1597#[derive(Debug, Clone, PartialEq)]
1598pub struct DescribeNode {
1600 pub id: String,
1602 pub name: Option<String>,
1604 pub factory: String,
1606 pub status: Status,
1608 pub value: Option<DescribeValue>,
1610 pub version: Option<NodeVersion>,
1612 pub deps: Vec<String>,
1614 pub meta: Option<BTreeMap<String, String>>,
1616}
1617
1618#[derive(Debug, Clone, PartialEq, Eq)]
1619pub struct DescribeEdge {
1621 pub from: String,
1623 pub to: String,
1625}
1626
1627#[derive(Debug, Clone, PartialEq)]
1628pub enum DescribeValue {
1630 Bool(bool),
1632 I64(i64),
1634 U64(u64),
1636 F64(f64),
1638 String(String),
1640 Opaque,
1642}
1643
1644#[derive(Clone)]
1645pub struct ObserveEvent {
1647 pub path: String,
1649 pub msg: ObserveMessage,
1651 pub tier: Tier,
1653 pub seq: u64,
1655}
1656
1657#[derive(Clone)]
1658pub enum ObserveMessage {
1660 Start,
1662 Pause(LockId),
1664 Resume(LockId),
1666 Pull(PullDemand),
1668 Dirty,
1670 Data(AnyValue),
1672 Resolved,
1674 Invalidate,
1676 Complete,
1678 Error(String),
1680 Teardown,
1682}
1683
1684impl ObserveMessage {
1685 fn from_message(msg: &Message<AnyValue>) -> Self {
1686 match msg {
1687 Message::Start => ObserveMessage::Start,
1688 Message::Pause(lock) => ObserveMessage::Pause(lock.clone()),
1689 Message::Resume(lock) => ObserveMessage::Resume(lock.clone()),
1690 Message::Pull(demand) => ObserveMessage::Pull(demand.clone()),
1691 Message::Dirty => ObserveMessage::Dirty,
1692 Message::Data(value) => ObserveMessage::Data(value.clone()),
1693 Message::Resolved => ObserveMessage::Resolved,
1694 Message::Invalidate => ObserveMessage::Invalidate,
1695 Message::Complete => ObserveMessage::Complete,
1696 Message::Error(error) => ObserveMessage::Error(format!("{error}")),
1697 Message::Teardown => ObserveMessage::Teardown,
1698 }
1699 }
1700
1701 pub fn kind(&self) -> &'static str {
1703 match self {
1704 ObserveMessage::Start => "START",
1705 ObserveMessage::Pause(_) => "PAUSE",
1706 ObserveMessage::Resume(_) => "RESUME",
1707 ObserveMessage::Pull(_) => "PULL",
1708 ObserveMessage::Dirty => "DIRTY",
1709 ObserveMessage::Data(_) => "DATA",
1710 ObserveMessage::Resolved => "RESOLVED",
1711 ObserveMessage::Invalidate => "INVALIDATE",
1712 ObserveMessage::Complete => "COMPLETE",
1713 ObserveMessage::Error(_) => "ERROR",
1714 ObserveMessage::Teardown => "TEARDOWN",
1715 }
1716 }
1717}
1718
1719#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1720pub enum TopologyEventKind {
1722 NodeRegistered,
1724 DepsChanged,
1726 NodeReleased,
1728 MountChanged,
1730}
1731
1732#[derive(Clone)]
1733pub struct TopologyEvent {
1735 pub kind: TopologyEventKind,
1737 pub path: String,
1739 pub deps: Vec<String>,
1741 pub prev_deps: Option<Vec<String>>,
1743 pub factory: Option<String>,
1745 pub seq: u64,
1747}
1748
1749pub struct GraphObserver {
1751 unsubs: Vec<Option<Box<dyn FnOnce()>>>,
1752}
1753
1754pub struct GraphTopologyObserver {
1756 graph: Graph,
1757 id: usize,
1758 active: bool,
1759}
1760
1761#[derive(Clone)]
1762pub struct ObserveStream {
1764 graph: Graph,
1765 path: Option<String>,
1766}
1767
1768#[derive(Clone)]
1769pub struct TopologyStream {
1771 graph: Graph,
1772 path: Option<String>,
1773}
1774
1775impl ObserveStream {
1776 pub fn subscribe(&self, sink: impl Fn(ObserveEvent) + 'static) -> GraphObserver {
1778 self.graph.observe_targets(self.path.as_deref(), sink)
1779 }
1780}
1781
1782impl TopologyStream {
1783 pub fn subscribe(&self, sink: impl Fn(TopologyEvent) + 'static) -> GraphTopologyObserver {
1785 let id = self
1786 .graph
1787 .inner
1788 .topology_observer_free
1789 .borrow_mut()
1790 .pop()
1791 .unwrap_or_else(|| {
1792 let mut observers = self.graph.inner.topology_observers.borrow_mut();
1793 let id = observers.len();
1794 observers.push(None);
1795 id
1796 });
1797 let mut observers = self.graph.inner.topology_observers.borrow_mut();
1798 observers[id] = Some(TopologyObserverEntry {
1799 path: self.path.clone(),
1800 sink: Rc::new(sink),
1801 });
1802 self.graph
1803 .inner
1804 .topology_observer_active
1805 .set(self.graph.inner.topology_observer_active.get() + 1);
1806 self.graph.ensure_mounted_topology_forwarders();
1807 GraphTopologyObserver {
1808 graph: self.graph.clone(),
1809 id,
1810 active: true,
1811 }
1812 }
1813}
1814
1815impl GraphObserver {
1816 pub fn unsubscribe(mut self) {
1818 self.unsubscribe_all();
1819 }
1820
1821 fn unsubscribe_all(&mut self) {
1822 for unsub in &mut self.unsubs {
1823 if let Some(unsub) = unsub.take() {
1824 unsub();
1825 }
1826 }
1827 }
1828}
1829
1830impl GraphTopologyObserver {
1831 pub fn unsubscribe(mut self) {
1833 self.unsubscribe_inner();
1834 }
1835
1836 fn unsubscribe_inner(&mut self) {
1837 if !self.active {
1838 return;
1839 }
1840 self.active = false;
1841 let mut release_forwarders = false;
1842 if let Some(slot) = self
1843 .graph
1844 .inner
1845 .topology_observers
1846 .borrow_mut()
1847 .get_mut(self.id)
1848 {
1849 if slot.take().is_some() {
1850 let active = self
1851 .graph
1852 .inner
1853 .topology_observer_active
1854 .get()
1855 .saturating_sub(1);
1856 self.graph.inner.topology_observer_active.set(active);
1857 release_forwarders = active == 0;
1858 self.graph
1859 .inner
1860 .topology_observer_free
1861 .borrow_mut()
1862 .push(self.id);
1863 }
1864 }
1865 if release_forwarders {
1866 self.graph.release_mounted_topology_forwarders();
1867 }
1868 }
1869}
1870
1871impl Drop for GraphTopologyObserver {
1872 fn drop(&mut self) {
1873 self.unsubscribe_inner();
1874 }
1875}
1876
1877impl Drop for GraphObserver {
1878 fn drop(&mut self) {
1879 self.unsubscribe_all();
1880 }
1881}
1882
1883#[derive(Debug, Clone, PartialEq, Eq)]
1884pub struct NodeProfile {
1886 pub invokes: u64,
1888 pub total_duration_ns: u128,
1890 pub last_duration_ns: u128,
1892 pub status: Status,
1894}
1895
1896#[derive(Debug, Clone, PartialEq, Eq)]
1897pub struct Profile {
1899 pub total_invokes: u64,
1901 pub nodes: BTreeMap<String, NodeProfile>,
1903}
1904
1905fn explain_subset(snapshot: DescribeSnapshot, explain: &Explain) -> DescribeSnapshot {
1906 let snapshot = flatten_describe_snapshot(snapshot);
1907 let mut fwd: HashMap<String, Vec<String>> = HashMap::new();
1908 let mut rev: HashMap<String, Vec<String>> = HashMap::new();
1909 for edge in &snapshot.edges {
1910 fwd.entry(edge.from.clone())
1911 .or_default()
1912 .push(edge.to.clone());
1913 rev.entry(edge.to.clone())
1914 .or_default()
1915 .push(edge.from.clone());
1916 }
1917 let from_reach = reachable(&explain.from, &fwd);
1918 let to_reach = reachable(&explain.to, &rev);
1919 let on_path: HashSet<String> = from_reach.intersection(&to_reach).cloned().collect();
1920 DescribeSnapshot {
1921 name: snapshot.name,
1922 nodes: snapshot
1923 .nodes
1924 .into_iter()
1925 .filter(|node| on_path.contains(&node.id))
1926 .collect(),
1927 edges: snapshot
1928 .edges
1929 .into_iter()
1930 .filter(|edge| on_path.contains(&edge.from) && on_path.contains(&edge.to))
1931 .collect(),
1932 subgraphs: None,
1933 }
1934}
1935
1936fn flatten_describe_snapshot(mut snapshot: DescribeSnapshot) -> DescribeSnapshot {
1937 if let Some(subgraphs) = snapshot.subgraphs.take() {
1938 for child in subgraphs {
1939 let mut child = flatten_describe_snapshot(child);
1940 snapshot.nodes.append(&mut child.nodes);
1941 snapshot.edges.append(&mut child.edges);
1942 }
1943 }
1944 snapshot
1945}
1946
1947fn reachable(start: &str, adj: &HashMap<String, Vec<String>>) -> HashSet<String> {
1948 let mut seen = HashSet::from([start.to_owned()]);
1949 let mut stack = vec![start.to_owned()];
1950 while let Some(current) = stack.pop() {
1951 for next in adj.get(¤t).into_iter().flatten() {
1952 if seen.insert(next.clone()) {
1953 stack.push(next.clone());
1954 }
1955 }
1956 }
1957 seen
1958}
1959
1960fn topology_path_matches(event_path: &str, path: &str) -> bool {
1961 event_path == path || event_path.starts_with(&format!("{path}::"))
1962}
1963
1964fn prefix_topology_path(prefix: &str, path: &str) -> String {
1965 format!("{prefix}::{path}")
1966}
1967
1968fn describe_value(value: &AnyValue) -> DescribeValue {
1969 let value = value.as_ref();
1970 if let Some(v) = value.downcast_ref::<bool>() {
1971 DescribeValue::Bool(*v)
1972 } else if let Some(v) = value.downcast_ref::<i8>() {
1973 DescribeValue::I64(i64::from(*v))
1974 } else if let Some(v) = value.downcast_ref::<i16>() {
1975 DescribeValue::I64(i64::from(*v))
1976 } else if let Some(v) = value.downcast_ref::<i32>() {
1977 DescribeValue::I64(i64::from(*v))
1978 } else if let Some(v) = value.downcast_ref::<i64>() {
1979 DescribeValue::I64(*v)
1980 } else if let Some(v) = value.downcast_ref::<isize>() {
1981 DescribeValue::I64(*v as i64)
1982 } else if let Some(v) = value.downcast_ref::<u8>() {
1983 DescribeValue::U64(u64::from(*v))
1984 } else if let Some(v) = value.downcast_ref::<u16>() {
1985 DescribeValue::U64(u64::from(*v))
1986 } else if let Some(v) = value.downcast_ref::<u32>() {
1987 DescribeValue::U64(u64::from(*v))
1988 } else if let Some(v) = value.downcast_ref::<u64>() {
1989 DescribeValue::U64(*v)
1990 } else if let Some(v) = value.downcast_ref::<usize>() {
1991 DescribeValue::U64(*v as u64)
1992 } else if let Some(v) = value.downcast_ref::<f32>() {
1993 DescribeValue::F64(f64::from(*v))
1994 } else if let Some(v) = value.downcast_ref::<f64>() {
1995 DescribeValue::F64(*v)
1996 } else if let Some(v) = value.downcast_ref::<String>() {
1997 DescribeValue::String(v.clone())
1998 } else if let Some(v) = value.downcast_ref::<&'static str>() {
1999 DescribeValue::String((*v).to_owned())
2000 } else {
2001 DescribeValue::Opaque
2002 }
2003}
2004
2005#[cfg(test)]
2006mod tests {
2007 use std::cell::{Cell, RefCell};
2008 use std::panic::{catch_unwind, AssertUnwindSafe};
2009 use std::rc::Rc;
2010
2011 use super::*;
2012
2013 #[test]
2014 fn release_nodes_hides_registry_and_releases_remaining_after_cleanup_panic() {
2015 let g = graph();
2016 let panic_node = g.state_opts(1, GraphNodeOpts::named("panic"));
2017 let later_node = g.state_opts(2, GraphNodeOpts::named("later"));
2018 let later_released = Rc::new(Cell::new(false));
2019 let flag = later_released.clone();
2020 panic_node
2021 .erased()
2022 .register_on_deactivation(Box::new(|| panic!("cleanup boom")));
2023 later_node
2024 .erased()
2025 .register_on_deactivation(Box::new(move || flag.set(true)));
2026
2027 let result = catch_unwind(AssertUnwindSafe(|| {
2028 g.release_nodes(&[later_node.erased(), panic_node.erased()], "test");
2029 }));
2030
2031 assert!(result.is_err());
2032 assert!(
2033 g.find("panic").is_none() && g.find("later").is_none(),
2034 "D122: cleanup panic must not leave graph registry pointing at released slots"
2035 );
2036 assert!(
2037 later_released.get(),
2038 "a cleanup panic in one released node must not skip later nodes in the group"
2039 );
2040 }
2041
2042 #[test]
2043 fn release_nodes_emits_node_released_even_when_cleanup_panics_after_commit() {
2044 let g = graph();
2045 let panic_node = g.state_opts(1, GraphNodeOpts::named("panic"));
2046 panic_node
2047 .erased()
2048 .register_on_deactivation(Box::new(|| panic!("cleanup boom")));
2049 let events = Rc::new(RefCell::new(Vec::new()));
2050 let event_sink = events.clone();
2051 let observer = g
2052 .observe_topology()
2053 .subscribe(move |event| event_sink.borrow_mut().push(event));
2054
2055 let result = catch_unwind(AssertUnwindSafe(|| {
2056 g.release_nodes(&[panic_node.erased()], "test");
2057 }));
2058
2059 observer.unsubscribe();
2060 assert!(result.is_err());
2061 let events = events.borrow();
2062 assert_eq!(events.len(), 1);
2063 assert_eq!(events[0].kind, TopologyEventKind::NodeReleased);
2064 assert_eq!(events[0].path, "panic");
2065 assert_eq!(events[0].factory.as_deref(), Some("state"));
2066 assert_eq!(events[0].deps, Vec::<String>::new());
2067 assert_eq!(events[0].seq, 0);
2068 assert!(g.find("panic").is_none());
2069 }
2070
2071 #[test]
2072 fn release_nodes_detaches_only_graph_observer_subscribers() {
2073 let g = graph();
2074 let group = g.topology_group_opts(TopologyGroupOptions::named("child"));
2075 let child = group.state_opts(1, GraphNodeOpts::named("child/value"));
2076 let unrelated = g.state_opts(10, GraphNodeOpts::named("unrelated"));
2077 let events = Rc::new(RefCell::new(Vec::new()));
2078 let event_sink = events.clone();
2079 let observer = g
2080 .observe()
2081 .subscribe(move |event| event_sink.borrow_mut().push(event));
2082
2083 group.release();
2084
2085 assert!(
2086 g.find("child/value").is_none(),
2087 "D122: released group nodes disappear from find"
2088 );
2089 assert!(
2090 !g.describe()
2091 .nodes
2092 .iter()
2093 .any(|entry| entry.id == "child/value"),
2094 "D122: released group nodes disappear from describe"
2095 );
2096
2097 unrelated.set(11);
2098 observer.unsubscribe();
2099
2100 let events = events.borrow();
2101 assert!(
2102 events
2103 .iter()
2104 .any(|event| event.path == "unrelated"
2105 && matches!(event.msg, ObserveMessage::Data(_))),
2106 "D557: graph observer remains an ordinary observer and keeps unrelated subscriptions"
2107 );
2108 assert!(
2109 events.iter().all(|event| event.path != "child/value"
2110 || !matches!(
2111 event.msg,
2112 ObserveMessage::Complete | ObserveMessage::Error(_) | ObserveMessage::Teardown
2113 )),
2114 "D527: release does not synthesize protocol terminals"
2115 );
2116 drop(child);
2117 }
2118
2119 #[test]
2120 fn release_nodes_still_rejects_public_subscribers() {
2121 let g = graph();
2122 let group = g.topology_group_opts(TopologyGroupOptions::named("child"));
2123 let child = group.state_opts(1, GraphNodeOpts::named("child/value"));
2124 let public_subscriber = child.subscribe(|_| {});
2125 let observer = g.observe().subscribe(|_| {});
2126
2127 let result = catch_unwind(AssertUnwindSafe(|| group.release()));
2128
2129 observer.unsubscribe();
2130 public_subscriber();
2131 group.release();
2132 assert!(result.is_err());
2133 assert!(
2134 g.find("child/value").is_none(),
2135 "after public subscriber detaches, release can proceed"
2136 );
2137 }
2138
2139 #[test]
2140 fn graph_observer_unsubscribe_after_release_is_stale_safe() {
2141 let g = graph();
2142 let group = g.topology_group_opts(TopologyGroupOptions::named("child"));
2143 let _child = group.state_opts(1, GraphNodeOpts::named("child/value"));
2144 let observer = g.observe().subscribe(|_| {});
2145
2146 group.release();
2147 observer.unsubscribe();
2148 }
2149}