a808c395f9
CI / benchmark-publish (pull_request) Has been skipped
CI / lint (pull_request) Successful in 15s
CI / build (pull_request) Successful in 15s
CI / quality (pull_request) Successful in 17s
CI / typecheck (pull_request) Successful in 35s
CI / security (pull_request) Successful in 50s
CI / unit_tests (pull_request) Successful in 2m46s
CI / integration_tests (pull_request) Successful in 3m16s
CI / docker (pull_request) Successful in 40s
CI / coverage (pull_request) Successful in 5m6s
CI / lint (push) Successful in 13s
CI / quality (push) Successful in 16s
CI / build (push) Successful in 18s
CI / security (push) Successful in 32s
CI / typecheck (push) Successful in 35s
CI / benchmark-regression (push) Has been skipped
CI / unit_tests (push) Successful in 2m52s
CI / integration_tests (push) Successful in 3m8s
CI / docker (push) Successful in 39s
CI / coverage (push) Successful in 5m53s
CI / benchmark-publish (push) Successful in 16m55s
CI / benchmark-regression (pull_request) Successful in 33m0s
Add 53 new .feature files and corresponding step definition files targeting uncovered lines identified in build/coverage.xml. Fix AmbiguousStep conflicts in 7 pre-existing step files by disambiguating step text. New tests cover: ACP clients/facade, actor CLI/config, application container, ACMS service/strategies, async worker, automation profile CLI, autonomy guardrail, bridge, change model, config CLI/service, context service, cross-plan correction, database models, decision service, decomposition clustering/service, discovery handler, langchain chat provider, langgraph nodes, materializers, multi-project service, plan apply/CLI/lifecycle/model/ preflight/resume/service, PostgreSQL analyzer, project CLI/context CLI, provider registry, reactive application/route, repositories, resolver handler, resource registry service, resume model, retry patterns, sandbox protocol, server CLI, skill CLI/service, skills registry, subplan execution/service, system CLI, UKO loader, UoW, and YAML template engine. Closes #645
471 lines
17 KiB
Python
471 lines
17 KiB
Python
"""Steps for remaining uncovered lines/branches in bridge.py.
|
|
|
|
Targets:
|
|
- Line 35: ``cancellation_reasons`` WeakKeyDictionary initialisation
|
|
- Line 68: ``task.cancel()`` inside ``cleanup_tasks_async``
|
|
- Lines 201-204: string-content branch in the inner ``execute_graph`` coroutine
|
|
- Branches 39->38 (__del__ normal exit), 67->68 (not-done cancel in async cleanup),
|
|
182->184 (valid graph in create_graph_stream), 260->258 (valid graph in checkpointer)
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
from typing import Any
|
|
from unittest.mock import AsyncMock, MagicMock
|
|
|
|
from behave import given, then, when
|
|
from behave.runner import Context
|
|
|
|
from cleveragents.langgraph.bridge import RxPyLangGraphBridge
|
|
from cleveragents.langgraph.state import GraphState
|
|
from cleveragents.reactive.stream_router import (
|
|
ReactiveStreamRouter,
|
|
StreamMessage,
|
|
)
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Helpers
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def _build_bridge() -> tuple[RxPyLangGraphBridge, MagicMock]:
|
|
"""Return a bridge backed by a mocked stream-router."""
|
|
scheduler = MagicMock()
|
|
router = ReactiveStreamRouter(scheduler=scheduler)
|
|
router.agents = {"agent": MagicMock()}
|
|
bridge = RxPyLangGraphBridge(router)
|
|
return bridge, scheduler
|
|
|
|
|
|
def _build_mock_graph(
|
|
name: str,
|
|
*,
|
|
messages: list[dict[str, Any]] | None = None,
|
|
state_dict: dict[str, Any] | None = None,
|
|
) -> MagicMock:
|
|
"""Build a mock LangGraph whose ``execute`` is an AsyncMock."""
|
|
graph = MagicMock()
|
|
graph.name = name
|
|
|
|
gs = MagicMock(spec=GraphState)
|
|
gs.messages = messages if messages is not None else []
|
|
gs.to_dict = MagicMock(return_value=state_dict or {"messages": gs.messages})
|
|
|
|
graph.execute = AsyncMock(return_value=gs)
|
|
graph.get_execution_history = MagicMock(return_value=["s1"])
|
|
graph.state_manager = MagicMock()
|
|
graph.state_manager.checkpoint_dir = None
|
|
graph.nodes = {}
|
|
return graph
|
|
|
|
|
|
# ===================================================================
|
|
# Scenario: Cancellation reason is recorded when cancelling a tracked task
|
|
# ===================================================================
|
|
|
|
|
|
@given("a bridge instance prepared for cancellation reason tracking")
|
|
def step_bridge_for_cancel_reason(context: Context) -> None:
|
|
bridge, _ = _build_bridge()
|
|
# Create an event loop and a dummy task
|
|
loop = asyncio.new_event_loop()
|
|
asyncio.set_event_loop(loop)
|
|
|
|
async def _forever() -> None:
|
|
await asyncio.sleep(9999)
|
|
|
|
task = loop.create_task(_forever())
|
|
bridge._active_tasks.add(task)
|
|
|
|
context.bridge = bridge
|
|
context.dummy_task = task
|
|
context.cancel_loop = loop
|
|
context.caught_error = None
|
|
|
|
def _cleanup() -> None:
|
|
# Cancel all pending tasks and close the loop to avoid __del__ errors
|
|
for t in list(bridge._active_tasks):
|
|
t.cancel()
|
|
bridge._active_tasks.clear()
|
|
if not loop.is_closed():
|
|
loop.run_until_complete(asyncio.sleep(0))
|
|
loop.close()
|
|
asyncio.set_event_loop(asyncio.new_event_loop())
|
|
|
|
context._cleanup_handlers.append(_cleanup)
|
|
|
|
|
|
@when('I cancel a tracked task providing reason "{reason}"')
|
|
def step_cancel_task_with_given_reason(context: Context, reason: str) -> None:
|
|
loop = context.cancel_loop
|
|
|
|
async def _do_cancel() -> None:
|
|
await context.bridge.cancel_task_with_reason(context.dummy_task, reason)
|
|
|
|
loop.run_until_complete(_do_cancel())
|
|
context.cancel_reason_used = reason
|
|
|
|
|
|
@then("the bridge cancellation_reasons mapping should contain the task")
|
|
def step_assert_cancellation_reasons_has_task(context: Context) -> None:
|
|
assert context.dummy_task in context.bridge.cancellation_reasons, (
|
|
"Expected task to be present in cancellation_reasons"
|
|
)
|
|
|
|
|
|
@then('the stored reason should equal "{expected}"')
|
|
def step_assert_stored_reason(context: Context, expected: str) -> None:
|
|
actual = context.bridge.cancellation_reasons.get(context.dummy_task)
|
|
assert actual == expected, f"Expected reason '{expected}', got '{actual}'"
|
|
# Cleanup
|
|
context.cancel_loop.close()
|
|
asyncio.set_event_loop(asyncio.new_event_loop())
|
|
|
|
|
|
# ===================================================================
|
|
# Scenario: Cancelling a None task raises ValueError
|
|
# ===================================================================
|
|
|
|
|
|
@when('I try cancelling a None task with reason "{reason}"')
|
|
def step_cancel_none_task(context: Context, reason: str) -> None:
|
|
loop = context.cancel_loop
|
|
|
|
async def _do_cancel() -> None:
|
|
await context.bridge.cancel_task_with_reason(None, reason) # type: ignore[arg-type]
|
|
|
|
try:
|
|
loop.run_until_complete(_do_cancel())
|
|
except ValueError as exc:
|
|
context.caught_error = exc
|
|
|
|
|
|
@then('a ValueError mentioning "{fragment}" should be raised')
|
|
def step_then_valueerror_mentioning(context: Context, fragment: str) -> None:
|
|
assert context.caught_error is not None, "Expected a ValueError but none was raised"
|
|
assert isinstance(context.caught_error, ValueError), (
|
|
f"Expected ValueError, got {type(context.caught_error).__name__}"
|
|
)
|
|
assert fragment in str(context.caught_error), (
|
|
f"Expected '{fragment}' in error message: {context.caught_error}"
|
|
)
|
|
|
|
|
|
# ===================================================================
|
|
# Scenario: Cancelling a task with empty reason raises ValueError
|
|
# ===================================================================
|
|
|
|
|
|
@when("I try cancelling a valid task with an empty reason string")
|
|
def step_cancel_task_empty_reason(context: Context) -> None:
|
|
loop = context.cancel_loop
|
|
|
|
async def _do_cancel() -> None:
|
|
await context.bridge.cancel_task_with_reason(context.dummy_task, "")
|
|
|
|
try:
|
|
loop.run_until_complete(_do_cancel())
|
|
except ValueError as exc:
|
|
context.caught_error = exc
|
|
|
|
|
|
# ===================================================================
|
|
# Scenario: Async cleanup cancels a not-yet-done task (line 68, branch 67->68)
|
|
# ===================================================================
|
|
|
|
|
|
@given("a bridge holding a deliberately stalled coroutine task")
|
|
def step_bridge_with_stalled_task(context: Context) -> None:
|
|
bridge, _ = _build_bridge()
|
|
loop = asyncio.new_event_loop()
|
|
asyncio.set_event_loop(loop)
|
|
|
|
async def _stall() -> None:
|
|
await asyncio.sleep(99999)
|
|
|
|
task = loop.create_task(_stall())
|
|
bridge._active_tasks.add(task)
|
|
|
|
context.bridge = bridge
|
|
context.stalled_task = task
|
|
context.async_cleanup_loop = loop
|
|
|
|
|
|
@when("I perform async cleanup with a 200ms timeout on stalled tasks")
|
|
def step_perform_async_cleanup(context: Context) -> None:
|
|
loop = context.async_cleanup_loop
|
|
loop.run_until_complete(context.bridge.cleanup_tasks_async(timeout=0.2))
|
|
|
|
|
|
@then("every stalled task should have been cancelled or finished")
|
|
def step_assert_stalled_tasks_cancelled(context: Context) -> None:
|
|
task = context.stalled_task
|
|
assert task.cancelled() or task.done(), (
|
|
f"Expected stalled task to be cancelled/done, state={task._state}"
|
|
)
|
|
|
|
|
|
@then("the bridge task tracking set should be empty after async cleanup")
|
|
def step_assert_tracking_set_empty(context: Context) -> None:
|
|
assert len(context.bridge._active_tasks) == 0, (
|
|
f"Expected empty _active_tasks, got {len(context.bridge._active_tasks)}"
|
|
)
|
|
# Cleanup loop
|
|
if (
|
|
hasattr(context, "async_cleanup_loop")
|
|
and not context.async_cleanup_loop.is_closed()
|
|
):
|
|
context.async_cleanup_loop.close()
|
|
asyncio.set_event_loop(asyncio.new_event_loop())
|
|
|
|
|
|
# ===================================================================
|
|
# Scenario: Async cleanup handles already-completed tasks
|
|
# ===================================================================
|
|
|
|
|
|
@given("a bridge holding an already-finished async task")
|
|
def step_bridge_with_finished_task(context: Context) -> None:
|
|
bridge, _ = _build_bridge()
|
|
loop = asyncio.new_event_loop()
|
|
asyncio.set_event_loop(loop)
|
|
|
|
async def _instant() -> str:
|
|
return "done"
|
|
|
|
task = loop.create_task(_instant())
|
|
# Let the task complete
|
|
loop.run_until_complete(task)
|
|
assert task.done(), "Task should be done before adding to bridge"
|
|
bridge._active_tasks.add(task)
|
|
|
|
context.bridge = bridge
|
|
context.async_cleanup_loop = loop
|
|
|
|
|
|
# ===================================================================
|
|
# Scenario: Graph executor inner coroutine with string content (lines 201-204)
|
|
# ===================================================================
|
|
|
|
|
|
@given("a bridge with a mock graph wired for direct coroutine invocation")
|
|
def step_bridge_for_direct_coroutine(context: Context) -> None:
|
|
bridge, _ = _build_bridge()
|
|
mock_graph = _build_mock_graph(
|
|
"direct_exec",
|
|
messages=[{"content": "coroutine-ok"}],
|
|
state_dict={"messages": [{"content": "coroutine-ok"}]},
|
|
)
|
|
bridge.graphs["direct_exec"] = mock_graph
|
|
context.bridge = bridge
|
|
context.direct_graph = mock_graph
|
|
context.exec_loop = asyncio.new_event_loop()
|
|
asyncio.set_event_loop(context.exec_loop)
|
|
|
|
|
|
@when("I invoke the executor coroutine with a plain string message body")
|
|
def step_invoke_executor_with_string(context: Context) -> None:
|
|
"""Push a string message through the executor and await the task directly."""
|
|
import rx # type: ignore
|
|
|
|
bridge = context.bridge
|
|
executor_op = bridge._create_graph_executor({"graph": "direct_exec"})
|
|
loop = context.exec_loop
|
|
|
|
results: list[Any] = []
|
|
errors: list[Any] = []
|
|
msg = StreamMessage(content="hello from string", metadata={"src": "test"})
|
|
|
|
rx.just(msg).pipe(executor_op).subscribe(
|
|
on_next=lambda x: results.append(x),
|
|
on_error=lambda e: errors.append(e),
|
|
)
|
|
|
|
# Drain all pending tasks by awaiting them directly
|
|
pending = list(bridge._active_tasks)
|
|
if pending:
|
|
loop.run_until_complete(asyncio.gather(*pending, return_exceptions=True))
|
|
else:
|
|
loop.run_until_complete(asyncio.sleep(0.05))
|
|
|
|
context.coroutine_results = results
|
|
context.coroutine_errors = errors
|
|
|
|
|
|
@then("the mock graph execute should have received a messages list with the string")
|
|
def step_assert_graph_received_string(context: Context) -> None:
|
|
assert len(context.coroutine_errors) == 0, (
|
|
f"Coroutine raised errors: {context.coroutine_errors}"
|
|
)
|
|
# Verify that graph.execute was called with input_data containing messages
|
|
call_args = context.direct_graph.execute.call_args
|
|
assert call_args is not None, "graph.execute was never called"
|
|
input_data = call_args[0][0] # First positional argument
|
|
assert "messages" in input_data, (
|
|
f"Expected 'messages' key in input_data, got keys: {list(input_data.keys())}"
|
|
)
|
|
assert input_data["messages"][0]["content"] == "hello from string", (
|
|
f"Expected string content in messages, got: {input_data['messages']}"
|
|
)
|
|
|
|
|
|
@then("the executor coroutine should return a StreamMessage with graph metadata")
|
|
def step_assert_coroutine_returns_stream_message(context: Context) -> None:
|
|
assert len(context.coroutine_results) > 0, "No results from executor coroutine"
|
|
result_msg = context.coroutine_results[0]
|
|
assert isinstance(result_msg, StreamMessage), (
|
|
f"Expected StreamMessage, got {type(result_msg).__name__}"
|
|
)
|
|
assert "graph" in result_msg.metadata, (
|
|
f"Expected 'graph' in metadata, keys: {list(result_msg.metadata.keys())}"
|
|
)
|
|
# Cleanup
|
|
if hasattr(context, "exec_loop") and not context.exec_loop.is_closed():
|
|
context.exec_loop.close()
|
|
asyncio.set_event_loop(asyncio.new_event_loop())
|
|
|
|
|
|
# ===================================================================
|
|
# Scenario: Graph executor inner coroutine with list content (else branch)
|
|
# ===================================================================
|
|
|
|
|
|
@when("I invoke the executor coroutine with a list payload as message body")
|
|
def step_invoke_executor_with_list(context: Context) -> None:
|
|
"""Push a list (non-str, non-dict) message through the executor."""
|
|
import rx # type: ignore
|
|
|
|
bridge = context.bridge
|
|
executor_op = bridge._create_graph_executor({"graph": "direct_exec"})
|
|
loop = context.exec_loop
|
|
|
|
results: list[Any] = []
|
|
errors: list[Any] = []
|
|
list_payload = ["item_a", "item_b", "item_c"]
|
|
msg = StreamMessage(content=list_payload, metadata={})
|
|
|
|
rx.just(msg).pipe(executor_op).subscribe(
|
|
on_next=lambda x: results.append(x),
|
|
on_error=lambda e: errors.append(e),
|
|
)
|
|
|
|
pending = list(bridge._active_tasks)
|
|
if pending:
|
|
loop.run_until_complete(asyncio.gather(*pending, return_exceptions=True))
|
|
else:
|
|
loop.run_until_complete(asyncio.sleep(0.05))
|
|
|
|
context.coroutine_results = results
|
|
context.coroutine_errors = errors
|
|
|
|
|
|
@then("the mock graph execute should have received a content key wrapping the list")
|
|
def step_assert_graph_received_list_wrapped(context: Context) -> None:
|
|
assert len(context.coroutine_errors) == 0, (
|
|
f"Coroutine raised errors: {context.coroutine_errors}"
|
|
)
|
|
call_args = context.direct_graph.execute.call_args
|
|
assert call_args is not None, "graph.execute was never called"
|
|
input_data = call_args[0][0]
|
|
assert "content" in input_data, (
|
|
f"Expected 'content' key in input_data for list fallback, got: {list(input_data.keys())}"
|
|
)
|
|
assert input_data["content"] == ["item_a", "item_b", "item_c"], (
|
|
f"Expected list content, got: {input_data['content']}"
|
|
)
|
|
# Cleanup
|
|
if hasattr(context, "exec_loop") and not context.exec_loop.is_closed():
|
|
context.exec_loop.close()
|
|
asyncio.set_event_loop(asyncio.new_event_loop())
|
|
|
|
|
|
# ===================================================================
|
|
# Scenario: Graph stream config with correct type/publication (branch 182->184)
|
|
# ===================================================================
|
|
|
|
|
|
@given('a bridge with a registered graph called "{name}"')
|
|
def step_bridge_with_named_registered_graph(context: Context, name: str) -> None:
|
|
bridge, _ = _build_bridge()
|
|
mock_graph = _build_mock_graph(name)
|
|
bridge.graphs[name] = mock_graph
|
|
context.bridge = bridge
|
|
|
|
|
|
@when('I obtain the stream configuration for graph "{name}"')
|
|
def step_obtain_stream_config(context: Context, name: str) -> None:
|
|
context.obtained_stream_config = context.bridge.create_graph_stream(name)
|
|
|
|
|
|
@then('the stream config name should equal "{expected}"')
|
|
def step_assert_stream_name_equals(context: Context, expected: str) -> None:
|
|
actual = context.obtained_stream_config.name
|
|
assert actual == expected, f"Expected stream name '{expected}', got '{actual}'"
|
|
|
|
|
|
@then('the stream config publications should include "{publication}"')
|
|
def step_assert_publications_include(context: Context, publication: str) -> None:
|
|
pubs = context.obtained_stream_config.publications
|
|
assert publication in pubs, f"Expected '{publication}' in publications {pubs}"
|
|
|
|
|
|
# ===================================================================
|
|
# Scenario: State checkpointer operator for valid graph (branch 260->258)
|
|
# ===================================================================
|
|
|
|
|
|
@given('the graph "{name}" has a checkpoint directory configured')
|
|
def step_configure_checkpoint_dir(context: Context, name: str) -> None:
|
|
graph = context.bridge.graphs[name]
|
|
graph.state_manager.checkpoint_dir = "/tmp/test_ckpt"
|
|
|
|
|
|
@when('I construct and apply the checkpointer operator for "{name}"')
|
|
def step_construct_and_apply_checkpointer(context: Context, name: str) -> None:
|
|
import rx # type: ignore
|
|
|
|
checkpointer_op = context.bridge._create_state_checkpointer({"graph": name})
|
|
msg = StreamMessage(content="ckpt-payload", metadata={"check": True})
|
|
results: list[Any] = []
|
|
rx.just(msg).pipe(checkpointer_op).subscribe(
|
|
on_next=lambda x: results.append(x),
|
|
)
|
|
context.ckpt_results = results
|
|
|
|
|
|
@then("the checkpointer should have invoked _save_checkpoint on the state manager")
|
|
def step_assert_save_checkpoint_called(context: Context) -> None:
|
|
assert len(context.ckpt_results) > 0, "Checkpointer produced no output"
|
|
# The graph "ckpt_valid" should have had _save_checkpoint called
|
|
graph = context.bridge.graphs.get("ckpt_valid")
|
|
if graph is not None:
|
|
graph.state_manager._save_checkpoint.assert_called()
|
|
|
|
|
|
# ===================================================================
|
|
# Scenario: __del__ completes normally (branch 39->38 normal exit)
|
|
# ===================================================================
|
|
|
|
|
|
@given("a fully functional bridge for destructor testing")
|
|
def step_functional_bridge_for_del(context: Context) -> None:
|
|
bridge, _ = _build_bridge()
|
|
context.bridge = bridge
|
|
context.del_call_error = None
|
|
|
|
|
|
@when("__del__ is called on the fully functional bridge")
|
|
def step_call_del_on_functional_bridge(context: Context) -> None:
|
|
try:
|
|
context.bridge.__del__()
|
|
except Exception as exc:
|
|
context.del_call_error = exc
|
|
|
|
|
|
@then("the __del__ call should complete without any error")
|
|
def step_assert_del_no_error(context: Context) -> None:
|
|
assert context.del_call_error is None, (
|
|
f"Expected no error from __del__, got {context.del_call_error!r}"
|
|
)
|