From f62df2e45f5fb6c5a5c5a19aa9d4918248b82921 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Mon, 28 Sep 2026 21:11:11 +0800 Subject: [PATCH 1/5] fix(decision-context): isolate capture capacity by source Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/capabilities/decision_context/README.md | 28 +++- .../decision_context/README.zh-CN.md | 19 ++- .../capabilities/decision_context/capture.py | 58 ++++++-- ...test_decision_context_capture_isolation.py | 136 ++++++++++++++++++ 4 files changed, 226 insertions(+), 15 deletions(-) create mode 100644 tests/capabilities/test_decision_context_capture_isolation.py diff --git a/loopx/capabilities/decision_context/README.md b/loopx/capabilities/decision_context/README.md index dc4f4490ba..a7682bb7c2 100644 --- a/loopx/capabilities/decision_context/README.md +++ b/loopx/capabilities/decision_context/README.md @@ -325,7 +325,14 @@ The mode-0600 SQLite spool binds to one goal/agent and records bounded public-sa scan receipts plus **private replay cursors**. It contains no source bodies. Ticks are serialized; batch insertion and capture cursor advancement commit together. Failed scans keep their cursor, and capacity exhaustion reports -`backpressure` without dropping pending batches. A changed source binding reports +`backpressure` without dropping pending batches. Each enrolled source now has +a reserved active window of `max(1, floor(max_pending_batches / source_count))`. +When that window is full, only that source reports `source_backpressure`; it is +not scanned and its capture cursor and last successful read stay unchanged. +Quiet sources keep their shares for later changes; idle shares and the integer +remainder are not lent to busy sources. This changes the previous global-only +admission default for multi-source profiles. Single-source behavior is unchanged. +A changed source binding reports `binding_changed`, requiring an explicit rebase or a separately scoped new spool. Do not store the spool or its journal in a public repository. @@ -416,9 +423,22 @@ has a separate cap of N; new hold/restart operations stop after 2N audit records or repeated capacity increase is performed. At that bound, preserve/export the private spool and make an explicit retention decision; increasing the limit is not evidence consumption. A hold is explicit, not an automatic fairness policy. -It can isolate a noisy source, but exhaustion can recur if other sources are -not reviewed. Backpressured sources now retry on the next tick when capacity is -available instead of waiting an additional scan interval. +It can isolate a noisy source; reserved source windows additionally prevent a +busy source from borrowing other sources' future capacity. Neither mechanism +replaces semantic review: a full source needs oldest-batch review or explicit +`capture-diagnose` and guarded recovery when replay is unavailable. Both pressure +statuses retry on the next tick when capacity is available instead of waiting +an additional scan interval. + +Existing over-share batches are retained, never evicted or silently held. On +upgrade they can still consume shared capacity until explicit review/recovery. +If the global capacity is smaller than the number of enrolled sources, isolation +cannot be guaranteed: status reports `reservation_capacity_sufficient=false` +and the global cap remains authoritative. `capture-status.capacity_policy` +reports both bounds and each source reports `pending_capacity`, `review_required` +and `recovery_diagnosis_required`. These are work hints, not proof that replay +failed, and do not grant recovery or settlement authority. The global cap +continues to count pending batches from sources no longer enrolled as well. `capture-status` separates active `pending_batch_count`, unresolved `held_batch_count`, per-source `acquisition_held` and diff --git a/loopx/capabilities/decision_context/README.zh-CN.md b/loopx/capabilities/decision_context/README.zh-CN.md index 4733e3d764..75c9bafa42 100644 --- a/loopx/capabilities/decision_context/README.zh-CN.md +++ b/loopx/capabilities/decision_context/README.zh-CN.md @@ -288,7 +288,11 @@ loopx decision-context capture --goal-id --agent-id \ 权限为 0600 的私有 SQLite spool 绑定单个 goal/agent,只保存有界 scan receipt 和私有回放游标,不保存正文。采集事务串行执行,批次和采集游标一起提交;失败不前移 -游标,容量耗尽报 `backpressure` 而不丢弃待审阅批次。来源绑定变化报 +游标,容量耗尽报 `backpressure` 而不丢弃待审阅批次。每个已登记采集来源现在保留 +`max(1, floor(max_pending_batches / 来源数))` 个活跃批次窗口。单来源窗口用完时, +只有该来源报 `source_backpressure`,不调用 provider、不推进采集游标或成功读取时间; +安静来源的空闲份额与整数余数不借给高频来源,保留给后来的独立变化。此项改变多来源 +profile 原先只有全局上限的默认准入,单来源行为不变。来源绑定变化报 `binding_changed`,需显式 rebase 或启用独立新 spool。数据库及 journal 均不得公开。 每个来源从状态中的 `next_batch_id` 开始回读: @@ -357,8 +361,17 @@ python3 -m pytest -q tests/capabilities/test_decision_context_capture.py `max_pending_batches=N` 继续限制活跃批次;另最多保留 N 条未解决历史, 2N 条审计记录后停止新增 hold/restart(每条适用回执最多再 rollback 一次)。 不会自动删除、压缩或无限扩容。达到上限需保留/导出私有 spool 后明确处理保留策略。 -这是显式来源隔离,不是默认公平调度;若其他来源长期不被审阅,仍可能再次背压。 -行为变化:曾背压的来源在容量释放后的下一 tick 可重试,不再多等一个扫描间隔。 +hold 是显式恢复操作,默认来源窗口则阻止高频来源借走其他来源的未来容量。 +两者都不能代替语义审阅:来源窗口满时须审阅最旧批次;无法回放时先 +`capture-diagnose`,再按授权走受保护恢复。两种背压状态在容量释放后的下一 tick +都可重试,不再多等一个扫描间隔。 + +升级前已有的超份额批次完整保留,不自动删除或转 held;它们仍占全局容量,直到显式 +审阅或恢复。若全局容量小于登记来源数,不能保证隔离,状态会明确报告 +`reservation_capacity_sufficient=false`,仍严格遵守全局上限。 +`capture-status.capacity_policy` 给出两层上限,各来源给出 `pending_capacity`、 +`review_required` 和 `recovery_diagnosis_required`。这些是工作提示,不代表已探测到 +回放失败,也不授予恢复或结算权限。已退出采集登记的来源留下的 pending 仍计入全局容量。 `capture-status` 分开报告 active pending、held 历史、每来源 acquisition hold, 并明确 `semantic_review_completion=not_inferred_from_capture`。`last_checked_at` diff --git a/loopx/capabilities/decision_context/capture.py b/loopx/capabilities/decision_context/capture.py index a525192ebb..ba06f95bc5 100644 --- a/loopx/capabilities/decision_context/capture.py +++ b/loopx/capabilities/decision_context/capture.py @@ -13,6 +13,7 @@ from contextlib import closing from dataclasses import replace from datetime import datetime, timezone +from enum import Enum from pathlib import Path from typing import Any @@ -36,6 +37,17 @@ _READ_FAILURE_STATUSES = frozenset({"provider_failed", "failed", "unavailable"}) +class CapturePressure(str, Enum): + GLOBAL = "backpressure" + SOURCE = "source_backpressure" + + +def _source_capacity(total: int, source_count: int) -> int: + # Do not lend idle shares: a quiet source may change after a busy source + # has filled its window. An undersized global cap still takes precedence. + return max(1, total // max(1, source_count)) + + class CaptureReplayError(ValueError): """Typed recovery diagnosis; never classify provider exception prose.""" @@ -113,6 +125,7 @@ def _status( sources: tuple[DecisionSourceSpec, ...], *, now: datetime, + capacity: int, ) -> dict[str, Any]: has_recovery = ( db.execute("SELECT 1 FROM sqlite_master WHERE name='capture_holds'").fetchone() @@ -121,6 +134,7 @@ def _status( columns = _source_columns(db) rows = [] freshness_rows = [] + source_capacity = _source_capacity(capacity, len(sources)) for spec in sources: source_id = spec.source_id source = db.execute( @@ -163,6 +177,9 @@ def _status( "status": source["status"] if source else "never_checked", "pending_batch_count": pending[0], "next_batch_id": pending[1], + "pending_capacity": source_capacity, + "review_required": bool(pending[0]), + "recovery_diagnosis_required": pending[0] >= source_capacity, "held_batch_count": db.execute( "SELECT count(*) FROM held_batches WHERE source_id=?", (source_id,) ).fetchone()[0] @@ -184,6 +201,15 @@ def _status( observed_at=now, rows=freshness_rows ), "pending_batch_count": db.execute("SELECT count(*) FROM batches").fetchone()[0], + "capacity_policy": { + "kind": "reserved_source_windows", + "max_pending_batches": capacity, + "max_pending_batches_per_source": source_capacity, + "enrolled_source_count": len(sources), + "reservation_capacity_sufficient": capacity >= len(sources), + "existing_batches_evicted": False, + "full_source_action": "review_oldest_or_explicit_capture_recovery", + }, "held_batch_count": db.execute("SELECT count(*) FROM held_batches").fetchone()[ 0 ] @@ -269,7 +295,9 @@ def capture_profile_sources( return { "activation": activation, "executed": False, - **_status(db, sources, now=now), + **_status( + db, sources, now=now, capacity=profile.capture_max_pending_batches + ), } def record_health(tick_status: str, freshness: Mapping[str, Any] | None) -> None: @@ -330,6 +358,9 @@ def _execute_capture_tick( db = _open_spool(spool_path, goal_id=goal_id, agent_id=agent_id) try: reviewed = load_private_decision_cursors(cursor_path, profile=profile) + source_capacity = _source_capacity( + profile.capture_max_pending_batches, len(sources) + ) for source in sources: if ( db.execute( @@ -377,17 +408,26 @@ def _execute_capture_tick( if ( row and row["checked_at"] - and row["status"] != "backpressure" + and row["status"] + not in (CapturePressure.GLOBAL, CapturePressure.SOURCE) and (now - datetime.fromisoformat(row["checked_at"])).total_seconds() < profile.capture_interval_seconds ): continue cursor = row["cursor"] if row else reviewed.get(source.source_id) - status = "backpressure" - if ( - db.execute("SELECT count(*) FROM batches").fetchone()[0] - < profile.capture_max_pending_batches - ): + pending_total = db.execute("SELECT count(*) FROM batches").fetchone()[0] + pending_source = db.execute( + "SELECT count(*) FROM batches WHERE source_id=?", (source.source_id,) + ).fetchone()[0] + pressure = ( + CapturePressure.GLOBAL + if pending_total >= profile.capture_max_pending_batches + else CapturePressure.SOURCE + if pending_source >= source_capacity + else None + ) + status = pressure.value if pressure else "provider_failed" + if pressure is None: try: scan = providers[source.provider_id].scan( source=source, @@ -460,7 +500,9 @@ def _execute_capture_tick( result = { "activation": activation, "executed": True, - **_status(db, sources, now=now), + **_status( + db, sources, now=now, capacity=profile.capture_max_pending_batches + ), } db.commit() return result diff --git a/tests/capabilities/test_decision_context_capture_isolation.py b/tests/capabilities/test_decision_context_capture_isolation.py new file mode 100644 index 0000000000..dc92427ea8 --- /dev/null +++ b/tests/capabilities/test_decision_context_capture_isolation.py @@ -0,0 +1,136 @@ +"""Capture capacity must isolate producers without inventing review authority.""" + +import json +import sqlite3 + +import pytest + +from loopx.capabilities.decision_context.capture import capture_profile_sources +from loopx.capabilities.decision_context.providers import ( + LocalFileDecisionSourceProvider, +) +from test_decision_context_capture import setup as capture_setup, settle_batch + + +@pytest.fixture +def pair(tmp_path): + args, payload, busy = capture_setup.__wrapped__(tmp_path) + quiet = tmp_path / "quiet.txt" + quiet.write_text("quiet baseline") + payload["sources"].append( + dict(payload["sources"][0], source_id="quiet", private_locator=str(quiet)) + ) + payload["automation"].update( + source_ids=[s["source_id"] for s in payload["sources"]], max_pending_batches=6 + ) + args["profile_path"].write_text(json.dumps(payload)) + return args, payload, busy, quiet + + +def make_due(args): + with sqlite3.connect(args["spool_path"]) as db: + db.execute("UPDATE sources SET checked_at=NULL") + + +def batches(args): + with sqlite3.connect(args["spool_path"]) as db: + return db.execute("SELECT * FROM batches ORDER BY id").fetchall() + + +def test_busy_source_cannot_borrow_quiet_sources_future_capacity(pair): + args, _, busy, quiet = pair + initial = capture_profile_sources(**args, execute=True) + initial_rows = batches(args) + for revision in range(8): + busy.write_text(f"busy revision {revision}") + make_due(args) + result = capture_profile_sources(**args, execute=True) + assert result["sources"][0]["status"] == "source_backpressure" + assert result["sources"][0]["pending_batch_count"] == 3 + assert result["sources"][1]["status"] == "no_change" + assert result["pending_batch_count"] == 4 + assert result["source_freshness"]["all_fresh"] is False + assert result["sources"][0]["recovery_diagnosis_required"] + last_read = result["sources"][0]["last_read_at"] + + quiet.write_text("late independent decision evidence") + make_due(args) + result = capture_profile_sources(**args, execute=True) + assert result["sources"][1]["pending_batch_count"] == 2 + assert result["sources"][0]["last_read_at"] == last_read + assert result["pending_batch_count"] == 5 + assert set(initial_rows).issubset(set(batches(args))) + assert not args["cursor_path"].exists() + assert result["held_batch_count"] == initial["held_batch_count"] == 0 + + +def test_legacy_overfull_source_preserved_without_provider_reads(pair): + args, payload, busy, _ = pair + # A one-source profile can legitimately predate enrollment of another. + original_sources = payload["automation"]["source_ids"] + payload["automation"]["source_ids"] = original_sources[:1] + args["profile_path"].write_text(json.dumps(payload)) + for revision in range(4): + busy.write_text(f"historical revision {revision}") + if args["spool_path"].exists(): + make_due(args) + capture_profile_sources(**args, execute=True) + original_rows = batches(args) + payload["automation"]["source_ids"] = original_sources + args["profile_path"].write_text(json.dumps(payload)) + make_due(args) + calls = [] + + class ObservedProvider(LocalFileDecisionSourceProvider): + def scan(self, **kwargs): + calls.append(kwargs["source"].source_id) + return super().scan(**kwargs) + + before = args["spool_path"].read_bytes() + preview = capture_profile_sources(**args) + assert args["spool_path"].read_bytes() == before + assert preview["capacity_policy"]["max_pending_batches_per_source"] == 3 + result = capture_profile_sources( + **args, + execute=True, + source_provider_overrides={ + "local-authority": ObservedProvider( + provider_id="local-authority", max_bytes=4096 + ) + }, + ) + assert calls == ["quiet"] + assert result["sources"][0]["pending_batch_count"] == 4 + assert result["sources"][0]["status"] == "source_backpressure" + assert result["sources"][1]["status"] == "completed" + assert set(original_rows).issubset(set(batches(args))) + assert not args["cursor_path"].exists() + + +def test_review_releases_source_window_without_additional_interval(pair): + args, payload, _, _ = pair + payload["automation"]["max_pending_batches"] = 2 + args["profile_path"].write_text(json.dumps(payload)) + first = capture_profile_sources(**args, execute=True) + # Settle only quiet's batch, so busy hits its source window with global + # room left. The busy batch is still exact-readable. + settle_batch(args, first["sources"][1]["next_batch_id"]) + make_due(args) + capture_profile_sources(**args, execute=True) + # Quiet was retired later in source order; the next tick sees global room. + pressure = capture_profile_sources(**args, execute=True) + assert pressure["sources"][0]["status"] == "source_backpressure" + settle_batch(args, first["sources"][0]["next_batch_id"]) + result = capture_profile_sources(**args, execute=True) + assert result["sources"][0]["status"] == "no_change" + assert result["pending_batch_count"] == 0 + + +def test_undersized_total_is_explicit_and_never_overflows(pair): + args, payload, _, _ = pair + payload["automation"]["max_pending_batches"] = 1 + args["profile_path"].write_text(json.dumps(payload)) + result = capture_profile_sources(**args, execute=True) + assert result["capacity_policy"]["reservation_capacity_sufficient"] is False + assert result["sources"][1]["status"] == "backpressure" + assert result["pending_batch_count"] == 1 From 48b74cab0e7927f17f20c66e064a25a1f2901d28 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Mon, 28 Sep 2026 21:23:23 +0800 Subject: [PATCH 2/5] fix(decision-context): bound capture ticks and resume oldest sources Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/capabilities/decision_context/README.md | 15 ++++- .../decision_context/README.zh-CN.md | 9 ++- .../capabilities/decision_context/capture.py | 27 +++++++- .../capabilities/decision_context/profile.py | 9 +++ ...test_decision_context_capture_isolation.py | 66 +++++++++++++++++++ 5 files changed, 122 insertions(+), 4 deletions(-) diff --git a/loopx/capabilities/decision_context/README.md b/loopx/capabilities/decision_context/README.md index a7682bb7c2..5f5e216563 100644 --- a/loopx/capabilities/decision_context/README.md +++ b/loopx/capabilities/decision_context/README.md @@ -299,7 +299,8 @@ Add these fields to an existing private profile's `automation` object: "fail_open": true, "source_ids": ["source:authority:baseline"], "interval_seconds": 900, - "max_pending_batches": 1000 + "max_pending_batches": 1000, + "max_sources_per_tick": 8 } ``` @@ -313,7 +314,17 @@ loopx decision-context capture --goal-id --agent-id \ --cursor-state --format json ``` -Add `--execute` for one tick. Use `capture-status` with the same arguments +Add `--execute` for one tick. + +Each tick attempts at most `automation.max_sources_per_tick` providers (default +8, integer 1–64), in oldest-attempt-first order. Failed calls count toward this +budget; held, pressured, and interval-skipped sources do not. Deferred sources +keep their cursors and freshness timestamps and resume on later host ticks. +`scan_budget` reports attempts and deferred source IDs. This bounds provider +calls, not wall time: providers must still honor their timeout and the host +must enforce its process deadline. Profile-race rollback remains atomic. + +Use `capture-status` with the same arguments and without `--execute` for readback. Configure a host scheduler to invoke the tick; the capability enforces `interval_seconds`, while the host owns process startup, an outer process timeout, and stop/uninstall. No model heartbeat is diff --git a/loopx/capabilities/decision_context/README.zh-CN.md b/loopx/capabilities/decision_context/README.zh-CN.md index 75c9bafa42..1a8855dd16 100644 --- a/loopx/capabilities/decision_context/README.zh-CN.md +++ b/loopx/capabilities/decision_context/README.zh-CN.md @@ -268,7 +268,8 @@ cursor。 "fail_open": true, "source_ids": ["source:authority:baseline"], "interval_seconds": 900, - "max_pending_batches": 1000 + "max_pending_batches": 1000, + "max_sources_per_tick": 8 } ``` @@ -281,6 +282,12 @@ loopx decision-context capture --goal-id --agent-id \ --cursor-state --format json ``` +每轮最多调用 `automation.max_sources_per_tick` 个 provider(默认 8,整数 1–64), +优先读取最久未尝试的来源。失败调用消耗名额;hold、背压和未到读取间隔的来源不消耗。 +延期来源保留采集游标和 freshness 时间,下轮继续;`scan_budget` 报告尝试数与延期来源。 +该预算限制调用次数,不保证墙钟耗时;provider 仍须遵守 timeout,宿主仍须施加进程期限。 +profile 并发修改时仍保持整轮原子回滚。 + `capture-status` 使用相同参数但不带 `--execute`,只读回查。宿主负责定时调用、 进程总超时和启动/卸载;capability 执行配置中的采集间隔,不创建模型 heartbeat。 私有接入方调用 `loopx.capabilities.decision_context.capture` 中的 diff --git a/loopx/capabilities/decision_context/capture.py b/loopx/capabilities/decision_context/capture.py index ba06f95bc5..9468636d31 100644 --- a/loopx/capabilities/decision_context/capture.py +++ b/loopx/capabilities/decision_context/capture.py @@ -361,7 +361,23 @@ def _execute_capture_tick( source_capacity = _source_capacity( profile.capture_max_pending_batches, len(sources) ) - for source in sources: + checked = { + row["source_id"]: row["checked_at"] + for row in db.execute("SELECT source_id, checked_at FROM sources") + } + # Stable oldest-attempt-first order lets subsequent ticks resume the + # remaining sources, including when a provider repeatedly fails. + scan_order = sorted( + sources, + key=lambda source: ( + datetime.fromisoformat(checked[source.source_id]).timestamp() + if checked.get(source.source_id) + else float("-inf") + ), + ) + attempted = 0 + deferred = [] + for source in scan_order: if ( db.execute( "SELECT 1 FROM sqlite_master WHERE name='capture_holds'" @@ -428,6 +444,10 @@ def _execute_capture_tick( ) status = pressure.value if pressure else "provider_failed" if pressure is None: + if attempted >= profile.capture_max_sources_per_tick: + deferred.append(source.source_id) + continue + attempted += 1 try: scan = providers[source.provider_id].scan( source=source, @@ -500,6 +520,11 @@ def _execute_capture_tick( result = { "activation": activation, "executed": True, + "scan_budget": { + "max_sources_per_tick": profile.capture_max_sources_per_tick, + "attempted_source_count": attempted, + "deferred_source_ids": deferred, + }, **_status( db, sources, now=now, capacity=profile.capture_max_pending_batches ), diff --git a/loopx/capabilities/decision_context/profile.py b/loopx/capabilities/decision_context/profile.py index 3e5fcc12d3..482edb3caa 100644 --- a/loopx/capabilities/decision_context/profile.py +++ b/loopx/capabilities/decision_context/profile.py @@ -59,6 +59,7 @@ "source_ids", "interval_seconds", "max_pending_batches", + "max_sources_per_tick", } @@ -172,6 +173,7 @@ class DecisionContextProfile: capture_source_ids: tuple[str, ...] = () capture_interval_seconds: int = 900 capture_max_pending_batches: int = 1000 + capture_max_sources_per_tick: int = 8 def provider_binding_map(self) -> dict[str, Mapping[str, Any]]: return { @@ -367,6 +369,11 @@ def normalize_decision_context_profile( field_name="max_pending_batches", maximum=10000, ) + capture_scan_limit = _positive_int( + automation.get("max_sources_per_tick", 8), + field_name="max_sources_per_tick", + maximum=MAX_DECISION_SOURCES, + ) if not fail_open: raise ValueError("decision-context providers must fail open") @@ -382,6 +389,7 @@ def normalize_decision_context_profile( capture_source_ids=tuple(capture_ids), capture_interval_seconds=capture_interval, capture_max_pending_batches=capture_capacity, + capture_max_sources_per_tick=capture_scan_limit, ) @@ -481,6 +489,7 @@ def resolve_decision_context_activation( "automatic_capture": config.automatic_capture, "capture_source_count": len(config.capture_source_ids), "capture_interval_seconds": config.capture_interval_seconds, + "capture_max_sources_per_tick": config.capture_max_sources_per_tick, } if not config.enabled: return status | {"status": "disabled"}, config diff --git a/tests/capabilities/test_decision_context_capture_isolation.py b/tests/capabilities/test_decision_context_capture_isolation.py index dc92427ea8..fd99259e59 100644 --- a/tests/capabilities/test_decision_context_capture_isolation.py +++ b/tests/capabilities/test_decision_context_capture_isolation.py @@ -134,3 +134,69 @@ def test_undersized_total_is_explicit_and_never_overflows(pair): assert result["capacity_policy"]["reservation_capacity_sufficient"] is False assert result["sources"][1]["status"] == "backpressure" assert result["pending_batch_count"] == 1 + + +def test_bounded_ticks_resume_deferred_sources_and_preserve_freshness(pair): + args, payload, _, _ = pair + payload["automation"]["max_sources_per_tick"] = 1 + args["profile_path"].write_text(json.dumps(payload)) + first = capture_profile_sources(**args, execute=True) + assert first["scan_budget"]["attempted_source_count"] == 1 + assert first["scan_budget"]["deferred_source_ids"] == ["quiet"] + assert first["sources"][1]["last_read_at"] is None + assert first["sources"][1]["last_checked_at"] is None + assert first["sources"][1]["status"] == "never_checked" + second = capture_profile_sources(**args, execute=True) + assert second["scan_budget"]["attempted_source_count"] == 1 + assert second["scan_budget"]["deferred_source_ids"] == [] + assert second["sources"][1]["status"] == "completed" + assert first["sources"][0] == second["sources"][0] + assert second["pending_batch_count"] == 2 + assert not args["cursor_path"].exists() + + +def test_failed_source_uses_budget_without_starving_oldest_source(pair): + args, payload, _, _ = pair + payload["automation"]["max_sources_per_tick"] = 1 + args["profile_path"].write_text(json.dumps(payload)) + calls = [] + + class FailingProvider(LocalFileDecisionSourceProvider): + def scan(self, **kwargs): + calls.append(kwargs["source"].source_id) + if len(calls) == 1: + raise RuntimeError("synthetic failure") + return super().scan(**kwargs) + + providers = { + "local-authority": FailingProvider( + provider_id="local-authority", max_bytes=4096 + ) + } + first = capture_profile_sources( + **args, execute=True, source_provider_overrides=providers + ) + assert first["sources"][0]["status"] == "provider_failed" + assert first["scan_budget"]["attempted_source_count"] == 1 + # Both are due, but the never-attempted source must win over the failed one. + with sqlite3.connect(args["spool_path"]) as db: + db.execute("UPDATE sources SET checked_at='2000-01-01T00:00:00+00:00'") + second = capture_profile_sources( + **args, execute=True, source_provider_overrides=providers + ) + assert calls == [payload["sources"][0]["source_id"], "quiet"] + assert second["sources"][0]["failure_streak"] == 1 + assert second["sources"][0]["last_read_at"] is None + assert second["scan_budget"]["deferred_source_ids"] == [calls[0]] + + +@pytest.mark.parametrize("invalid", [0, -1, 65, True, "8", 1.5]) +def test_capture_tick_budget_rejects_invalid_configuration(pair, invalid): + from loopx.capabilities.decision_context.profile import ( + normalize_decision_context_profile, + ) + + _, payload, _, _ = pair + payload["automation"]["max_sources_per_tick"] = invalid + with pytest.raises(ValueError, match="max_sources_per_tick"): + normalize_decision_context_profile(payload) From 711ef645eb6aa355a145df13c8f92150e534b278 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Mon, 28 Sep 2026 21:38:50 +0800 Subject: [PATCH 3/5] test(decision-context): assert stable deferred-source state Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- tests/capabilities/test_decision_context_capture_isolation.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/tests/capabilities/test_decision_context_capture_isolation.py b/tests/capabilities/test_decision_context_capture_isolation.py index fd99259e59..3f47dffe85 100644 --- a/tests/capabilities/test_decision_context_capture_isolation.py +++ b/tests/capabilities/test_decision_context_capture_isolation.py @@ -150,7 +150,8 @@ def test_bounded_ticks_resume_deferred_sources_and_preserve_freshness(pair): assert second["scan_budget"]["attempted_source_count"] == 1 assert second["scan_budget"]["deferred_source_ids"] == [] assert second["sources"][1]["status"] == "completed" - assert first["sources"][0] == second["sources"][0] + for key in ("last_read_at", "last_checked_at", "status", "pending_batch_count"): + assert first["sources"][0][key] == second["sources"][0][key] assert second["pending_batch_count"] == 2 assert not args["cursor_path"].exists() From b7b9280d639926d4fe7256be26188b486bcb17a3 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Tue, 29 Sep 2026 00:09:06 +0800 Subject: [PATCH 4/5] fix(decision-context): scope scan-budget metadata to capture results Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../capabilities/decision_context/capture.py | 3 ++ .../capabilities/decision_context/profile.py | 1 - ...test_decision_context_capture_isolation.py | 50 +++++++++++++++++++ 3 files changed, 53 insertions(+), 1 deletion(-) diff --git a/loopx/capabilities/decision_context/capture.py b/loopx/capabilities/decision_context/capture.py index 9468636d31..0afa785b53 100644 --- a/loopx/capabilities/decision_context/capture.py +++ b/loopx/capabilities/decision_context/capture.py @@ -267,6 +267,9 @@ def capture_profile_sources( "status": "capture_disabled", "executed": False, } + activation = activation | { + "capture_max_sources_per_tick": profile.capture_max_sources_per_tick, + } if private_file_digest(profile_path) != digest_before: raise ValueError("capture profile changed during activation") if spool_path.resolve() == profile_path.resolve() or ( diff --git a/loopx/capabilities/decision_context/profile.py b/loopx/capabilities/decision_context/profile.py index 482edb3caa..27da5ce34d 100644 --- a/loopx/capabilities/decision_context/profile.py +++ b/loopx/capabilities/decision_context/profile.py @@ -489,7 +489,6 @@ def resolve_decision_context_activation( "automatic_capture": config.automatic_capture, "capture_source_count": len(config.capture_source_ids), "capture_interval_seconds": config.capture_interval_seconds, - "capture_max_sources_per_tick": config.capture_max_sources_per_tick, } if not config.enabled: return status | {"status": "disabled"}, config diff --git a/tests/capabilities/test_decision_context_capture_isolation.py b/tests/capabilities/test_decision_context_capture_isolation.py index 3f47dffe85..eb2217fa53 100644 --- a/tests/capabilities/test_decision_context_capture_isolation.py +++ b/tests/capabilities/test_decision_context_capture_isolation.py @@ -9,6 +9,9 @@ from loopx.capabilities.decision_context.providers import ( LocalFileDecisionSourceProvider, ) +from loopx.capabilities.decision_context.profile import ( + resolve_decision_context_activation, +) from test_decision_context_capture import setup as capture_setup, settle_batch @@ -191,6 +194,53 @@ def scan(self, **kwargs): assert second["scan_budget"]["deferred_source_ids"] == [calls[0]] +@pytest.mark.parametrize( + "disabled_scope", ["profile", "capture", "unlisted-agent", "new-agent"] +) +def test_capture_budget_metadata_stays_inside_enabled_scope(pair, disabled_scope): + args, payload, _, _ = pair + if disabled_scope == "profile": + payload["enabled"] = False + elif disabled_scope == "capture": + payload["automation"]["automatic_capture"] = False + else: + args = {**args, "agent_id": disabled_scope} + args["profile_path"].write_text(json.dumps(payload)) + activation, _ = resolve_decision_context_activation( + goal_id=args["goal_id"], + agent_id=args["agent_id"], + profile_path=args["profile_path"], + ) + assert "capture_max_sources_per_tick" not in activation + for execute in (False, True): + assert capture_profile_sources(**args, execute=execute) == { + "activation": activation, + "status": "capture_disabled", + "executed": False, + } + assert not args["spool_path"].exists() + assert not args["cursor_path"].exists() + + +def test_enabled_capture_projects_budget_without_changing_shared_activation(pair): + args, payload, _, _ = pair + payload["automation"]["max_sources_per_tick"] = 1 + args["profile_path"].write_text(json.dumps(payload)) + activation, _ = resolve_decision_context_activation( + goal_id=args["goal_id"], + agent_id=args["agent_id"], + profile_path=args["profile_path"], + ) + assert activation["available"] is True + assert "capture_max_sources_per_tick" not in activation + result = capture_profile_sources(**args, execute=True) + assert result["activation"]["capture_max_sources_per_tick"] == 1 + assert result["scan_budget"]["attempted_source_count"] == 1 + before = args["spool_path"].read_bytes() + assert capture_profile_sources(**args)["activation"] == result["activation"] + assert args["spool_path"].read_bytes() == before + + @pytest.mark.parametrize("invalid", [0, -1, 65, True, "8", 1.5]) def test_capture_tick_budget_rejects_invalid_configuration(pair, invalid): from loopx.capabilities.decision_context.profile import ( From 769993a1b199045f5bf48bb98ce0496833d5406e Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Tue, 29 Sep 2026 00:09:06 +0800 Subject: [PATCH 5/5] docs(decision-context): clarify scoped budget readback and rollback Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/capabilities/decision_context/README.md | 11 ++++++++--- loopx/capabilities/decision_context/README.zh-CN.md | 10 +++++++--- 2 files changed, 15 insertions(+), 6 deletions(-) diff --git a/loopx/capabilities/decision_context/README.md b/loopx/capabilities/decision_context/README.md index 5f5e216563..57dc2a12c6 100644 --- a/loopx/capabilities/decision_context/README.md +++ b/loopx/capabilities/decision_context/README.md @@ -306,7 +306,10 @@ Add these fields to an existing private profile's `automation` object: Every listed source must already be enabled, incremental, and exact-readable. On-demand sources are never enrolled implicitly. The same goal/agent activation -checks apply. Preview does not call providers or create the spool: +checks apply. Capture results include `activation.capture_max_sources_per_tick` +only after profile, agent and automatic-capture activation. Shared activation, +`inspect-profile` and evidence assembly do not project this capture budget. +Preview does not call providers or create the spool: ```bash loopx decision-context capture --goal-id --agent-id \ @@ -382,8 +385,10 @@ replay is impossible. Capture health is not proof of complete decision coverage. To stop collection, set `automatic_capture=false` and unload the host scheduler. Existing reviewable batches remain private and can still be prepared. To roll -back to an older release, also remove the three new automation fields; retain -the spool as a private checkpoint rather than deleting unreviewed work. +back bounded ticks, remove `max_sources_per_tick`. A release that predates +reference capture also requires removing `source_ids`, `interval_seconds` and +`max_pending_batches`. Retain the spool as a private checkpoint rather than +deleting unreviewed work. #### Recover an unreplayable source without discarding history diff --git a/loopx/capabilities/decision_context/README.zh-CN.md b/loopx/capabilities/decision_context/README.zh-CN.md index 1a8855dd16..f16a23e731 100644 --- a/loopx/capabilities/decision_context/README.zh-CN.md +++ b/loopx/capabilities/decision_context/README.zh-CN.md @@ -274,7 +274,10 @@ cursor。 ``` 白名单只能包含已启用、支持 exact read 的 incremental source,不会隐式纳入 -on-demand 来源;goal/agent 的启用边界不变。先预览,再添加 `--execute` 执行一次: +on-demand 来源;goal/agent 的启用边界不变。只有 profile、当前 agent 和自动采集 +均已启用,capture 返回才包含 `activation.capture_max_sources_per_tick`。 +共用 activation、`inspect-profile` 和 evidence assembly 不投影此采集预算。 +先预览,再添加 `--execute` 执行一次: ```bash loopx decision-context capture --goal-id --agent-id \ @@ -329,8 +332,9 @@ loopx decision-context prepare-captured --goal-id --agent-id