diff --git a/docs/channel-plugin-guide.md b/docs/channel-plugin-guide.md index 179f2c09..9fbdf978 100644 --- a/docs/channel-plugin-guide.md +++ b/docs/channel-plugin-guide.md @@ -262,7 +262,7 @@ class OutboundMessage: event: object | None # typed runtime/UI event; usually inspect with isinstance() ``` -Runtime/UI semantics live on `msg.event`. Plugin-authored outbound messages should not rely on legacy metadata flags such as `_progress`, `_stream_delta`, `_stream_end`, `_reasoning_delta`, `_turn_end`, or `_goal_status`; they are treated as ordinary metadata, not parsed as runtime events. +Runtime/UI semantics live on `msg.event`. Plugin-authored outbound messages should use typed events instead of legacy metadata flags such as `_progress`, `_stream_delta`, `_stream_end`, `_reasoning_delta`, `_turn_end`, or `_goal_status`. nanobot still accepts those old flags as a compatibility bridge for existing in-process extensions, but new plugin code should not add fresh dependencies on them. ## Streaming Support diff --git a/nanobot/bus/outbound_events.py b/nanobot/bus/outbound_events.py index 2a5a5ba0..ed582c35 100644 --- a/nanobot/bus/outbound_events.py +++ b/nanobot/bus/outbound_events.py @@ -103,7 +103,9 @@ def outbound_message_for_event( def outbound_event_from_message(msg: OutboundMessage) -> OutboundEvent | None: """Return the typed outbound event carried by *msg*, if any.""" - return msg.event + if msg.event is not None: + return msg.event + return _legacy_event_from_metadata(msg) def replace_outbound_event( @@ -125,3 +127,100 @@ def _event_content(event: OutboundEvent) -> str: if isinstance(event, ProgressEvent | RetryWaitEvent | StreamDeltaEvent | StreamEndEvent): return event.content return "" + + +def _legacy_event_from_metadata(msg: OutboundMessage) -> OutboundEvent | None: + """Bridge pre-typed outbound metadata flags into typed events. + + New code should set ``OutboundMessage.event`` directly. The fallback keeps + older in-process extensions and channel plugins from losing runtime events + while they migrate off reserved metadata flags. + """ + + meta = msg.metadata or {} + if meta.get("_runtime_model_updated"): + return RuntimeModelUpdatedEvent( + model=_metadata_str(meta, "model"), + model_preset=_metadata_str(meta, "model_preset"), + ) + if meta.get("_goal_state_sync"): + goal_state = meta.get("goal_state") + return GoalStateSyncEvent(goal_state if isinstance(goal_state, dict) else {"active": False}) + if meta.get("_goal_status"): + status = meta.get("goal_status") + if not isinstance(status, str) or not status: + return None + return GoalStatusEvent( + status=status, + started_at=_metadata_float(meta, "started_at", "goal_started_at"), + ) + if meta.get("_turn_end"): + goal_state = meta.get("goal_state") + return TurnEndEvent( + latency_ms=_metadata_int(meta, "latency_ms"), + goal_state=goal_state if isinstance(goal_state, dict) else None, + ) + if meta.get("_session_updated"): + return SessionUpdatedEvent(scope=_metadata_str(meta, "_session_update_scope")) + if meta.get("_retry_wait"): + return RetryWaitEvent(content=msg.content) + if meta.get("_stream_end"): + return StreamEndEvent( + content=msg.content, + stream_id=_metadata_str(meta, "_stream_id"), + resuming=bool(meta.get("_resuming")), + ) + if meta.get("_stream_delta"): + return StreamDeltaEvent( + content=msg.content, + stream_id=_metadata_str(meta, "_stream_id"), + ) + if meta.get("_streamed"): + return StreamedResponseEvent() + if ( + meta.get("_progress") + or meta.get("_reasoning_delta") + or meta.get("_reasoning_end") + or meta.get("_reasoning") + or meta.get("_file_edit_events") + or meta.get("_tool_events") + ): + tool_events = meta.get("_tool_events") + file_edit_events = meta.get("_file_edit_events") + return ProgressEvent( + content=msg.content, + tool_hint=bool(meta.get("_tool_hint")), + reasoning=bool(meta.get("_reasoning")), + reasoning_delta=bool(meta.get("_reasoning_delta")), + reasoning_end=bool(meta.get("_reasoning_end")), + stream_id=_metadata_str(meta, "_stream_id"), + tool_events=tool_events if isinstance(tool_events, list) else None, + file_edit_events=file_edit_events if isinstance(file_edit_events, list) else None, + ) + return None + + +def _metadata_str(meta: Mapping[str, Any], key: str) -> str | None: + value = meta.get(key) + return value if isinstance(value, str) and value else None + + +def _metadata_int(meta: Mapping[str, Any], key: str) -> int | None: + value = meta.get(key) + if isinstance(value, bool): + return None + if isinstance(value, int): + return value + if isinstance(value, float) and value.is_integer(): + return int(value) + return None + + +def _metadata_float(meta: Mapping[str, Any], *keys: str) -> float | None: + for key in keys: + value = meta.get(key) + if isinstance(value, bool): + continue + if isinstance(value, int | float): + return float(value) + return None diff --git a/tests/bus/test_outbound_events.py b/tests/bus/test_outbound_events.py index d7658b17..0beb7e60 100644 --- a/tests/bus/test_outbound_events.py +++ b/tests/bus/test_outbound_events.py @@ -2,10 +2,16 @@ from __future__ import annotations from nanobot.bus.events import OutboundMessage from nanobot.bus.outbound_events import ( + GoalStateSyncEvent, + GoalStatusEvent, ProgressEvent, + RetryWaitEvent, + RuntimeModelUpdatedEvent, + SessionUpdatedEvent, StreamDeltaEvent, StreamedResponseEvent, StreamEndEvent, + TurnEndEvent, outbound_event_from_message, outbound_message_for_event, replace_outbound_event, @@ -49,20 +55,160 @@ def test_normal_outbound_message_has_no_runtime_event() -> None: assert outbound_event_from_message(msg) is None -def test_metadata_flags_do_not_create_runtime_events() -> None: +def test_legacy_progress_metadata_flags_create_runtime_event() -> None: + tool_events = [{"phase": "start", "name": "read_file"}] + file_edit_events = [{"phase": "end", "path": "app.py"}] msg = OutboundMessage( channel="websocket", chat_id="chat-1", content="legacy progress", metadata={ "_progress": True, - "_stream_delta": True, - "_goal_status": True, + "_tool_hint": True, + "_reasoning_delta": True, + "_stream_id": "r1", + "_tool_events": tool_events, + "_file_edit_events": file_edit_events, "message_id": "platform-routing-context", }, ) - assert outbound_event_from_message(msg) is None + event = outbound_event_from_message(msg) + assert isinstance(event, ProgressEvent) + assert event.content == "legacy progress" + assert event.tool_hint is True + assert event.reasoning_delta is True + assert event.stream_id == "r1" + assert event.tool_events == tool_events + assert event.file_edit_events == file_edit_events + + +def test_legacy_stream_metadata_flags_create_runtime_events() -> None: + delta = OutboundMessage( + channel="websocket", + chat_id="chat-1", + content="hello", + metadata={"_stream_delta": True, "_stream_id": "s1"}, + ) + end = OutboundMessage( + channel="websocket", + chat_id="chat-1", + content="", + metadata={"_stream_end": True, "_stream_id": "s1", "_resuming": True}, + ) + + delta_event = outbound_event_from_message(delta) + assert isinstance(delta_event, StreamDeltaEvent) + assert delta_event.content == "hello" + assert delta_event.stream_id == "s1" + + end_event = outbound_event_from_message(end) + assert isinstance(end_event, StreamEndEvent) + assert end_event.stream_id == "s1" + assert end_event.resuming is True + + +def test_legacy_webui_runtime_metadata_flags_create_runtime_events() -> None: + runtime = OutboundMessage( + channel="websocket", + chat_id="*", + content="", + metadata={ + "_runtime_model_updated": True, + "model": "gpt-5.5", + "model_preset": "high", + }, + ) + goal_state = OutboundMessage( + channel="websocket", + chat_id="chat-1", + content="", + metadata={"_goal_state_sync": True, "goal_state": {"active": True}}, + ) + goal_status = OutboundMessage( + channel="websocket", + chat_id="chat-1", + content="", + metadata={"_goal_status": True, "goal_status": "running", "started_at": 1.25}, + ) + turn_end = OutboundMessage( + channel="websocket", + chat_id="chat-1", + content="", + metadata={"_turn_end": True, "latency_ms": 42.0, "goal_state": {"active": False}}, + ) + session_updated = OutboundMessage( + channel="websocket", + chat_id="chat-1", + content="", + metadata={"_session_updated": True, "_session_update_scope": "metadata"}, + ) + + runtime_event = outbound_event_from_message(runtime) + assert isinstance(runtime_event, RuntimeModelUpdatedEvent) + assert runtime_event.model == "gpt-5.5" + assert runtime_event.model_preset == "high" + + goal_state_event = outbound_event_from_message(goal_state) + assert isinstance(goal_state_event, GoalStateSyncEvent) + assert goal_state_event.goal_state == {"active": True} + + goal_status_event = outbound_event_from_message(goal_status) + assert isinstance(goal_status_event, GoalStatusEvent) + assert goal_status_event.status == "running" + assert goal_status_event.started_at == 1.25 + + turn_end_event = outbound_event_from_message(turn_end) + assert isinstance(turn_end_event, TurnEndEvent) + assert turn_end_event.latency_ms == 42 + assert turn_end_event.goal_state == {"active": False} + + session_updated_event = outbound_event_from_message(session_updated) + assert isinstance(session_updated_event, SessionUpdatedEvent) + assert session_updated_event.scope == "metadata" + + +def test_legacy_metadata_numbers_ignore_bool_values() -> None: + goal_status = OutboundMessage( + channel="websocket", + chat_id="chat-1", + content="", + metadata={"_goal_status": True, "goal_status": "running", "started_at": True}, + ) + turn_end = OutboundMessage( + channel="websocket", + chat_id="chat-1", + content="", + metadata={"_turn_end": True, "latency_ms": True}, + ) + + goal_status_event = outbound_event_from_message(goal_status) + assert isinstance(goal_status_event, GoalStatusEvent) + assert goal_status_event.started_at is None + + turn_end_event = outbound_event_from_message(turn_end) + assert isinstance(turn_end_event, TurnEndEvent) + assert turn_end_event.latency_ms is None + + +def test_legacy_retry_wait_and_streamed_flags_create_runtime_events() -> None: + retry = OutboundMessage( + channel="cli", + chat_id="direct", + content="waiting", + metadata={"_retry_wait": True}, + ) + streamed = OutboundMessage( + channel="cli", + chat_id="direct", + content="final answer", + metadata={"_streamed": True}, + ) + + retry_event = outbound_event_from_message(retry) + assert isinstance(retry_event, RetryWaitEvent) + assert retry_event.content == "waiting" + assert isinstance(outbound_event_from_message(streamed), StreamedResponseEvent) def test_replace_outbound_event_keeps_routing_metadata() -> None: diff --git a/tests/channels/test_channel_manager_delta_coalescing.py b/tests/channels/test_channel_manager_delta_coalescing.py index 5224c803..688e133f 100644 --- a/tests/channels/test_channel_manager_delta_coalescing.py +++ b/tests/channels/test_channel_manager_delta_coalescing.py @@ -343,7 +343,7 @@ class TestProgressFiltering: assert send_mock.await_args_list[0].args[0].content == "final answer" @pytest.mark.asyncio - async def test_metadata_only_progress_flag_is_not_runtime_progress(self, manager, bus): + async def test_legacy_progress_flag_uses_runtime_progress_filter(self, manager, bus): manager.channels["mock"].send_progress = False await bus.publish_outbound(OutboundMessage( channel="mock", @@ -365,11 +365,7 @@ class TestProgressFiltering: except asyncio.CancelledError: pass - send_mock = manager.channels["mock"]._send_mock - assert send_mock.await_count == 1 - sent = send_mock.await_args_list[0].args[0] - assert sent.content == "legacy progress-shaped message" - assert sent.event is None + assert manager.channels["mock"]._send_mock.await_count == 0 @pytest.mark.asyncio async def test_channel_override_can_enable_tool_hints(self, manager, bus):