fix(dream): gate cursor on run completion

This commit is contained in:
chengyongru
2026-08-22 02:05:40 +08:00
committed by chengyongru
parent d853ac239f
commit 1fe36d5dec
5 changed files with 73 additions and 168 deletions
+3 -56
View File
@@ -53,42 +53,6 @@ if TYPE_CHECKING:
# ---------------------------------------------------------------------------
class DreamRunProgress:
"""Track tool failures that make a nominally completed Dream run unsafe to advance.
A failure in an earlier tool round is tolerated when the model observed
the error, retried, and the final tool round ran clean. Only failures the
model never got to correct invalidate the run.
"""
def __init__(self) -> None:
self.had_tool_errors = False
self.last_tool_round_had_errors: bool | None = None
@property
def recovered_from_tool_errors(self) -> bool:
"""True when errors occurred but the final tool round finished clean."""
return self.had_tool_errors and self.last_tool_round_had_errors is False
async def __call__(
self,
*_args: Any,
tool_events: list[dict[str, Any]] | None = None,
**_kwargs: Any,
) -> None:
events = [
event for event in tool_events or ()
if isinstance(cast(object, event), dict)
]
round_had_errors = any(event.get("phase") == "error" for event in events)
if round_had_errors:
self.had_tool_errors = True
# Terminal payloads ("end"/"error") close a tool round; the most
# recent closed round decides whether earlier failures were recovered.
if any(event.get("phase") in ("end", "error") for event in events):
self.last_tool_round_had_errors = round_had_errors
class MemoryStore:
"""Pure file I/O for memory files: MEMORY.md, history.jsonl, SOUL.md, USER.md."""
@@ -704,30 +668,16 @@ class MemoryStore:
@staticmethod
def dream_run_completed(
resp: object | None,
*,
had_tool_errors: bool = False,
recovered_tool_errors: bool = False,
) -> bool:
"""Return True only when a Dream turn finished cleanly enough to advance.
Tool failures from earlier rounds are acceptable when the final tool
round ran clean (``recovered_tool_errors``): the model observed the
failure, corrected it, and produced a consistent final state. Failures
in the final round still invalidate the run.
"""
"""Return True when the Dream agent reached a normal terminal response."""
metadata = getattr(resp, "metadata", None)
if not isinstance(metadata, dict):
return False
if cast(dict[str, Any], metadata).get("_stop_reason") != "completed":
return False
return not had_tool_errors or recovered_tool_errors
return cast(dict[str, Any], metadata).get("_stop_reason") == "completed"
@staticmethod
def dream_incompletion_reason(
resp: object | None,
*,
had_tool_errors: bool = False,
recovered_tool_errors: bool = False,
) -> str:
"""Human-readable explanation of why a Dream run cannot advance."""
metadata = getattr(resp, "metadata", None)
@@ -735,10 +685,7 @@ class MemoryStore:
stop_reason = cast(dict[str, Any], metadata).get("_stop_reason", "unknown")
else:
stop_reason = "missing response metadata"
parts = [f"stop_reason: {stop_reason}"]
if had_tool_errors and not recovered_tool_errors:
parts.append("unrecovered tool errors")
return ", ".join(parts)
return f"stop_reason: {stop_reason}"
# -- message formatting utility ------------------------------------------
+5 -14
View File
@@ -504,13 +504,12 @@ def _run_gateway(
# Dream is an internal job — run directly, not through the agent loop.
if job.name == "dream":
from nanobot.agent.memory import DreamRunProgress, MemoryStore
from nanobot.agent.memory import MemoryStore
dream_session_key = MemoryStore.dream_session_key
prune_dream_sessions = MemoryStore.prune_dream_sessions
store = agent.context.memory
progress = DreamRunProgress()
resp = None
diff_body = ""
try:
@@ -527,17 +526,13 @@ def _run_gateway(
session_key=key,
ephemeral=True,
tools=store.build_dream_tools(),
on_progress=progress,
on_progress=_silent,
runtime=dream_runtime,
)
# The real file delta grounds the audit record; clean completion
# The real file delta grounds the audit record; normal completion
# decides whether this history batch has finished processing.
diff_body = store.dream_content_diff()
completed = MemoryStore.dream_run_completed(
resp,
had_tool_errors=progress.had_tool_errors,
recovered_tool_errors=progress.recovered_from_tool_errors,
)
completed = MemoryStore.dream_run_completed(resp)
if completed:
store.set_last_dream_cursor(last_cursor)
if diff_body:
@@ -554,11 +549,7 @@ def _run_gateway(
else:
logger.warning(
"Dream cron job did not complete ({}); cursor remains at {}",
MemoryStore.dream_incompletion_reason(
resp,
had_tool_errors=progress.had_tool_errors,
recovered_tool_errors=progress.recovered_from_tool_errors,
),
MemoryStore.dream_incompletion_reason(resp),
store.get_last_dream_cursor(),
)
except Exception:
+8 -14
View File
@@ -423,14 +423,16 @@ async def cmd_dream(ctx: CommandContext) -> OutboundMessage:
msg = ctx.msg
async def _run_dream():
from nanobot.agent.memory import DreamRunProgress, MemoryStore
from nanobot.agent.memory import MemoryStore
async def _silent(*_args: Any, **_kwargs: Any) -> None:
pass
dream_session_key = MemoryStore.dream_session_key
build_dream_commit_message = MemoryStore.build_dream_commit_message
prune_dream_sessions = MemoryStore.prune_dream_sessions
store = loop.context.memory
progress = DreamRunProgress()
content = ""
resp = None
diff_body = ""
@@ -452,18 +454,14 @@ async def cmd_dream(ctx: CommandContext) -> OutboundMessage:
session_key=key,
ephemeral=True,
tools=store.build_dream_tools(),
on_progress=progress,
on_progress=_silent,
runtime=dream_runtime,
)
elapsed = time.monotonic() - t0
# The real file delta grounds the audit record; clean completion
# The real file delta grounds the audit record; normal completion
# decides whether this history batch has finished processing.
diff_body = store.dream_content_diff()
completed = MemoryStore.dream_run_completed(
resp,
had_tool_errors=progress.had_tool_errors,
recovered_tool_errors=progress.recovered_from_tool_errors,
)
completed = MemoryStore.dream_run_completed(resp)
if completed:
store.set_last_dream_cursor(last_cursor)
if diff_body:
@@ -471,11 +469,7 @@ async def cmd_dream(ctx: CommandContext) -> OutboundMessage:
else:
content = f"Dream completed in {elapsed:.1f}s; no memory changes."
else:
reason = MemoryStore.dream_incompletion_reason(
resp,
had_tool_errors=progress.had_tool_errors,
recovered_tool_errors=progress.recovered_from_tool_errors,
)
reason = MemoryStore.dream_incompletion_reason(resp)
content = (
f"Dream did not complete after {elapsed:.1f}s ({reason}); "
"memory cursor was not advanced."