diff --git a/.agents/skills/cc-remote-deploy/SKILL.md b/.agents/skills/cc-remote-deploy/SKILL.md index 93082e04..a33e90b1 100644 --- a/.agents/skills/cc-remote-deploy/SKILL.md +++ b/.agents/skills/cc-remote-deploy/SKILL.md @@ -37,6 +37,27 @@ migration, remaining in-process `/btw` turns, and deferred queues separately; wait for them to drain instead of interrupting work. Do not restart or replace an active SDK service to satisfy a version/readiness check. +## Retain rollback and all live dependencies + +Apply [deployment backup retention](../../../deploy/README.md#deployment-backup-retention) +on every in-scope host/install. Use the repository's `deploy/cleanup.py` with a +private, reviewed inventory; do not write an ad hoc `rm`/`rmtree` retention script. +Keep the active installation, one complete previous rollback generation, and +every live dependency. Dependency protection overrides the generation count. +Include all unresolved transactions and service dependencies in the inventory; +unknown provenance or incomplete process visibility means deferred cleanup. + +Follow the documented preview, quarantine, fresh acceptance and final removal +steps. Process argv alone is not evidence that a release is unused: cwd, +executables, open files and retained symlink targets count too. A cleanup that +fails or loses control remains an unresolved journal, not a reason to retry. +After cleanup, repeat live config checks as each Codex service user and verify +the Wrapper session route and public health. A pre-cleanup readiness receipt is +not final acceptance. Report retained/deferred paths and removed allocated bytes; +do not equate those bytes with actual free-space gain. No daemon or active task +may be stopped solely to make an artifact deletable. Installers do not prune +automatically. + ## Codex CLI sharing is an acceptance check For every enabled Codex **Code** account, follow diff --git a/AGENTS.md b/AGENTS.md index 3ae3a769..adfda933 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -31,6 +31,10 @@ transactions. A dropped SSH/control connection is an unknown result: inspect the original transaction and live state before deciding whether a retry is safe. Deployment is complete only after protocol/build identity, service stability, public health, and expected Wrapper connectivity are verified. +Use `deploy/cleanup.py` for reviewed retention inventories, never an ad hoc +deletion script. Preserve process cwd/open-file and retained runtime dependencies +even beyond the rollback generation limit. Repeat acceptance after cleanup, +including live Codex config reads; incomplete visibility means retain the files. For Codex Code, also follow `deploy/README.md`'s shared-control-plane acceptance: verify each account's daily CLI and Wrapper connect to the same official app-server, not a private stdio fallback. Do not force takeover or kill a live @@ -106,7 +110,7 @@ attachment and optional App-control MCP tools are separate user choices. transport, never the caller's Origin. Uvicorn trusts forwarded transport metadata only from loopback Caddy. Never put tokens in URLs or protocol message bodies; logging redacts token/password fields. -- **Protocol version gate**: current wire protocol v73 is declared by +- **Protocol version gate**: current wire protocol v74 is declared by `PROTOCOL_VERSION` in both `protocol.py` and `web/src/protocol.ts`. `deserialize` hard-rejects a version mismatch, and `_Base` is `extra="forbid"`, so ANY protocol change must be deployed to all diff --git a/CHANGELOG.md b/CHANGELOG.md index 4e6f0655..bc019913 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,14 @@ [中文](CHANGELOG_zh.md) +## Unreleased + +- Display native Codex cross-session messages and open their source conversations + in a scoped read-only history view. Wire protocol v74 requires a coordinated + Relay, Web, and Wrapper upgrade. +- Keep at most one complete previous deployment rollback generation, preserving + active dependencies and unresolved transaction recovery files. + ## v4.0.7 Add model-specific Codex speed selection and repair Claude recovery and private diff --git a/CHANGELOG_zh.md b/CHANGELOG_zh.md index ac653ccf..3d45815c 100644 --- a/CHANGELOG_zh.md +++ b/CHANGELOG_zh.md @@ -2,6 +2,12 @@ [English](CHANGELOG.md) +## 未发布 + +- 展示 Codex 原生跨会话消息,可在只读历史面板中查看对应来源会话,并保留设备与账号边界。 + 通信协议升级到 v74,需协调更新 Relay、Web 和 Wrapper。 +- 每次部署至多保留一套完整的上一版回滚备份;仍被运行服务引用的文件和未完成事务的恢复文件除外。 + ## v4.0.7 新增按模型选择 Codex 速度,修复 Claude 会话恢复和临时侧聊清理。 diff --git a/CLAUDE.md b/CLAUDE.md index 91248f8a..cb210aec 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -18,6 +18,9 @@ For deployment, upgrade, verification or recovery, read the repository skill at even if this client does not discover `.agents/skills` automatically. It routes to the maintained [`deploy/README.md`](deploy/README.md) automation contract, installation paths and shared-control acceptance; do not invent another flow. +Use `deploy/cleanup.py` for retention; preserve live cwd/open-file and runtime +dependencies beyond the rollback count. Incomplete visibility means retain the +files. Final acceptance, including live Codex config reads, runs after cleanup. Codex Code acceptance requires the daily CLI and Wrapper to use the same official daemon for each account. An online Web UI alone is insufficient. Never @@ -97,7 +100,7 @@ a separate choice; sharing alone does not authorize them. `useLayoutEffect` is deliberately dependency-free — late virtualizer/image measurements settle without a React render, and constraining it to its read set reintroduces a full-viewport jump on touch release. -- **Protocol version gate**: current wire protocol v73 is declared by +- **Protocol version gate**: current wire protocol v74 is declared by `PROTOCOL_VERSION` in both `protocol.py` and `web/src/protocol.ts`. `deserialize` hard-rejects a version mismatch, and `_Base` is `extra="forbid"`, so ANY protocol change must be deployed to all diff --git a/README.md b/README.md index edcfa048..8e64734e 100644 --- a/README.md +++ b/README.md @@ -4,7 +4,7 @@ 自托管 · 多会话 · 多设备 · 实时工具过程 · Code / Work · Web / PWA / TUI -**产品版本:v4.0.7** · Wire protocol v73 +**产品版本:v4.0.7** · 当前开发版 Wire protocol v74 [English](README_en.md) · [功能对照](#引擎与功能) · [快速开始](#快速开始) · [终端工作台](#terminal-workspace) · [安装与升级](#安装与升级) · [文档](#文档) · [更新记录](CHANGELOG_zh.md) diff --git a/README_en.md b/README_en.md index 702771ff..82ac3482 100644 --- a/README_en.md +++ b/README_en.md @@ -4,7 +4,7 @@ Self-hosted · Multiple sessions and devices · Live tool activity · Code / Work · Web / PWA / TUI -**Product version: v4.0.7** · Wire protocol v73 +**Product version: v4.0.7** · Current development wire protocol v74 [中文](README.md) · [Engine comparison](#engines-and-features) · [Quick start](#quick-start) · [Terminal workspace](#terminal-workspace) · [Install and upgrade](#install-and-upgrade) · [Documentation](#documentation) · [Changelog](CHANGELOG.md) diff --git a/cc_remote/protocol.py b/cc_remote/protocol.py index 170281a3..b64b1595 100644 --- a/cc_remote/protocol.py +++ b/cc_remote/protocol.py @@ -28,7 +28,7 @@ MAX_SINGLE_ATTACHMENT_BYTES, ) -PROTOCOL_VERSION = 73 +PROTOCOL_VERSION = 74 # Codex Desktop renders a 53-week daily token-activity calendar. Keep the wire # payload to that same bounded window so an account response can never turn a @@ -889,6 +889,7 @@ class UserMsg(_Base): # the later live echo. client_msg_id: Optional[WireId] = None timed_task: Optional[TimedMessage] = None + source_thread_id: Optional[WireId] = None prompt: str images: Optional[list[QueryImage]] = Field(default=None, max_length=MAX_ATTACHMENT_COUNT) # Metadata only: file bodies stay out of replay/cache, while names remain @@ -901,6 +902,7 @@ class TurnSteered(_Base): type: Literal["turn_steered"] = "turn_steered" msg_id: WireId turn_id: WireId + source_thread_id: Optional[WireId] = None prompt: str images: Optional[list[QueryImage]] = Field( default=None, max_length=MAX_ATTACHMENT_COUNT) @@ -2474,12 +2476,21 @@ class GetHistory(_Command): detail: Literal["summary", "full"] = "full" +class SessionMessageReceipt(BaseModel): + model_config = ConfigDict(extra="forbid") + itemId: WireId + threadId: WireId + status: Literal["sending", "sent", "failed"] + + class ConversationTurn(BaseModel): """Canonical lightweight turn rendered without replaying raw events.""" model_config = ConfigDict(extra="forbid") id: WireId clientMsgId: Optional[WireId] = None timedTask: Optional[TimedMessage] = None + sourceThreadId: Optional[WireId] = None + sessionMessages: Optional[list[SessionMessageReceipt]] = Field(default=None, max_length=16) prompt: str = Field(default="", max_length=128 * 1024) blocks: list[dict[str, Any]] = Field(default_factory=list, max_length=32) done: bool = False diff --git a/cc_remote/wrapper/codex_daemon.py b/cc_remote/wrapper/codex_daemon.py index 12dae6a3..32034f06 100644 --- a/cc_remote/wrapper/codex_daemon.py +++ b/cc_remote/wrapper/codex_daemon.py @@ -102,7 +102,13 @@ class _CommandResult: def _run_command( argv: tuple[str, ...], env: Mapping[str, str], timeout: float, ) -> _CommandResult: - """Blocking subprocess boundary, kept separate for deterministic tests.""" + """Run native lifecycle commands outside the disposable Wrapper release.""" + # Native start/restart inherits this cwd into the durable app-server. A + # release can be retired while that server survives several Wrapper upgrades. + # Never fall back to the caller's cwd, including for read-only lifecycle probes. + home = env.get("HOME") or str(Path.home()) + if not os.path.isabs(home) or not os.path.isdir(home): + return _CommandResult(127, b"", b"InvalidDaemonWorkingDirectory") command_argv = argv if os.name == "posix" and any( argv[index:index + 2] == ("app-server", "daemon") @@ -126,6 +132,7 @@ def _run_command( stdout=subprocess.PIPE, stderr=subprocess.PIPE, env=dict(env), + cwd=home, timeout=timeout, check=False, ) diff --git a/cc_remote/wrapper/codex_delegation.py b/cc_remote/wrapper/codex_delegation.py new file mode 100644 index 00000000..652f822b --- /dev/null +++ b/cc_remote/wrapper/codex_delegation.py @@ -0,0 +1,68 @@ +"""Codex App cross-thread input envelopes (not sub-agent collaboration). + +Only decode the native envelope. A source id is display/navigation metadata, +never authority to change accounts, machines or execute an action. +""" +from __future__ import annotations + +import re +from dataclasses import dataclass + +_NATIVE_ID = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$") +_ENVELOPE = re.compile( + r"\s*([^<>]+)" + r"\s*([^<>]*)\s*", + re.DOTALL, +) +_TOOLS = frozenset({"create_thread", "send_message_to_thread", "handoff_thread"}) + + +@dataclass(frozen=True) +class CodexDelegation: + source_thread_id: str + prompt: str + + +def parse_codex_delegation(text: object) -> CodexDelegation | None: + if not isinstance(text, str) or len(text) > 1024 * 1024: + return None + match = _ENVELOPE.fullmatch(text.strip()) + if match is None or not _NATIVE_ID.fullmatch(match[1].strip()): + return None + # Match the official App's escaping exactly; do not parse XML/entities. + prompt = match[2].strip().replace("<", "<").replace(">", ">").replace("&", "&") + if not prompt: + return None + return CodexDelegation(match[1].strip(), prompt) + + +def is_codex_delegation_output(item: object) -> bool: + return (isinstance(item, dict) + and str(item.get("type") or "").lower() == "functioncalloutput" + and item.get("namespace") == "codex_app" + and isinstance(item.get("name"), str) + and item.get("name") in _TOOLS) + + +def normalize_codex_delegation_item(item: dict) -> dict: + """Project native turnToolOutput as the same user item used by the App. + + Other function outputs must stay tools. Keep the exact item id for live / + persisted-history reconciliation; never infer identity from equal text. + """ + if (not is_codex_delegation_output(item) + or parse_codex_delegation(item.get("output")) is None): + return item + return {"type": "userMessage", "id": item.get("id"), "clientId": None, + "content": [{"type": "text", "text": item["output"]}]} + + +def codex_message_target(tool: object, arguments: object, server: object = None) -> str | None: + if not isinstance(tool, str) or not isinstance(arguments, dict): + return None + native_tool = tool in {"send_message_to_thread", "codex_app.send_message_to_thread"} + native_server = server == "codex_app" or arguments.get("namespace") == "codex_app" + if not ((native_tool and native_server) or tool == "mcp__codex_app__send_message_to_thread"): + return None + target = arguments.get("threadId") + return target if isinstance(target, str) and _NATIVE_ID.fullmatch(target) else None diff --git a/cc_remote/wrapper/codex_external.py b/cc_remote/wrapper/codex_external.py index a9ba4a62..bb2e66ae 100644 --- a/cc_remote/wrapper/codex_external.py +++ b/cc_remote/wrapper/codex_external.py @@ -20,6 +20,10 @@ from pathlib import Path from typing import Callable, Iterable, Mapping +from cc_remote.wrapper.codex_delegation import ( + normalize_codex_delegation_item, parse_codex_delegation, +) + from cc_remote.protocol import ( MAX_SAFE_WIRE_INTEGER, MAX_SAFE_WIRE_TIMESTAMP_SECONDS, @@ -98,6 +102,7 @@ class CodexRolloutUserMessage: raw_text: str prompt: str | None + source_thread_id: str | None = None message_id: str | None = None client_id: str | None = None turn_id: str | None = None @@ -1043,6 +1048,9 @@ def visible_codex_user_message(message: object) -> str | None: text = message.strip() if not text: return None + delegation = parse_codex_delegation(text) + if delegation is not None: + return delegation.prompt marker = text.rfind(_CODEX_REQUEST_MARKER) if marker >= 0: request = text[marker + len(_CODEX_REQUEST_MARKER):].strip() @@ -1089,6 +1097,8 @@ def codex_rollout_user_message( raw_text = payload.get("message") elif payload_type == "item_completed": item = payload.get("item") + if isinstance(item, dict): + item = normalize_codex_delegation_item(item) if ( not isinstance(item, dict) or str(item.get("type") or "").lower() != "usermessage" @@ -1101,7 +1111,9 @@ def codex_rollout_user_message( return None if not isinstance(raw_text, str): return None + delegation = parse_codex_delegation(raw_text) return CodexRolloutUserMessage( + source_thread_id=delegation.source_thread_id if delegation else None, raw_text=raw_text, prompt=visible_codex_user_message(raw_text), message_id=message_id if isinstance(message_id, str) else None, diff --git a/cc_remote/wrapper/codex_history.py b/cc_remote/wrapper/codex_history.py index b644a8d5..8327b557 100644 --- a/cc_remote/wrapper/codex_history.py +++ b/cc_remote/wrapper/codex_history.py @@ -16,6 +16,10 @@ from dataclasses import dataclass, field, replace from typing import Any, Awaitable, Callable +from cc_remote.wrapper.codex_delegation import ( + normalize_codex_delegation_item, parse_codex_delegation, +) + from cc_remote.protocol import TurnEnd, TurnResult, UserMsg from cc_remote.wrapper.claude_compaction import compact_completion_events from cc_remote.wrapper.codex_history_prefetch import CodexHistoryPrefetch @@ -224,6 +228,7 @@ def _validated_turn( "invalid Codex user content") normalized = dict(value) + normalized["items"] = [normalize_codex_delegation_item(item) for item in items] normalized["startedAt"] = _optional_nonnegative_int( value.get("startedAt"), "startedAt") normalized["completedAt"] = _optional_nonnegative_int( @@ -265,11 +270,14 @@ def _user_message(item: dict[str, Any], *, ts: float | None) -> UserMsg: "invalid Codex image data") images.append({"media_type": media_type, "data": data}) + prompt = "".join(prompt_parts) + delegation = parse_codex_delegation(prompt) kwargs: dict[str, Any] = { + "source_thread_id": delegation.source_thread_id if delegation else item.get("_ccRemoteSourceThreadId"), "msg_id": _wire_id(item.get("id"), "user"), "client_msg_id": _optional_wire_id( item.get("clientId"), "client-message"), - "prompt": "".join(prompt_parts), + "prompt": delegation.prompt if delegation else prompt, "images": images or None, } if ts is not None: @@ -1017,6 +1025,7 @@ async def summary_page( "type": "text", "text": recovered.prompt, }], + "_ccRemoteSourceThreadId": recovered.source_thread_id, "_ccRemoteImages": [ dict(image) for image in recovered.images or [] ], diff --git a/cc_remote/wrapper/codex_readiness.py b/cc_remote/wrapper/codex_readiness.py index 979b3594..3f89b4bd 100644 --- a/cc_remote/wrapper/codex_readiness.py +++ b/cc_remote/wrapper/codex_readiness.py @@ -18,6 +18,7 @@ from typing import Any from websockets.client import ClientProtocol +from websockets.asyncio.client import unix_connect from websockets.frames import Frame, Opcode from websockets.http11 import Response from websockets.uri import parse_uri @@ -31,6 +32,43 @@ _TIMEOUT = 8.0 +async def probe_config(socket_path: str) -> None: + """Read live config without starting/resuming a thread or an account CLI. + + A listening daemon can still have a deleted working directory. Initialization + alone doesn't exercise the config loader used to start the next turn. Never + expose the config or raw native errors in this acceptance check. + """ + before = socket_identity(socket_path, owner_uid=os.geteuid()) + async with asyncio.timeout(_TIMEOUT): + async with unix_connect( + socket_path, uri="ws://localhost/", proxy=None, + open_timeout=_TIMEOUT, close_timeout=1, max_size=2 * 1024 * 1024, + ) as connection: + async def request(request_id: int, method: str, params: dict) -> dict: + await connection.send(json.dumps({ + "id": request_id, "method": method, "params": params, + })) + while True: + response = json.loads(await connection.recv()) + if not isinstance(response, dict) or response.get("id") != request_id: + continue + if "error" in response or not isinstance(response.get("result"), dict): + raise RuntimeError("Codex live configuration check failed") + return response["result"] + + await request(1, "initialize", { + "clientInfo": {"name": "cc-remote-readiness", "version": __version__}, + }) + await connection.send('{"method":"initialized"}') + # No cwd override: validate the daemon's own configuration base. + result = await request(2, "config/read", {"includeLayers": False}) + if not isinstance(result.get("config"), dict): + raise RuntimeError("Codex live configuration response is invalid") + if socket_identity(socket_path, owner_uid=os.geteuid()) != before: + raise RuntimeError("Codex daemon changed during configuration check") + + async def probe_proxy(binary: str, env: dict[str, str], socket_path: str) -> None: """Initialize through the official raw WebSocket proxy, without a thread.""" process = await asyncio.create_subprocess_exec( diff --git a/cc_remote/wrapper/codex_sessions.py b/cc_remote/wrapper/codex_sessions.py index 81122c51..29815b14 100644 --- a/cc_remote/wrapper/codex_sessions.py +++ b/cc_remote/wrapper/codex_sessions.py @@ -955,8 +955,8 @@ def codex_session_settings( which 0.144.1 does not include in that response. Config.toml is never a valid resume source because it holds only fresh-thread global defaults. - Returns {} when the rollout is missing/unreadable; the caller falls back to the - config defaults (correct for a brand-new session). + Returns {} when the bounded tail has no readable settings. This is read + uncertainty, not evidence that a resumed thread should use global defaults. """ path = ( _rollout_path(session_id) @@ -977,23 +977,33 @@ def codex_session_settings( tail_bytes = max(1, int(max_bytes)) start = max(0, size - tail_bytes) with open(path, "rb") as f: + def read_chunk(): + remaining = size - f.tell() + return f.readline(min(MAX_JSONL_RECORD_BYTES + 1, remaining)) \ + if remaining > 0 else b"" + if start: f.seek(start - 1) starts_at_record = f.read(1) == b"\n" f.seek(start) if not starts_at_record: - discarded = f.readline(MAX_JSONL_RECORD_BYTES + 1) - if not discarded.endswith(b"\n"): + # The tail may begin inside a large image/tool record. + # Recover its newline in bounded chunks inside this snapshot + # instead of losing all newer settings behind that record. + discarded = read_chunk() + while discarded and not discarded.endswith(b"\n"): + discarded = read_chunk() + if not discarded: return {} while True: - raw = f.readline(MAX_JSONL_RECORD_BYTES + 1) + raw = read_chunk() if not raw: break if len(raw) > MAX_JSONL_RECORD_BYTES: - if not raw.endswith(b"\n"): - # The remainder is still the same oversized record. Stop: - # a boundary cannot be recovered without exceeding our cap. - break + # Skip one oversized record, retaining both the per-record + # allocation cap and the total tail scan budget. + while raw and not raw.endswith(b"\n"): + raw = read_chunk() continue try: line = raw.decode("utf-8") diff --git a/cc_remote/wrapper/codex_stream.py b/cc_remote/wrapper/codex_stream.py index ce638d11..3df6b1e7 100644 --- a/cc_remote/wrapper/codex_stream.py +++ b/cc_remote/wrapper/codex_stream.py @@ -21,6 +21,10 @@ from pydantic import ValidationError +from cc_remote.wrapper.codex_delegation import ( + normalize_codex_delegation_item, parse_codex_delegation, +) + from cc_remote.attachments import ( ALLOWED_IMAGE_TYPES, MAX_IMAGE_DIMENSION, @@ -136,6 +140,7 @@ class CodexLiveUserMessage: turn_id: str prompt: str client_id: str | None = None + source_thread_id: str | None = None def codex_live_user_message(message: object) -> CodexLiveUserMessage | None: @@ -146,6 +151,8 @@ def codex_live_user_message(message: object) -> CodexLiveUserMessage | None: return None params = message.get("params") item = params.get("item") if isinstance(params, dict) else None + if isinstance(item, dict): + item = normalize_codex_delegation_item(item) if not isinstance(item, dict) or item.get("type") != "userMessage": return None message_id = item.get("id") @@ -163,7 +170,9 @@ def codex_live_user_message(message: object) -> CodexLiveUserMessage | None: client_id = item.get("clientId") if not isinstance(client_id, str) or not _SAFE_WIRE_ID.fullmatch(client_id): client_id = None + delegation = parse_codex_delegation(codex_user_item_text(item)) return CodexLiveUserMessage( + source_thread_id=delegation.source_thread_id if delegation else None, message_id=message_id, turn_id=turn_id, prompt=prompt, @@ -675,6 +684,8 @@ def _history_user_cursors( or ( b'"user_message"' not in line and b'"UserMessage"' not in line + and b'"FunctionCallOutput"' not in line + and b'"functionCallOutput"' not in line )): return None try: @@ -1631,7 +1642,7 @@ def codex_history_process_append( if (any(marker in line for marker in ( b'"task_started"', b'"user_message"', b'"session_meta"', b'"thread_goal_updated"', b'"thread_goal_cleared"', - )) or re.search(rb'"usermessage"|"role"\s*:\s*"user"', line, re.I)): + )) or re.search(rb'"usermessage"|"functioncalloutput"|"role"\s*:\s*"user"', line, re.I)): return None if _history_generated_image_record(line): process.observe(None) @@ -2029,6 +2040,7 @@ def codex_history_boundary_user( msg_id=message_id, client_msg_id=client_id, prompt=user.prompt, + source_thread_id=user.source_thread_id, ts=0, ) if pending_images: @@ -4810,6 +4822,7 @@ def close_turn( um = UserMsg( msg_id=uid, client_msg_id=user_client_id, + source_thread_id=user_record.source_thread_id, prompt=msg, ) if pending_images: diff --git a/cc_remote/wrapper/history_store.py b/cc_remote/wrapper/history_store.py index 0ede588c..669515b6 100644 --- a/cc_remote/wrapper/history_store.py +++ b/cc_remote/wrapper/history_store.py @@ -18,6 +18,8 @@ from pathlib import Path from typing import Any, Callable +from cc_remote.wrapper.codex_delegation import codex_message_target + from cc_remote.attachments import ( ALLOWED_IMAGE_TYPES, MAX_IMAGE_DIMENSION, @@ -72,7 +74,8 @@ # v45 restores public AgentMessage records in new Codex rollouts and the # native owner of source-window tails. Old tools-only projections must rebuild. # v46 restores native commands, source clocks and closed segment envelopes. -_SCHEMA_VERSION = 46 +# v47 preserves native cross-thread provenance and outgoing message receipts. +_SCHEMA_VERSION = 47 _FINGERPRINT_SAMPLE_BYTES = 64 * 1024 _DEFAULT_MAX_ENTRIES = 128 _DEFAULT_MAX_BYTES = 64 * 1024 * 1024 @@ -563,6 +566,8 @@ def materialize_history_turns( prompt = "" has_user = False client_msg_id = None + source_thread_id = None + session_messages: dict[str, dict] = {} prompt_truncated = False image_refs: list[dict[str, Any]] = [] deferred_image_count = 0 @@ -697,6 +702,7 @@ def touch_process( started_ms = _event_ms(event.get("ts")) if event_type == "user_msg": has_user = True + source_thread_id = event.get("source_thread_id") if isinstance(event.get("client_msg_id"), str): client_msg_id = event["client_msg_id"] if isinstance(event.get("prompt"), str): @@ -798,6 +804,18 @@ def touch_process( elif event_type == "error": if isinstance(event.get("message"), str): error = _historical_turn_failure(event["message"]) + if event_type == "tool_use" and len(session_messages) < 16: + target = codex_message_target(event.get("tool"), event.get("input"), event.get("server")) + item_id = event.get("tool_use_id") + if target and isinstance(item_id, str): + session_messages.setdefault(item_id, { + "itemId": item_id, "threadId": target, "status": "sending", + }) + elif event_type == "tool_result" and event.get("tool_use_id") in session_messages: + session_messages[event["tool_use_id"]]["status"] = ( + "failed" if event.get("is_error") or event.get("status") in { + "failed", "cancelled", "declined", "interrupted", + } else "sent") if include_live_detail and event_type == "tool_use": tool_id = event.get("tool_use_id") message_id = event.get("message_id") @@ -1217,6 +1235,8 @@ def touch_process( } optional = { "clientMsgId": client_msg_id, + "sourceThreadId": source_thread_id, + "sessionMessages": list(session_messages.values()) or None, "forkPointId": fork_point, "checkpointId": checkpoint_id, "imageRefs": image_refs or None, @@ -1438,6 +1458,10 @@ def _ensure_schema(self) -> None: for table in ("history_pages", "history_turn_details"): connection.execute( f"DELETE FROM {table} WHERE engine='codex'") + if 0 < current < 47: + # v47 projects native Codex cross-thread input envelopes. + for table in ("history_pages", "history_turn_details"): + connection.execute(f"DELETE FROM {table} WHERE engine='codex'") if current in (10, 11, 12, 13, 14, 15, 16): # v16 makes browser/native ownership durable; v17 reuses the # adjacent native response-item id for legacy Codex user rows. @@ -1468,7 +1492,7 @@ def _ensure_schema(self) -> None: for table in ("history_pages", "history_turn_details"): connection.execute( f"DELETE FROM {table} WHERE engine='codex'") - elif current in (21, 22, 23, 24, 25, 26, 27, 28, 29, 30, 31, 32, 33, 34, 35, 36, 37, 38, 39, 40, 41, 42, 43, 44, 45): + elif current in (21, 22, 23, 24, 25, 26, 27, 28, 29, 30, 31, 32, 33, 34, 35, 36, 37, 38, 39, 40, 41, 42, 43, 44, 45, 46): # The independent v22-v44 invalidations above suffice. pass elif current not in (0, _SCHEMA_VERSION): diff --git a/cc_remote/wrapper/machine.py b/cc_remote/wrapper/machine.py index f48670a7..8b8b0241 100644 --- a/cc_remote/wrapper/machine.py +++ b/cc_remote/wrapper/machine.py @@ -304,6 +304,7 @@ CodexProcessClockStoreError, ) from cc_remote.wrapper.codex_permissions import codex_permission_profiles +from cc_remote.wrapper.codex_delegation import is_codex_delegation_output from cc_remote.wrapper.codex_stream import ( CodexHistoryImageView, CodexHistoryNativeWitness, CodexHistoryProcessPageWitness, CodexLiveUserMessage, @@ -22402,6 +22403,7 @@ async def publish_live_user(raw: dict) -> bool: msg_id=user.message_id, client_msg_id=user.client_id, prompt=user.prompt, + source_thread_id=user.source_thread_id, )) await self._emit(ctx, TurnBinding( msg_id=user.message_id, @@ -22433,6 +22435,7 @@ async def publish_live_user(raw: dict) -> bool: msg_id=user.message_id, client_msg_id=user.client_id, prompt=user.prompt, + source_thread_id=user.source_thread_id, )) await self._emit(ctx, TurnBinding( msg_id=user.client_id, @@ -22454,6 +22457,7 @@ async def publish_live_user(raw: dict) -> bool: msg_id=user.message_id, turn_id=current_turn_id, prompt=user.prompt, + source_thread_id=user.source_thread_id, )) return True @@ -22462,7 +22466,8 @@ def raw_is_user_item(raw: dict) -> bool: return False params = raw.get("params") item = params.get("item") if isinstance(params, dict) else None - return isinstance(item, dict) and item.get("type") == "userMessage" + return (isinstance(item, dict) and item.get("type") == "userMessage" + or is_codex_delegation_output(item)) def raw_proves_automatic_output(raw: dict) -> bool: method = raw.get("method") @@ -26103,10 +26108,6 @@ async def read_profile(profile: ClaudeProfile): c.key for c in self.sessions.values() if c.key and c.session_id and c.engine == "claude" } - resident_state = { - c.key: c.state for c in self.sessions.values() - if c.key and c.session_id and c.engine == "claude" - } work_records = await asyncio.to_thread( self._work.for_engine("claude").records_by_profile_session) pinned_ids = (self._session_pins.ids("claude") @@ -26158,7 +26159,6 @@ async def read_profile(profile: ClaudeProfile): tag=("archived" if record and record.archived else (info.tag or "")[:128] or None), pinned=wire_sid in pinned_ids, - state=resident_state.get(wire_sid), engine="claude", space=space, work_id=record.work_id if record else None, native_session_id=info.session_id, @@ -26204,14 +26204,23 @@ async def read_profile(profile: ClaudeProfile): session_id=broker_sid, summary="Claude Remote", cwd=broker_cwd[:4096], - state=resident_state.get(broker_sid, "idle"), + state="idle", pinned=broker_sid in pinned_ids, engine="claude", space="code", **self._session_presentation_fields("claude", broker_sid), )) known.add(broker_sid) + # Catalog metadata reads above can yield across a turn boundary. + # Sample activity only after the final read, so this later list + # cannot reintroduce running after the already-emitted idle frame + # (or erase a new running frame). Keep profile-qualified identities. + resident_state = { + c.key: c.state for c in self.sessions.values() + if c.key and c.session_id and c.engine == "claude" + } for session in sessions: + session.state = resident_state.get(session.session_id, session.state) self._remember_notification_title( session.session_id, session.summary or session.first_prompt) event = SessionList( @@ -36081,7 +36090,7 @@ async def codex_profile_allowed(profile_id: str) -> bool: if mode in CODEX_COLLABORATION_MODES: sdk.collaboration_mode = mode try: - resolved_model, model_replaced = ( + resolved_model, _ = ( await self._resolve_codex_profile_model( codex_profile, model, @@ -36097,7 +36106,7 @@ async def codex_profile_allowed(profile_id: str) -> bool: ) return None model = resolved_model - if model and (model_replaced or explicit_codex_model): + if model and explicit_codex_model: codex_resume_model_reconcile = model if model: sdk.model = model @@ -36222,6 +36231,18 @@ async def codex_profile_allowed(profile_id: str) -> bool: await ctx.sdk.connect( **codex_connect_options, ) + if resume_id and not explicit_codex_model: + # A missing/stale rollout tail cannot retire a live choice. + # Validate the authoritative resume model before deciding + # whether this account's advertised default must replace it. + native_model = getattr(ctx.sdk, "model", None) + resolved_model, model_replaced = ( + await self._resolve_codex_profile_model( + codex_profile, native_model, + ) + ) + if native_model and model_replaced: + codex_resume_model_reconcile = resolved_model if ( resume_id and codex_resume_model_reconcile @@ -37925,7 +37946,8 @@ async def publish_codex_user(raw: dict) -> None: initial = bool(codex_initial_msg_id and codex_initial_msg_id in { user.message_id, user.client_id, }) - if not initial and not user.client_id and not codex_initial_user_seen: + if (not initial and not user.client_id and not user.source_thread_id + and not codex_initial_user_seen): # A missed initial echo cannot turn an unlabelled first item # into a second input. Never guess from equal prompt text. return @@ -37951,6 +37973,7 @@ async def publish_codex_user(raw: dict) -> None: msg_id=msg_id, turn_id=user.turn_id, prompt=user.prompt, + source_thread_id=user.source_thread_id, )) if user.client_id is not None: # Reconcile the canonical history id only after the boundary; @@ -37959,6 +37982,7 @@ async def publish_codex_user(raw: dict) -> None: msg_id=user.message_id, client_msg_id=user.client_id, prompt=user.prompt, + source_thread_id=user.source_thread_id, )) async def emit_codex_event(event) -> None: diff --git a/deploy/README.md b/deploy/README.md index fa6a697e..1b213266 100644 --- a/deploy/README.md +++ b/deploy/README.md @@ -24,6 +24,8 @@ Before changing a live service: 1. Inspect the source worktree, target installation, current release, service manager, and health. Preserve unrelated changes; do not normalize a dirty worktree or silently replace a custom installation layout. + Apply [backup retention](#deployment-backup-retention) before creating another + deployment backup or staging copy. 2. Select the matching supported path. Prefer a tested source snapshot for current features using the [source deployment guide](../docs/installation_en.md#source-install). Use `install.sh` when the operator selects a published release that includes @@ -69,15 +71,166 @@ Success requires all of the following: the expected immutable releases are active, Python and served Web build metadata report the same protocol/product, services have stable PIDs without restart loops, the public health endpoint is healthy, expected Wrappers reconnect, and recent logs contain no new fatal -errors. Installations using Codex Code must also verify the +errors. These checks must pass **after the final cleanup**, not just before it. +Installations using Codex Code must also verify the [shared CLI control plane](#codex-code-shared-control-plane-acceptance); an online Wrapper alone does not prove bidirectional CLI access. After these checks, offer the [optional Codex App attachment](#optional-codex-app-attachment) on eligible desktops. Its consent/availability is reported separately and never turns a healthy core deployment into a failure. On failure, use the installer-owned rollback or the retained previous -release and matching state snapshot; do not delete old releases during the -deployment. +release and matching state snapshot. Never prune the active transaction's +rollback set during activation or recovery; remove older unreferenced generations +beforehand and finalize retention after coordinated acceptance as specified below. + +### Deployment backup retention + +Agent-led deployments retain **the active installation, one complete previous +rollback generation, and every still-referenced runtime dependency** on each +in-scope host. Runtime dependency protection always overrides the generation +count; an older release is not disposable merely because two newer ones exist. +A generation includes the matching release code/runtime, configuration copies +and private state snapshot needed to restore it; these are one recovery set, +not separate allowances for multiple historical copies. An intact immutable +release can serve as the code backup without another archive of the same tree. + +1. Before creating the next backup or staging copy, identify the active release, + the newest complete known-good rollback set, and any unresolved deployment + transaction. Remove only confirmed older, superseded deployment backups and + unused duplicate uploads/archives using `deploy/cleanup.py`. Determine + generations from release and transaction records, not filename age alone. + Preview an explicit private inventory before applying it; do not generate a + separate retention script or bypass a deferred result with `rm`/`rmtree`. +2. Create and validate the new pre-upgrade snapshot using the normal transaction + procedure. Keep the existing valid rollback set until coordinated acceptance + succeeds. Temporary coexistence during this transaction must not become + permanent retention; never delete the sole usable backup to make room for + an unverified replacement. +3. After all protocol tiers pass acceptance and the transaction is committed, + retain only the version just superseded and its matching recovery files. + Remove the older rollback generation and completed, unneeded staging/upload + copies. Quarantine candidates, repeat acceptance while their bytes remain + recoverable, then delete and verify again. On failure or unknown outcome, + preserve the exact transaction's recovery set and settle it before further + cleanup or deployment attempts. Deployment is complete only after this final + acceptance; earlier health checks are provisional. + +Do not delete native transcripts, credentials, current private state, project +files or unrelated user backups under this policy. A release still referenced +by a running service (including the independent Claude service), a process cwd, +an executable, an open/mapped file, a shared venv, an active job or the retained +rollback set is a live dependency. Checking `ps ... args` alone misses processes +whose argv is independent of their cwd. Include dormant service definitions and +configuration-dependent runtime paths explicitly as protected dependencies too. +Incomplete process visibility or uncertain provenance defers cleanup. Do not +stop a daemon or active work just to meet the retention count. + +Codex lifecycle commands use the service user's stable home directory, never a +Wrapper release, as their working directory. Existing daemons are not restarted +to apply this: retain any old release they still use until their normal lifecycle +has moved them away. This startup rule and the cleanup checks protect different +boundaries; neither replaces the other. + +Report retained rollback paths, deleted generations, removed allocated bytes +and any deferred paths. Allocated bytes are not a measurement of free-space gain +(for example, shared filesystem blocks may remain in use). Moving old copies +into another backup directory or Trash +does not reclaim disk space or satisfy this policy. This is an agent-operated +cleanup step: the existing installers and `cc-remote update` do not automatically +prune releases or backups. + +#### Repository cleanup command + +`cleanup.py` ships in both role bundles. It acquires the same `.update.lock` as +the installers, verifies the bound `current` link and completed transaction +records, and examines process cwd/executable/open-file references with `lsof` +plus argv references with `ps`. It also traces transitive symlink dependencies of +retained trees and live candidates, including shared venvs. Broken links, +unreadable paths or an exceeded scan bound stop cleanup. `lsof` is required; +missing tools, warnings and incomplete scans stop the +operation. User-owned private installations inspect that user's processes; +shared/system installations must run as root to cover every service user. + +Create a mode-0600 JSON inventory **outside the source repository**, from the +actual installation and its transaction records. All paths must be absolute, +canonical, and owned/maintained within the operator's deployment scope. The tool +does not discover unknown external transactions, infer generation age, or select +files for you. List all relevant transaction records; if any outcome is unresolved, +settle it first and retain its recovery set. Never list credentials, native +sessions, live private state, project files or unrelated backups as candidates. + +Inventory schema (replace the illustrative paths): + +```json +{ + "schema": 1, + "installation_root": "/opt/cc-remote", + "current_release": "/opt/cc-remote/releases/release-new", + "rollback_paths": [ + "/opt/cc-remote/releases/release-previous", + "/opt/cc-remote/rollback-data/before-new" + ], + "protected_paths": [], + "cleanup_roots": ["/opt/cc-remote/releases"], + "candidates": ["/opt/cc-remote/releases/release-old"], + "transactions": [ + {"path": "/opt/cc-remote/transactions/activation.json", "field": "phase", "equals": "committed"} + ], + "checks": [ + {"name": "release and public health", "argv": ["/absolute/path/to/read-only-health-check"], "cwd": "/"} + ] +} +``` + +`rollback_paths` describes one complete recovery generation, including its +configuration/state snapshots; it must contain exactly one previous release. +Use `protected_paths` for additional service/runtime dependencies. Candidates +must be direct children of explicitly named, dedicated `cleanup_roots`; nested +candidates, symlink boundaries and mounts are rejected. Transaction expectations +accept only completed states (`committed`, `complete`, `deployed_verified`, +`ready`), and the records must remain unchanged throughout cleanup. Unknown +layouts remain retained until their provenance is established. + +`checks` are operator-reviewed **read-only** argv arrays (no shell interpolation), +run before retirement, after quarantine, and after removal. They must cover the +actual role's release identity, stable services and health. For a Codex Wrapper, +also include a fresh configuration check, executed **as the Wrapper service user**: + +```bash + deploy/check_codex_readiness.py \ + --home --release --after \ + --live-config --wait 0 +# Add the installation's existing --plist or --env-file selector when needed. +``` + +On a root-managed Linux cleanup, use `runuser`/`sudo -u` in that check's argv to +select the real service user. `--live-config` sends only initialize and +`config/read` to existing account sockets: it neither starts a daemon nor +creates/resumes a thread or sends a model message, and never prints configuration +contents. It catches a daemon that accepts connections but can no longer load +configuration. Verify an existing idle session's Wrapper status route separately +as part of final acceptance; do not send a model turn without authorization. + +```bash + deploy/cleanup.py /private/path/cleanup-inventory.json + deploy/cleanup.py /private/path/cleanup-inventory.json --apply +``` + +Preview does not run acceptance commands or remove artifacts. Apply rescans live +references immediately before each rename and permanent deletion. Candidates are +temporarily renamed beside their original path; a failed quarantine check restores +the intact directory. A new live reference defers deletion and restores that path. +Checks are observations, not a lock on arbitrary external programs: operators must +not launch jobs against retired paths during cleanup. + +The private `/.cleanup-transaction.json` records intent before +each mutation. Exit 0 means success (or preview), 2 means applied cleanup completed +with retained/deferred artifacts, and 1 means failure. A lost connection, partial +deletion, or failed final check requires inspection of this exact journal before +retrying. Do not remove the journal or rerun the command to hide an unknown result. +Quarantine is temporary transaction state, not another retained backup or Trash. + +### Deployment entrypoints - `install.sh` — versioned GitHub Release bootstrap. It requires an explicit `relay` or `wrapper` role, detects OS/CPU, downloads that one role archive, @@ -193,11 +346,11 @@ deployment. migration transaction, restores matching pre-release data before an older wrapper is restarted, and verifies both engines' Work ownership backfills. -Protocol v73 is a coordinated upgrade: publish freshly built Relay/Web and +Protocol v74 is a coordinated upgrade: publish freshly built Relay/Web and Wrapper artifacts from the same tagged commit. The strict protocol gate is intentional and mixed protocol versions will not communicate. `setup-vps.sh` rejects a missing or mismatched web build manifest. Stop the wrapper first; -activate the v73 relay/web release; then start the v73 wrapper. +activate the v74 relay/web release; then start the v74 wrapper. The wrapper installer treats local Work data and versioned private control state as part of the release @@ -210,8 +363,8 @@ the previous code. If data restoration fails, it leaves the wrapper stopped instead of running old code against a new schema. A manual or legacy-layout deployment must use the same order: stop the wrapper, run `work_registry_snapshot.py snapshot` from the new staging tree, activate and -verify v73, and retain that snapshot with the previous release. To roll back, -stop v73, run `work_registry_snapshot.py restore`, then switch and start the old +verify v74, and retain that snapshot with the previous release. To roll back, +stop v74, run `work_registry_snapshot.py restore`, then switch and start the old release. Never copy only `registry.sqlite3` while the wrapper is live because committed state may still be in its WAL file. Restoring a pre-release snapshot also restores pre-release Work metadata: sessions, projects, or schedule state diff --git a/deploy/build_release.py b/deploy/build_release.py index d70c0f24..0f5cecce 100755 --- a/deploy/build_release.py +++ b/deploy/build_release.py @@ -35,6 +35,7 @@ class BuildError(ValueError): "Caddyfile.insecure", "Caddyfile.viewer.example", "caddy_managed_block.py", + "cleanup.py", "cc-remote-relay.service", "env.relay.example", "install-relay.sh", @@ -49,6 +50,7 @@ class BuildError(ValueError): _WRAPPER_DEPLOY = ( "atomic_symlink.py", "check_codex_readiness.py", + "cleanup.py", "cc-remote-wrapper.service", "com.muggle.cc-remote.wrapper.plist.in", "env.wrapper.example", diff --git a/deploy/check_codex_readiness.py b/deploy/check_codex_readiness.py index 97343352..e5b78c40 100644 --- a/deploy/check_codex_readiness.py +++ b/deploy/check_codex_readiness.py @@ -2,6 +2,7 @@ from __future__ import annotations import argparse +import asyncio import json import os from pathlib import Path @@ -12,7 +13,7 @@ sys.path.insert(0, str(Path(__file__).resolve().parents[1])) -from cc_remote.wrapper.codex_readiness import REPORT_NAME +from cc_remote.wrapper.codex_readiness import REPORT_NAME, probe_config from cc_remote.wrapper.codex_daemon import socket_identity from cc_remote.wrapper.process_scan import ProcessIdentity, process_identity, process_owner_uid from deploy.work_registry_snapshot import resolve_wrapper_state_dir @@ -25,9 +26,31 @@ "version_mismatch": "Codex CLI 与运行中的服务版本不一致;当前任务保留,结束后再更新或重开 Codex。", "daemon_changed": "检查期间 Codex 服务发生变化,本次未确认连接。", "connection_failed": "连接检查未通过;现有任务保留,请检查 Wrapper 日志。", + "configuration_failed": "实时配置读取失败;连接就绪不足以确认可用,请检查服务工作目录及配置。", + "wrong_probe_user": "实时配置检查须以 Wrapper 服务用户执行,不可代用 root 或其他账号。", } +async def refresh_configuration(report: dict) -> None: + """Supplement the activation receipt with fresh, non-mutating native RPCs.""" + owner = process_owner_uid(report["wrapper"]["pid"]) + for row in report["profiles"]: + if row["status"] != "ready": + continue + if owner != os.geteuid(): + row.update(status="unavailable", reason="wrong_probe_user") + continue + try: + expected = str(Path(row["home"]) / "app-server-control/app-server-control.sock") + if row["socket"] != expected: + raise ValueError("account endpoint mismatch") + await probe_config(expected) + row["live_config_verified"] = True + except Exception as exc: + row.update(status="unavailable", reason="configuration_failed", + error_type=type(exc).__name__) + + def socket_still_ready(row: dict, owner_uid: int) -> bool: try: path = row["socket"] @@ -87,7 +110,9 @@ def describe(report: dict) -> bool: # Values come from private configuration but must not inject terminal controls. profile = json.dumps(row["profile"], ensure_ascii=False) if row["status"] == "ready": - print(f"Codex {profile}: 共享连接已就绪(CLI 与 cc-remote 的连接检查通过)。") + detail = ("共享连接与实时配置读取检查通过。" if row.get("live_config_verified") + else "共享连接已就绪(CLI 与 cc-remote 的连接检查通过)。") + print(f"Codex {profile}: {detail}") elif row["status"] == "disabled": print(f"Codex {profile}: 保留已有的关闭设置,未启用共享连接。") else: @@ -107,6 +132,8 @@ def main(argv: list[str] | None = None) -> int: parser.add_argument("--env-file", type=Path) parser.add_argument("--plist", type=Path) parser.add_argument("--wait", type=float, default=45) + parser.add_argument("--live-config", action="store_true", + help="also read config from each live daemon as its service user; no model turn") args = parser.parse_args(argv) try: state = resolve_wrapper_state_dir(args.home, env_file=args.env_file, plist=args.plist) @@ -114,6 +141,8 @@ def main(argv: list[str] | None = None) -> int: while True: report = read_receipt(state / REPORT_NAME, args.release, args.after) if report is not None: + if args.live_config: + asyncio.run(refresh_configuration(report)) return 0 if describe(report) else 1 if time.monotonic() >= deadline: print("Codex: 未收到本次启动的连接检查结果,请检查 Wrapper 日志;不能据此确认已共享。") diff --git a/deploy/cleanup.py b/deploy/cleanup.py new file mode 100644 index 00000000..b68f0d1a --- /dev/null +++ b/deploy/cleanup.py @@ -0,0 +1,478 @@ +"""Preview or retire explicitly inventoried deployment artifacts, never by age. + +The private inventory supplies transaction provenance, the complete rollback +set, service dependencies and read-only acceptance commands. Live references +are checked independently. See deploy/README.md#deployment-backup-retention. +""" +from __future__ import annotations + +import argparse +import hashlib +import json +import os +from pathlib import Path +import shutil +import stat +import subprocess +import sys +import time +import uuid + +sys.path.insert(0, str(Path(__file__).resolve().parents[1])) +from deploy.install_lock import acquire_install_lock + +JOURNAL = ".cleanup-transaction.json" +MAX_JSON = 1024 * 1024 +MAX_SCAN = 32 * 1024 * 1024 +MAX_DEPENDENCY_ENTRIES = 250_000 +COMPLETE = {"committed", "complete", "deployed_verified", "ready"} + + +class CleanupError(ValueError): + pass + + +def read_json(path: Path) -> tuple[dict, str]: + fd = os.open(path, os.O_RDONLY | os.O_NOFOLLOW | os.O_NONBLOCK) + with os.fdopen(fd, "rb") as stream: + info = os.fstat(stream.fileno()) + if (not stat.S_ISREG(info.st_mode) or info.st_size > MAX_JSON + or info.st_uid not in {0, os.geteuid()} or info.st_mode & 0o022): + raise CleanupError("unsafe or oversized inventory/transaction record") + raw = stream.read(MAX_JSON + 1) + if len(raw) > MAX_JSON: + raise CleanupError("oversized inventory/transaction record") + data = json.loads(raw) + if not isinstance(data, dict): + raise CleanupError("inventory/transaction record must be an object") + return data, hashlib.sha256(raw).hexdigest() + + +def absolute_path(value: str) -> Path: + if not isinstance(value, str) or not value or "\0" in value: + raise CleanupError("expected an absolute path") + path = Path(value) + if not path.is_absolute() or path != Path(os.path.normpath(value)): + raise CleanupError("paths must be absolute and normalized") + # Reject symlinks in ancestors too. A symlink *inside* an artifact is removed + # as a link by fd-based rmtree; it is never recursively followed. + if path.resolve() != path: + raise CleanupError(f"symlink/alias is not a cleanup boundary: {path}") + return path + + +def identity(path: Path) -> tuple[int, int]: + info = path.lstat() + if not (stat.S_ISDIR(info.st_mode) or stat.S_ISREG(info.st_mode)): + raise CleanupError(f"not a regular artifact: {path}") + return info.st_dev, info.st_ino + + +def overlaps(first: Path, second: Path) -> bool: + return first == second or first.is_relative_to(second) or second.is_relative_to(first) + + +def load_inventory(path: Path) -> dict: + data, digest = read_json(path) + required = {"schema", "installation_root", "current_release", "rollback_paths", + "cleanup_roots", "candidates", "transactions", "checks"} + if (data.keys() - required - {"protected_paths"} or required - data.keys() + or data["schema"] != 1): + raise CleanupError("invalid cleanup inventory schema") + for key in ["rollback_paths", "cleanup_roots", "candidates", "transactions", "checks"]: + if not isinstance(data[key], list) or not 0 < len(data[key]) <= 128: + raise CleanupError(f"{key} must be a nonempty bounded list") + root = absolute_path(data["installation_root"]) + if os.geteuid() != 0 and root.stat().st_uid != os.geteuid(): + raise CleanupError("shared installations require root process visibility") + current = absolute_path(data["current_release"]) + if current.parent != root / "releases": + raise CleanupError("current release must be a direct child of installation releases") + rollback = [absolute_path(p) for p in data["rollback_paths"]] + if sum(p.parent == root / "releases" for p in rollback) != 1 or current in rollback: + raise CleanupError("inventory must retain exactly one previous release and its rollback files") + protected = data.get("protected_paths", []) + if not isinstance(protected, list) or len(protected) > 128: + raise CleanupError("invalid protected paths") + protected = [current, *rollback, *(absolute_path(p) for p in protected)] + for p in protected: + identity(p) + roots = [absolute_path(p) for p in data["cleanup_roots"]] + for p in roots: + if (not p.is_dir() or p in {Path("/"), Path.home(), root} + or p.name in {".codex", ".claude", ".cc-remote"}): + raise CleanupError("cleanup roots must be dedicated artifact directories") + if os.geteuid() != 0 and p.stat().st_uid != os.geteuid(): + raise CleanupError("shared artifact roots require root process visibility") + candidates = [absolute_path(p) for p in data["candidates"]] + if len(set(candidates)) != len(candidates): + raise CleanupError("duplicate cleanup candidate") + for p in candidates: + if p.parent not in roots or any(p != q and overlaps(p, q) for q in candidates): + raise CleanupError("candidates must be non-overlapping direct children of cleanup roots") + for check in data["checks"]: + if (not isinstance(check, dict) or set(check) != {"name", "argv", "cwd"} + or not isinstance(check["name"], str) or not check["name"] + or not isinstance(check["argv"], list) or not check["argv"] + or not all(isinstance(s, str) and s and "\0" not in s for s in check["argv"]) + or not os.path.isabs(check["argv"][0])): + raise CleanupError("checks require a name, absolute executable argv and stable cwd") + cwd = absolute_path(check["cwd"]) + if not cwd.is_dir() or any(cwd == p or cwd.is_relative_to(p) for p in candidates): + raise CleanupError("acceptance check cwd must be outside cleanup candidates") + transactions = [] + for entry in data["transactions"]: + if (not isinstance(entry, dict) or set(entry) != {"path", "field", "equals"} + or entry["field"] not in {"phase", "status"} or entry["equals"] not in COMPLETE): + raise CleanupError("transaction requires a known completed phase/status") + record = absolute_path(entry["path"]) + value, checksum = read_json(record) + if value.get(entry["field"]) != entry["equals"]: + raise CleanupError(f"unfinished or unknown transaction: {record}") + transactions.append((record, checksum)) + protected.append(record) + protected.append(path.resolve()) + protected_identities = {p: identity(p) for p in protected} + protected.extend([root / JOURNAL, root / ".update.lock"]) + return dict(data, root=root, current=current, protected=protected, candidates=candidates, + transactions=transactions, digest=digest, inventory=path.resolve(), + protected_identities=protected_identities) + + +def check_bindings(plan: dict) -> None: + root = plan["root"] + if not (root / "current").is_symlink() or (root / "current").resolve(strict=True) != plan["current"]: + raise CleanupError("active release changed or current is not a managed symlink") + if read_json(plan["inventory"])[1] != plan["digest"]: + raise CleanupError("inventory changed during cleanup") + for path, checksum in plan["transactions"]: + if read_json(path)[1] != checksum: + raise CleanupError("transaction changed during cleanup") + for path, expected in plan["protected_identities"].items(): + if identity(absolute_path(str(path))) != expected: + raise CleanupError(f"protected rollback/dependency replaced: {path}") + + +def scan_command(argv: list[str]) -> bytes: + # Never print process argv or lsof output: either can contain private paths + # or arguments. Only matched PID/reference categories enter the report. + result = subprocess.run(argv, capture_output=True, timeout=30, check=False) + if result.returncode or result.stderr.strip() or len(result.stdout) > MAX_SCAN: + raise CleanupError("process reference scan incomplete; no deletion is safe") + return result.stdout + + +def parse_lsof(raw: bytes) -> list[tuple[int, str, Path]]: + pid = None + descriptor = None + references = [] + for field in raw.split(b"\0"): + field = field.lstrip(b"\n") + if not field: + continue + tag, value = field[:1], field[1:] + if tag == b"p": + pid = int(value) + descriptor = None + elif tag == b"f": + descriptor = os.fsdecode(value) + if descriptor in {"NOFD", "err"}: + raise CleanupError("process files are not fully observable") + elif tag == b"n": + if pid is None or descriptor is None: + raise CleanupError("incomplete lsof process/file record") + name = os.fsdecode(value) + if "Permission denied" in name or "Operation not permitted" in name: + raise CleanupError("process files are not fully observable") + if name.startswith("/"): + references.append((pid, descriptor, Path(name.removesuffix(" (deleted)")))) + if pid is None: + raise CleanupError("empty process reference inventory") + return references + + +def process_references(candidates: list[Path]) -> dict[Path, list[str]]: + lsof = shutil.which("lsof") + if not lsof: + raise CleanupError("lsof is required; process arguments alone are insufficient") + # User installations inspect their owner's processes; shared/system installs + # must run as root to cover all service users. Never silently ignore stderr. + argv = [lsof, "-nP", "-F0pfn"] + if os.geteuid() != 0: + argv.extend(["-a", "-u", str(os.geteuid())]) + found = {p: [] for p in candidates} + for pid, descriptor, path in parse_lsof(scan_command(argv)): + for candidate in candidates: + if path == candidate or path.is_relative_to(candidate): + found[candidate].append(f"pid {pid}: {descriptor}") + ps = scan_command(["ps", "-axo", "pid=,uid=,args="]).decode(errors="surrogateescape") + for line in ps.splitlines(): + parts = line.strip().split(None, 2) + if len(parts) < 2: + raise CleanupError("incomplete process argument inventory") + try: + pid, uid = map(int, parts[:2]) + except ValueError: + raise CleanupError("incomplete process argument inventory") from None + if os.geteuid() != 0 and uid != os.geteuid(): + continue + arguments = parts[2] if len(parts) == 3 else "" + for candidate in candidates: + if str(candidate) in arguments: + found[candidate].append(f"pid {pid}: argv") + return {p: sorted(set(reasons)) for p, reasons in found.items()} + + +def artifact_entries(root: Path): + """Do not follow links or cross mounts while scanning an artifact tree.""" + device = root.lstat().st_dev + pending = [root] + while pending: + path = pending.pop() + info = path.lstat() + if info.st_dev != device or os.path.ismount(path): + raise CleanupError(f"mounted artifact must be retained: {path}") + yield path, info + if stat.S_ISDIR(info.st_mode): + with os.scandir(path) as entries: + pending.extend(Path(entry.path) for entry in entries) + + +def symlink_dependencies(roots: list[Path], *, candidates: list[Path] | None = None) -> set[Path]: + """Read a bounded dependency closure, including links through shared venvs. + + This traversal only reads. Deletion uses artifact_entries, which never + follows links. Cycles and shared directories are visited just once. + """ + pending = list(roots) + queued_artifacts = set(roots) + visited = set() + targets = set() + while pending: + root = pending.pop() + device = root.lstat().st_dev + entries = [root] + while entries: + path = entries.pop() + info = path.lstat() + # Hard-linked relative symlinks can resolve differently by parent. + key = (info.st_dev, info.st_ino, path if stat.S_ISLNK(info.st_mode) else None) + if key in visited: + continue + visited.add(key) + if len(visited) > MAX_DEPENDENCY_ENTRIES: + raise CleanupError("dependency scan exceeds limit; retain artifacts") + if info.st_dev != device or os.path.ismount(path): + raise CleanupError("dependency scan crosses a mount; retain artifacts") + if stat.S_ISLNK(info.st_mode): + target = path.resolve(strict=True) + targets.add(target) + pending.append(target) + # Retaining a child retains the entire candidate. Its sibling + # links may protect further artifacts outside the target tree. + for candidate in candidates or []: + if (candidate not in queued_artifacts and overlaps(target, candidate) + and candidate.exists()): + queued_artifacts.add(candidate) + pending.append(candidate) + elif stat.S_ISDIR(info.st_mode): + with os.scandir(path) as children: + entries.extend(Path(child.path) for child in children) + return targets + + +def protections(plan: dict, paths: list[Path]) -> dict[Path, list[str]]: + check_bindings(plan) + observed = process_references(list(dict.fromkeys([ + *plan["candidates"], *plan.get("quarantined_paths", []), *paths, + ]))) + found = {p: list(reasons) for p, reasons in observed.items()} + for protected in plan["protected"]: + for candidate in found: + if overlaps(candidate, protected): + found[candidate].append("active release, rollback, transaction or explicit dependency") + # Every retained candidate protects its complete dependency closure, even + # when this is a single-path rescan or only a nested file is protected. + roots = [p for p in plan["protected"] if p.exists()] + roots.extend(p for p, reasons in found.items() if reasons and p.exists()) + for target in symlink_dependencies(roots, candidates=list(found)): + for candidate in paths: + if overlaps(target, candidate): + found[candidate].append("symlink dependency of retained artifact") + return {p: found[p] for p in paths} + + +def run_checks(plan: dict) -> None: + for check in plan["checks"]: + # Inventory commands are operator-owned read-only checks, with no shell + # interpolation. Their output stays private; failures identify the check. + try: + result = subprocess.run(check["argv"], cwd=check["cwd"], timeout=60, + stdin=subprocess.DEVNULL, stdout=subprocess.DEVNULL, + stderr=subprocess.DEVNULL, check=False) + except (OSError, subprocess.SubprocessError) as exc: + # TimeoutExpired includes full argv by default (possibly auth flags). + raise CleanupError( + f"acceptance check did not complete: {check['name']} ({type(exc).__name__})" + ) from None + if result.returncode: + raise CleanupError(f"acceptance check failed: {check['name']}") + + +def sync_directory(path: Path) -> None: + directory = os.open(path, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW) + try: + os.fsync(directory) + finally: + os.close(directory) + + +def write_journal(path: Path, report: dict) -> None: + temporary = path.with_name(f"{path.name}.{uuid.uuid4().hex}.tmp") + fd = os.open(temporary, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o600) + try: + with os.fdopen(fd, "w") as stream: + json.dump(report, stream, indent=2, ensure_ascii=True) + stream.write("\n") + stream.flush() + os.fsync(stream.fileno()) + os.replace(temporary, path) + sync_directory(path.parent) + finally: + temporary.unlink(missing_ok=True) + + +def cleanup(inventory: Path, *, apply: bool = False) -> dict: + plan = load_inventory(inventory) + journal = plan["root"] / JOURNAL + # Share the installer's lock. A dropped control connection cannot overlap a + # second installer/cleanup; leftover journals are inspected, never replayed. + lock = acquire_install_lock(plan["root"]) + moved: list[tuple[Path, Path, dict]] = [] + report = {"schema": 1, "inventory_sha256": plan["digest"], "started_at": time.time(), + "mode": "apply" if apply else "preview", "status": "preview", + "process_scope": "all users" if os.geteuid() == 0 else f"uid {os.geteuid()}", + "retained": [str(p) for p in plan["protected"]], "candidates": [], + "removed_allocated_bytes": 0} + try: + if journal.exists(): + previous, _ = read_json(journal) + if previous.get("status") not in {"complete", "deferred", "aborted"}: + raise CleanupError("unfinished cleanup journal; inspect original operation before retrying") + existing = [p for p in plan["candidates"] if p.exists()] + guarded = protections(plan, existing) + for path in plan["candidates"]: + row = {"path": str(path), "state": "absent" if not path.exists() else "planned", + "reasons": guarded.get(path, [])} + if row["reasons"]: + row["state"] = "deferred" + elif path.exists(): + row["identity"] = list(identity(path)) + row["allocated_bytes"] = sum(info.st_blocks * 512 for _, info in artifact_entries(path)) + report["candidates"].append(row) + if not apply: + return report + if not shutil.rmtree.avoids_symlink_attacks: + raise CleanupError("fd-safe tree removal is unavailable") + run_checks(plan) + report["status"] = "quarantining" + write_journal(journal, report) + for row in report["candidates"]: + if row["state"] != "planned": + continue + path = absolute_path(row["path"]) + guarded = protections(plan, [path])[path] + if guarded: + row.update(state="deferred", reasons=guarded) + continue + if list(identity(path)) != row["identity"]: + raise CleanupError("candidate replaced since preview") + quarantine = path.with_name(f".cc-remote-retired-{uuid.uuid4().hex}") + row.update(state="rename_pending", quarantine=str(quarantine)) + write_journal(journal, report) + path.rename(quarantine) + moved.append((path, quarantine, row)) + row["state"] = "quarantined" + plan.setdefault("quarantined_paths", []).append(quarantine) + sync_directory(path.parent) + write_journal(journal, report) + # Keep bytes recoverable until the now-retired paths pass fresh checks. + run_checks(plan) + report["status"] = "removing" + write_journal(journal, report) + for path, quarantine, row in moved: + guarded = protections(plan, [path, quarantine]) + reasons = guarded[path] + guarded[quarantine] + if reasons: + if path.exists() or path.is_symlink(): + raise CleanupError("original path reappeared; inspect before restoring") + quarantine.rename(path) + row.update(state="deferred", reasons=reasons) + else: + if list(identity(quarantine)) != row["identity"]: + raise CleanupError("quarantined artifact changed identity") + # Recheck for mounts before a recursive, symlink-safe removal. + for _ in artifact_entries(quarantine): + pass + parent = os.open(absolute_path(quarantine.parent.as_posix()), + os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW) + try: + row["state"] = "remove_pending" + write_journal(journal, report) + if quarantine.is_dir(): + shutil.rmtree(quarantine.name, dir_fd=parent) + else: + os.unlink(quarantine.name, dir_fd=parent) + finally: + os.close(parent) + row["state"] = "removed" + report["removed_allocated_bytes"] += row["allocated_bytes"] + sync_directory(path.parent) + write_journal(journal, report) + run_checks(plan) + check_bindings(plan) + report["status"] = ("deferred" if any(r["state"] == "deferred" for r in report["candidates"]) + else "complete") + report["verified_at"] = time.time() + write_journal(journal, report) + return report + except BaseException: + # Restore quarantined *whole* artifacts. A partially deleted artifact is + # never relabeled as restored; its journal remains an unresolved outcome. + if apply and report["status"] != "preview": + restored = report["status"] == "quarantining" + for path, quarantine, row in reversed(moved): + if (row["state"] == "quarantined" and quarantine.exists() + and not path.exists() and not path.is_symlink()): + try: + quarantine.rename(path) + sync_directory(path.parent) + row["state"] = "restored" + except OSError: + restored = False + if any(row["state"] not in {"restored", "planned", "deferred", "absent"} + for row in report["candidates"]): + restored = False + report["status"] = "aborted" if restored else "requires_inspection" + write_journal(journal, report) + raise + finally: + os.close(lock) + + +def main(argv: list[str] | None = None) -> int: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("inventory", type=Path, help="private, operator-reviewed JSON inventory") + parser.add_argument("--apply", action="store_true", help="quarantine, validate, then remove eligible artifacts") + args = parser.parse_args(argv) + try: + report = cleanup(args.inventory, apply=args.apply) + print(json.dumps(report, indent=2, ensure_ascii=True)) + return 2 if report["status"] == "deferred" else 0 + except (OSError, ValueError, subprocess.SubprocessError) as exc: + print(f"Cleanup stopped ({type(exc).__name__}): {exc}", file=sys.stderr) + return 1 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/docs/codex-session-messages.md b/docs/codex-session-messages.md new file mode 100644 index 00000000..e9e810e3 --- /dev/null +++ b/docs/codex-session-messages.md @@ -0,0 +1,45 @@ +# Codex cross-session messages + +cc-remote displays the native Codex App cross-thread envelope. Sending continues +through the optional official `codex_app.send_message_to_thread` tool; this +feature does not install/enable the adapter, grant tool permissions, invent a +second message transport, or submit prompts itself. See +[codex-app-tools.md](codex-app-tools.md) for the shared account/daemon prerequisite. + +Incoming messages show **来自会话 · name**. The Wrapper recognizes both the older +`userMessage` envelope and the newer `functionCallOutput` from the `codex_app` +namespace (`create_thread`, `send_message_to_thread`, `handoff_thread`). It keeps +the native item id and source thread id across live delivery, steering, official +history, rollout compatibility reads and browser caches. Unknown/malformed +markup remains ordinary text. Provenance is native display metadata, not an +independent authentication claim. + +Native send calls show a compact outgoing receipt with the actual tool outcome. +A failed call does not become “sent”. The small receipt survives summary +materialization when the source includes that tool; for an older summary that +omits tool items, expand the process to read them. Tool bodies stay deferred. + +A source/recipient link resolves against the current device's Codex Code catalog +and the current account routing prefix. Missing, deleted, cross-account and +Work targets remain labelled but disabled. The link opens a read-only history +view in the same window. It only uses `GetHistory` / `GetTurnDetail`; it does not +change Wrapper focus, resume an engine, send a message or navigate Codex App. +The original conversation stays mounted with output-follow paused. Closing the +view (Back or Escape) restores focus without scrolling the original chat. +The preview reads bounded pages, rejects mismatched cursors/revisions, and offers +retry after timeout or history invalidation. + +The native envelope supplies a source thread id, not the sending call/turn id. +The link therefore opens that session, without guessing a corresponding record +from matching text. Codex App keeps its own native rendering and source links. +End-to-end sending still requires a running eligible App and its tool adapter. + +This change uses protocol **v74**. Deploy Web, Relay and Wrapper together. +Derived Codex history pages and browser projections rebuild on upgrade; native +transcripts, credentials and original tool outputs are unchanged. + +Validation uses native-shape fixtures without model calls: incoming format and +escaping, exact ids, live/history replay races, same-text separate messages, +steering, sender failure receipts, account/device routing, history pagination, +and desktop/mobile browser return-position checks. Live App/MCP delivery is a +separate acceptance check; fixture success alone does not prove it. diff --git a/docs/installation.md b/docs/installation.md index 4032ddb2..4914cf76 100644 --- a/docs/installation.md +++ b/docs/installation.md @@ -236,7 +236,7 @@ npm --prefix web run build # 产出 web/dist/ 网页构建不需要任何登录密钥。 **所有目标先 staging,再改动线上服务。** 下文分别描述 Relay 和 Wrapper, -不能在 Wrapper staging 未验证时先激活 Relay。协议 v73 不允许混用旧客户端: +不能在 Wrapper staging 未验证时先激活 Relay。协议 v74 不允许混用旧客户端: 停止不兼容的旧 Wrapper,激活 Relay + Web,再激活 Wrapper 并硬刷新网页。 Wrapper 激活须通过 `deploy/work_registry_snapshot.py` 保存 Work SQLite 与私有账号 控制状态,不再按“是否来自某个旧协议”决定是否保护。回滚先恢复匹配状态,再启动 @@ -294,7 +294,7 @@ sudo bash ~/cc-remote-upload/deploy/setup-vps.sh \ 脚本会:装 `python3-venv` + Caddy、建 `ccremote` 系统用户、创建不可变 release 和 release-local venv、合并 Caddy 配置、原子切换 `current`,再重启 relay。若新 relay 重启或健康检查失败,`current`、Caddyfile、systemd unit 会作为一个事务全部 -恢复,并验证旧 release 的 `/healthz`。成功后再启动 v73 wrapper。 +恢复,并验证旧 release 的 `/healthz`。成功后再启动 v74 wrapper。 验证: diff --git a/docs/installation_en.md b/docs/installation_en.md index 97f5be79..7b4970a1 100644 --- a/docs/installation_en.md +++ b/docs/installation_en.md @@ -276,7 +276,7 @@ as described in the deployment contract. No browser secret is needed for a build **Stage every target before changing live services.** The commands below describe the Relay and Wrapper separately; do not activate Relay until every Wrapper stage -has passed validation. Protocol v73 cannot be mixed with older clients. Stop old +has passed validation. Protocol v74 cannot be mixed with older clients. Stop old incompatible Wrappers, activate Relay + Web, then activate Wrappers and hard-refresh browser tabs. Wrapper activation must snapshot Work SQLite and private profile control state with `deploy/work_registry_snapshot.py`; this is not limited to @@ -338,7 +338,7 @@ The script installs `python3-venv` + Caddy, creates the `ccremote` service user, builds an immutable release and its venv, merges Caddy configuration, atomically switches `current`, and restarts the relay. If restart/readiness fails, `current`, the Caddyfile, and the systemd unit roll back as one transaction and the previous -release's `/healthz` is verified. Start the v73 wrapper after success. +release's `/healthz` is verified. Start the v74 wrapper after success. Verify: diff --git a/tests/test_claude_profiles.py b/tests/test_claude_profiles.py index 2b4869f2..8d546fad 100644 --- a/tests/test_claude_profiles.py +++ b/tests/test_claude_profiles.py @@ -33,6 +33,7 @@ from cc_remote.wrapper.machine import WrapperMachine from cc_remote.wrapper.session_presentation import SessionPresentationStore from cc_remote.viewer_pages import PageRef, PageScope, ViewerPageStore +from tests.test_multisession import _mk_ctx NATIVE_ID = "11111111-1111-4111-8111-111111111111" @@ -845,6 +846,80 @@ def test_claude_profile_transitions_migrate_viewer_scopes_once( assert multi.sessions == recovered.sessions == single.sessions == {} +@pytest.mark.parametrize("late_read", ["presentation", "broker"]) +@pytest.mark.parametrize("next_state", ["idle", "running", None]) +def test_claude_catalog_samples_activity_after_async_reads( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + late_read: str, + next_state: str | None, +) -> None: + personal, company = tmp_path / "personal", tmp_path / "company" + _write_transcript(personal, "personal prompt") + _write_transcript(company, "company prompt") + monkeypatch.setenv("CLAUDE_CONFIG_DIR", str(personal)) + cfg = WrapperConfig() + cfg.state_dir = tmp_path / "state" + cfg.claude_work_root = tmp_path / "work" / "claude" + cfg.codex_work_root = tmp_path / "work" / "codex" + # Exercise both namespaced accounts and the legacy broker's final await. + cfg.claude_profiles_json = ( + _profiles(personal, company) if late_read == "presentation" else "" + ) + transport = _StubTransport() + machine = WrapperMachine(cfg, transport) + + async def run(): + sid = machine._claude_wire_sid(machine._claude_profiles.default, NATIVE_ID) + ctx = _mk_ctx(sid, NATIVE_ID) + ctx.claude_profile_id = machine._claude_profiles.default.id + ctx.state = "idle" if next_state == "running" else "running" + machine.sessions[sid] = ctx + entered, release = asyncio.Event(), asyncio.Event() + + async def delayed_read(*_args): + entered.set() + await release.wait() + return {"sessions": []} + + if late_read == "presentation": + monkeypatch.setattr( + machine, "_claim_legacy_presentation_from_claude_catalog", + delayed_read, + ) + else: + machine._claude_broker_enabled = True + machine._claude_broker = SimpleNamespace(list=delayed_read) + + pending = asyncio.create_task(machine._handle_list_sessions(ListSessions( + engine="claude", space="code", client_id="viewer"))) + await asyncio.wait_for(entered.wait(), timeout=5) + try: + # The terminal/start/eviction reaches the viewer while the older + # catalog request is still waiting on unrelated metadata I/O. + if next_state is None: + del machine.sessions[sid] + else: + await machine._set_state(ctx, next_state) + finally: + release.set() + listing = await asyncio.wait_for(pending, timeout=5) + assert isinstance(listing, SessionList) + row = next(row for row in listing.sessions if row.session_id == sid) + assert row.state == next_state + assert transport.sent[-1] is listing + if next_state is not None: + assert transport.sent[-2].type == "state" + assert transport.sent[-2].sid == sid + assert transport.sent[-2].state == next_state + if late_read == "presentation": + other = next(row for row in listing.sessions + if row.session_id == f"company@{NATIVE_ID}") + assert other.state is None + + asyncio.run(run()) + + def test_claude_session_list_fails_closed_during_profile_migration() -> None: machine = WrapperMachine.__new__(WrapperMachine) transport = _StubTransport() diff --git a/tests/test_codex_controls.py b/tests/test_codex_controls.py index 2a7b50ed..a6519c9e 100644 --- a/tests/test_codex_controls.py +++ b/tests/test_codex_controls.py @@ -3406,8 +3406,20 @@ async def churning_effort(_ctx): asyncio.run(run()) -def test_codex_resume_restores_rollout_effort_after_nullable_resume( - monkeypatch, tmp_path): +@pytest.mark.parametrize( + "rollout_model,native_model,explicit_model,expected_model,expected_updates", + [ + ("retired-model", "retired-model", None, "current-model", ["current-model"]), + (None, "selected-model", None, "selected-model", []), + ("retired-model", "selected-model", None, "selected-model", []), + ("current-model", "selected-model", None, "selected-model", []), + (None, "retired-model", None, "current-model", ["current-model"]), + (None, "selected-model", "current-model", "current-model", ["current-model"]), + ], +) +def test_codex_resume_preserves_native_model_and_restores_nullable_effort( + monkeypatch, tmp_path, rollout_model, native_model, explicit_model, + expected_model, expected_updates): class FakeCodexHandle: def __init__(self, _cfg, cwd=None, daemon_mode=None, daemon_manager=None): @@ -3452,9 +3464,9 @@ async def connect(self, **kwargs): self.thread_id = kwargs["resume_id"] self.preconnect_effort = self.effort self.preconnect_model = self.model - # Model the resume response echoing the retired native value over - # the wrapper's pre-connect catalog replacement. - self.model = "retired-model" + # The native selection can be newer than a stale or unreadable + # rollout tail; availability must be decided after this receipt. + self.model = native_model # Some app-server versions omit the effective override on resume. self.effort = None self.applied_effort = None @@ -3487,7 +3499,7 @@ async def run(): machine_module, "codex_session_settings", lambda *_args, **_kwargs: { - "model": "retired-model", + "model": rollout_model, "effort": "high", }, ) @@ -3498,6 +3510,10 @@ async def catalog(*, codex_home=None): "efforts": ["high"], "default_effort": "high", "is_default": True, + }, { + "id": "selected-model", + "efforts": ["high"], + "default_effort": "high", }] monkeypatch.setattr(machine_module, "codex_catalog", catalog) @@ -3517,12 +3533,13 @@ async def unchanged_effort(_model, effort, **_kwargs): resume_id=thread_id, engine="codex", space="code", + model=explicit_model, ) assert ctx is not None assert ctx.sdk.preconnect_model == "current-model" - assert ctx.sdk.set_model_calls == ["current-model"] - assert ctx.sdk.model == "current-model" + assert ctx.sdk.set_model_calls == expected_updates + assert ctx.sdk.model == expected_model assert ctx.sdk.preconnect_effort == "high" assert ctx.sdk.effort == "high" assert ctx.sdk.applied_effort == "high" diff --git a/tests/test_codex_daemon.py b/tests/test_codex_daemon.py index d3eeb8be..1a27118d 100644 --- a/tests/test_codex_daemon.py +++ b/tests/test_codex_daemon.py @@ -179,6 +179,31 @@ def run(argv, **_kwargs): ) +def test_daemon_command_does_not_inherit_retired_release_cwd(tmp_path, monkeypatch): + release = tmp_path / "release" + home = tmp_path / "account-home" + release.mkdir() + home.mkdir() + monkeypatch.chdir(release) + # Exercise a real child, without invoking Codex or starting a daemon. + result = daemon_module._run_command( + (daemon_module.sys.executable, "-c", "import os; print(os.getcwd())"), + {"HOME": str(home)}, 5, + ) + assert result.returncode == 0 + assert Path(result.stdout.decode().strip()) == home.resolve() + + +@pytest.mark.parametrize("home", ["relative-home", "/missing-cc-remote-test-home"]) +def test_daemon_command_refuses_unsafe_working_directory(monkeypatch, home): + def unexpected(*args, **kwargs): + pytest.fail("must not fall back to disposable caller cwd") + monkeypatch.setattr(daemon_module.subprocess, "run", unexpected) + result = daemon_module._run_command(("codex", "app-server", "daemon", "start"), + {"HOME": home}, 1) + assert result.returncode == 127 + + @pytest.mark.parametrize("listener,allow_local,accepted", [ ("remote", False, True), ("local", True, True), ("default", True, True), ("other", True, False), ("local", False, False), ("ambiguous", True, False), diff --git a/tests/test_codex_delegation.py b/tests/test_codex_delegation.py new file mode 100644 index 00000000..08c3b563 --- /dev/null +++ b/tests/test_codex_delegation.py @@ -0,0 +1,165 @@ +"""Native Codex App cross-thread inputs: live, official history and rollout.""" +import asyncio +import json +from types import SimpleNamespace + +import pytest + +from cc_remote.protocol import UserMsg, TurnBinding, ConversationTurn, serialize, deserialize +from cc_remote.wrapper.codex_delegation import parse_codex_delegation, codex_message_target +from cc_remote.wrapper.codex_external import visible_codex_user_message +from cc_remote.wrapper.codex_handle import CodexHandle +from cc_remote.wrapper.codex_history import CodexOfficialHistory +from cc_remote.wrapper.codex_stream import ( + codex_live_user_message, codex_translate_history, codex_history_window_info, +) +from cc_remote.wrapper.history_store import materialize_history_turns +from tests.test_codex_history import _agent, _turn +from tests.test_codex_spontaneous_stream import _notification +from tests.test_multisession import _mk_ctx, _mk_machine + +SOURCE = "01a01234-1234-7890-abcd-123456789abc" +ENVELOPE = (f"\n{SOURCE}\n" + "Check <code> & &lt;literal&gt;\n") +PROMPT = "Check & <literal>" + + +def _input(kind, item_id="incoming"): + if kind == "legacy": + return {"type": "userMessage", "id": item_id, + "content": [{"type": "text", "text": ENVELOPE}]} + return {"type": "functionCallOutput", "id": item_id, "namespace": "codex_app", + "name": "send_message_to_thread", "output": ENVELOPE} + + +@pytest.mark.parametrize("kind", ["legacy", "native"]) +def test_same_native_identity_and_provenance_in_live_history_and_rollout(tmp_path, kind): + item = _input(kind) + user = codex_live_user_message(_notification("item/completed", "task", item=item)) + assert (user.message_id, user.prompt, user.source_thread_id) == ("incoming", PROMPT, SOURCE) + + async def run(): + async def rpc(method, params, *_): + assert method == "thread/turns/list" + return {"data": [_turn("task", [item, _agent("answer", "Checked.")])], "nextCursor": None} + page = await CodexOfficialHistory(65536, rpc=rpc).summary_page( + "receiver", before=None, limit=4) + assert len(page.turns) == 1 + assert page.turns[0]["id"] == user.message_id + assert page.turns[0]["sourceThreadId"] == SOURCE + assert page.turns[0]["prompt"] == PROMPT + ConversationTurn.model_validate(page.turns[0]) + asyncio.run(run()) + + rollout_item = {**item, "type": "UserMessage" if kind == "legacy" else "FunctionCallOutput"} + rows = [ + {"type": "event_msg", "payload": {"type": "task_started", "turn_id": "task"}}, + {"type": "event_msg", "payload": {"type": "item_completed", "turn_id": "task", "item": rollout_item}}, + {"type": "event_msg", "payload": {"type": "agent_message", "message": "Checked."}}, + {"type": "event_msg", "payload": {"type": "task_complete", "turn_id": "task"}}, + ] + path = tmp_path / "rollout.jsonl" + path.write_text("".join(json.dumps(row) + "\n" for row in rows)) + events, _ = codex_translate_history(str(path), 65536) + users = [e for e in events if isinstance(e, UserMsg)] + assert [(u.msg_id, u.prompt, u.source_thread_id) for u in users] == [("incoming", PROMPT, SOURCE)] + turn = materialize_history_turns([e.model_dump(mode="json") for e in events])[0] + assert turn["sourceThreadId"] == SOURCE + assert deserialize(serialize(users[0])).source_thread_id == SOURCE + window = codex_history_window_info(str(path), before=None, limit=1) + assert window is not None + + +@pytest.mark.parametrize("kind", ["legacy", "native", "native-partial"]) +def test_spontaneous_delivery_publishes_one_native_user_with_no_placeholder(kind): + async def run(): + machine, transport = _mk_machine() + ctx = _mk_ctx("thread-spontaneous", "thread-spontaneous") + ctx.engine = "codex" + handle = CodexHandle(machine.cfg) + handle.thread_id = ctx.session_id + handle.proc = SimpleNamespace(returncode=None) + ctx.sdk = handle + machine.sessions[ctx.key] = ctx + handle.turn_lifecycle_callback = lambda phase, turn_id: machine._on_codex_turn_lifecycle(ctx, phase, turn_id) + item = _input(kind) + for message in [ + _notification("turn/started", "task", turn={"id": "task"}), + _notification("item/started", "task", item={**item, "output": ""} if kind == "native-partial" else item), + _notification("item/completed", "task", item=item), + _notification("item/completed", "task", item=_agent("answer", "Checked.")), + _notification("turn/completed", "task", turn={"id": "task", "status": "completed"}), + ]: + await handle._dispatch(message) + await asyncio.wait_for(ctx.codex_spontaneous_task, 1) + users = [e for e in transport.sent if isinstance(e, UserMsg)] + assert [(u.msg_id, u.prompt, u.source_thread_id) for u in users] == [("incoming", PROMPT, SOURCE)] + assert [(e.msg_id, e.turn_id) for e in transport.sent if isinstance(e, TurnBinding)] == [("incoming", "task")] + asyncio.run(run()) + + +def test_unknown_markup_and_foreign_tool_outputs_remain_unattributed(): + for text in [ENVELOPE + "extra", ENVELOPE.replace(SOURCE, "other@thread"), + ENVELOPE.replace("<code>", ""), "literal"]: + assert parse_codex_delegation(text) is None + assert visible_codex_user_message(text) == text + for override in [{"namespace": "foreign"}, {"name": "read_thread"}, {"name": []}, {"id": "../bad"}]: + assert codex_live_user_message(_notification("item/completed", "task", item={**_input("native"), **override})) is None + assert parse_codex_delegation( + f"{SOURCE}" + + " " * (512 * 1024)) is None + + +def test_sender_receipt_requires_native_tool_and_retains_failure(): + assert codex_message_target("send_message_to_thread", {"threadId": SOURCE}, "foreign") is None + events = [ + {"type": "user_msg", "msg_id": "user", "prompt": "review"}, + {"type": "tool_use", "message_id": "assistant", "tool_use_id": "send", "tool": "send_message_to_thread", + "server": "codex_app", "input": {"threadId": SOURCE, "input": "private tool body"}}, + {"type": "tool_result", "tool_use_id": "send", "content": "denied", "is_error": True}, + ] + turn = materialize_history_turns(events)[0] + assert turn["sessionMessages"] == [{"itemId": "send", "threadId": SOURCE, "status": "failed"}] + assert "private tool body" not in json.dumps(turn) + + +@pytest.mark.parametrize("kind", ["mcpToolCall", "dynamicToolCall"]) +@pytest.mark.parametrize("success", [True, False]) +def test_official_sender_tool_outcome_becomes_summary_receipt(kind, success): + tool = { + "id": "send-call", "type": kind, "tool": "send_message_to_thread", + "arguments": {"threadId": SOURCE, "input": "Check the code"}, + "status": "completed" if success else "failed", + } + if kind == "mcpToolCall": + tool.update(server="codex_app", result={"content": []} if success else None, + error=None if success else {"message": "Target unavailable"}) + else: + tool.update(namespace="codex_app", success=success, contentItems=[]) + + async def run(): + async def rpc(method, params, *_): + return {"data": [_turn("task", [ + _input("legacy"), tool, _agent("answer", "Finished."), + ])], "nextCursor": None} + page = await CodexOfficialHistory(65536, rpc=rpc).summary_page( + "receiver", before=None, limit=4) + assert page.turns[0]["sessionMessages"] == [{ + "itemId": "send-call", "threadId": SOURCE, + "status": "sent" if success else "failed", + }] + asyncio.run(run()) + + +def test_native_steers_keep_separate_item_identities_and_sources(): + async def run(): + async def rpc(method, params, *_): + return {"data": [_turn("task", [ + _input("legacy", "first"), _agent("comment", "Working.", phase="commentary"), + _input("native", "second"), _agent("answer", "Done."), + ])], "nextCursor": None} + page = await CodexOfficialHistory(65536, rpc=rpc).summary_page("receiver", before=None, limit=4) + assert [turn["id"] for turn in page.turns] == ["first", "second"] + assert [turn["prompt"] for turn in page.turns] == [PROMPT, PROMPT] + assert [turn["sourceThreadId"] for turn in page.turns] == [SOURCE, SOURCE] + asyncio.run(run()) diff --git a/tests/test_codex_readiness.py b/tests/test_codex_readiness.py index e17613b5..f07a2cf6 100644 --- a/tests/test_codex_readiness.py +++ b/tests/test_codex_readiness.py @@ -309,6 +309,62 @@ def test_install_report_keeps_cli_discovery_unverified_and_errors_separate(capsy assert "版本不一致" in output +@pytest.mark.parametrize("mode", ["ready", "deleted_cwd", "invalid", "hang", "swapped"]) +def test_live_config_probe_rejects_healthy_handshake_but_unusable_config(monkeypatch, mode): + calls = [] + monkeypatch.setattr(readiness, "_TIMEOUT", 0.15 if mode == "hang" else 3) + async def handler(ws): + request = json.loads(await ws.recv()) + calls.append(request["method"]) + await ws.send(json.dumps({"id": 1, "result": {}})) + calls.append(json.loads(await ws.recv())["method"]) + request = json.loads(await ws.recv()) + calls.append(request["method"]) + assert request["params"] == {"includeLayers": False} + if mode == "hang": + await ws.wait_closed() + return + result = {"config": {"private-value": "must-not-be-printed"}} + if mode == "invalid": + result = {} + response = ({"id": 2, "error": {"message": "private config path: ENOENT"}} + if mode == "deleted_cwd" else {"id": 2, "result": result}) + await ws.send(json.dumps(response)) + await ws.wait_closed() + async def check(): + with tempfile.TemporaryDirectory(prefix="cc-config-", dir="/tmp") as directory: + path = str(Path(directory).resolve() / "socket") + async with unix_serve(handler, path, close_timeout=0.1): + os.chmod(path, 0o600) + if mode == "swapped": + identities = iter([(1, 2, 3), (1, 4, 5)]) + monkeypatch.setattr(readiness, "socket_identity", lambda *a, **kw: next(identities)) + if mode == "ready": + await readiness.probe_config(path) + else: + with pytest.raises((RuntimeError, TimeoutError)) as error: + await readiness.probe_config(path) + assert "private" not in str(error.value) + asyncio.run(check()) + assert calls == ["initialize", "initialized", "config/read"] + + +@pytest.mark.parametrize("wrong_user", [False, True]) +def test_live_installer_check_is_fresh_and_runs_only_as_service_owner(monkeypatch, wrong_user): + calls = [] + async def probe(path): + calls.append(path) + raise RuntimeError("secret diagnostic") + monkeypatch.setattr(installer, "probe_config", probe) + monkeypatch.setattr(installer, "process_owner_uid", lambda pid: os.geteuid() + int(wrong_user)) + report = {"wrapper": {"pid": 123}, "profiles": [{"profile": "account", "status": "ready", + "home": "/account", "socket": "/account/app-server-control/app-server-control.sock"}]} + asyncio.run(installer.refresh_configuration(report)) + assert report["profiles"][0]["reason"] == ("wrong_probe_user" if wrong_user else "configuration_failed") + assert len(calls) == int(not wrong_user) + assert "secret" not in json.dumps(report) + + def test_missing_receipt_does_not_pass_acceptance(tmp_path, capsys): assert installer.main([ "--home", str(tmp_path), "--release", str(tmp_path), "--after", "0", "--wait", "0", diff --git a/tests/test_codex_session_migration.py b/tests/test_codex_session_migration.py index 8fe93173..36a38da5 100644 --- a/tests/test_codex_session_migration.py +++ b/tests/test_codex_session_migration.py @@ -91,7 +91,7 @@ async def list_sessions(_cmd): def test_session_migration_protocol_roundtrips_as_control_frames(): - assert PROTOCOL_VERSION == 73 + assert PROTOCOL_VERSION == 74 command = deserialize(serialize(_command("/tmp/new-cwd"))) assert command.type == "migrate_session" assert command.session_id == "thread-1" diff --git a/tests/test_deploy_cleanup.py b/tests/test_deploy_cleanup.py new file mode 100644 index 00000000..097b1d8a --- /dev/null +++ b/tests/test_deploy_cleanup.py @@ -0,0 +1,388 @@ +"""Cleanup protects live dependencies and validates retirement before deletion.""" +from __future__ import annotations + +import json +import os +from pathlib import Path +import shutil +import subprocess +import sys + +import pytest + +from deploy import cleanup as module +from deploy.install_lock import acquire_install_lock + +REAL_PROCESS_REFERENCES = module.process_references + + +@pytest.fixture +def installation(tmp_path, monkeypatch): + root = tmp_path.resolve() / "installation" + releases = root / "releases" + releases.mkdir(parents=True) + current, previous, old = [releases / name for name in ["current-build", "previous-build", "old-build"]] + for path in [current, previous, old]: + path.mkdir() + (path / "payload").write_text("must not be lost") + (root / "current").symlink_to(current) + transaction = root / "activation.json" + transaction.write_text(json.dumps({"phase": "committed"})) + inventory = root / "cleanup.json" + data = { + "schema": 1, "installation_root": str(root), "current_release": str(current), + "rollback_paths": [str(previous)], "cleanup_roots": [str(releases)], + "candidates": [str(old)], + "transactions": [{"path": str(transaction), "field": "phase", "equals": "committed"}], + "checks": [{"name": "fixture health", "argv": [sys.executable, "-c", "pass"], "cwd": str(root)}], + } + inventory.write_text(json.dumps(data)) + monkeypatch.setattr(module, "process_references", lambda paths: {p: [] for p in paths}) + return root, current, previous, old, inventory, data + + +def update_inventory(installation, **changes): + *_, inventory, data = installation + data.update(changes) + inventory.write_text(json.dumps(data)) + + +def test_preview_never_moves_deletes_or_runs_acceptance(installation, monkeypatch): + root, _, _, old, inventory, _ = installation + monkeypatch.setattr(module, "run_checks", lambda plan: pytest.fail("preview must not run commands")) + result = module.cleanup(inventory) + assert result["status"] == "preview" + assert result["candidates"][0]["state"] == "planned" + assert old.is_dir() + assert not (root / module.JOURNAL).exists() + + +def test_success_retains_complete_rollback_and_checks_after_removal(installation, monkeypatch): + root, current, previous, old, inventory, _ = installation + observed = [] + def checks(plan): + observed.append((old.exists(), bool(list(old.parent.glob(".cc-remote-retired-*"))))) + monkeypatch.setattr(module, "run_checks", checks) + result = module.cleanup(inventory, apply=True) + assert observed == [(True, False), (False, True), (False, False)] + assert result["status"] == "complete" + assert result["removed_allocated_bytes"] > 0 + assert current.is_dir() and previous.is_dir() and not old.exists() + assert json.loads((root / module.JOURNAL).read_text())["status"] == "complete" + assert (root / module.JOURNAL).stat().st_mode & 0o777 == 0o600 + + +@pytest.mark.parametrize("protected", ["current", "rollback", "dependency"]) +def test_protected_artifact_cannot_be_retired(installation, protected): + _, current, previous, old, inventory, data = installation + if protected == "current": + update_inventory(installation, candidates=[str(current)]) + elif protected == "rollback": + update_inventory(installation, candidates=[str(previous)]) + else: + update_inventory(installation, protected_paths=[str(old / "payload")]) + result = module.cleanup(inventory, apply=True) + assert result["status"] == "deferred" + assert all(Path(p).exists() for p in data["candidates"]) + + +def test_retained_venv_symlink_protects_an_older_release(installation): + _, current, _, old, inventory, _ = installation + (current / ".venv").symlink_to(old) + result = module.cleanup(inventory, apply=True) + assert result["status"] == "deferred" + assert "symlink dependency" in " ".join(result["candidates"][0]["reasons"]) + assert old.is_dir() + + +def test_shared_venv_dependencies_are_transitive_and_cycles_are_bounded(installation): + root, current, _, old, inventory, _ = installation + shared = root / "shared-runtime" + shared.mkdir() + (current / ".venv").symlink_to(shared) + (shared / "package").symlink_to(old) + (old / "cycle").symlink_to(shared) + assert module.cleanup(inventory, apply=True)["status"] == "deferred" + assert old.is_dir() + + +def test_nested_protection_keeps_sibling_dependency_during_single_candidate_rescan(installation): + _, _, _, old, inventory, _ = installation + retained = old.parent / "retained-build" + retained.mkdir() + protected = retained / "service-config" + protected.write_text("protected") + (retained / ".venv").symlink_to(old) + update_inventory(installation, candidates=[str(old), str(retained)], + protected_paths=[str(protected)]) + # Pre-rename checks request one candidate, but another retained candidate + # still protects its dependencies even without a live process reference. + reasons = module.protections(module.load_inventory(inventory), [old])[old] + assert "symlink dependency" in " ".join(reasons) + assert module.cleanup(inventory, apply=True)["status"] == "deferred" + assert old.is_dir() and retained.is_dir() + + +@pytest.mark.parametrize("apply", [False, True]) +@pytest.mark.parametrize("reverse", [False, True]) +def test_dependency_on_candidate_child_retains_whole_dependency_closure(installation, apply, reverse): + root, current, _, old, inventory, _ = installation + retained = old.parent / "retained-build" + retained.mkdir() + (retained / "package").mkdir() + (current / "package").symlink_to(retained / "package") + (retained / ".venv").symlink_to(old) + (old / "cycle").symlink_to(retained) + candidates = [old, retained] + if reverse: + candidates.reverse() + update_inventory(installation, candidates=[str(p) for p in candidates]) + report = module.cleanup(inventory, apply=apply) + assert all(row["state"] == "deferred" for row in report["candidates"]) + assert old.is_dir() and retained.is_dir() + assert (retained / ".venv").resolve(strict=True) == old + assert not list(old.parent.glob(".cc-remote-retired-*")) + if apply: + assert json.loads((root / module.JOURNAL).read_text())["status"] == "deferred" + + +def test_newly_busy_candidate_protects_its_dependency_before_other_candidate_moves(installation, monkeypatch): + _, _, _, old, inventory, _ = installation + busy = old.parent / "busy-build" + busy.mkdir() + (busy / ".venv").symlink_to(old) + update_inventory(installation, candidates=[str(old), str(busy)]) + calls = 0 + def scan(paths): + nonlocal calls + calls += 1 + return {p: (["pid 42: cwd"] if p == busy and calls >= 2 else []) for p in paths} + monkeypatch.setattr(module, "process_references", scan) + report = module.cleanup(inventory, apply=True) + assert all(row["state"] == "deferred" for row in report["candidates"]) + assert old.is_dir() and busy.is_dir() + + +@pytest.mark.parametrize("fault", ["broken", "limit"]) +def test_incomplete_dependency_scan_never_deletes(installation, monkeypatch, fault): + _, current, _, old, inventory, _ = installation + if fault == "broken": + (current / ".venv").symlink_to(current / "missing") + else: + monkeypatch.setattr(module, "MAX_DEPENDENCY_ENTRIES", 1) + with pytest.raises((OSError, module.CleanupError)): + module.cleanup(inventory, apply=True) + assert old.is_dir() + + +@pytest.mark.parametrize("when", [1, 2, 3]) +def test_new_live_reference_defers_and_restores_candidate(installation, monkeypatch, when): + _, _, _, old, inventory, _ = installation + calls = 0 + def scan(paths): + nonlocal calls + calls += 1 + return {p: (["pid 42: cwd"] if calls == when else []) for p in paths} + monkeypatch.setattr(module, "process_references", scan) + result = module.cleanup(inventory, apply=True) + assert result["status"] == "deferred" + assert (old / "payload").read_text() == "must not be lost" + assert not list(old.parent.glob(".cc-remote-retired-*")) + + +def test_config_failure_after_quarantine_restores_original_path(installation, monkeypatch): + root, _, _, old, inventory, _ = installation + def config_read(plan): + if not old.exists(): + raise module.CleanupError("configuration dependency became unavailable") + monkeypatch.setattr(module, "run_checks", config_read) + with pytest.raises(module.CleanupError, match="configuration dependency"): + module.cleanup(inventory, apply=True) + assert (old / "payload").read_text() == "must not be lost" + report = json.loads((root / module.JOURNAL).read_text()) + assert report["status"] == "aborted" + assert report["candidates"][0]["state"] == "restored" + + +def test_partial_delete_is_not_misreported_as_restored(installation, monkeypatch): + root, _, _, old, inventory, _ = installation + def fail_remove(path, *, dir_fd): + os.unlink(f"{path}/payload", dir_fd=dir_fd) + raise OSError("injected interrupted removal") + fail_remove.avoids_symlink_attacks = True + monkeypatch.setattr(module.shutil, "rmtree", fail_remove) + with pytest.raises(OSError, match="interrupted removal"): + module.cleanup(inventory, apply=True) + report = json.loads((root / module.JOURNAL).read_text()) + assert report["status"] == "requires_inspection" + assert report["candidates"][0]["state"] == "remove_pending" + assert not old.exists() + with pytest.raises(module.CleanupError, match="unfinished cleanup"): + module.cleanup(inventory, apply=True) + + +@pytest.mark.parametrize("fault", ["transaction", "current", "candidate", "symlink", "scan", "lock", "journal"]) +def test_uncertain_state_never_deletes(installation, monkeypatch, fault): + root, _, previous, old, inventory, _ = installation + lock = None + if fault == "transaction": + (root / "activation.json").write_text('{"phase":"activating"}') + elif fault == "current": + (root / "current").unlink() + (root / "current").symlink_to(previous) + elif fault == "candidate": + update_inventory(installation, candidates=[str(root)]) + elif fault == "symlink": + alias = old.parent / "alias" + alias.symlink_to(old) + update_inventory(installation, candidates=[str(alias)]) + elif fault == "scan": + def scan(paths): + raise module.CleanupError("incomplete scan") + monkeypatch.setattr(module, "process_references", scan) + elif fault == "lock": + lock = acquire_install_lock(root) + elif fault == "journal": + (root / module.JOURNAL).write_text('{"status":"quarantining"}') + try: + with pytest.raises((ValueError, OSError)): + module.cleanup(inventory, apply=True) + assert old.is_dir() + finally: + if lock is not None: + os.close(lock) + + +def test_transaction_change_after_precheck_stops_before_retirement(installation, monkeypatch): + root, _, _, old, inventory, _ = installation + def checks(plan): + (root / "activation.json").write_text('{"phase":"committed","changed":true}') + monkeypatch.setattr(module, "run_checks", checks) + with pytest.raises(module.CleanupError, match="transaction changed"): + module.cleanup(inventory, apply=True) + assert old.is_dir() + + +@pytest.mark.parametrize("fault", ["inventory", "rollback"]) +def test_replaced_provenance_stops_cleanup(installation, monkeypatch, fault): + _, _, previous, old, inventory, _ = installation + def checks(plan): + if fault == "inventory": + inventory.write_text(inventory.read_text() + "\n") + else: + previous.rename(previous.with_name("moved-rollback")) + previous.mkdir() + monkeypatch.setattr(module, "run_checks", checks) + with pytest.raises(module.CleanupError, match="changed|replaced"): + module.cleanup(inventory, apply=True) + assert old.is_dir() + + +def test_lsof_parser_preserves_space_paths_and_cwd_without_argv(): + assert module.parse_lsof(b'p42\0\nfcwd\0n/opt/old release\0\nf9\0n/opt/file name\0\n') == [ + (42, "cwd", Path("/opt/old release")), (42, "9", Path("/opt/file name")), + ] + + +@pytest.mark.parametrize("record", [ + b'p42\0\nfNOFD\0nPermission denied\0', + b'p42\0\nfcwd\0n/proc/42/cwd (readlink: Permission denied)\0', + b'p42\0\nfrtd\0n/proc/42/root (readlink: Operation not permitted)\0', +]) +def test_unobservable_process_record_stops_cleanup(installation, monkeypatch, record): + root, _, _, old, inventory, _ = installation + monkeypatch.setattr(module, "scan_command", lambda argv: record) + monkeypatch.setattr(module.shutil, "which", lambda name: "/usr/bin/lsof") + monkeypatch.setattr(module, "process_references", REAL_PROCESS_REFERENCES) + with pytest.raises(module.CleanupError, match="not fully observable"): + module.cleanup(inventory, apply=True) + assert old.is_dir() + assert not list(old.parent.glob(".cc-remote-retired-*")) + assert not (root / module.JOURNAL).exists() + + +@pytest.mark.parametrize("failure", ["missing", "warning", "exit"]) +def test_incomplete_lsof_never_becomes_empty_safe_result(monkeypatch, failure): + monkeypatch.setattr(module.shutil, "which", lambda name: None if failure == "missing" else "/usr/bin/lsof") + monkeypatch.setattr(module.subprocess, "run", lambda *a, **kw: + subprocess.CompletedProcess(a, int(failure == "exit"), b'p1\0fcwd\0n/\0', + b'cannot stat' if failure == "warning" else b'')) + with pytest.raises(module.CleanupError): + module.process_references([Path("/candidate")]) + + +def test_check_timeout_never_discloses_private_arguments(installation, monkeypatch): + _, _, _, old, inventory, _ = installation + def timeout(*args, **kwargs): + raise subprocess.TimeoutExpired(["check", "secret-authorization-value"], 60) + monkeypatch.setattr(module.subprocess, "run", timeout) + with pytest.raises(module.CleanupError, match="did not complete") as error: + module.cleanup(inventory, apply=True) + assert "secret" not in str(error.value) + assert old.is_dir() + + +@pytest.mark.skipif(not shutil.which("lsof"), reason="live process regression requires lsof") +def test_cli_previews_and_applies_isolated_installation_with_real_scans(installation): + root, current, previous, old, inventory, _ = installation + command = [sys.executable, str(Path(module.__file__).resolve()), str(inventory)] + for apply in (False, True): + result = subprocess.run([*command, *(["--apply"] if apply else [])], + capture_output=True, text=True, timeout=60) + if result.returncode: + # Hosted Linux runners can have same-user nondumpable processes. + # The real CLI must refuse cleanup in that environment, not ignore + # denied /proc records or require elevated test-suite privileges. + assert result.returncode == 1 + assert result.stderr.strip() == ( + "Cleanup stopped (CleanupError): process files are not fully observable") + assert not result.stdout.strip() + assert old.is_dir() + assert not (root / module.JOURNAL).exists() + else: + report = json.loads(result.stdout) + assert report["status"] == ("complete" if apply else "preview") + assert report["candidates"][0]["state"] == ("removed" if apply else "planned") + assert old.exists() is not apply + assert current.is_dir() and previous.is_dir() + assert not list(old.parent.glob(".cc-remote-retired-*")) + + +@pytest.mark.skipif(not shutil.which("lsof"), reason="live process regression requires lsof") +@pytest.mark.parametrize("reference", ["cwd", "file"]) +def test_real_process_reference_absent_from_command_line_is_protected(installation, monkeypatch, reference): + _, _, _, old, inventory, _ = installation + code = ("import os,time; " + "f=open(os.environ['CLEANUP_TEST_FILE']) if os.environ['CLEANUP_TEST_KIND']=='file' else None; " + "print('ready',flush=True); time.sleep(30)") + proc = subprocess.Popen([sys.executable, "-c", code], + cwd=old if reference == "cwd" else old.parent, + env={**os.environ, "CLEANUP_TEST_FILE": str(old / "payload"), "CLEANUP_TEST_KIND": reference}, + stdout=subprocess.PIPE, text=True) + try: + assert proc.stdout.readline().strip() == "ready" + scan_command = module.scan_command + def scan_fixture_process(argv): + # Exercise real lsof and ps against this fixture's owned process, + # independently of unrelated runner processes with restricted /proc. + if Path(argv[0]).name == "lsof": + argv = [*argv, "-a", "-p", str(proc.pid)] + return scan_command(argv) + monkeypatch.setattr(module, "scan_command", scan_fixture_process) + found = REAL_PROCESS_REFERENCES([old])[old] + assert any(reason.startswith(f"pid {proc.pid}:") for reason in found) + assert f"pid {proc.pid}: argv" not in found + finally: + proc.terminate() + proc.wait(timeout=5) + + +def test_post_delete_health_failure_never_reports_success(installation, monkeypatch): + root, _, _, old, inventory, _ = installation + def check(plan): + if not old.exists() and not list(old.parent.glob(".cc-remote-retired-*")): + raise module.CleanupError("late health failure") + monkeypatch.setattr(module, "run_checks", check) + with pytest.raises(module.CleanupError, match="late health failure"): + module.cleanup(inventory, apply=True) + assert json.loads((root / module.JOURNAL).read_text())["status"] == "requires_inspection" diff --git a/tests/test_history_store.py b/tests/test_history_store.py index 44190822..25ef83c5 100644 --- a/tests/test_history_store.py +++ b/tests/test_history_store.py @@ -42,7 +42,7 @@ def test_public_codex_item_migration_invalidates_only_codex_projections(tmp_path connection.execute(f"PRAGMA user_version={old_version}") migrated = HistoryIndexStore(tmp_path / "state") with sqlite3.connect(migrated.path) as connection: - assert connection.execute("PRAGMA user_version").fetchone()[0] == 46 + assert connection.execute("PRAGMA user_version").fetchone()[0] == 47 for table in ("history_pages", "history_turn_details"): assert sorted(connection.execute(f"SELECT engine FROM {table}").fetchall()) == [ ("claude",), ("dsh",)] @@ -485,7 +485,7 @@ def test_paged_file_migration_rebuilds_summaries_once_without_removing_assets(tm connection.execute(f"PRAGMA user_version={old_version}") migrated = HistoryIndexStore(tmp_path / "state") with sqlite3.connect(migrated.path) as connection: - assert connection.execute("PRAGMA user_version").fetchone()[0] == 46 + assert connection.execute("PRAGMA user_version").fetchone()[0] == 47 for table in ("history_pages", "history_turn_details"): assert connection.execute( f"SELECT COUNT(*) FROM {table} WHERE engine='codex'" @@ -514,7 +514,7 @@ def test_compact_migration_rebuilds_claude_details_and_preserves_assets(tmp_path connection.execute(f"PRAGMA user_version={old_version}") migrated = HistoryIndexStore(tmp_path / "state") with sqlite3.connect(migrated.path) as connection: - assert connection.execute("PRAGMA user_version").fetchone()[0] == 46 + assert connection.execute("PRAGMA user_version").fetchone()[0] == 47 assert connection.execute( "SELECT engine FROM history_pages").fetchall() == [] assert connection.execute( @@ -589,7 +589,7 @@ def test_v19_migration_rebuilds_history_and_adds_agent_details(tmp_path): assert migrated.get_page( "session-1", "claude", source, before=None, limit=4) is None with sqlite3.connect(migrated.path) as connection: - assert connection.execute("PRAGMA user_version").fetchone()[0] == 46 + assert connection.execute("PRAGMA user_version").fetchone()[0] == 47 assert connection.execute( "SELECT COUNT(*) FROM history_agent_details").fetchone()[0] == 0 @@ -617,7 +617,7 @@ def test_v20_migration_rebuilds_codex_and_claude_identity_projections(tmp_path): migrated = HistoryIndexStore(state_dir) with sqlite3.connect(migrated.path) as connection: - assert connection.execute("PRAGMA user_version").fetchone()[0] == 46 + assert connection.execute("PRAGMA user_version").fetchone()[0] == 47 for table in ("history_pages", "history_turn_details"): assert connection.execute( f"SELECT COUNT(*) FROM {table} WHERE engine='claude'" @@ -678,7 +678,7 @@ def test_v21_migration_rebuilds_claude_alias_and_codex_process_projections( migrated = HistoryIndexStore(state_dir) with sqlite3.connect(migrated.path) as connection: - assert connection.execute("PRAGMA user_version").fetchone()[0] == 46 + assert connection.execute("PRAGMA user_version").fetchone()[0] == 47 for table in ("history_pages", "history_turn_details"): assert connection.execute( f"SELECT COUNT(*) FROM {table} WHERE engine='claude'" @@ -735,7 +735,7 @@ def test_v22_migration_applies_codex_and_claude_projection_repairs( migrated = HistoryIndexStore(state_dir) with sqlite3.connect(migrated.path) as connection: - assert connection.execute("PRAGMA user_version").fetchone()[0] == 46 + assert connection.execute("PRAGMA user_version").fetchone()[0] == 47 for table in ("history_pages", "history_turn_details"): assert connection.execute( f"SELECT COUNT(*) FROM {table} WHERE engine='claude'" @@ -797,7 +797,7 @@ def test_recent_migration_applies_codex_and_claude_projection_repairs( migrated = HistoryIndexStore(state_dir) with sqlite3.connect(migrated.path) as connection: - assert connection.execute("PRAGMA user_version").fetchone()[0] == 46 + assert connection.execute("PRAGMA user_version").fetchone()[0] == 47 for table in ("history_pages", "history_turn_details"): assert connection.execute( f"SELECT COUNT(*) FROM {table} WHERE engine='claude'" @@ -868,7 +868,7 @@ def test_async_question_migration_preserves_assets( migrated = HistoryIndexStore(state_dir) with sqlite3.connect(migrated.path) as connection: - assert connection.execute("PRAGMA user_version").fetchone()[0] == 46 + assert connection.execute("PRAGMA user_version").fetchone()[0] == 47 for table in ("history_pages", "history_turn_details"): assert connection.execute( f"SELECT COUNT(*) FROM {table} WHERE engine='claude'" @@ -977,7 +977,7 @@ def test_legacy_migration_rebuilds_all_derived_history_rows( migrated = HistoryIndexStore(state_dir) with sqlite3.connect(migrated.path) as connection: - assert connection.execute("PRAGMA user_version").fetchone()[0] == 46 + assert connection.execute("PRAGMA user_version").fetchone()[0] == 47 for table in ( "history_pages", "history_turn_details", @@ -1024,7 +1024,7 @@ def test_v10_migration_invalidates_changed_projection_rows(tmp_path): migrated = HistoryIndexStore(state_dir) with sqlite3.connect(migrated.path) as connection: - assert connection.execute("PRAGMA user_version").fetchone()[0] == 46 + assert connection.execute("PRAGMA user_version").fetchone()[0] == 47 for table in ( "history_pages", "history_turn_details", "history_image_assets", ): @@ -1075,7 +1075,7 @@ def test_v11_migration_invalidates_claude_pages_and_adds_compact_index( "claude-session", "claude", source, before=None, limit=4, ) is None with sqlite3.connect(migrated.path) as connection: - assert connection.execute("PRAGMA user_version").fetchone()[0] == 46 + assert connection.execute("PRAGMA user_version").fetchone()[0] == 47 tables = { row[0] for row in connection.execute( "SELECT name FROM sqlite_master WHERE type='table'" @@ -1121,7 +1121,7 @@ def test_recent_migration_invalidates_changed_projection_rows( migrated = HistoryIndexStore(state_dir) with sqlite3.connect(migrated.path) as connection: - assert connection.execute("PRAGMA user_version").fetchone()[0] == 46 + assert connection.execute("PRAGMA user_version").fetchone()[0] == 47 for table in ( "history_pages", "history_turn_details", "history_image_assets", ): @@ -1163,7 +1163,7 @@ def test_owner_and_interrupt_alias_migration_invalidates_both_projections( migrated = HistoryIndexStore(state_dir) with sqlite3.connect(migrated.path) as connection: - assert connection.execute("PRAGMA user_version").fetchone()[0] == 46 + assert connection.execute("PRAGMA user_version").fetchone()[0] == 47 for table in ("history_pages", "history_turn_details"): assert connection.execute( f"SELECT COUNT(*) FROM {table} WHERE engine='claude'" @@ -1227,7 +1227,7 @@ def test_recent_summary_migration_rebuilds_pages_but_preserves_source_assets( migrated = HistoryIndexStore(state_dir) with sqlite3.connect(migrated.path) as connection: - assert connection.execute("PRAGMA user_version").fetchone()[0] == 46 + assert connection.execute("PRAGMA user_version").fetchone()[0] == 47 assert connection.execute( "SELECT COUNT(*) FROM history_pages" ).fetchone()[0] == 0 diff --git a/tests/test_release_distribution.py b/tests/test_release_distribution.py index 572a5c58..693499af 100644 --- a/tests/test_release_distribution.py +++ b/tests/test_release_distribution.py @@ -158,6 +158,7 @@ def test_release_bundles_are_deterministic_and_role_scoped( assert f"{prefix}/bin/cc-remote" in members assert f"{prefix}/cc_remote/__main__.py" in members assert f"{prefix}/deploy/install_cli.py" in members + assert f"{prefix}/deploy/cleanup.py" in members assert f"{prefix}/licenses/uv-LICENSE-MIT" in members assert f"{prefix}/cc_remote/protocol.py" in members assert not any("/tests/" in name for name in members) diff --git a/tests/test_wrapper_core_fixes.py b/tests/test_wrapper_core_fixes.py index 0a59b035..eeae3622 100644 --- a/tests/test_wrapper_core_fixes.py +++ b/tests/test_wrapper_core_fixes.py @@ -1918,6 +1918,54 @@ def test_codex_session_settings_restores_last_valid_collaboration_mode( } +@pytest.mark.parametrize("tail_starts_inside_record", [False, True]) +def test_codex_settings_skip_oversized_record_and_read_newer_selection( + monkeypatch, tmp_path, tail_starts_inside_record): + rollout = tmp_path / "rollout-session-1.jsonl" + old = json.dumps({ + "type": "turn_context", "payload": {"model": "gpt-old"}, + }) + "\n" + latest = json.dumps({ + "type": "event_msg", "payload": { + "type": "thread_settings_applied", + "thread_settings": { + "model": "gpt-selected", "reasoning_effort": "xhigh", + "service_tier": "priority", + }, + }, + }) + "\n" + rollout.write_text(old + ("x" * 6000) + "\n" + latest) + monkeypatch.setattr(codex_sessions_module, "MAX_JSONL_RECORD_BYTES", 256) + monkeypatch.setattr( + codex_sessions_module, "_rollout_path", lambda _sid: str(rollout)) + + assert codex_session_settings( + "session-1", max_bytes=2000 if tail_starts_inside_record else 10000, + ) == { + "model": "gpt-selected", "effort": "xhigh", + "service_tier": "priority", + } + + +def test_codex_settings_tail_stops_at_captured_source_size( + monkeypatch, tmp_path): + rollout = tmp_path / "rollout-session-1.jsonl" + first = json.dumps({ + "type": "turn_context", "payload": {"model": "gpt-snapshot"}, + }) + "\n" + appended = json.dumps({ + "type": "turn_context", "payload": {"model": "gpt-after-snapshot"}, + }) + "\n" + rollout.write_text(first + appended) + monkeypatch.setattr( + codex_sessions_module, "_rollout_path", lambda _sid: str(rollout)) + monkeypatch.setattr(codex_sessions_module.os.path, "getsize", lambda _p: len(first)) + + assert codex_session_settings("session-1", max_bytes=len(first)) == { + "model": "gpt-snapshot", + } + + def test_codex_session_settings_restores_applied_update_before_next_turn( monkeypatch, tmp_path): rollout = tmp_path / "rollout-session-1.jsonl" diff --git a/web/playwright.config.ts b/web/playwright.config.ts index 43fc82f4..0d70eb4d 100644 --- a/web/playwright.config.ts +++ b/web/playwright.config.ts @@ -25,7 +25,7 @@ const WEBKIT_GENERAL_EXCLUSIONS = [ export default defineConfig({ testDir: "./tests", - testMatch: ["history-browser.spec.ts", "background-tasks.spec.ts", "themes.spec.ts"], + testMatch: ["history-browser.spec.ts", "background-tasks.spec.ts", "themes.spec.ts", "mobile-sidebar.spec.ts"], fullyParallel: false, workers: 1, retries: process.env.CI ? 2 : 0, diff --git a/web/public/cc-remote-build.json b/web/public/cc-remote-build.json index 5ad04af6..bed7db75 100644 --- a/web/public/cc-remote-build.json +++ b/web/public/cc-remote-build.json @@ -1,4 +1,4 @@ { "version": "4.0.7", - "protocol": 73 + "protocol": 74 } diff --git a/web/scripts/check-bundle-budget.mjs b/web/scripts/check-bundle-budget.mjs index fbdf438e..a8732d51 100644 --- a/web/scripts/check-bundle-budget.mjs +++ b/web/scripts/check-bundle-budget.mjs @@ -33,7 +33,9 @@ const DIST = resolve(import.meta.dirname, "../dist"); // startup JS. Keep entry, compressed-size and request-count limits unchanged. const MAX_ENTRY_BYTES = 537 * 1024; const MAX_INITIAL_BYTES = 938 * 1024; -const MAX_INITIAL_GZIP_BYTES = 280 * 1024; +// Cross-thread provenance, scoped navigation and receipts add <2 KiB gzip. +// The history reader and link UI load on demand; other caps stay unchanged. +const MAX_INITIAL_GZIP_BYTES = 282 * 1024; const MAX_INITIAL_JS_FILES = 4; const html = readFileSync(resolve(DIST, "index.html"), "utf8"); diff --git a/web/src/App.tsx b/web/src/App.tsx index 7df2055c..01209889 100644 --- a/web/src/App.tsx +++ b/web/src/App.tsx @@ -1,3 +1,5 @@ +import type { ServerEvent } from "./protocol"; +import { resolveRelatedSession, relatedSessionTitle } from "./session-messages"; import { lazy, Suspense, @@ -6,7 +8,6 @@ import { useReducer, useRef, useState, - type TouchEvent, } from "react"; import { RelayWs, sessionScopeKey, type EventOwnership } from "./ws"; import type { QueryAcceptanceResult } from "./outbox"; @@ -32,6 +33,10 @@ import { TurnFilePageRequests, type LoadTurnFilePage } from "./turn-file-pages"; import { Icon } from "./icons"; import { ChatView } from "./components/ChatView"; import { Composer } from "./components/Composer"; +import { usePanelWidthPreference } from "./use-panel-width"; +import { SidebarToggle } from "./components/SidebarToggle"; +const SessionMessagePreview = lazy(() => import("./components/SessionMessagePreview").then((m) => ({ default: m.SessionMessagePreview }))); + const BackgroundTaskControl = lazy(() => import("./components/BackgroundTaskControl")); import type { QueuedQueryEditor } from "./components/QueuedQueryDialog"; import { ReconnectBanner } from "./components/ReconnectBanner"; @@ -125,7 +130,6 @@ import { matchesBtwRequest, import type { EngineCapabilities, EngineCapabilityItem, EngineCapabilityKind, WorkArtifactInfo, WorkDashboard } from "./protocol"; import { isMarkdownPath } from "./preview-path"; import { parseGitDiff } from "./diff"; -import { resolveSidebarSwipe } from "./responsive-layout"; import { bumpSessionActivity, mergeSessionActivityState, @@ -312,8 +316,8 @@ interface QueuedQueryEditorState extends QueuedQueryEditor { const MAX_TERMINAL_HISTORY_REPAIR_ATTEMPTS = 2; -// The sidebar is an overlay on mobile (<980px, matches index.css) but a -// persistent grid column on desktop. So auto-close it after picking a session +// The sidebar pushes the page on mobile (<980px, matches index.css) but is a +// persistent column on desktop. So auto-close it after picking a session // ONLY on mobile; on desktop keep it open. const isMobile = () => window.matchMedia("(max-width: 979px)").matches; @@ -334,6 +338,7 @@ function catalogForEngineProfile( } export default function App() { + usePanelWidthPreference(); const initialEngineRef = useRef(normalizeEngine(localStorage.getItem(ENGINE_KEY))); const initialSpacesRef = useRef(readEngineSpaces(localStorage, initialEngineRef.current)); const [engine, setEngine] = useState(initialEngineRef.current); @@ -477,6 +482,10 @@ export default function App() { const stateRef = useRef(state); stateRef.current = state; const wsRef = useRef(null); + const [sessionMessagePreview, setSessionMessagePreview] = useState<{ + scope: string; originSid: string; session: SessionInfo; returnLabel: string; + } | null>(null); + const sessionMessagePreviewEvents = useRef<((event: ServerEvent) => boolean) | null>(null); const goalApiRef = useRef(null); const getGoalApi = useCallback(async () => { const transport = wsRef.current; @@ -651,9 +660,6 @@ export default function App() { const notificationListRequestRef = useRef(null); const notificationOriginRef = useRef(null); const pendingNotificationErrorRef = useRef(null); - const touchStartX = useRef(0); - const touchStartY = useRef(0); - const touchSwipeLocked = useRef(false); const artifactDirtyRef = useRef(false); const setArtifactDirty = useCallback((dirty: boolean) => { artifactDirtyRef.current = dirty; @@ -1619,29 +1625,6 @@ export default function App() { ); }, [authed, machineId, notificationMode]); - // Swipe right -> open sidebar, swipe left -> close (mobile). Interactive - // vertical scrollers opt out so a diagonal scroll never becomes navigation. - const onTouchStart = (e: TouchEvent) => { - const touch = e.touches[0]; - touchStartX.current = touch.clientX; - touchStartY.current = touch.clientY; - touchSwipeLocked.current = e.target instanceof Element - && !!e.target.closest("[data-lock-horizontal-swipe]"); - }; - const onTouchEnd = (e: TouchEvent) => { - const touch = e.changedTouches[0]; - const action = resolveSidebarSwipe( - touchStartX.current, - touchStartY.current, - touch.clientX, - touch.clientY, - window.innerWidth, - touchSwipeLocked.current, - ); - if (action === "open") setSidebarOpen(true); - else if (action === "close") setSidebarOpen(false); - }; - useEffect(() => { try { sessionStorage.setItem( @@ -2208,6 +2191,7 @@ export default function App() { const ws = new RelayWs({ onEvent: (msg, ownership) => { if (!acceptsLifecycle()) return; + if (sessionMessagePreviewEvents.current?.(msg)) return; if (turnFileRequestsRef.current.accept(msg)) return; if (goalApiRef.current?.accept(msg)) return; const settlesContextRequest = !!( @@ -5558,9 +5542,10 @@ export default function App() { -
+
- + {focusedWorkProfile && ( ) : ( <> + {sessionMessagePreview?.scope === activeScopeKey + && sessionMessagePreview.originSid === focusedSid + && state.sessions.some((s) => s.session_id === sessionMessagePreview.session.session_id) + && setSessionMessagePreview(null)} + api={{ receive: sessionMessagePreviewEvents, + history: (sid, before, cwd) => wsRef.current?.sendGetHistory(sid, before, 12, cwd) ?? false, + detail: (sid, turnId, revision, before) => wsRef.current?.sendGetTurnDetail(sid, turnId, revision, before) ?? false, + }} />} { + const session = resolveRelatedSession(nativeId, focusedSid, state.sessions); + return { title: relatedSessionTitle(session, nativeId), available: !!session }; + }} + onOpenSession={focusedEngine === "codex" && space === "code" ? (nativeId) => { + const session = resolveRelatedSession(nativeId, focusedSid, state.sessions); + if (!session || !focusedSid) return; + setSessionMessagePreview({ scope: activeScopeKey, originSid: focusedSid, + session, returnLabel: relatedSessionTitle(focusedSession, focusedSid) }); + } : undefined} loading={!!rt.loading || historyView.recovering} surface={space} engine={focusedEngine} forkingPointId={forkingPointId} diff --git a/web/src/cache.ts b/web/src/cache.ts index 756172a8..bc0e9ade 100644 --- a/web/src/cache.ts +++ b/web/src/cache.ts @@ -68,7 +68,8 @@ const SCHEMA = 1; // v25 reprojects native async questions instead of preserving plain-answer shells. // v26 discards summaries where those questions hid ordinary unphased replies. // v27 removes Claude turns split by native isMeta recovery prompts. -const CACHE_VER = 27; +// v28 discards projections missing native cross-thread message provenance. +const CACHE_VER = 28; const MAX_CACHE_SESSIONS = 64; const MAX_CACHE_TURNS = 100; const MAX_CACHE_BYTES = 2 * 1024 * 1024; diff --git a/web/src/components/ChatView.tsx b/web/src/components/ChatView.tsx index 0d132ea3..a5894341 100644 --- a/web/src/components/ChatView.tsx +++ b/web/src/components/ChatView.tsx @@ -88,6 +88,7 @@ import { mergeDetailWithLiveTail } from "../history-merge"; import { presentAsyncQuestionReplies } from "../async-question-presentation"; import type { QueryAcceptanceResult } from "../outbox"; +const SessionMessageLinks = lazy(() => import("./SessionMessageLinks")); const TurnUsageIndicator = lazy(() => import("./TurnUsageIndicator").then(m => ({ default: m.TurnUsageIndicator }))); const AsyncQuestionCard = lazy(() => import("./AsyncQuestionCard")); const AsyncQuestionHost = lazy(() => import("./AsyncQuestionDialog")); @@ -326,9 +327,12 @@ export function ChatView({ sid, turnUsage, turns: incomingTurns, engine = "claud onTextSelectionGuardChange, externalPlanProgress, onOpenAgent, + sessionLink, onOpenSession, activeTurnId = null, ambiguousActiveTurnIds = [], surface = "code" }: { + sessionLink?: (nativeId: string) => { title: string; available: boolean }; + onOpenSession?: (nativeId: string) => void; turnUsage?: TurnUsageReadings; sid: string | null; turns: Turn[]; @@ -2928,6 +2932,10 @@ export function ChatView({ sid, turnUsage, turns: incomingTurns, engine = "claud {(t.prompt || (t.images && t.images.length) || (t.imageRefs && t.imageRefs.length) || (t.files && t.files.length)) && (
{t.prompt &&
+ {engine === "codex" && t.sourceThreadId && + { pauseOutputFollow(); onOpenSession(id); } : undefined} /> + } {t.timedTask && } {supplemental.replies.has(t.id) ?
@@ -3002,6 +3010,11 @@ export function ChatView({ sid, turnUsage, turns: incomingTurns, engine = "claud
)} {showProcessTimeline && renderProcess()} + {engine === "codex" && (!!t.sessionMessages?.length || t.blocks.some((b) => + b.kind === "tool" && b.tool.endsWith("send_message_to_thread"))) && + { pauseOutputFollow(); onOpenSession(id); } : undefined} /> + } {modelNotices.map((notice) =>
{notice.summary} diff --git a/web/src/components/PanelResizer.tsx b/web/src/components/PanelResizer.tsx index 09351ad3..6ac6ed25 100644 --- a/web/src/components/PanelResizer.tsx +++ b/web/src/components/PanelResizer.tsx @@ -7,11 +7,7 @@ import { } from "react"; import { clampPanelWidth } from "../responsive-layout"; - -// Keep the historical key so an existing preview-panel preference also -// applies to BTW and Agent detail panels. -const PANEL_WIDTH_KEY = "cc_remote_artifact_panel_width"; -const DESKTOP_PANEL_QUERY = "(min-width: 981px)"; +import { DESKTOP_PANEL_QUERY, PANEL_WIDTH_KEY } from "../use-panel-width"; export function PanelResizer({ ariaLabel }: { ariaLabel: string }) { const handleRef = useRef(null); @@ -33,24 +29,10 @@ export function PanelResizer({ ariaLabel }: { ariaLabel: string }) { return width; }, []); - useEffect(() => { - if (!window.matchMedia(DESKTOP_PANEL_QUERY).matches) return; - const saved = Number.parseFloat( - localStorage.getItem(PANEL_WIDTH_KEY) || "", - ); - if (Number.isFinite(saved)) applyPanelWidth(saved); - const fitPanel = () => { - if (!window.matchMedia(DESKTOP_PANEL_QUERY).matches) return; - const current = panelElement()?.getBoundingClientRect().width; - if (current) applyPanelWidth(current); - }; - window.addEventListener("resize", fitPanel); - return () => { - window.removeEventListener("resize", fitPanel); - resizeRef.current = null; - document.documentElement.classList.remove("panel-resizing"); - }; - }, [applyPanelWidth, panelElement]); + useEffect(() => () => { + resizeRef.current = null; + document.documentElement.classList.remove("panel-resizing"); + }, []); const startResize = (event: ReactPointerEvent) => { const panel = panelElement(); diff --git a/web/src/components/SessionMessageLinks.tsx b/web/src/components/SessionMessageLinks.tsx new file mode 100644 index 00000000..650caebf --- /dev/null +++ b/web/src/components/SessionMessageLinks.tsx @@ -0,0 +1,32 @@ +import type { Turn } from "../domain/conversation"; +import { outgoingSessionMessages } from "../session-messages"; +import { Icon } from "../icons"; +import "./session-messages.css"; + +export default function SessionMessageLinks({ turn, source, resolve, onOpen }: { + turn: Turn; + source?: boolean; + resolve?: (nativeId: string) => { title: string; available: boolean }; + onOpen?: (nativeId: string) => void; +}) { + const link = (id: string, label: string, key: string) => { + const target = resolve?.(id); + const enabled = !!onOpen && !!target?.available; + return ; + }; + if (source) return turn.sourceThreadId + ? link(turn.sourceThreadId, "来自会话 ·", "source") : null; + return outgoingSessionMessages(turn).map((receipt) => link( + receipt.threadId, + receipt.status === "sent" ? "已发送给" : receipt.status === "failed" + ? "未能发送给" : turn.done ? "发送状态未确认 ·" : "正在发送给", + receipt.itemId, + )); +} diff --git a/web/src/components/SessionMessagePreview.tsx b/web/src/components/SessionMessagePreview.tsx new file mode 100644 index 00000000..333671af --- /dev/null +++ b/web/src/components/SessionMessagePreview.tsx @@ -0,0 +1,109 @@ +import { useCallback, useEffect, useReducer, useRef } from "react"; +import type { ServerEvent, SessionInfo } from "../protocol"; +import { SessionMessageReader } from "../session-message-reader"; +import { relatedSessionTitle } from "../session-messages"; +import { ChatView } from "./ChatView"; +import { Icon } from "../icons"; +import "./session-messages.css"; + +export interface SessionMessagePreviewApi { + history: (sid: string, before: string | null, cwd?: string | null) => boolean; + detail: (sid: string, turnId: string, revision: string, before: string | null) => boolean; + receive: { current: ((event: ServerEvent) => boolean) | null }; +} + +export function SessionMessagePreview({ session, returnSid, returnLabel, api, onClose }: { + session: SessionInfo; + returnLabel: string; + returnSid: string; + api: SessionMessagePreviewApi; + onClose: () => void; +}) { + const dialog = useRef(null); + const readerRef = useRef(null); + if (!readerRef.current) readerRef.current = new SessionMessageReader(session.session_id); + const reader = readerRef.current; + const [, repaint] = useReducer((n: number) => n + 1, 0); + const timers = useRef(new Set>()); + const apiRef = useRef(api); + apiRef.current = api; + + const armTimeout = useCallback(() => { + const timer = setTimeout(() => { + timers.current.delete(timer); + reader.fail("会话读取超时,请重试。"); + repaint(); + }, 15_000); + timers.current.add(timer); + }, [reader]); + const clearTimers = useCallback(() => { + for (const timer of timers.current) clearTimeout(timer); + timers.current.clear(); + }, []); + const load = useCallback((page: number) => { + const request = reader.requestPage(page); + if (!request) return; + clearTimers(); + if (!apiRef.current.history(reader.sid, request.before, session.cwd)) reader.fail(); + else armTimeout(); + repaint(); + }, [reader, session.cwd, clearTimers, armTimeout]); + useEffect(() => { + const node = dialog.current; + const previousFocus = document.activeElement instanceof HTMLElement + ? document.activeElement : null; + node?.showModal(); + const receive = (event: ServerEvent) => { + if (!reader.accept(event)) return false; + clearTimers(); + repaint(); + return event.type !== "history_invalidated"; + }; + const receiveRef = apiRef.current.receive; + receiveRef.current = receive; + load(0); + return () => { + clearTimers(); + // Also cancel the private pending state: StrictMode may immediately + // set this effect up again with the same reader instance. + reader.fail(); + if (receiveRef.current === receive) receiveRef.current = null; + node?.close(); + previousFocus?.focus({ preventScroll: true }); + }; + }, [reader, load, clearTimers]); + + return event.stopPropagation()} onCancel={(event) => { event.preventDefault(); event.stopPropagation(); onClose(); }}> +
+ + {relatedSessionTitle(session, session.session_id)} +
+ {reader.error &&
+ {reader.error} +
} + id === returnSid.slice(returnSid.indexOf("@") + 1) + ? { title: returnLabel, available: true } + : { title: id.slice(0, 8), available: false }} + onOpenSession={onClose} + historyRevision={reader.revision} historyGeneration={reader.generation} + onLoadDetail={(turnId, before) => { + if (!reader.requestDetail(turnId, before ?? null)) return false; + if (!apiRef.current.detail(reader.sid, turnId, reader.revision!, before ?? null)) reader.fail(); + else armTimeout(); + repaint(); + return true; + }} /> +
+ + 只读查看 + +
+
; +} diff --git a/web/src/components/SessionsSidebar.tsx b/web/src/components/SessionsSidebar.tsx index 6a860167..a6b6b906 100644 --- a/web/src/components/SessionsSidebar.tsx +++ b/web/src/components/SessionsSidebar.tsx @@ -19,6 +19,7 @@ import { codexProfilePresentation } from "../codex-profile-presentation"; import { newWorkProfileForSidebarFilter } from "../work-profile-selection"; import { manualUnreadKey } from "../manual-unread"; import { useManualUnread } from "../use-manual-unread"; +import { useMobileSidebar } from "../use-mobile-sidebar"; const SessionCardMenu = lazy(() => import("./SessionCardMenu").then(module => ({ default: module.SessionCardMenu }))); const TimedTaskIndicator = lazy(() => import("./TimedTaskIndicator").then(module => ({ default: module.TimedTaskIndicator }))); @@ -43,6 +44,7 @@ interface Props { onNew: (profileId?: string) => void; onNewInDir: (cwd: string) => void; onClose: () => void; + onOpenChange?: (open: boolean) => void; onRename: (id: string, title: string) => void; onArchive: (id: string, archived: boolean) => void; onPin: (session: SessionInfo, pinned: boolean) => void; @@ -77,8 +79,9 @@ export function SessionsSidebar({ open, engine, space, profileScopeKey, machineId, claudeProfiles = [], defaultClaudeProfileId, codexProfiles = [], defaultCodexProfileId, onSpaceChange, sessions, liveStates, - completionBadges, activeSessionId, onSelect, onNew, onNewInDir, onClose, + completionBadges, activeSessionId, onSelect, onNew, onNewInDir, onClose, onOpenChange, onRename, onArchive, onPin, onDelete, onForkWorktree, onMigrate }: Props) { + const sidebarRef = useMobileSidebar(open, onOpenChange); const manualUnread = useManualUnread(); const [q, setQ] = useState(""); const [menuCardId, setMenuCardId] = useState(null); @@ -457,9 +460,8 @@ export function SessionsSidebar({ open, engine, space, return ( <>
-