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