From c4bf90a6043b28e0608fbd9746c3d88c7402e92d Mon Sep 17 00:00:00 2001 From: huashen <2494946808@qq.com> Date: Sat, 26 Sep 2026 13:51:42 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E5=A4=8D=20Host=20Bridge=20=E7=9F=AD?= =?UTF-8?q?=E6=9A=82=E5=A4=B1=E8=81=94=E6=81=A2=E5=A4=8D=E4=B8=8E=E6=89=A7?= =?UTF-8?q?=E8=A1=8C=E7=A7=9F=E7=BA=A6=E9=9A=94=E7=A6=BB?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- agent/host_bridge/client.py | 93 +++++- agent/host_bridge/filesystem.py | 68 +++- agent/host_bridge/host_bridge.proto | 1 + agent/host_bridge/host_bridge_pb2.py | 4 +- agent/host_bridge/host_bridge_pb2_grpc.py | 43 +++ agent/host_bridge/monitor.py | 52 ++- agent/host_bridge/server.py | 41 ++- bootstrap/app.py | 6 +- bootstrap/dashboard_api.py | 14 +- bootstrap/web_shell.py | 4 + docker/debug/host_bridge_notice.mjs | 54 ++++ docker/debug/host_bridge_reliability.py | 303 ++++++++++++++++++ docs/INDEX.md | 3 +- .../0075-host-bridge-runtime-recovery.md | 41 +++ docs/decisions/README.md | 1 + docs/design/host-bridge-protocol-v2.md | 20 +- docs/design/host-bridge-reliability.md | 87 +++++ docs/projectneed.md | 2 +- frontend/chat/src/desktop-chat-view.tsx | 2 + frontend/chat/src/host-bridge-notice.tsx | 33 ++ 20 files changed, 822 insertions(+), 50 deletions(-) create mode 100644 docker/debug/host_bridge_notice.mjs create mode 100644 docker/debug/host_bridge_reliability.py create mode 100644 docs/decisions/0075-host-bridge-runtime-recovery.md create mode 100644 docs/design/host-bridge-reliability.md create mode 100644 frontend/chat/src/host-bridge-notice.tsx diff --git a/agent/host_bridge/client.py b/agent/host_bridge/client.py index df145ef21..1231e1158 100644 --- a/agent/host_bridge/client.py +++ b/agent/host_bridge/client.py @@ -2,6 +2,7 @@ import asyncio import contextlib +import logging import uuid from dataclasses import dataclass from pathlib import Path @@ -27,6 +28,26 @@ from core.common.diagnostic_log import current_diagnostic_context _HEARTBEAT_INTERVAL_S = 2.0 +logger = logging.getLogger(__name__) + + +class HostBridgeRpcError(RuntimeError): + """保留传输状态;只有明确的暂时失联允许恢复探测和心跳。""" + + def __init__(self, method: str, code: grpc.StatusCode, detail: str | None) -> None: + self.method = method + self.code = code + uncertainty = ( + ";操作可能已生效,不得自动重发" + if method in {"Exec", "WriteStdin", "FileTool"} + else "" + ) + super().__init__(f"Host Bridge {method} 失败: {code.name}: {detail}{uncertainty}") + + @property + def transient(self) -> bool: + return self.code in {grpc.StatusCode.UNAVAILABLE, grpc.StatusCode.DEADLINE_EXCEEDED} + @dataclass(frozen=True) @@ -121,6 +142,8 @@ def __init__( self._stub = rpc.HostBridgeStub(self._channel) self._heartbeat_task: asyncio.Task[None] | None = None self._lease_error: Exception | None = None + self._opened = False + self._open_lock = asyncio.Lock() self._unconfirmed_owners: dict[str, str] = {} self._closed = False @@ -142,6 +165,7 @@ async def probe(self) -> dict[str, Any]: pb.ContextRequest(context=self._request_context()), method="Probe", timeout=5, + lease=False, ) return self._identity_reply(reply) @@ -150,6 +174,7 @@ async def inspect(self) -> dict[str, Any]: self._stub.Inspect, pb.ContextRequest(context=self._request_context()), method="Inspect", + timeout=5, lease=False, ) return self._identity_reply(reply) @@ -288,6 +313,9 @@ async def terminate_owner( async def shutdown(self) -> ExecutionCleanupReport: if self._closed: return ExecutionCleanupReport((), (), ()) + if not self._opened: + await self.close_transport() + return ExecutionCleanupReport((), (), ()) await self._stop_heartbeat() reply: pb.CleanupReply = await self._call( self._stub.ShutdownManager, @@ -358,10 +386,10 @@ async def _call( """发起一次 RPC;失败或取消均不重放可能已生效的操作。""" if self._closed: raise RuntimeError("Host Bridge manager 已关闭") - if method not in {"Heartbeat", "ShutdownManager"} and self._lease_error is not None: - raise RuntimeError(f"Host Bridge lease 已失效: {self._lease_error}") + if lease and self._lease_error is not None: + raise self._lease_error if lease: - self._ensure_heartbeat() + await self._open_manager() try: return await call( request, @@ -369,14 +397,31 @@ async def _call( metadata=(("authorization", f"Bearer {self._token}"),), ) except grpc.aio.AioRpcError as exc: - uncertainty = ( - ";操作可能已生效,不得自动重发" - if method in {"Exec", "WriteStdin", "FileTool"} - else "" + error = HostBridgeRpcError(method, exc.code(), exc.details()) + if self._opened and error.code in { + grpc.StatusCode.NOT_FOUND, grpc.StatusCode.PERMISSION_DENIED, + grpc.StatusCode.UNAUTHENTICATED, grpc.StatusCode.FAILED_PRECONDITION, + }: + self._lease_error = error + raise error from exc + + async def _open_manager(self) -> None: + """业务调用前只登记一次;失联续期不能重新创建已丢失的 manager。""" + async with self._open_lock: + if self._opened: + return + reply: pb.HeartbeatReply = await self._call( + self._stub.OpenManager, + pb.ContextRequest(context=self._request_context()), + method="OpenManager", + lease=False, + timeout=5, ) - raise RuntimeError( - f"Host Bridge {method} 失败: {exc.code().name}: {exc.details()}{uncertainty}" - ) from exc + require_fields(reply, "alive") + if not reply.alive: + raise RuntimeError("Host Bridge 未确认 manager 登记") + self._opened = True + self._ensure_heartbeat() def _ensure_heartbeat(self) -> None: if self._heartbeat_task is None: @@ -385,22 +430,36 @@ def _ensure_heartbeat(self) -> None: ) async def _heartbeat_loop(self) -> None: + """暂时传输失败继续续期;租约丢失和身份错误终结旧 manager。""" + failures = 0 try: while True: - await asyncio.sleep(_HEARTBEAT_INTERVAL_S) - reply: pb.HeartbeatReply = await self._call( - self._stub.Heartbeat, - pb.ContextRequest(context=self._request_context()), - method="Heartbeat", - timeout=5, - ) + await asyncio.sleep(min(_HEARTBEAT_INTERVAL_S * (2 ** min(failures, 3)), 10)) + try: + reply: pb.HeartbeatReply = await self._call( + self._stub.Heartbeat, + pb.ContextRequest(context=self._request_context()), + method="Heartbeat", + lease=False, + timeout=5, + ) + except HostBridgeRpcError as exc: + if not exc.transient: + raise + failures += 1 + logger.warning("Host Bridge 心跳暂时失败,继续探测: %s", exc) + continue require_fields(reply, "alive") if not reply.alive: raise RuntimeError("Host Bridge 未确认 lease 存活") + if failures: + logger.info("Host Bridge 心跳恢复") + failures = 0 except asyncio.CancelledError: raise except Exception as exc: self._lease_error = exc + logger.error("Host Bridge manager 已失效: %s", exc) def _check_client_identity(socket_path: Path, boot_id: str, token: str) -> None: diff --git a/agent/host_bridge/filesystem.py b/agent/host_bridge/filesystem.py index c92f0fff0..716d29f35 100644 --- a/agent/host_bridge/filesystem.py +++ b/agent/host_bridge/filesystem.py @@ -7,8 +7,8 @@ import difflib import logging import os -from collections.abc import Awaitable, Callable -from dataclasses import dataclass +from collections.abc import Callable +from dataclasses import dataclass, field from pathlib import Path from typing import TYPE_CHECKING, Any, TypeVar @@ -32,6 +32,50 @@ class _FileMutationState: _FILE_MUTATION_LOCKS: dict[str, _FileMutationState] = {} +@dataclass +class _FileIoState: + slots: asyncio.Semaphore = field(default_factory=lambda: asyncio.Semaphore(4)) + users: int = 0 + + +_FILE_IO_SLOTS: dict[asyncio.AbstractEventLoop, _FileIoState] = {} + + +async def _run_file_io(fn: Callable[[], T]) -> T: + """最多四个磁盘操作并行;取消后仍等物理工作结束才归还锁与 owner。""" + # 1. 等待名额时可以取消;线程启动后不能把取消当作工作已结束。 + loop = asyncio.get_running_loop() + state = _FILE_IO_SLOTS.setdefault(loop, _FileIoState()) + state.users += 1 + try: + async with state.slots: + work = asyncio.create_task(asyncio.to_thread(fn)) + cancelled: asyncio.CancelledError | None = None + while not work.done(): + try: + await asyncio.shield(work) + except asyncio.CancelledError as exc: + cancelled = exc + except Exception: + # 实际错误由 result 取回;同时发生取消时保留两种失败。 + break + # 2. 到这里线程已结束,外层才可以释放文件锁和 manager operation。 + try: + result = work.result() + except Exception as exc: + if cancelled is not None: + raise BaseExceptionGroup("文件操作取消且物理工作失败", [cancelled, exc]) from None + raise + if cancelled is not None: + raise cancelled + return result + finally: + state.users -= 1 + if state.users == 0: + del _FILE_IO_SLOTS[loop] + + + def _is_inside(path: Path, allowed_dir: Path) -> bool: try: _ = path.relative_to(allowed_dir) @@ -106,12 +150,12 @@ def _get_file_mutation_key(file_path: Path) -> str: async def _run_with_file_mutation_lock( - file_path: Path, fn: Callable[[], Awaitable[T]] + file_path: Path, fn: Callable[[], T] ) -> T: """按规范化路径串行执行文件变更,并在异常或取消后回收锁状态。""" # 1. 登记当前调用,等待者也必须计入生命周期 - key = _get_file_mutation_key(file_path) + key = await _run_file_io(lambda: _get_file_mutation_key(file_path)) state = _FILE_MUTATION_LOCKS.get(key) if state is None: state = _FileMutationState(lock=asyncio.Lock()) @@ -121,7 +165,7 @@ async def _run_with_file_mutation_lock( try: # 2. 同一文件串行执行,取消也由 async with 释放底层锁 async with state.lock: - return await fn() + return await _run_file_io(fn) finally: # 3. 最后一个持有者或等待者退出后再移除路径映射 state.users -= 1 @@ -271,7 +315,7 @@ async def read_raw(self, path: str, **kwargs: Any) -> str | ToolResult: allowed_dir=self._allowed_dir, arguments={"path": path, **kwargs}, ) - return self.read_from_disk(path, **kwargs) + return await _run_file_io(lambda: self.read_from_disk(path, **kwargs)) def read_from_disk(self, path: str, **kwargs: Any) -> str | ToolResult: """Read host bytes without applying the current Turn model projection.""" @@ -363,9 +407,9 @@ async def execute(self, path: str, content: str, **kwargs: Any) -> str | ToolRes ) return result try: - file_path = _resolve_path(path, self._allowed_dir) + file_path = await _run_file_io(lambda: _resolve_path(path, self._allowed_dir)) - async def _write() -> str | ToolResult: + def _write() -> str | ToolResult: if file_path.exists() and file_path.is_dir(): return ToolResult( text=f"写入文件失败:目标路径是目录:{path}", is_error=True @@ -401,9 +445,9 @@ async def execute( ) return result try: - file_path = _resolve_path(path, self._allowed_dir) + file_path = await _run_file_io(lambda: _resolve_path(path, self._allowed_dir)) - async def _edit() -> str | ToolResult: + def _edit() -> str | ToolResult: if not file_path.exists(): return ToolResult(text=f"错误:文件不存在:{path}", is_error=True) if not file_path.is_file(): @@ -468,6 +512,10 @@ async def execute(self, path: str, **kwargs: Any) -> str | ToolResult: arguments={"path": path, **kwargs}, ) return result + return await _run_file_io(lambda: self._list_from_disk(path)) + + def _list_from_disk(self, path: str) -> str | ToolResult: + """在线程中完成路径解析、目录遍历和文件类型查询。""" try: dir_path = _resolve_path(path, self._allowed_dir) if not dir_path.exists(): diff --git a/agent/host_bridge/host_bridge.proto b/agent/host_bridge/host_bridge.proto index 9e7cbc79d..7ae858fa2 100644 --- a/agent/host_bridge/host_bridge.proto +++ b/agent/host_bridge/host_bridge.proto @@ -7,6 +7,7 @@ service HostBridge { rpc Inspect(ContextRequest) returns (IdentityReply); rpc ClaimBoot(ContextRequest) returns (ClaimBootReply); rpc Probe(ContextRequest) returns (IdentityReply); + rpc OpenManager(ContextRequest) returns (HeartbeatReply); rpc Heartbeat(ContextRequest) returns (HeartbeatReply); rpc Exec(ExecRequest) returns (ExecutionReply); rpc WriteStdin(WriteStdinRequest) returns (ExecutionReply); diff --git a/agent/host_bridge/host_bridge_pb2.py b/agent/host_bridge/host_bridge_pb2.py index 2ceb2c822..2abdeb264 100644 --- a/agent/host_bridge/host_bridge_pb2.py +++ b/agent/host_bridge/host_bridge_pb2.py @@ -24,7 +24,7 @@ -DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n#agent/host_bridge/host_bridge.proto\x12\x0f\x61kashic.host.v2\"\xd9\x01\n\x0eRequestContext\x12\x0f\n\x07\x62oot_id\x18\x01 \x01(\t\x12\x12\n\nmanager_id\x18\x02 \x01(\t\x12\x12\n\nrequest_id\x18\x03 \x01(\t\x12\x1f\n\x17\x65xpected_release_commit\x18\x04 \x01(\t\x12!\n\x19\x65xpected_toolchain_digest\x18\x05 \x01(\t\x12\x18\n\x0bsession_ref\x18\x06 \x01(\tH\x00\x88\x01\x01\x12\x14\n\x07turn_id\x18\x07 \x01(\tH\x01\x88\x01\x01\x42\x0e\n\x0c_session_refB\n\n\x08_turn_id\"B\n\x0e\x43ontextRequest\x12\x30\n\x07\x63ontext\x18\x01 \x01(\x0b\x32\x1f.akashic.host.v2.RequestContext\"W\n\rIdentityReply\x12\x16\n\x0erelease_commit\x18\x01 \x01(\t\x12\x18\n\x10toolchain_digest\x18\x02 \x01(\t\x12\x14\n\x0c\x63\x61pabilities\x18\x03 \x03(\t\"\xdb\x01\n\x0e\x43laimBootReply\x12\x15\n\rowner_boot_id\x18\x01 \x01(\t\x12\x1d\n\x10previous_boot_id\x18\x02 \x01(\tH\x00\x88\x01\x01\x12\"\n\x15\x63leaned_manager_count\x18\x03 \x01(\x04H\x01\x88\x01\x01\x12$\n\x17\x63leaned_execution_count\x18\x04 \x01(\x04H\x02\x88\x01\x01\x42\x13\n\x11_previous_boot_idB\x18\n\x16_cleaned_manager_countB\x1a\n\x18_cleaned_execution_count\".\n\x0eHeartbeatReply\x12\x12\n\x05\x61live\x18\x01 \x01(\x08H\x00\x88\x01\x01\x42\x08\n\x06_alive\"\xa1\x03\n\x0b\x45xecRequest\x12\x30\n\x07\x63ontext\x18\x01 \x01(\x0b\x32\x1f.akashic.host.v2.RequestContext\x12\x0f\n\x07\x63ommand\x18\x02 \x01(\t\x12\x0c\n\x04\x61rgv\x18\x03 \x03(\t\x12\x10\n\x03\x63wd\x18\x04 \x01(\tH\x00\x88\x01\x01\x12\x32\n\x03\x65nv\x18\x05 \x03(\x0b\x32%.akashic.host.v2.ExecRequest.EnvEntry\x12\x10\n\x03tty\x18\x06 \x01(\x08H\x01\x88\x01\x01\x12\x1a\n\ryield_time_ms\x18\x07 \x01(\x03H\x02\x88\x01\x01\x12\x1e\n\x11max_output_tokens\x18\x08 \x01(\x03H\x03\x88\x01\x01\x12\x1b\n\x0ehard_timeout_s\x18\t \x01(\x03H\x04\x88\x01\x01\x12\x19\n\x11owner_session_key\x18\n \x01(\t\x1a*\n\x08\x45nvEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\r\n\x05value\x18\x02 \x01(\t:\x02\x38\x01\x42\x06\n\x04_cwdB\x06\n\x04_ttyB\x10\n\x0e_yield_time_msB\x14\n\x12_max_output_tokensB\x11\n\x0f_hard_timeout_s\"\x8e\x02\n\x11WriteStdinRequest\x12\x30\n\x07\x63ontext\x18\x01 \x01(\x0b\x32\x1f.akashic.host.v2.RequestContext\x12\x19\n\x0c\x65xecution_id\x18\x02 \x01(\x03H\x00\x88\x01\x01\x12\x12\n\x05\x63hars\x18\x03 \x01(\tH\x01\x88\x01\x01\x12\x1a\n\ryield_time_ms\x18\x04 \x01(\x03H\x02\x88\x01\x01\x12\x1e\n\x11max_output_tokens\x18\x05 \x01(\x03H\x03\x88\x01\x01\x12\x19\n\x11owner_session_key\x18\x06 \x01(\tB\x0f\n\r_execution_idB\x08\n\x06_charsB\x10\n\x0e_yield_time_msB\x14\n\x12_max_output_tokens\"\xcc\x02\n\x0e\x45xecutionReply\x12\x13\n\x06output\x18\x01 \x01(\x0cH\x01\x88\x01\x01\x12\x19\n\x0cwall_time_ms\x18\x02 \x01(\x03H\x02\x88\x01\x01\x12!\n\x14original_token_count\x18\x03 \x01(\x03H\x03\x88\x01\x01\x12!\n\x14output_omitted_bytes\x18\x04 \x01(\x03H\x04\x88\x01\x01\x12\x16\n\x0c\x65xecution_id\x18\x05 \x01(\x03H\x00\x12\x13\n\texit_code\x18\x06 \x01(\x11H\x00\x12\x18\n\x0boutput_path\x18\x07 \x01(\tH\x05\x88\x01\x01\x12\x15\n\rfinish_reason\x18\x08 \x01(\tB\x08\n\x06resultB\t\n\x07_outputB\x0f\n\r_wall_time_msB\x17\n\x15_original_token_countB\x17\n\x15_output_omitted_bytesB\x0e\n\x0c_output_path\"\x86\x01\n\x0bStopRequest\x12\x30\n\x07\x63ontext\x18\x01 \x01(\x0b\x32\x1f.akashic.host.v2.RequestContext\x12\x19\n\x0c\x65xecution_id\x18\x02 \x01(\x03H\x00\x88\x01\x01\x12\x19\n\x11owner_session_key\x18\x03 \x01(\tB\x0f\n\r_execution_id\"-\n\tStopReply\x12\x14\n\x07stopped\x18\x01 \x01(\x08H\x00\x88\x01\x01\x42\n\n\x08_stopped\"[\n\x0cOwnerRequest\x12\x30\n\x07\x63ontext\x18\x01 \x01(\x0b\x32\x1f.akashic.host.v2.RequestContext\x12\x19\n\x11owner_session_key\x18\x02 \x01(\t\"a\n\x0e\x43leanupFailure\x12\x19\n\x0c\x65xecution_id\x18\x01 \x01(\x03H\x00\x88\x01\x01\x12\x12\n\nerror_type\x18\x02 \x01(\t\x12\x0f\n\x07message\x18\x03 \x01(\tB\x0f\n\r_execution_id\"e\n\x0c\x43leanupReply\x12\x11\n\tattempted\x18\x01 \x03(\x03\x12\x0f\n\x07\x63leaned\x18\x02 \x03(\x03\x12\x31\n\x08\x66\x61ilures\x18\x03 \x03(\x0b\x32\x1f.akashic.host.v2.CleanupFailure\".\n\x15\x41\x63tiveExecutionsReply\x12\x15\n\rexecution_ids\x18\x01 \x03(\x03\"\xa3\x02\n\x0b\x46ileRequest\x12\x30\n\x07\x63ontext\x18\x01 \x01(\x0b\x32\x1f.akashic.host.v2.RequestContext\x12\x18\n\x0b\x61llowed_dir\x18\x02 \x01(\tH\x01\x88\x01\x01\x12)\n\x04read\x18\x03 \x01(\x0b\x32\x19.akashic.host.v2.ReadFileH\x00\x12+\n\x05write\x18\x04 \x01(\x0b\x32\x1a.akashic.host.v2.WriteFileH\x00\x12)\n\x04\x65\x64it\x18\x05 \x01(\x0b\x32\x19.akashic.host.v2.EditFileH\x00\x12(\n\x04list\x18\x06 \x01(\x0b\x32\x18.akashic.host.v2.ListDirH\x00\x42\x0b\n\toperationB\x0e\n\x0c_allowed_dir\"d\n\x08ReadFile\x12\x11\n\x04path\x18\x01 \x01(\tH\x00\x88\x01\x01\x12\x13\n\x06offset\x18\x02 \x01(\x03H\x01\x88\x01\x01\x12\x12\n\x05limit\x18\x03 \x01(\x03H\x02\x88\x01\x01\x42\x07\n\x05_pathB\t\n\x07_offsetB\x08\n\x06_limit\"I\n\tWriteFile\x12\x11\n\x04path\x18\x01 \x01(\tH\x00\x88\x01\x01\x12\x14\n\x07\x63ontent\x18\x02 \x01(\tH\x01\x88\x01\x01\x42\x07\n\x05_pathB\n\n\x08_content\"\x98\x01\n\x08\x45\x64itFile\x12\x11\n\x04path\x18\x01 \x01(\tH\x00\x88\x01\x01\x12\x15\n\x08old_text\x18\x02 \x01(\tH\x01\x88\x01\x01\x12\x15\n\x08new_text\x18\x03 \x01(\tH\x02\x88\x01\x01\x12\x18\n\x0breplace_all\x18\x04 \x01(\x08H\x03\x88\x01\x01\x42\x07\n\x05_pathB\x0b\n\t_old_textB\x0b\n\t_new_textB\x0e\n\x0c_replace_all\"%\n\x07ListDir\x12\x11\n\x04path\x18\x01 \x01(\tH\x00\x88\x01\x01\x42\x07\n\x05_path\"\x7f\n\tFileReply\x12\x0e\n\x04text\x18\x01 \x01(\tH\x00\x12+\n\x05image\x18\x02 \x01(\x0b\x32\x1a.akashic.host.v2.FileImageH\x00\x12+\n\x05\x65rror\x18\x03 \x01(\x0b\x32\x1a.akashic.host.v2.FileErrorH\x00\x42\x08\n\x06result\"K\n\tFileError\x12\x11\n\x04text\x18\x01 \x01(\tH\x00\x88\x01\x01\x12\x15\n\x08is_error\x18\x02 \x01(\x08H\x01\x88\x01\x01\x42\x07\n\x05_textB\x0b\n\t_is_error\"f\n\tFileImage\x12\x11\n\x04text\x18\x01 \x01(\tH\x00\x88\x01\x01\x12\x11\n\tmime_type\x18\x02 \x01(\t\x12\x11\n\x04\x64\x61ta\x18\x03 \x01(\x0cH\x01\x88\x01\x01\x12\x0e\n\x06\x64\x65tail\x18\x04 \x01(\tB\x07\n\x05_textB\x07\n\x05_data\"g\n\x18SkillRequirementsRequest\x12\x30\n\x07\x63ontext\x18\x01 \x01(\x0b\x32\x1f.akashic.host.v2.RequestContext\x12\x0c\n\x04\x62ins\x18\x02 \x03(\t\x12\x0b\n\x03\x65nv\x18\x03 \x03(\t\"-\n\x10RequirementNames\x12\x0c\n\x04\x62ins\x18\x01 \x03(\t\x12\x0b\n\x03\x65nv\x18\x02 \x03(\t\"\x82\x01\n\x16SkillRequirementsReply\x12\x34\n\tavailable\x18\x01 \x01(\x0b\x32!.akashic.host.v2.RequirementNames\x12\x32\n\x07missing\x18\x02 \x01(\x0b\x32!.akashic.host.v2.RequirementNames2\xcb\x07\n\nHostBridge\x12J\n\x07Inspect\x12\x1f.akashic.host.v2.ContextRequest\x1a\x1e.akashic.host.v2.IdentityReply\x12M\n\tClaimBoot\x12\x1f.akashic.host.v2.ContextRequest\x1a\x1f.akashic.host.v2.ClaimBootReply\x12H\n\x05Probe\x12\x1f.akashic.host.v2.ContextRequest\x1a\x1e.akashic.host.v2.IdentityReply\x12M\n\tHeartbeat\x12\x1f.akashic.host.v2.ContextRequest\x1a\x1f.akashic.host.v2.HeartbeatReply\x12\x45\n\x04\x45xec\x12\x1c.akashic.host.v2.ExecRequest\x1a\x1f.akashic.host.v2.ExecutionReply\x12Q\n\nWriteStdin\x12\".akashic.host.v2.WriteStdinRequest\x1a\x1f.akashic.host.v2.ExecutionReply\x12@\n\x04Stop\x12\x1c.akashic.host.v2.StopRequest\x1a\x1a.akashic.host.v2.StopReply\x12N\n\x0eTerminateOwner\x12\x1d.akashic.host.v2.OwnerRequest\x1a\x1d.akashic.host.v2.CleanupReply\x12Q\n\x0fShutdownManager\x12\x1f.akashic.host.v2.ContextRequest\x1a\x1d.akashic.host.v2.CleanupReply\x12[\n\x10\x41\x63tiveExecutions\x12\x1f.akashic.host.v2.ContextRequest\x1a&.akashic.host.v2.ActiveExecutionsReply\x12\x44\n\x08\x46ileTool\x12\x1c.akashic.host.v2.FileRequest\x1a\x1a.akashic.host.v2.FileReply\x12g\n\x11SkillRequirements\x12).akashic.host.v2.SkillRequirementsRequest\x1a\'.akashic.host.v2.SkillRequirementsReplyb\x06proto3') +DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n#agent/host_bridge/host_bridge.proto\x12\x0f\x61kashic.host.v2\"\xd9\x01\n\x0eRequestContext\x12\x0f\n\x07\x62oot_id\x18\x01 \x01(\t\x12\x12\n\nmanager_id\x18\x02 \x01(\t\x12\x12\n\nrequest_id\x18\x03 \x01(\t\x12\x1f\n\x17\x65xpected_release_commit\x18\x04 \x01(\t\x12!\n\x19\x65xpected_toolchain_digest\x18\x05 \x01(\t\x12\x18\n\x0bsession_ref\x18\x06 \x01(\tH\x00\x88\x01\x01\x12\x14\n\x07turn_id\x18\x07 \x01(\tH\x01\x88\x01\x01\x42\x0e\n\x0c_session_refB\n\n\x08_turn_id\"B\n\x0e\x43ontextRequest\x12\x30\n\x07\x63ontext\x18\x01 \x01(\x0b\x32\x1f.akashic.host.v2.RequestContext\"W\n\rIdentityReply\x12\x16\n\x0erelease_commit\x18\x01 \x01(\t\x12\x18\n\x10toolchain_digest\x18\x02 \x01(\t\x12\x14\n\x0c\x63\x61pabilities\x18\x03 \x03(\t\"\xdb\x01\n\x0e\x43laimBootReply\x12\x15\n\rowner_boot_id\x18\x01 \x01(\t\x12\x1d\n\x10previous_boot_id\x18\x02 \x01(\tH\x00\x88\x01\x01\x12\"\n\x15\x63leaned_manager_count\x18\x03 \x01(\x04H\x01\x88\x01\x01\x12$\n\x17\x63leaned_execution_count\x18\x04 \x01(\x04H\x02\x88\x01\x01\x42\x13\n\x11_previous_boot_idB\x18\n\x16_cleaned_manager_countB\x1a\n\x18_cleaned_execution_count\".\n\x0eHeartbeatReply\x12\x12\n\x05\x61live\x18\x01 \x01(\x08H\x00\x88\x01\x01\x42\x08\n\x06_alive\"\xa1\x03\n\x0b\x45xecRequest\x12\x30\n\x07\x63ontext\x18\x01 \x01(\x0b\x32\x1f.akashic.host.v2.RequestContext\x12\x0f\n\x07\x63ommand\x18\x02 \x01(\t\x12\x0c\n\x04\x61rgv\x18\x03 \x03(\t\x12\x10\n\x03\x63wd\x18\x04 \x01(\tH\x00\x88\x01\x01\x12\x32\n\x03\x65nv\x18\x05 \x03(\x0b\x32%.akashic.host.v2.ExecRequest.EnvEntry\x12\x10\n\x03tty\x18\x06 \x01(\x08H\x01\x88\x01\x01\x12\x1a\n\ryield_time_ms\x18\x07 \x01(\x03H\x02\x88\x01\x01\x12\x1e\n\x11max_output_tokens\x18\x08 \x01(\x03H\x03\x88\x01\x01\x12\x1b\n\x0ehard_timeout_s\x18\t \x01(\x03H\x04\x88\x01\x01\x12\x19\n\x11owner_session_key\x18\n \x01(\t\x1a*\n\x08\x45nvEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\r\n\x05value\x18\x02 \x01(\t:\x02\x38\x01\x42\x06\n\x04_cwdB\x06\n\x04_ttyB\x10\n\x0e_yield_time_msB\x14\n\x12_max_output_tokensB\x11\n\x0f_hard_timeout_s\"\x8e\x02\n\x11WriteStdinRequest\x12\x30\n\x07\x63ontext\x18\x01 \x01(\x0b\x32\x1f.akashic.host.v2.RequestContext\x12\x19\n\x0c\x65xecution_id\x18\x02 \x01(\x03H\x00\x88\x01\x01\x12\x12\n\x05\x63hars\x18\x03 \x01(\tH\x01\x88\x01\x01\x12\x1a\n\ryield_time_ms\x18\x04 \x01(\x03H\x02\x88\x01\x01\x12\x1e\n\x11max_output_tokens\x18\x05 \x01(\x03H\x03\x88\x01\x01\x12\x19\n\x11owner_session_key\x18\x06 \x01(\tB\x0f\n\r_execution_idB\x08\n\x06_charsB\x10\n\x0e_yield_time_msB\x14\n\x12_max_output_tokens\"\xcc\x02\n\x0e\x45xecutionReply\x12\x13\n\x06output\x18\x01 \x01(\x0cH\x01\x88\x01\x01\x12\x19\n\x0cwall_time_ms\x18\x02 \x01(\x03H\x02\x88\x01\x01\x12!\n\x14original_token_count\x18\x03 \x01(\x03H\x03\x88\x01\x01\x12!\n\x14output_omitted_bytes\x18\x04 \x01(\x03H\x04\x88\x01\x01\x12\x16\n\x0c\x65xecution_id\x18\x05 \x01(\x03H\x00\x12\x13\n\texit_code\x18\x06 \x01(\x11H\x00\x12\x18\n\x0boutput_path\x18\x07 \x01(\tH\x05\x88\x01\x01\x12\x15\n\rfinish_reason\x18\x08 \x01(\tB\x08\n\x06resultB\t\n\x07_outputB\x0f\n\r_wall_time_msB\x17\n\x15_original_token_countB\x17\n\x15_output_omitted_bytesB\x0e\n\x0c_output_path\"\x86\x01\n\x0bStopRequest\x12\x30\n\x07\x63ontext\x18\x01 \x01(\x0b\x32\x1f.akashic.host.v2.RequestContext\x12\x19\n\x0c\x65xecution_id\x18\x02 \x01(\x03H\x00\x88\x01\x01\x12\x19\n\x11owner_session_key\x18\x03 \x01(\tB\x0f\n\r_execution_id\"-\n\tStopReply\x12\x14\n\x07stopped\x18\x01 \x01(\x08H\x00\x88\x01\x01\x42\n\n\x08_stopped\"[\n\x0cOwnerRequest\x12\x30\n\x07\x63ontext\x18\x01 \x01(\x0b\x32\x1f.akashic.host.v2.RequestContext\x12\x19\n\x11owner_session_key\x18\x02 \x01(\t\"a\n\x0e\x43leanupFailure\x12\x19\n\x0c\x65xecution_id\x18\x01 \x01(\x03H\x00\x88\x01\x01\x12\x12\n\nerror_type\x18\x02 \x01(\t\x12\x0f\n\x07message\x18\x03 \x01(\tB\x0f\n\r_execution_id\"e\n\x0c\x43leanupReply\x12\x11\n\tattempted\x18\x01 \x03(\x03\x12\x0f\n\x07\x63leaned\x18\x02 \x03(\x03\x12\x31\n\x08\x66\x61ilures\x18\x03 \x03(\x0b\x32\x1f.akashic.host.v2.CleanupFailure\".\n\x15\x41\x63tiveExecutionsReply\x12\x15\n\rexecution_ids\x18\x01 \x03(\x03\"\xa3\x02\n\x0b\x46ileRequest\x12\x30\n\x07\x63ontext\x18\x01 \x01(\x0b\x32\x1f.akashic.host.v2.RequestContext\x12\x18\n\x0b\x61llowed_dir\x18\x02 \x01(\tH\x01\x88\x01\x01\x12)\n\x04read\x18\x03 \x01(\x0b\x32\x19.akashic.host.v2.ReadFileH\x00\x12+\n\x05write\x18\x04 \x01(\x0b\x32\x1a.akashic.host.v2.WriteFileH\x00\x12)\n\x04\x65\x64it\x18\x05 \x01(\x0b\x32\x19.akashic.host.v2.EditFileH\x00\x12(\n\x04list\x18\x06 \x01(\x0b\x32\x18.akashic.host.v2.ListDirH\x00\x42\x0b\n\toperationB\x0e\n\x0c_allowed_dir\"d\n\x08ReadFile\x12\x11\n\x04path\x18\x01 \x01(\tH\x00\x88\x01\x01\x12\x13\n\x06offset\x18\x02 \x01(\x03H\x01\x88\x01\x01\x12\x12\n\x05limit\x18\x03 \x01(\x03H\x02\x88\x01\x01\x42\x07\n\x05_pathB\t\n\x07_offsetB\x08\n\x06_limit\"I\n\tWriteFile\x12\x11\n\x04path\x18\x01 \x01(\tH\x00\x88\x01\x01\x12\x14\n\x07\x63ontent\x18\x02 \x01(\tH\x01\x88\x01\x01\x42\x07\n\x05_pathB\n\n\x08_content\"\x98\x01\n\x08\x45\x64itFile\x12\x11\n\x04path\x18\x01 \x01(\tH\x00\x88\x01\x01\x12\x15\n\x08old_text\x18\x02 \x01(\tH\x01\x88\x01\x01\x12\x15\n\x08new_text\x18\x03 \x01(\tH\x02\x88\x01\x01\x12\x18\n\x0breplace_all\x18\x04 \x01(\x08H\x03\x88\x01\x01\x42\x07\n\x05_pathB\x0b\n\t_old_textB\x0b\n\t_new_textB\x0e\n\x0c_replace_all\"%\n\x07ListDir\x12\x11\n\x04path\x18\x01 \x01(\tH\x00\x88\x01\x01\x42\x07\n\x05_path\"\x7f\n\tFileReply\x12\x0e\n\x04text\x18\x01 \x01(\tH\x00\x12+\n\x05image\x18\x02 \x01(\x0b\x32\x1a.akashic.host.v2.FileImageH\x00\x12+\n\x05\x65rror\x18\x03 \x01(\x0b\x32\x1a.akashic.host.v2.FileErrorH\x00\x42\x08\n\x06result\"K\n\tFileError\x12\x11\n\x04text\x18\x01 \x01(\tH\x00\x88\x01\x01\x12\x15\n\x08is_error\x18\x02 \x01(\x08H\x01\x88\x01\x01\x42\x07\n\x05_textB\x0b\n\t_is_error\"f\n\tFileImage\x12\x11\n\x04text\x18\x01 \x01(\tH\x00\x88\x01\x01\x12\x11\n\tmime_type\x18\x02 \x01(\t\x12\x11\n\x04\x64\x61ta\x18\x03 \x01(\x0cH\x01\x88\x01\x01\x12\x0e\n\x06\x64\x65tail\x18\x04 \x01(\tB\x07\n\x05_textB\x07\n\x05_data\"g\n\x18SkillRequirementsRequest\x12\x30\n\x07\x63ontext\x18\x01 \x01(\x0b\x32\x1f.akashic.host.v2.RequestContext\x12\x0c\n\x04\x62ins\x18\x02 \x03(\t\x12\x0b\n\x03\x65nv\x18\x03 \x03(\t\"-\n\x10RequirementNames\x12\x0c\n\x04\x62ins\x18\x01 \x03(\t\x12\x0b\n\x03\x65nv\x18\x02 \x03(\t\"\x82\x01\n\x16SkillRequirementsReply\x12\x34\n\tavailable\x18\x01 \x01(\x0b\x32!.akashic.host.v2.RequirementNames\x12\x32\n\x07missing\x18\x02 \x01(\x0b\x32!.akashic.host.v2.RequirementNames2\x9c\x08\n\nHostBridge\x12J\n\x07Inspect\x12\x1f.akashic.host.v2.ContextRequest\x1a\x1e.akashic.host.v2.IdentityReply\x12M\n\tClaimBoot\x12\x1f.akashic.host.v2.ContextRequest\x1a\x1f.akashic.host.v2.ClaimBootReply\x12H\n\x05Probe\x12\x1f.akashic.host.v2.ContextRequest\x1a\x1e.akashic.host.v2.IdentityReply\x12O\n\x0bOpenManager\x12\x1f.akashic.host.v2.ContextRequest\x1a\x1f.akashic.host.v2.HeartbeatReply\x12M\n\tHeartbeat\x12\x1f.akashic.host.v2.ContextRequest\x1a\x1f.akashic.host.v2.HeartbeatReply\x12\x45\n\x04\x45xec\x12\x1c.akashic.host.v2.ExecRequest\x1a\x1f.akashic.host.v2.ExecutionReply\x12Q\n\nWriteStdin\x12\".akashic.host.v2.WriteStdinRequest\x1a\x1f.akashic.host.v2.ExecutionReply\x12@\n\x04Stop\x12\x1c.akashic.host.v2.StopRequest\x1a\x1a.akashic.host.v2.StopReply\x12N\n\x0eTerminateOwner\x12\x1d.akashic.host.v2.OwnerRequest\x1a\x1d.akashic.host.v2.CleanupReply\x12Q\n\x0fShutdownManager\x12\x1f.akashic.host.v2.ContextRequest\x1a\x1d.akashic.host.v2.CleanupReply\x12[\n\x10\x41\x63tiveExecutions\x12\x1f.akashic.host.v2.ContextRequest\x1a&.akashic.host.v2.ActiveExecutionsReply\x12\x44\n\x08\x46ileTool\x12\x1c.akashic.host.v2.FileRequest\x1a\x1a.akashic.host.v2.FileReply\x12g\n\x11SkillRequirements\x12).akashic.host.v2.SkillRequirementsRequest\x1a\'.akashic.host.v2.SkillRequirementsReplyb\x06proto3') _globals = globals() _builder.BuildMessageAndEnumDescriptors(DESCRIPTOR, _globals) @@ -86,5 +86,5 @@ _globals['_SKILLREQUIREMENTSREPLY']._serialized_start=3386 _globals['_SKILLREQUIREMENTSREPLY']._serialized_end=3516 _globals['_HOSTBRIDGE']._serialized_start=3519 - _globals['_HOSTBRIDGE']._serialized_end=4490 + _globals['_HOSTBRIDGE']._serialized_end=4571 # @@protoc_insertion_point(module_scope) diff --git a/agent/host_bridge/host_bridge_pb2_grpc.py b/agent/host_bridge/host_bridge_pb2_grpc.py index 9494bbcce..cff32b327 100644 --- a/agent/host_bridge/host_bridge_pb2_grpc.py +++ b/agent/host_bridge/host_bridge_pb2_grpc.py @@ -50,6 +50,11 @@ def __init__(self, channel): request_serializer=agent_dot_host__bridge_dot_host__bridge__pb2.ContextRequest.SerializeToString, response_deserializer=agent_dot_host__bridge_dot_host__bridge__pb2.IdentityReply.FromString, _registered_method=True) + self.OpenManager = channel.unary_unary( + '/akashic.host.v2.HostBridge/OpenManager', + request_serializer=agent_dot_host__bridge_dot_host__bridge__pb2.ContextRequest.SerializeToString, + response_deserializer=agent_dot_host__bridge_dot_host__bridge__pb2.HeartbeatReply.FromString, + _registered_method=True) self.Heartbeat = channel.unary_unary( '/akashic.host.v2.HostBridge/Heartbeat', request_serializer=agent_dot_host__bridge_dot_host__bridge__pb2.ContextRequest.SerializeToString, @@ -119,6 +124,12 @@ def Probe(self, request, context): context.set_details('Method not implemented!') raise NotImplementedError('Method not implemented!') + def OpenManager(self, request, context): + """Missing associated documentation comment in .proto file.""" + context.set_code(grpc.StatusCode.UNIMPLEMENTED) + context.set_details('Method not implemented!') + raise NotImplementedError('Method not implemented!') + def Heartbeat(self, request, context): """Missing associated documentation comment in .proto file.""" context.set_code(grpc.StatusCode.UNIMPLEMENTED) @@ -191,6 +202,11 @@ def add_HostBridgeServicer_to_server(servicer, server): request_deserializer=agent_dot_host__bridge_dot_host__bridge__pb2.ContextRequest.FromString, response_serializer=agent_dot_host__bridge_dot_host__bridge__pb2.IdentityReply.SerializeToString, ), + 'OpenManager': grpc.unary_unary_rpc_method_handler( + servicer.OpenManager, + request_deserializer=agent_dot_host__bridge_dot_host__bridge__pb2.ContextRequest.FromString, + response_serializer=agent_dot_host__bridge_dot_host__bridge__pb2.HeartbeatReply.SerializeToString, + ), 'Heartbeat': grpc.unary_unary_rpc_method_handler( servicer.Heartbeat, request_deserializer=agent_dot_host__bridge_dot_host__bridge__pb2.ContextRequest.FromString, @@ -329,6 +345,33 @@ def Probe(request, metadata, _registered_method=True) + @staticmethod + def OpenManager(request, + target, + options=(), + channel_credentials=None, + call_credentials=None, + insecure=False, + compression=None, + wait_for_ready=None, + timeout=None, + metadata=None): + return grpc.experimental.unary_unary( + request, + target, + '/akashic.host.v2.HostBridge/OpenManager', + agent_dot_host__bridge_dot_host__bridge__pb2.ContextRequest.SerializeToString, + agent_dot_host__bridge_dot_host__bridge__pb2.HeartbeatReply.FromString, + options, + channel_credentials, + insecure, + call_credentials, + compression, + wait_for_ready, + timeout, + metadata, + _registered_method=True) + @staticmethod def Heartbeat(request, target, diff --git a/agent/host_bridge/monitor.py b/agent/host_bridge/monitor.py index 5b010486d..a5e91a3cc 100644 --- a/agent/host_bridge/monitor.py +++ b/agent/host_bridge/monitor.py @@ -1,14 +1,33 @@ from __future__ import annotations import asyncio +import logging import os from collections.abc import Coroutine +from dataclasses import asdict, dataclass +from datetime import UTC, datetime from pathlib import Path -from typing import Any +from typing import Any, Literal -from agent.host_bridge.client import HostBridgeShellProcessManager +from core.common.diagnostic_log import log_event + +from agent.host_bridge.client import HostBridgeRpcError, HostBridgeShellProcessManager _MONITOR_INTERVAL_S = 2.0 +logger = logging.getLogger(__name__) + + +@dataclass +class HostBridgeStatus: + """由 App 监控任务更新的短命状态;只描述连接,不承诺旧 manager 仍存活。""" + + state: Literal["disabled", "checking", "healthy", "degraded"] = "disabled" + failures: int = 0 + code: str | None = None + checked_at: str | None = None + + def snapshot(self) -> dict[str, Any]: + return asdict(self) async def claim_host_bridge_boot() -> dict[str, Any] | None: @@ -24,13 +43,14 @@ async def claim_host_bridge_boot() -> dict[str, Any] | None: await manager.close_transport() -def build_host_bridge_monitor() -> Coroutine[Any, Any, None] | None: +def build_host_bridge_monitor(status: HostBridgeStatus) -> Coroutine[Any, Any, None] | None: """Build the required Core liveness monitor for host-bridge mode.""" identity = _configured_bridge_identity() if identity is None: return None - return _monitor(*identity) + status.state = "checking" + return _monitor(*identity, status=status) def _configured_bridge_identity() -> tuple[Path, str, str, str, str] | None: @@ -62,6 +82,8 @@ async def _monitor( token: str, release_commit: str, toolchain_digest: str, + *, + status: HostBridgeStatus, ) -> None: manager = HostBridgeShellProcessManager( socket_path, @@ -72,7 +94,25 @@ async def _monitor( ) try: while True: - await manager.probe() - await asyncio.sleep(_MONITOR_INTERVAL_S) + try: + await manager.probe() + except HostBridgeRpcError as exc: + status.checked_at = datetime.now(UTC).isoformat() + status.state = "degraded" + status.failures += 1 + status.code = exc.code.name + log_event(logger, logging.WARNING, "host_bridge.degraded", + reason=exc.code.name, counts=f"failures:{status.failures}") + if not exc.transient: + raise + else: + if status.failures: + log_event(logger, logging.INFO, "host_bridge.recovered", + counts=f"failures:{status.failures}") + status.checked_at = datetime.now(UTC).isoformat() + status.state = "healthy" + status.failures = 0 + status.code = None + await asyncio.sleep(min(_MONITOR_INTERVAL_S * (2 ** min(status.failures, 3)), 10)) finally: await manager.close_transport() diff --git a/agent/host_bridge/server.py b/agent/host_bridge/server.py index 2ee42d459..f14214387 100644 --- a/agent/host_bridge/server.py +++ b/agent/host_bridge/server.py @@ -67,6 +67,14 @@ def __post_init__(self) -> None: self.operations_drained.set() +class _ManagerUnavailable(Exception): + """manager 已停止接纳操作,不能继续复用。""" + + +class _ManagerNotFound(Exception): + """已登记的 manager 不再存在,旧执行句柄不能恢复使用。""" + + def _rpc[Request: Message, Reply: Message]( handler: Callable[["HostBridgeService", Request], Awaitable[Reply]], ) -> Callable[ @@ -114,6 +122,14 @@ async def run( except asyncio.CancelledError: # 2. 取消只结束本次 RPC 等待,不承诺进程未执行或输入未写入。 raise + except _ManagerUnavailable as exc: + self._log_rpc_failure(method, identity.request_id, identity.boot_id, + identity.manager_id, started, exc, "manager_unavailable") + await context.abort(grpc.StatusCode.FAILED_PRECONDITION, str(exc)) + except _ManagerNotFound as exc: + self._log_rpc_failure(method, identity.request_id, identity.boot_id, + identity.manager_id, started, exc, "manager_not_found") + await context.abort(grpc.StatusCode.NOT_FOUND, str(exc)) except (KeyError, TypeError, ValueError) as exc: self._log_rpc_failure( method, @@ -324,7 +340,8 @@ async def ClaimBoot(self, request: pb.ContextRequest) -> pb.ClaimBootReply: @_rpc async def Probe(self, request: pb.ContextRequest) -> pb.IdentityReply: - _ = await self._lease(request.context) + async with self._lock: + self._assert_active_boot(request.context.boot_id) return self._probe_payload() def _probe_payload(self) -> pb.IdentityReply: @@ -344,6 +361,12 @@ def _probe_payload(self) -> pb.IdentityReply: ], ) + @_rpc + async def OpenManager(self, request: pb.ContextRequest) -> pb.HeartbeatReply: + """唯一的首次登记入口,业务调用和心跳都不能创建 manager。""" + _ = await self._lease(request.context, create=True) + return pb.HeartbeatReply(alive=True) + @_rpc async def Heartbeat(self, request: pb.ContextRequest) -> pb.HeartbeatReply: _ = await self._lease(request.context) @@ -440,7 +463,7 @@ async def ShutdownManager(self, request: pb.ContextRequest) -> pb.CleanupReply: self._assert_active_boot(key[0]) lease = self._managers.get(key) if lease is None: - return encode_cleanup(ExecutionCleanupReport((), (), ())) + raise _ManagerNotFound("Host Bridge manager 已不存在,无法确认本次清理") lease.reaping = True await lease.operations_drained.wait() report = await lease.manager.shutdown() @@ -479,9 +502,9 @@ async def FileTool(self, request: pb.FileRequest) -> pb.FileReply: if read.HasField("limit"): require_positive(read.limit, "limit") async with self._manager_operation(request.context): - result = ReadFileOperation( + result = await ReadFileOperation( allowed_dir=allowed_dir, enable_bridge=False - ).read_from_disk( + ).read_raw( read.path, offset=read.offset, limit=read.limit if read.HasField("limit") else None, @@ -563,12 +586,14 @@ def _log_rpc_failure( exc_info=True, ) - async def _lease(self, context: pb.RequestContext) -> _ManagerLease: + async def _lease(self, context: pb.RequestContext, *, create: bool = False) -> _ManagerLease: key = (context.boot_id, context.manager_id) async with self._lock: self._assert_active_boot(key[0]) lease = self._managers.get(key) if lease is None: + if not create: + raise _ManagerNotFound("Host Bridge manager 已不存在,旧执行句柄已失效") manager_root = self._artifact_root / key[0] / key[1] lease = _ManagerLease( ShellProcessManager(output_dir=manager_root), @@ -577,9 +602,9 @@ async def _lease(self, context: pb.RequestContext) -> _ManagerLease: self._managers[key] = lease else: if lease.cleanup_failure is not None: - raise RuntimeError("Host Bridge manager cleanup 未确认,拒绝复用") + raise _ManagerUnavailable("Host Bridge manager cleanup 未确认,拒绝复用") if lease.reaping: - raise RuntimeError("Host Bridge manager lease 正在回收,拒绝复用") + raise _ManagerUnavailable("Host Bridge manager lease 正在回收,拒绝复用") lease.last_seen = time.monotonic() return lease @@ -595,7 +620,7 @@ async def _manager_operation( async with self._lock: self._assert_active_boot(key[0]) if self._managers.get(key) is not lease or lease.reaping: - raise RuntimeError("Host Bridge manager admission 已关闭,拒绝执行") + raise _ManagerUnavailable("Host Bridge manager admission 已关闭,拒绝执行") lease.active_operations += 1 lease.operations_drained.clear() try: diff --git a/bootstrap/app.py b/bootstrap/app.py index 15135e373..a6f8104fa 100644 --- a/bootstrap/app.py +++ b/bootstrap/app.py @@ -13,7 +13,7 @@ from agent.config import resolve_app_server_endpoint from agent.control.service import ControlService -from agent.host_bridge.monitor import build_host_bridge_monitor +from agent.host_bridge.monitor import HostBridgeStatus, build_host_bridge_monitor from agent.host_bridge.monitor import claim_host_bridge_boot from agent.restart import RestartGate from agent.config_models import Config @@ -200,6 +200,7 @@ def __init__( ) self.restart_gate = restart_gate self.readiness = readiness + self.host_bridge_status = HostBridgeStatus() self.http_resources = SharedHttpResources() self.app_server: SocketAppServer | None = None self.control_service: ControlService | None = None @@ -245,6 +246,7 @@ async def start(self) -> None: self.dashboard_server = build_dashboard_server( workspace=self.workspace, plugin_manager=manager, + host_bridge_status=self.host_bridge_status.snapshot, ) await self.core.start() if self.readiness is not None: @@ -284,7 +286,7 @@ async def start(self) -> None: self.readiness.mark_stage("channels.ready") if plugin_manager is None: raise RuntimeError("插件 Runtime 不可用") - host_bridge_monitor = build_host_bridge_monitor() + host_bridge_monitor = build_host_bridge_monitor(self.host_bridge_status) self.tasks = [] if host_bridge_monitor is not None: self.tasks.append(host_bridge_monitor) diff --git a/bootstrap/dashboard_api.py b/bootstrap/dashboard_api.py index 739118cc0..ac67c86e6 100644 --- a/bootstrap/dashboard_api.py +++ b/bootstrap/dashboard_api.py @@ -1,8 +1,9 @@ from __future__ import annotations import logging +from collections.abc import Callable from pathlib import Path -from typing import TYPE_CHECKING, cast +from typing import TYPE_CHECKING, Any, cast if TYPE_CHECKING: from agent.plugins.manager import PluginManager @@ -58,12 +59,19 @@ def create_dashboard_app( workspace: Path, *, plugin_manager: object | None = None, + host_bridge_status: Callable[[], dict[str, Any]] | None = None, ) -> FastAPI: workspace.mkdir(parents=True, exist_ok=True) project_root = Path(__file__).resolve().parent.parent static_dir = project_root / "static" / "dashboard" app = FastAPI(title="Akashic Dashboard API") + if host_bridge_status is not None: + + @app.get("/api/runtime/host-bridge") + async def read_host_bridge_status() -> dict[str, Any]: + return host_bridge_status() + # Vite 构建产物被 gitignore,新 clone 或 CI 环境可能没有该目录。 # 预先创建目录并在挂载时关闭目录检查,避免 app 创建依赖构建是否执行; # dashboard_index() 会在入口文件缺失时报告错误。 @@ -125,11 +133,13 @@ def _build_dashboard_uvicorn_config( port: int | None, uds: str | None = None, plugin_manager: object | None = None, + host_bridge_status: Callable[[], dict[str, Any]] | None = None, ) -> uvicorn.Config: config = uvicorn.Config( create_dashboard_app( workspace, plugin_manager=plugin_manager, + host_bridge_status=host_bridge_status, ), host=host or "127.0.0.1", port=port or 2236, @@ -147,6 +157,7 @@ def build_dashboard_server( port: int | None = None, uds: str | None = None, plugin_manager: object | None = None, + host_bridge_status: Callable[[], dict[str, Any]] | None = None, ) -> uvicorn.Server: config = _build_dashboard_uvicorn_config( workspace=workspace, @@ -154,5 +165,6 @@ def build_dashboard_server( port=port, uds=uds, plugin_manager=plugin_manager, + host_bridge_status=host_bridge_status, ) return uvicorn.Server(config) diff --git a/bootstrap/web_shell.py b/bootstrap/web_shell.py index 873ae5c04..084a131e4 100644 --- a/bootstrap/web_shell.py +++ b/bootstrap/web_shell.py @@ -133,6 +133,10 @@ async def shell_state() -> dict[str, object]: "chatReady": chat_ready, } + @app.get("/api/runtime/host-bridge") + async def proxy_host_bridge_status(request: Request) -> Response: + return await _proxy_http(request, dashboard_socket, "/api/runtime/host-bridge") + @app.api_route( "/api/chat/{proxy_path:path}", methods=["GET", "HEAD", "POST", "PUT", "PATCH", "DELETE", "OPTIONS"], diff --git a/docker/debug/host_bridge_notice.mjs b/docker/debug/host_bridge_notice.mjs new file mode 100644 index 000000000..6b371c6f8 --- /dev/null +++ b/docker/debug/host_bridge_notice.mjs @@ -0,0 +1,54 @@ +// 用真实 React 组件验证故障、恢复和未知状态,HTTP 响应由实验控制。 +import { build } from "esbuild"; +import { chromium } from "playwright-core"; +import { createServer } from "node:http"; +import { mkdtemp, readFile, rm } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { resolve } from "node:path"; + +const output = await mkdtemp(resolve(tmpdir(), "bridge-notice-")); +let state = "degraded"; +let httpFailure = false; +let browser; +const server = createServer(async (request, response) => { + if (request.url === "/api/runtime/host-bridge") { + response.writeHead(httpFailure ? 503 : 200, { "Content-Type": "application/json" }); + response.end(JSON.stringify({ state })); + } else if (request.url === "/app.js") { + response.writeHead(200, { "Content-Type": "text/javascript" }); + response.end(await readFile(resolve(output, "app.js"))); + } else { + response.writeHead(200, { "Content-Type": "text/html; charset=utf-8" }); + response.end('
'); + } +}); +try { + await build({ + stdin: { + contents: 'import {createRoot} from "react-dom/client"; import {HostBridgeNotice} from "./frontend/chat/src/host-bridge-notice"; createRoot(document.getElementById("root")).render(当前未加载回复插件
: null} +{notice}
: null; +}