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
1 change: 1 addition & 0 deletions .github/workflows/python-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -107,6 +107,7 @@ jobs:
run: |
./node_modules/.bin/playwright install --with-deps chromium
npm run smoke:personal-workspace-packaged
npm run smoke:chat-turn-acceptance-retry
npm run smoke:chat-upgrade
- uses: actions/upload-artifact@v7
with:
Expand Down
1 change: 1 addition & 0 deletions apps/presentation/dashboard/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
"smoke:benchmark-study-browser": "node ../../../examples/dashboard-benchmark-study-browser-smoke.mjs",
"smoke:capability-configuration": "rm -rf node_modules/.cache/loopx-capability-configuration-smoke && tsc --ignoreConfig --target ES2022 --module NodeNext --moduleResolution NodeNext --types node --skipLibCheck --strict --outDir node_modules/.cache/loopx-capability-configuration-smoke smoke/capability-configuration-smoke.ts src/data/capability-configuration.ts && node node_modules/.cache/loopx-capability-configuration-smoke/smoke/capability-configuration-smoke.js",
"smoke:chat-route": "tsc --ignoreConfig --target ES2022 --module ES2022 --moduleResolution Bundler --ignoreDeprecations 6.0 --skipLibCheck --strict --outDir node_modules/.cache/loopx-chat-route-smoke smoke/chat-route-smoke.ts src/data/chat-model.ts src/vite-env.d.ts && node node_modules/.cache/loopx-chat-route-smoke/smoke/chat-route-smoke.js",
"smoke:chat-turn-acceptance-retry": "vite build --ssr smoke/chat-turn-acceptance-retry-smoke.ts --outDir node_modules/.cache/loopx-chat-turn-acceptance-retry --emptyOutDir && node node_modules/.cache/loopx-chat-turn-acceptance-retry/chat-turn-acceptance-retry-smoke.js",
"smoke:demo-readiness": "bash ../../../scripts/loopx-python.sh --exec ../../../examples/dashboard-demo-readiness-smoke.py",
"smoke:frontstage-browser": "node ../../../examples/dashboard-frontstage-browser-smoke.mjs",
"smoke:frontstage-design-baseline": "node ../../../examples/dashboard-frontstage-design-baseline-smoke.mjs",
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,224 @@
from __future__ import annotations

import json
from pathlib import Path
import sys
import tempfile
import threading
from typing import Any

sys.path.insert(0, str(Path(__file__).resolve().parents[4]))

from loopx.chat_runtime import ChatRuntimeController
from loopx.chat_server import ChatHTTPServer, ChatRequestHandler
from loopx.chat_store import ChatSessionStore


class _HealthyAdapter:
upstream_thread_id = "fixture-upstream"

def capabilities(self) -> dict[str, Any]:
return {}

def start_turn(
self,
message: str,
event_sink: Any,
) -> dict[str, Any]:
del message, event_sink
raise AssertionError("the fixture replaces runtime dispatch")

def interrupt_turn(self, turn_id: str | None = None) -> None:
del turn_id

def healthcheck(self) -> bool:
return True

def close_session(self) -> None:
return None


def _turn_count(store: ChatSessionStore, session_id: str) -> int:
return sum(
1
for path in (store.sessions_root / session_id / "turns").glob("*.json")
if not path.name.endswith(".events.json")
)


def main() -> int:
scenario = sys.argv[1] if len(sys.argv) > 1 else ""
if scenario not in {"before_transcript", "after_queued"}:
raise ValueError("unknown acceptance fault scenario")

root = Path(tempfile.mkdtemp(prefix="loopx-chat-acceptance-http-"))
project = root / "project"
project.mkdir()
registry_path = root / "registry.json"
registry_path.write_text(
json.dumps(
{
"schema_version": "0.1",
"goals": [
{
"id": "goal-one",
"repo": str(project),
"status": "active",
}
],
}
),
encoding="utf-8",
)
store = ChatSessionStore(root / "runtime")
runtime = ChatRuntimeController(
store=store,
codex_bin="missing-codex",
registry_path=registry_path,
)
session_id = str(
store.create_session(
goal_id="goal-one",
agent_id="codex",
executor_endpoint_id="codex",
adapter_kind="codex_app_server",
upstream_thread_id="fixture-upstream",
upstream_mode="chat",
codex_home=str(runtime.codex_home),
)["session_id"]
)
mismatched_session_id = str(
store.create_session(
goal_id="goal-one",
agent_id="codex",
executor_endpoint_id="codex",
adapter_kind="codex_app_server",
upstream_thread_id="other-upstream",
upstream_mode="chat",
codex_home=str(root / "other-codex-home"),
)["session_id"]
)
runtime.adapters[session_id] = _HealthyAdapter()

dispatch_started = threading.Event()
dispatch_release = threading.Event()
dispatch_count = 0

def blocked_run_turn(**kwargs: object) -> None:
nonlocal dispatch_count
dispatch_count += 1
dispatch_started.set()
dispatch_release.wait(timeout=10)
done_event = kwargs["done_event"]
if not isinstance(done_event, threading.Event):
raise TypeError("turn dispatch must carry a completion event")
with runtime.lock:
runtime.turn_done_events.pop(
(session_id, str(kwargs["turn_id"])),
None,
)
done_event.set()

runtime._run_turn = blocked_run_turn

failed = False
if scenario == "before_transcript":
original_append_message = store.append_message

def fail_before_transcript(*args: Any, **kwargs: Any) -> dict[str, Any]:
nonlocal failed
if not failed and kwargs.get("role") == "user":
failed = True
raise OSError("private transcript fault detail")
result = original_append_message(*args, **kwargs)
if not isinstance(result, dict):
raise TypeError("append_message returned an invalid result")
return result

store.append_message = fail_before_transcript
else:
original_append_event = store.append_event

def fail_after_queued(*args: Any, **kwargs: Any) -> dict[str, Any]:
nonlocal failed
result = original_append_event(*args, **kwargs)
if not isinstance(result, dict):
raise TypeError("append_event returned an invalid result")
if not failed and kwargs.get("kind") == "turn.queued":
failed = True
raise OSError("private queued event fault detail")
return result

store.append_event = fail_after_queued

server = ChatHTTPServer(("127.0.0.1", 0), ChatRequestHandler)
server.verbose = False
server.registry_path = registry_path
server.chat_store = store
server.runtime_controller = runtime
server_thread = threading.Thread(target=server.serve_forever, daemon=True)
server_thread.start()
print(
json.dumps(
{
"origin": f"http://127.0.0.1:{server.server_port}",
"session_id": session_id,
"mismatched_session_id": mismatched_session_id,
}
),
flush=True,
)

try:
if sys.stdin.readline().strip() != "inspect":
raise ValueError("fixture did not receive an inspect request")
if not dispatch_started.wait(timeout=5):
raise TimeoutError("accepted turn was not dispatched")
turn = store.turn_for_client(session_id, "recoverable-request")
if turn is None:
raise AssertionError("accepted turn is missing")
queued_events = [
event
for event in store.events_after(
session_id,
str(turn["turn_id"]),
None,
)
if event.get("kind") == "turn.queued"
]
user_messages = [
message
for message in store.messages(session_id)
if message.get("role") == "user"
]
mismatch_turns = _turn_count(store, mismatched_session_id)
mismatch_messages = store.messages(mismatched_session_id)
print(
json.dumps(
{
"fault_injected": failed,
"turn_count": _turn_count(store, session_id),
"user_message_count": len(user_messages),
"queued_event_count": len(queued_events),
"dispatch_count": dispatch_count,
"acceptance_capsule_present": "_acceptance" in turn,
"mismatched_turn_count": mismatch_turns,
"mismatched_message_count": len(mismatch_messages),
"mismatched_session_status": (
store.load_session(mismatched_session_id) or {}
).get("status"),
}
),
flush=True,
)
finally:
dispatch_release.set()
server.shutdown()
server.server_close()
server_thread.join(timeout=5)
runtime.close()
return 0


if __name__ == "__main__":
raise SystemExit(main())
Loading
Loading