diff --git a/loopx/capabilities/decision_context/README.md b/loopx/capabilities/decision_context/README.md index dc4f4490ba..57dc2a12c6 100644 --- a/loopx/capabilities/decision_context/README.md +++ b/loopx/capabilities/decision_context/README.md @@ -299,13 +299,17 @@ 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 } ``` 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 \ @@ -313,7 +317,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 @@ -325,7 +339,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. @@ -364,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 @@ -416,9 +439,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..f16a23e731 100644 --- a/loopx/capabilities/decision_context/README.zh-CN.md +++ b/loopx/capabilities/decision_context/README.zh-CN.md @@ -268,12 +268,16 @@ 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 } ``` 白名单只能包含已启用、支持 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 \ @@ -281,6 +285,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` 中的 @@ -288,7 +298,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` 开始回读: @@ -318,8 +332,9 @@ loopx decision-context prepare-captured --goal-id --agent-id 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 ] @@ -241,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 ( @@ -269,7 +298,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,7 +361,26 @@ 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) - for source in sources: + source_capacity = _source_capacity( + profile.capture_max_pending_batches, len(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'" @@ -377,17 +427,30 @@ 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: + 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, @@ -460,7 +523,14 @@ def _execute_capture_tick( result = { "activation": activation, "executed": True, - **_status(db, sources, now=now), + "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 + ), } db.commit() return result diff --git a/loopx/capabilities/decision_context/profile.py b/loopx/capabilities/decision_context/profile.py index 3e5fcc12d3..27da5ce34d 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, ) 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..eb2217fa53 --- /dev/null +++ b/tests/capabilities/test_decision_context_capture_isolation.py @@ -0,0 +1,253 @@ +"""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 loopx.capabilities.decision_context.profile import ( + resolve_decision_context_activation, +) +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 + + +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" + 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() + + +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( + "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 ( + 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)