diff --git a/docs/web/control-ui.md b/docs/web/control-ui.md index a7ddef9d38fa..f436bee7b369 100644 --- a/docs/web/control-ui.md +++ b/docs/web/control-ui.md @@ -43,6 +43,12 @@ Hover the row or focus it with the keyboard for a tooltip explaining the exact status. Reduced motion keeps the claw still. Tasks without a display title keep the generic **Subagent** label. Select a row to open its details. +Dragging a session between sidebar groups updates its placement immediately. A successful +save keeps that placement even if the subsequent list refresh fails; the UI reports +the refresh error separately. If a connection failure leaves the save unconfirmed, +refresh and check the session's group before retrying. Other clients' newer group +changes still reconcile through session events. + The sidebar keeps unread child failures visible on their ancestors. These warnings name the child session that failed, even when its parent has finished or continues working. Select the warning to open the child session and acknowledge its failure; diff --git a/src/gateway/server-methods/session-change-event.test.ts b/src/gateway/server-methods/session-change-event.test.ts index 61957599c2c9..7b6581398ad9 100644 --- a/src/gateway/server-methods/session-change-event.test.ts +++ b/src/gateway/server-methods/session-change-event.test.ts @@ -13,6 +13,7 @@ import { clearAgentRunContext, registerAgentRunContext, } from "../../infra/agent-run-registry.js"; +import { sessionChanges } from "../../sessions/session-row-changes.js"; import type { ChatAbortControllerEntry } from "../chat-abort.js"; import { readGatewayAccessRevision } from "../gateway-access-revision.js"; import { @@ -154,6 +155,34 @@ afterEach(async () => { }); describe("sessions.changed coalescing", () => { + it("publishes catalog-only changes without invalidating session projections or access", async () => { + const context = createContext(); + const changed = vi.fn(); + const unsubscribe = sessionChanges.subscribe(changed); + onTestFinished(unsubscribe); + const initialAccessRevision = readGatewayAccessRevision(); + + await emitAndSettleLeading(context, { reason: "groups" }, { catalogOnly: true }); + + expect(changed).not.toHaveBeenCalled(); + expect(mocks.invalidate).not.toHaveBeenCalled(); + expect(context.mentionInbox?.invalidate).not.toHaveBeenCalled(); + expect(readGatewayAccessRevision()).toBe(initialAccessRevision); + expect(context.broadcastToConnIds).toHaveBeenCalledWith( + "sessions.changed", + expect.objectContaining({ reason: "groups" }), + expect.any(Set), + expect.any(Object), + ); + + // Rename/delete use the same public reason but can change member rows. + await emitAndSettleLeading(context, { reason: "groups" }); + expect(changed).toHaveBeenCalledWith({ all: true, scope: "sessions" }); + expect(mocks.invalidate).toHaveBeenCalledOnce(); + expect(context.mentionInbox?.invalidate).toHaveBeenCalledOnce(); + expect(readGatewayAccessRevision()).toBe(initialAccessRevision + 1); + }); + it("publishes the latest placement through coalesced unrelated mutations and clears it explicitly", async () => { const context = createContext(); const sessionKey = "agent:main:cloud"; diff --git a/src/gateway/server-methods/session-change-event.ts b/src/gateway/server-methods/session-change-event.ts index 708ce7091eea..066185c5f147 100644 --- a/src/gateway/server-methods/session-change-event.ts +++ b/src/gateway/server-methods/session-change-event.ts @@ -203,9 +203,12 @@ export async function flushPendingSessionsChangedEvents(context?: object): Promi export function emitSessionsChanged( context: SessionChangeContext, payload: SessionChangedPayload, - options: { accessChanged?: boolean; preparedPublication?: boolean } = {}, + options: { accessChanged?: boolean; preparedPublication?: boolean; catalogOnly?: boolean } = {}, ): void { - if (!options.preparedPublication) { + // Catalog absorption changes no session facts. Rename/delete callers retain + // normal invalidation because their sweeps can have committed member changes. + const catalogOnly = options.catalogOnly && payload.reason === "groups" && !payload.sessionKey; + if (!options.preparedPublication && !catalogOnly) { sessionChanges.emit( payload.sessionKey ? { @@ -216,12 +219,14 @@ export function emitSessionsChanged( ); } // Only a committed producer may certify unchanged access; unknown changes stay conservative. - if (options.accessChanged !== false) { + if (!catalogOnly && options.accessChanged !== false) { bumpGatewayAccessRevision(); } - invalidateSessionSharingSnapshot(payload.sessionKey); - // Inbox subscriptions are independent of session-list subscriptions, including a closed sidebar. - context.mentionInbox?.invalidate(); + if (!catalogOnly) { + invalidateSessionSharingSnapshot(payload.sessionKey); + // Inbox subscriptions are independent of session-list subscriptions, including a closed sidebar. + context.mentionInbox?.invalidate(); + } const connIds = context.getSessionEventSubscriberConnIds(); if (!hasSessionChangeReceivers(connIds)) { return; diff --git a/src/gateway/server-methods/sessions-patch-effects.test.ts b/src/gateway/server-methods/sessions-patch-effects.test.ts new file mode 100644 index 000000000000..480880338790 --- /dev/null +++ b/src/gateway/server-methods/sessions-patch-effects.test.ts @@ -0,0 +1,89 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; +import { disableCronJobsBoundToSessions } from "../../cron/job-session-bindings.js"; +import { ensureSessionGroupRegistered } from "../session-groups.js"; +import { emitSessionsChanged } from "./session-change-event.js"; +import { publishSessionPatchEffects } from "./sessions-patch-effects.js"; +import { sessionLog } from "./sessions-shared.js"; +import type { GatewayRequestContext } from "./types.js"; + +vi.mock("../../cron/job-session-bindings.js", () => ({ disableCronJobsBoundToSessions: vi.fn() })); +vi.mock("../session-groups.js", () => ({ ensureSessionGroupRegistered: vi.fn() })); +vi.mock("../session-patch-hooks.js", () => ({ triggerSessionPatchHook: vi.fn() })); +vi.mock("./session-change-event.js", () => ({ emitSessionsChanged: vi.fn() })); +vi.mock("./sessions-patch-model-selection.js", () => ({ + persistSessionPatchModelSelection: vi.fn(), +})); +vi.mock("./sessions-shared.js", () => ({ sessionLog: { warn: vi.fn(), info: vi.fn() } })); + +beforeEach(() => { + vi.resetAllMocks(); + vi.mocked(disableCronJobsBoundToSessions).mockResolvedValue(new Map()); +}); + +describe("committed category patch effects", () => { + function params(): Parameters[0] { + return { + cfg: {}, + context: { cron: {} } as GatewayRequestContext, + callerScopes: [], + callerCanManageCron: true, + category: "Travel", + targets: [ + { + accessChanged: false, + entry: { sessionId: "saved", updatedAt: 1, category: "Travel", archivedAt: 1 }, + target: { + canonicalKey: "agent:main:travel", + targetAgentId: "main", + fullPatch: { key: "agent:main:travel", category: "Travel", archived: true }, + }, + }, + ], + }; + } + + it("preserves the committed patch and remaining effects when catalog registration fails", async () => { + vi.mocked(ensureSessionGroupRegistered).mockImplementationOnce(() => { + throw new Error("catalog unavailable"); + }); + const patch = params(); + + await expect(publishSessionPatchEffects(patch)).resolves.toBeUndefined(); + + expect(patch.targets[0]?.entry.category).toBe("Travel"); + expect(emitSessionsChanged).toHaveBeenCalledWith( + patch.context, + { sessionKey: "agent:main:travel", reason: "patch" }, + { accessChanged: false }, + ); + expect(emitSessionsChanged).toHaveBeenCalledWith( + patch.context, + { reason: "groups" }, + { catalogOnly: true }, + ); + expect(sessionLog.warn).toHaveBeenCalledWith( + expect.stringContaining("retry the same category assignment"), + ); + expect(disableCronJobsBoundToSessions).toHaveBeenCalledOnce(); + + // A repeated assignment still invokes the catalog owner and publishes recovery. + vi.mocked(ensureSessionGroupRegistered).mockReturnValueOnce(true); + await publishSessionPatchEffects(patch); + expect(ensureSessionGroupRegistered).toHaveBeenCalledTimes(2); + expect( + vi.mocked(emitSessionsChanged).mock.calls.filter(([, event]) => event.reason === "groups"), + ).toHaveLength(2); + }); + + it("does not publish a catalog change when the group already exists", async () => { + vi.mocked(ensureSessionGroupRegistered).mockReturnValue(false); + await publishSessionPatchEffects(params()); + expect(emitSessionsChanged).toHaveBeenCalledOnce(); + }); + + it("does not register a category when no target committed", async () => { + await publishSessionPatchEffects({ ...params(), targets: [] }); + expect(ensureSessionGroupRegistered).not.toHaveBeenCalled(); + expect(emitSessionsChanged).not.toHaveBeenCalled(); + }); +}); diff --git a/src/gateway/server-methods/sessions-patch-effects.ts b/src/gateway/server-methods/sessions-patch-effects.ts index 2852f1e15aa9..f230e07be476 100644 --- a/src/gateway/server-methods/sessions-patch-effects.ts +++ b/src/gateway/server-methods/sessions-patch-effects.ts @@ -72,8 +72,21 @@ export async function publishSessionPatchEffects(params: { if (params.targets.length > 0 && typeof category === "string" && category.trim()) { // A first-use category is a group-catalog mutation: clients reload the // catalog only on reason "groups" (the sessions.groups.* siblings emit it). - if (ensureSessionGroupRegistered(category)) { - emitSessionsChanged(params.context, { reason: "groups" }, { accessChanged: false }); + let catalogChanged: boolean; + try { + catalogChanged = ensureSessionGroupRegistered(category); + } catch (error) { + // The session category is already durable. Preserve that outcome and the + // existing same-category patch recovery instead of asking clients to undo it. + sessionLog.warn( + `sessions.patch: category ${JSON.stringify(category)} was saved, but group registration failed; retry the same category assignment to repair the catalog: ${formatErrorMessage(error)}`, + ); + // Registration may have committed before cleanup failed. Reload the catalog + // on uncertain outcomes too, without invalidating unrelated session rows. + catalogChanged = true; + } + if (catalogChanged) { + emitSessionsChanged(params.context, { reason: "groups" }, { catalogOnly: true }); } } if (params.callerCanManageCron && archivedSessionKeys.size > 0) { diff --git a/src/gateway/session-groups.test.ts b/src/gateway/session-groups.test.ts index 87130288167f..f50254890867 100644 --- a/src/gateway/session-groups.test.ts +++ b/src/gateway/session-groups.test.ts @@ -12,6 +12,7 @@ import { closeOpenClawAgentDatabasesForTest, runOpenClawAgentWriteTransaction, } from "../state/openclaw-agent-db.js"; +import * as stateDatabase from "../state/openclaw-state-db.js"; import { closeOpenClawStateDatabaseForTest, openOpenClawStateDatabase, @@ -41,6 +42,7 @@ describe("session groups catalog", () => { }); afterEach(async () => { + vi.restoreAllMocks(); closeOpenClawAgentDatabasesForTest(); closeOpenClawStateDatabaseForTest(); await fs.rm(root, { recursive: true, force: true }); @@ -287,6 +289,43 @@ describe("session groups catalog", () => { ]); }); + it("does not admit a write transaction for an existing normalized category", () => { + putSessionGroups({ cfg, names: ["Work"], env }); + const transaction = vi.spyOn(stateDatabase, "runOpenClawStateWriteTransaction"); + + expect(ensureSessionGroupRegistered(" Work ", env)).toBe(false); + + expect(transaction).not.toHaveBeenCalled(); + expect(listSessionGroups(env)).toEqual([{ name: "Work", position: 0 }]); + }); + + it("rechecks a missing category after another writer registers it", () => { + putSessionGroups({ cfg, names: ["Work"], env }); + const originalTransaction = stateDatabase.runOpenClawStateWriteTransaction; + vi.spyOn(stateDatabase, "runOpenClawStateWriteTransaction").mockImplementationOnce( + (operation, options, transactionOptions) => { + // Commit a competing registration between the optimistic read and admission. + originalTransaction( + ({ db }) => { + db.prepare( + "INSERT INTO session_groups (name, position, created_at) VALUES (?, ?, ?)", + ).run("Travel", 1, 123); + }, + { env }, + ); + return originalTransaction(operation, options, transactionOptions); + }, + ); + + expect(ensureSessionGroupRegistered("Travel", env)).toBe(false); + expect(listSessionGroups(env)).toEqual([ + { name: "Work", position: 0 }, + { name: "Travel", position: 1 }, + ]); + expect(ensureSessionGroupRegistered("Later", env)).toBe(true); + expect(listSessionGroups(env).at(-1)).toEqual({ name: "Later", position: 2 }); + }); + it("renames a group and repoints member categories without bumping updatedAt", async () => { putSessionGroups({ cfg, diff --git a/src/gateway/session-groups.ts b/src/gateway/session-groups.ts index 85d1c427f125..5ece07cecb97 100644 --- a/src/gateway/session-groups.ts +++ b/src/gateway/session-groups.ts @@ -293,6 +293,21 @@ export function ensureSessionGroupRegistered( if (!normalized) { return false; } + // Existing categories need no writer admission. A missing name is only a + // hint: another writer can register it before our transaction is admitted. + const readDb = dbFor(env); + if ( + executeSqliteQuerySync( + readDb, + kyselyFor(readDb) + .selectFrom("session_groups") + .select("name") + .where("name", "=", normalized) + .limit(1), + ).rows[0] + ) { + return false; + } let inserted = false; runOpenClawStateWriteTransaction( ({ db }) => { diff --git a/ui/src/api/types.ts b/ui/src/api/types.ts index ffbc31c34fb6..0d777fd0e83b 100644 --- a/ui/src/api/types.ts +++ b/ui/src/api/types.ts @@ -279,6 +279,7 @@ export type SessionsBranchesSwitchResult = export type SessionsPatchResult = SessionsPatchResultBase<{ sessionId: string; + category?: GatewaySessionRow["category"]; updatedAt?: number; createdAt?: number; pinnedAt?: number; diff --git a/ui/src/i18n/locales/en.ts b/ui/src/i18n/locales/en.ts index 80412f712736..39fd511a37e1 100644 --- a/ui/src/i18n/locales/en.ts +++ b/ui/src/i18n/locales/en.ts @@ -3224,6 +3224,9 @@ export const en: TranslationMap & { actionsUnavailable: "Actions are unavailable while the Gateway reconnects.", settingsChangesUnavailable: "Changes to settings are disabled while the Gateway is reconnecting.", + sessionMoveRefreshFailed: "The session move was saved, but refreshing the list failed: {error}", + sessionMoveUncertain: + "The session move could not be confirmed. Refresh and check its group before retrying. {error}", sessionOperationCompletedPreviousConnection: "The session operation completed on the previous connection. Check the current session list before continuing.", sessionOperationCompletedPreviousConnectionWithRefreshError: diff --git a/ui/src/lib/sessions/index.ts b/ui/src/lib/sessions/index.ts index 90110d72fb2f..96816a1c6a3c 100644 --- a/ui/src/lib/sessions/index.ts +++ b/ui/src/lib/sessions/index.ts @@ -247,7 +247,7 @@ export function createSessionCapability( if (source) { mutations.observePendingFields( source.row, - source.select(row, ["pinned", "pinnedAt", "unread"]), + source.select(row, ["pinned", "pinnedAt", "unread", "category"]), agentId, ); } diff --git a/ui/src/lib/sessions/session-category-mutations.test.ts b/ui/src/lib/sessions/session-category-mutations.test.ts new file mode 100644 index 000000000000..82e99d750fa9 --- /dev/null +++ b/ui/src/lib/sessions/session-category-mutations.test.ts @@ -0,0 +1,377 @@ +// @vitest-environment node +import { describe, expect, it, vi } from "vitest"; +import { createDeferred } from "../../../../test/helpers/promise.js"; +import { GatewayRequestError } from "../../api/gateway.ts"; +import type { GatewaySessionRow } from "../../api/types.ts"; +import { createTestGatewayClient } from "../../test-helpers/gateway-client.ts"; +import { + createGatewayHarness, + createTestSessionCapability, + sessionsResult, +} from "./session-capability.test-support.ts"; + +const initial: GatewaySessionRow = { + key: "agent:main:category-test", + sessionId: "category-incarnation", + kind: "direct", + updatedAt: 1, + category: "Alpha", + pinned: false, +}; + +function setup() { + const replies = [createDeferred(), createDeferred()]; + const list = createDeferred>(); + const listStarted = createDeferred(); + let reads = 0; + let writes = 0; + const client = createTestGatewayClient(async (method) => { + if (method === "sessions.subscribe") { + return { subscribed: true }; + } + if (method === "sessions.patch") { + return replies[writes++]!.promise; + } + if (method === "sessions.list") { + if (++reads === 1) { + return sessionsResult([initial], 1); + } + listStarted.resolve(); + return list.promise; + } + throw new Error(`Unexpected method: ${method}`); + }); + const harness = createGatewayHarness(client); + const sessions = createTestSessionCapability(harness.gateway); + const confirm = (index: number, category: string | undefined, updatedAt = index + 2) => + replies[index]!.resolve({ + ok: true, + key: initial.key, + path: "", + entry: { ...initial, category, updatedAt }, + }); + return { + ...harness, + sessions, + replies, + list, + listStarted, + confirm, + category: () => sessions.state.result?.sessions[0]?.category, + }; +} + +const options = { agentId: "main", expectedSessionId: initial.sessionId }; + +describe("session category mutations", () => { + it.each(["Beta", null])( + "projects %s immediately and preserves its receipt when refresh fails", + async (category) => { + const h = setup(); + try { + await h.sessions.refresh({ agentId: "main", force: true }); + const pending = h.sessions.patch(initial.key, { category }, options); + expect(h.category()).toBe(category ?? undefined); + h.confirm(0, category ?? undefined); + await h.listStarted.promise; + expect(h.category()).toBe(category ?? undefined); + await expect(pending).resolves.toMatchObject({ ok: true }); + h.list.reject(new Error("injected list failure")); + await h.list.promise.catch(() => undefined); + await new Promise((resolve) => { + setImmediate(resolve); + }); + expect(h.category()).toBe(category ?? undefined); + expect(h.sessions.state.error).toContain("The session move was saved"); + expect(h.sessions.state.error).toContain("injected list failure"); + } finally { + h.sessions.dispose(); + } + }, + ); + + it("keeps the latest A→B→A intent through out-of-order receipts and stale events", async () => { + const h = setup(); + try { + await h.sessions.refresh({ agentId: "main", force: true }); + const first = h.sessions.patch( + initial.key, + { category: "Beta" }, + { ...options, deferListRefresh: true }, + ); + const second = h.sessions.patch( + initial.key, + { category: "Alpha" }, + { ...options, deferListRefresh: true }, + ); + expect(h.category()).toBe("Alpha"); + h.confirm(1, "Alpha", 3); + await second; + h.confirm(0, "Beta", 2); + await first; + expect(h.category()).toBe("Alpha"); + h.emitEvent({ + type: "event", + event: "sessions.changed", + payload: { + ...initial, + sessionKey: initial.key, + reason: "patch", + category: "Beta", + updatedAt: 2, + archived: false, + session: { ...initial, category: "Beta", updatedAt: 2, archived: false }, + }, + }); + expect(h.category()).toBe("Alpha"); + } finally { + h.sessions.dispose(); + } + }); + + it("rolls back a rejected newest move to the earlier confirmed category", async () => { + const h = setup(); + try { + await h.sessions.refresh({ agentId: "main", force: true }); + const first = h.sessions.patch( + initial.key, + { category: "Beta" }, + { ...options, deferListRefresh: true }, + ); + const second = h.sessions.patch( + initial.key, + { category: null }, + { ...options, deferListRefresh: true }, + ); + h.confirm(0, "Beta"); + await first; + expect(h.category()).toBeUndefined(); + h.replies[1]!.reject( + new GatewayRequestError({ code: "INVALID_REQUEST", message: "move rejected" }), + ); + await expect(second).rejects.toThrow("move rejected"); + expect(h.category()).toBe("Beta"); + } finally { + h.sessions.dispose(); + } + }); + it("does not roll back an uncertain transport failure and never retries the write", async () => { + const h = setup(); + try { + await h.sessions.refresh({ agentId: "main", force: true }); + const pending = h.sessions.patch(initial.key, { category: "Beta" }, options); + h.replies[0]!.reject(new Error("connection lost")); + await expect(pending).rejects.toThrow("could not be confirmed"); + expect(h.category()).toBe("Beta"); + await h.listStarted.promise; + h.list.resolve(sessionsResult([{ ...initial, category: "Beta", updatedAt: 2 }], 2)); + await h.list.promise; + } finally { + h.sessions.dispose(); + } + }); + + it("keeps category and pin receipts together without waiting for the list", async () => { + const h = setup(); + try { + await h.sessions.refresh({ agentId: "main", force: true }); + const pending = h.sessions.patch(initial.key, { category: null, pinned: true }, options); + expect(h.sessions.state.result?.sessions[0]).toMatchObject({ pinned: true }); + expect(h.category()).toBeUndefined(); + h.replies[0]!.resolve({ + ok: true, + path: "", + key: initial.key, + entry: { ...initial, category: undefined, pinnedAt: 7, updatedAt: 7 }, + }); + await pending; + expect(h.sessions.state.result?.sessions[0]).toMatchObject({ pinned: true, pinnedAt: 7 }); + h.list.resolve( + sessionsResult( + [{ ...initial, category: undefined, pinned: true, pinnedAt: 7, updatedAt: 7 }], + 7, + ), + ); + } finally { + h.sessions.dispose(); + } + }); + + it("does not restore an old intent over a newer pending move", async () => { + const h = setup(); + try { + await h.sessions.refresh({ agentId: "main", force: true }); + const first = h.sessions.patch( + initial.key, + { category: "Beta" }, + { ...options, deferListRefresh: true }, + ); + const second = h.sessions.patch( + initial.key, + { category: "Gamma" }, + { ...options, deferListRefresh: true }, + ); + h.replies[0]!.reject(new GatewayRequestError({ code: "FORBIDDEN", message: "first denied" })); + await expect(first).rejects.toThrow("first denied"); + expect(h.category()).toBe("Gamma"); + h.confirm(1, "Gamma"); + await second; + expect(h.category()).toBe("Gamma"); + } finally { + h.sessions.dispose(); + } + }); + + it("does not project a late receipt into a replacement session", async () => { + const h = setup(); + try { + await h.sessions.refresh({ agentId: "main", force: true }); + const pending = h.sessions.patch( + initial.key, + { category: "Beta" }, + { ...options, deferListRefresh: true }, + ); + const reading = h.sessions.refresh({ agentId: "main", force: true }); + h.list.resolve( + sessionsResult( + [{ ...initial, sessionId: "replacement", category: "Replacement", updatedAt: 9 }], + 9, + ), + ); + await reading; + h.confirm(0, "Beta"); + await pending; + expect(h.category()).toBe("Replacement"); + } finally { + h.sessions.dispose(); + } + }); + + it("retires placement intent when the connection is replaced", async () => { + const h = setup(); + try { + await h.sessions.refresh({ agentId: "main", force: true }); + const pending = h.sessions.patch( + initial.key, + { category: "Beta" }, + { ...options, deferListRefresh: true }, + ); + h.publish(false, null); + h.confirm(0, "Beta"); + await expect(pending).resolves.toBeNull(); + expect(h.sessions.state.error).toBeNull(); + } finally { + h.sessions.dispose(); + } + }); + it("keeps a confirmed category ahead of a read started before the write", async () => { + const h = setup(); + try { + await h.sessions.refresh({ agentId: "main", force: true }); + const reading = h.sessions.refresh({ agentId: "main", force: true }); + const pending = h.sessions.patch( + initial.key, + { category: "Beta" }, + { ...options, deferListRefresh: true }, + ); + h.confirm(0, "Beta"); + await pending; + h.list.resolve(sessionsResult([{ ...initial }], 1)); + await reading; + expect(h.category()).toBe("Beta"); + } finally { + h.sessions.dispose(); + } + }); + + it("continues admitting newer authoritative category events", async () => { + const h = setup(); + try { + await h.sessions.refresh({ agentId: "main", force: true }); + const pending = h.sessions.patch( + initial.key, + { category: "Beta" }, + { ...options, deferListRefresh: true }, + ); + h.confirm(0, "Beta"); + await pending; + h.emitEvent({ + type: "event", + event: "sessions.changed", + payload: { + ...initial, + sessionKey: initial.key, + reason: "patch", + category: "External", + updatedAt: 9, + archived: false, + session: { ...initial, category: "External", updatedAt: 9, archived: false }, + }, + }); + expect(h.category()).toBe("External"); + } finally { + h.sessions.dispose(); + } + }); +}); + +it.each(["different", "returned"] as const)( + "does not publish an old scoped category refresh failure into the %s foreground", + async (selection) => { + const reply = createDeferred(); + const oldRead = createDeferred>(); + let committed = false; + let scopedReads = 0; + const writer = { + ...initial, + key: "agent:writer:main", + sessionId: "writer", + category: "Writer", + }; + const client = createTestGatewayClient(async (method, raw) => { + if (method === "sessions.subscribe") { + return { subscribed: true }; + } + if (method === "sessions.patch") { + return reply.promise; + } + if (method === "sessions.list") { + if ((raw as { agentId?: string }).agentId === "writer") { + return sessionsResult([writer], 3); + } + if (committed && ++scopedReads === 1) { + return oldRead.promise; + } + return sessionsResult([{ ...initial, category: committed ? "External" : "Alpha" }], 4); + } + throw new Error(`Unexpected method: ${method}`); + }); + const h = createGatewayHarness(client); + const sessions = createTestSessionCapability(h.gateway); + try { + await sessions.refresh({ agentId: "main", force: true }); + const move = sessions.patch(initial.key, { category: "Beta" }, options); + await sessions.refresh({ agentId: "writer", force: true }); + committed = true; + reply.resolve({ + ok: true, + key: initial.key, + entry: { ...initial, category: "Beta", updatedAt: 2 }, + }); + await move; + await vi.waitFor(() => expect(scopedReads).toBe(1)); + if (selection === "returned") { + await sessions.refresh({ agentId: "main", force: true }); + } + oldRead.reject(new Error("old agent read failed")); + // Let the completed read propagate through the refresh promise chain. + await new Promise((resolve) => { + setImmediate(resolve); + }); + expect(sessions.state.agentId).toBe(selection === "returned" ? "main" : "writer"); + expect(sessions.state.error).toBeNull(); + } finally { + sessions.dispose(); + } + }, +); diff --git a/ui/src/lib/sessions/session-mutation-refresh.ts b/ui/src/lib/sessions/session-mutation-refresh.ts new file mode 100644 index 000000000000..bb1b0dbaddbf --- /dev/null +++ b/ui/src/lib/sessions/session-mutation-refresh.ts @@ -0,0 +1,109 @@ +import { ErrorCodes, GatewayProtocolRequestError } from "@openclaw/gateway-client/browser"; +import { t } from "../../i18n/index.ts"; +import { formatUiError } from "../format-error.ts"; +import type { + SessionConnectionOwner, + SessionConnectionScope, + SessionRefreshOutcome, + SessionState, +} from "./session-capability.ts"; +import { sessionListAgentMatcher } from "./session-list-query.ts"; + +type Host = { + connection: SessionConnectionOwner; + readState: () => SessionState; + publish: (state: SessionState, errorSource?: "session-observer" | "operation") => void; + reconcileMutation: ( + agentId?: string | null, + isErrorCurrent?: () => boolean, + ) => Promise; +}; + +/** Reconcile committed mutations without confusing refresh failures with failed writes. */ +export function createSessionMutationRefresh(host: Host) { + const reconcileConfirmedPreviousConnection = async ( + scope: SessionConnectionScope, + agentId?: string | null, + ): Promise => { + const replacement = host.connection.capture(); + if (!replacement || replacement.client !== scope.client) { + return false; + } + let refreshError: string | undefined; + try { + const outcome = await host.reconcileMutation(agentId); + refreshError = outcome.status === "failed" ? outcome.error : undefined; + } catch (error) { + refreshError = formatUiError(error); + } + if (!host.connection.isCurrent(replacement)) { + return false; + } + host.publish( + { + ...host.readState(), + error: refreshError + ? t("connection.sessionOperationCompletedPreviousConnectionWithRefreshError", { + error: refreshError, + }) + : t("connection.sessionOperationCompletedPreviousConnection"), + }, + "operation", + ); + return true; + }; + + // Placement receipts complete immediately; the roster owner reconciles asynchronously. + const refreshCategory = (scope: SessionConnectionScope, agentId?: string | null): void => { + // A scoped read cannot acquire error ownership by later selecting its agent. + // A replacement foreground snapshot retires this warning, not row reconciliation. + const foreground = host.readState(); + const ownsForeground = sessionListAgentMatcher(agentId)(foreground.agentId ?? undefined); + const isErrorCurrent = () => { + const current = host.readState(); + return ( + ownsForeground && + host.connection.isCurrent(scope) && + current.agentId === foreground.agentId && + current.result === foreground.result && + current.resultCached === foreground.resultCached + ); + }; + const report = (error: string) => { + if (isErrorCurrent()) { + host.publish( + { ...host.readState(), error: t("connection.sessionMoveRefreshFailed", { error }) }, + "operation", + ); + } + }; + void host + .reconcileMutation(agentId, isErrorCurrent) + .then((outcome) => { + if (outcome.status === "failed") { + report(outcome.error); + } + }) + .catch((error: unknown) => report(formatUiError(error))); + }; + // Reconcile the original owner after transport loss without retrying its write. + const reportUncertainCategory = (error: unknown, agentId?: string | null): Error => { + void host.reconcileMutation(agentId).catch(() => undefined); + const uncertainty = new Error( + t("connection.sessionMoveUncertain", { error: formatUiError(error) }), + ); + host.publish({ ...host.readState(), error: uncertainty.message }, "operation"); + return uncertainty; + }; + return { reconcileConfirmedPreviousConnection, refreshCategory, reportUncertainCategory }; +} + +/** Only definitive protocol rejection proves that the requested mutation did not commit. */ +export function isRejectedSessionMutation(error: unknown): boolean { + return ( + error instanceof GatewayProtocolRequestError && + (error.gatewayCode === ErrorCodes.INVALID_REQUEST || + error.gatewayCode === ErrorCodes.FORBIDDEN || + error.gatewayCode === ErrorCodes.APPROVAL_NOT_FOUND) + ); +} diff --git a/ui/src/lib/sessions/session-mutation-selection.test.ts b/ui/src/lib/sessions/session-mutation-selection.test.ts index 60c89507404a..52484e6a7386 100644 --- a/ui/src/lib/sessions/session-mutation-selection.test.ts +++ b/ui/src/lib/sessions/session-mutation-selection.test.ts @@ -9,7 +9,7 @@ import { sessionsResult, } from "./session-capability.test-support.ts"; -it.each(["rename", "archive", "delete"] as const)( +it.each(["rename", "archive", "delete", "category"] as const)( "%s completion preserves loaded and queued foreground queries and reconciles affected managed lists", async (operation) => { for (const foreground of ["loaded", "queued", "global"] as const) { @@ -49,6 +49,7 @@ it.each(["rename", "archive", "delete"] as const)( ...row, label: committed && operation === "rename" ? "Renamed" : row.label, archived: committed && operation === "archive", + category: committed && operation === "category" ? "Moved" : undefined, }, ]; return sessionsResult(params.agentId === "writer" ? [writer] : mainRows, committed ? 2 : 1); @@ -69,7 +70,11 @@ it.each(["rename", "archive", "delete"] as const)( ? sessions.delete(row.key, options) : sessions.patch( row.key, - operation === "archive" ? { archived: true } : { label: "Renamed" }, + operation === "archive" + ? { archived: true } + : operation === "category" + ? { category: "Moved" } + : { label: "Renamed" }, options, ); if (foreground === "queued") { @@ -90,7 +95,12 @@ it.each(["rename", "archive", "delete"] as const)( : { ok: true, key: row.key, - entry: { ...row, label: "Renamed", archived: operation === "archive" }, + entry: { + ...row, + label: "Renamed", + archived: operation === "archive", + category: operation === "category" ? "Moved" : undefined, + }, }, ); if (foreground === "queued") { @@ -112,7 +122,11 @@ it.each(["rename", "archive", "delete"] as const)( expect(affected).toEqual([]); } else { expect(affected[0]).toMatchObject( - operation === "rename" ? { label: "Renamed" } : { archived: true }, + operation === "rename" + ? { label: "Renamed" } + : operation === "category" + ? { category: "Moved" } + : { archived: true }, ); } expect( diff --git a/ui/src/lib/sessions/session-mutations.ts b/ui/src/lib/sessions/session-mutations.ts index d77cb7b65823..5e5b7e4dfc80 100644 --- a/ui/src/lib/sessions/session-mutations.ts +++ b/ui/src/lib/sessions/session-mutations.ts @@ -4,7 +4,6 @@ import type { SessionsAssignOwnerResult, } from "../../../../packages/gateway-protocol/src/index.js"; import type { GatewaySessionRow, SessionsListResult } from "../../api/types.ts"; -import { t } from "../../i18n/index.ts"; import { formatUiError } from "../format-error.ts"; import { requestSessionCreate, @@ -18,7 +17,6 @@ import { createSessionArchiveState, projectSessionArchiveFields } from "./sessio import type { SessionCapability, SessionConnectionOwner, - SessionConnectionScope, SessionCreateReconciliation, SessionRefreshOutcome, SessionResetOptions, @@ -26,6 +24,10 @@ import type { SessionState, } from "./session-capability.ts"; import { areUiSessionKeysEquivalent } from "./session-key.ts"; +import { + createSessionMutationRefresh, + isRejectedSessionMutation, +} from "./session-mutation-refresh.ts"; import { projectSessionPatchRowFields } from "./session-patch-row-facts.ts"; import { createOptimisticRowPatches, @@ -168,6 +170,11 @@ export function createSessionMutations(host: SessionMutationsHost) { pinnedAt: names.includes("pinnedAt") ? row.pinnedAt : previous.pinnedAt, }), }); + const optimisticCategories = createOptimisticRowPatches(host, { + read: (row) => row.category, + write: (row, category) => (row.category === category ? row : host.copyRow(row, { category })), + observe: (previous, row, names) => (names.includes("category") ? row.category : previous), + }); const optimisticUnread = createOptimisticRowPatches(host, { read: (row) => row.unread, write: (row, unread) => (row.unread === unread ? row : host.copyRow(row, { unread })), @@ -177,7 +184,10 @@ export function createSessionMutations(host: SessionMutationsHost) { row: GatewaySessionRow, sourceAgentId?: string | null, ): GatewaySessionRow => - optimisticUnread.applyRow(optimisticPins.applyRow(row, sourceAgentId), sourceAgentId); + optimisticCategories.applyRow( + optimisticUnread.applyRow(optimisticPins.applyRow(row, sourceAgentId), sourceAgentId), + sourceAgentId, + ); const retireModelOverride = (key: string) => { const normalizedKey = key.trim(); @@ -188,37 +198,8 @@ export function createSessionMutations(host: SessionMutationsHost) { setModelOverride(normalizedKey, undefined); }; - const reconcileConfirmedPreviousConnection = async ( - scope: SessionConnectionScope, - agentId?: string | null, - ): Promise => { - const replacement = host.connection.capture(); - if (!replacement || replacement.client !== scope.client) { - return false; - } - let refreshError: string | undefined; - try { - const outcome = await host.reconcileMutation(agentId); - refreshError = outcome.status === "failed" ? outcome.error : undefined; - } catch (error) { - refreshError = formatUiError(error); - } - if (!host.connection.isCurrent(replacement)) { - return false; - } - host.publish( - { - ...host.readState(), - error: refreshError - ? t("connection.sessionOperationCompletedPreviousConnectionWithRefreshError", { - error: refreshError, - }) - : t("connection.sessionOperationCompletedPreviousConnection"), - }, - "operation", - ); - return true; - }; + const { reconcileConfirmedPreviousConnection, refreshCategory, reportUncertainCategory } = + createSessionMutationRefresh(host); const createResult = async ( params: SessionCreateParams = {}, @@ -289,6 +270,7 @@ export function createSessionMutations(host: SessionMutationsHost) { const patchSnapshot = host.snapshot(); const pendingConversation = managesModelOverride || + patchParams.category !== undefined || patchParams.pinned !== undefined || patchParams.unread === false || patchParams.archived !== undefined || @@ -312,6 +294,7 @@ export function createSessionMutations(host: SessionMutationsHost) { ? { ...pendingConversation, sessionId: pendingSessionId } : null; let rowPatchConfirmed = false; + let writeConfirmed = false; let modelPatchStarted = false; let modelPatchRevision = 0; const modelPatchToken = Symbol("session-model-patch"); @@ -357,7 +340,14 @@ export function createSessionMutations(host: SessionMutationsHost) { } unreadPatchToken = optimisticUnread.start(pendingTarget, () => false); }; + let categoryPatchToken: symbol | null = null; const startOptimisticPatch = () => { + if (patchParams.category !== undefined && !categoryPatchToken && pendingTarget) { + categoryPatchToken = optimisticCategories.start( + pendingTarget, + () => patchParams.category?.trim() || undefined, + ); + } startModelPatch(); startPinPatch(); startUnreadPatch(); @@ -426,6 +416,14 @@ export function createSessionMutations(host: SessionMutationsHost) { settleModelOverride(completed); settlePinPatch(completed); settleUnreadPatch(completed); + if (categoryPatchToken && pendingTarget) { + optimisticCategories.settle( + pendingTarget, + categoryPatchToken, + completed && rowPatchConfirmed, + host.connection.isCurrent(scope), + ); + } }; try { if (options.waitFor) { @@ -445,6 +443,7 @@ export function createSessionMutations(host: SessionMutationsHost) { } const confirmFields = pendingTarget ? host.capturePatchFields(pendingTarget) : undefined; const result = await requestSessionPatch(scope.client, key, patchParams, options); + writeConfirmed = true; if (!host.connection.isCurrent(scope)) { settleOptimisticPatch(false); return (await reconcileConfirmedPreviousConnection(scope, options.agentId)) ? result : null; @@ -495,6 +494,18 @@ export function createSessionMutations(host: SessionMutationsHost) { ); } } + // Placement is settled by its durable receipt, not a roster round trip. + // The existing refresh owner still reconciles membership in the background. + if ( + patchParams.category !== undefined && + Object.keys(patchParams).every((name) => name === "category" || name === "pinned") + ) { + settleOptimisticPatch(true); + if (!options.deferListRefresh) { + refreshCategory(scope, options.agentId); + } + return result; + } // Commit and list reconciliation are separate outcomes. Callers must not // turn a failed refresh into an apparent rollback of the committed patch. let refreshOutcome: SessionRefreshOutcome = { status: "refreshed" }; @@ -523,10 +534,25 @@ export function createSessionMutations(host: SessionMutationsHost) { ? { ...result, listRefreshError: refreshOutcome.error } : result; } catch (error) { - settleOptimisticPatch(false); + // Transport loss is not evidence of rollback. Release the tentative value + // without restoring its predecessor, then reconcile the original owner. + const uncertainCategory = + patchParams.category !== undefined && !writeConfirmed && !isRejectedSessionMutation(error); + if (uncertainCategory && categoryPatchToken && pendingTarget) { + optimisticCategories.abandon(pendingTarget, categoryPatchToken); + categoryPatchToken = null; + if (pinPatchToken) { + optimisticPins.abandon(pendingTarget, pinPatchToken); + pinPatchToken = null; + } + } + settleOptimisticPatch(writeConfirmed); if (!host.connection.isCurrent(scope)) { return null; } + if (uncertainCategory) { + throw reportUncertainCategory(error, options.agentId); + } if (ownsModelOverride()) { host.publish({ ...host.readState(), error: formatUiError(error) }, "operation"); } @@ -658,6 +684,7 @@ export function createSessionMutations(host: SessionMutationsHost) { names: readonly string[], sourceAgentId?: string | null, ) { + optimisticCategories.observe(row, names, sourceAgentId); optimisticPins.observe(row, names, sourceAgentId); optimisticUnread.observe(row, names, sourceAgentId); }, @@ -665,7 +692,12 @@ export function createSessionMutations(host: SessionMutationsHost) { result: SessionsListResult | null, sourceAgentId?: string | null, ): SessionsListResult | null { - if (!result || (!optimisticPins.hasPending() && !optimisticUnread.hasPending())) { + if ( + !result || + (!optimisticPins.hasPending() && + !optimisticUnread.hasPending() && + !optimisticCategories.hasPending()) + ) { return result; } return projectSessionResultRows( @@ -698,6 +730,7 @@ export function createSessionMutations(host: SessionMutationsHost) { // Row intents live inside `result`, which the replacement connection // rehydrates wholesale; only the model-override side map outlives that // replacement, so it is the one that needs an explicit rollback below. + optimisticCategories.clear(); optimisticPins.clear(); optimisticUnread.clear(); archiveState.clearAll(); @@ -708,6 +741,7 @@ export function createSessionMutations(host: SessionMutationsHost) { dispose() { pendingCreatedModelOverrides.clear(); pendingModelPatches.clear(); + optimisticCategories.clear(); optimisticPins.clear(); optimisticUnread.clear(); archiveState.clearAll(); diff --git a/ui/src/lib/sessions/session-patch-row-facts.ts b/ui/src/lib/sessions/session-patch-row-facts.ts index 4da8cbec517a..b474a623dab0 100644 --- a/ui/src/lib/sessions/session-patch-row-facts.ts +++ b/ui/src/lib/sessions/session-patch-row-facts.ts @@ -24,6 +24,9 @@ export function projectSessionPatchRowFields( if (typeof patch.archived === "boolean") { fields.push(projectSessionArchiveFields(patch.archived, entry)); } + if (patch.category !== undefined) { + fields.push({ category: entry.category }); + } if (patch.boardPresentation !== undefined) { fields.push({ boardPresentation: entry.boardPresentation }); } diff --git a/ui/src/lib/sessions/session-pending-rows.ts b/ui/src/lib/sessions/session-pending-rows.ts index 96e6c9c0ad0e..11c8cbd8d0bf 100644 --- a/ui/src/lib/sessions/session-pending-rows.ts +++ b/ui/src/lib/sessions/session-pending-rows.ts @@ -39,6 +39,7 @@ export type SessionPatchRowFact = { updatedAt: number | null; readCutoff?: number; fields: + | { category: GatewaySessionRow["category"] } | SessionPinFields | { pinned: true } | SessionReadFields @@ -168,6 +169,13 @@ export function createOptimisticRowPatches( pending.delete(target.identity); } }, + /** An uncertain write releases its overlay without claiming a rollback. */ + abandon(target: PendingRowTarget, token: symbol): void { + const current = pending.get(target.identity); + if (current?.token === token && current.sessionId === target.sessionId) { + pending.delete(target.identity); + } + }, applyRow(row: GatewaySessionRow, sourceAgentId?: string | null): GatewaySessionRow { if (pending.size === 0) { return row; diff --git a/ui/src/lib/sessions/session-reconciliation.ts b/ui/src/lib/sessions/session-reconciliation.ts index e7aa2b56b083..8e7c19df7f3e 100644 --- a/ui/src/lib/sessions/session-reconciliation.ts +++ b/ui/src/lib/sessions/session-reconciliation.ts @@ -86,7 +86,7 @@ type Host = { }; export function createSessionReconciliation(host: Host) { - const pendingFields = ["pinned", "pinnedAt", "unread"] as const; + const pendingFields = ["pinned", "pinnedAt", "unread", "category"] as const; const projectRowFields = ( row: GatewaySessionRow, agentId?: string | null,