@@ -82,6 +82,11 @@ let count_changes entries =
8282 entries;
8383 (! adds, ! removes)
8484
85+ (* * {1 Debug} *)
86+
87+ let debug_enabled = ref false
88+ let set_debug b = debug_enabled := b
89+
8590(* * {1 Node Registry} *)
8691
8792module Registry = struct
@@ -247,6 +252,26 @@ module Registry = struct
247252 let print_stats () =
248253 let all = Hashtbl. fold (fun _ info acc -> info :: acc) nodes [] in
249254 let sorted = List. sort (fun a b -> compare a.level b.level) all in
255+ let by_time =
256+ List. sort
257+ (fun a b ->
258+ Int64. compare b.stats.process_time_ns a.stats.process_time_ns)
259+ all
260+ in
261+ let top =
262+ by_time
263+ |> List. filter (fun info -> info.stats.process_time_ns <> 0L )
264+ |> List. filteri (fun i _ -> i < 5 )
265+ in
266+ if top <> [] then (
267+ Printf. eprintf " Top nodes by process time:\n " ;
268+ List. iter
269+ (fun info ->
270+ let time_ms = Int64. to_float info.stats.process_time_ns /. 1e6 in
271+ Printf. eprintf " - %s (L%d): %.2fms (runs=%d)\n " info.name info.level
272+ time_ms info.stats.process_count)
273+ top;
274+ Printf. eprintf " \n " );
250275 Printf. eprintf " Node statistics:\n " ;
251276 Printf. eprintf " %-30s | %8s %8s %5s %5s | %8s %8s %5s %5s | %5s %8s\n "
252277 " name" " d_recv" " e_recv" " +in" " -in" " d_emit" " e_emit" " +out" " -out"
@@ -273,13 +298,59 @@ module Scheduler = struct
273298
274299 let is_propagating () = ! propagating
275300
301+ type stats_snapshot = {
302+ deltas_received : int ;
303+ entries_received : int ;
304+ adds_received : int ;
305+ removes_received : int ;
306+ deltas_emitted : int ;
307+ entries_emitted : int ;
308+ adds_emitted : int ;
309+ removes_emitted : int ;
310+ process_count : int ;
311+ process_time_ns : int64 ;
312+ }
313+
314+ let snapshot_stats (s : stats ) : stats_snapshot =
315+ {
316+ deltas_received = s.deltas_received;
317+ entries_received = s.entries_received;
318+ adds_received = s.adds_received;
319+ removes_received = s.removes_received;
320+ deltas_emitted = s.deltas_emitted;
321+ entries_emitted = s.entries_emitted;
322+ adds_emitted = s.adds_emitted;
323+ removes_emitted = s.removes_emitted;
324+ process_count = s.process_count;
325+ process_time_ns = s.process_time_ns;
326+ }
327+
328+ let diff_stats (before : stats_snapshot ) (after_ : stats ) =
329+ let d_int x y = x - y in
330+ let d_time x y = Int64. sub x y in
331+ ( d_int after_.deltas_received before.deltas_received,
332+ d_int after_.entries_received before.entries_received,
333+ d_int after_.adds_received before.adds_received,
334+ d_int after_.removes_received before.removes_received,
335+ d_int after_.deltas_emitted before.deltas_emitted,
336+ d_int after_.entries_emitted before.entries_emitted,
337+ d_int after_.adds_emitted before.adds_emitted,
338+ d_int after_.removes_emitted before.removes_emitted,
339+ d_int after_.process_count before.process_count,
340+ d_time after_.process_time_ns before.process_time_ns )
341+
276342 (* * Process all dirty nodes in level order *)
277343 let propagate () =
278344 if ! propagating then
279345 failwith " Scheduler.propagate: already propagating (nested call)"
280346 else (
281347 propagating := true ;
282348 incr wave_counter;
349+ let wave_id = ! wave_counter in
350+ let wave_start = Unix. gettimeofday () in
351+ let processed_nodes = ref 0 in
352+ if ! debug_enabled then
353+ Printf. eprintf " \n === Reactive wave %d ===\n %!" wave_id;
283354
284355 while ! Registry. dirty_nodes <> [] do
285356 (* Get all dirty nodes, sort by level *)
@@ -319,17 +390,51 @@ module Scheduler = struct
319390 List. iter
320391 (fun (_ , _ , info ) ->
321392 info.Registry. dirty < - false ;
393+ let before =
394+ if ! debug_enabled then Some (snapshot_stats info.stats)
395+ else None
396+ in
322397 let start = Sys. time () in
323398 info.Registry. process () ;
324399 let elapsed = Sys. time () -. start in
325400 info.Registry. stats.process_time_ns < -
326401 Int64. add info.Registry. stats.process_time_ns
327402 (Int64. of_float (elapsed *. 1e9 ));
328403 info.Registry. stats.process_count < -
329- info.Registry. stats.process_count + 1 )
404+ info.Registry. stats.process_count + 1 ;
405+ if ! debug_enabled then (
406+ incr processed_nodes;
407+ match before with
408+ | None -> ()
409+ | Some b ->
410+ let ( d_recv,
411+ e_recv,
412+ add_in,
413+ rem_in,
414+ d_emit,
415+ e_emit,
416+ add_out,
417+ rem_out,
418+ runs,
419+ dt_ns ) =
420+ diff_stats b info.Registry. stats
421+ in
422+ (* runs should always be 1 here, but keep the check defensive *)
423+ if runs <> 0 then
424+ Printf. eprintf
425+ " %-30s (L%d): recv d/e/+/-=%d/%d/%d/%d emit \
426+ d/e/+/-=%d/%d/%d/%d time=%.2fms\n \
427+ %!"
428+ info.Registry. name info.Registry. level d_recv e_recv
429+ add_in rem_in d_emit e_emit add_out rem_out
430+ (Int64. to_float dt_ns /. 1e6 )))
330431 at_level
331432 done ;
332433
434+ (if ! debug_enabled then
435+ let wave_elapsed_ms = (Unix. gettimeofday () -. wave_start) *. 1000.0 in
436+ Printf. eprintf " Wave %d: processed_nodes=%d wall=%.2fms\n %!" wave_id
437+ ! processed_nodes wave_elapsed_ms);
333438 propagating := false )
334439
335440 let wave_count () = ! wave_counter
@@ -1181,5 +1286,6 @@ let fixpoint ~name ~(init : ('k, unit) t) ~(edges : ('k, 'k list) t) () :
11811286
11821287let to_mermaid () = Registry. to_mermaid ()
11831288let print_stats () = Registry. print_stats ()
1289+ let set_debug = set_debug
11841290let reset () = Registry. clear ()
11851291let reset_stats () = Registry. reset_stats ()
0 commit comments