1use std::rc::Rc;
2
3use web_time::Instant;
4
5use crate::{
6 Applier, ApplierGuard, ApplierHost, CommandQueue, Composer, CompositionPassDebugStats,
7 ConcreteApplierHost, DefaultScheduler, Key, NodeError, NodeId, RecomposeScope, RetentionPolicy,
8 Runtime, RuntimeHandle, ScopeId, SlotDebugSnapshot, SlotTable, SlotTableDebugStats, SlotsHost,
9 SnapshotStateObserver, collections::map::HashMap, debug_scope_invalidation_sources,
10 debug_scope_label, runtime, scheduler_ref, snapshot_state_observer,
11};
12
13pub struct Composition<A: Applier + 'static> {
14 pub(crate) composer_state: Rc<crate::composer::ComposerRuntimeState>,
15 pub(crate) slots: Rc<SlotsHost>,
16 pub(crate) applier: Rc<ConcreteApplierHost<A>>,
17 pub(crate) runtime: Runtime,
18 pub(crate) observer: SnapshotStateObserver,
19 pub(crate) root: Option<NodeId>,
20 pub(crate) root_key: Option<Key>,
21 pub(crate) root_render_requested: bool,
22 pub(crate) last_pass_stats: CompositionPassDebugStats,
23 teardown: Option<runtime::StateTeardownScope>,
24}
25
26pub const ROOT_RENDER_REPLAY_LIMIT: usize = 100;
42
43fn recompose_scope_telemetry_threshold_ms() -> Option<f64> {
44 std::env::var("CRANPOSE_RECOMPOSE_SCOPE_TELEMETRY_MS")
45 .ok()
46 .and_then(|value| value.parse::<f64>().ok())
47 .filter(|value| value.is_finite() && *value >= 0.0)
48}
49
50impl<A: Applier + 'static> Composition<A> {
51 pub fn new(applier: A) -> Self {
52 Self::with_runtime(applier, Runtime::new(scheduler_ref(DefaultScheduler)))
53 }
54
55 pub fn with_runtime(applier: A, runtime: Runtime) -> Self {
56 let composer_state = Rc::new(crate::composer::ComposerRuntimeState::default());
57 let slots = Rc::new(SlotsHost::new(SlotTable::new()));
58 let applier = Rc::new(ConcreteApplierHost::new(applier));
59 let observer_handle = runtime.handle();
60 let observer = SnapshotStateObserver::new(move |callback| {
61 observer_handle.enqueue_ui_task(callback);
62 });
63 observer.start();
64 Self {
65 composer_state,
66 slots,
67 applier,
68 runtime,
69 observer,
70 root: None,
71 root_key: None,
72 root_render_requested: false,
73 last_pass_stats: CompositionPassDebugStats::default(),
74 teardown: None,
75 }
76 }
77
78 pub fn root_key(&self) -> Option<Key> {
81 self.root_key
82 }
83
84 pub fn set_retention_policy(&self, policy: RetentionPolicy) {
85 self.composer_state.set_retention_policy(policy);
86 }
87
88 fn slots_host(&self) -> Rc<SlotsHost> {
89 Rc::clone(&self.slots)
90 }
91
92 fn applier_host(&self) -> Rc<dyn ApplierHost> {
93 self.applier.clone()
94 }
95
96 fn reset_last_pass_stats(&mut self) {
97 self.last_pass_stats = CompositionPassDebugStats::default();
98 }
99
100 fn maybe_dump_slot_table(&self, label: &str) {
101 if !crate::env_flag!("COMPOSE_DEBUG_SLOT_TABLE") {
102 return;
103 }
104 eprintln!(
105 "[COMPOSE_DEBUG_SLOT_TABLE] {label}\n{:#?}",
106 self.debug_slot_snapshot()
107 );
108 }
109
110 pub fn take_root_render_request(&mut self) -> bool {
111 std::mem::take(&mut self.root_render_requested)
112 }
113
114 pub fn request_root_render(&mut self) {
115 self.root_render_requested = true;
116 self.runtime.handle().schedule();
117 }
118
119 fn record_pass_stats(
120 &mut self,
121 commands: &CommandQueue,
122 side_effects: &Vec<Box<dyn FnOnce()>>,
123 ) {
124 self.last_pass_stats.commands_len = self.last_pass_stats.commands_len.max(commands.len());
125 self.last_pass_stats.commands_cap =
126 self.last_pass_stats.commands_cap.max(commands.capacity());
127 self.last_pass_stats.command_payload_len_bytes = self
128 .last_pass_stats
129 .command_payload_len_bytes
130 .max(commands.payload_len_bytes());
131 self.last_pass_stats.command_payload_cap_bytes = self
132 .last_pass_stats
133 .command_payload_cap_bytes
134 .max(commands.payload_capacity_bytes());
135 self.last_pass_stats.sync_children_len = self
136 .last_pass_stats
137 .sync_children_len
138 .max(commands.sync_children.len());
139 self.last_pass_stats.sync_children_cap = self
140 .last_pass_stats
141 .sync_children_cap
142 .max(commands.sync_children.capacity());
143 self.last_pass_stats.sync_child_ids_len = self
144 .last_pass_stats
145 .sync_child_ids_len
146 .max(commands.sync_child_ids.len());
147 self.last_pass_stats.sync_child_ids_cap = self
148 .last_pass_stats
149 .sync_child_ids_cap
150 .max(commands.sync_child_ids.capacity());
151 self.last_pass_stats.side_effects_len = self
152 .last_pass_stats
153 .side_effects_len
154 .max(side_effects.len());
155 self.last_pass_stats.side_effects_cap = self
156 .last_pass_stats
157 .side_effects_cap
158 .max(side_effects.capacity());
159 }
160
161 fn finalize_runtime_state(&mut self) {
162 let runtime_handle = self.runtime_handle();
163 self.observer.prune_dead_scopes();
164 if !self.runtime.has_updates()
165 && !runtime_handle.has_invalid_scopes()
166 && !runtime_handle.has_frame_callbacks()
167 && !runtime_handle.has_pending_ui()
168 {
169 self.runtime.set_needs_frame(false);
170 }
171 }
172
173 fn abandon_host_after_apply_failure(&mut self, host: &Rc<SlotsHost>) {
174 host.abandon_after_apply_failure();
175 if Rc::ptr_eq(host, &self.slots) {
176 self.root = None;
177 }
178 self.root_render_requested = true;
179 self.finalize_runtime_state();
180 }
181
182 fn apply_commands_and_updates_for_host(
183 &mut self,
184 host: &Rc<SlotsHost>,
185 runtime_handle: &RuntimeHandle,
186 commands: CommandQueue,
187 ) -> Result<(), NodeError> {
188 let result = {
189 let mut applier = self.applier.borrow_dyn();
190 let mut result = commands.apply(&mut *applier);
191 if result.is_ok() {
192 for update in runtime_handle.take_updates() {
193 if let Err(err) = update.apply(&mut *applier) {
194 result = Err(err);
195 break;
196 }
197 }
198 }
199 result
200 };
201 if result.is_err() {
202 self.abandon_host_after_apply_failure(host);
203 }
204 result
205 }
206
207 fn render_root_pass(&mut self, key: Key, content: &mut dyn FnMut()) -> Result<(), NodeError> {
208 self.root_key = Some(key);
209 self.root_render_requested = false;
210 let runtime_handle = self.runtime_handle();
211 runtime_handle.drain_ui();
212 let side_effects = {
213 let _teardown = runtime::enter_state_teardown_scope();
214 let composer = Composer::new_with_shared_state(
215 Rc::clone(&self.composer_state),
216 Rc::clone(&self.slots),
217 self.applier.clone(),
218 runtime_handle.clone(),
219 self.observer.clone(),
220 self.root,
221 );
222 self.observer.begin_frame();
223 let (root, commands, side_effects, compact_applier) = composer.install(|composer| {
224 let ((), outcome) = composer.try_with_slot_host_pass(
225 Rc::clone(&self.slots),
226 crate::slot::SlotPassMode::Compose,
227 |composer| composer.with_group(key, |_| content()),
228 )?;
229 let root = composer.root();
230 let commands = composer.take_commands();
231 let side_effects = composer.take_side_effects();
232 Ok((root, commands, side_effects, outcome.compacted))
233 })?;
234 self.record_pass_stats(&commands, &side_effects);
235 self.apply_commands_and_updates_for_host(
236 &Rc::clone(&self.slots),
237 &runtime_handle,
238 commands,
239 )?;
240 if compact_applier {
241 self.applier.compact();
242 self.applier.borrow_dyn().clear_recycled_nodes();
243 }
244
245 self.root = root;
246 side_effects
247 };
248 runtime_handle.drain_ui();
249 for effect in side_effects {
250 effect();
251 }
252 runtime_handle.drain_ui();
253 self.maybe_dump_slot_table("root_render_pass");
254 Ok(())
255 }
256
257 fn reconcile_with_content(
258 &mut self,
259 key: Key,
260 content: &mut dyn FnMut(),
261 ) -> Result<bool, NodeError> {
262 self.root_key = Some(key);
263 let mut did_work = false;
264 let mut root_render_replays = 0usize;
265 loop {
266 did_work |= self.process_invalid_scopes_until_root_request()?;
267 if !self.take_root_render_request() {
268 return Ok(did_work);
269 }
270
271 root_render_replays += 1;
272 if root_render_replays > ROOT_RENDER_REPLAY_LIMIT {
273 log::error!(
274 "root render replay looped past {ROOT_RENDER_REPLAY_LIMIT} iterations; breaking to keep UI responsive"
275 );
276 return Err(NodeError::RecompositionLimitExceeded {
277 operation: "root render replay",
278 limit: ROOT_RENDER_REPLAY_LIMIT,
279 });
280 }
281
282 self.render_root_pass(key, content)?;
283 did_work = true;
284 }
285 }
286
287 pub fn render(&mut self, key: Key, mut content: impl FnMut()) -> Result<(), NodeError> {
288 self.reset_last_pass_stats();
289 self.render_root_pass(key, &mut content)?;
290 let _ = self.process_invalid_scopes()?;
291 Ok(())
292 }
293
294 pub fn render_stable(&mut self, key: Key, mut content: impl FnMut()) -> Result<(), NodeError> {
297 self.reset_last_pass_stats();
298 self.render_root_pass(key, &mut content)?;
299 let _ = self.reconcile_with_content(key, &mut content)?;
300 Ok(())
301 }
302
303 pub fn reconcile(&mut self, key: Key, mut content: impl FnMut()) -> Result<bool, NodeError> {
306 self.reconcile_with_content(key, &mut content)
307 }
308
309 pub fn should_render(&self) -> bool {
318 self.root_render_requested || self.runtime.needs_frame() || self.runtime.has_updates()
319 }
320
321 pub fn should_recompose(&self) -> bool {
337 self.root_render_requested || self.runtime.has_updates()
338 }
339
340 pub fn runtime_handle(&self) -> RuntimeHandle {
341 self.runtime.handle()
342 }
343
344 pub fn applier_mut(&mut self) -> ApplierGuard<'_, A> {
345 ApplierGuard::new(self.applier.borrow_typed())
346 }
347
348 pub fn root(&self) -> Option<NodeId> {
349 self.root
350 }
351
352 pub fn debug_dump_slot_table_groups(&self) -> Vec<(usize, Key, Option<ScopeId>, usize)> {
353 self.slots.borrow().debug_dump_groups()
354 }
355
356 pub fn debug_dump_slot_entries(&self) -> Vec<crate::SlotDebugEntry> {
357 self.slots.borrow().debug_dump_slot_entries()
358 }
359
360 pub fn slot_table_heap_bytes(&self) -> usize {
361 self.slots.borrow().heap_bytes()
362 }
363
364 pub fn debug_slot_table_stats(&self) -> SlotTableDebugStats {
365 self.slots.debug_stats()
366 }
367
368 pub fn debug_slot_snapshot(&self) -> SlotDebugSnapshot {
369 self.slots.debug_snapshot()
370 }
371
372 pub fn debug_observer_stats(&self) -> snapshot_state_observer::SnapshotStateObserverDebugStats {
373 self.observer.debug_stats()
374 }
375
376 pub fn debug_last_pass_stats(&self) -> CompositionPassDebugStats {
377 self.last_pass_stats
378 }
379
380 #[cfg(test)]
381 pub(crate) fn debug_validate_slots(&self) -> Result<(), crate::slot::SlotInvariantError> {
382 let table = self.slots.borrow();
383 table.validate()?;
384 self.composer_state
385 .validate_host_retention(self.slots.as_ref(), &table)
386 }
387
388 fn process_invalid_scopes_until_root_request(&mut self) -> Result<bool, NodeError> {
389 let runtime_handle = self.runtime_handle();
390 let mut did_recompose = false;
391 let mut loop_count = 0;
392 loop {
393 loop_count += 1;
394 if loop_count > ROOT_RENDER_REPLAY_LIMIT {
395 log::error!(
396 "process_invalid_scopes looped past {ROOT_RENDER_REPLAY_LIMIT} iterations; breaking to keep UI responsive"
397 );
398 return Err(NodeError::RecompositionLimitExceeded {
399 operation: "process_invalid_scopes",
400 limit: ROOT_RENDER_REPLAY_LIMIT,
401 });
402 }
403 runtime_handle.drain_ui();
404 self.dispose_forgotten_movables()?;
405 let Some(scopes) = live_invalidated_scopes(&runtime_handle) else {
406 break;
407 };
408 if scopes.is_empty() {
409 continue;
410 }
411 did_recompose |= recomposes_content(&scopes);
412 let runtime_clone = runtime_handle.clone();
413 let root_host = self.slots_host();
414 let mut scope_groups: Vec<(Rc<SlotsHost>, Vec<RecomposeScope>)> = Vec::new();
415 let mut scope_group_index: HashMap<usize, usize> = HashMap::default();
416 for scope in scopes {
417 let host = scope
418 .slots_runtime_state()
419 .and_then(|state| {
420 scope
421 .slots_storage_key()
422 .and_then(|storage_key| state.host_for_storage_key(storage_key))
423 })
424 .or_else(|| {
425 scope.slots_storage_key().and_then(|storage_key| {
426 self.composer_state.host_for_storage_key(storage_key)
427 })
428 })
429 .unwrap_or_else(|| Rc::clone(&root_host));
430 let host_key = host.storage_key();
431 if let Some(index) = scope_group_index.get(&host_key).copied() {
432 scope_groups[index].1.push(scope);
433 } else {
434 scope_group_index.insert(host_key, scope_groups.len());
435 scope_groups.push((host, vec![scope]));
436 }
437 }
438 let mut host_group_index = 0usize;
439 while host_group_index < scope_groups.len() {
440 let (host, scopes) = &scope_groups[host_group_index];
441 let scope_telemetry_threshold_ms = recompose_scope_telemetry_threshold_ms();
442 let shared_state = host
443 .runtime_state()
444 .or_else(|| scopes.first().and_then(RecomposeScope::slots_runtime_state))
445 .unwrap_or_else(|| Rc::clone(&self.composer_state));
446 let side_effects = {
447 let _teardown = runtime::enter_state_teardown_scope();
448 let composer = Composer::new_with_shared_state(
449 shared_state,
450 Rc::clone(host),
451 self.applier_host(),
452 runtime_clone.clone(),
453 self.observer.clone(),
454 self.root,
455 );
456 composer.parent_stack().clear();
457 self.observer.begin_frame();
458 let (root, commands, side_effects, requested_root_render, compact_applier) =
459 composer.install(|composer| {
460 let ((), outcome) = composer.try_with_slot_host_pass(
461 Rc::clone(host),
462 crate::slot::SlotPassMode::Recompose,
463 |composer| {
464 for scope in scopes {
465 if let Some(threshold_ms) = scope_telemetry_threshold_ms {
466 let start = Instant::now();
467 composer.recompose_group(scope);
468 let elapsed_ms = start.elapsed().as_secs_f64() * 1000.0;
469 if elapsed_ms >= threshold_ms {
470 eprintln!(
471 "[recompose-scope-telemetry] scope_id={} label={:?} elapsed_ms={elapsed_ms:.3} invalidation_sources={:?}",
472 scope.id(),
473 debug_scope_label(scope.id()),
474 debug_scope_invalidation_sources(scope.id())
475 );
476 }
477 } else {
478 composer.recompose_group(scope);
479 }
480 }
481 },
482 )?;
483 let root = composer.root();
484 let commands = composer.take_commands();
485 let side_effects = composer.take_side_effects();
486 let requested_root_render = composer.take_root_render_request();
487 Ok((
488 root,
489 commands,
490 side_effects,
491 requested_root_render,
492 outcome.compacted,
493 ))
494 })?;
495 self.record_pass_stats(&commands, &side_effects);
496 self.apply_commands_and_updates_for_host(host, &runtime_handle, commands)?;
497 if compact_applier {
498 self.applier.compact();
499 self.applier.borrow_dyn().clear_recycled_nodes();
500 }
501 if root.is_some() {
502 self.root = root;
503 }
504 if requested_root_render {
505 self.root_render_requested = true;
506 }
507 side_effects
508 };
509 runtime_handle.drain_ui();
510 for effect in side_effects {
511 effect();
512 }
513 runtime_handle.drain_ui();
514 self.maybe_dump_slot_table("recompose_pass");
515 if self.root_render_requested {
516 for (_, remaining_scopes) in scope_groups.iter().skip(host_group_index + 1) {
517 for scope in remaining_scopes {
518 runtime_handle.requeue_invalid_scope(scope.id(), scope.downgrade());
519 }
520 }
521 break;
522 }
523 host_group_index += 1;
524 }
525 if self.root_render_requested {
526 break;
527 }
528 }
529 self.finalize_runtime_state();
530 Ok(did_recompose)
531 }
532
533 pub fn process_invalid_scopes(&mut self) -> Result<bool, NodeError> {
534 self.process_invalid_scopes_until_root_request()
535 }
536
537 fn dispose_forgotten_movables(&mut self) -> Result<(), NodeError> {
538 let runtime_handle = self.runtime_handle();
539 let ids = runtime_handle.take_forgotten_movables();
540 if ids.is_empty() {
541 return Ok(());
542 }
543 let host = self.slots_host();
544 let composer = Composer::new_with_shared_state(
545 Rc::clone(&self.composer_state),
546 Rc::clone(&host),
547 self.applier_host(),
548 runtime_handle.clone(),
549 self.observer.clone(),
550 self.root,
551 );
552 let commands = composer.install(|composer| {
553 composer.forget_movables(&ids)?;
554 Ok::<_, NodeError>(composer.take_commands())
555 })?;
556 self.apply_commands_and_updates_for_host(&host, &runtime_handle, commands)
557 }
558}
559
560fn recomposes_content(scopes: &[RecomposeScope]) -> bool {
564 scopes.iter().any(|scope| !scope.is_derivation())
565}
566
567fn live_invalidated_scopes(runtime_handle: &RuntimeHandle) -> Option<Vec<RecomposeScope>> {
568 let pending = runtime_handle.take_invalidated_scopes();
569 if pending.is_empty() {
570 return None;
571 }
572 let mut scopes = Vec::with_capacity(pending.len());
573 for (id, weak) in pending {
574 if let Some(inner) = weak.upgrade() {
575 scopes.push(RecomposeScope { inner });
576 } else {
577 runtime_handle.mark_scope_recomposed(id);
578 }
579 }
580 Some(scopes)
581}
582
583impl<A: Applier + 'static> Composition<A> {
584 pub fn flush_pending_node_updates(&mut self) -> Result<(), NodeError> {
585 let updates = self.runtime_handle().take_updates();
586 let mut applier = self.applier.borrow_dyn();
587 for update in updates {
588 update.apply(&mut *applier)?;
589 }
590 Ok(())
591 }
592}
593
594impl<A: Applier + 'static> Drop for Composition<A> {
595 fn drop(&mut self) {
596 self.observer.stop();
597 self.teardown = Some(runtime::enter_state_teardown_scope());
598 }
599}