Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 4 additions & 3 deletions loopx/capabilities/manager_context/roundtrip.py
Original file line number Diff line number Diff line change
Expand Up @@ -1355,7 +1355,7 @@ def __init__(self, root, registry, store, external_sender):
target=self.run, daemon=True, name="loopx-manager-returns"
)

def start(self):
def start(self) -> None:
self.thread.start()

def run(self):
Expand All @@ -1368,6 +1368,7 @@ def run(self):
)
self.stop.wait(3)

def close(self):
def close(self) -> None:
self.stop.set()
self.thread.join(timeout=3)
if self.thread.ident is not None:
self.thread.join(timeout=3)
11 changes: 9 additions & 2 deletions loopx/chat_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -52,7 +52,10 @@
from .extensions.lark.conversation_identity import observe_lark_conversation_identity
from .extensions.lark.private_conversations import LarkPrivateConversations
from .chat_loopx_mode import handle_loopx_request
from .capabilities.manager_context.roundtrip import project_chat_session_snapshot
from .capabilities.manager_context.roundtrip import (
ReturnService,
project_chat_session_snapshot,
)
from .control_plane.goals.active_state_metadata import active_state_section_text
from .control_plane.coordination.local_authority import LocalCoordinationAuthorityUnavailable
from .control_plane.status.ssh_host_catalog import (
Expand Down Expand Up @@ -420,6 +423,7 @@ class ChatHTTPServer(ThreadingHTTPServer):
action_store: ChatActionStore
action_service: ChatActionService
runtime_controller: ChatRuntimeController
manager_return_service: ReturnService
lark_runner: CommandRunner
lark_cli_resolution: LarkCliResolution
lark_app_setup_manager: LarkAppSetupManager
Expand Down Expand Up @@ -1690,7 +1694,9 @@ def _lark_snapshot():
)
server.lark_goal_topic_runtime.start()
from .extensions.lark.manager_returns import start_return_service
server.manager_return_service = start_return_service(server, server.runtime_controller.coordination_runtime_root)
server.manager_return_service = start_return_service(
server, server.runtime_controller.coordination_runtime_root, start_service=False
)
from .chat_loopx_mode import DelegationWakeService

def _wake_goal_context(session):
Expand All @@ -1717,6 +1723,7 @@ def _wake_goal_context(session):
try:
if external_conversation_factories:
server.conversation_transports.start(server)
server.manager_return_service.start()
server.serve_forever()
except KeyboardInterrupt:
print("Stopping LoopX Chat", flush=True)
Expand Down
5 changes: 3 additions & 2 deletions loopx/extensions/lark/manager_returns.py
Original file line number Diff line number Diff line change
Expand Up @@ -320,7 +320,7 @@ def verify(self, route: dict[str, Any], session: dict[str, Any], turn: dict[str,
)


def start_return_service(server: Any, runtime_root: Path) -> Any:
def start_return_service(server: Any, runtime_root: Path, *, start_service: bool = True) -> Any:
"""Compose the return pump at the existing Chat/Lark service boundary."""
from ...capabilities.manager_context.roundtrip import ReturnService

Expand All @@ -333,5 +333,6 @@ def start_return_service(server: Any, runtime_root: Path) -> Any:
runtime_root, server.registry_path, server.chat_store, transport
)
transport.cancelled = service.stop.is_set
service.start()
if start_service:
service.start()
return service
8 changes: 4 additions & 4 deletions loopx/semantics/project_registry_io_manifest_v1.json
Original file line number Diff line number Diff line change
Expand Up @@ -495,31 +495,31 @@
},
{
"site": "loopx/chat_server.py::<module>.ChatRequestHandler._goal_channel_extension_ready::codec_read:load_registry#1",
"line": 1009,
"line": 1013,
"column": 24,
"kind": "codec_read",
"api": "load_registry",
"classification": "codec_api"
},
{
"site": "loopx/chat_server.py::<module>.ChatRequestHandler._registry_and_goal::codec_read:load_registry#1",
"line": 536,
"line": 540,
"column": 20,
"kind": "codec_read",
"api": "load_registry",
"classification": "codec_api"
},
{
"site": "loopx/chat_server.py::<module>.serve_chat::codec_read:load_registry#1",
"line": 1588,
"line": 1592,
"column": 16,
"kind": "codec_read",
"api": "load_registry",
"classification": "codec_api"
},
{
"site": "loopx/chat_server.py::<module>.serve_chat._wake_goal_context::codec_read:load_registry#1",
"line": 1697,
"line": 1703,
"column": 20,
"kind": "codec_read",
"api": "load_registry",
Expand Down
13 changes: 13 additions & 0 deletions tests/test_chat_transport_composition.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ def __init__(self, ref):
self.transport_ref = ref
self.calls = []
self.available = True
self.started = False

def observe(self):
self.calls.append("observe")
Expand All @@ -33,6 +34,7 @@ def verify(self, route, session, turn, text, attempt):
return {"reply_verified": False, "verification_performed": True}

def start(self, server):
self.started = True
self.calls.append(("start", server))

def close(self):
Expand Down Expand Up @@ -184,6 +186,16 @@ def test_real_chat_entrypoint_composes_one_store_controller_and_return_service(t
entered = threading.Event()
captured = []
external = Transport("external-owner")
return_service_start_observations = []
if installed:
from loopx.capabilities.manager_context import roundtrip
original_return_service_start = roundtrip.ReturnService.start

def record_return_service_start(service):
return_service_start_observations.append(external.started)
original_return_service_start(service)

monkeypatch.setattr(roundtrip.ReturnService, "start", record_return_service_start)
class Server(chat.ChatHTTPServer):
def __init__(self, *args, **kwargs):
super().__init__(*args, **kwargs)
Expand All @@ -204,6 +216,7 @@ def __init__(self, *args, **kwargs):
if installed:
assert server.manager_return_service.args[3] is server.conversation_transports
assert server.conversation_transports.bindings is server.runtime_controller.project_contexts.conversation_bindings
assert return_service_start_observations == [True]
else:
from loopx.extensions.lark.manager_returns import LarkManagerReturnTransport
assert not hasattr(server, "conversation_transports")
Expand Down
10 changes: 10 additions & 0 deletions tests/test_manager_context_roundtrip.py
Original file line number Diff line number Diff line change
Expand Up @@ -855,6 +855,16 @@ def transport(*_):
assert observations == [] # Legacy delivery is outside the exact-source telemetry contract.


def test_return_service_can_close_before_start(flow):
root, registry, store, _ = flow
service = ReturnService(root, registry, store, lambda *_: {"reply_verified": True})

service.close()

assert service.stop.is_set()
assert not service.thread.is_alive()


@pytest.mark.parametrize("project", [False, True], ids=["steward", "project"])
def test_registration_revocation_blocks_return_without_retargeting(flow, project):
root, registry, store, create = flow
Expand Down
Loading