Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 24 additions & 0 deletions docs/reference/reward-memory-decision-consumption.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 遥测,不重复查询。这不是自动跨进程存储或新的缓存。
Expand All @@ -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
Expand Down
66 changes: 54 additions & 12 deletions loopx/capabilities/reward_memory/decision.py
Original file line number Diff line number Diff line change
Expand Up @@ -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."""
Expand All @@ -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(
Expand Down Expand Up @@ -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,
*,
Expand Down Expand Up @@ -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)
Expand All @@ -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),
Expand All @@ -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",
Expand All @@ -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)
Loading
Loading