diff --git a/.github/workflows/plugin-api-v3.yml b/.github/workflows/plugin-api-v3.yml index ff6a040..d460cb4 100644 --- a/.github/workflows/plugin-api-v3.yml +++ b/.github/workflows/plugin-api-v3.yml @@ -36,7 +36,7 @@ jobs: - uses: actions/checkout@v4 with: repository: kachofugetsu09/akashic-agent - ref: 3005f838bcd96e2cbc58616aede46e4f39df4523 + ref: 065dfb7fb37534c57ed500b1ed9b5deb7090cdb0 path: .akashic-core - uses: actions/setup-python@v5 with: @@ -56,14 +56,15 @@ jobs: - name: Run Feed unit tests env: AKASHIC_AGENT_ROOT: .akashic-core + AKASHIC_PLUGIN_FIXTURE_PYTHON: ${{ github.workspace }}/mcp/.venv/bin/python PYTHONPATH: .akashic-core:mcp/.venv/lib/python3.13/site-packages run: mcp/.venv/bin/python -m pytest -q mcp/tests tests - name: Run pyright env: AKASHIC_AGENT_ROOT: .akashic-core PYTHONPATH: .akashic-core:mcp/.venv/lib/python3.13/site-packages - run: mcp/.venv/bin/pyright plugin.py mcp/run_mcp.py mcp/src mcp/scripts scripts + run: mcp/.venv/bin/pyright plugin.py content_source.py legacy_handoff.py feed_runtime mcp/run_mcp.py mcp/src mcp/scripts scripts - name: Compile Python sources - run: python -m compileall -q plugin.py mcp/run_mcp.py mcp/src mcp/scripts scripts tests + run: python -m compileall -q plugin.py content_source.py legacy_handoff.py feed_runtime mcp/run_mcp.py mcp/src mcp/scripts scripts tests - name: Check diff formatting run: git diff --check diff --git a/README.md b/README.md index 7d25fca..165fad4 100644 --- a/README.md +++ b/README.md @@ -1,71 +1,61 @@ # feed-mcp -Feed 是一个 Akashic Plugin API v3 插件,提供 RSS 订阅管理、Feed MCP 和 -`subscriptions` 主动内容源,同时保留 `skills/` 下的 Feed 技能。 +Feed 是一个 Akashic Plugin API v3 插件。它用现有的三个正交能力组合出完整 +Feed 链路: -`feed_query` 只读缓存,不主动触发拉取;缓存 freshness 由 MCP lifespan 的后台 -`FeedPoller` 按 `poll_ttl_seconds`(默认 300s)周期刷新,或通过 `poll_feeds` -显式触发。 +```text +Core Timer ──触发──> Feed source ──submit──> Content inbox + │ │ + └──精确 revision ACK <───┘ + +用户 Turn ──调用──> Feed MCP ──只做──> 订阅管理 / 缓存查询 +``` + +## 能力与 owner -目录结构: +- `MCP_SERVERS`:只注册用户主动调用的 `feed_manage`、`feed_query`。MCP + lifespan 不启动后台轮询,也没有主动内容专用工具。 +- `TIMERS`:正式稳定 Root 拥有唯一轮询 Timer。每次只注册一个 one-shot + deadline;完成 settlement、拉取、提交和 deadline 持久化后再注册下一次。 +- `content.source.v1`:以 `feed-subscriptions` 身份提交完整 Feed item,并接收 + Content 的 delivery settlement。ACK 精确绑定 `event_id + content_hash`;旧 + revision 的 ACK 不会误确认新内容。 + +Feed 的 SQLite 是 provider 事实 owner:订阅、完整 item、当前 poll state、ACK +和每个已导出 revision 的冻结 payload 都留在 `feed_mcp.sqlite3`。Content 是待处理 +与投递状态 owner。纯诊断日志固定为 5 MiB、最多 3 个备份;空轮询只更新当前 +deadline,不提交 Content,也不制造持久内容历史。 + +候选版本允许自己的 MCP managed-process 完成 readiness/handshake,但 +`candidate_read_only_tools = []`,且 recording backend 不接触 Feed 数据。候选 Root +不会收到 `runtime.started`,因此不会注册 Timer、访问外部 Feed 或写正式数据库。 + +## 目录 ```text feed-mcp ├─ akashic.plugin.toml -├─ plugin.py -├─ scripts/migrate_v2_data.py +├─ plugin.py # 组合 MCP、Timer、Content source +├─ content_source.py # Feed source 私有生命周期 +├─ feed_runtime/backend.py # MCP 与 source 共享的唯一 Feed domain 实现 ├─ skills/ └─ mcp/ ├─ run_mcp.py - └─ src/ + └─ src/mcp_bridge.py # 薄 MCP adapter ``` -`plugin.py` 只执行 `apply(ctx, config)` 声明,不启动进程、不访问网络、不读写 -插件数据。Core 从静态 manifest 准备 MCP runtime,并按配置注册 -`subscriptions` 主动事件源;`skill_roots = ("skills",)` 保持原有技能装载路径。 - -正式运行数据位于 Core 分配的 `plugin-data/feed-/`: - -- `feed_mcp.sqlite3`:订阅、条目、确认和轮询状态 -- `source_scores.json`、`feed_cache.db`:v2 历史运行数据(如存在) -- 运行日志只通过 MCP stderr 输出,不创建 `feed_mcp.runtime.log` - -候选验证使用 `FEED_BACKEND=recording`: - -- `get_proactive_events` 固定返回 `{"status":"empty"}` -- 不启动 `FeedPoller`,不访问 RSS/RSSHub、不连接 SQLite -- `acknowledge_events` 在 recording 后端 fail-loud -- candidate 只开放只读的 `get_proactive_events` - -正式主动端口使用明确的 typed 结果:拉取返回 `empty` 或 `items`,确认只有全部 -请求 ID 持久成功时才返回 `committed`;异常和部分确认返回 `failure`,不会伪装 -为成功。 - -## 从 v2 迁移 - -先停止占用 workspace 的 Akashic runtime,再运行: - -```bash -PYTHONPATH=/path/to/akashic-agent \ -python scripts/migrate_v2_data.py \ - --workspace /path/to/workspace \ - --marketplace github -``` +正式数据由 Core 分配到 `plugin-data/feed-/`。从 v2 迁移仍使用 +`scripts/migrate_v2_data.py`;它保留源数据并产生可核对的 migration receipt。 -迁移脚本持有 workspace 独占锁,按 `mcp/feed-mcp/` primary、再按 -`backups/feed-plugin-migration-*/feed-mcp/` 最新备份顺序选择第一个含数据的源, -并保留源目录。`feed_mcp.sqlite3` 和其他 SQLite 数据使用在线 backup 后执行 -integrity check;目标存在不同内容时直接失败。 +旧 proactive island 正式切换使用 +`scripts/retire_legacy_feed_backlog.py`。默认 `--plan` 只读输出 Core inventory 与 +Feed 当前完整 backlog 的 count/digest;显式 `--apply` 必须回传这些 exact 值,并在 +Core H2 handoff 前把整批 pre-cutover provider revision 原子标记为 +`cutover_superseded`。该路径不向 Content 提交旧条目,重复执行只接受同一 batch receipt。 -最终 receipt 写入 -`plugin-data/feed-/.feed-v2-migration.json`,逐文件记录 -`source_missing`、`target_only`、`verified` 或 `copied`、源路径、SHA-256、大小和 -SQLite integrity。进程内发布失败会回滚本次新增文件;进程崩溃后重跑会清理残留 -staging、核对同内容目标并完成发布。源数据保留作为 recovery source;receipt -不属于候选验证输入。 +## 验证 -完整外网 RSS E2E 不属于本插件工作流。v3 workflow 固定 Core -`78e50d4dfb3f4348fff37d55d9c9bdd0e002164d` 与 contracts -`4dd69dd621e029e51e99aa428443fa3a4ec1f6cf`,执行插件单元测试、pyright、 -`compileall` 和 `git diff --check`,并以空订阅库走真实 Manager、stdio MCP、 -committed proactive source lease 与 terminate cleanup。 +CI 固定 Core `065dfb7fb37534c57ed500b1ed9b5deb7090cdb0`,运行单元测试、真实 +Manager + stdio MCP + Content + Timer fixture、pyright、compileall 和 +`git diff --check`。Manager fixture 还证明:候选零 Timer/零正式写,发布时旧 +Timer 已取消后新稳定 Root 才接班。 diff --git a/akashic.plugin.toml b/akashic.plugin.toml index e18dcbd..15f59d6 100644 --- a/akashic.plugin.toml +++ b/akashic.plugin.toml @@ -1,6 +1,6 @@ schema_version = 1 name = "feed" -version = "3.0.0" +version = "3.1.0" api_version = 3 entrypoint = "plugin.py" @@ -16,12 +16,20 @@ exclude_data_paths = [ "feed_cache.db", "feed_cache.db-wal", "feed_cache.db-shm", + "feed_mcp.runtime.log", + "feed_mcp.runtime.log.1", + "feed_mcp.runtime.log.2", + "feed_mcp.runtime.log.3", + "feed_source.runtime.log", + "feed_source.runtime.log.1", + "feed_source.runtime.log.2", + "feed_source.runtime.log.3", ".feed-v2-migration.json", ] [[mcp]] name = "feed" command = ["python", "mcp/run_mcp.py"] -required_tools = ["get_proactive_events", "acknowledge_events"] -candidate_read_only_tools = ["get_proactive_events"] +required_tools = ["feed_manage", "feed_query"] +candidate_read_only_tools = [] candidate_env = {FEED_BACKEND = "recording"} diff --git a/content_source.py b/content_source.py new file mode 100644 index 0000000..ac4364b --- /dev/null +++ b/content_source.py @@ -0,0 +1,241 @@ +from __future__ import annotations + +import asyncio +import hashlib +import json +import logging +from collections.abc import Callable, Mapping, Sequence +from datetime import UTC, datetime, timedelta +from logging.handlers import RotatingFileHandler +from pathlib import Path +from typing import Protocol, cast + +from agent.control.timer import TimerHandle, TimerStatus +from agent.plugin_composition import PluginTimers + +from feed_runtime import backend + + +CONTENT_SOURCE_ID = "feed-subscriptions" + + +class BoundContentSource(Protocol): + def submit( + self, batch_id: str, items: Sequence[Mapping[str, object]] + ) -> Mapping[str, object]: ... + + def unsettled(self, limit: int = 100) -> tuple[Mapping[str, object], ...]: ... + + def ack(self, settlement_ref: str) -> Mapping[str, object]: ... + + +class ContentSourceServices(Protocol): + def bind(self, source_id: str) -> BoundContentSource: ... + + +class FeedContentRuntime: + """轮询 Feed、提交精确 revision,并收束 provider ACK。""" + + def __init__( + self, + data_root: Path, + timers: PluginTimers, + content: BoundContentSource, + *, + now: Callable[[], datetime] | None = None, + ) -> None: + self._data_root = data_root + self._timers = timers + self._content = content + self._now = now or (lambda: datetime.now(UTC)) + self._handle: TimerHandle | None = None + self._task: asyncio.Task[None] | None = None + self._closed = False + self._log = logging.Logger("feed-content-source", level=logging.INFO) + self._log.propagate = False + self._log_handler: RotatingFileHandler | None = None + + async def start(self) -> None: + """恢复 source deadline,并只注册一个 Timer。""" + + if self._closed: + raise RuntimeError("Feed Content runtime 已关闭") + if self._handle is not None: + return + self._start_diagnostics() + deadline = await asyncio.to_thread( + backend.content_source_deadline, + data_root=self._data_root, + now=self._aware_now(), + ) + self._arm(deadline) + + async def close(self) -> None: + """收束进行中的轮询、取消等待并释放诊断日志。""" + + if self._closed: + return + self._closed = True + handle = self._handle + task = self._task + if handle is not None: + _ = await handle.cancel() + if task is not None and task is not asyncio.current_task(): + await task + if handle is not None: + await handle.cleanup() + self._handle = None + self._task = None + self._stop_diagnostics() + + def _arm(self, deadline: datetime) -> None: + if self._closed or self._handle is not None: + return + if deadline.tzinfo is None: + raise ValueError("Feed Content deadline 必须包含时区") + handle = self._timers.schedule(deadline) + self._handle = handle + self._task = asyncio.create_task( + self._wait_poll_rearm(handle), name="feed-content-source:poll" + ) + + async def _wait_poll_rearm(self, handle: TimerHandle) -> None: + """消费一个 Timer,完成一次 source 事务后再注册下一次。""" + + try: + receipt = await handle.result() + if receipt.status is TimerStatus.CANCELLED or self._closed: + return + + # 1. 先完成 provider ACK,再获取下一份待处理快照。 + settled = await self._settle_delivered() + + # 2. 由唯一外部轮询 owner 刷新持久 Feed 缓存。 + await asyncio.to_thread( + backend.poll_feeds_only, data_root=self._data_root + ) + items = await asyncio.to_thread( + backend.prepare_content_items, data_root=self._data_root + ) + + # 3. 非空 Content 提交成功后才能推进 source deadline。 + submitted = 0 + if items: + result = self._content.submit(_batch_id(items), items) + inserted = result["inserted"] + if not isinstance(inserted, list): + raise TypeError("Feed Content submit receipt inserted 必须是 list") + submitted = len(inserted) + config = await asyncio.to_thread( + backend.load_config, self._data_root + ) + next_due = self._aware_now() + timedelta( + seconds=config.poll_ttl_seconds + ) + await asyncio.to_thread( + backend.commit_content_source_deadline, + data_root=self._data_root, + deadline=next_due, + ) + self._log.info( + "poll committed items=%d inserted=%d settled=%d next_due=%s", + len(items), + submitted, + settled, + next_due.isoformat(), + ) + finally: + self._handle = None + self._task = None + await handle.cleanup() + if not self._closed: + deadline = await asyncio.to_thread( + backend.content_source_deadline, + data_root=self._data_root, + now=self._aware_now(), + ) + self._arm(deadline) + + async def _settle_delivered(self) -> int: + """用精确 Feed revision 收束每条已投递 Content receipt。""" + + settled = 0 + while rows := self._content.unsettled(100): + for row in rows: + ref = _mapping(row["ref"], "Feed unsettled ref") + result = await asyncio.to_thread( + backend.settle_content_item, + _string(ref["item_id"], "Feed item_id"), + _string(ref["revision"], "Feed revision"), + data_root=self._data_root, + ) + if result.get("status") != "committed" or result.get( + "disposition" + ) not in {"acknowledged", "obsolete_revision", "not_pending"}: + raise RuntimeError(f"Feed provider ACK 未提交: {result!r}") + settlement_ref = _string( + row["settlement_ref"], "Feed settlement_ref" + ) + receipt = self._content.ack(settlement_ref) + if receipt.get("changed") is not True: + raise RuntimeError(f"Feed Content ACK 未提交: {receipt!r}") + settled += 1 + return settled + + def _aware_now(self) -> datetime: + value = self._now() + if value.tzinfo is None: + raise ValueError("Feed Content clock 必须包含时区") + return value.astimezone(UTC) + + def _start_diagnostics(self) -> None: + if self._log_handler is not None: + return + self._data_root.mkdir(parents=True, exist_ok=True) + handler = RotatingFileHandler( + self._data_root / "feed_source.runtime.log", + maxBytes=5 * 1024 * 1024, + backupCount=3, + encoding="utf-8", + ) + handler.setFormatter( + logging.Formatter( + "%(asctime)s %(levelname)-8s %(name)s | %(message)s" + ) + ) + self._log.addHandler(handler) + self._log_handler = handler + + def _stop_diagnostics(self) -> None: + handler = self._log_handler + if handler is None: + return + self._log.removeHandler(handler) + handler.close() + self._log_handler = None + + +def _batch_id(items: Sequence[Mapping[str, object]]) -> str: + identities = [ + { + "item_id": _string(item["item_id"], "Feed item_id"), + "revision": _string(item["revision"], "Feed revision"), + } + for item in items + ] + encoded = json.dumps( + identities, sort_keys=True, separators=(",", ":") + ).encode("utf-8") + return f"feed-content:{hashlib.sha256(encoded).hexdigest()}" + + +def _mapping(value: object, label: str) -> Mapping[str, object]: + if not isinstance(value, Mapping): + raise TypeError(f"{label} 必须是 Mapping") + return cast(Mapping[str, object], value) + + +def _string(value: object, label: str) -> str: + if not isinstance(value, str) or not value: + raise TypeError(f"{label} 必须是非空字符串") + return value diff --git a/feed_runtime/__init__.py b/feed_runtime/__init__.py new file mode 100644 index 0000000..0b23ad9 --- /dev/null +++ b/feed_runtime/__init__.py @@ -0,0 +1 @@ +"""Feed domain runtime shared by ordinary plugin adapters.""" diff --git a/mcp/src/feed_backend.py b/feed_runtime/backend.py similarity index 73% rename from mcp/src/feed_backend.py rename to feed_runtime/backend.py index c94fbf1..0951454 100644 --- a/mcp/src/feed_backend.py +++ b/feed_runtime/backend.py @@ -1,10 +1,10 @@ """ -feed-mcp backend +Feed domain backend shared by the MCP adapter and Content source runtime. 最小实现: 1. 用 sqlite 管理订阅源和条目 2. 按需轮询 RSS/Atom -3. 对外提供 proactive content 事件与基础管理查询能力 +3. 对外提供 Content 条目与基础管理查询能力 """ from __future__ import annotations @@ -60,27 +60,46 @@ class FeedMcpConfig: def _config_path() -> Path: - return Path(__file__).resolve().parent.parent / "feed_mcp.json" + return Path(__file__).resolve().parent.parent / "mcp" / "feed_mcp.json" -def _runtime_root() -> Path: - raw = os.environ.get("AKA_PLUGIN_DATA_DIR", "").strip() - if not raw: - raise RuntimeError("feed backend 缺少 AKA_PLUGIN_DATA_DIR") - path = Path(raw).expanduser() +def _runtime_root(data_root: Path | None = None) -> Path: + if data_root is None: + raw = os.environ.get("AKA_PLUGIN_DATA_DIR", "").strip() + if not raw: + raise RuntimeError("feed backend 缺少 AKA_PLUGIN_DATA_DIR") + path = Path(raw).expanduser() + else: + path = data_root.expanduser() path.mkdir(parents=True, exist_ok=True) return path -def load_config() -> FeedMcpConfig: - runtime_root = _runtime_root() +def _config_values() -> dict[str, Any]: raw = dict(_DEFAULT_CONFIG) path = _config_path() if path.exists(): raw.update(json.loads(path.read_text())) + return raw + + +def _database_path(data_root: Path, raw: dict[str, Any]) -> Path: db_path = Path(str(raw["db_path"])) if not db_path.is_absolute(): - db_path = (runtime_root / db_path).resolve() + db_path = (data_root.expanduser() / db_path).resolve() + return db_path + + +def provider_database_path(data_root: Path) -> Path: + """Resolve the configured Feed database without creating runtime state.""" + + return _database_path(data_root, _config_values()) + + +def load_config(data_root: Path | None = None) -> FeedMcpConfig: + runtime_root = _runtime_root(data_root) + raw = _config_values() + db_path = _database_path(runtime_root, raw) return FeedMcpConfig( db_path=db_path, poll_ttl_seconds=max(60, int(raw["poll_ttl_seconds"])), @@ -211,6 +230,17 @@ def _connect(cfg: FeedMcpConfig) -> sqlite3.Connection: ) """ ) + conn.execute( + """ + CREATE TABLE IF NOT EXISTS content_exports ( + event_id TEXT NOT NULL, + content_hash TEXT NOT NULL, + payload_json TEXT NOT NULL, + created_at TEXT NOT NULL, + PRIMARY KEY(event_id, content_hash) + ) + """ + ) _ensure_rank_stats(conn) _ensure_rank_model(conn, cfg) _normalize_existing_source_urls(conn) @@ -1034,8 +1064,7 @@ def feed_query( conn = _connect(cfg) try: trace_id = _trace_id_short() - # 1. 只读缓存查询:不主动拉取(拉取由 MCP lifespan 的后台 FeedPoller - # 按 poll_ttl_seconds 周期执行,或通过 poll_feeds 显式触发)。 + # 1. 只读缓存查询:不主动拉取(拉取由插件 Timer 统一触发)。 enabled_rows = conn.execute( "SELECT * FROM sources WHERE enabled = 1 ORDER BY added_at DESC" ).fetchall() @@ -1168,10 +1197,10 @@ def feed_query( def _build_display_text(row: sqlite3.Row) -> str: - """为单条 item 生成预格式化展示文本,供 proactive LLM 直接使用。 + """为单条 item 生成稳定展示文本,供 Content 消费方直接使用。 - MCP 侧控制格式,proactive 侧无需感知内容类型细节。 - url 同时保留在独立字段,proactive 会兜底追加以保证溯源链完整。 + Feed 侧控制格式,消费方无需感知内容类型细节。 + url 同时保留在独立字段,保证溯源链完整。 """ source = (row["source_name"] or "").strip() title = (row["title"] or "(无标题)").strip() @@ -1195,7 +1224,7 @@ def _build_display_text(row: sqlite3.Row) -> str: _SOURCE_DIVERSITY_DECAY = 0.5 _SOURCE_DIVERSITY_FLOOR = 0.05 _MISSING_PUBLICATION_CONFIDENCE = 0.03 -_WAKE_ADMISSION_FLOOR = 0.02 +_CONTENT_ADMISSION_FLOOR = 0.02 def _tokenize_rank_text(*parts: object) -> list[str]: @@ -1395,7 +1424,7 @@ def _freshness_score(row: sqlite3.Row, cfg: FeedMcpConfig, now: datetime) -> tup ) -def _wake_admission_score(features: dict[str, float]) -> float: +def _content_admission_score(features: dict[str, float]) -> float: interest = min(0.999, max(0.0, float(features.get("interest", 0.0)))) freshness = min(1.0, max(0.0, float(features.get("freshness", 0.0)))) return -math.log1p(-interest) * freshness @@ -1758,13 +1787,10 @@ def _record_rank_impressions( ) -def poll_feeds_only() -> None: - """按需轮询所有启用源(尊重 poll_ttl_seconds),不返回内容。 - 由 MCP lifespan 按固定周期调用,与 get_proactive_events 完全解耦。 - 单源失败已在 _poll_rows 内部隔离;系统级异常(DB 不可用、配置损坏等)直接上抛, - 由调用方决定如何处理,避免故障被静默吞掉。 - """ - cfg = load_config() +def poll_feeds_only(*, data_root: Path | None = None) -> None: + """Poll every due Feed source once and preserve per-source outcomes.""" + + cfg = load_config(data_root) conn = _connect(cfg) try: trace_id = _trace_id_short() @@ -1786,11 +1812,13 @@ def poll_feeds_only() -> None: conn.close() -def get_proactive_events(*, offset: int = 0, limit: int = 50) -> list[dict[str, Any]]: - cfg = load_config() +def prepare_content_items(*, data_root: Path) -> tuple[dict[str, object], ...]: + """Freeze complete current revisions for idempotent Content submission.""" + + cfg = load_config(data_root) conn = _connect(cfg) try: - # 纯 DB 查询,不触发轮询。MCP lifespan 独立维护缓存 freshness。 + # 1. Rank every current unacknowledged item. now = _now() rows = conn.execute( """ @@ -1804,7 +1832,8 @@ def get_proactive_events(*, offset: int = 0, limit: int = 50) -> list[dict[str, i.url, i.author, i.published_at, - i.first_seen_at + i.first_seen_at, + i.content_hash FROM items i LEFT JOIN acked_items a ON a.event_id = i.event_id WHERE a.event_id IS NULL @@ -1819,97 +1848,588 @@ def get_proactive_events(*, offset: int = 0, limit: int = 50) -> list[dict[str, admitted = [] for row in rows: score, features = ranked_by_id.get(str(row["event_id"]), (0.0, {})) - admission_score = _wake_admission_score(features) - if admission_score < _WAKE_ADMISSION_FLOOR: + admission_score = _content_admission_score(features) + if admission_score < _CONTENT_ADMISSION_FLOOR: continue visible_features = dict(features) - visible_features["wake_admission_score"] = admission_score + visible_features["content_admission_score"] = admission_score admitted.append((row, score, visible_features)) - selected = admitted[max(0, offset):max(0, offset) + max(1, limit)] - return [ - { - "event_id": row["event_id"], + + # 2. Freeze one complete payload for each exact content revision. + prepared: list[dict[str, object]] = [] + for row, score, features in admitted: + event_id = str(row["event_id"]) + revision = str(row["content_hash"]) + payload = { "kind": "content", "source_type": row["source_type"], "source_id": row["source_id"], "source_name": row["source_name"], "title": row["title"], + "content": row["content"], "url": row["url"], + "author": row["author"], "published_at": row["published_at"], "first_seen_at": row["first_seen_at"], "preprocess_score": round(score, 6), "preprocess_features": features, } - for row, score, features in selected - ] + conn.execute( + """ + INSERT OR IGNORE INTO content_exports( + event_id, content_hash, payload_json, created_at + ) VALUES (?, ?, ?, ?) + """, + ( + event_id, + revision, + json.dumps(payload, sort_keys=True, separators=(",", ":")), + now.isoformat(), + ), + ) + frozen = conn.execute( + """ + SELECT payload_json FROM content_exports + WHERE event_id = ? AND content_hash = ? + """, + (event_id, revision), + ).fetchone() + if frozen is None: + raise RuntimeError("Feed Content export 冻结失败") + frozen_payload = json.loads(str(frozen["payload_json"])) + prepared.append( + { + "item_id": event_id, + "revision": revision, + "payload": frozen_payload, + "not_before": str( + frozen_payload.get("published_at") + or frozen_payload["first_seen_at"] + ), + "requires_ack": True, + } + ) + conn.commit() + return tuple(prepared) finally: conn.close() -def _interest_ok_from_feedback(feedback: str | None) -> int | None: - if feedback is None or feedback == "consumed": - return None - if feedback == "interesting": - return 1 - if feedback == "not_interesting": - return 0 - raise ValueError(f"invalid feedback: {feedback}") +def plan_content_backlog(*, data_root: Path) -> tuple[dict[str, str], ...]: + """Read the exact currently admissible Feed backlog without creating state.""" + raw = _config_values() + cfg = FeedMcpConfig( + db_path=provider_database_path(data_root), + poll_ttl_seconds=max(60, int(raw["poll_ttl_seconds"])), + item_retention_hours=max(1, int(raw["item_retention_hours"])), + max_items_per_source=max(1, int(raw["max_items_per_source"])), + max_content_events=max(1, int(raw["max_content_events"])), + rank_mode=str(raw["rank_mode"]), + rank_impression_limit=max(1, int(raw["rank_impression_limit"])), + rank_model_learning_rate=max(0.001, float(raw["rank_model_learning_rate"])), + ) + wal = cfg.db_path.with_name(cfg.db_path.name + "-wal") + if wal.is_file() and wal.stat().st_size > 0: + raise RuntimeError("Feed provider backlog plan requires a checkpoint") + uri = cfg.db_path.resolve().as_uri() + "?mode=ro&immutable=1" + conn = sqlite3.connect(uri, uri=True) + conn.row_factory = sqlite3.Row + try: + # 1. Reuse the production admission calculation over a frozen provider DB. + _ = conn.execute("PRAGMA query_only = ON") + now = _now() + rows = conn.execute( + """ + SELECT + i.event_id, i.source_id, i.source_type, i.source_name, + i.title, i.content, i.url, i.author, i.published_at, + i.first_seen_at, i.content_hash + FROM items i + LEFT JOIN acked_items a ON a.event_id = i.event_id + WHERE a.event_id IS NULL + ORDER BY i.source_name ASC, + i.published_at IS NULL ASC, + i.published_at DESC, + i.first_seen_at DESC + """ + ).fetchall() + ranked = _rank_rows(conn, cfg, list(rows), now) + ranked_by_id = { + str(row["event_id"]): features for row, _score, features in ranked + } -def acknowledge_events( - event_ids: list[str], - feedback: str | None = None, -) -> dict[str, list[str]]: - cfg = load_config() + # 2. Return only stable provider identities in deterministic order. + planned = [ + { + "event_id": str(row["event_id"]), + "revision": str(row["content_hash"]), + } + for row in rows + if _content_admission_score( + ranked_by_id.get(str(row["event_id"]), {}) + ) + >= _CONTENT_ADMISSION_FLOOR + ] + return tuple(sorted(planned, key=lambda item: item["event_id"])) + finally: + conn.close() + + +def supersede_content_backlog( + batch_id: str, + items: tuple[dict[str, str], ...], + *, + data_root: Path, +) -> dict[str, object]: + """Atomically ACK one frozen pre-cutover backlog and retain its receipt.""" + + if not batch_id or not items: + raise ValueError("Feed cutover supersession requires a non-empty batch") + normalized = tuple( + sorted( + ( + {"event_id": item["event_id"], "revision": item["revision"]} + for item in items + ), + key=lambda item: item["event_id"], + ) + ) + items_json = json.dumps(normalized, sort_keys=True, separators=(",", ":")) + items_digest = hashlib.sha256(items_json.encode()).hexdigest() + cfg = load_config(data_root) conn = _connect(cfg) now = _now() - acked: list[str] = [] - failed: list[str] = [] try: - interest_ok = _interest_ok_from_feedback(feedback) - for event_id in event_ids: - try: - conn.execute( - """ - INSERT INTO acked_items (event_id, acked_at, expires_at) - VALUES (?, ?, ?) - ON CONFLICT(event_id) DO UPDATE SET - acked_at=excluded.acked_at, - expires_at=excluded.expires_at - """, - ( - event_id, - now.isoformat(), - (now + timedelta(hours=cfg.item_retention_hours)).isoformat(), - ), + # 1. Resume only the same completed batch identity. + _ensure_legacy_backlog_receipts(conn) + existing = conn.execute( + "SELECT * FROM legacy_backlog_supersession_receipts WHERE batch_id = ?", + (batch_id,), + ).fetchone() + if existing is not None: + if ( + str(existing["items_digest"]) != items_digest + or int(existing["item_count"]) != len(normalized) + or str(existing["items_json"]) != items_json + ): + raise RuntimeError("Feed cutover supersession batch conflict") + receipt = _legacy_backlog_receipt(existing) + if not _backlog_rows_are_acked(conn, normalized): + raise RuntimeError("Feed cutover supersession receipt is unverified") + return receipt + + # 2. Lock and reverify every exact provider revision before ACK writes. + conn.commit() + conn.execute("BEGIN IMMEDIATE") + for item in normalized: + current = conn.execute( + "SELECT content_hash FROM items WHERE event_id = ?", + (item["event_id"],), + ).fetchone() + if current is None or str(current["content_hash"]) != item["revision"]: + raise RuntimeError( + f"Feed cutover provider revision changed: {item['event_id']}" ) - if interest_ok is not None: - conn.execute( - """ - UPDATE items - SET interest_ok = ?, interest_scored_at = ? - WHERE event_id = ? - """, - (interest_ok, now.isoformat(), event_id), - ) - row = conn.execute( - """ - SELECT - event_id, source_id, source_name, author, title, content, published_at, - first_seen_at, interest_ok, interest_scored_at - FROM items - WHERE event_id = ? - """, - (event_id,), - ).fetchone() - if row is not None and interest_ok is not None: - _update_rank_stats_for_row(conn, row, interest_ok, now.isoformat()) - _update_rank_model_for_row(conn, cfg, row, interest_ok, now.isoformat()) - acked.append(event_id) - except Exception: - logger.exception("feed ack failed: %s", event_id) - failed.append(event_id) + + # 3. ACK the complete set and commit its durable batch receipt atomically. + acked_at = now.isoformat() + expires_at = (now + timedelta(hours=cfg.item_retention_hours)).isoformat() + conn.executemany( + """ + INSERT INTO acked_items(event_id, acked_at, expires_at) + VALUES (?, ?, ?) + ON CONFLICT(event_id) DO UPDATE SET + acked_at=excluded.acked_at, + expires_at=excluded.expires_at + """, + ( + (item["event_id"], acked_at, expires_at) + for item in normalized + ), + ) + receipt_id = f"feed-cutover:{items_digest}" + conn.execute( + """ + INSERT INTO legacy_backlog_supersession_receipts( + batch_id, receipt_id, items_digest, item_count, + items_json, committed_at + ) VALUES (?, ?, ?, ?, ?, ?) + """, + ( + batch_id, + receipt_id, + items_digest, + len(normalized), + items_json, + now.isoformat(), + ), + ) + conn.commit() + row = conn.execute( + "SELECT * FROM legacy_backlog_supersession_receipts WHERE batch_id = ?", + (batch_id,), + ).fetchone() + if row is None: + raise RuntimeError("Feed cutover supersession receipt commit missing") + return _legacy_backlog_receipt(row) + finally: + conn.close() + + +def verify_content_backlog_supersession( + batch_id: str, + items: tuple[dict[str, str], ...], + *, + data_root: Path, +) -> bool: + """Verify the exact batch receipt and every provider suppression row read-only.""" + + path = provider_database_path(data_root) + wal = path.with_name(path.name + "-wal") + if wal.is_file() and wal.stat().st_size > 0: + return False + uri = path.resolve().as_uri() + "?mode=ro&immutable=1" + conn = sqlite3.connect(uri, uri=True) + conn.row_factory = sqlite3.Row + try: + _ = conn.execute("PRAGMA query_only = ON") + table = conn.execute( + "SELECT 1 FROM sqlite_master WHERE type='table' " + "AND name='legacy_backlog_supersession_receipts'" + ).fetchone() + if table is None: + return False + row = conn.execute( + "SELECT * FROM legacy_backlog_supersession_receipts WHERE batch_id = ?", + (batch_id,), + ).fetchone() + if row is None: + return False + normalized = tuple(sorted(items, key=lambda item: item["event_id"])) + items_json = json.dumps(normalized, sort_keys=True, separators=(",", ":")) + if ( + str(row["items_json"]) != items_json + or str(row["items_digest"]) + != hashlib.sha256(items_json.encode()).hexdigest() + or int(row["item_count"]) != len(normalized) + ): + return False + return _backlog_rows_are_acked(conn, normalized) + finally: + conn.close() + + +def read_content_backlog_supersession( + batch_id: str, + *, + data_root: Path, +) -> tuple[dict[str, object], tuple[dict[str, str], ...]] | None: + """Read one completed cutover batch and its exact item identities.""" + + path = provider_database_path(data_root) + wal = path.with_name(path.name + "-wal") + if wal.is_file() and wal.stat().st_size > 0: + raise RuntimeError("Feed cutover receipt read requires a checkpoint") + uri = path.resolve().as_uri() + "?mode=ro&immutable=1" + conn = sqlite3.connect(uri, uri=True) + conn.row_factory = sqlite3.Row + try: + _ = conn.execute("PRAGMA query_only = ON") + table = conn.execute( + "SELECT 1 FROM sqlite_master WHERE type='table' " + "AND name='legacy_backlog_supersession_receipts'" + ).fetchone() + if table is None: + return None + row = conn.execute( + "SELECT * FROM legacy_backlog_supersession_receipts WHERE batch_id = ?", + (batch_id,), + ).fetchone() + if row is None: + return None + decoded = json.loads(str(row["items_json"])) + if not isinstance(decoded, list): + raise RuntimeError("Feed cutover receipt items must be a list") + parsed: list[dict[str, str]] = [] + for item in decoded: + if not isinstance(item, dict): + raise RuntimeError("Feed cutover receipt item must be an object") + event_id = item.get("event_id") + revision = item.get("revision") + if not isinstance(event_id, str) or not isinstance(revision, str): + raise RuntimeError("Feed cutover receipt item identity is invalid") + parsed.append({"event_id": event_id, "revision": revision}) + return _legacy_backlog_receipt(row), tuple(parsed) + finally: + conn.close() + + +def _ensure_legacy_backlog_receipts(conn: sqlite3.Connection) -> None: + conn.execute( + """ + CREATE TABLE IF NOT EXISTS legacy_backlog_supersession_receipts( + batch_id TEXT PRIMARY KEY, + receipt_id TEXT NOT NULL UNIQUE, + items_digest TEXT NOT NULL, + item_count INTEGER NOT NULL, + items_json TEXT NOT NULL, + committed_at TEXT NOT NULL + ) + """ + ) + + +def _legacy_backlog_receipt(row: sqlite3.Row) -> dict[str, object]: + return { + "batch_id": str(row["batch_id"]), + "receipt_id": str(row["receipt_id"]), + "items_digest": str(row["items_digest"]), + "item_count": int(row["item_count"]), + "committed_at": str(row["committed_at"]), + } + + +def _backlog_rows_are_acked( + conn: sqlite3.Connection, + items: tuple[dict[str, str], ...], +) -> bool: + return all( + conn.execute( + "SELECT 1 FROM acked_items a JOIN items i USING(event_id) " + "WHERE a.event_id = ? AND i.content_hash = ?", + (item["event_id"], item["revision"]), + ).fetchone() + is not None + for item in items + ) + + +def settle_content_item( + event_id: str, + revision: str, + *, + data_root: Path, +) -> dict[str, str]: + """Ensure one exact exported revision is no longer pending in Feed.""" + + cfg = load_config(data_root) + conn = _connect(cfg) + now = _now() + try: + current = conn.execute( + "SELECT content_hash FROM items WHERE event_id = ?", (event_id,) + ).fetchone() + if current is None: + return {"status": "committed", "disposition": "not_pending"} + if str(current["content_hash"]) != revision: + return {"status": "committed", "disposition": "obsolete_revision"} + conn.execute( + """ + INSERT INTO acked_items (event_id, acked_at, expires_at) + VALUES (?, ?, ?) + ON CONFLICT(event_id) DO UPDATE SET + acked_at=excluded.acked_at, + expires_at=excluded.expires_at + """, + ( + event_id, + now.isoformat(), + (now + timedelta(hours=cfg.item_retention_hours)).isoformat(), + ), + ) + conn.commit() + return {"status": "committed", "disposition": "acknowledged"} + finally: + conn.close() + + +def settle_legacy_ack( + event_id: str, + revision: str, + action: str, + source_digest: str, + *, + data_root: Path, +) -> dict[str, str]: + """Commit one legacy Wake ACK and retain its target-owned receipt.""" + + cfg = load_config(data_root) + conn = _connect(cfg) + now = _now() + receipt_id = f"feed-legacy-ack:{source_digest}" + try: + # 1. Reuse a completed handoff without extending the provider ACK. + _ensure_legacy_ack_receipts(conn) + existing = conn.execute( + "SELECT * FROM legacy_ack_handoff_receipts WHERE receipt_id = ?", + (receipt_id,), + ).fetchone() + if existing is not None: + identity = tuple( + str(existing[field]) + for field in ("source_digest", "event_id", "revision", "action") + ) + if identity != (source_digest, event_id, revision, action): + raise RuntimeError("Feed legacy ACK receipt identity conflict") + return _legacy_ack_receipt(existing) + + # 2. Commit one exact provider ACK and its durable receipt atomically. + acked_at, expires_at = _commit_legacy_provider_ack( + conn, cfg, event_id, revision, now + ) + _insert_legacy_ack_receipt( + conn, + receipt_id, + source_digest, + event_id, + revision, + action, + acked_at, + expires_at, + now, + ) + conn.commit() + row = conn.execute( + "SELECT * FROM legacy_ack_handoff_receipts WHERE receipt_id = ?", + (receipt_id,), + ).fetchone() + if row is None: + raise RuntimeError("Feed legacy ACK receipt commit missing") + return _legacy_ack_receipt(row) + finally: + conn.close() + + +def _ensure_legacy_ack_receipts(conn: sqlite3.Connection) -> None: + conn.execute( + """ + CREATE TABLE IF NOT EXISTS legacy_ack_handoff_receipts( + receipt_id TEXT PRIMARY KEY, + source_digest TEXT NOT NULL UNIQUE, + event_id TEXT NOT NULL, + revision TEXT NOT NULL, + action TEXT NOT NULL, + acked_at TEXT NOT NULL, + expires_at TEXT NOT NULL, + committed_at TEXT NOT NULL + ) + """ + ) + + +def _commit_legacy_provider_ack( + conn: sqlite3.Connection, + cfg: FeedMcpConfig, + event_id: str, + revision: str, + now: datetime, +) -> tuple[str, str]: + """Preserve a live provider ACK or establish one new retention window.""" + + current = conn.execute( + "SELECT content_hash FROM items WHERE event_id = ?", (event_id,) + ).fetchone() + if current is None: + raise RuntimeError(f"Feed legacy ACK provider item missing: {event_id}") + if str(current["content_hash"]) != revision: + raise RuntimeError(f"Feed legacy ACK revision changed: {event_id}") + acknowledgement = conn.execute( + "SELECT acked_at, expires_at FROM acked_items WHERE event_id = ?", + (event_id,), + ).fetchone() + if acknowledgement is not None and datetime.fromisoformat( + str(acknowledgement["expires_at"]) + ) > now: + return str(acknowledgement["acked_at"]), str(acknowledgement["expires_at"]) + acked_at = now.isoformat() + expires_at = (now + timedelta(hours=cfg.item_retention_hours)).isoformat() + conn.execute( + """ + INSERT INTO acked_items(event_id, acked_at, expires_at) + VALUES (?, ?, ?) + ON CONFLICT(event_id) DO UPDATE SET + acked_at=excluded.acked_at, + expires_at=excluded.expires_at + """, + (event_id, acked_at, expires_at), + ) + return acked_at, expires_at + + +def _insert_legacy_ack_receipt( + conn: sqlite3.Connection, + receipt_id: str, + source_digest: str, + event_id: str, + revision: str, + action: str, + acked_at: str, + expires_at: str, + now: datetime, +) -> None: + conn.execute( + """ + INSERT INTO legacy_ack_handoff_receipts( + receipt_id, source_digest, event_id, revision, action, + acked_at, expires_at, committed_at + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?) + """, + ( + receipt_id, + source_digest, + event_id, + revision, + action, + acked_at, + expires_at, + now.isoformat(), + ), + ) + + +def _legacy_ack_receipt(row: sqlite3.Row) -> dict[str, str]: + return { + field: str(row[field]) + for field in ( + "receipt_id", + "source_digest", + "event_id", + "revision", + "action", + "acked_at", + "expires_at", + ) + } + + +def content_source_deadline(*, data_root: Path, now: datetime) -> datetime: + """Return the durable next source deadline, defaulting to immediate work.""" + + cfg = load_config(data_root) + conn = _connect(cfg) + try: + row = conn.execute( + "SELECT value FROM metadata WHERE key = 'content_source_next_due'" + ).fetchone() + return now if row is None else datetime.fromisoformat(str(row["value"])) + finally: + conn.close() + + +def commit_content_source_deadline(*, data_root: Path, deadline: datetime) -> None: + """Commit the next source-owned deadline without appending poll history.""" + + cfg = load_config(data_root) + conn = _connect(cfg) + try: + conn.execute( + """ + INSERT INTO metadata(key, value) VALUES('content_source_next_due', ?) + ON CONFLICT(key) DO UPDATE SET value=excluded.value + """, + (deadline.isoformat(),), + ) conn.commit() - return {"acknowledged": acked, "failed": failed} finally: conn.close() diff --git a/legacy_handoff.py b/legacy_handoff.py new file mode 100644 index 0000000..0decb0f --- /dev/null +++ b/legacy_handoff.py @@ -0,0 +1,380 @@ +from __future__ import annotations + +import hashlib +import json +import sqlite3 +from collections.abc import Mapping, Sequence +from contextlib import closing +from pathlib import Path +from typing import Protocol, cast + +from agent.migrations.proactive_island import ( + AdapterPlan, + HandoffBlocked, + LegacyFact, + LegacyFactKind, + TargetReceipt, +) +from agent.migrations.proactive_island.handoff import receipt_digest + +from content_source import CONTENT_SOURCE_ID +from feed_runtime import backend + + +LEGACY_SOURCE_ID = "feed@github:subscriptions" + + +class BoundContentSource(Protocol): + def submit( + self, batch_id: str, items: Sequence[Mapping[str, object]] + ) -> Mapping[str, object]: ... + + def read_submission(self, batch_id: str) -> Mapping[str, object] | None: ... + + def read_revision( + self, item_id: str, revision: str + ) -> Mapping[str, object] | None: ... + + +class FeedLegacyHandoffAdapter: + """Move exact legacy Feed reservoir facts into the existing Content source.""" + + def __init__( + self, + feed_data_root: Path, + content: BoundContentSource | None, + *, + supersede_backlog: bool = False, + ) -> None: + self._data_root = feed_data_root + self._provider_db = backend.provider_database_path(feed_data_root) + self._content = content + self._supersede_backlog = supersede_backlog + + @classmethod + def for_cutover(cls, feed_data_root: Path) -> FeedLegacyHandoffAdapter: + """Build the explicit pre-cutover backlog supersession owner.""" + + return cls(feed_data_root, None, supersede_backlog=True) + + def accepts(self, fact: LegacyFact) -> bool: + return ( + fact.kind in {LegacyFactKind.WAKE_SOURCE_ITEM, LegacyFactKind.WAKE_ACK} + and fact.source_identity == LEGACY_SOURCE_ID + ) + + def plan(self, fact: LegacyFact) -> AdapterPlan: + """Resolve one provider-owned revision without mounting or writing Content.""" + + row = _fact_row(fact) + provider = self._provider_item(_text(row, "source_event_id")) + return AdapterPlan( + _target_identity( + fact.kind, + provider, + supersede_backlog=self._supersede_backlog, + ) + ) + + def apply(self, fact: LegacyFact, plan: AdapterPlan) -> TargetReceipt: + """Submit the exact planned target and return its normalized durable receipt.""" + + # 1. Re-read the owner row and reject a revision change after planning. + row = _fact_row(fact) + provider = self._provider_item(_text(row, "source_event_id")) + target_identity = _target_identity( + fact.kind, + provider, + supersede_backlog=self._supersede_backlog, + ) + if plan.target_identity != target_identity: + raise RuntimeError("Feed handoff target identity drift after plan") + + # 2. Each fact kind commits through its target owner's durable primitive. + if fact.kind is LegacyFactKind.WAKE_ACK: + acknowledgement = backend.settle_legacy_ack( + _text(provider, "event_id"), + _text(provider, "content_hash"), + _text(row, "action"), + fact.source_digest, + data_root=self._data_root, + ) + return _ack_receipt(fact, target_identity, acknowledgement) + + # 3. Explicit cutover supersession ACKs the provider without creating Content. + if self._supersede_backlog: + acknowledgement = backend.settle_legacy_ack( + _text(provider, "event_id"), + _text(provider, "content_hash"), + "cutover_superseded", + fact.source_digest, + data_root=self._data_root, + ) + return _ack_receipt(fact, target_identity, acknowledgement) + + # 4. A fact-stable batch makes target-before-marker replay idempotent. + if self._content is None: + raise RuntimeError("Feed Content target is unavailable") + batch_id = _batch_id(fact) + item = _content_item(row, provider) + content_receipt = self._content.submit(batch_id, (item,)) + receipt_id = _text(content_receipt, "receipt_id") + normalized = _receipt_payload(fact, target_identity, content_receipt) + return TargetReceipt( + receipt_id=receipt_id, + receipt_digest=receipt_digest(normalized), + target_identity=target_identity, + ) + + def verify(self, fact: LegacyFact, receipt: TargetReceipt) -> bool: + """Verify provider lineage and checkpointed Content facts without writing.""" + + try: + plan = self.plan(fact) + except HandoffBlocked: + return False + if receipt.target_identity != plan.target_identity: + return False + row = _fact_row(fact) + provider = self._provider_item(_text(row, "source_event_id")) + if fact.kind is LegacyFactKind.WAKE_ACK or self._supersede_backlog: + acknowledgement = self._provider_ack(fact.source_digest) + if acknowledgement is None: + return False + if self._supersede_backlog and not self._provider_is_acked( + _text(provider, "event_id") + ): + return False + expected = _ack_receipt(fact, plan.target_identity, acknowledgement) + return receipt == expected + if self._content is None: + raise RuntimeError("Feed Content target is unavailable") + item = _content_item(row, provider) + batch_id = _batch_id(fact) + submission = self._content.read_submission(batch_id) + revision = self._content.read_revision( + _text(item, "item_id"), _text(item, "revision") + ) + if submission is None or revision is None: + return False + normalized = _receipt_payload(fact, plan.target_identity, submission) + return ( + receipt.receipt_id == submission.get("receipt_id") + and receipt.receipt_digest == receipt_digest(normalized) + and _revision_matches(revision, item) + ) + + def _provider_ack(self, source_digest: str) -> dict[str, object] | None: + """Read one retained legacy ACK receipt without opening a writer.""" + + wal = self._provider_db.with_name(self._provider_db.name + "-wal") + if wal.is_file() and wal.stat().st_size > 0: + raise HandoffBlocked("feed_provider_checkpoint_required") + uri = self._provider_db.resolve().as_uri() + "?mode=ro&immutable=1" + with closing(sqlite3.connect(uri, uri=True)) as connection: + connection.row_factory = sqlite3.Row + _ = connection.execute("PRAGMA query_only = ON") + table = connection.execute( + "SELECT 1 FROM sqlite_master " + "WHERE type='table' AND name='legacy_ack_handoff_receipts'" + ).fetchone() + if table is None: + return None + result = connection.execute( + "SELECT * FROM legacy_ack_handoff_receipts WHERE source_digest = ?", + (source_digest,), + ).fetchone() + return None if result is None else {key: result[key] for key in result.keys()} + + def _provider_item(self, event_id: str) -> dict[str, object]: + """Read one exact Feed row through a query-only SQLite connection.""" + + if not self._provider_db.is_file(): + raise HandoffBlocked("feed_provider_database_missing") + wal = self._provider_db.with_name(self._provider_db.name + "-wal") + if wal.is_file() and wal.stat().st_size > 0: + raise HandoffBlocked("feed_provider_checkpoint_required") + uri = self._provider_db.resolve().as_uri() + "?mode=ro&immutable=1" + with closing(sqlite3.connect(uri, uri=True)) as connection: + connection.row_factory = sqlite3.Row + _ = connection.execute("PRAGMA query_only = ON") + result = connection.execute( + """ + SELECT event_id, source_id, source_type, source_name, title, + content, url, author, published_at, first_seen_at, + content_hash + FROM items WHERE event_id = ? + """, + (event_id,), + ).fetchone() + if result is None: + raise HandoffBlocked(f"feed_provider_item_missing:{event_id}") + return {key: result[key] for key in result.keys()} + + def _provider_is_acked(self, event_id: str) -> bool: + """Verify the cutover suppression row without opening a writer.""" + + uri = self._provider_db.resolve().as_uri() + "?mode=ro&immutable=1" + with closing(sqlite3.connect(uri, uri=True)) as connection: + _ = connection.execute("PRAGMA query_only = ON") + result = connection.execute( + "SELECT 1 FROM acked_items WHERE event_id = ?", (event_id,) + ).fetchone() + return result is not None + + +def _fact_row(fact: LegacyFact) -> dict[str, object]: + if fact.kind not in {LegacyFactKind.WAKE_SOURCE_ITEM, LegacyFactKind.WAKE_ACK}: + raise TypeError("Feed handoff received another legacy fact kind") + if fact.source_identity != LEGACY_SOURCE_ID: + raise TypeError("Feed handoff received another legacy source owner") + if hashlib.sha256(fact.opaque).hexdigest() != fact.source_digest: + raise RuntimeError("Feed legacy source digest mismatch") + decoded = json.loads(fact.opaque) + if not isinstance(decoded, dict): + raise TypeError("Feed legacy reservoir row must be an object") + row = cast(dict[str, object], decoded) + event_id = _text(row, "source_event_id") + if not fact.locator.endswith(f":{_text(row, 'item_id')}"): + raise RuntimeError("Feed legacy locator does not match item_id") + if fact.kind is LegacyFactKind.WAKE_SOURCE_ITEM: + payload = _payload(row) + if payload.get("event_id") != event_id or payload.get("kind") != "content": + raise RuntimeError("Feed legacy payload identity mismatch") + elif _text(row, "source_id") != LEGACY_SOURCE_ID or _text(row, "action") not in { + "consume", + "expire", + }: + raise RuntimeError("Feed legacy ACK identity mismatch") + return row + + +def _content_item( + legacy: Mapping[str, object], provider: Mapping[str, object] +) -> dict[str, object]: + event_id = _text(provider, "event_id") + if _text(legacy, "source_event_id") != event_id: + raise RuntimeError("Feed provider join identity mismatch") + source_payload = _payload(legacy) + payload: dict[str, object] = { + "kind": "content", + "source_type": provider["source_type"], + "source_id": provider["source_id"], + "source_name": provider["source_name"], + "title": provider["title"], + "content": provider["content"], + "url": provider["url"], + "author": provider["author"], + "published_at": provider["published_at"], + "first_seen_at": provider["first_seen_at"], + "preprocess_score": source_payload.get( + "preprocess_score", legacy["preprocess_score"] + ), + "preprocess_features": source_payload.get("preprocess_features", {}), + } + return { + "item_id": event_id, + "revision": _text(provider, "content_hash"), + "payload": payload, + "not_before": str(provider["published_at"] or provider["first_seen_at"]), + "requires_ack": True, + } + + +def _target_identity( + kind: LegacyFactKind, + provider: Mapping[str, object], + *, + supersede_backlog: bool = False, +) -> str: + if kind is LegacyFactKind.WAKE_SOURCE_ITEM and supersede_backlog: + prefix = "feed-cutover-superseded" + else: + prefix = "content" if kind is LegacyFactKind.WAKE_SOURCE_ITEM else "feed-ack" + return ( + f"{prefix}:{CONTENT_SOURCE_ID}:{_text(provider, 'event_id')}:" + f"{_text(provider, 'content_hash')}" + ) + + +def _ack_receipt( + fact: LegacyFact, + target_identity: str, + acknowledgement: Mapping[str, object], +) -> TargetReceipt: + normalized = { + "legacy_locator": fact.locator, + "legacy_source_digest": fact.source_digest, + "target_identity": target_identity, + "acknowledgement": { + key: acknowledgement[key] + for key in ( + "receipt_id", + "source_digest", + "event_id", + "revision", + "action", + "acked_at", + "expires_at", + ) + }, + } + return TargetReceipt( + receipt_id=_text(acknowledgement, "receipt_id"), + receipt_digest=receipt_digest(normalized), + target_identity=target_identity, + ) + + +def _batch_id(fact: LegacyFact) -> str: + encoded = f"{fact.locator}\x00{fact.source_digest}".encode("utf-8") + return f"feed-legacy:{hashlib.sha256(encoded).hexdigest()}" + + +def _receipt_payload( + fact: LegacyFact, + target_identity: str, + content_receipt: Mapping[str, object], +) -> dict[str, object]: + return { + "schema_version": 1, + "legacy_locator": fact.locator, + "legacy_source_digest": fact.source_digest, + "legacy_source_identity": fact.source_identity, + "target_identity": target_identity, + "content_receipt": dict(content_receipt), + } + + +def _revision_matches( + revision: Mapping[str, object], item: Mapping[str, object] +) -> bool: + return ( + revision.get("ref") + == { + "source_id": CONTENT_SOURCE_ID, + "item_id": item["item_id"], + "revision": item["revision"], + } + and revision.get("payload") == item["payload"] + and revision.get("not_before") == item["not_before"] + and revision.get("requires_ack") is True + ) + + +def _payload(row: Mapping[str, object]) -> dict[str, object]: + raw = _text(row, "payload_json") + value = json.loads(raw) + if not isinstance(value, dict): + raise TypeError("Feed legacy payload_json must contain an object") + return cast(dict[str, object], value) + + +def _text(row: Mapping[str, object], field: str) -> str: + value = row[field] + if not isinstance(value, str) or not value: + raise TypeError(f"Feed {field} must be a non-empty string") + return value + + +__all__ = ["FeedLegacyHandoffAdapter"] diff --git a/mcp/run_mcp.py b/mcp/run_mcp.py index 3a2a6c7..391bfda 100755 --- a/mcp/run_mcp.py +++ b/mcp/run_mcp.py @@ -4,6 +4,7 @@ import logging import os import sys +from logging.handlers import RotatingFileHandler from pathlib import Path @@ -16,12 +17,21 @@ def _runtime_dir() -> Path: return Path(raw).expanduser() -def _setup_logging() -> None: - """把 MCP 日志绑定到 stderr,不创建运行日志文件。""" +def _setup_logging(runtime_dir: Path) -> None: + """将诊断写入 stderr 和三个有界本地轮转文件。""" + runtime_dir.mkdir(parents=True, exist_ok=True) formatter = logging.Formatter( "%(asctime)s %(levelname)-8s %(name)s | %(message)s" ) + file_handler = RotatingFileHandler( + runtime_dir / "feed_mcp.runtime.log", + maxBytes=5 * 1024 * 1024, + backupCount=3, + encoding="utf-8", + ) + file_handler.setLevel(logging.INFO) + file_handler.setFormatter(formatter) stream_handler = logging.StreamHandler(sys.stderr) stream_handler.setLevel(logging.INFO) stream_handler.setFormatter(formatter) @@ -29,21 +39,20 @@ def _setup_logging() -> None: root = logging.getLogger() root.setLevel(logging.INFO) root.handlers.clear() + root.addHandler(file_handler) root.addHandler(stream_handler) def main() -> None: - # 1. 校验 Core 注入的数据目录变量,但不创建或读取运行态文件。 - _runtime_dir() - - # 2. 切换到脚本目录,保证相对代码路径稳定。 + # 1. 暴露插件根目录中的共享 Feed domain package。 script_dir = Path(__file__).resolve().parent os.chdir(script_dir) - if str(script_dir) not in sys.path: - sys.path.insert(0, str(script_dir)) + for path in (script_dir.parent, script_dir): + if str(path) not in sys.path: + sys.path.insert(0, str(path)) - # 3. 启动 MCP stdio 服务;日志只经过 stderr。 - _setup_logging() + # 2. 用有界诊断日志启动用户驱动的 MCP adapter。 + _setup_logging(_runtime_dir()) from src.mcp_bridge import create_mcp_server create_mcp_server().run(transport="stdio") diff --git a/mcp/scripts/backfill_x_published_at.py b/mcp/scripts/backfill_x_published_at.py index 804f2fc..826dd72 100644 --- a/mcp/scripts/backfill_x_published_at.py +++ b/mcp/scripts/backfill_x_published_at.py @@ -8,9 +8,10 @@ from pathlib import Path sys.path.insert(0, str(Path(__file__).resolve().parents[1])) +sys.path.insert(0, str(Path(__file__).resolve().parents[2])) from scripts.path_config import resolve_feed_db -from src import feed_backend +from feed_runtime import backend as feed_backend _TWITTER_EPOCH_MS = 1288834974657 diff --git a/mcp/scripts/causal_offline_scorer.py b/mcp/scripts/causal_offline_scorer.py index 6c15370..a9cb86e 100644 --- a/mcp/scripts/causal_offline_scorer.py +++ b/mcp/scripts/causal_offline_scorer.py @@ -13,8 +13,9 @@ from typing import Any sys.path.insert(0, str(Path(__file__).resolve().parents[1])) +sys.path.insert(0, str(Path(__file__).resolve().parents[2])) -from src import feed_backend +from feed_runtime import backend as feed_backend from scripts.recover_published_at import recovery_key diff --git a/mcp/scripts/evaluate_ranker.py b/mcp/scripts/evaluate_ranker.py index fce6bb2..0ca5f9c 100644 --- a/mcp/scripts/evaluate_ranker.py +++ b/mcp/scripts/evaluate_ranker.py @@ -9,9 +9,10 @@ from pathlib import Path sys.path.insert(0, str(Path(__file__).resolve().parents[1])) +sys.path.insert(0, str(Path(__file__).resolve().parents[2])) from scripts.path_config import resolve_feed_db, resolve_workspace_path -from src import feed_backend +from feed_runtime import backend as feed_backend def _precision(labels: list[int], k: int) -> float: diff --git a/mcp/scripts/recover_published_at.py b/mcp/scripts/recover_published_at.py index a24aaf6..ff9c191 100644 --- a/mcp/scripts/recover_published_at.py +++ b/mcp/scripts/recover_published_at.py @@ -16,8 +16,9 @@ import requests sys.path.insert(0, str(Path(__file__).resolve().parents[1])) +sys.path.insert(0, str(Path(__file__).resolve().parents[2])) -from src import feed_backend +from feed_runtime import backend as feed_backend _ARXIV_ID_RE = re.compile(r"^(?:[a-z0-9._-]+/\d{7}|\d{4}\.\d{4,5})(?:v\d+)?$", re.I) diff --git a/mcp/scripts/replay_recent_ticks.py b/mcp/scripts/replay_recent_ticks.py index fb8ff53..f5911a8 100644 --- a/mcp/scripts/replay_recent_ticks.py +++ b/mcp/scripts/replay_recent_ticks.py @@ -8,10 +8,11 @@ from pathlib import Path sys.path.insert(0, str(Path(__file__).resolve().parents[1])) +sys.path.insert(0, str(Path(__file__).resolve().parents[2])) from scripts.evaluate_ranker import _labels_from_proactive_feedback from scripts.path_config import resolve_feed_db, resolve_workspace_path -from src import feed_backend +from feed_runtime import backend as feed_backend def _label(event_id: str, fallback_labels: dict[str, int], feedback_labels: dict[str, int] | None) -> int | None: diff --git a/mcp/src/mcp_bridge.py b/mcp/src/mcp_bridge.py index 61e52a7..0c0d858 100644 --- a/mcp/src/mcp_bridge.py +++ b/mcp/src/mcp_bridge.py @@ -1,16 +1,10 @@ from __future__ import annotations -import asyncio -import json -import logging import os -from contextlib import asynccontextmanager -from typing import Any, AsyncIterator +from typing import Any from mcp.server.fastmcp import FastMCP -logger = logging.getLogger(__name__) - def _recording_backend() -> bool: return os.environ.get("FEED_BACKEND", "").strip().lower() == "recording" @@ -19,113 +13,15 @@ def _recording_backend() -> bool: def _live_backend() -> Any: if _recording_backend(): raise RuntimeError("feed recording backend 不允许访问正式 Feed 后端") - from src import feed_backend - - return feed_backend - - -class FeedPoller: - """在正式运行时后台维护 Feed 缓存刷新。""" - - def __init__(self) -> None: - self._lock = asyncio.Lock() - self._stop = asyncio.Event() - self._task: asyncio.Task[None] | None = None - - async def start(self) -> None: - if self._task is not None: - raise RuntimeError("Feed poller 已启动") - self._task = asyncio.create_task(self._run(), name="feed-cache-poller") - - async def stop(self) -> None: - self._stop.set() - if self._task is None: - return - task = self._task - task.cancel() - caller_cancelled = False - while not task.done(): - try: - await asyncio.shield(task) - except asyncio.CancelledError: - if not task.done(): - caller_cancelled = True - except Exception: - break - self._task = None - if not task.cancelled(): - task.result() - if caller_cancelled: - raise asyncio.CancelledError - - async def poll_now(self) -> None: - async with self._lock: - worker = asyncio.create_task( - asyncio.to_thread(_live_backend().poll_feeds_only), - name="feed-cache-poll-worker", - ) - try: - await asyncio.shield(worker) - except asyncio.CancelledError: - while not worker.done(): - try: - await asyncio.shield(worker) - except asyncio.CancelledError: - continue - except Exception: - break - error = worker.exception() - if error is not None: - logger.error( - "[feed] poll worker 在取消收束期间失败", - exc_info=(type(error), error, error.__traceback__), - ) - raise - - async def _run(self) -> None: - """首次立即刷新,随后按缓存 TTL 持续刷新。""" - - # 1. 首次刷新失败必须暴露,同时保留后续重试能力。 - try: - await self.poll_now() - except Exception: - logger.exception("[feed] 首次缓存刷新失败") + from feed_runtime import backend - # 2. 正式后端按配置周期刷新,recording 不会进入此生命周期。 - while not self._stop.is_set(): - try: - interval = _live_backend().load_config().poll_ttl_seconds - except Exception: - logger.exception("[feed] 读取轮询配置失败") - interval = 60 - try: - await asyncio.wait_for(self._stop.wait(), timeout=interval) - return - except TimeoutError: - pass - try: - await self.poll_now() - except Exception: - logger.exception("[feed] 后台缓存刷新失败") + return backend def create_mcp_server() -> FastMCP: - """创建 Feed MCP,并让 recording 生命周期保持零轮询、零数据库访问。""" + """暴露用户驱动的 Feed 工具,但不拥有后台轮询。""" - poller = None if _recording_backend() else FeedPoller() - - @asynccontextmanager - async def lifespan(_: FastMCP) -> AsyncIterator[None]: - if poller is None: - yield None - return - await poller.start() - try: - yield None - finally: - await poller.stop() - - mcp = FastMCP("feed-mcp", lifespan=lifespan) + mcp = FastMCP("feed-mcp") @mcp.tool() def feed_manage( @@ -154,7 +50,7 @@ def feed_query( page: int = 1, page_size: int = 20, ) -> str: - """查询 RSS 订阅内容。""" + """查询由唯一 Timer owner 维护的 RSS 缓存。""" return _live_backend().feed_query( action=action, @@ -165,109 +61,4 @@ def feed_query( page_size=page_size, ) - @mcp.tool() - async def poll_feeds() -> str: - if poller is None: - raise RuntimeError("feed recording backend 不允许轮询") - await poller.poll_now() - return "ok" - - @mcp.tool() - async def get_proactive_events( - offset: int = 0, - limit: int = 50, - cursor: str | None = None, - ) -> str: - events = await asyncio.to_thread( - _fetch_proactive_events, - offset=offset, - limit=limit, - cursor=cursor, - ) - return json.dumps(events, ensure_ascii=False) - - @mcp.tool() - def acknowledge_events( - event_ids: list[str], feedback: str | None = None - ) -> str: - return json.dumps( - _acknowledge_proactive_events(event_ids, feedback=feedback), - ensure_ascii=False, - ) - return mcp - - -def _proactive_fetch_payload( - events: list[dict[str, Any]], - *, - cursor: str | None = None, -) -> dict[str, Any]: - """把正式后端结果编码成 Core 可识别的 typed empty/items。""" - - if not events: - return {"status": "empty"} - payload: dict[str, Any] = {"status": "items", "items": events} - if cursor is not None: - payload["cursor"] = cursor - return payload - - -def _fetch_proactive_events( - *, - offset: int = 0, - limit: int = 50, - cursor: str | None = None, -) -> dict[str, Any]: - """recording 固定返回 typed empty,正式运行才读取 Feed 数据库。""" - - if _recording_backend(): - return {"status": "empty"} - if limit < 1: - raise ValueError("Feed proactive limit 必须大于零") - if cursor is not None: - if offset != 0: - raise ValueError("Feed proactive cursor 不能与 offset 同时使用") - prefix = "feed-offset:" - if not cursor.startswith(prefix) or not cursor[len(prefix) :].isdigit(): - raise ValueError("Feed proactive cursor 无效") - offset = int(cursor[len(prefix) :]) - backend = _live_backend() - events = backend.get_proactive_events(offset=offset, limit=limit + 1) - has_more = len(events) > limit - return _proactive_fetch_payload( - events[:limit], - cursor=f"feed-offset:{offset + limit}" if has_more else None, - ) - - -def _acknowledge_proactive_events( - requested: list[str], *, feedback: str | None = None -) -> dict[str, Any]: - """只有全部请求 ID 持久确认后才编码 committed。""" - - if not requested: - return {"status": "skipped", "reason": "no_ids"} - if _recording_backend(): - raise RuntimeError("feed recording backend 不允许确认事件") - result = _live_backend().acknowledge_events(requested, feedback=feedback) - return _proactive_ack_payload(requested, result) - - -def _proactive_ack_payload( - requested: list[str], result: dict[str, list[str]] -) -> dict[str, Any]: - """把 Feed ack 结果转换为完整 committed 或明确 failure。""" - - if not requested: - return {"status": "skipped", "reason": "no_ids"} - acknowledged = list(result.get("acknowledged", [])) - failed = list(result.get("failed", [])) - if failed or acknowledged != requested: - return { - "status": "failure", - "error": "feed ack 未完整提交", - "retryable": True, - "failed_ids": failed, - } - return {"status": "committed", "ids": acknowledged} diff --git a/mcp/tests/test_causal_offline_scorer.py b/mcp/tests/test_causal_offline_scorer.py index 3dd304d..96ba86e 100644 --- a/mcp/tests/test_causal_offline_scorer.py +++ b/mcp/tests/test_causal_offline_scorer.py @@ -9,7 +9,7 @@ sys.path.insert(0, str(Path(__file__).resolve().parents[1])) from scripts.causal_offline_scorer import score_databases -from src import feed_backend +from feed_runtime import backend as feed_backend def _create_db(path: Path) -> sqlite3.Connection: diff --git a/mcp/tests/test_wake_contract.py b/mcp/tests/test_content_contract.py similarity index 67% rename from mcp/tests/test_wake_contract.py rename to mcp/tests/test_content_contract.py index bddc56e..16337ca 100644 --- a/mcp/tests/test_wake_contract.py +++ b/mcp/tests/test_content_contract.py @@ -7,7 +7,7 @@ sys.path.insert(0, str(Path(__file__).resolve().parents[1])) -from src import feed_backend +from feed_runtime import backend as feed_backend def _insert_item(conn, *, index: int, source: str, published_at: datetime) -> str: @@ -37,13 +37,13 @@ def _insert_item(conn, *, index: int, source: str, published_at: datetime) -> st return event_id -def test_wake_fetch_returns_all_unread_grouped_by_source_and_consumption_is_not_feedback( +def test_content_export_preserves_full_payload_and_exact_ack( tmp_path, monkeypatch ): now = datetime(2026, 7, 12, 12, tzinfo=UTC) monkeypatch.setenv("AKA_PLUGIN_DATA_DIR", str(tmp_path)) monkeypatch.setattr(feed_backend, "_now", lambda: now) - cfg = feed_backend.load_config() + cfg = feed_backend.load_config(tmp_path) conn = feed_backend._connect(cfg) expected = [] try: @@ -61,19 +61,24 @@ def test_wake_fetch_returns_all_unread_grouped_by_source_and_consumption_is_not_ finally: conn.close() - events = feed_backend.get_proactive_events(limit=100) + events = feed_backend.prepare_content_items(data_root=tmp_path) assert len(events) == 60 - assert [event["source_name"] for event in events] == ["Alpha"] * 30 + ["Beta"] * 30 + payloads = [event["payload"] for event in events] + assert [event["source_name"] for event in payloads] == ["Alpha"] * 30 + ["Beta"] * 30 for source in ("Alpha", "Beta"): timestamps = [ - event["published_at"] for event in events if event["source_name"] == source + event["published_at"] for event in payloads if event["source_name"] == source ] assert timestamps == sorted(timestamps, reverse=True) - assert all("preprocess_score" in event for event in events) + assert all("preprocess_score" in event for event in payloads) + assert payloads[0]["content"].startswith("content-") consumed = expected[0] - feed_backend.acknowledge_events([consumed]) + result = feed_backend.settle_content_item( + consumed, "hash-0", data_root=tmp_path + ) + assert result == {"status": "committed", "disposition": "acknowledged"} conn = feed_backend._connect(cfg) try: row = conn.execute( @@ -90,19 +95,19 @@ def test_wake_fetch_returns_all_unread_grouped_by_source_and_consumption_is_not_ assert row["interest_scored_at"] is None assert updates == 0 - monkeypatch.setattr(feed_backend, "_now", lambda: now + timedelta(days=10)) assert consumed not in { - event["event_id"] for event in feed_backend.get_proactive_events(limit=100) + event["item_id"] + for event in feed_backend.prepare_content_items(data_root=tmp_path) } -def test_wake_fetch_admits_only_content_above_decayed_mass_floor( +def test_content_export_admits_only_items_above_decayed_mass_floor( tmp_path, monkeypatch ): now = datetime(2026, 7, 12, 12, tzinfo=UTC) monkeypatch.setenv("AKA_PLUGIN_DATA_DIR", str(tmp_path)) monkeypatch.setattr(feed_backend, "_now", lambda: now) - cfg = feed_backend.load_config() + cfg = feed_backend.load_config(tmp_path) conn = feed_backend._connect(cfg) try: conn.executemany( @@ -129,9 +134,9 @@ def test_wake_fetch_admits_only_content_above_decayed_mass_floor( finally: conn.close() - events = feed_backend.get_proactive_events() + events = feed_backend.prepare_content_items(data_root=tmp_path) - assert [event["event_id"] for event in events] == ["fresh"] + assert [event["item_id"] for event in events] == ["fresh"] conn = feed_backend._connect(cfg) try: @@ -141,13 +146,47 @@ def test_wake_fetch_admits_only_content_above_decayed_mass_floor( assert remaining == 3 +def test_content_export_freezes_rank_payload_and_obsolete_ack_keeps_new_revision( + tmp_path, monkeypatch +): + now = datetime(2026, 7, 12, 12, tzinfo=UTC) + monkeypatch.setattr(feed_backend, "_now", lambda: now) + cfg = feed_backend.load_config(tmp_path) + conn = feed_backend._connect(cfg) + try: + _insert_item(conn, index=1, source="Alpha", published_at=now) + conn.commit() + finally: + conn.close() + + first = feed_backend.prepare_content_items(data_root=tmp_path) + monkeypatch.setattr(feed_backend, "_now", lambda: now + timedelta(hours=12)) + repeated = feed_backend.prepare_content_items(data_root=tmp_path) + assert repeated == first + + conn = feed_backend._connect(cfg) + try: + conn.execute( + "UPDATE items SET content='changed', content_hash='hash-new' WHERE event_id='event-001'" + ) + conn.execute("DELETE FROM acked_items WHERE event_id='event-001'") + conn.commit() + finally: + conn.close() + assert feed_backend.settle_content_item( + "event-001", "hash-1", data_root=tmp_path + ) == {"status": "committed", "disposition": "obsolete_revision"} + current = feed_backend.prepare_content_items(data_root=tmp_path) + assert [item["revision"] for item in current] == ["hash-new"] + + def test_missing_publication_requires_strong_interest_to_enter_transport(): - assert feed_backend._wake_admission_score( + assert feed_backend._content_admission_score( {"interest": 0.45, "freshness": 0.03} - ) < feed_backend._WAKE_ADMISSION_FLOOR - assert feed_backend._wake_admission_score( + ) < feed_backend._CONTENT_ADMISSION_FLOOR + assert feed_backend._content_admission_score( {"interest": 0.9, "freshness": 0.03} - ) > feed_backend._WAKE_ADMISSION_FLOOR + ) > feed_backend._CONTENT_ADMISSION_FLOOR def test_feedparser_reads_rfc2822_namespace_date_and_stable_entry_id(): diff --git a/mcp/tests/test_poll_lifecycle.py b/mcp/tests/test_poll_lifecycle.py deleted file mode 100644 index ca27dd1..0000000 --- a/mcp/tests/test_poll_lifecycle.py +++ /dev/null @@ -1,132 +0,0 @@ -from __future__ import annotations - -import asyncio -import threading -from types import SimpleNamespace - -import pytest - -from src import mcp_bridge - - -def test_poller_refreshes_immediately_and_continues(monkeypatch) -> None: - calls: list[int] = [] - - backend = SimpleNamespace( - poll_feeds_only=lambda: calls.append(len(calls) + 1), - load_config=lambda: SimpleNamespace(poll_ttl_seconds=0.01), - ) - monkeypatch.setattr(mcp_bridge, "_live_backend", lambda: backend) - - async def scenario() -> None: - poller = mcp_bridge.FeedPoller() - await poller.start() - try: - for _ in range(100): - if calls: - break - await asyncio.sleep(0.01) - assert calls == [1] - for _ in range(100): - if len(calls) >= 2: - break - await asyncio.sleep(0.01) - assert len(calls) >= 2 - finally: - await poller.stop() - - asyncio.run(scenario()) - - -def test_poller_logs_refresh_failure_and_retries(monkeypatch) -> None: - attempts = 0 - - def poll() -> None: - nonlocal attempts - attempts += 1 - if attempts == 1: - raise OSError("feed database unavailable") - - backend = SimpleNamespace( - poll_feeds_only=poll, - load_config=lambda: SimpleNamespace(poll_ttl_seconds=0.01), - ) - monkeypatch.setattr(mcp_bridge, "_live_backend", lambda: backend) - - async def scenario() -> None: - poller = mcp_bridge.FeedPoller() - await poller.start() - try: - for _ in range(100): - if attempts >= 2: - break - await asyncio.sleep(0.01) - assert attempts >= 2 - finally: - await poller.stop() - - asyncio.run(scenario()) - - -def test_poller_stop_waits_for_inflight_thread(monkeypatch) -> None: - entered = threading.Event() - release = threading.Event() - finished = threading.Event() - - def poll() -> None: - entered.set() - release.wait(timeout=5) - finished.set() - - backend = SimpleNamespace( - poll_feeds_only=poll, - load_config=lambda: SimpleNamespace(poll_ttl_seconds=60), - ) - monkeypatch.setattr(mcp_bridge, "_live_backend", lambda: backend) - - async def scenario() -> None: - poller = mcp_bridge.FeedPoller() - await poller.start() - await asyncio.to_thread(entered.wait, 5) - stop = asyncio.create_task(poller.stop()) - await asyncio.sleep(0) - assert not stop.done() - assert not finished.is_set() - release.set() - await stop - assert finished.is_set() - - asyncio.run(scenario()) - - -def test_poller_stop_finishes_worker_before_restoring_cancellation(monkeypatch) -> None: - entered = threading.Event() - release = threading.Event() - finished = threading.Event() - - def poll() -> None: - entered.set() - release.wait(timeout=5) - finished.set() - - backend = SimpleNamespace( - poll_feeds_only=poll, - load_config=lambda: SimpleNamespace(poll_ttl_seconds=60), - ) - monkeypatch.setattr(mcp_bridge, "_live_backend", lambda: backend) - - async def scenario() -> None: - poller = mcp_bridge.FeedPoller() - await poller.start() - await asyncio.to_thread(entered.wait, 5) - stop = asyncio.create_task(poller.stop()) - await asyncio.sleep(0) - stop.cancel() - await asyncio.sleep(0) - assert not stop.done() - release.set() - with pytest.raises(asyncio.CancelledError): - await stop - assert finished.is_set() - - asyncio.run(scenario()) diff --git a/mcp/tests/test_recover_published_at.py b/mcp/tests/test_recover_published_at.py index 7f06f97..00a1ae0 100644 --- a/mcp/tests/test_recover_published_at.py +++ b/mcp/tests/test_recover_published_at.py @@ -9,7 +9,7 @@ sys.path.insert(0, str(Path(__file__).resolve().parents[1])) from scripts import recover_published_at -from src import feed_backend +from feed_runtime import backend as feed_backend class _Response: diff --git a/mcp/tests/test_runtime_paths.py b/mcp/tests/test_runtime_paths.py index 5ac411d..93fda8b 100644 --- a/mcp/tests/test_runtime_paths.py +++ b/mcp/tests/test_runtime_paths.py @@ -8,8 +8,8 @@ import plugin from run_mcp import _runtime_dir -from src import feed_backend -from src.feed_backend import _runtime_root, load_config +from feed_runtime import backend as feed_backend +from feed_runtime.backend import _config_path, _runtime_root, load_config def test_runtime_entrypoints_reject_missing_data_dir( @@ -27,8 +27,9 @@ def test_runtime_entrypoints_reject_missing_data_dir( def test_v3_module_keeps_skill_root_and_identity_exports() -> None: assert plugin.api_version == 3 assert plugin.name == "feed" - assert plugin.version == "3.0.0" + assert plugin.version == "3.1.0" assert plugin.skill_roots == ("skills",) + assert _config_path() == Path(__file__).resolve().parents[1] / "feed_mcp.json" def test_concurrent_legacy_connections_share_one_schema_migration( diff --git a/plugin.py b/plugin.py index 860057d..6cf0fee 100644 --- a/plugin.py +++ b/plugin.py @@ -1,61 +1,63 @@ from __future__ import annotations -from pydantic import BaseModel, Field +from pydantic import BaseModel from agent.plugin_composition import ( MCP_SERVERS, - PROACTIVE_COMPONENTS, + RUNTIME_STARTED, + RUNTIME_STOPPING, + TIMERS, Context, McpServerDefinition, - ProactiveSourceDefinition, + ServiceKey, ) - -class FeedProactiveConfig(BaseModel): - enabled: bool = True +from content_source import CONTENT_SOURCE_ID, ContentSourceServices, FeedContentRuntime class FeedConfig(BaseModel): - proactive: FeedProactiveConfig = Field(default_factory=FeedProactiveConfig) + pass + +CONTENT_SOURCE = ServiceKey[ContentSourceServices]("content.source.v1") api_version = 3 name = "feed" -version = "3.0.0" -desc = "Feed MCP plugin" +version = "3.1.0" +desc = "由 Timer 驱动的 Feed Content source 与用户 MCP" Config = FeedConfig -inject = (MCP_SERVERS, PROACTIVE_COMPONENTS) +inject = (MCP_SERVERS, TIMERS, CONTENT_SOURCE) skill_roots = ("skills",) async def apply(ctx: Context, config: object) -> None: - """注册 Feed MCP 与可选的订阅主动事件源。""" + """注册用户 MCP 工具和一个普通 Timer 驱动的 Content source。""" if not isinstance(config, FeedConfig): raise TypeError("feed config 必须是 FeedConfig") - # 1. 只声明 MCP;apply 本身不启动进程、不访问网络或插件数据。 + # 1. MCP 只拥有用户触发的订阅管理和缓存查询。 await ctx.require(MCP_SERVERS).register( ctx, McpServerDefinition( name="feed", command=("python", "mcp/run_mcp.py"), - required_tools=("get_proactive_events", "acknowledge_events"), - candidate_read_only_tools=("get_proactive_events",), + required_tools=("feed_manage", "feed_query"), + candidate_read_only_tools=(), candidate_env={"FEED_BACKEND": "recording"}, ), ) - # 2. 主动能力由用户配置决定是否发布。 - if config.proactive.enabled: - await ctx.require(PROACTIVE_COMPONENTS).register( - ctx, - ProactiveSourceDefinition( - name="subscriptions", - channels=("content",), - mcp_server="feed", - fetch_tool="get_proactive_events", - ack_tool="acknowledge_events", - fetch_page_size=50, - ), - ) + # 2. 正式 Root 独占外部轮询与 Content ACK。 + runtime = FeedContentRuntime( + ctx.data_root, + ctx.require(TIMERS), + ctx.require(CONTENT_SOURCE).bind(CONTENT_SOURCE_ID), + ) + + def setup() -> object: + return runtime.close + + _ = await ctx.effect(setup, label="feed-content-source-runtime") + _ = await ctx.on(RUNTIME_STARTED, lambda _: runtime.start()) + _ = await ctx.on(RUNTIME_STOPPING, lambda _: runtime.close()) diff --git a/pyrightconfig.json b/pyrightconfig.json index 17e8cd5..e676354 100644 --- a/pyrightconfig.json +++ b/pyrightconfig.json @@ -1,6 +1,9 @@ { "include": [ "plugin.py", + "content_source.py", + "legacy_handoff.py", + "feed_runtime", "mcp/run_mcp.py", "mcp/src", "scripts", diff --git a/scripts/retire_legacy_feed_backlog.py b/scripts/retire_legacy_feed_backlog.py new file mode 100755 index 0000000..148cd80 --- /dev/null +++ b/scripts/retire_legacy_feed_backlog.py @@ -0,0 +1,139 @@ +#!/usr/bin/env python3 +"""Retire an exact legacy Feed backlog without submitting it to Content.""" + +from __future__ import annotations + +import argparse +import hashlib +import json +import sys +from pathlib import Path + +PLUGIN_ROOT = Path(__file__).resolve().parents[1] +if str(PLUGIN_ROOT) not in sys.path: + sys.path.insert(0, str(PLUGIN_ROOT)) + +from agent.migrations.proactive_island.cli import report_payload, retire +from agent.migrations.proactive_island.inventory import ( + inventory_digest, + inventory_workspace, +) +from legacy_handoff import FeedLegacyHandoffAdapter +from feed_runtime import backend + + +def _backlog_digest(items: tuple[dict[str, str], ...]) -> str: + payload = json.dumps(items, sort_keys=True, separators=(",", ":")).encode() + return hashlib.sha256(payload).hexdigest() + + +def main() -> int: + parser = argparse.ArgumentParser(description=__doc__) + _ = parser.add_argument("--workspace", required=True, type=Path) + _ = parser.add_argument("--feed-data-root", required=True, type=Path) + mode = parser.add_mutually_exclusive_group() + _ = mode.add_argument("--apply", action="store_true") + _ = mode.add_argument("--plan", action="store_true") + _ = parser.add_argument("--backup-root", type=Path) + _ = parser.add_argument("--expected-inventory-sha256") + _ = parser.add_argument("--expected-provider-backlog-sha256") + _ = parser.add_argument("--expected-provider-backlog-count", type=int) + args = parser.parse_args() + + inventory = inventory_workspace(args.workspace) + inventory_sha256 = inventory_digest(inventory) + backlog = backend.plan_content_backlog(data_root=args.feed_data_root) + backlog_sha256 = _backlog_digest(backlog) + if not args.apply: + print( + json.dumps( + { + "status": "plan", + "inventory_sha256": inventory_sha256, + "provider_backlog_sha256": backlog_sha256, + "provider_backlog_count": len(backlog), + }, + ensure_ascii=False, + indent=2, + ) + ) + return 0 + if any( + value is None + for value in ( + args.backup_root, + args.expected_inventory_sha256, + args.expected_provider_backlog_sha256, + args.expected_provider_backlog_count, + ) + ): + parser.error( + "--apply requires --backup-root and all expected inventory/backlog fields" + ) + expected_inventory = str(args.expected_inventory_sha256) + expected_backlog = str(args.expected_provider_backlog_sha256) + expected_count = int(args.expected_provider_backlog_count) + batch_id = f"proactive-island-cutover:{expected_inventory}" + completed = backend.read_content_backlog_supersession( + batch_id, + data_root=args.feed_data_root, + ) + if completed is not None: + completed_receipt, completed_items = completed + if ( + completed_receipt["items_digest"] != expected_backlog + or completed_receipt["item_count"] != expected_count + or not backend.verify_content_backlog_supersession( + batch_id, completed_items, data_root=args.feed_data_root + ) + ): + raise RuntimeError( + "Feed completed cutover receipt conflicts with expected plan" + ) + backlog = completed_items + backlog_sha256 = _backlog_digest(backlog) + + if ( + inventory_sha256 != expected_inventory + or backlog_sha256 != expected_backlog + or len(backlog) != expected_count + ): + print( + json.dumps( + { + "status": "block", + "reason": "cutover_plan_drift", + "inventory_sha256": inventory_sha256, + "provider_backlog_sha256": backlog_sha256, + "provider_backlog_count": len(backlog), + }, + ensure_ascii=False, + indent=2, + ) + ) + return 2 + + provider_receipt = backend.supersede_content_backlog( + batch_id, + backlog, + data_root=args.feed_data_root, + ) + adapter = FeedLegacyHandoffAdapter.for_cutover(args.feed_data_root) + report = retire( + args.workspace, + args.backup_root, + inventory_sha256, + (adapter,), + ) + payload = report_payload(report) + payload["inventory_sha256"] = inventory_sha256 + payload["provider_backlog"] = provider_receipt + payload["provider_backlog_remaining"] = len( + backend.plan_content_backlog(data_root=args.feed_data_root) + ) + print(json.dumps(payload, ensure_ascii=False, indent=2)) + return 0 if report.status.value != "block" else 2 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tests/fixtures/legacy_feed_owner/akashic.plugin.toml b/tests/fixtures/legacy_feed_owner/akashic.plugin.toml new file mode 100644 index 0000000..f649afc --- /dev/null +++ b/tests/fixtures/legacy_feed_owner/akashic.plugin.toml @@ -0,0 +1,14 @@ +schema_version = 1 +name = "feed" +version = "3.0.0" +api_version = 3 +entrypoint = "plugin.py" + +[[python]] +requirements = "mcp/requirements.txt" + +[[mcp]] +name = "feed" +command = ["python", "mcp/run_mcp.py"] +required_tools = ["legacy_status"] +candidate_read_only_tools = ["legacy_status"] diff --git a/tests/fixtures/legacy_feed_owner/mcp/requirements.txt b/tests/fixtures/legacy_feed_owner/mcp/requirements.txt new file mode 100644 index 0000000..3533ce3 --- /dev/null +++ b/tests/fixtures/legacy_feed_owner/mcp/requirements.txt @@ -0,0 +1 @@ +mcp==1.28.1 diff --git a/tests/fixtures/legacy_feed_owner/mcp/run_mcp.py b/tests/fixtures/legacy_feed_owner/mcp/run_mcp.py new file mode 100644 index 0000000..8689f14 --- /dev/null +++ b/tests/fixtures/legacy_feed_owner/mcp/run_mcp.py @@ -0,0 +1,38 @@ +from __future__ import annotations + +import json +import os +import time +from collections.abc import AsyncIterator +from contextlib import asynccontextmanager +from pathlib import Path + +from mcp.server.fastmcp import FastMCP + + +def _record(event: str) -> None: + root = Path(os.environ["AKA_PLUGIN_DATA_DIR"]) + root.mkdir(parents=True, exist_ok=True) + with (root / "legacy-owner.jsonl").open("a", encoding="utf-8") as handle: + handle.write(json.dumps({"event": event, "time_ns": time.time_ns()}) + "\n") + + +@asynccontextmanager +async def _lifespan(_: FastMCP) -> AsyncIterator[None]: + _record("started") + try: + yield None + finally: + _record("stopped") + + +server = FastMCP("legacy-feed-owner", lifespan=_lifespan) + + +@server.tool() +def legacy_status() -> str: + return "ready" + + +if __name__ == "__main__": + server.run(transport="stdio") diff --git a/tests/fixtures/legacy_feed_owner/plugin.py b/tests/fixtures/legacy_feed_owner/plugin.py new file mode 100644 index 0000000..5792dcd --- /dev/null +++ b/tests/fixtures/legacy_feed_owner/plugin.py @@ -0,0 +1,26 @@ +from __future__ import annotations + +from agent.plugin_composition import MCP_SERVERS, Context, McpServerDefinition + + +api_version = 3 +name = "feed" +version = "3.0.0" +desc = "旧 lifespan 轮询 owner 换班 fixture" +inject = (MCP_SERVERS,) +skill_roots = () + + +async def apply(ctx: Context, config: object) -> None: + """注册一个由 lifespan 拥有后台轮询的旧 MCP。""" + + _ = config + await ctx.require(MCP_SERVERS).register( + ctx, + McpServerDefinition( + name="feed", + command=("python", "mcp/run_mcp.py"), + required_tools=("legacy_status",), + candidate_read_only_tools=("legacy_status",), + ), + ) diff --git a/tests/fixtures/legacy_feed_reservoir_rows.json b/tests/fixtures/legacy_feed_reservoir_rows.json new file mode 100644 index 0000000..6f6b555 --- /dev/null +++ b/tests/fixtures/legacy_feed_reservoir_rows.json @@ -0,0 +1,17 @@ +[ + {"item_id":"feed@github:subscriptions:event-01","kind":"content","source_id":"feed@github:subscriptions","original_source_id":"source-01","ack_source_id":"feed@github:subscriptions","source_event_id":"event-01","published_at":"2026-08-01T01:00:00+00:00","first_seen_at":"2026-08-01T01:01:00+00:00","preprocess_score":0.51,"payload_json":"{\"event_id\":\"event-01\",\"kind\":\"content\",\"preprocess_score\":0.51,\"preprocess_features\":{\"freshness\":0.91}}","embedding_json":null,"status":"unread","consumed_at":null}, + {"item_id":"feed@github:subscriptions:event-02","kind":"content","source_id":"feed@github:subscriptions","original_source_id":"source-02","ack_source_id":"feed@github:subscriptions","source_event_id":"event-02","published_at":"2026-08-01T02:00:00+00:00","first_seen_at":"2026-08-01T02:01:00+00:00","preprocess_score":0.52,"payload_json":"{\"event_id\":\"event-02\",\"kind\":\"content\",\"preprocess_score\":0.52,\"preprocess_features\":{\"freshness\":0.92}}","embedding_json":"[0.2,0.8]","status":"unread","consumed_at":null}, + {"item_id":"feed@github:subscriptions:event-03","kind":"content","source_id":"feed@github:subscriptions","original_source_id":"source-03","ack_source_id":"feed@github:subscriptions","source_event_id":"event-03","published_at":"2026-08-01T03:00:00+00:00","first_seen_at":"2026-08-01T03:01:00+00:00","preprocess_score":0.53,"payload_json":"{\"event_id\":\"event-03\",\"kind\":\"content\",\"preprocess_score\":0.53,\"preprocess_features\":{\"freshness\":0.93}}","embedding_json":null,"status":"unread","consumed_at":null}, + {"item_id":"feed@github:subscriptions:event-04","kind":"content","source_id":"feed@github:subscriptions","original_source_id":"source-04","ack_source_id":"feed@github:subscriptions","source_event_id":"event-04","published_at":"2026-08-01T04:00:00+00:00","first_seen_at":"2026-08-01T04:01:00+00:00","preprocess_score":0.54,"payload_json":"{\"event_id\":\"event-04\",\"kind\":\"content\",\"preprocess_score\":0.54,\"preprocess_features\":{\"freshness\":0.94}}","embedding_json":"[0.4,0.6]","status":"unread","consumed_at":null}, + {"item_id":"feed@github:subscriptions:event-05","kind":"content","source_id":"feed@github:subscriptions","original_source_id":"source-05","ack_source_id":"feed@github:subscriptions","source_event_id":"event-05","published_at":"2026-08-01T05:00:00+00:00","first_seen_at":"2026-08-01T05:01:00+00:00","preprocess_score":0.55,"payload_json":"{\"event_id\":\"event-05\",\"kind\":\"content\",\"preprocess_score\":0.55,\"preprocess_features\":{\"freshness\":0.95}}","embedding_json":null,"status":"unread","consumed_at":null}, + {"item_id":"feed@github:subscriptions:event-06","kind":"content","source_id":"feed@github:subscriptions","original_source_id":"source-06","ack_source_id":"feed@github:subscriptions","source_event_id":"event-06","published_at":"2026-08-01T06:00:00+00:00","first_seen_at":"2026-08-01T06:01:00+00:00","preprocess_score":0.56,"payload_json":"{\"event_id\":\"event-06\",\"kind\":\"content\",\"preprocess_score\":0.56,\"preprocess_features\":{\"freshness\":0.96}}","embedding_json":"[0.6,0.4]","status":"unread","consumed_at":null}, + {"item_id":"feed@github:subscriptions:event-07","kind":"content","source_id":"feed@github:subscriptions","original_source_id":"source-07","ack_source_id":"feed@github:subscriptions","source_event_id":"event-07","published_at":"2026-08-01T07:00:00+00:00","first_seen_at":"2026-08-01T07:01:00+00:00","preprocess_score":0.57,"payload_json":"{\"event_id\":\"event-07\",\"kind\":\"content\",\"preprocess_score\":0.57,\"preprocess_features\":{\"freshness\":0.97}}","embedding_json":null,"status":"unread","consumed_at":null}, + {"item_id":"feed@github:subscriptions:event-08","kind":"content","source_id":"feed@github:subscriptions","original_source_id":"source-08","ack_source_id":"feed@github:subscriptions","source_event_id":"event-08","published_at":"2026-08-01T08:00:00+00:00","first_seen_at":"2026-08-01T08:01:00+00:00","preprocess_score":0.58,"payload_json":"{\"event_id\":\"event-08\",\"kind\":\"content\",\"preprocess_score\":0.58,\"preprocess_features\":{\"freshness\":0.98}}","embedding_json":"[0.8,0.2]","status":"unread","consumed_at":null}, + {"item_id":"feed@github:subscriptions:event-09","kind":"content","source_id":"feed@github:subscriptions","original_source_id":"source-09","ack_source_id":"feed@github:subscriptions","source_event_id":"event-09","published_at":"2026-08-01T09:00:00+00:00","first_seen_at":"2026-08-01T09:01:00+00:00","preprocess_score":0.59,"payload_json":"{\"event_id\":\"event-09\",\"kind\":\"content\",\"preprocess_score\":0.59,\"preprocess_features\":{\"freshness\":0.99}}","embedding_json":null,"status":"unread","consumed_at":null}, + {"item_id":"feed@github:subscriptions:event-10","kind":"content","source_id":"feed@github:subscriptions","original_source_id":"source-10","ack_source_id":"feed@github:subscriptions","source_event_id":"event-10","published_at":"2026-08-01T10:00:00+00:00","first_seen_at":"2026-08-01T10:01:00+00:00","preprocess_score":0.60,"payload_json":"{\"event_id\":\"event-10\",\"kind\":\"content\",\"preprocess_score\":0.6,\"preprocess_features\":{\"freshness\":1.0}}","embedding_json":"[1.0,0.0]","status":"unread","consumed_at":null}, + {"item_id":"feed@github:subscriptions:event-11","kind":"content","source_id":"feed@github:subscriptions","original_source_id":"source-11","ack_source_id":"feed@github:subscriptions","source_event_id":"event-11","published_at":"2026-08-01T11:00:00+00:00","first_seen_at":"2026-08-01T11:01:00+00:00","preprocess_score":0.61,"payload_json":"{\"event_id\":\"event-11\",\"kind\":\"content\",\"preprocess_score\":0.61,\"preprocess_features\":{\"freshness\":0.89}}","embedding_json":null,"status":"unread","consumed_at":null}, + {"item_id":"feed@github:subscriptions:event-12","kind":"content","source_id":"feed@github:subscriptions","original_source_id":"source-12","ack_source_id":"feed@github:subscriptions","source_event_id":"event-12","published_at":"2026-08-01T12:00:00+00:00","first_seen_at":"2026-08-01T12:01:00+00:00","preprocess_score":0.62,"payload_json":"{\"event_id\":\"event-12\",\"kind\":\"content\",\"preprocess_score\":0.62,\"preprocess_features\":{\"freshness\":0.88}}","embedding_json":"[0.3,0.7]","status":"unread","consumed_at":null}, + {"item_id":"feed@github:subscriptions:event-13","kind":"content","source_id":"feed@github:subscriptions","original_source_id":"source-13","ack_source_id":"feed@github:subscriptions","source_event_id":"event-13","published_at":"2026-08-01T13:00:00+00:00","first_seen_at":"2026-08-01T13:01:00+00:00","preprocess_score":0.63,"payload_json":"{\"event_id\":\"event-13\",\"kind\":\"content\",\"preprocess_score\":0.63,\"preprocess_features\":{\"freshness\":0.87}}","embedding_json":null,"status":"unread","consumed_at":null}, + {"item_id":"feed@github:subscriptions:event-14","kind":"content","source_id":"feed@github:subscriptions","original_source_id":"source-14","ack_source_id":"feed@github:subscriptions","source_event_id":"event-14","published_at":"2026-08-01T14:00:00+00:00","first_seen_at":"2026-08-01T14:01:00+00:00","preprocess_score":0.64,"payload_json":"{\"event_id\":\"event-14\",\"kind\":\"content\",\"preprocess_score\":0.64,\"preprocess_features\":{\"freshness\":0.86}}","embedding_json":"[0.7,0.3]","status":"unread","consumed_at":null}, + {"item_id":"feed@github:subscriptions:event-15","kind":"content","source_id":"feed@github:subscriptions","original_source_id":"source-15","ack_source_id":"feed@github:subscriptions","source_event_id":"event-15","published_at":"2026-08-01T15:00:00+00:00","first_seen_at":"2026-08-01T15:01:00+00:00","preprocess_score":0.65,"payload_json":"{\"event_id\":\"event-15\",\"kind\":\"content\",\"preprocess_score\":0.65,\"preprocess_features\":{\"freshness\":0.85}}","embedding_json":null,"status":"unread","consumed_at":null} +] diff --git a/tests/test_content_source.py b/tests/test_content_source.py new file mode 100644 index 0000000..c6a6d9a --- /dev/null +++ b/tests/test_content_source.py @@ -0,0 +1,179 @@ +from __future__ import annotations + +import asyncio +import sqlite3 +from collections.abc import Mapping +from datetime import UTC, datetime +from pathlib import Path +from typing import cast + +import pytest + +from agent.control.timer import TimerReceipt, TimerStatus +from agent.plugin_composition import PluginTimers +from content_source import FeedContentRuntime +from feed_runtime import backend + + +class _TimerHandle: + def __init__(self, deadline: datetime, now: datetime) -> None: + self.id = "timer:feed" + self.deadline = deadline + self.now = now + self.future: asyncio.Future[TimerReceipt] = ( + asyncio.get_running_loop().create_future() + ) + + async def result(self) -> TimerReceipt: + return await self.future + + async def cancel(self) -> TimerReceipt: + if not self.future.done(): + self.future.set_result(self._receipt(TimerStatus.CANCELLED)) + return await self.future + + async def cleanup(self) -> None: + return None + + def fire(self) -> None: + self.future.set_result(self._receipt(TimerStatus.FIRED)) + + def _receipt(self, status: TimerStatus) -> TimerReceipt: + return TimerReceipt(self.id, self.deadline, self.now, status) + + +class _Timer: + def __init__(self, now: datetime) -> None: + self.now = now + self.handles: list[_TimerHandle] = [] + + def schedule(self, deadline: datetime) -> _TimerHandle: + handle = _TimerHandle(deadline, self.now) + self.handles.append(handle) + return handle + + +class _Content: + def __init__(self) -> None: + self.submissions: list[tuple[str, tuple[dict[str, object], ...]]] = [] + self.rows: list[dict[str, object]] = [] + self.fail_ack_once = False + + def submit(self, batch_id, items): + frozen = tuple(dict(item) for item in items) + self.submissions.append((batch_id, frozen)) + return {"inserted": [item["item_id"] for item in frozen]} + + def unsettled(self, limit=100): + return tuple(self.rows[:limit]) + + def ack(self, settlement_ref): + if self.fail_ack_once: + self.fail_ack_once = False + raise RuntimeError("crash after provider ACK") + self.rows = [ + row for row in self.rows if row["settlement_ref"] != settlement_ref + ] + return {"changed": True} + + +def _seed_item(data_root: Path, now: datetime) -> None: + cfg = backend.load_config(data_root) + connection = backend._connect(cfg) + try: + connection.execute( + """ + INSERT INTO items( + event_id, source_id, source_name, source_type, title, content, + url, author, published_at, first_seen_at, last_seen_at, + emitted_at, content_hash + ) VALUES('event-1', 'source', 'Source', 'rss', 'Title', 'Body', + 'https://example.com/1', 'Author', ?, ?, ?, NULL, 'revision-1') + """, + (now.isoformat(), now.isoformat(), now.isoformat()), + ) + connection.commit() + finally: + connection.close() + + +@pytest.mark.asyncio +async def test_timer_poll_submits_nonempty_once_and_empty_poll_has_no_history( + tmp_path: Path, +) -> None: + now = datetime(2026, 8, 23, 10, tzinfo=UTC) + timer = _Timer(now) + content = _Content() + _seed_item(tmp_path, now) + runtime = FeedContentRuntime( + tmp_path, + PluginTimers(timer), + content, + now=lambda: now, + ) + + await runtime.start() + assert len(timer.handles) == 1 + handler = runtime._log_handler # pyright: ignore[reportPrivateUsage] + assert handler is not None + assert handler.backupCount == 3 + assert handler.maxBytes == 5 * 1024 * 1024 + timer.handles[0].fire() + for _ in range(200): + if content.submissions and len(timer.handles) == 2: + break + await asyncio.sleep(0.01) + assert len(content.submissions) == 1 + payload = cast(Mapping[str, object], content.submissions[0][1][0]["payload"]) + assert payload["content"] == "Body" + + # ACK the only item, then the next empty poll only updates current deadline. + _ = backend.settle_content_item( + "event-1", "revision-1", data_root=tmp_path + ) + timer.handles[1].fire() + for _ in range(200): + if len(timer.handles) == 3: + break + await asyncio.sleep(0.01) + assert len(content.submissions) == 1 + with sqlite3.connect(tmp_path / "feed_mcp.sqlite3") as connection: + assert connection.execute( + "SELECT COUNT(*) FROM content_exports" + ).fetchone() == (1,) + assert connection.execute( + "SELECT COUNT(*) FROM metadata WHERE key='content_source_next_due'" + ).fetchone() == (1,) + await runtime.close() + + +@pytest.mark.asyncio +async def test_provider_ack_retries_after_crash_before_content_ack( + tmp_path: Path, +) -> None: + now = datetime(2026, 8, 23, 10, tzinfo=UTC) + _seed_item(tmp_path, now) + content = _Content() + content.rows = [ + { + "ref": {"item_id": "event-1", "revision": "revision-1"}, + "settlement_ref": "delivery:1", + } + ] + content.fail_ack_once = True + runtime = FeedContentRuntime( + tmp_path, + PluginTimers(_Timer(now)), + content, + now=lambda: now, + ) + + with pytest.raises(RuntimeError, match="crash after provider ACK"): + await runtime._settle_delivered() # pyright: ignore[reportPrivateUsage] + assert content.rows + assert await runtime._settle_delivered() == 1 # pyright: ignore[reportPrivateUsage] + assert content.rows == [] + with sqlite3.connect(tmp_path / "feed_mcp.sqlite3") as connection: + assert connection.execute( + "SELECT event_id FROM acked_items" + ).fetchall() == [("event-1",)] diff --git a/tests/test_legacy_handoff.py b/tests/test_legacy_handoff.py new file mode 100644 index 0000000..e897475 --- /dev/null +++ b/tests/test_legacy_handoff.py @@ -0,0 +1,512 @@ +from __future__ import annotations + +import hashlib +import json +import sqlite3 +from collections.abc import Mapping, Sequence +from contextlib import closing +from datetime import UTC, datetime +from pathlib import Path +from typing import cast + +import pytest + +from agent.migrations.proactive_island import ( + HandoffBlocked, + HandoffStatus, + Inventory, + LegacyFact, + LegacyFactKind, + apply_handoff, +) +from plugins.content.store import ContentIdentityConflict, ContentStore + +from feed_runtime import backend +from legacy_handoff import FeedLegacyHandoffAdapter, LEGACY_SOURCE_ID + + +ROOT = Path(__file__).resolve().parent +TARGET_SOURCE = "feed-subscriptions" + + +class _BoundContent: + def __init__(self, store: ContentStore, source_id: str = TARGET_SOURCE) -> None: + self.store = store + self.source_id = source_id + + def submit( + self, batch_id: str, items: Sequence[Mapping[str, object]] + ) -> Mapping[str, object]: + return self.store.submit(self.source_id, batch_id, items) + + def read_submission(self, batch_id: str) -> Mapping[str, object] | None: + return self.store.read_submission(self.source_id, batch_id) + + def read_revision( + self, item_id: str, revision: str + ) -> Mapping[str, object] | None: + return self.store.read_revision(self.source_id, item_id, revision) + + def unsettled(self, limit: int = 100) -> tuple[Mapping[str, object], ...]: + return self.store.unsettled(self.source_id, limit) + + def ack(self, settlement_ref: str) -> Mapping[str, object]: + return self.store.ack(self.source_id, settlement_ref) + + +def _legacy_rows() -> list[dict[str, object]]: + value = json.loads( + (ROOT / "fixtures" / "legacy_feed_reservoir_rows.json").read_text() + ) + assert isinstance(value, list) and len(value) == 15 + return cast(list[dict[str, object]], value) + + +def _fact(row: Mapping[str, object]) -> LegacyFact: + opaque = json.dumps( + row, ensure_ascii=False, sort_keys=True, separators=(",", ":") + ).encode("utf-8") + item_id = cast(str, row["item_id"]) + return LegacyFact( + kind=LegacyFactKind.WAKE_SOURCE_ITEM, + locator=f"wake:reservoir_events:{item_id}", + source_digest=hashlib.sha256(opaque).hexdigest(), + source_identity="feed@github:subscriptions", + opaque=opaque, + ) + + +def _ack_fact(row: Mapping[str, object], action: str = "consume") -> LegacyFact: + event_id = cast(str, row["source_event_id"]) + item_id = cast(str, row["item_id"]) + acknowledgement = { + "source_id": LEGACY_SOURCE_ID, + "source_event_id": event_id, + "item_id": item_id, + "action": action, + "queued_at": "2026-08-23T10:00:00+00:00", + } + opaque = json.dumps( + acknowledgement, ensure_ascii=False, sort_keys=True, separators=(",", ":") + ).encode("utf-8") + return LegacyFact( + kind=LegacyFactKind.WAKE_ACK, + locator=( + "wake:pending_acknowledgements:" + f"{LEGACY_SOURCE_ID}:{event_id}:{item_id}" + ), + source_digest=hashlib.sha256(opaque).hexdigest(), + source_identity=LEGACY_SOURCE_ID, + opaque=opaque, + ) + + +def _seed_provider(data_root: Path, rows: Sequence[Mapping[str, object]]) -> None: + config = backend.load_config(data_root) + connection = backend._connect(config) + try: + for index, row in enumerate(rows, start=1): + event_id = cast(str, row["source_event_id"]) + connection.execute( + """ + INSERT INTO items( + event_id, source_id, source_name, source_type, title, + content, url, author, published_at, first_seen_at, + last_seen_at, emitted_at, content_hash + ) VALUES(?, ?, ?, 'rss', ?, ?, ?, ?, ?, ?, ?, NULL, ?) + """, + ( + event_id, + f"source-{index:02d}", + f"Source {index:02d}", + f"Title {index:02d}", + f"Body {index:02d}", + f"https://example.com/{index:02d}", + None if index % 3 == 0 else f"Author {index:02d}", + row["published_at"], + row["first_seen_at"], + row["first_seen_at"], + f"revision-{index:02d}", + ), + ) + connection.commit() + finally: + connection.close() + + +def _fixture( + tmp_path: Path, +) -> tuple[Path, ContentStore, _BoundContent, FeedLegacyHandoffAdapter]: + data_root = tmp_path / "feed-data" + _seed_provider(data_root, _legacy_rows()) + store = ContentStore(tmp_path / "content.sqlite3") + store.initialize() + bound = _BoundContent(store) + return data_root, store, bound, FeedLegacyHandoffAdapter(data_root, bound) + + +def _tree_state(root: Path) -> tuple[tuple[str, int, int, str], ...]: + return tuple( + ( + str(path.relative_to(root)), + path.stat().st_size, + path.stat().st_mtime_ns, + hashlib.sha256(path.read_bytes()).hexdigest(), + ) + for path in sorted(root.rglob("*")) + if path.is_file() + ) + + +def test_real_shape_fifteen_rows_plan_apply_and_verify(tmp_path: Path) -> None: + data_root, store, _bound, adapter = _fixture(tmp_path) + facts = tuple(_fact(row) for row in _legacy_rows()) + before = _tree_state(tmp_path) + + plans = tuple(adapter.plan(fact) for fact in facts) + + assert _tree_state(tmp_path) == before + assert len({plan.target_identity for plan in plans}) == 15 + receipts = tuple( + adapter.apply(fact, plan) for fact, plan in zip(facts, plans, strict=True) + ) + after_apply = _tree_state(tmp_path) + assert all( + adapter.verify(fact, receipt) + for fact, receipt in zip(facts, receipts, strict=True) + ) + assert _tree_state(tmp_path) == after_apply + assert len({receipt.receipt_id for receipt in receipts}) == 15 + assert store.state_counts() == {"pending": 15} + with closing(sqlite3.connect(store.path)) as connection: + assert connection.execute( + "SELECT DISTINCT source_id FROM items" + ).fetchall() == [(TARGET_SOURCE,)] + assert not (data_root / "legacy_handoff.log").exists() + + +def test_adapter_accepts_only_legacy_feed_content_fact(tmp_path: Path) -> None: + _data_root, _store, _bound, adapter = _fixture(tmp_path) + fact = _fact(_legacy_rows()[0]) + + assert adapter.accepts(fact) is True + assert fact.source_identity == LEGACY_SOURCE_ID + assert adapter.accepts( + LegacyFact( + fact.kind, + fact.locator, + fact.source_digest, + "calendar@github:upcoming", + fact.opaque, + ) + ) is False + + +def test_missing_provider_item_blocks_without_writes(tmp_path: Path) -> None: + data_root, store, _bound, adapter = _fixture(tmp_path) + config = backend.load_config(data_root) + connection = backend._connect(config) + connection.execute("DELETE FROM items WHERE event_id='event-15'") + connection.commit() + connection.close() + before = _tree_state(tmp_path) + + with pytest.raises(HandoffBlocked, match="feed_provider_item_missing:event-15"): + adapter.plan(_fact(_legacy_rows()[-1])) + + assert _tree_state(tmp_path) == before + assert store.state_counts() == {} + + +def test_target_before_marker_replay_returns_same_receipt(tmp_path: Path) -> None: + _data_root, store, _bound, adapter = _fixture(tmp_path) + fact = _fact(_legacy_rows()[0]) + plan = adapter.plan(fact) + + first = adapter.apply(fact, plan) + repeated = adapter.apply(fact, plan) + + assert repeated == first + assert adapter.verify(fact, repeated) is True + assert store.state_counts() == {"pending": 1} + + +def test_cutover_supersedes_backlog_without_submitting_content(tmp_path: Path) -> None: + data_root, store, _bound, _adapter = _fixture(tmp_path) + adapter = FeedLegacyHandoffAdapter.for_cutover(data_root) + facts = tuple(_fact(row) for row in _legacy_rows()) + + receipts = tuple(adapter.apply(fact, adapter.plan(fact)) for fact in facts) + + assert all( + adapter.verify(fact, receipt) + for fact, receipt in zip(facts, receipts, strict=True) + ) + assert store.state_counts() == {} + assert backend.prepare_content_items(data_root=data_root) == () + config = backend.load_config(data_root) + with closing(backend._connect(config)) as connection: + assert connection.execute("SELECT count(*) FROM acked_items").fetchone()[0] == 15 + assert connection.execute( + "SELECT count(*) FROM legacy_ack_handoff_receipts " + "WHERE action='cutover_superseded'" + ).fetchone()[0] == 15 + + +def test_cutover_replay_uses_same_receipt_and_requires_provider_ack( + tmp_path: Path, +) -> None: + data_root, store, _bound, _adapter = _fixture(tmp_path) + adapter = FeedLegacyHandoffAdapter.for_cutover(data_root) + fact = _fact(_legacy_rows()[0]) + plan = adapter.plan(fact) + + first = adapter.apply(fact, plan) + repeated = adapter.apply(fact, plan) + + assert repeated == first + assert adapter.verify(fact, repeated) is True + assert store.state_counts() == {} + config = backend.load_config(data_root) + with closing(backend._connect(config)) as connection: + connection.execute("DELETE FROM acked_items") + connection.commit() + assert adapter.verify(fact, repeated) is False + + +def test_provider_backlog_plan_is_read_only_and_batch_supersession_is_exact( + tmp_path: Path, +) -> None: + data_root, store, _bound, _adapter = _fixture(tmp_path) + config = backend.load_config(data_root) + with closing(backend._connect(config)) as connection: + now = datetime.now(UTC).isoformat() + connection.execute( + "UPDATE items SET published_at = ?, first_seen_at = ?", (now, now) + ) + connection.commit() + before = _tree_state(tmp_path) + + planned = backend.plan_content_backlog(data_root=data_root) + + assert len(planned) == 15 + assert _tree_state(tmp_path) == before + receipt = backend.supersede_content_backlog( + "cutover:fixture", + planned, + data_root=data_root, + ) + repeated = backend.supersede_content_backlog( + "cutover:fixture", + planned, + data_root=data_root, + ) + assert repeated == receipt + assert receipt["item_count"] == 15 + assert backend.verify_content_backlog_supersession( + "cutover:fixture", planned, data_root=data_root + ) + assert backend.plan_content_backlog(data_root=data_root) == () + assert store.state_counts() == {} + + +def test_provider_backlog_supersession_rejects_revision_drift(tmp_path: Path) -> None: + data_root, _store, _bound, _adapter = _fixture(tmp_path) + config = backend.load_config(data_root) + with closing(backend._connect(config)) as connection: + now = datetime.now(UTC).isoformat() + connection.execute( + "UPDATE items SET published_at = ?, first_seen_at = ?", (now, now) + ) + connection.commit() + planned = backend.plan_content_backlog(data_root=data_root) + with closing(backend._connect(config)) as connection: + connection.execute( + "UPDATE items SET content_hash='changed' WHERE event_id='event-01'" + ) + connection.commit() + + with pytest.raises(RuntimeError, match="provider revision changed"): + backend.supersede_content_backlog( + "cutover:fixture", planned, data_root=data_root + ) + + +def test_core_replays_after_target_receipt_before_lineage_marker( + tmp_path: Path, +) -> None: + _data_root, store, _bound, adapter = _fixture(tmp_path) + fact = _fact(_legacy_rows()[0]) + inventory = Inventory((fact,), ()) + workspace = tmp_path / "workspace" + + def crash_after_target(_fact: LegacyFact, _receipt: object) -> None: + raise RuntimeError("crash before central marker") + + with pytest.raises(RuntimeError, match="crash before central marker"): + apply_handoff( + workspace, + inventory, + (adapter,), + after_target=crash_after_target, + ) + assert store.state_counts() == {"pending": 1} + + recovered = apply_handoff(workspace, inventory, (adapter,)) + + assert recovered.status is HandoffStatus.APPLIED + assert recovered.items[0].state == "applied" + assert store.state_counts() == {"pending": 1} + + +def test_revision_change_after_target_is_a_batch_conflict(tmp_path: Path) -> None: + data_root, store, _bound, adapter = _fixture(tmp_path) + fact = _fact(_legacy_rows()[0]) + original = adapter.plan(fact) + _ = adapter.apply(fact, original) + config = backend.load_config(data_root) + connection = backend._connect(config) + connection.execute( + "UPDATE items SET content_hash='revision-changed' WHERE event_id='event-01'" + ) + connection.commit() + connection.close() + changed = adapter.plan(fact) + + with pytest.raises(ContentIdentityConflict, match="batch identity conflict"): + adapter.apply(fact, changed) + + assert store.state_counts() == {"pending": 1} + + +def test_plan_revision_change_before_apply_fails_before_submit(tmp_path: Path) -> None: + data_root, store, _bound, adapter = _fixture(tmp_path) + fact = _fact(_legacy_rows()[0]) + plan = adapter.plan(fact) + config = backend.load_config(data_root) + connection = backend._connect(config) + connection.execute( + "UPDATE items SET content_hash='revision-new' WHERE event_id='event-01'" + ) + connection.commit() + connection.close() + + with pytest.raises(RuntimeError, match="target identity drift"): + adapter.apply(fact, plan) + + assert store.state_counts() == {} + + +def test_already_acked_target_replays_and_source_ack_is_bound_once( + tmp_path: Path, +) -> None: + data_root, store, bound, adapter = _fixture(tmp_path) + now = datetime(2026, 8, 23, 10, tzinfo=UTC) + fact = _fact(_legacy_rows()[0]) + plan = adapter.plan(fact) + receipt = adapter.apply(fact, plan) + snapshot = store.snapshot(now) + candidate = cast(tuple[dict[str, object], ...], snapshot["items"])[0] + selected = store.select( + cast(Mapping[str, object], candidate["ref"]), + cast(int, snapshot["snapshot_seq"]), + {"session_id": "wake:fixture", "turn_id": "turn:feed-legacy"}, + now, + ) + token = cast(str, selected["selection_token"]) + _ = store.transition(token, "ready_for_delivery") + _ = store.transition(token, "delivered", settlement_ref="delivery:feed-legacy") + + assert _BoundContent(store, "another-source").unsettled() == () + assert len(bound.unsettled()) == 1 + assert backend.settle_content_item( + "event-01", "revision-01", data_root=data_root + )["disposition"] == "acknowledged" + assert bound.ack("delivery:feed-legacy") == { + "settled": True, + "duplicate": False, + } + assert bound.ack("delivery:feed-legacy") == { + "settled": True, + "duplicate": True, + } + assert adapter.apply(fact, plan) == receipt + assert adapter.verify(fact, receipt) is True + assert store.state_counts() == {"settled": 1} + + +@pytest.mark.parametrize("action", ["consume", "expire"]) +def test_core_hands_off_pending_ack_once_after_target_replay( + tmp_path: Path, + action: str, +) -> None: + data_root, store, _bound, adapter = _fixture(tmp_path) + fact = _ack_fact(_legacy_rows()[0], action) + inventory = Inventory((fact,), ()) + workspace = tmp_path / "workspace" + + def crash_after_target(_fact: LegacyFact, _receipt: object) -> None: + raise RuntimeError("crash before central ACK marker") + + with pytest.raises(RuntimeError, match="crash before central ACK marker"): + apply_handoff( + workspace, + inventory, + (adapter,), + after_target=crash_after_target, + ) + + config = backend.load_config(data_root) + with closing(backend._connect(config)) as connection: + assert connection.execute("SELECT count(*) FROM acked_items").fetchone()[0] == 1 + assert connection.execute( + "SELECT count(*) FROM legacy_ack_handoff_receipts" + ).fetchone()[0] == 1 + connection.execute("DELETE FROM acked_items") + connection.commit() + + recovered = apply_handoff(workspace, inventory, (adapter,)) + + assert recovered.status is HandoffStatus.APPLIED + assert recovered.items[0].state == "applied" + assert store.state_counts() == {} + with closing(backend._connect(config)) as connection: + assert connection.execute("SELECT count(*) FROM acked_items").fetchone()[0] == 0 + assert [ + tuple(row) + for row in connection.execute( + "SELECT event_id, revision, action FROM legacy_ack_handoff_receipts" + ).fetchall() + ] == [("event-01", "revision-01", action)] + + +def test_provider_plan_requires_checkpoint_and_creates_no_files( + tmp_path: Path, +) -> None: + data_root, store, _bound, adapter = _fixture(tmp_path) + wal = data_root / "feed_mcp.sqlite3-wal" + wal.write_bytes(b"not-checkpointed") + before = _tree_state(tmp_path) + + with pytest.raises(HandoffBlocked, match="feed_provider_checkpoint_required"): + adapter.plan(_fact(_legacy_rows()[0])) + + assert _tree_state(tmp_path) == before + assert store.state_counts() == {} + + +def test_pending_ack_uses_original_root_for_nested_provider_path( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + config = backend._config_values() + config["db_path"] = "state/feed.sqlite3" + monkeypatch.setattr(backend, "_config_values", lambda: config) + data_root, _store, _bound, adapter = _fixture(tmp_path) + fact = _ack_fact(_legacy_rows()[0]) + + receipt = adapter.apply(fact, adapter.plan(fact)) + + assert adapter.verify(fact, receipt) is True + assert (data_root / "state" / "feed.sqlite3").is_file() + assert not (data_root / "state" / "state" / "feed.sqlite3").exists() diff --git a/tests/test_manager_integration.py b/tests/test_manager_integration.py index 4f46b3c..2d68ecf 100644 --- a/tests/test_manager_integration.py +++ b/tests/test_manager_integration.py @@ -1,98 +1,357 @@ from __future__ import annotations +import asyncio +import hashlib +import json +import os import shutil import sqlite3 -import sys +import time +from datetime import UTC, datetime from pathlib import Path import pytest -from agent.plugin_composition.proactive import FetchEmpty -from agent.plugins.generation_activity_host import ActivityHost -from agent.plugins.generation_proactive_host import ( - ProactiveActivityAdapter, - ProactiveRuntimeBinding, -) +import agent.plugins.manager as plugin_manager_module +from agent.control.timer import TimerReceipt, TimerStatus from agent.plugins.manager import PluginManager from bus.event_bus import EventBus +from plugins.content.store import ContentStore + +from feed_runtime import backend ROOT = Path(__file__).resolve().parents[1] +CORE_ROOT = Path(os.environ["AKASHIC_AGENT_ROOT"]) + + +def _fixture_runtime() -> Path: + artifact_python = Path(os.environ["AKASHIC_PLUGIN_FIXTURE_PYTHON"]) + return artifact_python.parent.parent + + +class _TimerHandle: + def __init__(self, timer_id: str, deadline: datetime, now: datetime) -> None: + self._id = timer_id + self.deadline = deadline + self.now = now + self.future: asyncio.Future[TimerReceipt] = ( + asyncio.get_running_loop().create_future() + ) + + @property + def id(self) -> str: + return self._id + + async def result(self) -> TimerReceipt: + return await asyncio.shield(self.future) + + async def cancel(self) -> TimerReceipt: + if not self.future.done(): + self.future.set_result(self._receipt(TimerStatus.CANCELLED)) + return await self.future + + async def cleanup(self) -> None: + _ = await self.cancel() + + def fire(self) -> None: + self.future.set_result(self._receipt(TimerStatus.FIRED)) + def _receipt(self, status: TimerStatus) -> TimerReceipt: + return TimerReceipt(self.id, self.deadline, self.now, status) -def _stage_plugin(tmp_path: Path) -> Path: - """复制可执行插件,并复用当前测试解释器的依赖环境。""" - source = tmp_path / "plugins" / "feed" +class _Timer: + def __init__(self, now: datetime) -> None: + self.now = now + self.handles: list[_TimerHandle] = [] + self.schedule_times: list[int] = [] + + def schedule(self, deadline: datetime) -> _TimerHandle: + self.schedule_times.append(time.time_ns()) + handle = _TimerHandle(f"timer:{len(self.handles)}", deadline, self.now) + self.handles.append(handle) + return handle + + +async def _eventually(predicate) -> None: + for _ in range(300): + if predicate(): + return + await asyncio.sleep(0.01) + raise AssertionError("condition did not settle") + + +def _stage_plugins(tmp_path: Path) -> tuple[Path, Path]: + """将真实 Content、Feed 插件装入临时测试环境。""" + + runtime = _fixture_runtime() + plugins = tmp_path / "plugins" + content = plugins / "content" + feed = plugins / "feed" + shutil.copytree(CORE_ROOT / "plugins" / "content", content) shutil.copytree( ROOT, - source, + feed, ignore=shutil.ignore_patterns( ".git", ".akashic-core", + ".plugin-contracts", ".pytest_cache", ".venv", "__pycache__", "tests", ), ) - runtime = Path(sys.executable).parent.parent - (source / "mcp" / ".venv").symlink_to(runtime, target_is_directory=True) - return source + (feed / "mcp" / ".venv").symlink_to(runtime, target_is_directory=True) + return content, feed + + +def _stage_legacy_plugins(tmp_path: Path) -> tuple[Path, Path]: + """装入 Content 和旧 lifespan Feed owner fixture。""" + + runtime = _fixture_runtime() + plugins = tmp_path / "plugins" + content = plugins / "content" + feed = plugins / "feed" + shutil.copytree(CORE_ROOT / "plugins" / "content", content) + shutil.copytree(ROOT / "tests" / "fixtures" / "legacy_feed_owner", feed) + (feed / "mcp" / ".venv").symlink_to(runtime, target_is_directory=True) + return content, feed + + +def _replace_with_current_feed(feed: Path) -> None: + """只将临时旧 Feed source 替换为当前插件树。""" + + runtime = _fixture_runtime() + shutil.rmtree(feed) + shutil.copytree( + ROOT, + feed, + ignore=shutil.ignore_patterns( + ".git", + ".akashic-core", + ".plugin-contracts", + ".pytest_cache", + ".venv", + "__pycache__", + "tests", + ), + ) + (feed / "mcp" / ".venv").symlink_to(runtime, target_is_directory=True) + + +def test_stage_plugins_uses_explicit_fixture_runtime( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + runtime = tmp_path / "artifact-runtime" + artifact_python = runtime / "bin" / "python" + artifact_python.parent.mkdir(parents=True) + artifact_python.touch() + monkeypatch.setenv("AKASHIC_PLUGIN_FIXTURE_PYTHON", str(artifact_python)) + + _content, feed = _stage_plugins(tmp_path / "stage") + + assert (feed / "mcp" / ".venv").resolve() == runtime.resolve() + + +def test_stage_plugins_requires_fixture_python_before_writes( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + monkeypatch.delenv("AKASHIC_PLUGIN_FIXTURE_PYTHON", raising=False) + stage = tmp_path / "missing-runtime" + + with pytest.raises(KeyError, match="AKASHIC_PLUGIN_FIXTURE_PYTHON"): + _stage_plugins(stage) + + assert not stage.exists() + + +def _seed_item(data_root: Path, now: datetime) -> None: + config = backend.load_config(data_root) + connection = backend._connect(config) + try: + connection.execute( + """ + INSERT INTO items( + event_id, source_id, source_name, source_type, title, content, + url, author, published_at, first_seen_at, last_seen_at, + emitted_at, content_hash + ) VALUES('event-1', 'source', 'Source', 'rss', 'Title', 'Body', + 'https://example.com/1', 'Author', ?, ?, ?, NULL, 'revision-1') + """, + (now.isoformat(), now.isoformat(), now.isoformat()), + ) + connection.commit() + finally: + connection.close() + + +def _sqlite_hashes(path: Path) -> dict[str, str]: + return { + candidate.name: hashlib.sha256(candidate.read_bytes()).hexdigest() + for candidate in sorted(path.parent.glob(path.name + "*")) + } @pytest.mark.asyncio -async def test_manager_boots_formal_feed_fetches_empty_and_drains( +async def test_manager_content_candidate_and_timer_handoff( tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, ) -> None: - """走真实 stdio 与 exact source lease,并证明空库运行不访问外部 Feed。""" + """证明唯一正式轮询 owner、静默候选和有序热更新。""" + + # 1. 加载真实插件,让稳定 Feed Root 提交一条完整 item。 + now = datetime(2026, 8, 23, 10, tzinfo=UTC) + timers: list[_Timer] = [] - # 1. staging 中没有订阅源,正式 poller 只会初始化临时数据库。 - plugin_root = _stage_plugin(tmp_path) + def timer_factory() -> _Timer: + timer = _Timer(now) + timers.append(timer) + return timer + + monkeypatch.setattr(plugin_manager_module, "AsyncioOneShotTimer", timer_factory) + content_dir, feed_dir = _stage_plugins(tmp_path) workspace = tmp_path / "workspace" manager = PluginManager( - plugin_dirs=[plugin_root.parent], + plugin_dirs=[content_dir, feed_dir], event_bus=EventBus(), tool_registry=None, workspace=workspace, installed_cache_root=tmp_path / "cache", ) - adapter = ProactiveActivityAdapter(manager.composition_generation_host) - activity = ActivityHost((adapter,)) - manager.bind_activity_host(activity) - - # 2. 通过 committed Activity binding 调用真实 MCP source。 - snapshot = None - generation_id = None - lease = None + await manager.load_all() + snapshot = manager.current_snapshot + assert snapshot is not None and snapshot.mcp_server_registry is not None + runtime = manager.composition_generation_host.get( + snapshot.generations["feed"].generation_id + ) + assert runtime is not None and runtime.mcp is not None + assert runtime.mcp.server("feed").tool_names == ( + "feed_manage", + "feed_query", + ) + feed_data = workspace / "plugin-data" / "feed-builtin" + _seed_item(feed_data, now) + lifecycle = asyncio.create_task(manager.run_runtime_services()) + formal_reader: sqlite3.Connection | None = None try: - await manager.load_all() - snapshot = manager.current_snapshot - assert snapshot is not None and snapshot.mcp_server_registry is not None - generation = next(iter(snapshot.generations.values())) - generation_id = generation.generation_id - runtime = manager.composition_generation_host.get(generation_id) - assert runtime is not None and runtime.mode == "formal" - assert runtime.mcp is not None and runtime.mcp.state == "ready" - assert "get_proactive_events" in runtime.mcp.server("feed").tool_names - - binding = activity.active - assert binding is not None - proactive = binding.child_bindings["proactive_components"] - assert isinstance(proactive, ProactiveRuntimeBinding) - lease = manager.snapshot_store.lease(snapshot.snapshot_id) - result = await proactive.source("subscriptions").fetch(lease) - assert isinstance(result, FetchEmpty) + await _eventually(lambda: sum(len(timer.handles) for timer in timers) == 1) + formal_timer = next(timer for timer in timers if timer.handles) + formal_timer.handles[0].fire() + content_path = ( + workspace / "plugin-data" / "content-builtin" / "content.sqlite3" + ) + content_store = ContentStore(content_path) + await _eventually( + lambda: content_store.state_counts().get("pending") == 1 + ) + await _eventually(lambda: len(formal_timer.handles) == 2) + with sqlite3.connect(feed_data / "feed_mcp.sqlite3") as connection: + assert connection.execute("PRAGMA integrity_check").fetchone() == ("ok",) + payload = connection.execute( + "SELECT payload_json FROM content_exports" + ).fetchone()[0] + assert '\"content\":\"Body\"' in payload + + # 2. 候选可以握手私有 MCP,但不能轮询或写正式状态。 + formal_reader = sqlite3.connect(content_path) + assert formal_reader.execute("SELECT COUNT(*) FROM items").fetchone() == (1,) + feed_hashes = _sqlite_hashes(feed_data / "feed_mcp.sqlite3") + content_hashes = _sqlite_hashes(content_path) + with (feed_dir / "plugin.py").open("a", encoding="utf-8") as handle: + handle.write("\n# candidate handoff fixture\n") + candidate = await manager.prepare_candidate("feed") + assert candidate is not None and candidate.runtime_snapshot is not None + candidate_root = candidate.runtime_snapshot.composition_root + assert candidate_root is not None + assert candidate_root.plugin_runtime("feed").data_dir != feed_data + assert sum(len(timer.handles) for timer in timers) == 2 + assert _sqlite_hashes(feed_data / "feed_mcp.sqlite3") == feed_hashes + assert _sqlite_hashes(content_path) == content_hashes + + # 3. 发布先取消旧等待,再由新稳定 Root 注册 Timer。 + result = await manager.publish_prepared("feed") + assert result["publication_state"] == "committed" + await _eventually(lambda: sum(len(timer.handles) for timer in timers) == 3) + assert content_store.state_counts() == {"pending": 1} + assert (await formal_timer.handles[1].result()).status is TimerStatus.CANCELLED + active = [ + handle + for timer in timers + for handle in timer.handles + if not handle.future.done() + ] + assert len(active) == 1 finally: - if lease is not None: - await lease.release() + if formal_reader is not None: + formal_reader.close() + lifecycle.cancel() + _ = await asyncio.gather(lifecycle, return_exceptions=True) await manager.terminate_all() - # 3. formal SQLite 完整,terminate 后 runtime、Root 与 activity 全释放。 - database_path = workspace / "plugin-data" / "feed-builtin" / "feed_mcp.sqlite3" - with sqlite3.connect(database_path) as database: - assert database.execute("PRAGMA integrity_check").fetchone() == ("ok",) - assert activity.active is None - assert manager.composition_generation_host.get(generation_id) is None - assert snapshot is not None and snapshot.composition_root is not None - assert snapshot.composition_root.receipt().effects == () - assert snapshot.composition_root.topology_view().listeners == () + assert all(handle.future.done() for timer in timers for handle in timer.handles) + + +@pytest.mark.asyncio +async def test_legacy_mcp_owner_stops_before_new_timer_starts( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + """证明旧 lifespan owner 退场后,Timer ownership 才开始。""" + + # 1. 启动真实 managed MCP,用 lifespan 表示旧轮询 owner。 + now = datetime(2026, 8, 23, 10, tzinfo=UTC) + timers: list[_Timer] = [] + + def timer_factory() -> _Timer: + timer = _Timer(now) + timers.append(timer) + return timer + + monkeypatch.setattr(plugin_manager_module, "AsyncioOneShotTimer", timer_factory) + content_dir, feed_dir = _stage_legacy_plugins(tmp_path) + workspace = tmp_path / "workspace" + manager = PluginManager( + plugin_dirs=[content_dir, feed_dir], + event_bus=EventBus(), + workspace=workspace, + installed_cache_root=tmp_path / "cache", + ) + await manager.load_all() + owner_log = ( + workspace / "plugin-data" / "feed-builtin" / "legacy-owner.jsonl" + ) + await _eventually(owner_log.is_file) + lifecycle = asyncio.create_task(manager.run_runtime_services()) + try: + assert sum(len(timer.handles) for timer in timers) == 0 + + # 2. 准备并发布真实 Timer + Content 实现。 + _replace_with_current_feed(feed_dir) + candidate = await manager.prepare_candidate("feed") + assert candidate is not None + assert sum(len(timer.handles) for timer in timers) == 0 + result = await manager.publish_prepared("feed") + assert result["publication_state"] == "committed" + await _eventually( + lambda: sum(len(timer.handles) for timer in timers) == 1 + ) + await _eventually( + lambda: '"event": "stopped"' in owner_log.read_text(encoding="utf-8") + ) + + # 3. 对比进程证据,而不只检查进程内对象状态。 + events = [ + json.loads(line) + for line in owner_log.read_text(encoding="utf-8").splitlines() + ] + assert [event["event"] for event in events] == ["started", "stopped"] + scheduled = [value for timer in timers for value in timer.schedule_times] + assert len(scheduled) == 1 + assert events[1]["time_ns"] <= scheduled[0] + finally: + lifecycle.cancel() + _ = await asyncio.gather(lifecycle, return_exceptions=True) + await manager.terminate_all() diff --git a/tests/test_mcp_v3.py b/tests/test_mcp_v3.py index b282ff0..55ee06a 100644 --- a/tests/test_mcp_v3.py +++ b/tests/test_mcp_v3.py @@ -1,10 +1,9 @@ from __future__ import annotations -import asyncio -import importlib.util -import logging +import json +import os from pathlib import Path -from types import SimpleNamespace +import subprocess import pytest @@ -13,106 +12,130 @@ RUN_MCP_PATH = Path(__file__).resolve().parents[1] / "mcp" / "run_mcp.py" -def _load_module(path: Path, name: str): - spec = importlib.util.spec_from_file_location(name, path) - assert spec is not None and spec.loader is not None +_ARTIFACT_PROBE = r""" +import importlib.util +import json +import logging +from logging.handlers import RotatingFileHandler +from pathlib import Path +import os +import sys + + +def load_module(path): + spec = importlib.util.spec_from_file_location("feed_artifact_probe", path) + if spec is None or spec.loader is None: + raise RuntimeError(f"cannot load probe module: {path}") module = importlib.util.module_from_spec(spec) spec.loader.exec_module(module) return module -def test_recording_fetch_ack_and_lifespan_are_zero_write( - tmp_path: Path, - monkeypatch: pytest.MonkeyPatch, -) -> None: - bridge = _load_module(MCP_BRIDGE_PATH, "feed_test_mcp_bridge") - monkeypatch.setenv("FEED_BACKEND", "recording") - monkeypatch.setenv("AKA_PLUGIN_DATA_DIR", str(tmp_path)) - - class UnexpectedPoller: - def __init__(self) -> None: - raise AssertionError("recording 不得创建 FeedPoller") - - monkeypatch.setattr(bridge, "FeedPoller", UnexpectedPoller) - monkeypatch.setattr( - bridge, - "_live_backend", - lambda: (_ for _ in ()).throw(AssertionError("recording 不得加载后端")), +action = sys.argv[1] +module = load_module(Path(sys.argv[2])) +if action == "tools": + os.environ["FEED_BACKEND"] = "recording" + server = module.create_mcp_server() + result = { + "tools": sorted(tool.name for tool in server._tool_manager.list_tools()) + } +elif action == "recording-error": + os.environ["FEED_BACKEND"] = "recording" + try: + module._live_backend() + except RuntimeError as error: + result = {"error_type": type(error).__name__, "message": str(error)} + else: + raise AssertionError("recording backend unexpectedly reached live backend") +elif action == "logging": + runtime = Path(sys.argv[3]) + module._setup_logging(runtime) + handlers = logging.getLogger().handlers + rotating = [ + handler for handler in handlers if isinstance(handler, RotatingFileHandler) + ] + result = { + "rotating_handlers": len(rotating), + "backup_count": rotating[0].backupCount, + "max_bytes": rotating[0].maxBytes, + } + for handler in handlers: + handler.close() + logging.getLogger().handlers.clear() +else: + raise ValueError(f"unknown artifact probe: {action}") +print(json.dumps(result, sort_keys=True)) +""" + + +def _run_artifact_probe(action: str, module: Path, *args: Path) -> dict[str, object]: + """Run one module oracle inside the explicitly supplied service artifact.""" + + # 1. Resolve the required artifact boundary without a pytest fallback. + artifact_python = Path(os.environ["AKASHIC_PLUGIN_FIXTURE_PYTHON"]) + environment = os.environ.copy() + environment.pop("PYTHONPATH", None) + environment["PYTHONDONTWRITEBYTECODE"] = "1" + + # 2. Execute the exact module and decode its fixed observable result. + completed = subprocess.run( + [ + str(artifact_python), + "-c", + _ARTIFACT_PROBE, + action, + str(module), + *(str(arg) for arg in args), + ], + check=True, + capture_output=True, + text=True, + env=environment, ) + result = json.loads(completed.stdout) + if not isinstance(result, dict): + raise TypeError("Feed artifact probe must return a JSON object") + return result - server = bridge.create_mcp_server() - assert server is not None - assert bridge._fetch_proactive_events() == {"status": "empty"} - with pytest.raises(RuntimeError, match="不允许确认"): - bridge._acknowledge_proactive_events(["event-1"]) - assert list(tmp_path.iterdir()) == [] - - -def test_live_results_are_explicit_typed_payloads(monkeypatch: pytest.MonkeyPatch) -> None: - bridge = _load_module(MCP_BRIDGE_PATH, "feed_test_mcp_bridge_live") - monkeypatch.delenv("FEED_BACKEND", raising=False) - backend = SimpleNamespace( - get_proactive_events=lambda **_: [], - acknowledge_events=lambda ids, feedback=None: { - "acknowledged": list(ids), - "failed": [], - }, - ) - monkeypatch.setattr(bridge, "_live_backend", lambda: backend) - assert bridge._fetch_proactive_events() == {"status": "empty"} - backend.get_proactive_events = lambda **_: [{"event_id": "one", "kind": "content"}] - assert bridge._fetch_proactive_events() == { - "status": "items", - "items": [{"event_id": "one", "kind": "content"}], +def test_recording_mcp_exposes_only_user_driven_tools() -> None: + assert _run_artifact_probe("tools", MCP_BRIDGE_PATH) == { + "tools": ["feed_manage", "feed_query"] } - assert bridge._acknowledge_proactive_events(["one"]) == { - "status": "committed", - "ids": ["one"], - } - assert bridge._proactive_ack_payload( - ["one", "two"], {"acknowledged": ["one"], "failed": ["two"]} - )["status"] == "failure" - assert bridge._proactive_ack_payload( - [], {"acknowledged": [], "failed": []} - ) == {"status": "skipped", "reason": "no_ids"} -def test_proactive_cursor_returns_every_event_once( - monkeypatch: pytest.MonkeyPatch, +def test_recording_user_tool_fails_before_backend_access() -> None: + result = _run_artifact_probe("recording-error", MCP_BRIDGE_PATH) + + assert result["error_type"] == "RuntimeError" + assert "recording backend" in str(result["message"]) + + +def test_runner_uses_three_bounded_log_rotations( + tmp_path: Path, ) -> None: - bridge = _load_module(MCP_BRIDGE_PATH, "feed_test_mcp_bridge_pages") - monkeypatch.delenv("FEED_BACKEND", raising=False) - events = [{"event_id": f"event-{index}", "kind": "content"} for index in range(51)] + assert _run_artifact_probe("logging", RUN_MCP_PATH, tmp_path) == { + "rotating_handlers": 1, + "backup_count": 3, + "max_bytes": 5 * 1024 * 1024, + } - def fetch(*, offset: int, limit: int): - return events[offset : offset + limit] - monkeypatch.setattr( - bridge, - "_live_backend", - lambda: SimpleNamespace(get_proactive_events=fetch), - ) +def test_artifact_probe_requires_fixture_python( + monkeypatch: pytest.MonkeyPatch, +) -> None: + monkeypatch.delenv("AKASHIC_PLUGIN_FIXTURE_PYTHON", raising=False) - first = bridge._fetch_proactive_events(limit=50) - assert first["cursor"] == "feed-offset:50" - second = bridge._fetch_proactive_events(limit=50, cursor=first["cursor"]) - combined = [*first["items"], *second["items"]] - assert [item["event_id"] for item in combined] == [ - f"event-{index}" for index in range(51) - ] - assert "cursor" not in second + with pytest.raises(KeyError, match="AKASHIC_PLUGIN_FIXTURE_PYTHON"): + _run_artifact_probe("tools", MCP_BRIDGE_PATH) -def test_runner_configures_stderr_without_runtime_log( +def test_artifact_probe_rejects_missing_interpreter( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, ) -> None: - runner = _load_module(RUN_MCP_PATH, "feed_test_run_mcp") - monkeypatch.setenv("AKA_PLUGIN_DATA_DIR", str(tmp_path)) - runner._setup_logging() - assert all( - not isinstance(handler, logging.FileHandler) - for handler in logging.getLogger().handlers - ) - assert list(tmp_path.iterdir()) == [] + missing = tmp_path / "missing-artifact" / "bin" / "python" + monkeypatch.setenv("AKASHIC_PLUGIN_FIXTURE_PYTHON", str(missing)) + + with pytest.raises(FileNotFoundError): + _run_artifact_probe("tools", MCP_BRIDGE_PATH) diff --git a/tests/test_plugin.py b/tests/test_plugin.py index 0ba824d..bd7df68 100644 --- a/tests/test_plugin.py +++ b/tests/test_plugin.py @@ -2,46 +2,75 @@ import inspect from pathlib import Path +from typing import cast import pytest + import plugin +from agent.control.timer import OneShotTimer from agent.plugin_composition import ( MCP_SERVERS, - PROACTIVE_COMPONENTS, + TIMERS, CompositionRoot, - PluginProactiveComponents, PluginRuntime, + PluginTimers, ) from agent.plugin_composition.mcp_slots import ( PluginMcpServers, _freeze_plugin_mcp_servers, ) -from agent.plugin_composition.proactive import _freeze_plugin_proactive_components -from agent.plugins.static_manifest import load_static_plugin_manifest from agent.plugins.composable import ComposablePlugin from agent.plugins.manager import _copy_validation_data -from plugin import FeedConfig, FeedProactiveConfig +from agent.plugins.static_manifest import load_static_plugin_manifest +from content_source import BoundContentSource, ContentSourceServices ROOT = Path(__file__).resolve().parents[1] +class _Content: + def submit(self, batch_id, items): + raise AssertionError((batch_id, items)) + + def unsettled(self, limit=100): + raise AssertionError(limit) + + def ack(self, settlement_ref): + raise AssertionError(settlement_ref) + + +class _Sources: + def __init__(self) -> None: + self.bound: list[str] = [] + + def bind(self, source_id: str) -> BoundContentSource: + self.bound.append(source_id) + return cast(BoundContentSource, _Content()) + + def test_pure_v3_exports_and_exact_apply() -> None: assert plugin.api_version == 3 assert plugin.name == "feed" - assert plugin.version == "3.0.0" + assert plugin.version == "3.1.0" assert plugin.skill_roots == ("skills",) assert tuple(inspect.signature(plugin.apply).parameters) == ("ctx", "config") assert ComposablePlugin.from_module(plugin).skill_roots == ("skills",) + assert "content.source.v1" in inspect.getsource(plugin) @pytest.mark.asyncio -async def test_apply_registers_mcp_and_source_without_data_writes(tmp_path: Path) -> None: +async def test_apply_registers_user_mcp_and_dormant_content_runtime( + tmp_path: Path, +) -> None: root = CompositionRoot("feed:test") servers = PluginMcpServers(root.instance_token) - components = PluginProactiveComponents(root.instance_token) + sources = _Sources() await root.context.provide(MCP_SERVERS, servers) - await root.context.provide(PROACTIVE_COMPONENTS, components) + await root.context.provide( + TIMERS, + PluginTimers(cast(OneShotTimer, object())), + ) + await root.context.provide(plugin.CONTENT_SOURCE, sources) data_dir = tmp_path / "plugin-data" await root.mount( ComposablePlugin.from_module(plugin), @@ -51,95 +80,48 @@ async def test_apply_registers_mcp_and_source_without_data_writes(tmp_path: Path plugin_dir=ROOT, data_dir=data_dir, workspace=tmp_path / "workspace", - config=FeedConfig( - proactive=FeedProactiveConfig(enabled=True), - ), + config=plugin.FeedConfig(), ), ) mcp = _freeze_plugin_mcp_servers(servers, root.instance_token)["feed"].definition - source = _freeze_plugin_proactive_components( - components, - root.instance_token, - {"feed": "feed:test"}, - ).source("subscriptions") + assert mcp.required_tools == ("feed_manage", "feed_query") + assert mcp.candidate_read_only_tools == () assert mcp.candidate_env == {"FEED_BACKEND": "recording"} - assert mcp.candidate_read_only_tools == ("get_proactive_events",) - assert source is not None - assert source.definition.channels == ("content",) - assert source.definition.mcp_server == "feed" + assert sources.bound == ["feed-subscriptions"] assert not data_dir.exists() - await root.dispose() - - -@pytest.mark.asyncio -async def test_disabled_proactive_omits_source(tmp_path: Path) -> None: - root = CompositionRoot("feed:disabled") - servers = PluginMcpServers(root.instance_token) - components = PluginProactiveComponents(root.instance_token) - await root.context.provide(MCP_SERVERS, servers) - await root.context.provide(PROACTIVE_COMPONENTS, components) - await root.mount( - ComposablePlugin.from_module(plugin), - name="feed", - runtime=PluginRuntime( - plugin_id="feed", - plugin_dir=ROOT, - data_dir=tmp_path / "plugin-data", - workspace=tmp_path / "workspace", - config=FeedConfig( - proactive=FeedProactiveConfig(enabled=False), - ), - ), + topology = root.topology_view() + assert topology.listeners == ( + "serial:runtime.started:feed", + "serial:runtime.stopping:feed", ) - catalog = _freeze_plugin_proactive_components( - components, - root.instance_token, - {"feed": "feed:disabled"}, - ) - assert catalog.sources == {} await root.dispose() -def test_static_manifest_freezes_recording_and_data_exclusions() -> None: - manifest = load_static_plugin_manifest(Path(__file__).resolve().parents[1]) +def test_static_manifest_freezes_tools_and_data_exclusions() -> None: + manifest = load_static_plugin_manifest(ROOT) assert manifest.name == "feed" - assert manifest.version == "3.0.0" + assert manifest.version == "3.1.0" assert manifest.api_version == 3 assert manifest.requirements == ("mcp/requirements.txt",) - assert manifest.exclude_data_paths == ( - "feed_mcp.sqlite3", - "feed_mcp.sqlite3-wal", - "feed_mcp.sqlite3-shm", - "source_scores.json", - "feed_cache.db", - "feed_cache.db-wal", - "feed_cache.db-shm", - ".feed-v2-migration.json", - ) - assert len(manifest.mcp_servers) == 1 + assert "feed_mcp.sqlite3" in manifest.exclude_data_paths + assert "feed_mcp.runtime.log.3" in manifest.exclude_data_paths + assert "feed_source.runtime.log.3" in manifest.exclude_data_paths server = manifest.mcp_servers[0] - assert server.required_tools == ("get_proactive_events", "acknowledge_events") - assert server.candidate_read_only_tools == ("get_proactive_events",) + assert server.required_tools == ("feed_manage", "feed_query") + assert server.candidate_read_only_tools == () assert server.candidate_env == (("FEED_BACKEND", "recording"),) -def test_candidate_copy_excludes_sqlite_and_sidecars(tmp_path: Path) -> None: +def test_candidate_copy_excludes_sqlite_logs_and_sidecars(tmp_path: Path) -> None: manifest = load_static_plugin_manifest(ROOT) source = tmp_path / "workspace" / "plugin-data" / "feed-builtin" source.mkdir(parents=True) - for name in ( - "feed_mcp.sqlite3", - "feed_mcp.sqlite3-wal", - "feed_mcp.sqlite3-shm", - "feed_cache.db", - "feed_cache.db-wal", - "feed_cache.db-shm", - "source_scores.json", - ".feed-v2-migration.json", - ): - (source / name).write_text(f"secret:{name}", encoding="utf-8") + for name in manifest.exclude_data_paths: + path = source / name + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(f"formal:{name}", encoding="utf-8") (source / "candidate-visible.txt").write_text("visible", encoding="utf-8") target = tmp_path / "validation" / "feed" @@ -150,4 +132,4 @@ def test_candidate_copy_excludes_sqlite_and_sidecars(tmp_path: Path) -> None: ) assert inventory == ("candidate-visible.txt",) - assert (target / "candidate-visible.txt").read_text(encoding="utf-8") == "visible" + assert (target / "candidate-visible.txt").read_text() == "visible"