From 2b7b0cbeeb6f6dc667064b4925ef2e049f0d8cb3 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Mon, 28 Sep 2026 07:38:10 +0800 Subject: [PATCH] fix(reward-memory): recover pending TS projections without repeating work Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../reward-memory-decision-consumption.md | 24 ++ loopx/capabilities/reward_memory/decision.py | 66 +++++- .../test_reward_memory_decision.py | 213 ++++++++++++++++++ 3 files changed, 291 insertions(+), 12 deletions(-) diff --git a/docs/reference/reward-memory-decision-consumption.md b/docs/reference/reward-memory-decision-consumption.md index 5a488f0c2..075a6e376 100644 --- a/docs/reference/reward-memory-decision-consumption.md +++ b/docs/reference/reward-memory-decision-consumption.md @@ -175,6 +175,21 @@ are separate facts. This flag attests the caller's verified SDK context callback **not** frontend/Lark transport delivery or model utility. Public packets alone cannot recreate that private lineage or upgrade historical receipts. +If the SDK provider/application has finished but TS result projection is +temporarily unavailable, the complete result retains its private observation +and pending output. An exact `previous_result` replay retries only the existing +TS projection. Assessment first recovers the projection; it may then make the +first semantic judgment after verified context delivery, but does not repeat +an already-attempted semantic callback, even when its evidence is invalid. +A later explicit assessment may correct that incomplete SDK evidence without +recall. While TS remains unavailable, the same +incomplete result and baseline remain; after recovery, TS revalidates the +original receipts before exposing the retained output or completion. Changed +configuration/question/scope/artifact still fails the exact request fence; +invalid application evidence is not upgraded. Pre-provider failures are not +automatically retried. No new provider permission, persistent store, retry loop +or cross-process restore API is introduced. + 仅 `public_packet` 用于展示,其余结果私有。通过既有执行上下文保留结果; `previous_result` 仅复用配置和输入均匹配的请求,变化则拒绝复用。后续判断使用 原条目和累计多 corpus 遥测,不重复查询。这不是自动跨进程存储或新的缓存。 @@ -188,6 +203,15 @@ EOF/重启丢失 recall_session 时,交付过上下文也不能完成 assessme 该标记证明调用方已验证的 SDK 上下文 callback,不证明前端/飞书传输或模型收益; 仅凭公开 packet 不能重建这条私有链路,也不追溯升级历史回执。 +若原 SDK 的 provider/应用已完成,但 TS 结果投影暂时不可用,完整 result 保留 +私有观察和待确认产物。相同请求的 previous_result 仅重试既有 TS 投影;assessment +先恢复投影,可在交付验证后作第一次语义判断,但不重复已尝试的语义 callback, +即使其证据无效。后续显式 assessment 可纠正不完整的 SDK 证据,但不重新召回。 +故障期间仍返回同一 incomplete 结果和 +基线;恢复后由 TS 重新核验原回执,再披露已保留产物或完成状态。配置、问题、范围 +或产物变化仍被精确请求 fence 拒绝,无效判断不能升级;provider 前的失败不自动 +重试。不新增 provider 权限、持久存储、自动重试循环或跨进程恢复 API。 + 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 diff --git a/loopx/capabilities/reward_memory/decision.py b/loopx/capabilities/reward_memory/decision.py index e719ad42a..f37ef30d7 100644 --- a/loopx/capabilities/reward_memory/decision.py +++ b/loopx/capabilities/reward_memory/decision.py @@ -20,6 +20,18 @@ from .runtime_hooks import run_reward_memory_automatic_recall_hook +@dataclass(frozen=True) +class _PendingDecisionProjection: + """Private SDK observation awaiting the existing TS projection, not new work.""" + + request: dict[str, Any] + status: str + telemetry: Mapping[str, Any] + receipt: Mapping[str, Any] | None + context_delivery_receipt: Mapping[str, Any] | None + output: Any + + @dataclass(frozen=True) class RewardMemoryDecisionResult: """Only public_packet is a projection. All other fields stay caller-private.""" @@ -33,6 +45,7 @@ class RewardMemoryDecisionResult: application_receipt: Mapping[str, Any] | None = None recall_telemetry: Mapping[str, Any] | None = None context_delivery_receipt: Mapping[str, Any] | None = None + pending_projection: _PendingDecisionProjection | None = None def _transport_failure( @@ -93,6 +106,25 @@ def _project( }) +def _recover_pending_projection(result: RewardMemoryDecisionResult) -> RewardMemoryDecisionResult: + pending = result.pending_projection + if pending is None: + return result + try: + packet = _project(pending.request, pending.status, pending.telemetry, + pending.receipt, pending.context_delivery_receipt) + delivery_receipt = pending.context_delivery_receipt + if delivery_receipt is None and packet["context_delivery_verified"]: + delivery_receipt = deepcopy(pending.receipt) + return replace(result, public_packet=packet, request=pending.request, + output=result.base_output if packet["preserve_base_output"] else pending.output, + application_receipt=pending.receipt, recall_telemetry=pending.telemetry, + context_delivery_receipt=delivery_receipt, pending_projection=None) + except (KeyError, OSError, RuntimeError, TypeError, ValueError): + # Keep the same private observation and fail-open baseline; do not rerun SDK work. + return result + + def run_reward_memory_decision( config: Mapping[str, Any] | None, *, @@ -131,7 +163,7 @@ def run_reward_memory_decision( return RewardMemoryDecisionResult( _transport_failure(request, "replay_request_mismatch"), base, base, digest, request, ) - return previous_result + return _recover_pending_projection(previous_result) plan = effect_runtime_result("reward_memory.decision.plan", request) if not plan["should_recall"]: return RewardMemoryDecisionResult(plan, base, base, digest, request) @@ -154,12 +186,13 @@ def apply(original: Any, items: tuple[RewardMemoryRecallItem, ...]) -> Mapping[s 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, - deepcopy(receipt) if packet["context_delivery_verified"] else None, + pending = _PendingDecisionProjection(deepcopy(request), hook["status"], + deepcopy(telemetry), deepcopy(receipt), None, deepcopy(hook["output"])) + result = RewardMemoryDecisionResult( + _transport_failure(request, "consumer_input_or_runtime_failed", telemetry), + base, base, digest, request, session, receipt, telemetry, pending_projection=pending, ) + return _recover_pending_projection(result) except (KeyError, OSError, RuntimeError, TypeError, ValueError): return RewardMemoryDecisionResult( _transport_failure(request, "consumer_input_or_runtime_failed", telemetry), @@ -176,8 +209,14 @@ def assess_reward_memory_decision( 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. + pending projection retries only TS readback; a completed assessment never + calls the model again. """ + recovering_assessment = (delivered.pending_projection is not None and + delivered.pending_projection.request.get("application_kind") == "semantic_application") + delivered = _recover_pending_projection(delivered) + if recovering_assessment or delivered.pending_projection is not None: + return delivered if delivered.public_packet.get("decision_consumption_complete") is True: return delivered request = {**delivered.request, "application_kind": "semantic_application", @@ -193,11 +232,14 @@ def assess_reward_memory_decision( ) 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"], delivered.context_delivery_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"]) + pending = _PendingDecisionProjection(deepcopy(request), application["status"], + deepcopy(delivered.recall_telemetry or {}), deepcopy(receipt), + deepcopy(delivered.context_delivery_receipt), deepcopy(application["output"])) + result = replace(delivered, + public_packet=_transport_failure(request, "consumer_input_or_runtime_failed", delivered.recall_telemetry), + request=request, output=delivered.base_output, application_receipt=receipt, + pending_projection=pending) + return _recover_pending_projection(result) 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/tests/capabilities/test_reward_memory_decision.py b/tests/capabilities/test_reward_memory_decision.py index 4705ceadf..d84ee367f 100644 --- a/tests/capabilities/test_reward_memory_decision.py +++ b/tests/capabilities/test_reward_memory_decision.py @@ -431,3 +431,216 @@ def reject_projection(*args): assert result.application_receipt["outcome"] == "ignored" assert result.public_packet["provider_call_count"] == provider.calls == 1 assert result.output == arguments["base_output"] + + +def test_delivery_projection_recovers_exact_result_without_recalling_or_reapplying(tmp_path, monkeypatch): + config, arguments, records = context(tmp_path) + provider = Provider(records) + transport = decision.effect_runtime_result + deliveries = [] + + def deliver(base, items): + deliveries.append(items) + return delivery(base, items) + + def reject_projection(method, params): + if method == "reward_memory.decision.project": + raise RuntimeError("private transport diagnostic") + return transport(method, params) + + monkeypatch.setattr(decision, "effect_runtime_result", reject_projection) + failed = run_reward_memory_decision(config, query_ready=True, application_kind="context_delivery", + apply_memory=deliver, provider=provider, **arguments) + assert failed.output == arguments["base_output"] + assert not failed.public_packet["context_delivery_verified"] + assert len(deliveries) == provider.calls == 1 + monkeypatch.setattr(decision, "effect_runtime_result", transport) + recovered = run_reward_memory_decision(config, query_ready=True, application_kind="context_delivery", + apply_memory=deliver, previous_result=failed, provider=provider, **arguments) + assert recovered.public_packet["status"] == "context_delivered" + assert recovered.public_packet["context_delivery_verified"] is True + assert not recovered.public_packet["decision_consumption_complete"] + assert recovered.output["private_context"] == ["Private reviewed summary language lesson."] + assert recovered.context_delivery_receipt == failed.application_receipt + assert len(deliveries) == provider.calls == recovered.public_packet["provider_call_count"] == 1 + assert "private transport diagnostic" not in json.dumps(recovered.public_packet) + + +@pytest.mark.parametrize("outcome", ["applied", "ignored", "refuted", "invalid"]) +def test_assessment_projection_recovers_original_output_without_second_judgment(tmp_path, monkeypatch, outcome): + config, arguments, records = context(tmp_path, two_corpora=True) + next(iter(records.values()))["corpus_id"] = "unrelated" + provider = Provider(records) + transport = decision.effect_runtime_result + delivered = run_reward_memory_decision(config, query_ready=True, application_kind="context_delivery", + apply_memory=delivery, provider=provider, **arguments) + assessments = [] + + def judge(base, items): + assessments.append(items) + return {"outcome": "ignored" if outcome == "invalid" else outcome, + "output": {"summary": "reviewed"} if outcome == "applied" else base, + "memory_refs": ["foreign"] if outcome == "invalid" else [item.memory_ref for item in items], + "current_artifact_verified": True, + "reasoning_summary": "Compared this exact retained context with the current artifact."} + + def reject_projection(*args): + raise RuntimeError("TS unavailable") + + monkeypatch.setattr(decision, "effect_runtime_result", reject_projection) + failed = assess_reward_memory_decision(delivered, apply_memory=judge) + assert failed.output == arguments["base_output"] + assert not failed.public_packet["decision_consumption_complete"] + assert len(assessments) == 1 + monkeypatch.setattr(decision, "effect_runtime_result", transport) + recovered = assess_reward_memory_decision(failed, apply_memory=lambda *args: pytest.fail("second model judgment")) + assert recovered.public_packet["decision_consumption_complete"] is (outcome != "invalid") + assert recovered.public_packet["context_delivery_verified"] is True + assert recovered.public_packet["semantic_disposition"] == (None if outcome == "invalid" else outcome) + assert recovered.output == ({"summary": "reviewed"} if outcome == "applied" else arguments["base_output"]) + assert recovered.public_packet["provider_call_count"] == provider.calls == 2 + assert recovered.public_packet["filtered_count"] == 1 + assert len(assessments) == 1 + assert recovered.public_packet["utility_verified"] is False + + +def test_assess_recovers_delivery_before_one_bound_judgment(tmp_path, monkeypatch): + config, arguments, records = context(tmp_path) + provider = Provider(records) + transport = decision.effect_runtime_result + + def reject_projection(method, params): + if method == "reward_memory.decision.project": + raise RuntimeError("TS unavailable") + return transport(method, params) + + monkeypatch.setattr(decision, "effect_runtime_result", reject_projection) + failed = run_reward_memory_decision(config, query_ready=True, application_kind="context_delivery", + apply_memory=delivery, provider=provider, **arguments) + assert assess_reward_memory_decision(failed, apply_memory=lambda *args: pytest.fail("unverified delivery")) is failed + monkeypatch.setattr(decision, "effect_runtime_result", transport) + judgments = [] + + def judge(base, items): + judgments.append(items) + return {"outcome": "ignored", "output": base, "memory_refs": [items[0].memory_ref], + "current_artifact_verified": True, "reasoning_summary": "The current artifact already covers this lesson."} + + result = assess_reward_memory_decision(failed, apply_memory=judge) + assert result.public_packet["decision_consumption_complete"] is True + assert result.public_packet["context_delivery_verified"] is True + assert result.output == arguments["base_output"] + assert len(judgments) == provider.calls == 1 + + +@pytest.mark.parametrize("change", ["artifact", "query", "scope", "configuration"]) +def test_pending_projection_keeps_exact_request_fence(tmp_path, monkeypatch, change): + config, arguments, records = context(tmp_path) + provider = Provider(records) + transport = decision.effect_runtime_result + + def reject_projection(method, params): + if method == "reward_memory.decision.project": + raise RuntimeError("TS unavailable") + return transport(method, params) + + monkeypatch.setattr(decision, "effect_runtime_result", reject_projection) + failed = run_reward_memory_decision(config, query_ready=True, application_kind="context_delivery", + apply_memory=delivery, provider=provider, **arguments) + changed_config, changed_arguments = copy.deepcopy(config), copy.deepcopy(arguments) + if change == "configuration": + changed_config["surfaces"][SURFACE]["recall_profile"]["limit"] = 1 + elif change == "artifact": + changed_arguments["artifact_ref"] = "artifact:new" + elif change == "query": + changed_arguments["queries"][0]["query"] = "A different question" + else: + changed_arguments["workspace_ref"] = "workspace:unrelated" + monkeypatch.setattr(decision, "effect_runtime_result", lambda *args: pytest.fail("mismatch TS call")) + mismatch = run_reward_memory_decision(changed_config, query_ready=True, application_kind="context_delivery", + apply_memory=delivery, previous_result=failed, provider=provider, **changed_arguments) + assert mismatch.public_packet["reason_code"] == "replay_request_mismatch" + assert mismatch.output == changed_arguments["base_output"] + assert provider.calls == 1 + monkeypatch.setattr(decision, "effect_runtime_result", transport) + recovered = run_reward_memory_decision(config, query_ready=True, application_kind="context_delivery", + apply_memory=delivery, previous_result=failed, provider=provider, **arguments) + assert recovered.public_packet["context_delivery_verified"] is True + assert provider.calls == 1 + + +def test_projection_recovery_cannot_upgrade_invalid_model_evidence(tmp_path, monkeypatch): + config, arguments, records = context(tmp_path) + provider = Provider(records) + transport = decision.effect_runtime_result + judgments = [] + + def invalid(base, items): + judgments.append(items) + return {"outcome": "applied", "output": {"summary": "unverified"}, "memory_refs": ["foreign"], + "current_artifact_verified": True, "reasoning_summary": "Foreign attribution is not valid."} + + def reject_projection(method, params): + if method == "reward_memory.decision.project": + raise RuntimeError("TS unavailable") + return transport(method, params) + + monkeypatch.setattr(decision, "effect_runtime_result", reject_projection) + failed = run_reward_memory_decision(config, query_ready=True, application_kind="semantic_application", + apply_memory=invalid, provider=provider, **arguments) + monkeypatch.setattr(decision, "effect_runtime_result", transport) + recovered = run_reward_memory_decision(config, query_ready=True, application_kind="semantic_application", + apply_memory=invalid, previous_result=failed, provider=provider, **arguments) + assert recovered.public_packet["status"] == "incomplete" + assert not recovered.public_packet["decision_consumption_complete"] + assert not recovered.public_packet["context_delivery_verified"] + assert recovered.output == arguments["base_output"] + assert len(judgments) == provider.calls == 1 + + +def test_pre_provider_transport_failure_does_not_implicitly_restart_work(tmp_path, monkeypatch): + config, arguments, records = context(tmp_path) + provider = Provider(records) + + def reject_transport(*args): + raise RuntimeError("TS unavailable") + + monkeypatch.setattr(decision, "effect_runtime_result", reject_transport) + failed = run_reward_memory_decision(config, query_ready=True, application_kind="context_delivery", + apply_memory=delivery, provider=provider, **arguments) + assert failed.public_packet["reason_code"] == "consumer_input_or_runtime_failed" + monkeypatch.setattr(decision, "effect_runtime_result", lambda *args: pytest.fail("implicit work restart")) + assert run_reward_memory_decision(config, query_ready=True, application_kind="context_delivery", + apply_memory=delivery, previous_result=failed, provider=provider, **arguments) is failed + assert provider.calls == 0 + + +def test_direct_semantic_projection_recovers_without_inventing_delivery(tmp_path, monkeypatch): + config, arguments, records = context(tmp_path) + provider = Provider(records) + transport = decision.effect_runtime_result + output = {"summary": "reviewed"} + judgments = [] + + def judge(base, items): + judgments.append(items) + return {"outcome": "applied", "output": output, "memory_refs": [items[0].memory_ref], + "current_artifact_verified": True, "reasoning_summary": "Applied this lesson to the current artifact."} + + def reject_projection(method, params): + if method == "reward_memory.decision.project": + raise RuntimeError("TS unavailable") + return transport(method, params) + + monkeypatch.setattr(decision, "effect_runtime_result", reject_projection) + failed = run_reward_memory_decision(config, query_ready=True, application_kind="semantic_application", + apply_memory=judge, provider=provider, **arguments) + output["summary"] = "later unrelated mutation" + monkeypatch.setattr(decision, "effect_runtime_result", transport) + recovered = run_reward_memory_decision(config, query_ready=True, application_kind="semantic_application", + apply_memory=judge, previous_result=failed, provider=provider, **arguments) + assert recovered.public_packet["decision_consumption_complete"] is True + assert recovered.public_packet["context_delivery_verified"] is False + assert recovered.output == {"summary": "reviewed"} + assert recovered.context_delivery_receipt is None + assert len(judgments) == provider.calls == 1