From 94436bb535be21cf2ec6abd07917d0ec1efaaf35 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=8F=AD=E6=89=AC?= Date: Fri, 7 Aug 2026 16:23:28 +0800 Subject: [PATCH 1/6] feat(stream): implement realtime streaming signal pipeline (Layer 0) - Add monitor signal_producer/event_bridge/signal_metrics for event-driven signal observation - Add backpressure (bounded queue + overflow accounting) to CompositeEventSource - Add capacity-bounded ReorderBuffer with drop accounting - Wire signal flow through event_bus, observers, normalizer, perception, causal pipeline - Extend monitor_coordinator/daemon service with signal watch RPC and event-driven watches - Add dashboard signals template, domain registration, and live signal rendering - Add tests/mock_signals simulator and new unit tests (event bridge, event-driven watch, buffer overflow, reorder capacity) - Bump version to 0.0.9 --- src/leapflow/cache/manager.py | 6 +- src/leapflow/causal/pipeline.py | 11 + src/leapflow/cli/commands/interactive.py | 3 + src/leapflow/cli/commands/registry.py | 1 + src/leapflow/cli/commands/slash_handlers.py | 26 +- src/leapflow/cli/tui_app/status.py | 7 +- src/leapflow/daemon/client.py | 10 + src/leapflow/daemon/monitor_coordinator.py | 128 ++++- src/leapflow/daemon/protocol.py | 1 + src/leapflow/daemon/service.py | 34 ++ src/leapflow/dashboard/server.py | 2 +- src/leapflow/dashboard/service.py | 51 +- src/leapflow/dashboard/static/app.js | 101 +++- src/leapflow/dashboard/static/styles.css | 20 + src/leapflow/dashboard/templates.py | 11 + src/leapflow/dashboard/templates/finance.yaml | 1 + .../dashboard/templates/research.yaml | 1 + .../dashboard/templates/sentiment.yaml | 1 + src/leapflow/dashboard/templates/signals.yaml | 90 ++++ src/leapflow/domain/events.py | 14 +- .../connectors/composite_event_source.py | 58 ++- .../gateway/connectors/event_sources.py | 43 +- src/leapflow/monitor/__init__.py | 4 + src/leapflow/monitor/event_bridge.py | 164 +++++++ src/leapflow/monitor/manager.py | 16 + src/leapflow/monitor/signal_metrics.py | 126 +++++ src/leapflow/monitor/signal_producer.py | 35 ++ src/leapflow/perception/signals.py | 14 +- src/leapflow/platform/event_bus.py | 36 +- src/leapflow/platform/normalizer.py | 33 +- src/leapflow/platform/observers/app_focus.py | 1 + src/leapflow/platform/observers/clipboard.py | 1 + src/leapflow/platform/observers/fs_watcher.py | 1 + src/leapflow/platform/observers/input_tap.py | 6 + src/leapflow/platform/reorder_buffer.py | 21 +- src/leapflow/scheduler/local_scheduler.py | 22 +- src/leapflow/version.py | 2 +- tests/README.md | 62 +++ tests/mock_signals/__init__.py | 29 ++ tests/mock_signals/__main__.py | 97 ++++ tests/mock_signals/generators.py | 291 ++++++++++++ tests/mock_signals/profiles.py | 74 +++ tests/mock_signals/runner.py | 446 ++++++++++++++++++ tests/test_cli_ndjson_event_source.py | 1 + tests/test_dashboard_view.py | 7 +- tests/test_dashboard_watch_rpc.py | 6 +- tests/test_event_bridge.py | 207 ++++++++ tests/test_event_driven_watch.py | 233 +++++++++ tests/test_gateway_consumer_loop.py | 82 ++++ tests/test_gateway_tool_e2e.py | 76 +++ tests/test_monitor_subsystem.py | 95 ++++ tests/test_reorder_buffer_capacity.py | 95 ++++ tests/test_signal_buffer_overflow.py | 58 +++ 53 files changed, 2912 insertions(+), 49 deletions(-) create mode 100644 src/leapflow/dashboard/templates/signals.yaml create mode 100644 src/leapflow/monitor/event_bridge.py create mode 100644 src/leapflow/monitor/signal_metrics.py create mode 100644 src/leapflow/monitor/signal_producer.py create mode 100644 tests/mock_signals/__init__.py create mode 100644 tests/mock_signals/__main__.py create mode 100644 tests/mock_signals/generators.py create mode 100644 tests/mock_signals/profiles.py create mode 100644 tests/mock_signals/runner.py create mode 100644 tests/test_event_bridge.py create mode 100644 tests/test_event_driven_watch.py create mode 100644 tests/test_reorder_buffer_capacity.py create mode 100644 tests/test_signal_buffer_overflow.py diff --git a/src/leapflow/cache/manager.py b/src/leapflow/cache/manager.py index 29c3549..aa65306 100644 --- a/src/leapflow/cache/manager.py +++ b/src/leapflow/cache/manager.py @@ -3,6 +3,7 @@ import hashlib import json +import threading import time from dataclasses import dataclass, field from enum import Enum @@ -51,6 +52,7 @@ class CacheManager: def __init__(self, layout: CacheLayout, *, profile_id: str) -> None: self._layout = layout self._profile_id = profile_id + self._connect_lock = threading.Lock() self._layout.ensure() self._init_schema() @@ -322,7 +324,9 @@ def _build_entry( ) def _connect(self): - return duckdb.connect(str(self._layout.index_path)) + """Create a DuckDB connection, serialized to protect against concurrent writes.""" + with self._connect_lock: + return duckdb.connect(str(self._layout.index_path)) def _is_managed_path(self, path: Path) -> bool: try: diff --git a/src/leapflow/causal/pipeline.py b/src/leapflow/causal/pipeline.py index bba0bc5..6064d07 100644 --- a/src/leapflow/causal/pipeline.py +++ b/src/leapflow/causal/pipeline.py @@ -57,6 +57,17 @@ class ReorderBuffer: Holds events for up to `window_s` seconds before releasing them in timestamp order. This handles platform-level delivery jitter (e.g., app_switch arriving 200ms after the keyboard shortcut that caused it). + + Timebase note: + This buffer sorts by ``CausalEvent.timestamp`` (wall-clock, ``time.time()`` + origin) because causal events are domain objects that may be serialized, + persisted, and correlated across sessions — monotonic clocks are not + comparable across processes or restarts. + + In contrast, :class:`~leapflow.platform.reorder_buffer.EventReorderBuffer` + sorts by ``payload["_mono_ts"]`` (``time.monotonic()``) because it operates + within a single process lifetime on raw observer events where monotonic + ordering is both available and more reliable than wall-clock. """ __slots__ = ("_window_s", "_buffer") diff --git a/src/leapflow/cli/commands/interactive.py b/src/leapflow/cli/commands/interactive.py index 63da6d8..cb68dac 100644 --- a/src/leapflow/cli/commands/interactive.py +++ b/src/leapflow/cli/commands/interactive.py @@ -1573,6 +1573,9 @@ async def _refresh_watch_count() -> None: app.invalidate() elif event_type == "watch.state": await _refresh_watch_count() + elif event_type == "signal.stream": + # Track signal flow activity for health visibility + status.increment_signal_stream() elif event_type == "monitor.error": logger.debug("monitor error notification: %s", payload) except (DaemonUnavailableError, OSError, asyncio.IncompleteReadError): diff --git a/src/leapflow/cli/commands/registry.py b/src/leapflow/cli/commands/registry.py index ca1dcda..f4ffdac 100644 --- a/src/leapflow/cli/commands/registry.py +++ b/src/leapflow/cli/commands/registry.py @@ -137,6 +137,7 @@ def supports_runtime(self, runtime: CommandRuntime) -> bool: # Board & Monitors (LeapBoard) — one analysis target (current session), # rendered through a selectable template lens. CommandDef("board", "Analyze the current session; optionally pick a template lens", "Board", args_hint="[