mirror of
https://github.com/openswarm-ai/openswarm.git
synced 2026-09-10 03:37:44 +02:00
[eric] agents: a session reaches disk the moment you send, so a backend death can no longer erase the conversation (ENG-313)
This commit is contained in:
@@ -17,7 +17,7 @@ from backend.apps.agents.core.models import (
|
||||
)
|
||||
from backend.apps.agents.core.ws_manager import ws_manager
|
||||
from backend.apps.settings.settings import load_settings
|
||||
from backend.apps.agents.manager.session.session_store import load_session_data
|
||||
from backend.apps.agents.manager.session.session_store import snapshot_session_now, load_session_data
|
||||
from backend.apps.agents.manager.session.apply_context_window import apply_context_window
|
||||
from backend.apps.agents.manager.session.workspace_git import (
|
||||
detect_git_identity,
|
||||
@@ -245,6 +245,7 @@ class AgentLaunch(AgentManagerProtocol):
|
||||
branch_id=fork.active_branch_id,
|
||||
)
|
||||
fork.messages.append(user_msg)
|
||||
snapshot_session_now(fork)
|
||||
await ws_manager.send_to_session(fork.id, "agent:message", {
|
||||
"session_id": fork.id,
|
||||
"message": user_msg.model_dump(mode="json"),
|
||||
|
||||
@@ -14,7 +14,7 @@ from backend.apps.agents.core.models import AgentSession, Message, MessageBranch
|
||||
from backend.apps.agents.core.ws_manager import ws_manager
|
||||
from backend.apps.settings.settings import load_settings
|
||||
from backend.apps.agents.manager.run_browser_fast_path import run_browser_fast_path
|
||||
from backend.apps.agents.manager.session.session_store import load_session_data
|
||||
from backend.apps.agents.manager.session.session_store import snapshot_session_now, load_session_data
|
||||
from backend.apps.agents.manager.session.apply_context_window import apply_context_window
|
||||
from backend.apps.agents.manager.prompt.tool_catalog import get_all_tool_names
|
||||
from backend.apps.agents.manager.prompt.prompt_context import resolve_mode
|
||||
@@ -128,6 +128,7 @@ class Messaging(AgentManagerProtocol):
|
||||
client_message_id=client_message_id,
|
||||
)
|
||||
session.messages.append(user_msg)
|
||||
snapshot_session_now(session)
|
||||
await ws_manager.send_to_session(session_id, "agent:message", {
|
||||
"session_id": session_id,
|
||||
"message": user_msg.model_dump(mode="json"),
|
||||
@@ -269,6 +270,7 @@ class Messaging(AgentManagerProtocol):
|
||||
attached_skills=target_msg.attached_skills,
|
||||
)
|
||||
session.messages.append(edited_msg)
|
||||
snapshot_session_now(session)
|
||||
|
||||
await ws_manager.send_to_session(session_id, "agent:message", {
|
||||
"session_id": session_id,
|
||||
|
||||
@@ -16,7 +16,7 @@ from backend.apps.agents.core.models import AgentSession, Message
|
||||
from backend.apps.agents.core.ws_manager import ws_manager
|
||||
from backend.apps.agents.manager.AgentManagerProtocol import AgentManagerProtocol
|
||||
from backend.apps.agents.manager.session.apply_context_window import apply_context_window
|
||||
from backend.apps.agents.manager.session.session_store import load_session_data
|
||||
from backend.apps.agents.manager.session.session_store import snapshot_session_now, load_session_data
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
@@ -80,6 +80,7 @@ class SpawnAgentRun(AgentManagerProtocol):
|
||||
branch_id=child.active_branch_id,
|
||||
)
|
||||
child.messages.append(user_msg)
|
||||
snapshot_session_now(child)
|
||||
await ws_manager.send_to_session(child.id, "agent:message", {
|
||||
"session_id": child.id,
|
||||
"message": user_msg.model_dump(mode="json"),
|
||||
|
||||
@@ -55,3 +55,19 @@ def build_search_text(session: AgentSession, max_len: int = 5000) -> str:
|
||||
parts.append(msg.content)
|
||||
text = " ".join(parts)
|
||||
return text[:max_len]
|
||||
|
||||
|
||||
@typechecked
|
||||
def snapshot_session_now(session: AgentSession) -> None:
|
||||
"""Write the session to disk immediately, swallowing failure. Called the moment a user message
|
||||
is appended: sessions used to reach disk only at turn END, so a backend death mid-turn destroyed
|
||||
the whole conversation (404, zero bytes) if it was the first turn, and ate the turn's user
|
||||
message otherwise (ENG-313). One bounded whole-file write per send, same cost as the existing
|
||||
end-of-turn save; a disk hiccup must never break the send itself."""
|
||||
try:
|
||||
doc = session.model_dump(mode="json")
|
||||
doc["search_text"] = build_search_text(session)
|
||||
save_session(session.id, doc)
|
||||
except Exception:
|
||||
import logging
|
||||
logging.getLogger(__name__).warning("snapshot_session_now failed for %s", session.id, exc_info=True)
|
||||
|
||||
@@ -0,0 +1,72 @@
|
||||
"""A session reaches disk the moment a user message is appended (ENG-313).
|
||||
|
||||
Sessions used to persist only at turn END (and orderly close/shutdown), so a backend death mid-turn
|
||||
destroyed the whole conversation if it was the first turn (404, zero bytes on disk, measured live on
|
||||
packaged exp.8) and silently ate the turn's user message otherwise. The boot-side restore machinery
|
||||
already handled files that exist; the write half was missing. These pin the write half.
|
||||
"""
|
||||
import os
|
||||
|
||||
from backend.apps.agents.core.models import AgentSession, Message
|
||||
from backend.apps.agents.manager.session.session_store import load_session_data, snapshot_session_now
|
||||
|
||||
|
||||
def p_make_session(sid: str) -> AgentSession:
|
||||
s = AgentSession(id=sid, name="snap-test", model="sonnet")
|
||||
s.status = "running"
|
||||
s.messages.append(Message(role="user", content="the message a crash must not eat", branch_id=s.active_branch_id))
|
||||
return s
|
||||
|
||||
|
||||
def test_snapshot_writes_a_loadable_file_with_the_user_message(tmp_path, monkeypatch):
|
||||
import backend.apps.agents.agent_manager as am
|
||||
|
||||
monkeypatch.setattr(am, "SESSIONS_DIR", str(tmp_path))
|
||||
s = p_make_session("snap-1")
|
||||
snapshot_session_now(s)
|
||||
data = load_session_data("snap-1")
|
||||
assert data is not None, "the whole bug: nothing on disk until turn end"
|
||||
assert data["messages"][0]["content"] == "the message a crash must not eat"
|
||||
assert data["status"] == "running"
|
||||
assert data["closed_at"] is None, "closed_at must stay unset or restore_all_sessions skips it"
|
||||
assert "search_text" in data, "history search must keep working on snapshotted sessions"
|
||||
|
||||
|
||||
def test_restore_marks_a_snapshotted_midturn_session_resumable(tmp_path, monkeypatch):
|
||||
# The restore half already existed; this proves the two halves meet: a snapshot taken mid-turn
|
||||
# (user message, no assistant reply) restores as "stopped", which is the resumable state.
|
||||
import backend.apps.agents.agent_manager as am
|
||||
|
||||
monkeypatch.setattr(am, "SESSIONS_DIR", str(tmp_path))
|
||||
s = p_make_session("snap-2")
|
||||
snapshot_session_now(s)
|
||||
data = load_session_data("snap-2")
|
||||
restored = AgentSession(**data)
|
||||
branch = restored.active_branch_id or "main"
|
||||
msgs = [m for m in restored.messages if (m.branch_id or "main") == branch]
|
||||
assert msgs and msgs[-1].role == "user"
|
||||
# Mirrors restore_all_sessions: last message is the user's -> agent was cut off owing a reply.
|
||||
expected = "completed" if msgs[-1].role == "assistant" else "stopped"
|
||||
assert expected == "stopped"
|
||||
|
||||
|
||||
def test_a_snapshot_failure_never_breaks_the_send(tmp_path, monkeypatch):
|
||||
import backend.apps.agents.agent_manager as am
|
||||
|
||||
monkeypatch.setattr(am, "SESSIONS_DIR", os.path.join(str(tmp_path), "no-such", "\0bad"))
|
||||
snapshot_session_now(p_make_session("snap-3")) # must not raise
|
||||
|
||||
|
||||
def test_every_user_message_append_site_snapshots():
|
||||
# The chokepoint audit: a new send path that forgets the snapshot reintroduces the bug for that
|
||||
# path only, which is exactly how the class comes back. Enumerate the sites.
|
||||
import backend.apps.agents.manager.AgentLaunch as launch
|
||||
import backend.apps.agents.manager.Messaging as messaging
|
||||
import backend.apps.agents.manager.SpawnAgentRun as spawn
|
||||
|
||||
for mod, expected_appends in ((messaging, 2), (launch, 1), (spawn, 1)):
|
||||
src = open(mod.__file__).read()
|
||||
appends = src.count('.messages.append(user_msg)') + src.count('.messages.append(edited_msg)')
|
||||
snaps = src.count('snapshot_session_now(')
|
||||
assert appends == expected_appends, f"{mod.__name__}: append sites moved; re-audit this test"
|
||||
assert snaps >= appends, f"{mod.__name__}: {appends} user-append site(s) but only {snaps} snapshot(s)"
|
||||
Reference in New Issue
Block a user