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
1 change: 1 addition & 0 deletions apps/presentation/dashboard/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,7 @@
"build:chat:vite": "tsc --noEmit && vite build --config vite.chat.config.ts",
"smoke:chat-upgrade": "LOOPX_PLAYWRIGHT_PACKAGE=\"$PWD/node_modules/playwright\" node ../../../examples/chat-bundle-upgrade-browser-smoke.mjs",
"test:conversation-returns": "node --experimental-strip-types src/data/conversation-returns.test.mjs",
"smoke:conversation-history": "vite build --ssr smoke/conversation-history-smoke.ts --outDir node_modules/.cache/loopx-conversation-history --emptyOutDir && node node_modules/.cache/loopx-conversation-history/conversation-history-smoke.js",
"smoke:recent-completions": "tsc --ignoreConfig --target ES2022 --module CommonJS --moduleResolution Node --ignoreDeprecations 6.0 --resolveJsonModule --esModuleInterop --jsx react-jsx --skipLibCheck --strict --types node --outDir /tmp/loopx-recent-completions-smoke smoke/recent-completions-smoke.ts && NODE_PATH=\"$PWD/node_modules\" node /tmp/loopx-recent-completions-smoke/apps/presentation/dashboard/smoke/recent-completions-smoke.js"
},
"dependencies": {
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
"""Disposable real Chat HTTP/store fixture; the fault changes reads only."""
from __future__ import annotations

import json
from pathlib import Path
import sys
import tempfile
import threading

sys.path.insert(0, str(Path(__file__).resolve().parents[4]))

from loopx.chat_runtime import ChatRuntimeController
from loopx.chat_server import ChatHTTPServer, ChatRequestHandler
from loopx.chat_store import ChatSessionStore


class HistoryReadFault(ChatRequestHandler):
def _session_snapshot(self, session_id: str) -> None:
if session_id == self.server.unavailable_session_id:
self._send_error("Synthetic history read unavailable", status=503)
else:
super()._session_snapshot(session_id)


def main() -> None:
with tempfile.TemporaryDirectory(prefix="loopx-history-http-") as directory:
root = Path(directory)
registry = root / "registry.json"
registry.write_text(json.dumps({"schema_version": "0.1", "goals": [
{"id": "research", "repo": str(root), "status": "active"},
]}), encoding="utf-8")
store = ChatSessionStore(root / "runtime")
runtime = ChatRuntimeController(store=store, codex_bin="missing-codex", registry_path=registry)
for session_id, channel, text in [
("old", "goal.research", "Earlier public report"),
("current", "goal.research", "Current public report"),
("other-channel", "manager", "Unrelated conversation"),
]:
store.create_session(goal_id="research", agent_id="codex", adapter_kind="codex_app_server",
upstream_thread_id=session_id, channel_id=channel, session_id=session_id)
store.append_message(session_id, role="agent", text=text, message_id="answer")
before = {str(path.relative_to(store.sessions_root)): path.read_bytes()
for path in store.sessions_root.rglob("*") if path.is_file()}
server = ChatHTTPServer(("127.0.0.1", 0), HistoryReadFault)
server.verbose = False
server.registry_path = registry
server.runtime_root = root / "runtime"
server.chat_store = store
server.runtime_controller = runtime
server.unavailable_session_id = "old"
thread = threading.Thread(target=server.serve_forever, daemon=True)
thread.start()
print(json.dumps({"origin": f"http://127.0.0.1:{server.server_port}"}), flush=True)
try:
if sys.stdin.readline().strip() != "recover":
raise ValueError("Expected read recovery")
server.unavailable_session_id = ""
print(json.dumps({"recovered": True}), flush=True)
if sys.stdin.readline().strip() != "inspect":
raise ValueError("Expected final inspection")
after = {str(path.relative_to(store.sessions_root)): path.read_bytes()
for path in store.sessions_root.rglob("*") if path.is_file()}
print(json.dumps({"store_unchanged": before == after,
"turn_count": sum(1 for _ in store.sessions_root.glob("*/turns/*.json"))}), flush=True)
finally:
server.shutdown()
server.server_close()
thread.join(timeout=5)
runtime.close()


if __name__ == "__main__":
main()
58 changes: 58 additions & 0 deletions apps/presentation/dashboard/smoke/conversation-history-smoke.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,58 @@
import assert from "node:assert/strict";
import { spawn } from "node:child_process";
import { once } from "node:events";
import { resolve } from "node:path";
import { createInterface } from "node:readline";
import { resolveTestPython } from "../../../../scripts/test-python.mjs";
import { fetchChatHistory } from "../src/data/chat";

const repoRoot = resolve(process.cwd(), "../../..");
const child = spawn(resolveTestPython({ repoRoot }), ["-u", "apps/presentation/dashboard/smoke/conversation-history-http-fixture.py"],
{ cwd: repoRoot, stdio: ["pipe", "pipe", "pipe"] });
const exited = once(child, "exit");
let stderr = "";
child.stderr.on("data", chunk => { stderr += String(chunk); });
const output = createInterface({ input: child.stdout });
const lines = output[Symbol.asyncIterator]();
const originalFetch = globalThis.fetch;
async function next() {
const line = await lines.next();
assert.equal(line.done, false, stderr);
return JSON.parse(line.value!);
}
try {
const { origin } = await next();
const requests: string[] = [];
globalThis.fetch = (input, init) => {
const path = String(input);
assert.ok(!init?.method || init.method === "GET", "History recovery is read-only");
requests.push(path);
return originalFetch(new URL(path, origin), init);
};
const options = { agentId: "codex", goalId: "research", channelId: "goal.research" };
const partial = await fetchChatHistory(options);
assert.deepEqual(partial.unavailableSessionIds, ["old"]);
assert.deepEqual(partial.messages.map(row => row.text), ["Current public report"]);
assert.equal(partial.sessions[0].session_id, "current");
assert.equal(partial.sessions.some(row => row.session_id === "other-channel"), false);
child.stdin.write("recover\n");
assert.equal((await next()).recovered, true);
requests.length = 0;
const restored = await fetchChatHistory(options, partial);
assert.deepEqual(requests, ["/api/chat/sessions/old"], "Recovery reads only the missing session");
assert.deepEqual(restored.unavailableSessionIds, []);
assert.deepEqual(restored.messages.map(row => row.text), ["Earlier public report", "Current public report"]);
assert.deepEqual(restored.messages.map(row => row.session_id), ["old", "current"], "Colliding IDs preserve both sessions");
requests.length = 0;
await fetchChatHistory(options, restored);
assert.deepEqual(requests, [], "A complete cached history does not poll");
child.stdin.end("inspect\n");
assert.deepEqual(await next(), { store_unchanged: true, turn_count: 0 });
const [exitCode] = await exited;
assert.equal(exitCode, 0, stderr);
console.log("conversation-history: passed (real HTTP/store, partial read, channel isolation, missing-only recovery, scoped identity, zero writes or Turns)");
} finally {
globalThis.fetch = originalFetch;
output.close();
if (child.exitCode === null) { child.kill(); await exited; }
}
38 changes: 28 additions & 10 deletions apps/presentation/dashboard/src/data/chat.ts
Original file line number Diff line number Diff line change
Expand Up @@ -624,7 +624,10 @@ export async function recordProjectionExchange(options: {
goalId?: string;
question: string;
}) {
return requestJson<{ ok: true; schema_version: "loopx_chat_projection_exchange_v1"; session_id: string }>(
return requestJson<{
ok: true; schema_version: "loopx_chat_projection_exchange_v1";
session_id: string; user_message_id: string; answer_message_id: string;
}>(
"/api/chat/projection-messages",
{
method: "POST",
Expand Down Expand Up @@ -744,14 +747,15 @@ export type ChatSessionSnapshot = {
active_turn: Record<string, unknown> | null;
};

export async function fetchChatSession(sessionId: string) {
return requestJson<ChatSessionSnapshot>(`/api/chat/sessions/${sessionId}`);
export async function fetchChatSession(sessionId: string, signal?: AbortSignal) {
return requestJson<ChatSessionSnapshot>(`/api/chat/sessions/${sessionId}`, { signal });
}

export async function fetchChatSessions(options: {
agentId?: string;
channelId?: string;
goalId?: string;
signal?: AbortSignal;
}) {
const query = new URLSearchParams();
if (options.agentId) query.set("agent_id", options.agentId);
Expand All @@ -761,14 +765,14 @@ export async function fetchChatSessions(options: {
ok: true;
schema_version: "loopx_chat_session_list_v1";
sessions: ChatSessionSummary[];
}>(`/api/chat/sessions?${query.toString()}`);
}>(`/api/chat/sessions?${query.toString()}`, { signal: options.signal });
}

export function mergeChatSessionMessages(snapshots: ChatSessionSnapshot[]) {
const messages = new Map<string, ChatVisibleMessage>();
for (const snapshot of snapshots) {
for (const message of snapshot.messages) {
messages.set(message.message_id, { ...message, session_id: snapshot.session.session_id });
messages.set(`${snapshot.session.session_id}:${message.message_id}`, { ...message, session_id: snapshot.session.session_id });
}
}
return [...messages.values()].sort((left, right) =>
Expand All @@ -777,22 +781,36 @@ export function mergeChatSessionMessages(snapshots: ChatSessionSnapshot[]) {
);
}

export type ChatHistory = {
messages: ChatVisibleMessage[];
sessions: ChatSessionSummary[];
snapshots: ChatSessionSnapshot[];
unavailableSessionIds: string[];
};

export async function fetchChatHistory(options: {
// An omitted ``agentId`` reads the whole channel transcript. The steward
// channel is one conversation across whatever executor it currently
// resolves, so the client must not filter it by its own assumed executor.
agentId?: string;
channelId: string;
goalId?: string;
}) {
const listed = await fetchChatSessions(options);
const snapshots = await Promise.all(
listed.sessions.map((session) => fetchChatSession(session.session_id)),
);
}, previous?: ChatHistory): Promise<ChatHistory> {
const listed = previous ?? await fetchChatSessions({ ...options, signal: AbortSignal.timeout(5000) });
const known = new Map(previous?.snapshots.map((snapshot) => [snapshot.session.session_id, snapshot]));
// A failed historical read is not an empty transcript. Retrying only the
// missing snapshots keeps this recovery read-only and bounds repeated work.
const missing = listed.sessions.filter((session) => !known.has(session.session_id));
const results = await Promise.allSettled(missing.map((session) => fetchChatSession(session.session_id, AbortSignal.timeout(5000))));
for (const result of results) {
if (result.status === "fulfilled") known.set(result.value.session.session_id, result.value);
}
const snapshots = listed.sessions.flatMap((session) => known.has(session.session_id) ? [known.get(session.session_id)!] : []);
return {
messages: mergeChatSessionMessages(snapshots),
sessions: listed.sessions,
snapshots,
unavailableSessionIds: listed.sessions.filter((session) => !known.has(session.session_id)).map((session) => session.session_id),
};
}

Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import assert from "node:assert/strict";
import { conversationReturnSessions, reconcileConversationReturns } from "./conversation-returns.ts";
import { conversationReturnSessions, reconcileConversationHistory, reconcileConversationReturns } from "./conversation-returns.ts";

const collaboration = { returns: [] };
const original = [
Expand Down Expand Up @@ -39,4 +39,32 @@ assert.equal(hydrated[0], original[0]);
assert.deepEqual(conversationReturnSessions(undefined, delivered), ["current"]);
assert.deepEqual(conversationReturnSessions(undefined, original), ["current", "old"]);
assert.deepEqual(conversationReturnSessions(undefined, [delivered[0], delivered.at(-1)]), []);
const historyRows = [
{ session_id: "old", message_id: "answer", role: "agent", turn_id: "old-turn", text: "Older answer", created_at: "2026-08-01" },
{ session_id: "current", message_id: "answer", role: "agent", turn_id: "current-turn", text: "Stored answer", created_at: "2026-08-02", collaboration },
];
const live = [{ sourceSessionId: "current", sourceTurnId: "current-turn", text: "Live answer", pending: true }];
const createHistory = (row) => ({ sourceSessionId: row.session_id, sourceMessageId: row.message_id,
sourceCreatedAt: row.created_at, text: row.text });
const recoveredHistory = reconcileConversationHistory(live, historyRows, createHistory);
assert.equal(recoveredHistory.length, 2, "Same Turn keeps its existing live answer");
assert.equal(recoveredHistory[0].text, "Older answer");
assert.equal(recoveredHistory[1].text, "Live answer");
assert.equal(recoveredHistory[1].pending, true);
assert.equal(recoveredHistory[1].collaboration, collaboration);
assert.equal(reconcileConversationHistory(recoveredHistory, historyRows, createHistory), recoveredHistory, "Repeated read is idempotent");
assert.equal(reconcileConversationHistory(recoveredHistory, [historyRows[0]], createHistory), recoveredHistory, "An incomplete read never retracts known history");
const bothRoles = [
{ sourceSessionId: "current", sourceTurnId: "same-turn", role: "user", text: "My request" },
{ sourceSessionId: "current", sourceTurnId: "same-turn", role: "assistant", text: "Live answer" },
];
const storedRoles = [
{ session_id: "current", message_id: "user-message", turn_id: "same-turn", role: "user", text: "My request" },
{ session_id: "current", message_id: "agent-message", turn_id: "same-turn", role: "agent", text: "Stored answer" },
];
const hydratedRoles = reconcileConversationHistory(bothRoles, storedRoles, createHistory);
assert.equal(hydratedRoles.length, 2);
assert.equal(hydratedRoles[0].sourceMessageId, "user-message");
assert.equal(hydratedRoles[1].sourceMessageId, "agent-message");
assert.equal(hydratedRoles[1].text, "Live answer");
console.log("conversation-returns: passed (session isolation, late return, deduplication, transport uncertainty, stream preservation and watch retirement)");
30 changes: 25 additions & 5 deletions apps/presentation/dashboard/src/data/conversation-returns.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,8 @@ type ConversationMessage = {
sourceSessionId?: string;
sourceMessageId?: string;
sourceTurnId?: string;
sourceCreatedAt?: string;
role?: string;
collaboration?: ChatVisibleMessage["collaboration"];
returnDelivery?: ChatVisibleMessage["return_delivery"];
};
Expand All @@ -29,23 +31,25 @@ export function reconcileConversationReturns<T extends ConversationMessage>(
createReply: (message: ChatVisibleMessage) => T,
): T[] {
const byId = new Map(messages.map((row) => [row.message_id, row]));
const byTurn = new Map(messages.filter((row) => row.role !== "user" && row.origin !== "manager_followup")
.map((row) => [row.turn_id, row]));
const byTurn = new Map(messages.filter((row) => row.origin !== "manager_followup")
.map((row) => [`${row.turn_id}:${row.role === "user" ? "user" : "assistant"}`, row]));
const seen = new Set(previous.filter((row) => row.sourceSessionId === sessionId).map((row) => row.sourceMessageId));
let changed = false;
const updated = previous.map((row) => {
if (row.sourceSessionId !== sessionId) return row;
const source = row.sourceMessageId ? byId.get(row.sourceMessageId) : row.sourceTurnId ? byTurn.get(row.sourceTurnId) : undefined;
const source = row.sourceMessageId ? byId.get(row.sourceMessageId) : row.sourceTurnId
? byTurn.get(`${row.sourceTurnId}:${row.role === "user" ? "user" : "assistant"}`) : undefined;
if (!source) return row;
// Projection absence is not a retraction: the backend can temporarily be
// unable to read collaboration metadata. Keep the last observed receipt and
// its outstanding read obligation until a newer observation arrives.
const returnDelivery = source.return_delivery ?? row.returnDelivery;
const collaboration = source.collaboration ?? row.collaboration;
const sourceCreatedAt = source.created_at ?? row.sourceCreatedAt;
if (row.sourceMessageId === source.message_id && JSON.stringify(row.returnDelivery) === JSON.stringify(returnDelivery)
&& JSON.stringify(row.collaboration) === JSON.stringify(collaboration)) return row;
&& JSON.stringify(row.collaboration) === JSON.stringify(collaboration) && row.sourceCreatedAt === sourceCreatedAt) return row;
changed = true;
return { ...row, sourceMessageId: source.message_id, returnDelivery, collaboration };
return { ...row, sourceMessageId: source.message_id, sourceCreatedAt, returnDelivery, collaboration };
});
for (const message of messages) {
if (message.origin !== "manager_followup" || seen.has(message.message_id)) continue;
Expand All @@ -55,3 +59,19 @@ export function reconcileConversationReturns<T extends ConversationMessage>(
}
return changed ? updated : previous;
}

/** Recover missing history without retracting messages or overwriting live text. */
export function reconcileConversationHistory<T extends ConversationMessage>(
previous: T[], messages: ChatVisibleMessage[], createMessage: (message: ChatVisibleMessage) => T,
): T[] {
let updated = previous;
const sessions = new Set(messages.map((message) => message.session_id).filter((id): id is string => Boolean(id)));
for (const sessionId of sessions) {
updated = reconcileConversationReturns(updated, sessionId, messages.filter((message) => message.session_id === sessionId), createMessage);
}
const seen = new Set(updated.map((message) => `${message.sourceSessionId}:${message.sourceMessageId}`));
const recovered = messages.filter((message) => !seen.has(`${message.session_id}:${message.message_id}`)).map(createMessage);
if (!recovered.length) return updated;
return [...updated, ...recovered].sort((left, right) => !left.sourceCreatedAt ? (right.sourceCreatedAt ? 1 : 0)
: !right.sourceCreatedAt ? -1 : left.sourceCreatedAt.localeCompare(right.sourceCreatedAt));
}
Loading
Loading