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
4 changes: 2 additions & 2 deletions apps/presentation/dashboard/src/data/chat.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1032,12 +1032,12 @@ export async function interruptChatTurn(sessionId: string, turnId: string) {
}

export async function steerChatTurn(sessionId: string, turnId: string, message: string, ingressId: string) {
const receipt = await requestJson<{ ok: boolean; session_id: string; turn_id: string; client_ingress_id: string; status: string }>(
const receipt = await requestJson<{ ok: boolean; session_id: string; turn_id: string; client_ingress_id: string; status: string; created: boolean }>(
`/api/chat/sessions/${sessionId}/turns/${turnId}/steer`,
{ method: "POST", body: JSON.stringify({ message, client_ingress_id: ingressId }) },
);
if (receipt.ok !== true || receipt.session_id !== sessionId || receipt.turn_id !== turnId
|| receipt.client_ingress_id !== ingressId || receipt.status !== "delivered") {
|| receipt.client_ingress_id !== ingressId || receipt.status !== "delivered" || typeof receipt.created !== "boolean") {
throw new ChatApiError("追加指令的回执不匹配,请保留草稿并检查当前状态。", { error_code: "steer_receipt_mismatch" });
}
return receipt;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -67,4 +67,16 @@ 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");
// The initial optimistic request is the first user message of its Turn.
// Later instructions must retain independent stored identities and remain visible.
const withInstructions = [storedRoles[0],
{ session_id: "current", message_id: "instruction-1", turn_id: "same-turn", role: "user", text: "Chinese first" },
{ session_id: "current", message_id: "instruction-2", turn_id: "same-turn", role: "user", text: "Do not publish" },
storedRoles[1]];
const recoveredInstructions = reconcileConversationHistory(bothRoles, withInstructions, createHistory);
assert.equal(recoveredInstructions.length, 4);
assert.equal(recoveredInstructions.find(row => row.text === "My request").sourceMessageId, "user-message");
assert.equal(recoveredInstructions.find(row => row.text === "Chinese first").sourceMessageId, "instruction-1");
assert.equal(recoveredInstructions.find(row => row.text === "Do not publish").sourceMessageId, "instruction-2");
assert.equal(reconcileConversationHistory(recoveredInstructions, withInstructions, createHistory), recoveredInstructions);
console.log("conversation-returns: passed (session isolation, late return, deduplication, transport uncertainty, stream preservation and watch retirement)");
10 changes: 8 additions & 2 deletions apps/presentation/dashboard/src/data/conversation-returns.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,8 +31,14 @@ 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.origin !== "manager_followup")
.map((row) => [`${row.turn_id}:${row.role === "user" ? "user" : "assistant"}`, row]));
const byTurn = new Map<string, ChatVisibleMessage>();
for (const row of messages) {
if (row.origin === "manager_followup") continue;
const key = `${row.turn_id}:${row.role === "user" ? "user" : "assistant"}`;
// The Turn starts with its user request. Later user instructions have their
// own message IDs; they must not hydrate that original optimistic request.
if (row.role !== "user" || !byTurn.has(key)) byTurn.set(key, row);
}
const seen = new Set(previous.filter((row) => row.sourceSessionId === sessionId).map((row) => row.sourceMessageId));
let changed = false;
const updated = previous.map((row) => {
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
/** Client retry identity only; the Chat ingress/store still owns delivery. */
export type ComposerSteeringRequest = {
sessionId: string;
turnId: string;
text: string;
id: string;
};

const storageKey = "loopx-pw-composer-steering";

export function readComposerSteeringRequests(): Map<string, ComposerSteeringRequest> {
try {
const raw = window.sessionStorage.getItem(storageKey);
const entries: unknown = raw ? JSON.parse(raw) : [];
if (!Array.isArray(entries)) return new Map();
return new Map(entries.flatMap((entry): [string, ComposerSteeringRequest][] => {
if (!Array.isArray(entry) || entry.length !== 2 || typeof entry[0] !== "string") return [];
const request = entry[1];
if (!request || typeof request !== "object" || Array.isArray(request)
|| ![request.sessionId, request.turnId, request.text, request.id]
.every(value => typeof value === "string" && value.length > 0)) return [];
return [[entry[0], { sessionId: request.sessionId, turnId: request.turnId, text: request.text, id: request.id }]];
}));
} catch {
return new Map();
}
}

export function persistComposerSteeringRequests(requests: ReadonlyMap<string, ComposerSteeringRequest>) {
try {
window.sessionStorage.setItem(storageKey, JSON.stringify([...requests]));
} catch {
// As with composer drafts, unavailable browser storage retains only memory.
// No provider action or delivery conclusion follows from this cache.
}
}
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import { goalCreateRequest } from "./goal-create-request";
import { persistComposerSteeringRequests, readComposerSteeringRequests, type ComposerSteeringRequest } from "./composer-steering-recovery";
import type { ConversationHistoryStatus } from "../../data/use-conversation-history";
import { GoalDraftCard } from "./goal-draft-card";
import type { GoalDraft } from "../../../../../../loopx/control_plane/collaboration/goal_draft.js";
Expand Down Expand Up @@ -810,7 +811,8 @@ export function PersonalWorkspacePage({
});
const [sending, setSending] = useState(false);
const [steering, setSteering] = useState(false);
const steeringRequests = useRef(new Map<string, { sessionId: string; turnId: string; text: string; id: string }>());
const [restoredSteeringRequests] = useState(readComposerSteeringRequests);
const steeringRequests = useRef(restoredSteeringRequests);
const [actionDraft, setActionDraft] = useState<WorkspaceActionDraft | null>(null);
const [loopxMode, setLoopxMode] = useState<LoopXModeSnapshot | null>(null);
const [loopxDelivery, setLoopxDelivery] = useState<"queue" | "inbox" | "steer">("queue");
Expand Down Expand Up @@ -873,6 +875,17 @@ export function PersonalWorkspacePage({
function setComposer(value: string) {
setComposerDraft(composerDraftKey, value);
}
function retainSteeringRequest(key: string, request: ComposerSteeringRequest) {
steeringRequests.current.set(key, request);
// Persist before sending: reload after provider acceptance must replay the
// original ingress, never silently submit another Turn.
persistComposerSteeringRequests(steeringRequests.current);
}
function retireSteeringRequest(key: string, id: string) {
if (steeringRequests.current.get(key)?.id !== id) return;
steeringRequests.current.delete(key);
persistComposerSteeringRequests(steeringRequests.current);
}
async function reviewGoalDraft(draft: GoalDraft, edit = false, draftId = "") {
// Source message + reviewed contents survive retry without merging distinct requests.
if (!edit && !draft.question && draft.completion_criteria.trim()) {
Expand Down Expand Up @@ -1641,20 +1654,20 @@ export function PersonalWorkspacePage({
return;
}
const request = retry ?? { sessionId: conversationSessionId, turnId: steeringTurnId!, text: message, id: crypto.randomUUID() };
steeringRequests.current.set(composerDraftKey, request);
retainSteeringRequest(composerDraftKey, request);
setSteering(true);
setActionFeedback(null);
setImageAttachmentError(null);
try {
await callbacks.onSteerConversationTurn(selectedGoalId ?? "manager", request.turnId, message, request.id);
steeringRequests.current.delete(composerDraftKey);
retireSteeringRequest(composerDraftKey, request.id);
if (!messageOverride) setComposerDraft(composerDraftKey, "", composer);
setActionFeedback(locale === "zh-CN" ? "执行器已接收本轮追加指令。" : "The executor accepted instructions for this turn.");
} catch (error) {
// Unknown delivery retries the original Turn even after it completes.
// A confirmed non-delivery may use a new ingress after recovery.
if (error instanceof ChatApiError && error.payload.delivery_state === "not_delivered") {
steeringRequests.current.delete(composerDraftKey);
retireSteeringRequest(composerDraftKey, request.id);
}
setActionFeedback(error instanceof Error ? error.message : t("feedback.sendGenericError"));
} finally { setSteering(false); }
Expand Down
8 changes: 7 additions & 1 deletion apps/presentation/dashboard/src/views/dashboard-page.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -2952,7 +2952,13 @@ function PersonalGoalHome({
if (!binding?.sessionId) throw new Error("当前会话不可用,追加指令未发送,草稿已保留。");
// The service owns exact-turn admission and durable retry. A delivered
// ingress may be read back after completion; never retarget it locally.
await steerChatTurn(binding.sessionId, turnId, message, ingressId);
const receipt = await steerChatTurn(binding.sessionId, turnId, message, ingressId);
if (receipt.created === false) {
// A replay reads an existing delivery; its message belongs to the
// stored transcript, not a second optimistic user bubble.
if (targetContextId === contextId) await conversationHistory.refresh();
return;
}
const id = managerMessageId.current++;
setMessagesByContext(current => {
const messages = current[targetContextId] ?? [];
Expand Down
11 changes: 10 additions & 1 deletion docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md
Original file line number Diff line number Diff line change
Expand Up @@ -133,7 +133,16 @@ receipt → observed work or actionable failure → readable answer in the same
idle composer and initial presentation are unchanged. This changes the former
managed-running Send lockout; LoopX mode retains its explicit delivery choice.
A lost or mismatched receipt preserves draft and ingress identity, including
retry after completion. Only a confirmed non-delivery permits a new ingress.
retry after completion or page reload in the same browser tab. The shared
composer persists the original Session/Turn/text/ingress before dispatch and
reuses it only for the matching Session and text; restoring the cache never
sends automatically. Only a confirmed non-delivery permits a new ingress.
Replaying a delivered receipt reads the current stored transcript rather than
reusing an earlier snapshot or inserting another local copy of the instruction.
History keeps the initial user request and later instructions as distinct messages
even while the initial request is still being reconciled with its stored identity.
Unavailable browser storage keeps current-page retry behavior but cannot
promise reload recovery. This tab-local cache is not delivery authority.
Acceptance means the executor received the instructions, not that it adopted
them or that delegated/team work stopped. Live adoption stays a release gate.
- The compact receipt and full conversation offer the same controls. Failure
Expand Down
2 changes: 1 addition & 1 deletion examples/personal-workspace-browser/chat-recovery.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -529,7 +529,7 @@ export const chatRecoveryScenario = {
const body = route.request().postDataJSON();
const [sessionId, turnId] = new URL(route.request().url()).pathname.match(/sessions\/([^/]+)\/turns\/([^/]+)\/steer/).slice(1);
steers.push({ sessionId, turnId });
await route.fulfill({ json: { ok: true, session_id: sessionId, turn_id: turnId, client_ingress_id: body.client_ingress_id, status: "delivered" } });
await route.fulfill({ json: { ok: true, session_id: sessionId, turn_id: turnId, client_ingress_id: body.client_ingress_id, status: "delivered", created: true } });
});
const interruptsBefore = api.interrupts.length;
await composerInput.fill(draft);
Expand Down
Loading
Loading