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
23 changes: 23 additions & 0 deletions docs/architecture/rfcs/capable-manager-semantic-handoff-v0.md
Original file line number Diff line number Diff line change
Expand Up @@ -595,6 +595,29 @@ Target an ingress receipt within two seconds on a healthy local service, indepen

Use existing service recovery and receipt pumps. No manager-specific business automation for each kind of request. Expose configuration and failures through the existing CLI, capability settings and manager conversation. Troubleshooting distinguishes model failure, tool/policy denial, state conflict, unreachable receiver and transport formatting/delivery failure.

**Accepted queue preparation failures (S1/S10, A12/A22/A23):** an accepted
request owns a terminal outcome even before an adapter starts. A missing runtime
asset, invalid workspace or failed session restoration must settle the affected
queued Turn through the shared Chat lifecycle and release its claim. Waiting for
an answer must observe that durable failure promptly, rather than wait for the
model timeout while leaving runnable work behind. Keep the original typed
provider failure where available; unexpected local preparation errors use
`runtime_unavailable`, with private diagnostics retained locally. Cancellation
and an already terminal result win over a late preparation error. Restoring the
runtime must not replay a failed request; the same ingress identity returns the
same failure, while a fresh explicit request can run after repair.

The bounded Python queue repair uses the existing store's fenced failure and
claim-release operations for all queue callers; Lark only translates the typed
outcome. It does not create a separate manager scheduler or new TS authority.
The TS turn-driver migration must preserve this pre-dispatch failure matrix
alongside accepted-request recovery. Validate with a removed-release fixture,
multiple queued requests, a stop race, same-identity redelivery and a fresh
request after recovery. These qualify the preparation boundary, not successful
owner selection, receiver adoption or the complete A24 journey. Operational
recovery must also verify the service's actual installed release: a healthy HTTP
listener alone does not prove its lazy-loaded runtime assets still exist.

## 11. Normative delivery plan

Implement coherent end-to-end slices, not one PR per incidental field. The manager engineering owner maintains canonical Todos and a private incident-to-acceptance map; PRs link this RFC milestone and acceptance IDs. Public progress updates contain only safe results. Milestone completion requires current deployment evidence, not merged PR count.
Expand Down
2 changes: 1 addition & 1 deletion docs/architecture/rfcs/loopx-overall-roadmap-v0.md
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,7 @@ P0 blocks correctness or continuity in the current user journey. P1 enables repe
| **S7 Budget, scheduling and fleet scale · P0 observation/P1–P2 expansion** | Quota/scheduler and partial usage aggregates exist; full provider cost, distributed reservations and hundred-Agent concurrency need evidence | Separate configured budget, admission, consumption and estimates; unknown is not zero and replay cannot double-charge. R7 pagination/bounded summaries and [complete-history transport](typescript-control-plane-migration-v0.md), including refresh/replay/single-debit evidence beyond the RPC limit; provider/host limits, fairness, backpressure, event wake and isolation; report registration/activity/throughput and cost per accepted outcome separately |
| **S8 Capabilities, extensions and domain integration · P1/P2** | Capability catalog, extension lifecycle, hooks, engineering/research/content/office capabilities and computer-use contracts exist | First exercise the shared control plane with existing issue-fix/PR-review and material/research callers. Every provider has readiness/version/permissions/default-off/uninstall/rollback/isolation and real-entry evidence. New domain effects start with one simulated operation, not a marketplace or workflow DSL |
| **S9 Identity, authority, privacy and trust · continuous P0/P1–P2 remote** | Public/private scope, capability gates, fencing and confirmation contracts belong to existing owners | R1/R3 cover sender/audience/artifact scope and stale authority; R6 authenticates tenant/Goal/actor/host, rotation/revocation and least privilege. Qualify credential custody, untrusted tool/document inputs, dependency supply chain, audit retention/deletion and vulnerability response through real paths; roles/messages/memory mint no write authority |
| **S10 Reliability, diagnostics and operations · P0/P1** | Recovery/canary, read-only diagnostics prototype and DSH event adapter exist; C0/C1, overhead and full operations qualification are open | Failure classification→observable state→recovery drill→regression prevention; process/storage/network/delivery failures and data growth. Use [bounded repair lookup and targeted diagnostics](../../../skills/loopx-self-repair/references/targeted-diagnostics.md) to reduce redundant reads above the provider boundary; measure backend-specific cold/warm reads, writes and lock waits separately. Freeze SLO/RPO/RTO/capacity/retention boundaries and measure before qualification. Runbooks include upgrade, restore, stop and human takeover; test counts do not prove recovery |
| **S10 Reliability, diagnostics and operations · P0/P1** | Recovery/canary, read-only diagnostics prototype and DSH event adapter exist; C0/C1, overhead and full operations qualification are open | Failure classification→observable state→recovery drill→regression prevention; process/storage/network/delivery failures and data growth. Accepted Chat requests must settle even when runtime preparation fails before dispatch; qualify missing runtime assets, stop races and recovery without replay under [the shared conversation operational contract](capable-manager-semantic-handoff-v0.md#10-operational-contract). Use [bounded repair lookup and targeted diagnostics](../../../skills/loopx-self-repair/references/targeted-diagnostics.md) to reduce redundant reads above the provider boundary; measure backend-specific cold/warm reads, writes and lock waits separately. Freeze SLO/RPO/RTO/capacity/retention boundaries and measure before qualification. Runbooks include upgrade, restore, stop and human takeover; test counts do not prove recovery |
| **S11 Evaluation and scientific research · continuous P1/P2 research** | Benchmark toolkit, Explore, long-horizon portfolio and ten frontier-science tracks have designs/partial implementations | Pin native/passive/governed arms, model/harness/budget/task split and evaluator; report native scores, cost, failures, attention and uncertainty. Prioritize sequential evidence, continuation and stride; memory, formal kernel, curriculum/evolution, active experiments and multiscale state follow T01–T10 gates without automatic production treatment |
| **S12 Release, developer experience and community governance · P0 hygiene/P1** | Install/source validation, registration, DCO/PR, test layers, contributor routes and bilingual docs exist | Qualify first work and upgrade/rollback from clean machines/release artifacts; host/OS support follows the release contract. Reduce localization/test/review effort for useful changes; preserve exact-head evidence, fixtures, compatibility, maintainer routing and contributor credit; retire duplicate protocols/stale evidence |
| **S13 Adoption, ecosystem and sustainability · P1 discovery/P2 pilots** | Public adoption loop, showcases, licensing/governance and observer-first product contract exist; paid PMF is unproven | Gather independent first/repeat usage and exit reasons; reproducible cases and pilots with fixed budgets/acceptance/rollback. Retain reusable adapters/delivery guides. Account for model/compute/storage/support and maintenance costs; only repeated demand justifies commercial hosting/support/distribution decisions, with no invented SLA or open-source-term change |
Expand Down
58 changes: 46 additions & 12 deletions loopx/chat_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@

from dataclasses import dataclass
import json
import logging
import os
from pathlib import Path
import threading
Expand Down Expand Up @@ -1114,12 +1115,16 @@ def _drain_session_queue(
with self.lock:
attached = self.adapters.get(session_id)
if attached is None or not attached.healthcheck():
self._ensure_adapter(
session,
work_dir=work_dir,
objective=objective,
)
continue
try:
self._ensure_adapter(
session,
work_dir=work_dir,
objective=objective,
)
except Exception as exc:
self._fail_queue_preparation(session_id, active_turn_id, exc)
else:
continue
active = self.store.load_turn(session_id, active_turn_id)
if active and active.get("status") not in TERMINAL_TURN_STATES:
self.closed.wait(0.05)
Expand All @@ -1132,11 +1137,17 @@ def _drain_session_queue(
refreshed = self.store.load_session(session_id)
if refreshed is None:
return
adapter = self._ensure_adapter(
refreshed,
work_dir=work_dir,
objective=objective,
)
preparation_error = None
try:
adapter = self._ensure_adapter(
refreshed,
work_dir=work_dir,
objective=objective,
)
except Exception as exc:
# An accepted queued request still owns an outcome when
# workspace preparation or adapter restoration fails.
preparation_error = exc
with self.lock:
self.session_queue_wakeups.discard(session_id)
turn = self.store.claim_next_queued_turn(session_id)
Expand All @@ -1150,6 +1161,9 @@ def _drain_session_queue(
self.session_queue_threads.pop(session_id, None)
return
turn_id = str(turn["turn_id"])
if preparation_error is not None:
self._fail_queue_preparation(session_id, turn_id, preparation_error)
continue
with self.lock:
self.turn_done_events[(session_id, turn_id)] = threading.Event()
self._run_turn(
Expand All @@ -1167,6 +1181,25 @@ def _drain_session_queue(
self.session_queue_threads.pop(session_id, None)
self.session_queue_wakeups.discard(session_id)

def _fail_queue_preparation(
self, session_id: str, turn_id: str, error: Exception,
) -> None:
logging.getLogger(__name__).error(
"Chat queue runtime preparation failed",
exc_info=(type(error), error, error.__traceback__),
)
self._fail_turn(
session_id, turn_id,
error.error_code if isinstance(error, CodexChatAgentError) else "runtime_unavailable",
str(error) if isinstance(error, CodexChatAgentError) else (
"The Agent runtime could not be prepared. Check the runtime installation "
"and configuration, then retry."
),
status="failed",
gate=error.gate if isinstance(error, CodexChatAgentError) else None,
expected_statuses={"queued", "starting", "running"},
)

def _run_turn(
self,
*,
Expand Down Expand Up @@ -1390,12 +1423,13 @@ def _fail_turn(
*,
status: str,
gate: dict[str, Any] | None = None,
expected_statuses: set[str] | None = None,
) -> None:
completed = utc_now()
failed = self.store.update_turn(
session_id,
turn_id,
expected_statuses={"starting", "running"},
expected_statuses=expected_statuses if expected_statuses is not None else {"starting", "running"},
status=status,
error_code=error_code,
error=message,
Expand Down
2 changes: 2 additions & 0 deletions loopx/extensions/lark/manager_context.py
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,8 @@ def manager_failure_reply(error: Exception) -> tuple[str, str]:
"host_gate": "执行器未能完成本次调用,请检查管家执行器的具体错误",
"upstream_invalid_request": "上游拒绝了请求参数,请检查所选模型与 Codex CLI、账户的兼容性",
"server_restarted": "管家服务在处理过程中重启",
"runtime_unavailable": "本地运行环境初始化失败,请检查 LoopX 安装与运行配置,修复后重试",
"resume_failed": "原 Agent 会话恢复失败,请检查执行器状态,修复后重试",
"manager_channel_executor_rebind_required": (
"管家的执行器已由本机设置更改,需要重新应用一次管家连接"
),
Expand Down
2 changes: 2 additions & 0 deletions tests/extensions/test_lark_goal_topic_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -2570,6 +2570,8 @@ def answer(_route, _text):
@pytest.mark.parametrize("error_code,label", [
("cyber_policy", "安全策略拦截"),
("rate_limit_exceeded", "请求频率限制"),
("runtime_unavailable", "本地运行环境初始化失败"),
("resume_failed", "原 Agent 会话恢复失败"),
("private-upstream-detail", "管家处理失败"),
])
@pytest.mark.parametrize("reply_ok", [True, False])
Expand Down
159 changes: 159 additions & 0 deletions tests/test_chat_queue_preparation.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,159 @@
"""Accepted requests must settle even when the runtime cannot be prepared."""

import threading

import pytest

from loopx.chat_agent import CodexChatAgentError
from loopx.chat_runtime import ChatRuntimeController
from loopx.chat_store import ChatSessionStore
from loopx.extensions.lark.goal_topic_runtime import (
LarkGoalTopicTurnFailed, answer_lark_goal_topic,
)
from loopx.extensions.lark.manager_context import manager_failure_reply


def session_runtime(tmp_path, *, channel="goal.public-research"):
store = ChatSessionStore(tmp_path)
session = store.create_session(
goal_id="public-research", agent_id="codex", channel_id=channel,
executor_endpoint_id="codex", adapter_kind="codex_app_server",
upstream_thread_id="public-thread", upstream_mode="chat",
)
return store, str(session["session_id"]), ChatRuntimeController(
store=store, codex_bin="unused-test-codex",
)


def drain(runtime, session_id, tmp_path):
runtime.resume_session_queue(session_id=session_id, work_dir=tmp_path, objective="Research")
# A worker that already retired has completed its queue cleanup.
with runtime.lock:
worker = runtime.session_queue_threads.get(session_id)
if worker:
worker.join(3)
assert not worker.is_alive()


@pytest.mark.parametrize("claimed", [False, True])
def test_missing_runtime_asset_settles_accepted_work_without_dispatch(
tmp_path, monkeypatch, claimed,
):
store, sid, runtime = session_runtime(tmp_path)
turns = [store.create_queued_turn(sid, client_turn_id=f"request-{i}", message="Research cash flow")[0]
for i in range(2)]
if claimed:
assert store.claim_next_queued_turn(sid)

def missing_asset(*args, **kwargs):
(tmp_path / "removed-release" / "SKILL.md").read_text()

monkeypatch.setattr(runtime, "_ensure_adapter", missing_asset)
drain(runtime, sid, tmp_path)
for turn in turns:
result = runtime.wait_for_turn(session_id=sid, turn_id=turn["turn_id"], timeout_sec=0.1)
assert result["status"] == "failed"
assert result["error_code"] == "runtime_unavailable"
assert result.get("started_at") is None
assert result.get("upstream_turn_id") is None
assert str(tmp_path) not in result["error"]
events = store.events_after(sid, turn["turn_id"], None)
assert [e["kind"] for e in events].count("turn.failed") == 1
assert not store.queued_turns(sid)
assert store.load_session(sid)["active_turn_id"] is None

class Adapter:
upstream_thread_id = "public-thread"
messages = []

def healthcheck(self):
return True

def start_turn(self, message, sink):
self.messages.append(message)
return {"message": "Recovered"}

adapter = Adapter()
monkeypatch.setattr(runtime, "_ensure_adapter", lambda *a, **kw: adapter)
retry, created = runtime.enqueue_turn(
session_id=sid, client_turn_id="request-0", message="Research cash flow",
work_dir=tmp_path, objective="Research",
)
assert not created
assert retry["status"] == "failed"
new, _ = runtime.enqueue_turn(
session_id=sid, client_turn_id="explicit-new-request", message="New request",
work_dir=tmp_path, objective="Research",
)
assert runtime.wait_for_turn(session_id=sid, turn_id=new["turn_id"], timeout_sec=3)["status"] == "completed"
assert adapter.messages == ["New request"]


def test_stop_wins_over_late_preparation_failure(tmp_path, monkeypatch):
store, sid, runtime = session_runtime(tmp_path)
entered, release = threading.Event(), threading.Event()

def preparation(*args, **kwargs):
entered.set()
assert release.wait(3)
raise OSError("runtime removed")

monkeypatch.setattr(runtime, "_ensure_adapter", preparation)
turn, _ = runtime.enqueue_turn(
session_id=sid, client_turn_id="stop-race", message="Research",
work_dir=tmp_path, objective="Research",
)
assert entered.wait(3)
assert runtime.interrupt_turn(session_id=sid, turn_id=turn["turn_id"])["status"] == "interrupted"
release.set()
drain(runtime, sid, tmp_path)
assert store.load_turn(sid, turn["turn_id"])["status"] == "interrupted"
assert not any(e["kind"] == "turn.failed" for e in store.events_after(sid, turn["turn_id"], None))


def test_removed_manager_release_fails_before_starting_provider(tmp_path, monkeypatch):
import loopx.chat_manager as manager

store, sid, runtime = session_runtime(tmp_path, channel="manager")
monkeypatch.setattr(manager, "__file__", str(tmp_path / "removed-release" / "chat_manager.py"))
starts = []
monkeypatch.setattr(runtime, "_start_adapter", lambda **kw: starts.append(kw))
turn, _ = runtime.enqueue_turn(
session_id=sid, client_turn_id="missing-manager-skill", message="Research cash flow",
work_dir=tmp_path, objective="Research",
)
result = runtime.wait_for_turn(session_id=sid, turn_id=turn["turn_id"], timeout_sec=3)
assert result["status"] == "failed"
assert result["error_code"] == "runtime_unavailable"
assert result.get("started_at") is None
assert starts == []
assert store.load_session(sid)["active_turn_id"] is None


@pytest.mark.parametrize("error,code", [
(OSError("private installation path"), "runtime_unavailable"),
(CodexChatAgentError("Cannot restore", error_code="resume_failed", gate=None), "resume_failed"),
])
def test_lark_wait_receives_typed_failure_without_model_or_timeout(tmp_path, monkeypatch, error, code):
store, sid, runtime = session_runtime(tmp_path)

def preparation(*args, **kwargs):
raise error

monkeypatch.setattr(runtime, "_ensure_adapter", preparation)
real_wait = runtime.wait_for_turn
monkeypatch.setattr(runtime, "wait_for_turn", lambda **kw: real_wait(**kw, timeout_sec=3))
with pytest.raises(LarkGoalTopicTurnFailed) as failure:
answer_lark_goal_topic(
route={"goal_id": "public-research", "session_id": sid,
"ingress_mode": "session_queue", "message_id": "om_public_question",
"topic_root_message_id": "om_public_root"},
text="Research Microsoft's cash flow", work_dir=tmp_path,
objective="Research", runtime_controller=runtime,
)
assert failure.value.error_code == code
reply_code, reply = manager_failure_reply(failure.value)
assert reply_code == code
assert "修复后重试" in reply
assert "private installation path" not in reply
assert not store.queued_turns(sid)
Loading