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
Original file line number Diff line number Diff line change
Expand Up @@ -80,6 +80,15 @@ approval store, scheduler or model-execution authority. Missing evidence is not
permission to restart a request in a fresh model thread, elevate host policy
or acknowledge work as completed.

Private DM result reconciliation uses the existing Chat store's terminal-state
owner, including `timed_out`, for ordinary replies and commission results.
Timeouts return an explicit failure and the original conversation's recovery
controls. Provider readback still gates delivery: restarting the transport or
redelivering the source reconciles the saved attempt without another send or
model execution. Synthetic host/provider checks cover idle/hard DM timeouts
and an injected commission timeout; a live-provider recovery drill remains
part of the switch gate below.

## Qualification before switching

The acceptance owner is the steward RFC's
Expand Down
11 changes: 7 additions & 4 deletions loopx/extensions/lark/private_conversations.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@
from typing import Any

from ...capabilities.native_chat.external_conversations import ChatExternalConversations
from ...chat_store import _atomic_write_json, _read_json
from ...chat_store import TERMINAL_TURN_STATES, _atomic_write_json, _read_json
from ...file_lock import exclusive_file_lock
from .conversation_identity import identity_ref, lark_private_source
from .event_inbox import acknowledge_lark_event_inbox, ingest_lark_event_inbox
Expand Down Expand Up @@ -261,10 +261,11 @@ def reconcile(self) -> int:
# independent of terminal execution and reply delivery.
self._deliver(path, record, "admission", "已持久受理到原 Agent 会话;等待原宿主领取。/status 查看持久队列,/project 返回普通项目对话。实时停止暂不支持,请在原宿主处理。" if native.get("agent_target") else "已持久受理;若已有执行,本条会排队。可发送 /status、/stop 或 /new。")
turn = self.core.controller.store.load_turn(record["session_id"], record["turn_id"])
if not turn or turn["status"] not in {"completed", "failed", "interrupted", "expired"}:
if not turn or turn["status"] not in TERMINAL_TURN_STATES:
continue
response = str((turn.get("response") or {}).get("message") or "") if turn["status"] == "completed" else (
"本次执行已停止。" if turn["status"] == "interrupted" else
"本次执行超时,原会话已保留;请发送 /status 查看状态后再决定是否重试。" if turn["status"] == "timed_out" else
"本次执行失败或已过期,原会话已保留;请发送 /status 后再决定是否重试。")
else:
response = _command_text(str(record.get("response_code") or "")) or str(record.get("response") or "")
Expand All @@ -274,12 +275,14 @@ def reconcile(self) -> int:
resources = record.get("commission_resources") or {}
if resources.get("session_id") and resources.get("turn_id"):
first_turn = self.core.controller.store.load_turn(resources["session_id"], resources["turn_id"])
if not first_turn or first_turn["status"] not in {"completed", "failed", "interrupted", "expired"}:
if not first_turn or first_turn["status"] not in TERMINAL_TURN_STATES:
record["status"] = "commission_running"
_atomic_write_json(path, record)
continue
result_text = str((first_turn.get("response") or {}).get("message") or "")
result_text = "委托执行结果:\n" + result_text if first_turn["status"] == "completed" else "委托首轮执行未完成;原 Goal 和回执已保留,请查看状态后决定恢复。"
result_text = ("委托执行结果:\n" + result_text if first_turn["status"] == "completed" else
"委托首轮执行超时;原 Goal 和回执已保留,请查看状态后决定恢复。" if first_turn["status"] == "timed_out" else
"委托首轮执行未完成;原 Goal 和回执已保留,请查看状态后决定恢复。")
proposal_id = native.get("proposal_id")
if proposal_id:
result_text += f"\n如需恢复暂停或额度受限的原执行:/resume-commission {proposal_id} --tokens N(N 为包含历史用量的总上限,须大于已用量;不会重开线程)。"
Expand Down
97 changes: 97 additions & 0 deletions tests/test_lark_private_timeouts.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,97 @@
"""Real native timeout settlement with synthetic host and provider transports."""
import json
import threading

import pytest
from test_chat_ordinary_project import ordinary # noqa: F401
from test_lark_private_conversations import connect
from test_native_steward_private import steward # noqa: F401

from loopx.chat_codex_goal import CodexGoalDriver
from loopx.extensions.lark.private_conversations import LarkPrivateConversations


@pytest.mark.parametrize("timeout_kind", ["idle", "hard"])
def test_private_timeout_returns_once_after_delivery_readback_recovery(ordinary, timeout_kind): # noqa: F811
store, runtime, provider, transport = connect(ordinary)
runtime.idle_timeout_sec = .1 if timeout_kind == "idle" else 30
runtime.hard_timeout_sec = .1 if timeout_kind == "hard" else 30
try:
event = provider.event("notes-app", "timeout", "wait for interrupt")
assert transport.admit("notes-app", event)["status"] == "durably_accepted"
row = transport.core.pending()[0]
sid, tid = row["session_id"], row["turn_id"]
terminal = runtime.wait_for_turn(session_id=sid, turn_id=tid, timeout_sec=10)
assert terminal["status"] == "timed_out"
assert terminal["error_code"] == f"{timeout_kind}_timeout"
original = store.load_session(sid)["upstream_thread_id"]

provider.verify_replies = False
assert transport.reconcile() == 0
assert any(profile == "notes-app" and "本次执行超时" in text for profile, text in provider.writes)
assert not transport.core.read_request(row["request_ref"]).get("delivery_verified")
count = len(provider.writes)
restarted = LarkPrivateConversations(controller=runtime, runtime_root=store.root.parent,
runner=provider, cli_bin="lark-cli")
assert restarted.reconcile() == 0
assert len(provider.writes) == count
provider.verify_replies = True
assert restarted.reconcile() == 1
assert restarted.core.read_request(row["request_ref"])["delivery_verified"] is True
assert restarted.health()[row["binding_id"]] == {"pending_count": 0, "recovery_count": 0}
assert restarted.admit("notes-app", {**event, "event_id": "redelivery"})["status"] == "durably_accepted"
restarted.reconcile()
assert len(provider.writes) == count
assert len(store.list_sessions()) == 1
assert store.load_session(sid)["upstream_thread_id"] == original
assert store.load_session(sid)["goal_id"] is None
assert store.load_turn(sid, tid)["status"] == "timed_out"
requests = [json.loads(line) for line in ordinary[4].read_text().splitlines()]
assert len([r for r in requests if r.get("method") == "turn/start"]) == 1
finally:
runtime.close()


def test_commission_timeout_returns_original_resources_without_restarting_work(steward, monkeypatch): # noqa: F811
store, runtime, provider, transport, binding, capture, _ = steward
started, release = threading.Event(), threading.Event()
# Inject a typed host timeout after Goal activation. Core persists the
# terminal Turn; this does not qualify the real native Goal deadline.
def timeout(self, emit):
started.set()
assert release.wait(10), "timeout fixture was not released"
raise self.session._timeout_error("hard_timeout", "Synthetic host timeout")
monkeypatch.setattr(CodexGoalDriver, "_observe", timeout)
event = provider.event("steward-app", "delegate", "/delegate --tokens 12000 wait for interrupt")
assert transport.admit("steward-app", event)["status"] == "command_recorded"
transport.reconcile()
proposal = transport.core.actions.store.list()[0]
confirm = provider.event("steward-app", "confirm", "/confirm " + proposal["proposal_id"])
assert transport.admit("steward-app", confirm)["status"] == "command_recorded"
transport.reconcile()
resources = transport.core.actions.load(proposal["proposal_id"])["receipt"]["resource_ids"]
assert started.wait(10)
provider.verify_replies = False
release.set()
turn = runtime.wait_for_turn(session_id=resources["session_id"], turn_id=resources["turn_id"], timeout_sec=10)
assert turn["status"] == "timed_out" and turn["error_code"] == "hard_timeout"
transport.reconcile()
assert any(profile == "steward-app" and "委托首轮执行超时" in text and
f"/resume-commission {proposal['proposal_id']}" in text for profile, text in provider.writes)
count = len(provider.writes)
before = capture.read_text()
restarted = LarkPrivateConversations(controller=runtime, runtime_root=store.root.parent,
runner=provider, cli_bin="lark-cli")
assert restarted.reconcile() == 0
assert len(provider.writes) == count
provider.verify_replies = True
assert restarted.reconcile() == 1
assert restarted.health()[binding["binding_id"]] == {"pending_count": 0, "recovery_count": 0}
restarted.admit("steward-app", {**confirm, "event_id": "redelivery"})
restarted.reconcile()
assert len(provider.writes) == count
assert transport.core.actions.load(proposal["proposal_id"])["receipt"]["resource_ids"] == resources
assert len(json.loads(runtime.registry_path.read_text())["goals"]) == 1
assert store.load_turn(resources["session_id"], resources["turn_id"])["status"] == "timed_out"
assert capture.read_text() == before
assert not any(profile == "notes-app" for profile, _ in provider.writes)
Loading