diff --git a/docs/architecture/rfcs/goal-scoped-capability-portfolio-v0.md b/docs/architecture/rfcs/goal-scoped-capability-portfolio-v0.md index 16f4ced27e..7211cf7129 100644 --- a/docs/architecture/rfcs/goal-scoped-capability-portfolio-v0.md +++ b/docs/architecture/rfcs/goal-scoped-capability-portfolio-v0.md @@ -534,6 +534,17 @@ an engineering and a research question; no host-session-file copying. The current continuation and Decision Context contracts own the work; missing provenance/update fields belong to their owners, not a new continuity store. +The optional [query-ready Reward Memory caller](../../reference/reward-memory-decision-consumption.md) +provides a bounded prerequisite: TS admission/completion, existing Python +provider/applier adaptation, context delivery distinct from semantic assessment, +and exact caller-retained replay. It does not complete the fresh-session journey, +automatically compose capabilities, or establish memory utility; those remain +subject to real caller and held-out outcome qualification. + +可选的 query-ready Reward Memory 调用方提供有界前置切片:TS 准入/完成、原 +Python provider/applier 适配、上下文交付与语义判断分开、调用方保留的精确复用。 +它不宣称已完成跨会话纠正、自动组合或记忆效果;真实调用和留出结果仍须验收。 + ## 6. Alternatives and design choices ### Domain skills organize capabilities diff --git a/docs/reference/reward-memory-decision-consumption.md b/docs/reference/reward-memory-decision-consumption.md new file mode 100644 index 0000000000..9bd24fc16e --- /dev/null +++ b/docs/reference/reward-memory-decision-consumption.md @@ -0,0 +1,158 @@ +# Query-ready Reward Memory consumption / 决策就绪的经验消费 + +This optional caller API closes **qualified recall → private context delivery → +explicit semantic assessment**. It does not generate a query, enable a provider, +write memory, or prove utility. Existing experiment configuration, corpus scope, +freshness/conflict checks and provider adapters remain their owners. + +此可选 API 连接「合格召回 → 私有上下文交付 → 显式语义判断」。它不生成问题、 +开启 provider、写入记忆或证明效果;配置、corpus 范围、时效/冲突和 provider +仍由既有 owner 持有。应先冻结具体问题和当前产物,再调用,不在每次心跳重复召回。 + +## Entry and ownership / 入口与归属 + +Import `run_reward_memory_decision` and `assess_reward_memory_decision` from +`loopx.capabilities.reward_memory`. Resolve the original normalized configuration +with `resolve_reward_memory_experiment`; an unavailable binding must not be +replaced with an unverified configuration. Pass the existing automatic-recall +hook arguments unchanged: surface, workspace/project, revision, bounded queries, +observation time, freshness/conflict, read-authority checkpoints and provider. +Add `query_ready`, the consumer mode and an explicit application strategy. +`application_id` and `artifact_ref` identify this question and current artifact, +not a fresh identity on every retry. The caller's baseline/arguments must be +JSON-compatible; non-serializable input fails before recall. + +Read-authority checkpoints must match the exact consumer surface and corpus; +a turn-admission checkpoint cannot authorize a different review surface. +`freshness_context.age_seconds`, when supplied, is a nonnegative integer. +Rejected requests expose only the original hook's allowlisted +`boundary_reason_code`, never exception text or private input values. + +使用上述导出入口和 `resolve_reward_memory_experiment` 的原配置读回;不可用时 +不能拿未经验证的配置替代。原 hook 的范围、revision、问题、时点、时效/冲突、 +读授权 checkpoint 和 provider 原样传入,另加 `query_ready`、模式和应用策略。 +`application_id` / `artifact_ref` 绑定同一问题和当前产物,不因重试换身份。 +基线和参数须可 JSON 序列化,非法输入在 provider 调用前 fail-open。 + +读授权 checkpoint 必须匹配本次 surface/corpus,不能拿 Turn 准入的 checkpoint +授权另一评审入口;age_seconds 如提供,须为非负整数。拒绝回执仅投影原 hook +白名单内的 boundary_reason_code,不暴露异常正文或私有参数。 + +TypeScript owns admission and completion (`reward_memory.decision.plan/project`); +Python adapts the existing provider/applier and retains transient private values. +TS receives only compact references, typed statuses, counts and receipt digests: +never the query, lesson, baseline, provider payload or model rationale. +No new store, SDK, key, enablement switch or action authority is introduced. + +TS 负责准入和完成语义;Python 只适配现有 provider/applier 并保留瞬时私有值。 +TS 只收到引用、状态、计数与摘要,不收到问题、经验正文、原产物或模型判断内容; +不新增存储、SDK、密钥、开关或行动授权。 + +| Mode / 模式 | Provider / 调用 | Meaning / 意义 | +| --- | --- | --- | +| Disabled/unconfigured / 未配置或关闭 | Zero; returns `None` / 零调用,无新 packet | Original path unchanged / 原路径不变 | +| `preview` | Zero / 零调用 | Admission preview, never adoption / 预检,不代表采用 | +| `recall_only` | Existing bounded route / 原有界路由 | Redacted recall observation, base unchanged; no semantic claim / 脱敏观察,不表示语义采用 | +| `execute` without strategy/artifact / 缺策略或产物 | Zero / 零调用 | `incomplete`, research may continue / 不完整,研究可继续 | +| `execute` + `context_delivery` | Existing bounded route / 原有界路由 | Private context available, **not** semantic application / 私有上下文可读,非语义采用 | +| `execute` + `semantic_application` | Existing bounded route / 原有界路由 | Exact attributed `applied/ignored/refuted` / 当前产物的归因判断 | + +## Minimal caller integration / 最小接入 + +`hook_arguments` below are the already-qualified arguments from the caller's +existing surface. The callback is the original SDK applier contract, not a +provider response interpreted as authority: + +下例 `hook_arguments` 来自既有已验证 surface;callback 复用原 SDK 的 applier +契约,不把 provider 内容解释成授权。 + +```python +from loopx.capabilities.reward_memory import ( + run_reward_memory_decision, assess_reward_memory_decision, +) + +def deliver_context(base, items): + return { + "outcome": "applied", # SDK delivery; NOT semantic-use evidence + "output": {"baseline": base, "private_context": items}, + "memory_refs": [item.memory_ref for item in items], + "current_artifact_verified": True, # caller must actually verify it + "reasoning_summary": "Qualified context delivered for separate review.", + } + +delivered = run_reward_memory_decision( + config, query_ready=True, mode="execute", + application_kind="context_delivery", apply_memory=deliver_context, + **hook_arguments, +) +if delivered is not None and delivered.public_packet["context_delivery_verified"]: + # The real caller/model reads delivered.output and the CURRENT artifact. + # judge returns output, outcome, current_artifact_verified, memory_refs + # from these exact items, and a bounded evidence-backed reasoning_summary. + assessed = assess_reward_memory_decision(delivered, apply_memory=judge) + public_receipt = assessed.public_packet +``` + +`judge` must represent actual comparison, not unconditional adoption. For +`ignored/refuted`, keep the original baseline. All semantic dispositions require +nonempty attribution to these exact recalled items and verified current artifact. +An applied lesson may preserve the decision; application is still not utility. + +`judge` 必须表达真实比较,不能无条件采用;ignored/refuted 保持原基线。 +所有语义判断都需引用本次实际条目并核验当前产物。经验被采用也可能不改变结论, +更不代表质量、成本、收益或 alpha 已改善。 + +Only `public_packet` is a display projection. Output, session, original +application receipt, baseline and request remain caller-private, never generic +frontend/Lark/registry payloads. Retain the result in the caller's existing +execution context. `previous_result=result` replays only an exact request digest +(configuration and input included) without another provider call; changed input +returns `replay_request_mismatch`. Reassessment uses retained qualified items and +the original **cumulative** multi-corpus counters, not a second query. This is +caller-retained replay, not automatic cross-process persistence or a new cache. + +仅 `public_packet` 用于展示,其余结果私有。通过既有执行上下文保留结果; +`previous_result` 仅复用配置和输入均匹配的请求,变化则拒绝复用。后续判断使用 +原条目和累计多 corpus 遥测,不重复查询。这不是自动跨进程存储或新的缓存。 + +Empty/filtered/unavailable and invalid model/transport results preserve the base +and allow ordinary research. Post-provider transport failure retains actual +call/filter counts and the original private receipt. The existing route still +owns its call cap: one query per corpus can mean multiple provider calls. +A caller with a one-call budget must use a qualified one-corpus route. + +空结果、过滤、故障和无效判断不阻塞独立研究。provider 后 TS 故障仍保留真实 +计数和原私有回执。预算由原路由负责;每 corpus 一次不等于全路由一次,单次预算 +须使用已验证的单 corpus 路由。 + +## Validation, surfaces and rollback / 验证、产品入口与回滚 + +```sh +uv run --extra test pytest -q tests/capabilities/test_reward_memory_decision.py +node --experimental-strip-types --test tests/control_plane_ts/reward_memory_decision.test.ts +``` + +This delivery adds a caller API for CLI/managed decision boundaries. Configuration +is unchanged: Dashboard's existing Reward Memory editor still controls +`config_path` and `enabled_agents` through the same preview/apply/readback owner; +Lark keeps its existing status-only coverage. No new frontend control or packaged +frontend change is needed for these unchanged settings. The new typed receipt +can be returned by an integrated caller; **universal frontend/Lark consumption +and cross-session persistence are not delivered by this API**. + +本次交付是 CLI/managed 决策边界的调用 API。配置未变,Dashboard 继续通过原 +preview/apply/readback 编辑 config_path 和 enabled_agents;Lark 仍为原状态 +投影。无需新增配置控件或前端包;接入的调用方可回传同一 typed receipt, +但此 API 不宣称已打通所有前端/Lark 或自动跨 session 持久化。 + +Rollback the caller to `run_reward_memory_automatic_recall_hook`, whose optional +callback/`available_not_applied` behavior is unchanged. To disable recall, use +the original configuration owner to set `automatic_recall=false` and requalify +the binding, or disable the experiment with `configure-goal --clear-reward-memory-config`. +No recall, disposition or receipt grants orders, signing, transfer, publishing +or legacy-record migration rights. + +回滚调用方至原 automatic hook 即可;其可选 callback 语义未变。关闭召回通过 +原配置 owner 设置 automatic_recall=false 后重新验证绑定,或用原 +configure-goal --clear-reward-memory-config 关闭实验。任何回执均不扩大交易、 +签名、转账、发布或旧记录迁移权限。 diff --git a/loopx/capabilities/reward_memory/README.md b/loopx/capabilities/reward_memory/README.md index 7bdac66dec..dfa0ebd0e1 100644 --- a/loopx/capabilities/reward_memory/README.md +++ b/loopx/capabilities/reward_memory/README.md @@ -650,6 +650,11 @@ performs no provider or external write. ## Stage 3 recall and application seam +For an explicit query-ready decision, use the optional +[decision-consumption caller API](../../../docs/reference/reward-memory-decision-consumption.md). +It distinguishes private context delivery from attributed semantic assessment, +reuses the existing provider/applier, and retains the original optional-callback SDK path. + Stage 3 accepts only an explicit `reward_memory_recall_request_v0` naming one registered corpus and one module-owned surface. The request carries a matching read-authority checkpoint and current freshness/conflict observations. A diff --git a/loopx/capabilities/reward_memory/__init__.py b/loopx/capabilities/reward_memory/__init__.py index ef5ee5179c..dd34eb42a5 100644 --- a/loopx/capabilities/reward_memory/__init__.py +++ b/loopx/capabilities/reward_memory/__init__.py @@ -36,6 +36,11 @@ normalize_reward_memory_standing_policy, ) from .evaluation import run_reward_memory_evaluation +from .decision import ( + RewardMemoryDecisionResult, + assess_reward_memory_decision, + run_reward_memory_decision, +) from .dogfood import ( build_reward_memory_dogfood_batch, build_reward_memory_dogfood_receipt, @@ -65,6 +70,9 @@ __all__ = [ "build_reward_memory_architecture_packet", + "RewardMemoryDecisionResult", + "assess_reward_memory_decision", + "run_reward_memory_decision", "RewardMemoryFilteredRecallItem", "RewardMemoryRecallItem", "RewardMemoryRecallSession", diff --git a/loopx/capabilities/reward_memory/decision.py b/loopx/capabilities/reward_memory/decision.py new file mode 100644 index 0000000000..b176b04fdc --- /dev/null +++ b/loopx/capabilities/reward_memory/decision.py @@ -0,0 +1,192 @@ +"""Query-ready caller boundary over the existing config, recall and applier owners. + +Python retains the provider/model adapters and transient private values. The +TypeScript owner decides admission and completion; no new store or opt-in. +""" +from __future__ import annotations + +import hashlib +import json +from collections.abc import Mapping +from copy import deepcopy +from dataclasses import dataclass, replace +from typing import Any + +from ...control_plane.effect_runtime import effect_runtime_result +from .application import ( + RewardMemoryApplier, RewardMemoryRecallItem, RewardMemoryRecallSession, + apply_reward_memory_recall, +) +from .runtime_hooks import run_reward_memory_automatic_recall_hook + + +@dataclass(frozen=True) +class RewardMemoryDecisionResult: + """Only public_packet is a projection. All other fields stay caller-private.""" + + public_packet: dict[str, Any] + output: Any + base_output: Any + request_digest: str + request: dict[str, Any] + recall_session: RewardMemoryRecallSession | None = None + application_receipt: Mapping[str, Any] | None = None + recall_telemetry: Mapping[str, Any] | None = None + + +def _transport_failure( + request: Mapping[str, Any], reason: str, telemetry: Mapping[str, Any] | None = None, +) -> dict[str, Any]: + # Transport cannot invent a successful TS decision when the kernel is unavailable. + return { + "schema_version": "reward_memory_decision_consumption_v0", + "status": "incomplete", "reason_code": reason, + **{key: request.get(key) for key in ("mode", "application_kind", "application_id", "artifact_ref", "surface_id")}, + "should_recall": False, "decision_consumption_complete": False, + "context_delivery_verified": False, "semantic_disposition": None, + "preserve_base_output": True, "research_may_continue": True, + "grants_new_action_authority": False, "utility_verified": False, + "external_writes_performed": False, "raw_content_captured": False, + **dict(telemetry or {"provider_call_count": 0, "filtered_count": 0, + "result_readback_verified": False}), + } + + +def _recall_telemetry(hook: Mapping[str, Any]) -> dict[str, Any]: + telemetry = hook.get("telemetry") or {} + attempts = hook.get("recall_attempts") or [] + return { + "result_readback_verified": telemetry.get("result_readback_verified", False), + "provider_call_count": telemetry.get("provider_call_count", 0), + "filtered_count": sum(item.get("filtered_item_count", 0) for item in attempts), + "recall_status": attempts[-1].get("status") if attempts else None, + "boundary_reason_code": hook.get("reason_code"), + } + + +def _project( + request: dict[str, Any], status: str, telemetry: Mapping[str, Any], + receipt: Mapping[str, Any] | None, +) -> dict[str, Any]: + return effect_runtime_result("reward_memory.decision.project", { + "request": request, + "observation": { + "status": status, **dict(telemetry), + # No query, summary, lesson, model rationale or base artifact crosses into TS. + "application_receipt": { + key: value for key, value in (receipt or {}).items() + if key in {"schema_version", "application_id", "artifact_ref", "surface_id", + "outcome", "memory_ref_digests", "current_artifact_verified", + "result_readback_verified"} + }, + }, + }) + + +def run_reward_memory_decision( + config: Mapping[str, Any] | None, + *, + query_ready: bool, + mode: str = "execute", + application_kind: str | None = None, + apply_memory: RewardMemoryApplier | None = None, + previous_result: RewardMemoryDecisionResult | None = None, + **hook_arguments: Any, +) -> RewardMemoryDecisionResult | None: + """Consume an explicit decision question using the existing automatic hook. + +Pass its existing keyword arguments unchanged. No configured automatic recall +means None (no added packet, TS call or provider call). Execute requires an +explicit context_delivery or semantic_application callback. Retain the result +privately to replay the exact request or assess delivered context without recall. +""" + automation = config.get("automation") if isinstance(config, Mapping) else None + if not isinstance(automation, Mapping) or automation.get("automatic_recall") is not True: + return None + base = hook_arguments.get("base_output") + request = {"mode": mode, "query_ready": query_ready, "application_kind": application_kind, + "has_applier": callable(apply_memory), + **{key: hook_arguments.get(key) for key in ("application_id", "artifact_ref", "surface_id")}} + digest = "" + session = None + receipt = None + telemetry = None + try: + digest = hashlib.sha256(json.dumps({ + "config": config, "request": request, + "context": {key: value for key, value in hook_arguments.items() if key != "provider"}, + }, sort_keys=True, ensure_ascii=False, allow_nan=False).encode()).hexdigest() + if previous_result is not None: + if previous_result.request_digest != digest: + return RewardMemoryDecisionResult( + _transport_failure(request, "replay_request_mismatch"), base, base, digest, request, + ) + return previous_result + plan = effect_runtime_result("reward_memory.decision.plan", request) + if not plan["should_recall"]: + return RewardMemoryDecisionResult(plan, base, base, digest, request) + captured: tuple[RewardMemoryRecallItem, ...] = () + + def apply(original: Any, items: tuple[RewardMemoryRecallItem, ...]) -> Mapping[str, Any]: + nonlocal captured + captured = items + assert apply_memory is not None + return apply_memory(original, items) + + # Fail-open preserves the original baseline even if a model mutates its input. + arguments = {**hook_arguments, "base_output": deepcopy(base)} + hook = run_reward_memory_automatic_recall_hook( + config, **arguments, + apply_memory=apply if mode == "execute" else None, + ) + application = hook.get("application") or {} + attempts = hook.get("recall_attempts") or [] + session = RewardMemoryRecallSession(attempts[-1], captured) if captured and attempts else None + receipt = application.get("receipt") + telemetry = _recall_telemetry(hook) + packet = _project(request, hook["status"], telemetry, receipt) + return RewardMemoryDecisionResult( + packet, base if packet["preserve_base_output"] else hook["output"], + base, digest, request, session, receipt, telemetry, + ) + except (KeyError, OSError, RuntimeError, TypeError, ValueError): + return RewardMemoryDecisionResult( + _transport_failure(request, "consumer_input_or_runtime_failed", telemetry), + base, base, digest, request, session, receipt, telemetry, + ) + + +def assess_reward_memory_decision( + delivered: RewardMemoryDecisionResult, + *, + apply_memory: RewardMemoryApplier, +) -> RewardMemoryDecisionResult: + """Actual caller reasoning over exact retained items; never queries a provider. + + The callback returns the existing SDK's output, applied/ignored/refuted, + current_artifact_verified, memory_refs and reasoning_summary fields. A + previously assessed result is already the receipt, not another model call. + """ + if delivered.public_packet.get("decision_consumption_complete") is True: + return delivered + request = {**delivered.request, "application_kind": "semantic_application", + "has_applier": callable(apply_memory)} + receipt = delivered.application_receipt + try: + if delivered.recall_session is None or delivered.public_packet.get("context_delivery_verified") is not True: + return replace(delivered, public_packet=_transport_failure(request, "verified_context_delivery_required", delivered.recall_telemetry), output=delivered.base_output) + application = apply_reward_memory_recall( + deepcopy(delivered.base_output), delivered.recall_session, + application_id=request["application_id"], artifact_ref=request["artifact_ref"], + apply_memory=apply_memory, + ) + receipt = application["receipt"] + # Reassessment is not recall: retain every corpus's original cumulative counters. + packet = _project(request, application["status"], delivered.recall_telemetry or {}, + application["receipt"]) + return replace(delivered, public_packet=packet, request=request, + output=delivered.base_output if packet["preserve_base_output"] else application["output"], + application_receipt=application["receipt"]) + except (KeyError, OSError, RuntimeError, TypeError, ValueError): + return replace(delivered, public_packet=_transport_failure(request, "consumer_input_or_runtime_failed", delivered.recall_telemetry), output=delivered.base_output, + request=request, application_receipt=receipt) diff --git a/loopx/control_plane/capabilities/reward_memory_decision.ts b/loopx/control_plane/capabilities/reward_memory_decision.ts new file mode 100644 index 0000000000..244df8e24b --- /dev/null +++ b/loopx/control_plane/capabilities/reward_memory_decision.ts @@ -0,0 +1,117 @@ +import type { JsonObject } from "../effect_program.ts"; +import { EffectRuntimeRequestError } from "../effect_runtime_errors.ts"; +import { requireJsonObject, requireStringLiteral } from "../runtime_decode.ts"; + +const MODES = ["execute", "preview", "recall_only"] as const; +const KINDS = ["context_delivery", "semantic_application"] as const; +const TOKEN = /^[A-Za-z0-9][A-Za-z0-9._:/#-]{0,199}$/; + +function token(value: unknown, name: string, optional = false): string | null { + if (optional && (value === null || value === undefined)) return null; + if (typeof value !== "string" || !TOKEN.test(value)) { + throw new EffectRuntimeRequestError(`${name} must be a compact reference`); + } + return value; +} + +function boolean(value: unknown, name: string): boolean { + if (typeof value !== "boolean") throw new EffectRuntimeRequestError(`${name} must be boolean`); + return value; +} + +function count(value: unknown, name: string): number { + if (!Number.isSafeInteger(value) || (value as number) < 0) { + throw new EffectRuntimeRequestError(`${name} must be a nonnegative integer`); + } + return value as number; +} + +/** Query-ready consumption policy; no config, provider content or model calls. */ +export function planRewardMemoryDecision(params: JsonObject): JsonObject { + const mode = requireStringLiteral(params.mode, MODES, "mode"); + const kind = params.application_kind === null || params.application_kind === undefined + ? null : requireStringLiteral(params.application_kind, KINDS, "application_kind"); + const ready = boolean(params.query_ready, "query_ready"); + const hasApplier = boolean(params.has_applier, "has_applier"); + const packet: JsonObject = { + schema_version: "reward_memory_decision_consumption_v0", + mode, application_kind: kind, + application_id: token(params.application_id, "application_id"), + artifact_ref: token(params.artifact_ref, "artifact_ref", true), + surface_id: token(params.surface_id, "surface_id"), + status: "incomplete", reason_code: null, should_recall: false, + decision_consumption_complete: false, context_delivery_verified: false, + semantic_disposition: null, result_readback_verified: false, + provider_call_count: 0, filtered_count: 0, preserve_base_output: true, + research_may_continue: true, grants_new_action_authority: false, + external_writes_performed: false, raw_content_captured: false, + utility_verified: false, + }; + if (!ready) return {...packet, reason_code: "query_not_ready"}; + if (mode === "preview") return {...packet, status: "preview"}; + if (mode === "execute") { + if (!kind || !hasApplier) return {...packet, reason_code: "application_strategy_required"}; + if (!packet.artifact_ref) return {...packet, reason_code: "current_artifact_binding_required"}; + } + return {...packet, status: "ready", should_recall: true}; +} + +/** Reduce original recall/application receipts, never interpret private lessons. */ +export function projectRewardMemoryDecision(params: JsonObject): JsonObject { + const plan = planRewardMemoryDecision(requireJsonObject(params.request, "request")); + if (!plan.should_recall) return plan; + const observation = requireJsonObject(params.observation, "observation"); + const hookStatus = requireStringLiteral(observation.status, [ + "disabled", "guard_rejected", "provider_unavailable", "not_available", + "available_not_applied", "failed", "applied", "ignored", "refuted", + ] as const, "observation.status"); + const readback = boolean(observation.result_readback_verified, "result_readback_verified"); + const recallStatus = observation.recall_status === null ? null + : requireStringLiteral(observation.recall_status, ["completed", "empty", "provider_unavailable", "guard_blocked"] as const, "recall_status"); + const packet: JsonObject = { + ...plan, should_recall: false, + boundary_reason_code: observation.boundary_reason_code == null ? null + : requireStringLiteral(observation.boundary_reason_code, [ + "automation_config_invalid", "surface_profile_or_query_invalid", + "exact_corpus_request_invalid", "surface_has_no_recall_corpus", + ] as const, "boundary_reason_code"), + provider_call_count: count(observation.provider_call_count, "provider_call_count"), + filtered_count: count(observation.filtered_count, "filtered_count"), + result_readback_verified: readback, + recall_status: recallStatus, + }; + if (hookStatus === "provider_unavailable") { + return {...packet, status: "provider_unavailable", reason_code: "provider_unavailable"}; + } + if (hookStatus === "guard_rejected" || hookStatus === "disabled" || recallStatus === "guard_blocked") { + return {...packet, status: "incomplete", reason_code: "recall_boundary_rejected"}; + } + if (!readback) return {...packet, status: "empty", reason_code: packet.filtered_count + ? "all_provider_items_filtered" : "provider_returned_no_items"}; + if (plan.mode === "recall_only") return {...packet, status: "recalled"}; + const receipt = requireJsonObject(observation.application_receipt, "application_receipt"); + const digests = receipt.memory_ref_digests; + const attributed = Array.isArray(digests) && digests.length > 0 && digests.length <= 8 && + digests.every((item) => typeof item === "string" && /^[0-9a-f]{16}$/.test(item)); + const bound = receipt.schema_version === "reward_memory_application_receipt_v0" && + receipt.application_id === plan.application_id && receipt.artifact_ref === plan.artifact_ref && + receipt.surface_id === plan.surface_id && receipt.outcome === hookStatus && + receipt.current_artifact_verified === true && receipt.result_readback_verified === true && attributed; + if (!bound || hookStatus === "failed" || hookStatus === "available_not_applied") { + return {...packet, status: "incomplete", reason_code: "application_evidence_incomplete"}; + } + // A delivered context is available for reasoning; it is not the reasoning disposition. + if (plan.application_kind === "context_delivery") { + return hookStatus === "applied" + ? {...packet, status: "context_delivered", memory_ref_digests: digests, + context_delivery_verified: true, preserve_base_output: false} + : {...packet, status: "incomplete", reason_code: "context_delivery_not_verified"}; + } + if (hookStatus !== "applied" && hookStatus !== "ignored" && hookStatus !== "refuted") { + return {...packet, status: "incomplete", reason_code: "semantic_disposition_required"}; + } + return { + ...packet, status: hookStatus, semantic_disposition: hookStatus, memory_ref_digests: digests, + decision_consumption_complete: true, preserve_base_output: hookStatus !== "applied", + }; +} diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index a774751aff..90b4049cd6 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -231,6 +231,10 @@ import { projectExternalEvidenceRetirement, recordExternalEvidenceReceiptObservation, } from "./capabilities/external_evidence.ts"; +import { + planRewardMemoryDecision, + projectRewardMemoryDecision, +} from "./capabilities/reward_memory_decision.ts"; type EffectRuntimeHandler = (params: JsonObject) => unknown | Promise; @@ -726,6 +730,8 @@ export function createEffectRuntimeHandlers( ["external_evidence.receipt", recordExternalEvidenceReceiptObservation], ["external_evidence.admit", evaluateExternalEvidenceAdmission], ["external_evidence.retire", projectExternalEvidenceRetirement], + ["reward_memory.decision.plan", planRewardMemoryDecision], + ["reward_memory.decision.project", projectRewardMemoryDecision], [ "manager.return_delivery.normalize_attempt", (params) => normalizeManagerReturnDeliveryAttempt(params.attempt), diff --git a/tests/capabilities/test_reward_memory_decision.py b/tests/capabilities/test_reward_memory_decision.py new file mode 100644 index 0000000000..8f9d9d32b0 --- /dev/null +++ b/tests/capabilities/test_reward_memory_decision.py @@ -0,0 +1,269 @@ +"""Production TS transport + original configured recall/application boundary.""" +from __future__ import annotations + +import copy +import hashlib +import json +from pathlib import Path +from typing import Any + +import pytest + +from loopx.capabilities.context_providers.base import ContextProviderItem, ContextProviderRetrieval +from loopx.capabilities.reward_memory import assess_reward_memory_decision, run_reward_memory_decision +from loopx.capabilities.reward_memory import decision +from loopx.capabilities.reward_memory.experiment import load_reward_memory_experiment_config +from loopx.capabilities.reward_memory.runtime_hooks import run_reward_memory_automatic_recall_hook + +SURFACE = "reviewer_artifact.summary" +REVISION = "revision:abc123" +FIXTURE = Path(__file__).resolve().parents[2] / "examples/fixtures/reward-memory-scoped-feedback-ingest.public.json" + + +class Provider: + provider_id = "openviking" + + def __init__(self, records: dict[str, Any] | None = None, *, unavailable: bool = False): + self.records = records or {} + self.unavailable = unavailable + self.calls = 0 + + def retrieve(self, **kwargs): + self.calls += 1 + if self.unavailable: + raise RuntimeError("provider unavailable") + scope = kwargs["scope_ref"] + record = self.records.get(scope) + return ContextProviderRetrieval( + provider=self.provider_id, namespace=kwargs["namespace"], status="completed", + query_summary=kwargs["query_summary"], observed_at=kwargs["observed_at"], + search_performed=True, read_performed=True, requested_limit=kwargs["max_results"], + items=(ContextProviderItem(resource_ref=f"{scope}/memory.json", summary="reviewed lesson", + content=json.dumps(record), score=0.9),) if record else (), + ) + + +def context(tmp_path: Path, *, two_corpora: bool = False): + fixture = json.loads(FIXTURE.read_text()) + entries, scopes, records, checkpoints = [], [], {}, {} + for name in (["primary", "overlay"] if two_corpora else ["primary"]): + corpus, policy = copy.deepcopy(fixture["corpus"]), copy.deepcopy(fixture["standing_policy"]) + scope = f"viking://resources/reward-memory/{name}" + corpus["corpus_id"] = name + corpus["provider_scope_ref_digest"] = hashlib.sha256(scope.encode()).hexdigest()[:16] + policy["policy_id"] = f"policy:example:{name}" + entries.append({"corpus": corpus, "standing_policy": policy}) + scopes.append({"corpus_id": name, "scope_ref": scope}) + checkpoints[name] = {"verified": True, "corpus_id": name, **{ + key: corpus["scope"][key] for key in ("workspace_ref", "project_ref")}, + "surface_id": SURFACE, "read_authority": corpus["read_authority"], + "source_ref": policy["authority_source_ref"]} + records[scope] = {"schema_version": "reward_memory_active_record_v0", "corpus_id": name, + "candidate_ref": f"candidate:{name}", "target_class": "hard_policy", + "content_summary": "Private reviewed summary language lesson.", + "scope": {**corpus["scope"], "revision_ref": REVISION}, "lifecycle": {"state": "active"}} + raw = {"schema_version": "reward_memory_experiment_config_v1", "corpora": entries, + "project_provider_binding": {"provider_id": "openviking", "namespace": "reward_memory", + "timeout_seconds": 30, "minimum_provider_version": "0.4.9", "corpus_scopes": scopes}, + "surfaces": [{"surface_id": SURFACE, "adapter": "scoped_feedback", + "corpus_ids": [item["corpus"]["corpus_id"] for item in entries], "ingest_corpus_id": "primary", + "recall_profile": {"profile_id": "summary", "mode": "function_boundary", "max_queries": 1, "limit": 4}}], + "automation": {"automatic_recall": True, "automatic_ingest": False, "fail_open": True}} + path = tmp_path / "config.json" + path.write_text(json.dumps(raw)) + config = load_reward_memory_experiment_config(project=tmp_path, config_path=path.name) + arguments = {"surface_id": SURFACE, "base_output": {"summary": "baseline"}, + "workspace_ref": "workspace:example", "project_ref": "repository:example", "revision_ref": REVISION, + "queries": [{"query": "Which reviewed language lesson applies?", "query_summary": "summary language"}], + "observed_at": "2026-07-17T03:00:00+08:00", "freshness_context": { + "source_truth_current": True, "source_revision": REVISION}, "conflict_state": "clear", + "read_authority_checkpoints": checkpoints, "application_id": "decision:summary", + "artifact_ref": "artifact:current:abc123"} + return config, arguments, records + + +def delivery(base, items): + return {"outcome": "applied", "output": {**base, "private_context": [item.content_summary for item in items]}, + "memory_refs": [item.memory_ref for item in items], "current_artifact_verified": True, + "reasoning_summary": "Delivered qualified context; semantic assessment is separate."} + + +def test_disabled_adds_no_packet_transport_or_provider(tmp_path, monkeypatch): + config, arguments, records = context(tmp_path) + config["automation"]["automatic_recall"] = False + before = copy.deepcopy(config) + provider = Provider(records) + monkeypatch.setattr(decision, "effect_runtime_result", lambda *args: pytest.fail("disabled TS call")) + for disabled in (None, {}, config): + assert run_reward_memory_decision(disabled, query_ready=True, provider=provider, **arguments) is None + assert config == before + assert provider.calls == 0 + + +@pytest.mark.parametrize("extra,reason", [ + ({"mode": "preview"}, None), ({"query_ready": False}, "query_not_ready"), + ({}, "application_strategy_required"), + ({"application_kind": "context_delivery"}, "application_strategy_required"), + ({"apply_memory": delivery}, "application_strategy_required"), + ({"application_kind": "context_delivery", "apply_memory": delivery, "artifact_ref": None}, "current_artifact_binding_required"), +]) +def test_pre_provider_admission(tmp_path, extra, reason): + config, arguments, records = context(tmp_path) + provider = Provider(records) + result = run_reward_memory_decision(config, **{**arguments, "query_ready": True, "provider": provider, **extra}) + assert result.public_packet["reason_code"] == reason + assert result.public_packet["status"] == ("preview" if extra.get("mode") == "preview" else "incomplete") + assert result.output == arguments["base_output"] + assert not result.public_packet["decision_consumption_complete"] + assert provider.calls == 0 + + +@pytest.mark.parametrize("case,expected,calls,filtered", [ + ("empty", "empty", 1, 0), ("filtered", "empty", 1, 1), + ("unavailable", "provider_unavailable", 1, 0), ("wrong_scope", "incomplete", 0, 0), +]) +def test_failure_keeps_base_research_and_truthful_counters(tmp_path, case, expected, calls, filtered): + config, arguments, records = context(tmp_path) + if case == "filtered": + next(iter(records.values()))["corpus_id"] = "unrelated" + if case == "wrong_scope": + arguments["workspace_ref"] = "workspace:unrelated" + provider = Provider({} if case == "empty" else records, unavailable=case == "unavailable") + result = run_reward_memory_decision(config, query_ready=True, application_kind="context_delivery", + apply_memory=delivery, provider=provider, **arguments) + assert result.public_packet["status"] == expected + assert result.public_packet["provider_call_count"] == provider.calls == calls + assert result.public_packet["filtered_count"] == filtered + assert result.output == arguments["base_output"] + assert result.public_packet["research_may_continue"] + assert not result.public_packet["decision_consumption_complete"] + + +def test_recall_only_and_old_optional_callback_sdk_remain_legitimate(tmp_path): + config, arguments, records = context(tmp_path) + provider = Provider(records) + result = run_reward_memory_decision(config, query_ready=True, mode="recall_only", provider=provider, **arguments) + assert result.public_packet["status"] == "recalled" + assert result.public_packet["semantic_disposition"] is None + assert not result.public_packet["decision_consumption_complete"] + assert result.output == arguments["base_output"] + original = run_reward_memory_automatic_recall_hook(config, provider=provider, **arguments) + assert original["status"] == "available_not_applied" + assert original["application"]["receipt"]["reasoning_summary"] == "model_application_callback_not_supplied" + + +def test_invalid_original_request_retains_safe_boundary_reason(tmp_path): + config, arguments, records = context(tmp_path) + arguments["freshness_context"]["age_seconds"] = 1.5 + provider = Provider(records) + result = run_reward_memory_decision(config, query_ready=True, application_kind="context_delivery", + apply_memory=delivery, provider=provider, **arguments) + assert result.public_packet["reason_code"] == "recall_boundary_rejected" + assert result.public_packet["boundary_reason_code"] == "exact_corpus_request_invalid" + assert result.public_packet["provider_call_count"] == provider.calls == 0 + assert "freshness_context" not in json.dumps(result.public_packet) + + +@pytest.mark.parametrize("disposition", ["applied", "applied_unchanged", "ignored", "refuted"]) +def test_delivery_then_actual_bound_assessment_and_exact_replay(tmp_path, disposition): + outcome = "applied" if disposition == "applied_unchanged" else disposition + config, arguments, records = context(tmp_path, two_corpora=True) + next(iter(records.values()))["corpus_id"] = "unrelated" + provider = Provider(records) + delivered = run_reward_memory_decision(config, query_ready=True, application_kind="context_delivery", + apply_memory=delivery, provider=provider, **arguments) + assert delivered.public_packet["status"] == "context_delivered" + assert delivered.output["private_context"] == ["Private reviewed summary language lesson."] + assert delivered.public_packet["semantic_disposition"] is None + assert not delivered.public_packet["decision_consumption_complete"] + assert delivered.public_packet["provider_call_count"] == 2 + assert delivered.public_packet["filtered_count"] == 1 + assessments = [] + + def judge(base, items): + assessments.append(items[0].content_summary) + return {"outcome": outcome, "output": {"summary": "reviewed"} if disposition == "applied" else base, + "memory_refs": [item.memory_ref for item in items], "current_artifact_verified": True, + "reasoning_summary": "Compared returned lesson with the current artifact, not injection."} + + assessed = assess_reward_memory_decision(delivered, apply_memory=judge) + assert assessed.public_packet["decision_consumption_complete"] + assert assessed.public_packet["semantic_disposition"] == outcome + assert assessed.public_packet["provider_call_count"] == provider.calls == 2 + assert assessed.public_packet["filtered_count"] == 1 + assert len(assessments) == 1 + assert assess_reward_memory_decision(assessed, apply_memory=judge) is assessed + assert run_reward_memory_decision(config, query_ready=True, application_kind="context_delivery", + apply_memory=delivery, previous_result=assessed, provider=provider, **arguments) is assessed + changed = run_reward_memory_decision(config, query_ready=True, application_kind="context_delivery", + apply_memory=delivery, previous_result=assessed, provider=provider, + **{**arguments, "artifact_ref": "artifact:new"}) + assert changed.public_packet["reason_code"] == "replay_request_mismatch" + assert provider.calls == 2 + packet = json.dumps(assessed.public_packet) + for private in ("Private reviewed", "Which reviewed", "viking://", "reasoning_summary", "has_applier"): + assert private not in packet + for flag in ("utility_verified", "grants_new_action_authority", "external_writes_performed", "raw_content_captured"): + assert assessed.public_packet[flag] is False + + +@pytest.mark.parametrize("bad", ["foreign_ref", "unverified_artifact", "unattributed", "throws"]) +def test_invalid_semantic_evidence_is_not_adoption(tmp_path, bad): + config, arguments, records = context(tmp_path) + provider = Provider(records) + + def invalid(base, items): + if bad == "throws": + base["summary"] = "mutated" + raise RuntimeError("model failure") + return {"outcome": "ignored", "output": base, + "memory_refs": ["foreign"] if bad == "foreign_ref" else [] if bad == "unattributed" else [items[0].memory_ref], + "current_artifact_verified": bad != "unverified_artifact", "reasoning_summary": "Not verified."} + + result = run_reward_memory_decision(config, query_ready=True, application_kind="semantic_application", + apply_memory=invalid, provider=provider, **arguments) + assert result.public_packet["status"] == "incomplete" + assert not result.public_packet["decision_consumption_complete"] + assert result.output == arguments["base_output"] == {"summary": "baseline"} + assert result.public_packet["provider_call_count"] == 1 + + +def test_ts_projection_failure_after_provider_retains_call_evidence(tmp_path, monkeypatch): + config, arguments, records = context(tmp_path) + provider = Provider(records) + transport = decision.effect_runtime_result + + def fail_projection(method, params): + if method == "reward_memory.decision.project": + assert "Which reviewed" not in json.dumps(params) + raise RuntimeError("TS unavailable") + return transport(method, params) + + monkeypatch.setattr(decision, "effect_runtime_result", fail_projection) + result = run_reward_memory_decision(config, query_ready=True, application_kind="context_delivery", + apply_memory=delivery, provider=provider, **arguments) + assert result.public_packet["status"] == "incomplete" + assert result.public_packet["provider_call_count"] == provider.calls == 1 + assert result.application_receipt["outcome"] == "applied" + assert result.output == arguments["base_output"] + assert run_reward_memory_decision(config, query_ready=True, application_kind="context_delivery", + apply_memory=delivery, previous_result=result, provider=provider, **arguments) is result + assert provider.calls == 1 + + +def test_semantic_receipt_survives_ts_failure_without_claiming_completion(tmp_path, monkeypatch): + config, arguments, records = context(tmp_path) + provider = Provider(records) + delivered = run_reward_memory_decision(config, query_ready=True, application_kind="context_delivery", + apply_memory=delivery, provider=provider, **arguments) + def reject_projection(*args): + raise RuntimeError("TS unavailable") + monkeypatch.setattr(decision, "effect_runtime_result", reject_projection) + result = assess_reward_memory_decision(delivered, apply_memory=lambda base, items: { + "outcome": "ignored", "output": base, "memory_refs": [items[0].memory_ref], + "current_artifact_verified": True, "reasoning_summary": "Current evidence does not need this lesson."}) + assert result.public_packet["status"] == "incomplete" + assert not result.public_packet["decision_consumption_complete"] + assert result.application_receipt["outcome"] == "ignored" + assert result.public_packet["provider_call_count"] == provider.calls == 1 + assert result.output == arguments["base_output"] diff --git a/tests/control_plane_ts/reward_memory_decision.test.ts b/tests/control_plane_ts/reward_memory_decision.test.ts new file mode 100644 index 0000000000..d82a012b0c --- /dev/null +++ b/tests/control_plane_ts/reward_memory_decision.test.ts @@ -0,0 +1,66 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import {planRewardMemoryDecision, projectRewardMemoryDecision} from "../../loopx/control_plane/capabilities/reward_memory_decision.ts"; + +const request = {mode: "execute", query_ready: true, application_kind: "semantic_application", + has_applier: true, application_id: "application:one", artifact_ref: "artifact:current", surface_id: "review.summary"}; +const observation = {status: "ignored", recall_status: "completed", provider_call_count: 2, filtered_count: 1, result_readback_verified: true, + application_receipt: {schema_version: "reward_memory_application_receipt_v0", application_id: "application:one", + artifact_ref: "artifact:current", surface_id: "review.summary", outcome: "ignored", + current_artifact_verified: true, result_readback_verified: true, memory_ref_digests: ["a".repeat(16)]}}; + +test("execute requires explicit strategy and current artifact before recall", () => { + assert.equal(planRewardMemoryDecision({...request, has_applier: false}).should_recall, false); + assert.equal(planRewardMemoryDecision({...request, application_kind: null}).reason_code, "application_strategy_required"); + assert.equal(planRewardMemoryDecision({...request, artifact_ref: null}).reason_code, "current_artifact_binding_required"); + assert.equal(planRewardMemoryDecision({...request, mode: "preview"}).should_recall, false); + assert.equal(planRewardMemoryDecision({...request, query_ready: false}).should_recall, false); +}); + +test("context delivery, recall-only and semantic judgment are different receipts", () => { + const delivered = projectRewardMemoryDecision({request: {...request, application_kind: "context_delivery"}, + observation: {...observation, status: "applied", application_receipt: {...observation.application_receipt, outcome: "applied"}}}); + assert.equal(delivered.status, "context_delivered"); + assert.equal(delivered.decision_consumption_complete, false); + assert.equal(delivered.semantic_disposition, null); + const recalled = projectRewardMemoryDecision({request: {...request, mode: "recall_only", has_applier: false}, observation}); + assert.equal(recalled.status, "recalled"); + assert.equal(recalled.decision_consumption_complete, false); + const ignored = projectRewardMemoryDecision({request, observation}); + assert.equal(ignored.decision_consumption_complete, true); + assert.equal(ignored.preserve_base_output, true); + assert.equal(ignored.utility_verified, false); + assert.equal(ignored.provider_call_count, 2); + assert.equal(ignored.filtered_count, 1); +}); + +test("stale/wrong attribution cannot complete semantic consumption", () => { + for (const patch of [{artifact_ref: "artifact:old"}, {surface_id: "other"}, {application_id: "other"}, + {memory_ref_digests: []}, {current_artifact_verified: false}, {result_readback_verified: false}, {outcome: "applied"}]) { + const result = projectRewardMemoryDecision({request, observation: {...observation, + application_receipt: {...observation.application_receipt, ...patch}}}); + assert.equal(result.status, "incomplete"); + assert.equal(result.decision_consumption_complete, false); + } +}); + +test("decode unknown transport values strictly", () => { + for (const patch of [{mode: "automatic"}, {application_kind: "inject_and_adopt"}, {query_ready: "true"}, + {artifact_ref: "private content with spaces"}]) { + assert.throws(() => planRewardMemoryDecision({...request, ...patch})); + } + for (const patch of [{provider_call_count: -1}, {filtered_count: 1.5}, {status: "success"}, + {boundary_reason_code: "private exception text"}]) { + assert.throws(() => projectRewardMemoryDecision({request, observation: {...observation, ...patch}})); + } +}); + +test("project only the original hook's safe typed boundary reason", () => { + const rejected = projectRewardMemoryDecision({request, observation: {...observation, + status: "guard_rejected", recall_status: null, provider_call_count: 0, + result_readback_verified: false, boundary_reason_code: "exact_corpus_request_invalid"}}); + assert.equal(rejected.status, "incomplete"); + assert.equal(rejected.reason_code, "recall_boundary_rejected"); + assert.equal(rejected.boundary_reason_code, "exact_corpus_request_invalid"); + assert.equal(rejected.provider_call_count, 0); +});