fix: session group moves lag while requests are pending (#151689)

Preserve category intent and committed placement independently of roster refresh. Keep saved membership distinct from catalog failure and fence background errors to their foreground owner.

Refs #151647.
This commit is contained in:
Josh Lehman 2026-09-18 12:58:27 -07:00 • committed by GitHub
parent 9d4b4b7cf5
commit 730b8e75be
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
17 changed files with 795 additions and 50 deletions

View file

@ -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;

View file

@ -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";

View file

@ -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;

View file

@ -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<typeof publishSessionPatchEffects>[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();
});
});

View file

@ -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) {

View file

@ -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,

View file

@ -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 }) => {

View file

@ -279,6 +279,7 @@ export type SessionsBranchesSwitchResult =
export type SessionsPatchResult = SessionsPatchResultBase<{
sessionId: string;
category?: GatewaySessionRow["category"];
updatedAt?: number;
createdAt?: number;
pinnedAt?: number;

View file

@ -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:

View file

@ -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,
);
}

View file

@ -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<unknown>(), createDeferred<unknown>()];
const list = createDeferred<ReturnType<typeof sessionsResult>>();
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<void>((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<unknown>();
const oldRead = createDeferred<ReturnType<typeof sessionsResult>>();
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<void>((resolve) => {
setImmediate(resolve);
});
expect(sessions.state.agentId).toBe(selection === "returned" ? "main" : "writer");
expect(sessions.state.error).toBeNull();
} finally {
sessions.dispose();
}
},
);

View file

@ -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<SessionRefreshOutcome>;
};
/** Reconcile committed mutations without confusing refresh failures with failed writes. */
export function createSessionMutationRefresh(host: Host) {
const reconcileConfirmedPreviousConnection = async (
scope: SessionConnectionScope,
agentId?: string | null,
): Promise<boolean> => {
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)
);
}

View file

@ -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(

View file

@ -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<boolean> => {
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();

View file

@ -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 });
}

View file

@ -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<T>(
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;

View file

@ -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,