diff --git a/app/acp_client.py b/app/acp_client.py new file mode 100644 index 0000000..cd6f6f9 --- /dev/null +++ b/app/acp_client.py @@ -0,0 +1,227 @@ +"""Bounded ACP channel for one WorkBuddy agent turn over streamable HTTP. + +WorkBuddy's console web app drives an agent turn with the Agent Client Protocol +transport: one long-lived ``GET`` channel carrying server-sent events, plus one +short JSON-RPC ``POST`` per method, both tied together by the ``Acp-Connection-Id`` +response header. Only enough of the protocol is implemented here to run a single +turn and report a bounded outcome -- turn results are never interpreted, because +completion is read back from the console API instead. + +``http.client`` is used deliberately rather than ``httpx``: the channel has to +stay open for the whole turn while console polling proceeds on separate +connections, and its bytes are drained opportunistically so a read can never +block the caller. That is the one place this module deviates from the ``httpx`` +style of the other automation modules. +""" +import json +import select +import time +from http.client import HTTPConnection, HTTPSConnection, HTTPSConnection as _HTTPS +from urllib.parse import urlsplit + +from .credits import BROWSER_UA + +PROTOCOL_VERSION = 1 +# A turn must never depend on client-provided files, terminals or permissions: +# declaring them would make the sandbox wait for callbacks it will not get. +CLIENT_CAPABILITIES = {"fs": {"readTextFile": False, "writeTextFile": False}, + "terminal": False} +MAX_CHANNEL_BYTES = 256 * 1024 +CONNECT_TIMEOUT = 10.0 +POST_TIMEOUT = 30.0 +DRAIN_STEP = 0.2 +READ_CHUNK = 8192 +_METHODS = ("initialize", "session/load", "session/prompt") + + +class AcpError(ValueError): + """One bounded failure; never carries upstream text or headers.""" + + def __init__(self, kind, http_status=None, code=None): + super().__init__("ACP turn was not confirmed") + self.diagnostics = {"error_kind": kind, "http_status": http_status, "code": code} + + +def _target(link): + parts = urlsplit(link if isinstance(link, str) else "") + if parts.scheme not in ("http", "https") or not parts.netloc or not parts.path: + raise AcpError("protocol") + # Keep the query string: dropping it would silently change the sandbox endpoint. + return parts, parts.path + (("?" + parts.query) if parts.query else "") + + +class AcpChannel: + """One SSE channel plus the JSON-RPC requests bound to its connection id.""" + + def __init__(self, link, token): + self.parts, self.path = _target(link) + token = token if isinstance(token, str) else "" + self.token = token + self.connection_id = None + self._connection = None + self._response = None + self._buffer = bytearray() + self._line = bytearray() + self._closed = False + + def open(self): + self._connection = self._connect() + try: + self._connection.putrequest("GET", self.path) + self._connection.putheader("Accept", "text/event-stream") + self._connection.putheader("Accept-Encoding", "identity") + self._connection.putheader("Authorization", "Bearer " + self.token) + self._connection.putheader("User-Agent", BROWSER_UA) + self._connection.endheaders() + self._response = self._connection.getresponse() + except (OSError, ValueError): + self.close() + raise AcpError("network") from None + if self._response.status != 200: + status = self._response.status + self.close() + raise AcpError("http", status) + connection_id = self._response.getheader("Acp-Connection-Id") + if not isinstance(connection_id, str) or not connection_id.strip(): + self.close() + raise AcpError("protocol") + self.connection_id = connection_id.strip() + return self + + def post(self, method, params, request_id): + """Send one JSON-RPC request on a fresh connection; its reply rides the channel.""" + if method not in _METHODS: + raise AcpError("protocol") + if self._closed or not self.connection_id: + raise AcpError("protocol") + body = json.dumps({"jsonrpc": "2.0", "id": request_id, "method": method, + "params": params}, allow_nan=False).encode() + connection = self._connect() + status = None + try: + connection.putrequest("POST", self.path) + connection.putheader("Content-Type", "application/json") + connection.putheader("Accept", "application/json, text/event-stream") + connection.putheader("Acp-Connection-Id", self.connection_id) + connection.putheader("Authorization", "Bearer " + self.token) + connection.putheader("User-Agent", BROWSER_UA) + connection.putheader("Content-Length", str(len(body))) + connection.endheaders(message_body=body) + response = connection.getresponse() + response.read() + status = response.status + except (OSError, ValueError): + raise AcpError("network") from None + finally: + self._drop(connection) + if status not in (200, 202): + raise AcpError("http", status) + + def drain(self, seconds): + """Read available channel bytes without blocking; return harvested usage updates.""" + updates = [] + if self._closed or self._response is None: + return updates + deadline = time.monotonic() + max(0.0, float(seconds)) + while len(self._buffer) <= MAX_CHANNEL_BYTES: + remaining = deadline - time.monotonic() + if remaining <= 0: + break + try: + readable, _, _ = select.select([self._response.fp], [], [], + min(DRAIN_STEP, remaining)) + except (OSError, ValueError, TypeError): + break + if not readable: + continue + try: + chunk = self._response.read1(READ_CHUNK) + except (OSError, ValueError): + break + if not chunk: + break + self._buffer.extend(chunk) + updates.extend(self._events()) + return updates + + def _events(self): + """Yield ACP ``usage_update`` notifications from the drained buffer.""" + found = [] + while True: + break_at = self._buffer.find(b"\n") + if break_at < 0: + break + line = bytes(self._buffer[:break_at]) + del self._buffer[:break_at + 1] + line = line.rstrip(b"\r") + if not line.startswith(b"data:"): + continue + payload = line[5:].strip() + if not payload or payload == b"[DONE]": + continue + try: + message = json.loads(payload) + except (ValueError, UnicodeError, RecursionError): + continue + if not isinstance(message, dict): + continue + update = message.get("params", {}) + if not isinstance(update, dict): + continue + update = update.get("update") + if (message.get("method") == "session/update" and isinstance(update, dict) + and update.get("sessionUpdate") == "usage_update"): + found.append(update) + return found + + def close(self): + self._closed = True + response, connection = self._response, self._connection + self._response = self._connection = None + if response is not None: + try: + response.close() + except (OSError, ValueError): + pass + if connection is not None: + self._drop(connection) + + def _connect(self): + host = self.parts.hostname + if not host: + raise AcpError("protocol") + try: + if self.parts.scheme == "https": + return _HTTPS(host, self.parts.port, timeout=CONNECT_TIMEOUT) + return HTTPConnection(host, self.parts.port, timeout=CONNECT_TIMEOUT) + except (OSError, ValueError): + raise AcpError("network") from None + + @staticmethod + def _drop(connection): + try: + connection.close() + except (OSError, ValueError): + pass + + +def run_turn(link, token, session_id, cwd, prompt, *, on_prompt=None): + """Open the channel and drive one turn; return it open for console polling. + + ``on_prompt`` runs immediately before the prompt is posted, so the caller can + record the turn as committed at the last moment where that is still knowable. + """ + channel = AcpChannel(link, token).open() + try: + channel.post("initialize", {"protocolVersion": PROTOCOL_VERSION, + "clientCapabilities": CLIENT_CAPABILITIES}, 1) + channel.post("session/load", {"sessionId": session_id, "cwd": cwd, + "mcpServers": []}, 2) + if on_prompt is not None: + on_prompt() + channel.post("session/prompt", {"sessionId": session_id, + "prompt": [{"type": "text", "text": prompt}]}, 3) + except Exception: + channel.close() + raise + return channel diff --git a/app/admin_api.py b/app/admin_api.py index 2bc0618..a64d242 100644 --- a/app/admin_api.py +++ b/app/admin_api.py @@ -60,6 +60,7 @@ def _public_credential(item): "fail_until", "cooldown_until", "cooldown_remaining", "last_failure_at", "catalog_ready", "bindings", "auto_checkin", "auto_travel", "travel_supported", "checkin", "travel", "trial_supported", "trial", + "daily_chat_supported", "auto_daily_chat", "daily_chat", "token_expired", "token_expires_at", "last_refresh_time", "sessions", "sticky_sessions", "last_error_code"} result = {key: value for key, value in item.items() if key in fields} identity = item.get("account_key") or item.get("id") @@ -368,8 +369,8 @@ def preview(): @route("PATCH", "/admin/credentials/{id}") async def credentials_patch(request): data = await _body(request) - if set(data) not in ({"enabled"}, {"auto_checkin"}, {"auto_travel"}) or any(type(value) is not bool for value in data.values()): - raise ValueError("仅接受一个布尔字段:enabled、auto_checkin 或 auto_travel") + if set(data) not in ({"enabled"}, {"auto_checkin"}, {"auto_travel"}, {"auto_daily_chat"}) or any(type(value) is not bool for value in data.values()): + raise ValueError("仅接受一个布尔字段:enabled、auto_checkin、auto_travel 或 auto_daily_chat") field, value = next(iter(data.items())) identity = request.path_params["id"] def apply(): @@ -381,6 +382,8 @@ def apply(): gateway.admin_set_credential_enabled(identity, value) elif field == "auto_checkin": gateway.admin_set_auto_checkin(identity, value) + elif field == "auto_daily_chat": + gateway.admin_set_auto_daily_chat(identity, value) else: gateway.admin_set_auto_travel(identity, value) event("credential." + field, {"credential": identity, field: value}) diff --git a/app/control_store.py b/app/control_store.py index b18c619..c4a6236 100644 --- a/app/control_store.py +++ b/app/control_store.py @@ -3,6 +3,7 @@ import copy import json +from datetime import date from pathlib import Path import sqlite3 import threading @@ -23,6 +24,23 @@ def _identifier(value, label): return value +def _day(value, label="日期"): + """Accept only a YYYY-MM-DD local calendar day; never a timestamp or free text. + + The shape check alone admits 2026-02-31 and 2026-99-99, and a day that never + happened would become a distinct key that the once-per-day guard could be + satisfied against instead of the real day. + """ + if (not isinstance(value, str) or len(value) != 10 or value[4] != "-" or value[7] != "-" + or not value.replace("-", "").isdigit()): + raise ValueError(f"{label} 无效") + try: + date.fromisoformat(value) + except ValueError: + raise ValueError(f"{label} 无效") from None + return value + + def buddy_claim_reserved(record): """Only pre-claim reservations may expire; a pending send can outlive its process.""" return bool(record and (record["claimed"] or record["stage"] not in {"reserved", "agree", "buddy_agree"})) @@ -107,6 +125,12 @@ def __init__(self, path): "accept_started INTEGER NOT NULL DEFAULT 0, chat_started INTEGER NOT NULL DEFAULT 0, " "completed INTEGER NOT NULL DEFAULT 0, model TEXT, chat_state TEXT NOT NULL DEFAULT 'pending', " "total_tokens INTEGER, updated_at REAL NOT NULL)") + self._db.execute("CREATE TABLE IF NOT EXISTS daily_chats (" + "account_key TEXT NOT NULL, day TEXT NOT NULL, attempt_id TEXT NOT NULL, " + "phase TEXT NOT NULL, conversation_id TEXT, sandbox_status TEXT, " + "usage_before REAL, usage_after REAL, acp_usage REAL, " + "reserved_at REAL NOT NULL, confirmed_at REAL, updated_at REAL NOT NULL, " + "PRIMARY KEY (account_key, day))") self._db.execute("CREATE TABLE IF NOT EXISTS runtime_state (name TEXT PRIMARY KEY, payload TEXT NOT NULL)") self._db.execute("CREATE TABLE IF NOT EXISTS state_imports (name TEXT PRIMARY KEY, imported INTEGER NOT NULL, migrated_at REAL NOT NULL)") self._db.execute("CREATE TABLE IF NOT EXISTS gateway_secrets (name TEXT PRIMARY KEY, value TEXT NOT NULL, announced INTEGER NOT NULL DEFAULT 0 CHECK(announced IN (0,1)))") @@ -134,9 +158,10 @@ def _load(self): validate_model(source, rule, data["models"], legacy_scopes=True) # Preserve legacy scope intersections. for identity, metadata in data["credentials"].items(): _identifier(identity, "账号指纹") - if (not isinstance(metadata, dict) or set(metadata) - {"enabled", "label", "auto_checkin", "auto_travel"} + if (not isinstance(metadata, dict) or set(metadata) - {"enabled", "label", "auto_checkin", "auto_travel", "auto_daily_chat"} or type(metadata.get("enabled")) is not bool - or any(key in metadata and type(metadata[key]) is not bool for key in ("auto_checkin", "auto_travel"))): + or any(key in metadata and type(metadata[key]) is not bool for key in + ("auto_checkin", "auto_travel", "auto_daily_chat"))): raise ValueError("管理数据库凭证元数据无效") return {"revision": row[0], **data} @@ -244,6 +269,13 @@ def set_auto_travel(self, account_key, enabled): return self._update(None, lambda state: state["credentials"].setdefault( account_key, {"enabled": True}).update(auto_travel=enabled)) + def set_auto_daily_chat(self, account_key, enabled): + _identifier(account_key, "账号指纹") + if type(enabled) is not bool: + raise ValueError("auto_daily_chat 必须为布尔值") + return self._update(None, lambda state: state["credentials"].setdefault( + account_key, {"enabled": True}).update(auto_daily_chat=enabled)) + def has_buddy_consent(self, identity, revision): _identifier(identity, "账号指纹") with self._lock: @@ -425,6 +457,121 @@ def buddy_task_checkpoint(self, identity, *, completed=False, chat_state=None, t (int(completed), chat_state, total_tokens, time.time(), identity)) + # -- International daily-activity turns ------------------------------------ + # One row per account per local day. A turn may consume credits on the account + # and is never replayed automatically, so only an unsent reservation is + # cancellable: a sent-but-unconfirmed one stays pending until the next day. + + _DAILY_CHAT_PHASES = ("reserved", "sent", "confirmed", "cancelled", "reconciled") + + _DAILY_CHAT_TRANSITIONS = {"sent": ("reserved",), + "confirmed": ("sent", "confirmed", "reconciled"), + "cancelled": ("reserved", "sent", "cancelled"), + "reconciled": ("reserved", "sent", "confirmed", "reconciled")} + _DAILY_CHAT_UPDATE_SQL = { + "sent": "UPDATE daily_chats SET phase=CASE WHEN phase='reconciled' AND ?='confirmed' THEN phase ELSE ? END, " + "confirmed_at=CASE WHEN ? IN ('confirmed','reconciled') THEN COALESCE(confirmed_at,?) ELSE confirmed_at END, " + "updated_at=? WHERE account_key=? AND day=? AND attempt_id=? AND phase IN (?)", + "confirmed": "UPDATE daily_chats SET phase=CASE WHEN phase='reconciled' AND ?='confirmed' THEN phase ELSE ? END, " + "confirmed_at=CASE WHEN ? IN ('confirmed','reconciled') THEN COALESCE(confirmed_at,?) ELSE confirmed_at END, " + "updated_at=? WHERE account_key=? AND day=? AND attempt_id=? AND phase IN (?,?,?)", + "cancelled": "UPDATE daily_chats SET phase=CASE WHEN phase='reconciled' AND ?='confirmed' THEN phase ELSE ? END, " + "confirmed_at=CASE WHEN ? IN ('confirmed','reconciled') THEN COALESCE(confirmed_at,?) ELSE confirmed_at END, " + "updated_at=? WHERE account_key=? AND day=? AND attempt_id=? AND phase IN (?,?,?)", + "reconciled": "UPDATE daily_chats SET phase=CASE WHEN phase='reconciled' AND ?='confirmed' THEN phase ELSE ? END, " + "confirmed_at=CASE WHEN ? IN ('confirmed','reconciled') THEN COALESCE(confirmed_at,?) ELSE confirmed_at END, " + "updated_at=? WHERE account_key=? AND day=? AND attempt_id=? AND phase IN (?,?,?,?)", + } + def daily_chat_record(self, identity, day): + _identifier(identity, "账号指纹") + _day(day) + with self._lock: + cursor = self._db.execute("SELECT * FROM daily_chats WHERE account_key=? AND day=?", + (identity, day)) + row = cursor.fetchone() + return dict(zip((column[0] for column in cursor.description), row)) if row else None + + def reserve_daily_chat(self, identity, day): + """Reserve the single daily turn; a completed or already-sent day is never reopened. + + A ``reserved`` row is retryable because it means nothing was written upstream + yet (the caller only marks ``sent`` once the turn can really run), so a crash + or a failed sandbox provisioning does not cost the day. + """ + _identifier(identity, "账号指纹") + _day(day) + with self._lock: + self._db.execute("BEGIN IMMEDIATE") + try: + previous = self.daily_chat_record(identity, day) + if previous and previous["phase"] not in {"reserved", "cancelled", "reconciled"}: + self._db.execute("COMMIT") + return None + attempt = uuid.uuid4().hex + self._db.execute( + "INSERT INTO daily_chats (account_key,day,attempt_id,phase,reserved_at,updated_at) " + "VALUES(?,?,?,'reserved',?,?) ON CONFLICT(account_key,day) DO UPDATE SET " + "attempt_id=excluded.attempt_id,phase='reserved',conversation_id=NULL," + "sandbox_status=NULL,usage_before=NULL,usage_after=NULL,acp_usage=NULL," + "reserved_at=excluded.reserved_at,confirmed_at=NULL,updated_at=excluded.updated_at", + (identity, day, attempt, time.time(), time.time())) + self._db.execute("COMMIT") + return {"attempt_id": attempt} + except Exception: + if self._db.in_transaction: + self._db.execute("ROLLBACK") + raise + + def transition_daily_chat(self, identity, day, attempt, phase): + """Late receipts cannot reopen or replace an already reconciled reservation.""" + _identifier(identity, "账号指纹") + _day(day) + _identifier(attempt, "尝试 ID") + if phase not in self._DAILY_CHAT_PHASES: + raise ValueError("打卡阶段无效") + expected = self._DAILY_CHAT_TRANSITIONS[phase] + with self._lock: + updated = self._db.execute( + self._DAILY_CHAT_UPDATE_SQL[phase], + (phase, phase, phase, time.time(), time.time(), identity, day, attempt, *expected)) + if updated.rowcount != 1: + raise ValueError("打卡预留已变化") + + def daily_chat_checkpoint(self, identity, day, *, sandbox_status=None, usage_before=None, + usage_after=None, acp_usage=None, conversation_id=None): + _identifier(identity, "账号指纹") + _day(day) + if sandbox_status is not None: + _identifier(sandbox_status, "会话状态") + if conversation_id is not None: + _identifier(conversation_id, "会话 ID") + for name, value in (("usage_before", usage_before), ("usage_after", usage_after), + ("acp_usage", acp_usage)): + if value is not None and (type(value) not in (int, float) or not 0 <= value <= 10**9): + raise ValueError(f"{name} 无效") + with self._lock: + self._db.execute( + "UPDATE daily_chats SET sandbox_status=COALESCE(?,sandbox_status), " + "conversation_id=COALESCE(?,conversation_id), usage_before=COALESCE(?,usage_before), " + "usage_after=COALESCE(?,usage_after), acp_usage=COALESCE(?,acp_usage), updated_at=? " + "WHERE account_key=? AND day=?", + (sandbox_status, conversation_id, usage_before, usage_after, acp_usage, + time.time(), identity, day)) + + def daily_chat_done(self, identity, day): + record = self.daily_chat_record(identity, day) + return bool(record and record["phase"] in {"confirmed", "reconciled"}) + + def prune_daily_chats(self, keep_days=30): + """Drop old rows; records are diagnostic, never a source of truth for routing.""" + keep_days = int(keep_days) + if not 1 <= keep_days <= 3650: + raise ValueError("保留天数无效") + cutoff = time.strftime("%Y-%m-%d", time.localtime(time.time() - keep_days * 86400)) + with self._lock: + self._db.execute("DELETE FROM daily_chats WHERE day= due_at(identity, day, window=window) + + +def _console_headers(token, uid="", domain=""): + """Console web headers: two credential headers only, no CLI X-IDE-* fingerprint.""" + return { + "accept": "application/json, text/plain, */*", + "content-type": "application/json", + "x-client-platform": "web", + "origin": HOST, + "referer": HOST + "/app", + "authorization": "Bearer " + (token if isinstance(token, str) else ""), + "x-user-id": str(uid or ""), + "x-domain": str(domain or ""), + "user-agent": BROWSER_UA, + } + + +def _read(client, method, url, headers, *, body=None): + """One request without redirects; the body is parsed here and never disclosed.""" + with client.stream(method, url, headers=headers, content=body, timeout=REQUEST_TIMEOUT, + follow_redirects=False) as response: + status = response.status_code + raw = bytearray() + started = time.monotonic() + for chunk in response.iter_bytes(): + if time.monotonic() - started > TURN_SECONDS or len(raw) + len(chunk) > MAX_BODY_BYTES: + raise ChatFailure("protocol_error", status) + raw.extend(chunk) + try: + payload = json.loads(bytes(raw).decode("utf-8")) + except (ValueError, UnicodeError, RecursionError): + raise ChatFailure("protocol_error" if status == 200 else "http_error", status) from None + return status, payload + + +def _check(status, payload): + if not isinstance(payload, dict): + raise ChatFailure("protocol_error" if status == 200 else "http_error", status) + code = payload.get("code") + code = code if type(code) is int and -(2**31) <= code < 2**31 else None + if status in (401, 403): + raise ChatFailure("auth_error", status, code) + if status not in (200, 201, 202): + raise ChatFailure("http_error", status, code) + if code not in (0, None): + raise ChatFailure("rejected", status, code) + return payload + + +def _post(client, url, headers, body): + status, payload = _read(client, "POST", url, headers, + body=json.dumps(body, allow_nan=False).encode()) + return _check(status, payload) + + +def _get(client, url, headers): + status, payload = _read(client, "GET", url, headers) + return _check(status, payload) + + +def _text(value, limit=200): + return value if isinstance(value, str) and 0 < len(value) <= limit else None + + +def create_conversation(client, headers): + """Queue one conversation; the returned session stays CREATING until attached.""" + payload = _post(client, HOST + CONVERSATIONS_PATH, headers, + {"prompt": PROMPT, "model": MODEL, + "conversationOrigin": CONVERSATION_ORIGIN, "plugins": PLUGINS}) + data = payload.get("data") + identity = _text(data.get("id"), 128) if isinstance(data, dict) else None + if identity is None: + raise ChatFailure("protocol_error") + return identity + + +def sandbox_of(client, headers, conversation_id): + """Read the sandbox endpoint and session identity the ACP turn needs.""" + payload = _get(client, HOST + CONVERSATIONS_PATH + conversation_id + SESSION_SUFFIX, headers) + data = payload.get("data") + if not isinstance(data, dict): + raise ChatFailure("protocol_error") + link = _text(data.get("link") or data.get("endpoint"), 2048) + token = _text(data.get("token"), 2048) + session_id = _text(data.get("sessionId") or data.get("session_id") or conversation_id, 256) + cwd = _text(data.get("cwd"), 512) or DEFAULT_CWD + if link is None or token is None: + # The sandbox is still being provisioned; the caller may poll again. + raise ChatFailure("sandbox_pending") + return {"link": link, "token": token, "session_id": session_id, "cwd": cwd} + + +def _await_sandbox(client, headers, conversation_id): + """Poll for a ready sandbox, bounded by the same wall clock as the turn.""" + deadline = time.monotonic() + SANDBOX_SECONDS + last = None + while True: + try: + return sandbox_of(client, headers, conversation_id) + except ChatFailure as error: + last = error + if error.diagnostics.get("error_kind") != "sandbox_pending": + raise + if time.monotonic() >= deadline: + raise ChatFailure("sandbox_timeout", None, None) from None + time.sleep(SANDBOX_POLL_SECONDS) + del last + + +def status_of(client, headers, conversation_id): + """Console business status: the only reliable completion signal for a turn.""" + payload = _get(client, HOST + CONVERSATIONS_PATH + conversation_id, headers) + data = payload.get("data") + return (_text(data.get("status"), 64) or "") if isinstance(data, dict) else "" + + +def _usage_credits(token, uid, domain): + """Best-effort same-day credit total; upstream usage lags by minutes, so None is normal.""" + try: + snapshot = credits.fetch_request_usage(token, days=1, uid=uid, domain=domain) + except Exception: + return None + rows = snapshot.get("by_day", {}).get(today(), {}) + return round(sum(rows.values()), 4) if isinstance(rows, dict) and rows else None + + +def perform(token, profile, *, uid="", domain="", can_write=lambda: True, + store=None, identity=None): + """Run one daily-activity turn under durable reservation and generation checks.""" + result = {"ok": False, "state": "unknown", "day": today(), + "phase": "verify", "conversation_id": None, "sandbox_status": None, + "usage_before": None, "usage_after": None, "acp_usage": None} + day, attempt_id = result["day"], None + conversation_id, sandbox_status, acp_usage = None, None, None + + def stop(reason, **fields): + result.update(ok=False, state=reason, message=_MESSAGES[reason], **fields) + return result + + def write_store(method, *args, **kwargs): + if store is None or not identity: + raise ChatFailure("storage_error") + try: + return getattr(store, method)(identity, *args, **kwargs) + except Exception: + raise ChatFailure("storage_error") from None + + def cancel_unless_sent(): + """Leave a never-posted turn retryable; never reopen one that was sent. + + 'sent' is written immediately before the prompt POST, so a record still in + 'reserved' is knowably unsent and replaying it cannot spend credits twice. + Without this an ACP failure would strand the day as an unretryable 'pending'. + """ + if attempt_id is None: + return + try: + current = store.daily_chat_record(identity, day) if store else None + except Exception: + current = None + if current and current.get("phase") == "sent": + return + try: + write_store("transition_daily_chat", day, attempt_id, "cancelled") + except ChatFailure: + pass + + if not supported(profile): + return unavailable() + try: + previous = write_store("daily_chat_record", day) + if previous: + phase = previous.get("phase") + if phase in {"confirmed", "reconciled"}: + # Like a check-in that already ran: the day's goal is met, just not again. + result.update(ok=True, state="done", skipped=True, message=_MESSAGES["done"], + conversation_id=previous.get("conversation_id"), + sandbox_status=previous.get("sandbox_status"), + usage_before=previous.get("usage_before"), + usage_after=previous.get("usage_after"), + acp_usage=previous.get("acp_usage")) + return result + if phase != "cancelled": + # A pending or unconfirmed send is never replayed automatically. + return stop("pending", skipped=True, conversation_id=previous.get("conversation_id")) + if not can_write(): + return stop("changed", skipped=True) + reserved = write_store("reserve_daily_chat", day) + if reserved is None: + return stop("pending", skipped=True) + attempt_id = reserved["attempt_id"] + if not can_write(): + write_store("transition_daily_chat", day, attempt_id, "cancelled") + return stop("changed", skipped=True) + result["phase"] = "conversation" + result["usage_before"] = _usage_credits(token, uid, domain) + with httpx.Client(follow_redirects=False) as client: + headers = _console_headers(token, uid, domain) + conversation_id = create_conversation(client, headers) + result["conversation_id"] = conversation_id + # Record the conversation immediately: if anything fails later, this is + # what distinguishes "the turn ran" from "nothing ever reached the agent". + write_store("daily_chat_checkpoint", day, conversation_id=conversation_id) + # The sandbox is provisioned asynchronously, so an unready response is + # retried briefly. Nothing has reached the agent yet, which is why this + # failure path stays retryable for the rest of the day. + sandbox = _await_sandbox(client, headers, conversation_id) + if not can_write(): + raise ChatFailure("changed") + result["phase"] = "turn" + + def committed(): + # Last point at which the day is knowably still unspent: the prompt + # POST is about to go out and the turn can really run. + write_store("transition_daily_chat", day, attempt_id, "sent") + + channel = acp_client.run_turn(sandbox["link"], sandbox["token"], + sandbox["session_id"], sandbox["cwd"], PROMPT, + on_prompt=committed) + try: + started = time.monotonic() + while True: + if time.monotonic() - started > TURN_SECONDS: + raise ChatFailure("timeout") + if not can_write(): + raise ChatFailure("changed") + for update in channel.drain(POLL_SECONDS): + cost = update.get("cost") if isinstance(update, dict) else None + amount = cost.get("amount") if isinstance(cost, dict) else None + if type(amount) in (int, float) and 0 <= amount <= 10**9: + acp_usage = round(float(amount), 6) + try: + sandbox_status = status_of(client, headers, conversation_id) + except ChatFailure as error: + # A transient status failure is retried on the next poll; an auth + # rejection is terminal because the turn cannot finish either. + if error.diagnostics.get("error_kind") == "auth_error": + raise + continue + result["sandbox_status"] = sandbox_status or None + if sandbox_status in TERMINAL_STATES: + break + if sandbox_status in FAILED_STATES: + raise ChatFailure("rejected") + finally: + channel.close() + result["phase"] = "usage" + result["usage_after"] = _usage_credits(token, uid, domain) + result["acp_usage"] = acp_usage + write_store("daily_chat_checkpoint", day, sandbox_status=sandbox_status, + usage_before=result["usage_before"], usage_after=result["usage_after"], + acp_usage=acp_usage, conversation_id=conversation_id) + write_store("transition_daily_chat", day, attempt_id, "confirmed") + result.update(ok=True, state="done", message=_MESSAGES["done"]) + return result + except ChatFailure as error: + fields = dict(error.diagnostics) + kind = fields.pop("error_kind", "network_error") + if kind == "network_error" and fields.get("http_status") in (401, 403): + kind = "auth_error" + cancel_unless_sent() + result.update(state=kind, message=_MESSAGES.get(kind, _MESSAGES["network_error"]), **fields) + return result + except acp_client.AcpError as error: + # Surface the channel's own reason instead of collapsing every ACP failure + # into one opaque state: the diagnostics carry only a kind and a status. + # A channel that never delivered the prompt is an unsent day, so it must be + # cancelled here too; otherwise the day reads as unretryable 'pending'. + fields = dict(error.diagnostics) + kind = "acp_" + str(fields.pop("error_kind", "protocol")) + cancel_unless_sent() + result.update(state=kind, message=_MESSAGES.get(kind, _MESSAGES["protocol_error"]), **fields) + return result + except (httpx.HTTPError, ValueError, TypeError): + cancel_unless_sent() + return stop("network_error" if result["phase"] in {"conversation", "usage"} else "protocol_error") + except Exception: + cancel_unless_sent() + return stop("storage_error") + + +def view(record): + """Public state for the inventory; never exposes paths, tokens or upstream text.""" + if not record: + return {"state": "available", "message": _MESSAGES["available"], "day": None} + state = record.get("state") + if state not in _MESSAGES: + state = "done" if record.get("phase") in {"confirmed", "reconciled"} else "unconfirmed" + return {"state": state, "message": _MESSAGES[state], "day": record.get("day"), + "at": record.get("confirmed_at"), "conversation_id": record.get("conversation_id"), + "usage_before": record.get("usage_before"), "usage_after": record.get("usage_after"), + "acp_usage": record.get("acp_usage")} diff --git a/app/gateway_management.py b/app/gateway_management.py index c21ee33..e73a7db 100644 --- a/app/gateway_management.py +++ b/app/gateway_management.py @@ -8,7 +8,7 @@ from fastapi.responses import FileResponse, JSONResponse, RedirectResponse from starlette.staticfiles import StaticFiles -from . import checkin, model_policy, travel, trial_management +from . import checkin, daily_chat, model_policy, travel, trial_management class Management: @@ -26,6 +26,7 @@ def admin_credential_inventory(self): pool._rescan() now = time.time() trials = trial_management.inventory(self.CONFIG.get("trial_ledger")) + daily_chats = self.CONFIG.get("control_store") with pool._lock: entries = {entry["id"]: entry for entry in pool.entries()} result = [] @@ -53,6 +54,9 @@ def admin_credential_inventory(self): trial=(trial_management.view(trials.get(identity), now=now) if trials is not None else trial_management.failure("storage_error")) if entry.get("profile") == "intl-work" else trial_management.failure("not_applicable"), + daily_chat_supported=daily_chat.supported(entry.get("profile")), + auto_daily_chat=model_policy.credential_auto_daily_chat(self.CONFIG, entry), + daily_chat=self._daily_chat_view(daily_chats, identity, now), fail_until=until, cooldown_until=until, cooldown_remaining=max(0, round(until - now)), last_error_code=("http_401" if entry.get("last_error") == "backend HTTP 401" else @@ -109,6 +113,30 @@ def admin_set_auto_travel(self, identity, enabled): raise HTTPException(400, "旅行仅适用于国内账号") self.CONFIG["control_store"].set_auto_travel(identity, enabled) + def admin_set_auto_daily_chat(self, identity, enabled): + pool = self.CONFIG.get("cred_pool") + if pool is None: + raise HTTPException(404, "凭证不存在") + with pool._lock: + entry = next((entry for entry in pool.entries() if entry.get("account_key") == identity), None) + if entry is None: + raise HTTPException(404, "凭证不存在") + if enabled and not daily_chat.supported(entry.get("profile")): + raise HTTPException(400, "活跃打卡仅适用于国际 WorkBuddy 账号") + self.CONFIG["control_store"].set_auto_daily_chat(identity, enabled) + # Persist only: enabling never launches a turn; the sweep decides when to run it. + + @staticmethod + def _daily_chat_view(store, identity, now): + """Today's turn state, or the availability prompt when nothing is recorded yet.""" + if store is None: + return daily_chat.failure_view() + try: + record = store.daily_chat_record(identity, daily_chat.today(now)) + except Exception: + return daily_chat.failure_view() + return daily_chat.view(record) + def admin_delete_guard(self, name, *, unbind=False): """Block a still-bound delete, or return the identity to unbind after removal.""" rows = self.admin_credential_inventory() diff --git a/app/model_policy.py b/app/model_policy.py index d61f608..66aa180 100644 --- a/app/model_policy.py +++ b/app/model_policy.py @@ -41,6 +41,12 @@ def credential_auto_travel(config, entry): return supported and snapshot(config)["credentials"].get(entry.get("account_key"), {}).get("auto_travel", True) +def credential_auto_daily_chat(config, entry): + """A turn spends credits on the account, so international accounts opt in explicitly.""" + supported = entry.get("profile") == "intl-work" + return supported and snapshot(config)["credentials"].get(entry.get("account_key"), {}).get("auto_daily_chat", False) + + def default_rule(source): return {"public_id": source, "upstream_id": source, "custom": False, "enabled": True, "keep_original": False, diff --git a/converter.py b/converter.py index 62e5936..11034d1 100644 --- a/converter.py +++ b/converter.py @@ -50,7 +50,7 @@ def desensitize_body(body, roles=("system",), desensitize_harness_user=False, from app import auth_oauth from app import trial_rewards -from app import buddy, checkin as checkin_service, model_policy, travel +from app import buddy, daily_chat, checkin as checkin_service, model_policy, travel from app.credential_cooldowns import CredentialCooldowns from app.model_blocks import ModelBlocks from app.usage_snapshots import UsageSnapshots @@ -1411,6 +1411,18 @@ def _sync_error(pool, ledger, entry, generation, phase, error): _log(f"[{phase}] {Path(entry['id']).name} 同步失败(保留旧数据): {message}") +def _daily_chat_done(entry, day): + """Today's turn already confirmed, so the sweep must not spend another one.""" + store = CONFIG.get("control_store") + if store is None or not entry.get("account_key"): + return True + try: + return store.daily_chat_done(entry["account_key"], day) + except Exception: + # An unreadable record is treated as done: never guess a turn was free to send. + return True + + def _buddy_context(entry, headers, consent_revision=None): def select_model(requested=None): from app.audit_store import safe_label @@ -1464,7 +1476,7 @@ def _sync_credits(pool, ledger, entry, *, checkin, failed, expected_identity=Non return None # A path now owned by another account must be rescheduled with its own preferences. site = site_for_headers(headers) token, uid, domain = _bearer_token(headers), headers.get("X-User-Id", ""), headers.get("X-Domain", "") - day = time.strftime("%Y-%m-%d") + day = daily_chat.today() if checkin and model_policy.credential_auto_checkin(CONFIG, entry) and not ledger.checkin_done(cid, day): try: def can_claim(): @@ -1500,6 +1512,24 @@ def can_travel(): _sync_error(pool, ledger, entry, generation, "travel", error) if not model_policy.credential_enabled(CONFIG, entry): return None + # A daily-activity turn spends credits on the account and is opt-in, so it runs + # before the balance refresh reads back whatever it cost. Each account waits for + # its own slot inside the day so the accounts never go out together. + if (checkin and model_policy.credential_auto_daily_chat(CONFIG, entry) + and not _daily_chat_done(entry, day) + and daily_chat.is_due(entry.get("account_key"), day)): + try: + def can_chat(): + return (model_policy.credential_auto_daily_chat(CONFIG, entry) + and pool.apply_if_current(cm, generation, lambda: None)) + turn = daily_chat.perform(token, profile, uid=uid, domain=domain, can_write=can_chat, + store=CONFIG.get("control_store"), + identity=entry.get("account_key")) + _log(f"[dailychat] {Path(cid).name}: state={turn.get('state')} ok={turn.get('ok')}") + except Exception as error: + _sync_error(pool, ledger, entry, generation, "dailychat", error) + if not model_policy.credential_enabled(CONFIG, entry): + return None balance = credits_mod.fetch_credits(token, uid=uid, domain=domain) if bool(balance.get("intl")) != (site == INTERNATIONAL): raise ValueError("积分响应与凭据站点不一致") @@ -1748,6 +1778,14 @@ def _housekeep_once(pool: CredentialPool, ledger, *, pending_only=False): pool._warn_storage("cooldown", pool._cooldowns.last_error is None) if CONFIG.get("usage_snapshots") is not None: CONFIG["usage_snapshots"].prune() + control_store = CONFIG.get("control_store") + if control_store is not None: + # Daily-activity rows are diagnostic only, so old ones are dropped rather + # than kept as history; the day itself is still read back from this table. + try: + control_store.prune_daily_chats() + except Exception as error: + _log(f"[housekeeper] 打卡记录清理失败: {_network_error_text(error)}") ids = pool.begin_sync(all_entries=not pending_only) failed = set() try: @@ -1889,7 +1927,7 @@ async def _protocol_http_exception(request: Request, exc: HTTPException): def _log(msg: str): """Persist allowlisted runtime events in SQLite, never free-form text or secrets.""" audit = CONFIG.get("audit_store") - component = re.match(r"\[(cred|credits|models|usage|trial|checkin|housekeeper)\]", msg) + component = re.match(r"\[(cred|credits|models|usage|trial|checkin|dailychat|housekeeper)\]", msg) if audit is not None and component: # Persist event codes, not free-form lines which may contain upstream data. code = "cooldown" if "熔断" in msg or "冷却" in msg else "failure" if "失败" in msg or "异常" in msg else "updated" diff --git a/docker-compose.yml b/docker-compose.yml index 601dc6d..604aeb2 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -15,6 +15,9 @@ services: environment: CODEBUDDY2API_KEY: ${CODEBUDDY2API_KEY:-} CODEBUDDY2API_ADMIN_CSRF: ${CODEBUDDY2API_ADMIN_CSRF:-true} + # Daily guards key on a local calendar day, so the container should share the + # operator's zone; a UTC container turns those guards over at 08:00 Beijing. + TZ: ${TZ:-UTC} # First-Buddy onboarding may send one bounded real WorkBuddy conversation per account. CODEBUDDY2API_AUTO_ACCEPT_BUDDY: ${CODEBUDDY2API_AUTO_ACCEPT_BUDDY:-false} # Unset variables remain configurable through the WebUI. diff --git a/tests/test_daily_chat.py b/tests/test_daily_chat.py new file mode 100644 index 0000000..3d9db11 --- /dev/null +++ b/tests/test_daily_chat.py @@ -0,0 +1,497 @@ +"""Offline contracts for the international daily-activity turn: no upstream calls. + +Everything here drives ``app.daily_chat`` and ``app.acp_client`` against fakes, so +nothing reaches WorkBuddy. The properties worth pinning are the safety ones: the +turn is one-shot per day, a sent-but-unconfirmed attempt is never replayed, no +upstream text escapes into results, and the ACP channel declares no client +capabilities (otherwise the sandbox would wait for callbacks that never come). +""" +import sys +from pathlib import Path +sys.path.insert(0, str(Path(__file__).resolve().parents[1])) + +import json +import time +import unittest +from unittest.mock import patch + +import converter +from app import acp_client, daily_chat +from app.control_store import ControlStore +from tests import test_credential_actions as fixtures + + +class FakeConsole: + """In-memory stand-in for the three console endpoints plus one sandbox turn.""" + + def __init__(self, *, conversation_id="conv-1", sandbox_status="completed", + link="https://sandbox.invalid/acp", token="sandbox-token", + session_id="sess-1", cwd="/workspace", fail_create=False): + self.conversation_id = conversation_id + self.sandbox_status = sandbox_status + self.session = {"link": link, "token": token, "session_id": session_id, "cwd": cwd} + self.fail_create = fail_create + self.created = [] + self.statuses = [] + self.posts = [] + self.methods = [] + self.channel_opens = 0 + + def install(self, test): + # Lambdas, not bound methods: patch.object installs the object as a class + # attribute, and a bound method is not a descriptor, so instance.open() + # would call it with no arguments instead of the channel. + console = self + test.enterContext(patch.object(daily_chat, "create_conversation", self.create)) + test.enterContext(patch.object(daily_chat, "sandbox_of", self.sandbox_of)) + test.enterContext(patch.object(daily_chat, "status_of", self.status_of)) + test.enterContext(patch.object(acp_client.AcpChannel, "open", + lambda channel: console.open(channel))) + test.enterContext(patch.object(acp_client.AcpChannel, "post", + lambda channel, method, params, request_id: + console.post(channel, method, params, request_id))) + test.enterContext(patch.object(acp_client.AcpChannel, "drain", + lambda channel, seconds: console.drain(channel, seconds))) + test.enterContext(patch.object(acp_client.AcpChannel, "close", + lambda channel: console.close(channel))) + test.enterContext(patch.object(daily_chat.credits, "fetch_request_usage", + side_effect=self.usage)) + + def usage(self, token, days=1, uid="", domain=""): + return {"by_day": {}, "total_credits": 0.0, "requests": 0, "partial": False} + + def create(self, client, headers): + self.created.append(headers) + if self.fail_create: + raise daily_chat.ChatFailure("http_error", 503) + return self.conversation_id + + def sandbox_of(self, client, headers, conversation_id): + return dict(self.session, conversation_id=conversation_id) + + def status_of(self, client, headers, conversation_id): + self.statuses.append(conversation_id) + return self.sandbox_status + + def open(self, channel_self): + self.channel_opens += 1 + channel_self.connection_id = "conn-1" + channel_self._closed = False + return channel_self + + def post(self, channel_self, method, params, request_id): + self.methods.append((method, request_id)) + self.posts.append({"method": method, "params": params, "id": request_id}) + + def drain(self, channel_self, seconds): + self.drained = getattr(self, "drained", 0) + 1 + if self.drained == 1: + return [{"sessionUpdate": "usage_update", + "cost": {"amount": 3, "currency": "credits"}}] + return [] + + def close(self, channel_self): + channel_self._closed = True + + def run_turn(self, link, token, session_id, cwd, prompt): + """Kept only so the old assertions on turn arguments still have a place.""" + self.channel_opens += 1 + self.posts.append({"link": link, "token": token, "session_id": session_id, + "cwd": cwd, "prompt": prompt}) + return FakeChannel(self) + + def channel_post(self, method, params, request_id): + self.posts.append({"method": method, "params": params, "id": request_id}) + + +class FakeChannel: + """Stands in for the ACP channel: one drained usage update, then done.""" + + def __init__(self, console): + self.console = console + self.closed = False + self.drained = 0 + + def drain(self, seconds): + self.drained += 1 + if self.drained == 1: + return [{"sessionUpdate": "usage_update", + "cost": {"amount": 3, "currency": "credits"}}] + return [] + + def close(self): + self.closed = True + + +class DailyChatTests(fixtures.CredentialActionTests): + add_account = fixtures.CredentialActionTests.add_account + configure = fixtures.CredentialActionTests.configure + handle_upstream = fixtures.CredentialActionTests.handle_upstream + + def setUp(self): + fixtures.CredentialActionTests.setUp(self) + self.intl = self.entries["intl-work"] + self.console = FakeConsole() + self.console.install(self) + + def turn(self, entry=None, **kwargs): + return daily_chat.perform("synthetic-token", (entry or self.intl).get("profile"), + uid="synthetic-uid", domain="workbuddy.ai", + store=self.control, identity=(entry or self.intl)["account_key"], **kwargs) + + # -- profile gating ------------------------------------------------------ + + def test_only_international_workbuddy_accounts_are_supported(self): + self.assertEqual(daily_chat.supported("intl-work"), True) + for profile in ("intl-cli", "cn-cli", "cn-work"): + self.assertFalse(daily_chat.supported(profile), profile) + self.assertTrue(daily_chat.unavailable()["skipped"]) + + def test_domestic_profile_never_reaches_the_console(self): + self.assertFalse(self.turn(self.entries["cn-work"])["ok"]) + self.assertEqual(self.console.created, []) + + # -- day boundary and staggering ----------------------------------------- + + def test_today_uses_the_process_local_calendar_day(self): + # The reward day is a local calendar day, so the guard must not key on UTC: + # a UTC container would roll it over at 08:00 Beijing and could visit one + # Beijing day twice. 16:30 UTC is already 10-01 in Beijing. + self.assertEqual(daily_chat.today(), time.strftime("%Y-%m-%d", time.localtime())) + self.assertEqual(daily_chat.today(time.mktime((2026, 10, 1, 0, 30, 0, 0, 0, -1))), "2026-10-01") + + def test_staggering_slots_are_stable_and_spread_across_accounts(self): + day = "2026-10-01" + keys = [f"account-{i}" for i in range(24)] + slots = [daily_chat.due_at(key, day) for key in keys] + # Stable for the same account and day: a sweep may run many times. + self.assertEqual(daily_chat.due_at(keys[0], day), slots[0]) + self.assertEqual(len(set(slots)), len(slots), "accounts must not share a slot") + for slot in slots: + self.assertTrue(0 <= slot < 6 * 3600, slot) + # A different day reshuffles the order, so no account is always first. + other = [daily_chat.due_at(key, "2026-10-02") for key in keys] + self.assertNotEqual(slots, other) + + def test_is_due_respects_the_account_slot(self): + day = "2026-10-01" + midnight = time.mktime((2026, 10, 1, 0, 0, 0, 0, 0, -1)) + key = "account-x" + slot = daily_chat.due_at(key, day) + self.assertFalse(daily_chat.is_due(key, day, midnight + max(0, slot - 1))) + self.assertTrue(daily_chat.is_due(key, day, midnight + slot)) + + def test_zero_window_is_immediately_due(self): + day = "2026-10-01" + self.assertEqual(daily_chat.due_at("account-x", day, window=0), 0.0) + self.assertTrue(daily_chat.is_due("account-x", day, 0, window=0)) + + def test_every_slot_is_reachable_by_an_hourly_sweep(self): + # The sweep runs hourly (HOUSEKEEP_INTERVAL), so a slot must never fall + # between two runs and leave the account unwoken for the day. The window is + # bounded well inside the day for exactly this reason. + day = "2026-10-01" + window = 6 * 3600 + for index in range(200): + key = f"account-{index}" + slot = daily_chat.due_at(key, day, window=window) + self.assertLess(slot, 24 * 3600, "a slot must stay inside the same day") + self.assertTrue(any(run >= slot for run in range(3600, 86400, 3600)), + f"{key} slot {slot} is unreachable by an hourly sweep") + + # -- happy path ---------------------------------------------------------- + + def test_completed_turn_is_confirmed_and_records_cost(self): + result = self.turn() + self.assertTrue(result["ok"], result) + self.assertEqual(result["state"], "done") + self.assertEqual(result["sandbox_status"], "completed") + self.assertEqual(result["conversation_id"], "conv-1") + # Cost comes from the protocol-native usage_update, not from a subtraction. + self.assertEqual(result["acp_usage"], 3.0) + self.assertEqual(len(self.console.created), 1) + record = self.control.daily_chat_record(self.intl["account_key"], result["day"]) + self.assertEqual(record["phase"], "confirmed") + self.assertTrue(self.control.daily_chat_done(self.intl["account_key"], result["day"])) + + def test_console_headers_carry_no_cli_fingerprint(self): + self.turn() + headers = self.console.created[0] + self.assertEqual(headers["x-user-id"], "synthetic-uid") + self.assertEqual(headers["referer"], daily_chat.HOST + "/app") + self.assertTrue(headers["authorization"].startswith("Bearer ")) + # A CLI X-IDE-* fingerprint is exactly what fails to earn the reward. + for name in headers: + self.assertFalse(name.lower().startswith("x-ide-"), name) + + def test_acp_sequence_uses_ordered_ids_and_closed_capabilities(self): + self.turn() + self.assertEqual([method for method, _ in self.console.methods], + ["initialize", "session/load", "session/prompt"]) + self.assertEqual([request_id for _, request_id in self.console.methods], [1, 2, 3]) + by_method = {method: params for method, params in self.console.methods + for params in [next(p["params"] for p in self.console.posts + if p["method"] == method)]} + self.assertEqual(by_method["initialize"]["protocolVersion"], 1) + capabilities = by_method["initialize"]["clientCapabilities"] + self.assertIs(capabilities["fs"]["readTextFile"], False) + self.assertIs(capabilities["fs"]["writeTextFile"], False) + self.assertIs(capabilities["terminal"], False) + self.assertEqual(by_method["session/load"]["sessionId"], "sess-1") + self.assertEqual(by_method["session/load"]["cwd"], "/workspace") + self.assertEqual(by_method["session/prompt"]["prompt"], + [{"type": "text", "text": daily_chat.PROMPT}]) + + # -- one-shot discipline ------------------------------------------------- + + def test_second_turn_the_same_day_sends_nothing(self): + self.turn() + self.assertEqual(len(self.console.created), 1) + self.assertEqual(len(self.console.methods), 3) + again = self.turn() + self.assertTrue(again["ok"]) + self.assertEqual(again["state"], "done") + self.assertTrue(again["skipped"]) + self.assertEqual(len(self.console.created), 1, "a confirmed day is never reopened") + self.assertEqual(len(self.console.methods), 3) + + def test_unconfirmed_send_is_never_replayed(self): + # Key the seeded record to the day perform() will actually use, not a fixed + # date: a hard-coded day silently stops matching once the clock moves past it. + day = daily_chat.today() + self.control.reserve_daily_chat(self.intl["account_key"], day) + self.control.transition_daily_chat(self.intl["account_key"], day, + self.control.daily_chat_record(self.intl["account_key"], day)["attempt_id"], "sent") + result = self.turn() + self.assertEqual(result["state"], "pending") + self.assertTrue(result["skipped"]) + self.assertEqual(self.console.created, [], "an unconfirmed send must not be replayed") + + def test_failed_today_allows_tomorrow(self): + self.console.sandbox_status = "failed" + result = self.turn() + self.assertFalse(result["ok"]) + self.assertEqual(result["state"], "rejected") + record = self.control.daily_chat_record(self.intl["account_key"], result["day"]) + # 'sent' means the turn really ran: it stays unconfirmed for the rest of the day. + self.assertEqual(record["phase"], "sent") + self.assertFalse(self.control.daily_chat_done(self.intl["account_key"], result["day"])) + + def test_unready_sandbox_retries_then_gives_up_without_sending(self): + # A sandbox that never becomes ready must not spend the day: the turn was + # never handed to the agent, so the reservation stays retryable. + with patch.object(daily_chat, "sandbox_of", + side_effect=daily_chat.ChatFailure("sandbox_pending")), \ + patch.object(daily_chat, "SANDBOX_SECONDS", 0.3), \ + patch.object(daily_chat, "SANDBOX_POLL_SECONDS", 0.1): + result = self.turn() + self.assertFalse(result["ok"]) + self.assertEqual(result["state"], "sandbox_timeout") + record = self.control.daily_chat_record(self.intl["account_key"], result["day"]) + # Cancelled, not sent: nothing reached the agent, so the day is not spent. + self.assertEqual(record["phase"], "cancelled") + self.assertFalse(self.control.daily_chat_done(self.intl["account_key"], result["day"])) + self.assertIsNotNone(self.control.reserve_daily_chat(self.intl["account_key"], result["day"])) + + def test_sandbox_readiness_is_retried_until_it_answers(self): + attempts = {"n": 0} + real = daily_chat.sandbox_of + + def flaky(client, headers, conversation_id): + attempts["n"] += 1 + if attempts["n"] < 3: + raise daily_chat.ChatFailure("sandbox_pending") + return real(client, headers, conversation_id) + + with patch.object(daily_chat, "sandbox_of", side_effect=flaky), \ + patch.object(daily_chat, "SANDBOX_POLL_SECONDS", 0.05): + result = self.turn() + self.assertTrue(result["ok"], result) + self.assertEqual(attempts["n"], 3) + + def test_cancelled_day_can_be_retried_manually(self): + day = "2026-09-30" + attempt = self.control.reserve_daily_chat(self.intl["account_key"], day)["attempt_id"] + self.control.transition_daily_chat(self.intl["account_key"], day, attempt, "cancelled") + reserved = self.control.reserve_daily_chat(self.intl["account_key"], day) + self.assertIsNotNone(reserved) + + # -- failure containment ------------------------------------------------- + + def test_upstream_http_failure_cancels_unsent_reservation(self): + self.console.fail_create = True + result = self.turn() + self.assertFalse(result["ok"]) + self.assertEqual(result["state"], "http_error") + record = self.control.daily_chat_record(self.intl["account_key"], result["day"]) + self.assertEqual(record["phase"], "cancelled") + self.assertEqual(self.console.channel_opens, 0, "no sandbox turn was opened") + + def test_acp_failure_cancels_unsent_reservation(self): + # The channel failed before the prompt was posted, so the day is knowably + # unspent and must stay retryable instead of stranding as 'pending'. + with patch.object(acp_client.AcpChannel, "open", + side_effect=acp_client.AcpError("network")): + result = self.turn() + self.assertFalse(result["ok"]) + self.assertEqual(result["state"], "acp_network") + record = self.control.daily_chat_record(self.intl["account_key"], result["day"]) + self.assertEqual(record["phase"], "cancelled") + self.assertFalse(self.control.daily_chat_done(self.intl["account_key"], result["day"])) + self.assertIsNotNone(self.control.reserve_daily_chat(self.intl["account_key"], result["day"])) + + def test_acp_failure_after_prompt_keeps_the_day_spent(self): + # Once the prompt is posted the turn may have run upstream, so the record + # stays 'sent' and is never replayed even though the channel then failed. + with patch.object(acp_client.AcpChannel, "post", + side_effect=[None, None, acp_client.AcpError("network")]): + result = self.turn() + self.assertFalse(result["ok"]) + self.assertEqual(result["state"], "acp_network") + record = self.control.daily_chat_record(self.intl["account_key"], result["day"]) + self.assertEqual(record["phase"], "sent") + self.assertIsNone(self.control.reserve_daily_chat(self.intl["account_key"], result["day"])) + + def test_timeout_leaves_the_reservation_uncancelled(self): + self.console.sandbox_status = "working" + with patch.object(daily_chat, "TURN_SECONDS", 1.0), patch.object(daily_chat, "POLL_SECONDS", 0.4): + result = self.turn() + self.assertFalse(result["ok"]) + self.assertEqual(result["state"], "timeout") + record = self.control.daily_chat_record(self.intl["account_key"], result["day"]) + # The turn was sent, so it stays unconfirmed; the next attempt is tomorrow. + self.assertEqual(record["phase"], "sent") + self.assertFalse(self.control.daily_chat_done(self.intl["account_key"], result["day"])) + + def test_unwritable_store_stops_before_any_request(self): + with patch.object(self.control, "daily_chat_record", side_effect=OSError("synthetic")): + result = self.turn() + self.assertEqual(result["state"], "storage_error") + self.assertEqual(self.console.created, []) + + def test_generation_change_cancels_before_sending(self): + result = self.turn(can_write=lambda: False) + self.assertEqual(result["state"], "changed") + self.assertTrue(result["skipped"]) + self.assertEqual(self.console.created, []) + + def test_no_upstream_text_reaches_the_result(self): + with patch.object(daily_chat, "create_conversation", + side_effect=daily_chat.ChatFailure("rejected", 429, 11140)): + result = self.turn() + # Keys such as sandbox_status are our own vocabulary; only *values* must be clean. + def values(node): + if isinstance(node, dict): + for value in node.values(): + yield from values(value) + elif isinstance(node, list): + for item in node: + yield from values(item) + elif isinstance(node, str): + yield node + for text in values(result): + for forbidden in ("synthetic-token", "synthetic-uid", "sandbox.invalid", "Bearer"): + self.assertNotIn(forbidden, text, forbidden) + self.assertEqual(result["code"], 11140) + self.assertEqual(result["http_status"], 429) + + # -- manual action and preference --------------------------------------- + + def test_manual_action_runs_one_turn_and_audits_it(self): + response = self.client.post("/admin/credentials/" + self.intl["account_key"] + "/daily-chat") + self.assertEqual(response.status_code, 200, response.text) + result = response.json()["results"][0] + self.assertTrue(result["ok"], result) + self.assertEqual(result["state"], "done") + self.assertEqual(len(self.console.created), 1) + + def test_preference_defaults_off_and_persists(self): + row = next(r for r in self.client.get("/admin/credentials").json()["credentials"] + if r["id"] == self.intl["account_key"]) + self.assertIs(row["auto_daily_chat"], False) + self.assertIs(row["daily_chat_supported"], True) + self.assertEqual(row["daily_chat"]["state"], "available") + response = self.client.patch("/admin/credentials/" + self.intl["account_key"], + json={"auto_daily_chat": True}) + self.assertEqual(response.status_code, 200, response.text) + self.assertTrue(converter.model_policy.credential_auto_daily_chat(converter.CONFIG, self.intl)) + + def test_preference_is_rejected_for_domestic_accounts(self): + response = self.client.patch("/admin/credentials/" + self.entries["cn-work"]["account_key"], + json={"auto_daily_chat": True}) + self.assertEqual(response.status_code, 400, response.text) + + def test_unknown_action_is_still_rejected(self): + response = self.client.post(self.url + "/daily-checkin") + self.assertEqual(response.status_code, 404, response.text) + + # -- store contracts ----------------------------------------------------- + + def test_reserve_rejects_a_second_reservation_the_same_day(self): + day = "2026-09-30" + first = self.control.reserve_daily_chat(self.intl["account_key"], day) + self.assertIsNotNone(first) + # Still 'reserved' (nothing sent yet), so re-reserving is allowed on purpose: + # an unsent reservation must not cost the day. + second = self.control.reserve_daily_chat(self.intl["account_key"], day) + self.assertIsNotNone(second) + self.assertNotEqual(first["attempt_id"], second["attempt_id"]) + # Once sent, the day is spent. + self.control.transition_daily_chat(self.intl["account_key"], day, + second["attempt_id"], "sent") + self.assertIsNone(self.control.reserve_daily_chat(self.intl["account_key"], day)) + + def test_transition_rejects_an_unknown_attempt(self): + day = "2026-09-30" + attempt = self.control.reserve_daily_chat(self.intl["account_key"], day)["attempt_id"] + with self.assertRaises(ValueError): + self.control.transition_daily_chat(self.intl["account_key"], day, "not-the-attempt", "sent") + self.assertIsNotNone(attempt) + + def test_day_must_be_a_plain_date(self): + for bad in ("2026-9-30", "20260930", 20260930, "2026-09-30T00:00:00", "", + # Syntactically well-shaped but not a real calendar day. + "2026-99-99", "2026-02-31", "2026-13-01", "2026-04-31", "0000-01-01"): + with self.assertRaises(ValueError, msg=repr(bad)): + self.control.daily_chat_record(self.intl["account_key"], bad) + + def test_prune_drops_only_old_days(self): + identity = self.intl["account_key"] + for day in ("2020-01-01", "2999-01-01"): + attempt = self.control.reserve_daily_chat(identity, day)["attempt_id"] + self.control.transition_daily_chat(identity, day, attempt, "sent") + self.control.transition_daily_chat(identity, day, attempt, "confirmed") + self.control.prune_daily_chats(keep_days=30) + self.assertIsNone(self.control.daily_chat_record(identity, "2020-01-01")) + self.assertIsNotNone(self.control.daily_chat_record(identity, "2999-01-01")) + + +class AcpClientTests(unittest.TestCase): + """Protocol-shape guards that need no network at all.""" + + def test_link_must_be_http_and_keeps_its_query(self): + for bad in ("", "ftp://x/y", "notaurl", "https:///path"): + with self.assertRaises(acp_client.AcpError, msg=bad): + acp_client.AcpChannel(bad, "t").open() + channel = acp_client.AcpChannel.__new__(acp_client.AcpChannel) + channel.parts, channel.path = acp_client._target("https://sandbox.invalid/acp?a=1&b=2") + self.assertEqual(channel.path, "/acp?a=1&b=2") + + def test_unknown_method_is_refused_before_any_connection(self): + channel = acp_client.AcpChannel.__new__(acp_client.AcpChannel) + channel._closed = False + channel.connection_id = "conn-1" + with self.assertRaises(acp_client.AcpError): + channel.post("session/delete", {}, 4) + + def test_drain_harvests_usage_updates_only(self): + channel = acp_client.AcpChannel.__new__(acp_client.AcpChannel) + channel._closed, channel._response, channel._buffer, channel._line = False, object(), bytearray(), bytearray() + channel._buffer.extend( + b"data: {\"jsonrpc\":\"2.0\",\"method\":\"session/update\",\"params\":{\"update\":" + b"{\"sessionUpdate\":\"usage_update\",\"cost\":{\"amount\":7,\"currency\":\"credits\"}}}}\n" + b"data: {\"method\":\"session/update\",\"params\":{\"update\":{\"sessionUpdate\":" + b"\"agent_message_chunk\",\"content\":{\"text\":\"hi\"}}}}\n" + b": heartbeat\n\n") + self.assertEqual(channel._events(), [{"sessionUpdate": "usage_update", + "cost": {"amount": 7, "currency": "credits"}}]) diff --git a/web/e2e/automation.spec.ts b/web/e2e/automation.spec.ts index 8017c36..2c5f5dc 100644 --- a/web/e2e/automation.spec.ts +++ b/web/e2e/automation.spec.ts @@ -70,7 +70,7 @@ test("per-account automation persists with scoped actions, dark mode and mobile await expect(page.getByRole("switch", { name: "自动签到 intl.info" })).not.toBeChecked(); await expect(page.getByRole("switch", { name: "自动旅行 intl.info" })).toBeDisabled(); await page.getByRole("switch", { name: "自动签到 intl.info" }).click(); - await expect(page.getByText(/保存不会立即领取/)).toBeVisible(); + await expect(page.getByText(/保存不会立即执行/)).toBeVisible(); expect(writes).toEqual(["PATCH /admin/credentials/intl"]); await page.reload(); await expect(page.getByRole("switch", { name: "自动签到 intl.info" })).toBeChecked(); diff --git a/web/src/api.ts b/web/src/api.ts index bb0552c..3b9ba95 100644 --- a/web/src/api.ts +++ b/web/src/api.ts @@ -227,10 +227,17 @@ export function credentialResponse(value: unknown): Credential[] { const id = c.account_key ?? c.id; if (typeof id !== "string" || !id) throw new Error("凭证缺少公开 account_key"); const name = c.name ?? c.filename; - for (const field of ["auto_checkin", "auto_travel", "travel_supported", "trial_supported"]) + for (const field of [ + "auto_checkin", + "auto_travel", + "auto_daily_chat", + "travel_supported", + "trial_supported", + "daily_chat_supported", + ]) if (c[field] !== undefined && typeof c[field] !== "boolean") throw new Error("自动任务状态必须为布尔值"); - for (const field of ["checkin", "travel", "trial"]) + for (const field of ["checkin", "travel", "trial", "daily_chat"]) if (c[field] !== undefined && c[field] !== null) object(c[field], "自动任务结果"); return { ...c, id, name: typeof name === "string" && !/[\\/]/.test(name) ? name : null }; }); diff --git a/web/src/automation.test.tsx b/web/src/automation.test.tsx index 8e158a0..98008ee 100644 --- a/web/src/automation.test.tsx +++ b/web/src/automation.test.tsx @@ -32,8 +32,11 @@ function fixture() { auto_checkin: false, auto_travel: false, travel_supported: false, + daily_chat_supported: true, + auto_daily_chat: false, checkin: { state: "inactive", date: "2026-09-15", message: "签到活动未开放或已结束" }, travel: { state: "unknown", message: "尚未查询" }, + daily_chat: { state: "idle", day: "2026-10-03", message: "尚未执行" }, }, ]; const reload = vi.fn(); @@ -61,6 +64,16 @@ it("renders server preferences, allows international checkin opt-in, but never o true, ); expect(screen.queryByRole("button", { name: "旅行领派 intl.info" })).toBeNull(); + // The daily turn is international-only, so the domestic switch is shown but unusable. + expect(screen.getByRole("switch", { name: "自动活跃打卡 cn.info" })).toHaveProperty( + "disabled", + true, + ); + expect(screen.getByRole("switch", { name: "自动活跃打卡 intl.info" })).toHaveProperty( + "checked", + false, + ); + expect(screen.getByText("今日打卡:尚未执行")).toBeTruthy(); await act(async () => fireEvent.click(screen.getByRole("switch", { name: "自动签到 intl.info" })), ); @@ -69,9 +82,15 @@ it("renders server preferences, allows international checkin opt-in, but never o "checked", true, ); - expect(screen.getByText(/保存不会立即领取/)).toBeTruthy(); + expect(screen.getByText(/保存不会立即执行/)).toBeTruthy(); expect(reload).toHaveBeenCalled(); post.mockClear(); + // Opting into the daily turn only persists; the sweep decides when the turn runs. + await act(async () => + fireEvent.click(screen.getByRole("switch", { name: "自动活跃打卡 intl.info" })), + ); + expect(patch).toHaveBeenLastCalledWith("/credentials/intl", { auto_daily_chat: true }); + expect(post).not.toHaveBeenCalled(); await act(async () => fireEvent.click(screen.getByRole("switch", { name: "自动旅行 cn.info" }))); expect(patch).toHaveBeenLastCalledWith("/credentials/cn", { auto_travel: false }); expect(screen.getByRole("switch", { name: "自动签到 cn.info" })).toHaveProperty("checked", true); diff --git a/web/src/pages/Credentials.tsx b/web/src/pages/Credentials.tsx index 3a87787..89622ab 100644 --- a/web/src/pages/Credentials.tsx +++ b/web/src/pages/Credentials.tsx @@ -122,7 +122,14 @@ export function Credentials() { const [notice, setNotice] = useState(null); const [maintenance, setMaintenance] = useState[]>([]); const maintain = ( - action: "refresh" | "checkin" | "sync" | "travel" | "travel-status" | "reset-cooldown", + action: + | "refresh" + | "checkin" + | "sync" + | "travel" + | "travel-status" + | "daily-chat" + | "reset-cooldown", credential?: Credential, ) => { if (busy) return; @@ -185,7 +192,7 @@ export function Credentials() { }; const preference = ( credential: Credential, - field: "auto_checkin" | "auto_travel", + field: "auto_checkin" | "auto_travel" | "auto_daily_chat", enabled: boolean, ) => { run(async () => { @@ -199,9 +206,13 @@ export function Credentials() { !Number.isInteger(saved.revision) ) throw new Error("设置保存结果未确认,请刷新列表核验"); - setNotice( - `${field === "auto_checkin" ? "自动签到" : "自动旅行"}已${enabled ? "开启" : "关闭"};保存不会立即领取,后续维护按新设置执行。`, - ); + const label = + field === "auto_checkin" + ? "自动签到" + : field === "auto_travel" + ? "自动旅行" + : "自动活跃打卡"; + setNotice(`${label}已${enabled ? "开启" : "关闭"};保存不会立即执行,后续维护按新设置执行。`); }); }; const liveSelected = selected.filter((id) => resource.data?.some((c) => c.id === id)); @@ -332,6 +343,10 @@ export function Credentials() { const checkin = c.checkin && typeof c.checkin === "object" ? object(c.checkin) : null; const trip = c.travel && typeof c.travel === "object" ? object(c.travel) : null; + const dailyChat = + c.daily_chat && typeof c.daily_chat === "object" + ? object(c.daily_chat) + : null; return ( @@ -406,6 +421,29 @@ export function Credentials() { ? `上次旅行:${trip ? text(trip.message) : "尚未查询"}` : "旅行仅适用于国内账号"} + + + {c.daily_chat_supported === true + ? `今日打卡:${dailyChat ? text(dailyChat.message) : "尚未执行"}` + : "活跃打卡仅适用于国际 WorkBuddy 账号"} + + {dailyChat?.acp_usage != null && ( + 本次会话成本:{text(dailyChat.acp_usage)} 积分 + )} {trip?.stale === true && 状态可能已变化,请先查询核验} {c.travel_supported === true && } {c.enabled === false && 账号停用期间不执行自动任务} @@ -504,6 +542,15 @@ export function Credentials() { 领取体验积分 )} + {c.daily_chat_supported === true && ( + + )} {c.travel_supported === true && ( <>