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
3 changes: 2 additions & 1 deletion apps/presentation/dashboard/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,8 @@
"smoke:team-artifact-comparison": "tsc --ignoreConfig --target ES2022 --module ES2022 --moduleResolution Bundler --skipLibCheck --strict --rootDir src --outDir node_modules/.cache/team-comparison src/features/personal-workspace/team-artifact-comparison.ts src/vite-env.d.ts && node smoke/team-artifact-comparison-smoke.mjs",
"smoke:team-report": "tsc --ignoreConfig --target ES2022 --module ES2022 --moduleResolution Bundler --jsx react-jsx --skipLibCheck --strict --rootDir src --outDir node_modules/.cache/team-report src/features/personal-workspace/team-artifact-content.tsx src/vite-env.d.ts && node smoke/team-report-smoke.mjs",
"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"
"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"
},
"dependencies": {
"@fontsource-variable/geist": "^5.3.0",
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
import assert from "node:assert/strict";
import { conversationReturnSessions, reconcileConversationReturns } from "./conversation-returns.ts";

const collaboration = { returns: [] };
const original = [
{ sourceSessionId: "old", sourceMessageId: "brief", collaboration, text: "Original brief" },
{ sourceSessionId: "current", sourceMessageId: "brief", text: "New conversation" },
{ sourceSessionId: "current", sourceTurnId: "running", text: "Streaming text", pending: true },
];
const conclusion = { message_id: "result", turn_id: "original", origin: "manager_followup", text: "Checked result",
return_delivery: { phase: "conclusion", status: "verification_required" } };
const createReply = (row) => ({ sourceSessionId: "old", sourceMessageId: row.message_id, text: row.text, returnDelivery: row.return_delivery });
assert.deepEqual(conversationReturnSessions("current", original), ["current", "old"]);
// Reading another session must not erase the old brief, even with colliding IDs.
const current = reconcileConversationReturns(original, "current", [{ message_id: "brief", text: "New conversation" }], createReply);
assert.equal(current, original);
assert.equal(reconcileConversationReturns(original, "old", [], createReply), original);
assert.equal(reconcileConversationReturns(original, "old", [{ message_id: "brief" }], createReply), original);
const arrived = reconcileConversationReturns(original, "old", [conclusion, conclusion], createReply);
assert.equal(arrived.length, original.length + 1);
assert.equal(arrived[2], original[2]);
assert.equal(reconcileConversationReturns(arrived, "old", [conclusion], createReply), arrived);
assert.equal(arrived.at(-1).text, "Checked result");
// The worker has concluded, but transport uncertainty still requires readback.
const settledBrief = { message_id: "brief", collaboration: { returns: [{ phase: "conclusion", status: "delivered" }] } };
const waiting = reconcileConversationReturns(arrived, "old", [settledBrief, conclusion], createReply);
assert.deepEqual(conversationReturnSessions("current", waiting), ["current", "old"]);
const delivered = reconcileConversationReturns(waiting, "old", [settledBrief, { ...conclusion,
return_delivery: { phase: "conclusion", status: "delivered" } }], createReply);
assert.deepEqual(conversationReturnSessions("current", delivered), ["current"]);
assert.equal(delivered.length, arrived.length);
// Recovered Turns acquire stored identity without replacing the live text.
const hydrated = reconcileConversationReturns(original, "current", [{ message_id: "answer", turn_id: "running",
role: "agent", text: "Stored text", collaboration }], createReply);
assert.equal(hydrated[2].sourceMessageId, "answer");
assert.equal(hydrated[2].text, "Streaming text");
assert.equal(hydrated[2].pending, true);
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)]), []);
console.log("conversation-returns: passed (session isolation, late return, deduplication, transport uncertainty, stream preservation and watch retirement)");
57 changes: 57 additions & 0 deletions apps/presentation/dashboard/src/data/conversation-returns.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,57 @@
import type { ChatVisibleMessage } from "./chat";

type ConversationMessage = {
sourceSessionId?: string;
sourceMessageId?: string;
sourceTurnId?: string;
collaboration?: ChatVisibleMessage["collaboration"];
returnDelivery?: ChatVisibleMessage["return_delivery"];
};

/** Old sessions remain relevant while they owe a result, not forever. */
export function conversationReturnSessions(activeSessionId: string | undefined, messages: ConversationMessage[]): string[] {
const sessions = new Set(activeSessionId ? [activeSessionId] : []);
for (const message of messages) {
const waitingForConclusion = message.collaboration && !message.collaboration.returns.some(
(reply) => reply.phase === "conclusion" && reply.status === "delivered",
);
const delivery = message.returnDelivery;
const waitingForDelivery = delivery && !["delivered", "superseded"].includes(delivery.status);
const waitingForTranscript = message.sourceTurnId && !message.sourceMessageId;
if (message.sourceSessionId && (waitingForConclusion || waitingForDelivery || waitingForTranscript)) sessions.add(message.sourceSessionId);
}
return [...sessions].sort();
}

/** A snapshot may refresh only its own session; it never replaces streamed text. */
export function reconcileConversationReturns<T extends ConversationMessage>(
previous: T[], sessionId: string, messages: ChatVisibleMessage[],
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 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;
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;
if (row.sourceMessageId === source.message_id && JSON.stringify(row.returnDelivery) === JSON.stringify(returnDelivery)
&& JSON.stringify(row.collaboration) === JSON.stringify(collaboration)) return row;
changed = true;
return { ...row, sourceMessageId: source.message_id, returnDelivery, collaboration };
});
for (const message of messages) {
if (message.origin !== "manager_followup" || seen.has(message.message_id)) continue;
seen.add(message.message_id);
changed = true;
updated.push(createReply(message));
}
return changed ? updated : previous;
}
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ const statusSourceSwitcher = source("./status-source-switcher.tsx");
const workspaceSettings = source("./workspace-settings-page.tsx");
const styles = source("./personal-workspace.css");
const dashboard = source("../../views/dashboard-page.tsx");
const conversationReturns = source("../../data/conversation-returns.ts");
const tasks = source("./goal-tasks-view.tsx");
const status = source("../../data/status.ts");
const chatData = source("../../data/chat.ts");
Expand Down Expand Up @@ -100,7 +101,8 @@ assert.doesNotMatch(page.match(/function operationProposalFields[\s\S]*?\n\}/)?.
assert.match(page, /t\("proposal\.primary\.operationGroup"\)/, "Operation confirmation routes users to the bound group");
assert.match(chatData, /result_delivery:/, "Dashboard retains operation result-delivery readback");
assert.match(chatData, /return_delivery\??:/, "Chat messages retain manager return-delivery readback");
assert.match(dashboard, /deliveryByMessage/, "Manager return polling refreshes delivery state after the message arrives");
assert.match(dashboard, /reconcileConversationReturns\(/, "Manager return polling uses the shared conversation-return read model");
assert.match(conversationReturns, /source\.return_delivery \?\? row\.returnDelivery/, "Manager return polling refreshes delivery state after the message arrives");
assert.match(timeline + page, /ReturnDeliveryStatus/, "Both manager conversation surfaces render return delivery state");
for (const state of ["delivered", "verification_required", "explicit_unverified"]) {
assert.match(returnDelivery, new RegExp(state), `Return delivery renders ${state}`);
Expand Down
67 changes: 28 additions & 39 deletions apps/presentation/dashboard/src/views/dashboard-page.tsx
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import { conversationReturnSessions, reconcileConversationReturns } from "../data/conversation-returns";
import {compactWorkspaceText as compactShareText} from "../features/personal-workspace/personal-workspace-model";
import type { GoalAcceptanceObservation } from "../data/goal-acceptance-observation";
import { attentionDetails, sourceAttention } from "../features/personal-workspace/attention-details";
Expand Down Expand Up @@ -1484,59 +1485,44 @@ function PersonalGoalHome({
statusSourceControl.activeSource.statusUrl,
]);

// Worker returns are transcript messages, not new model turns. Keep an open
// conversation current without replacing in-flight user/agent text.
const conversationReturnSessionId = runtimeBindings[contextId]?.sessionId;
// Read the active session plus older sessions that still owe a result. The
// stable key changes only when that set changes, never on each stream delta.
const conversationReturnSessionKey = JSON.stringify(conversationReturnSessions(
runtimeBindings[contextId]?.sessionId, messagesByContext[contextId] ?? [],
));
useEffect(() => {
if (readOnly || !conversationReturnSessionId) return;
if (readOnly) return;
const sessionIds: string[] = JSON.parse(conversationReturnSessionKey);
let cancelled = false;
let timer: ReturnType<typeof setTimeout> | undefined;
const receive = async () => {
const timers = new Set<ReturnType<typeof setTimeout>>();
const receive = async (sessionId: string) => {
try {
const snapshot = await fetchChatSession(conversationReturnSessionId);
const snapshot = await fetchChatSession(sessionId);
if (cancelled) return;
const replies = snapshot.messages.filter((row) => row.origin === "manager_followup");
setMessagesByContext((current) => {
const previous = current[contextId] ?? [];
const seen = new Set(previous.map((row) => row.sourceMessageId));
const fresh = replies.filter((row) => !seen.has(row.message_id));
const deliveryByMessage = new Map(
replies.map((row) => [row.message_id, row.return_delivery]),
);
const byTurn = new Map(snapshot.messages.filter((row) => row.role !== "user" && row.origin !== "manager_followup").map((row) => [row.turn_id, row]));
const collaborationByMessage = new Map(snapshot.messages.map((row) => [row.message_id, row.collaboration]));
let deliveryChanged = false;
const updated = previous.map((row) => {
const delivery = row.sourceMessageId
? deliveryByMessage.get(row.sourceMessageId)
: undefined;
const source = row.sourceTurnId ? byTurn.get(row.sourceTurnId) : undefined;
const collaboration = row.sourceMessageId ? collaborationByMessage.get(row.sourceMessageId) : source?.collaboration;
if (JSON.stringify(delivery) === JSON.stringify(row.returnDelivery) && JSON.stringify(collaboration) === JSON.stringify(row.collaboration)) return row;
deliveryChanged = true;
return { ...row, sourceMessageId: row.sourceMessageId ?? source?.message_id,
sourceSessionId: row.sourceSessionId ?? (source?.message_id ? conversationReturnSessionId : undefined),
returnDelivery: delivery, collaboration };
});
if (!fresh.length && !deliveryChanged) return current;
return { ...current, [contextId]: [...updated, ...fresh.map((row) => ({
const updated = reconcileConversationReturns(previous, sessionId, snapshot.messages, (row) => ({
id: managerMessageId.current++, sourceMessageId: row.message_id,
sourceSessionId: conversationReturnSessionId,
sourceSessionId: sessionId,
role: "assistant" as const,
agentLabel: "协作回执",
sourceLabel: "协作回执", text: visibleAgentMessage(row.text), lines: [],
returnDelivery: row.return_delivery,
}))] };
returnDelivery: row.return_delivery, collaboration: row.collaboration,
}));
return updated === previous ? current : { ...current, [contextId]: updated };
});
} catch {
// The durable transcript is retried after reconnection; no model replay.
// Retry this transcript read independently; never replay the model.
} finally {
if (!cancelled) timer = setTimeout(receive, 3000);
if (!cancelled) {
const timer = setTimeout(() => { timers.delete(timer); void receive(sessionId); }, 3000);
timers.add(timer);
}
}
};
void receive();
return () => { cancelled = true; if (timer) clearTimeout(timer); };
}, [readOnly, conversationReturnSessionId, contextId, selectedAgent.label]);
sessionIds.forEach((sessionId) => { void receive(sessionId); });
return () => { cancelled = true; timers.forEach(clearTimeout); };
}, [readOnly, conversationReturnSessionKey, contextId]);

function recordRuntimeBinding(targetContextId: string, binding: PersonalRuntimeBinding | null) {
setRuntimeBindings((current) => {
Expand Down Expand Up @@ -1703,6 +1689,7 @@ function PersonalGoalHome({
const streamingMessageId = appendManagerAssistantMessage(targetContextId, {
activity: ["正在恢复进行中的 Agent 回合"],
sourceTurnId: activeTurnId,
sourceSessionId: created.session_id,
agentLabel: answerIdentityLabel(targetContextId, selectedAgent.label),
lines: [],
pending: true,
Expand Down Expand Up @@ -2154,7 +2141,7 @@ function PersonalGoalHome({
},
onPhase: (_phase: string, turnId: string) => {
submittedTurnId = turnId;
if (streamingMessageId !== null) updateManagerAssistantMessage(targetContextId, streamingMessageId, { sourceTurnId: turnId });
if (streamingMessageId !== null) updateManagerAssistantMessage(targetContextId, streamingMessageId, { sourceTurnId: turnId, sourceSessionId: sessionId });
activeTurnIds.current.set(targetContextId, turnId);
recordRuntimeBinding(targetContextId, {
agentId: selectedRoute.agentId,
Expand Down Expand Up @@ -2191,6 +2178,7 @@ function PersonalGoalHome({
item.turn_id === streamed.turnId && ["agent", "assistant"].includes(item.role));
if (answer) updateManagerAssistantMessage(targetContextId, completedMessageId, {
sourceMessageId: answer.message_id, sourceSessionId: sessionId,
collaboration: answer.collaboration, returnDelivery: answer.return_delivery,
});
}).catch(() => { /* The original conversation remains readable. */ });
}
Expand Down Expand Up @@ -2705,6 +2693,7 @@ function PersonalGoalHome({
const messageId = appendManagerAssistantMessage(run.goalId, {
activity: ["正在把纠偏送入原执行 Session"],
sourceTurnId: turnId,
sourceSessionId: run.sessionId,
agentLabel: run.agentLabel,
lines: [],
pending: true,
Expand Down
29 changes: 22 additions & 7 deletions docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md
Original file line number Diff line number Diff line change
Expand Up @@ -228,7 +228,7 @@ separate; a browser fixture cannot qualify a real attached host.
| Scope | Navigate A→B→A, late response, old subscription terminal event: update the original source/session/Turn only. Full snapshots and delta streams have different merge rules |
| Stream | Duplicate/late events and hydrate overlap preserve one logical answer. A new event does not force scrolling while the user reads history |
| Correction/stop | During tool execution and at completion: actual receiver adopts the latest scope or reports queued/unsupported. Stop targets the original Turn, never its successor or all peers implicitly |
| Result | Missing file retries read only; v1 review cannot certify v2; opening a report is not adoption. Lost return ACK reconciles before another send |
| Result | Replacing the active session must not hide a result owed by an older session in the same conversation. Readback updates only its own session, preserves streamed text and retires old reads after verified return. Missing file retries read only; v1 review cannot certify v2; opening a report is not adoption. Lost return ACK reconciles before another send |
| Attached | Native host offline, stale binding, unsupported steering, next-Turn-only adapter and restart: request remains visible; no guessed success or competing driver |
| Managed | Runtime start failure, quota denial, missing login and stop/restart: effective profile and actual condition readable; no silent model/account substitution |
| Authority | Revoked access or source change rejects stale effects; unrelated permitted branches continue. No private history enters a shared audience |
Expand All @@ -241,12 +241,27 @@ comparison; no measured improvement is claimed by this proposal.

## Delivery boundary

The delivered conversation-entry repair removes all browser free-text action classification in the App and checks
its ordinary Chat path plus explicit scheduling controls. It changes no authority
or stored message schema. Managed/attached conversation continuity, generic TS
inbox extraction and live two-cycle small-team acceptance remain planned until
their own evidence is recorded. The entry repair can roll back as an App routing
change; later persisted-contract migrations need their own compatibility plan.
The delivered conversation-entry repair removes browser free-text action
classification. Shared queue preparation now settles accepted start failures
instead of leaving requests indefinitely queued. Neither change proves a whole
managed or attached journey.

The next qualified frontend slice retains late worker returns after active-session
replacement. The existing Chat snapshot supplies session/message lineage to a
shared TypeScript read model; only sessions still owing a conclusion or delivery
verification remain alongside the active session. No model replay, execution
driver, persisted schema or new inbox owner is introduced. Browser acceptance
covers a replacement session, one lost old-session read, automatic return and
cross-session isolation. Real host execution and receiver adoption retain their
separate acceptance requirements.

Keep work in this order: qualify the installed App's existing-owner-to-original-
conversation journey (GQ02–04, with entry/recovery companions); then G1's two real
small-team cycles (GQ05/GQ11–13, with correction and interruption). Shared report
polish, materials and attention summaries follow; promotional film and scale
follow product evidence. Reuse pending acceptance-recovery and GoalRef work
rather than implement a competing session or inbox lifecycle. Generic TS inbox
extraction remains incremental within those journeys, not their prerequisite.

### Entry behavior compatibility

Expand Down
Loading
Loading