From 5fc006ea7fa275fcf9e3dac152ed51a06dddccec Mon Sep 17 00:00:00 2001 From: Jeffrey Phillips Freeman Date: Sun, 1 Feb 2026 21:08:28 -0500 Subject: [PATCH] Added some additional stream router coverage --- .../steps/stream_router_coverage_steps.py | 219 ++++++++++++++++++ features/stream_router_coverage.feature | 167 +++++++++++++ 2 files changed, 386 insertions(+) create mode 100644 features/steps/stream_router_coverage_steps.py create mode 100644 features/stream_router_coverage.feature diff --git a/features/steps/stream_router_coverage_steps.py b/features/steps/stream_router_coverage_steps.py new file mode 100644 index 000000000..88f921d67 --- /dev/null +++ b/features/steps/stream_router_coverage_steps.py @@ -0,0 +1,219 @@ +"""Steps for stream router coverage (actor-first). + +These steps focus on creation/error cases and lightweight operator helpers +compatible with the current ReactiveStreamRouter implementation. +""" + +from __future__ import annotations + +import json + +from behave import given, when, then + +from cleveragents.reactive.stream_router import ( + ReactiveStreamRouter, + StreamConfig, + StreamMessage, + StreamType, +) +from cleveragents.reactive.route_bridge import RouteBridge +from cleveragents.core.exceptions import StreamRoutingError + + +class DummyAgent: + def __init__(self, name: str): + self.name = name + + def process_message_sync(self, message, metadata=None): # type: ignore[override] + return f"echo:{message}" + + def get_capabilities(self): # pragma: no cover - not used in assertions + return [] + + +@given("I have a stream message for router testing") +def step_stream_message_testing(context): + context.stream_message = StreamMessage( + content="orig", metadata={"a": 1}, source_stream="s1", timestamp=1.0 + ) + context.original_msg = context.stream_message + + +@given("a reactive stream router harness") +def step_router_harness(context): + context.router = ReactiveStreamRouter() + context.errors = {} + context.last_messages = [] + + +@given("I have a replay stream configuration:") +@given("I have a hot stream configuration:") +@given("I have a stream configuration:") +@given("I have a stream configuration with agent mapper:") +@given("I have a stream configuration with missing agent:") +@given("I have a stream configuration with extract transform:") +@given("I have a stream configuration with unknown operator:") +@given("I have a stream configuration with LangGraph operator:") +@given("I have a stream configuration with equals filter:") +@given("I have a stream configuration with field condition:") +@given("I have a stream configuration with buffer and window:") +@given("I have a stream configuration with take skip merge:") +@given("I have a stream configuration with accumulate:") +@given("I have a stream configuration with switch operator:") +def step_load_stream_config(context): + cfg_text = context.text.strip() + cfg_dict = json.loads(cfg_text) + cfg_dict.setdefault("operators", []) + if "type" in cfg_dict: + cfg_dict["type"] = StreamType(cfg_dict["type"].lower()) + context.stream_config = StreamConfig(**cfg_dict) + + +@given("I have a stream for router testing") +def step_stream_for_testing(context): + context.stream_config = StreamConfig(name="test_stream") + context.router.create_stream(context.stream_config) + + +@given('I have a tool agent named "{name}" with echo tool') +def step_tool_agent(context, name): + agent = DummyAgent(name) + context.router.register_agent(name, agent) + + +@when("I create the reactive stream") +@given("I create the reactive stream") +def step_create_stream(context): + try: + context.router.create_stream(context.stream_config) + context.errors["create_stream"] = None + except Exception as exc: # pylint: disable=broad-except + context.errors["create_stream"] = exc + + +@when("I try to create the reactive stream") +@when("I try to create another stream with the same name") +def step_create_stream_error(context): + try: + context.router.create_stream(context.stream_config) + context.errors["create_stream"] = None + except Exception as exc: # pylint: disable=broad-except + context.errors["create_stream"] = exc + + +@when('I send router message "{message}" to stream "{name}"') +def step_send_message(context, message, name): + if name not in context.router.streams: + context.router.create_stream(StreamConfig(name=name)) + collected = [] + context.router.streams[name].subscribe(collected.append) + context.router.send_message(name, message) + context.last_messages = collected + + +@when('I send router None message to stream "{name}"') +def step_send_none_message(context, name): + step_send_message(context, None, name) + + +@when('I send a dictionary message with field "{field}"') +def step_send_dict_message(context, field): + name = context.stream_config.name + if name not in context.router.streams: + context.router.create_stream(context.stream_config) + collected = [] + context.router.streams[name].subscribe(collected.append) + context.router.send_message(name, {field: "extracted"}) + context.last_messages = collected + + +@when("I try to send message to missing stream") +def step_send_missing_stream(context): + try: + context.router.send_message("missing_stream", "data") + context.errors["send"] = None + except Exception as exc: # pylint: disable=broad-except + context.errors["send"] = exc + + +@when("I copy the message with modifications") +def step_copy_message(context): + msg = StreamMessage( + content="orig", metadata={"a": 1}, source_stream="s1", timestamp=1.0 + ) + context.original_msg = msg + context.copied_msg = msg.copy_with(content="new", metadata={"b": 2}) + + +@when("I dispose of the stream router") +def step_dispose_router(context): + context.router.dispose() + + +@when("I send a message to check timestamp") +def step_send_timestamp(context): + name = "ts_stream" + if name not in context.router.streams: + context.router.create_stream(StreamConfig(name=name)) + collected = [] + context.router.streams[name].subscribe(collected.append) + context.router.send_message(name, "payload") + context.last_messages = collected + + +@then("the stream should be created successfully") +def step_stream_created(context): + assert context.stream_config.name in context.router.streams + assert context.errors.get("create_stream") is None + + +@then("I should get a stream router coverage error") +@then("I should get a stream router coverage error about missing agent") +@then("I should get a stream router coverage error about unknown operator") +@then("I should get a stream router coverage error about LangGraph bridge") +def step_stream_error(context): + err = context.errors.get("create_stream") + assert isinstance(err, StreamRoutingError) + + +@then("the field value should be extracted") +def step_field_extracted(context): + # Directly apply transform to a dict message + transform = context.stream_config.operators[0]["params"]["transform"] + result = context.router._apply_transform({"data": "value"}, transform) # pylint: disable=protected-access + assert result == "value" + + +@then("the new message should have the modifications") +def step_new_message_has_mods(context): + assert context.copied_msg.content == "new" + assert context.copied_msg.metadata == {"b": 2} + + +@then("the original message should be unchanged") +def step_original_message_unchanged(context): + assert context.original_msg.content == "orig" + assert context.original_msg.metadata == {"a": 1} + + +@then("the agent should handle None gracefully") +def step_agent_handles_none(context): + # Agent mapper returns string even on None + mapper = context.router._create_agent_mapper(DummyAgent("none_agent")) # pylint: disable=protected-access + result = mapper(None) + assert isinstance(result, StreamMessage) + assert "echo:" in result.content + + +@then("the stream router should be disposed properly") +def step_disposed(context): + assert context.router.streams == {} + assert context.router.observables == {} + assert context.router.stream_configs == {} + + +@then("the message should have a timestamp") +def step_timestamp_present(context): + assert context.last_messages + msg = context.last_messages[-1] + assert getattr(msg, "timestamp", None) is not None diff --git a/features/stream_router_coverage.feature b/features/stream_router_coverage.feature new file mode 100644 index 000000000..2143fd8b6 --- /dev/null +++ b/features/stream_router_coverage.feature @@ -0,0 +1,167 @@ +Feature: Stream Router Coverage + As a developer + I want to exercise reactive stream router behaviors + So that stream_router.py is covered in actor-first mode + + Background: + Given a reactive stream router harness + + Scenario: Create replay stream with buffer size + Given I have a replay stream configuration: + """ + {"name": "replay_stream", "type": "replay", "buffer_size": 5, "operators": []} + """ + When I create the reactive stream + Then the stream should be created successfully + + Scenario: Create hot stream with initial value + Given I have a hot stream configuration: + """ + {"name": "hot_stream", "type": "hot", "initial_value": "initial_data", "operators": []} + """ + When I create the reactive stream + Then the stream should be created successfully + + Scenario: Stream creation with existing name should fail + Given I have a stream configuration: + """ + {"name": "existing_stream", "type": "cold", "operators": []} + """ + When I create the reactive stream + And I try to create another stream with the same name + Then I should get a stream router coverage error + + Scenario: Create agent mapper with tool agent + Given I have a tool agent named "echo_agent" with echo tool + And I have a stream configuration with agent mapper: + """ + {"name": "agent_stream", "type": "cold", "operators": [{"type": "map", "params": {"agent": "echo_agent"}}]} + """ + When I create the reactive stream + And I send router message "hello" to stream "agent_stream" + Then the stream should be created successfully + + + Scenario: Create agent mapper with missing agent should fail + Given I have a stream configuration with missing agent: + """ + {"name": "missing_agent_stream", "type": "cold", "operators": [{"type": "map", "params": {"agent": "nonexistent_agent"}}]} + """ + When I try to create the reactive stream + Then I should get a stream router coverage error about missing agent + + Scenario: Transform operator with extract_field transform + Given I have a stream configuration with extract transform: + """ + {"name": "extract_stream", "type": "cold", "operators": [{"type": "map", "params": {"transform": {"type": "extract_field", "field": "data"}}}]} + """ + When I create the reactive stream + And I send a dictionary message with field "data" + Then the field value should be extracted + + Scenario: Unknown operator should fail + Given I have a stream configuration with unknown operator: + """ + {"name": "unknown_stream", "type": "cold", "operators": [{"type": "unknown_operator"}]} + """ + When I try to create the reactive stream + Then I should get a stream router coverage error about unknown operator + + Scenario: LangGraph operators without bridge should fail + Given I have a stream configuration with LangGraph operator: + """ + {"name": "langgraph_stream", "type": "cold", "operators": [{"type": "graph_execute", "params": {}}]} + """ + When I try to create the reactive stream + Then I should get a stream router coverage error about LangGraph bridge + + Scenario: Stream message copy_with functionality + Given I have a stream message for router testing + When I copy the message with modifications + Then the new message should have the modifications + And the original message should be unchanged + + Scenario: Agent mapper with None message handling + Given I have a tool agent named "none_agent" with echo tool + And I have a stream configuration with agent mapper: + """ + {"name": "none_stream", "type": "cold", "operators": [{"type": "map", "params": {"agent": "none_agent"}}]} + """ + When I create the reactive stream + And I send router None message to stream "none_stream" + Then the agent should handle None gracefully + + + Scenario: Stream disposal + Given I have a stream for router testing + When I dispose of the stream router + Then the stream router should be disposed properly + + Scenario: Message with timestamp + Given I have a stream for router testing + When I send a message to check timestamp + Then the message should have a timestamp + + Scenario: Filter operator with equals condition + Given I have a stream configuration with equals filter: + """ + {"name": "equals_filter", "type": "cold", "operators": [{"type": "filter", "params": {"condition": {"equals": "match"}}}]} + """ + When I create the reactive stream + Then the stream should be created successfully + + Scenario: Filter operator with field presence condition + Given I have a stream configuration with field condition: + """ + {"name": "field_filter", "type": "cold", "operators": [{"type": "filter", "params": {"condition": {"field": "test_key"}}}]} + """ + When I create the reactive stream + Then the stream should be created successfully + + Scenario: Buffer and window operators + Given I have a stream configuration with buffer and window: + """ + {"name": "buffer_window", "type": "cold", "operators": [ + {"type": "buffer", "params": {"count": 2, "timeout": 0.01}}, + {"type": "window", "params": {"count": 1}} + ]} + """ + When I create the reactive stream + Then the stream should be created successfully + + Scenario: Take, skip, and merge operators + Given I have a stream configuration with take skip merge: + """ + {"name": "take_skip_merge", "type": "cold", "operators": [ + {"type": "take", "params": {"count": 1}}, + {"type": "skip", "params": {"count": 0}}, + {"type": "merge", "params": {"streams": ["__input__"]}} + ]} + """ + When I create the reactive stream + Then the stream should be created successfully + + Scenario: Accumulate operator + Given I have a stream configuration with accumulate: + """ + {"name": "accumulate_stream", "type": "cold", "operators": [ + {"type": "accumulate", "params": {"accumulator": {"op": "sum"}}} + ]} + """ + When I create the reactive stream + Then the stream should be created successfully + + Scenario: Switch/conditional routing operator + Given I have a stream configuration with switch operator: + """ + {"name": "switch_stream", "type": "cold", "operators": [ + {"type": "switch", "params": { + "cases": [ + {"condition": {"equals": "yes"}, "target": "__output__"} + ], + "default": "__output__" + }} + ]} + """ + When I create the reactive stream + Then the stream should be created successfully