Skip to content

Commit 94aa9af

Browse files
committed
feat: consolidate inspection tools
1 parent f47a8b4 commit 94aa9af

28 files changed

Lines changed: 735 additions & 329 deletions

docs/optimizations.md

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,10 @@ Cross-cutting rules for reactive/async integration (especially `patterns.ai`, LL
3434
| **`Node` resolution without `get()`** | When blocking until first `DATA`, prefer `node.get()` when it already holds a settled value, then subscribe only if still pending — avoids hangs when the node does not replay `DATA` to new subscribers. ||
3535
| **Passing plain strings through `fromAny` (TypeScript)** | `fromAny` treats strings as iterables (one `DATA` per character). For tool handlers that return plain strings, return the string directly; use `fromAny` only for `Node` / `AsyncIterable` / Promise-like after await. ||
3636

37+
- **`batch_id` in `ObserveEvent` timeline fields (decided 2026-04-08):** Added `batch_id?: number` / `batch_id: int` to `ObserveEvent` in both repos. Increments once per subscribe-callback invocation; all messages in one delivery share the same `batch_id`. Useful for correlating events that arrived together. TS: added to `_createObserveResult`, `_createObserveResultForAll`, and both fallback paths. PY: already present; now documented.
38+
39+
- **`ObserveResult.completedCleanly` ambiguous in graph-wide mode (noted 2026-04-08, carried from TS):** Same issue as TS. In graph-wide observation, `completed_cleanly` and `errored` can both be `True` simultaneously if one node completes cleanly before another errors. Single-node mode is unaffected. Options: (A) rename to `any_completed_cleanly` / `any_errored`; (B) add `all_completed_cleanly`; (C) reset on any ERROR. Pending decision.
40+
3741
---
3842

3943
## Deferred follow-ups

docs/roadmap.md

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -47,6 +47,24 @@ TS has this shipped; Python needs parity for production reduction pipelines.
4747
- [ ] PyPI publish: `graphrefly-py` (7 → 9.3)
4848
- [ ] Docs site at `py.graphrefly.dev` (7 → 9.3)
4949

50+
### Inspection Tool Consolidation (cross-cutting, TS + PY)
51+
52+
PY consolidation — reduce inspection surface from 14+ tools to 9 with clear, non-overlapping responsibilities.
53+
54+
**Design reference:** `archive/docs/SESSION-inspection-consolidation.md` (in graphrefly-ts)
55+
56+
#### PY consolidation (DONE 2026-04-08)
57+
58+
- [x] Merge `spy()` into `observe(format=)``format="pretty"|"json"` on `observe()` replaces `spy()`
59+
- [x] Add `trace()` (write + read), absorb `annotate()` + `trace_log()`
60+
- [x] Unexport `describe_node`, `meta_snapshot` from public API (internal to `core.meta`)
61+
- [x] `harness_trace()` pipeline stage tracer (`patterns/harness/trace.py`)
62+
- [x] Runner `__repr__` with `_scheduled`/`_completed` counters on `AsyncioRunner` and `TrioRunner`
63+
64+
#### PY parity items already done
65+
66+
- [x] `Graph.diff()` — already implemented
67+
5068
### Deferred (post-Wave 3)
5169

5270
- §7.2 Showcase demos (Pyodide/WASM lab) — after TS demos prove the pattern

src/graphrefly/__init__.py

Lines changed: 0 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -37,7 +37,6 @@
3737
defer_down,
3838
defer_set,
3939
derived,
40-
describe_node,
4140
dispatch_messages,
4241
down_with_batch,
4342
dynamic_node,
@@ -47,7 +46,6 @@
4746
is_batching,
4847
is_phase2_message,
4948
is_v1,
50-
meta_snapshot,
5149
monotonic_ns,
5250
node,
5351
normalize_actor,
@@ -74,7 +72,6 @@
7472
GraphDiffResult,
7573
GraphObserveSource,
7674
ObserveResult,
77-
SpyHandle,
7875
TraceEntry,
7976
reachable,
8077
)
@@ -93,7 +90,6 @@
9390
"NodeVersionInfo",
9491
"ObserveResult",
9592
"PATH_SEP",
96-
"SpyHandle",
9793
"TraceEntry",
9894
"V0",
9995
"V1",
@@ -138,13 +134,11 @@
138134
"compose_guards",
139135
"defer_down",
140136
"defer_set",
141-
"describe_node",
142137
"dispatch_messages",
143138
"down_with_batch",
144139
"ensure_registered",
145140
"is_batching",
146141
"is_phase2_message",
147-
"meta_snapshot",
148142
"node",
149143
"normalize_actor",
150144
"partition_for_batch",

src/graphrefly/compat/asyncio_runner.py

Lines changed: 26 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@
33
from __future__ import annotations
44

55
import asyncio
6+
import threading
67
from typing import TYPE_CHECKING, Any
78

89
if TYPE_CHECKING:
@@ -29,10 +30,12 @@ async def main():
2930
asyncio.run(main())
3031
"""
3132

32-
__slots__ = ("_loop",)
33+
__slots__ = ("_loop", "_scheduled", "_completed")
3334

3435
def __init__(self, loop: asyncio.AbstractEventLoop) -> None:
3536
self._loop = loop
37+
self._scheduled = 0
38+
self._completed = 0
3639

3740
@classmethod
3841
def from_running(cls) -> AsyncioRunner:
@@ -50,11 +53,12 @@ def schedule(
5053
on_error: Callable[[BaseException], None],
5154
) -> Callable[[], None]:
5255
task: asyncio.Task[Any] | None = None
53-
cancelled = False
56+
cancelled = threading.Event()
5457

5558
def _create_task() -> None:
5659
nonlocal task
57-
if cancelled:
60+
if cancelled.is_set():
61+
self._completed += 1
5862
coro.close()
5963
return
6064

@@ -71,19 +75,34 @@ async def _wrapper() -> None:
7175
on_error(err)
7276
else:
7377
on_result(result)
78+
finally:
79+
self._completed += 1
7480

7581
task = self._loop.create_task(_wrapper())
7682

77-
# Thread-safe: schedule task creation on the event loop.
78-
self._loop.call_soon_threadsafe(_create_task)
83+
# Increment eagerly on the calling thread so __repr__ is always consistent.
84+
self._scheduled += 1
85+
try:
86+
self._loop.call_soon_threadsafe(_create_task)
87+
except RuntimeError:
88+
# Loop closed — cannot schedule; close the coroutine and balance the counter.
89+
self._completed += 1
90+
coro.close()
7991

8092
def cancel() -> None:
81-
nonlocal cancelled
82-
cancelled = True
93+
cancelled.set()
8394
if task is not None:
8495
task.cancel()
8596

8697
return cancel
8798

99+
def __repr__(self) -> str:
100+
pending = self._scheduled - self._completed
101+
running = self._loop.is_running()
102+
return (
103+
f"AsyncioRunner(scheduled={self._scheduled}, completed={self._completed}, "
104+
f"pending={pending}, loop_running={running})"
105+
)
106+
88107

89108
__all__ = ["AsyncioRunner"]

src/graphrefly/compat/trio_runner.py

Lines changed: 19 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@
88
from graphrefly.compat.trio_runner import TrioRunner
99
from graphrefly.core.runner import set_default_runner
1010
11+
1112
async def main():
1213
async with trio.open_nursery() as nursery:
1314
runner = TrioRunner(nursery)
@@ -19,6 +20,7 @@ async def main():
1920

2021
from __future__ import annotations
2122

23+
import threading
2224
from typing import TYPE_CHECKING, Any
2325

2426
if TYPE_CHECKING:
@@ -34,10 +36,12 @@ class TrioRunner:
3436
Cancel scopes provide best-effort cancellation.
3537
"""
3638

37-
__slots__ = ("_nursery",)
39+
__slots__ = ("_nursery", "_scheduled", "_completed")
3840

3941
def __init__(self, nursery: trio.Nursery) -> None:
4042
self._nursery = nursery
43+
self._scheduled = 0
44+
self._completed = 0
4145

4246
def schedule(
4347
self,
@@ -48,11 +52,12 @@ def schedule(
4852
import trio as _trio
4953

5054
cancel_scope = _trio.CancelScope()
51-
cancelled = False
55+
cancelled = threading.Event()
5256

5357
async def _wrapper() -> None:
5458
with cancel_scope:
55-
if cancelled:
59+
if cancelled.is_set():
60+
self._completed += 1
5661
coro.close()
5762
return
5863
try:
@@ -67,15 +72,24 @@ async def _wrapper() -> None:
6772
on_error(err)
6873
else:
6974
on_result(result)
75+
finally:
76+
self._completed += 1
7077

78+
self._scheduled += 1
7179
self._nursery.start_soon(_wrapper)
7280

7381
def cancel() -> None:
74-
nonlocal cancelled
75-
cancelled = True
82+
cancelled.set()
7683
cancel_scope.cancel()
7784

7885
return cancel
7986

87+
def __repr__(self) -> str:
88+
pending = self._scheduled - self._completed
89+
return (
90+
f"TrioRunner(scheduled={self._scheduled}, completed={self._completed}, "
91+
f"pending={pending})"
92+
)
93+
8094

8195
__all__ = ["TrioRunner"]

src/graphrefly/core/__init__.py

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,6 @@
1717
record_mutation,
1818
system_actor,
1919
)
20-
from graphrefly.core.meta import describe_node, meta_snapshot
2120
from graphrefly.core.node import (
2221
NO_VALUE,
2322
Node,
@@ -112,7 +111,6 @@
112111
"compose_guards",
113112
"defer_down",
114113
"defer_set",
115-
"describe_node",
116114
"dispatch_messages",
117115
"down_with_batch",
118116
"ensure_registered",
@@ -121,7 +119,6 @@
121119
"is_terminal_message",
122120
"message_tier",
123121
"propagates_to_meta",
124-
"meta_snapshot",
125122
"node",
126123
"normalize_actor",
127124
"partition_for_batch",

src/graphrefly/core/node.py

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -394,11 +394,11 @@ def _handle_local_lifecycle(self, messages: Messages) -> None:
394394
continue
395395
if lock is not None:
396396
with lock:
397-
self._cached = m[1] # type: ignore[misc]
397+
self._cached = m[1]
398398
else:
399-
self._cached = m[1] # type: ignore[misc]
399+
self._cached = m[1]
400400
if self._versioning is not None:
401-
advance_version(self._versioning, m[1], self._hash_fn) # type: ignore[misc]
401+
advance_version(self._versioning, m[1], self._hash_fn)
402402
if t is MessageType.INVALIDATE:
403403
# GRAPHREFLY-SPEC §1.2: clear cached state; do not auto-emit from here.
404404
if self._cleanup is not None:

src/graphrefly/extra/cascading_cache.py

Lines changed: 17 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@
33
from __future__ import annotations
44

55
from collections import OrderedDict
6-
from typing import TYPE_CHECKING, Any, Generic, Protocol, TypeVar, runtime_checkable
6+
from typing import TYPE_CHECKING, Any, Generic, Protocol, TypeVar, cast, runtime_checkable
77

88
if TYPE_CHECKING:
99
from collections.abc import Sequence
@@ -137,10 +137,20 @@ def _tier_has_save(tier: CacheTier) -> bool:
137137
return callable(getattr(tier, "save", None))
138138

139139

140+
def _tier_save(tier: CacheTier, key: str, value: Any) -> None:
141+
"""Call ``tier.save(key, value)`` — caller must guard with ``_tier_has_save`` first."""
142+
cast("Any", tier).save(key, value)
143+
144+
140145
def _tier_has_clear(tier: CacheTier) -> bool:
141146
return callable(getattr(tier, "clear", None))
142147

143148

149+
def _tier_clear(tier: CacheTier, key: str) -> None:
150+
"""Call ``tier.clear(key)`` — caller must guard with ``_tier_has_clear`` first."""
151+
cast("Any", tier).clear(key)
152+
153+
144154
class CascadingCache(Generic[V]): # noqa: UP046
145155
"""Multi-tier cache where each entry is a ``state()`` node.
146156
@@ -174,7 +184,7 @@ def _promote(self, key: str, value: V, hit_tier: int) -> None:
174184
for i in range(hit_tier):
175185
tier = self._tiers[i]
176186
if _tier_has_save(tier):
177-
tier.save(key, value)
187+
_tier_save(tier, key, value)
178188

179189
def _cascade(self, key: str, nd: Node[Any]) -> None:
180190
for tier_index, tier in enumerate(self._tiers):
@@ -202,10 +212,10 @@ def _evict_if_needed(self) -> None:
202212
# Demote to deepest tier with save before evicting
203213
for i in range(len(self._tiers) - 1, -1, -1):
204214
if _tier_has_save(self._tiers[i]):
205-
self._tiers[i].save(victim, value)
215+
_tier_save(self._tiers[i], victim, value)
206216
for j in range(i):
207217
if _tier_has_clear(self._tiers[j]):
208-
self._tiers[j].clear(victim)
218+
_tier_clear(self._tiers[j], victim)
209219
break
210220
nd.down([(MessageType.TEARDOWN,)])
211221
del self._entries[victim]
@@ -239,9 +249,9 @@ def save(self, key: str, value: V) -> None:
239249
if self._write_through:
240250
for tier in self._tiers:
241251
if _tier_has_save(tier):
242-
tier.save(key, value)
252+
_tier_save(tier, key, value)
243253
elif self._tiers and _tier_has_save(self._tiers[0]):
244-
self._tiers[0].save(key, value)
254+
_tier_save(self._tiers[0], key, value)
245255
if key in self._entries:
246256
self._entries[key].down([(MessageType.DATA, value)])
247257
if self._eviction is not None:
@@ -273,7 +283,7 @@ def delete(self, key: str) -> None:
273283
self._eviction.delete(key)
274284
for tier in self._tiers:
275285
if _tier_has_clear(tier):
276-
tier.clear(key)
286+
_tier_clear(tier, key)
277287

278288
def has(self, key: str) -> bool:
279289
"""Check if a key is in the in-memory entries."""

src/graphrefly/extra/data_structures.py

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -145,9 +145,7 @@ def _prune_expired_internal(self) -> bool:
145145
"""Remove expired keys from ``_store``; return whether anything changed."""
146146
now = _mono_now()
147147
expired = [
148-
k
149-
for k, (_, exp) in self._store.entries.items()
150-
if exp is not None and exp <= now
148+
k for k, (_, exp) in self._store.entries.items() if exp is not None and exp <= now
151149
]
152150
if not expired:
153151
return False

src/graphrefly/graph/__init__.py

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -25,7 +25,6 @@
2525
GraphDiffResult,
2626
GraphObserveSource,
2727
ObserveResult,
28-
SpyHandle,
2928
TraceEntry,
3029
reachable,
3130
)
@@ -56,7 +55,6 @@
5655
"NodeProfile",
5756
"ObserveResult",
5857
"PATH_SEP",
59-
"SpyHandle",
6058
"TraceEntry",
6159
"WALEntry",
6260
"create_dag_cbor_codec",

0 commit comments

Comments
 (0)