mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 09:39:25 +00:00
* test: align frozen-base fixtures with current lifecycle contracts Apply the prerequisite fixture repairs already landed in0b6f02faae(#159836) so the narration comparison base can complete validation. Co-authored-by: IWhatsskill <284122573+IWhatsskill@users.noreply.github.com> * perf(gateway): throttle background sidebar narration with digests Declare narration intent on targeted session subscriptions and pace bounded text snapshots on the server. Preserve full streams whenever any SDK owner renders a foreground transcript, and flush corrected final text before terminal events. Update bundled clients, protocol artifacts, and the event contract. * fix(gateway): preserve foreground ownership across narration addresses Track independent SDK observer IDs at the server so every full-stream interest survives alias collisions, rollback, and unrelated releases. Keep global requests bound to their acknowledged owner and pace distinct logical sessions separately. Wait for foreground admission before history reads, retire stale admissions, and load subscription runtime on demand. Share the existing browser-safe UUID implementation and extend protocol and lifecycle regressions. * test(gateway): expect resolved agent ownership in shared approval replay ACKs * test(gateway): isolate the shared subscription lookup assertion Narration also reads subscription intent for eligible recipients, so the visible chat path legitimately performs three registry lookups. Check the single admission lookup after visibility is revoked, when narration is not entered. Preserve delivery, unrelated-recipient, visibility-revocation, and event-sequencing assertions. * ci: split UI component tests into their own core type shard Apply the upstream partition fix so the narration merge retains complete root coverage within the unchanged local shard-size budget. (cherry picked from commit52ae7e314d) * test(ui): align narration fixtures with stream admission --------- Co-authored-by: IWhatsskill <284122573+IWhatsskill@users.noreply.github.com>
226 lines
8.3 KiB
TypeScript
226 lines
8.3 KiB
TypeScript
import { GatewayProtocolRequestTimeoutError } from "@openclaw/gateway-client/browser";
|
|
import type { GatewayBrowserClient } from "../../api/gateway.ts";
|
|
import { requestSessionRecovery } from "./recover.ts";
|
|
import type {
|
|
SessionCompactResult,
|
|
SessionCapability,
|
|
SessionConnectionOwner,
|
|
SessionMessageSubscription,
|
|
SessionRefreshOutcome,
|
|
} from "./session-capability.ts";
|
|
import { areUiSessionKeysEquivalent, normalizeAgentId } from "./session-key.ts";
|
|
import {
|
|
requestSessionBranchSwitch,
|
|
requestSessionBranches,
|
|
requestSessionCompact,
|
|
requestSessionFile,
|
|
requestSessionFilesList,
|
|
requestSessionFileSet,
|
|
requestSessionFork,
|
|
requestSessionRewind,
|
|
} from "./session-requests.ts";
|
|
|
|
type SessionScopedOperationsHost = {
|
|
connection: SessionConnectionOwner;
|
|
reconcileMutation: (agentId?: string | null) => Promise<SessionRefreshOutcome>;
|
|
notifyCreated: (key: string) => void;
|
|
reportError: (error: unknown) => void;
|
|
};
|
|
|
|
const retiredFailedSubscriptionRecoveries = new WeakSet<AggregateError>();
|
|
|
|
export function createSessionScopedOperations(host: SessionScopedOperationsHost) {
|
|
const ownedSubscriptions = new Set<SessionMessageSubscription>();
|
|
type SubscriptionRuntime = typeof import("./session-message-subscriptions.runtime.ts");
|
|
let subscriptionRuntime: SubscriptionRuntime | undefined;
|
|
let subscriptionRuntimeLoading: Promise<SubscriptionRuntime> | undefined;
|
|
let disposed = false;
|
|
const loadSubscriptionRuntime = () =>
|
|
(subscriptionRuntimeLoading ??= import("./session-message-subscriptions.runtime.ts").then(
|
|
(runtime) => (subscriptionRuntime = runtime),
|
|
(error: unknown) => {
|
|
subscriptionRuntimeLoading = undefined;
|
|
throw error;
|
|
},
|
|
));
|
|
|
|
const recover = async (params: { key: string; agentId?: string }) => {
|
|
const scope = host.connection.capture();
|
|
if (!scope) {
|
|
return null;
|
|
}
|
|
try {
|
|
const result = await requestSessionRecovery(scope.client, params);
|
|
if (!host.connection.isCurrent(scope)) {
|
|
return null;
|
|
}
|
|
host.notifyCreated(result.key);
|
|
await host.reconcileMutation(params.agentId);
|
|
return host.connection.isCurrent(scope) ? result : null;
|
|
} catch (error) {
|
|
if (host.connection.isCurrent(scope)) {
|
|
host.reportError(error);
|
|
}
|
|
return null;
|
|
}
|
|
};
|
|
|
|
const compact = async (
|
|
key: string,
|
|
options: { agentId?: string | null } = {},
|
|
): Promise<SessionCompactResult> => {
|
|
const scope = host.connection.capture();
|
|
if (!scope) {
|
|
throw new Error("Session compaction requires an active Gateway connection");
|
|
}
|
|
const result = await requestSessionCompact(scope.client, key, options);
|
|
if (!host.connection.isCurrent(scope)) {
|
|
throw new Error("Session compaction completed on a replaced Gateway connection");
|
|
}
|
|
return result;
|
|
};
|
|
|
|
const requestCurrent = async <T>(
|
|
request: (client: GatewayBrowserClient) => Promise<T>,
|
|
): Promise<T | null> => {
|
|
const scope = host.connection.capture();
|
|
if (!scope) {
|
|
return null;
|
|
}
|
|
const result = await request(scope.client);
|
|
return host.connection.isCurrent(scope) ? result : null;
|
|
};
|
|
|
|
const listFiles: SessionCapability["listFiles"] = (key, options = {}) =>
|
|
requestCurrent((client) => requestSessionFilesList(client, key, options));
|
|
|
|
const getFile: SessionCapability["getFile"] = (key, path, options = {}) =>
|
|
requestCurrent((client) => requestSessionFile(client, key, path, options));
|
|
|
|
const setFile: SessionCapability["setFile"] = (key, path, content, options) =>
|
|
requestCurrent((client) => requestSessionFileSet(client, key, path, content, options));
|
|
|
|
const unsubscribeMessages = async (subscription: SessionMessageSubscription): Promise<void> => {
|
|
const runtime = subscriptionRuntime ?? (await loadSubscriptionRuntime());
|
|
await runtime.releaseGatewaySessionMessageSubscription(subscription);
|
|
ownedSubscriptions.delete(subscription);
|
|
};
|
|
|
|
const subscribeMessages = async (
|
|
key: string,
|
|
options: NonNullable<Parameters<SessionCapability["subscribeMessages"]>[1]> = {},
|
|
): Promise<SessionMessageSubscription> => {
|
|
const scope = host.connection.capture();
|
|
if (!scope || disposed) {
|
|
throw new Error("Session message subscription requires an active Gateway connection");
|
|
}
|
|
const normalizedKey = key.trim();
|
|
const agentId = options.agentId?.trim() ? normalizeAgentId(options.agentId) : null;
|
|
const { mode, includeApprovals } = options;
|
|
const runtime = subscriptionRuntime ?? (await loadSubscriptionRuntime());
|
|
if (disposed || !host.connection.isCurrent(scope)) {
|
|
throw new Error("Session message subscription completed on a replaced Gateway connection");
|
|
}
|
|
const subscription = await runtime
|
|
.getGatewaySessionMessageSubscriptionCoordinator(scope.client, {
|
|
keysEquivalent: areUiSessionKeysEquivalent,
|
|
})
|
|
.acquire(normalizedKey, {
|
|
agentId,
|
|
...(includeApprovals ? { includeApprovals: true } : {}),
|
|
...(mode ? { mode } : {}),
|
|
})
|
|
.catch((error: unknown) => {
|
|
if (
|
|
error instanceof AggregateError &&
|
|
error.errors[0] instanceof GatewayProtocolRequestTimeoutError &&
|
|
error.errors[0].requestSent &&
|
|
!disposed &&
|
|
host.connection.isCurrent(scope) &&
|
|
!retiredFailedSubscriptionRecoveries.has(error)
|
|
) {
|
|
// Failed compensation cannot prove privileged observers were removed;
|
|
// closing their owning socket invokes authoritative Gateway cleanup.
|
|
retiredFailedSubscriptionRecoveries.add(error);
|
|
scope.client.forceReconnect("session subscription recovery failed");
|
|
}
|
|
throw error;
|
|
});
|
|
ownedSubscriptions.add(subscription);
|
|
if (disposed || !host.connection.isCurrent(scope)) {
|
|
await unsubscribeMessages(subscription).catch(() => undefined);
|
|
throw new Error("Session message subscription completed on a replaced Gateway connection");
|
|
}
|
|
return subscription;
|
|
};
|
|
|
|
const requestCommittedMutation = async <T>(
|
|
disconnectedError: string,
|
|
request: (client: GatewayBrowserClient) => Promise<T>,
|
|
agentId?: string | null,
|
|
): Promise<T> => {
|
|
const scope = host.connection.capture();
|
|
if (!scope) {
|
|
throw new Error(disconnectedError);
|
|
}
|
|
const result = await request(scope.client);
|
|
// The gateway response commits destructive work; refresh is connection-scoped
|
|
// best effort and must never turn that commit into uncertainty or a retry.
|
|
if (host.connection.isCurrent(scope)) {
|
|
await host.reconcileMutation(agentId).catch(() => {});
|
|
}
|
|
return result;
|
|
};
|
|
|
|
const rewind: SessionCapability["rewind"] = (key, entryId, options = {}) =>
|
|
requestCommittedMutation(
|
|
"Session rewind requires an active Gateway connection",
|
|
(client) => requestSessionRewind(client, key, entryId, options),
|
|
options.agentId,
|
|
);
|
|
|
|
const forkAtMessage: SessionCapability["forkAtMessage"] = (key, entryId, options = {}) =>
|
|
requestCommittedMutation(
|
|
"Session fork requires an active Gateway connection",
|
|
(client) => requestSessionFork(client, key, entryId, options),
|
|
options.agentId,
|
|
);
|
|
|
|
const listBranches: SessionCapability["listBranches"] = async (key, options = {}) =>
|
|
(await requestCurrent((client) => requestSessionBranches(client, key, options))) ?? [];
|
|
|
|
const switchBranch: SessionCapability["switchBranch"] = (key, leafEntryId, options = {}) =>
|
|
requestCommittedMutation(
|
|
"Session branch switch requires an active Gateway connection",
|
|
(client) => requestSessionBranchSwitch(client, key, leafEntryId, options),
|
|
options.agentId,
|
|
);
|
|
|
|
return {
|
|
compact,
|
|
forkAtMessage,
|
|
getFile,
|
|
listBranches,
|
|
listFiles,
|
|
recover,
|
|
rewind,
|
|
setFile,
|
|
subscribeMessages,
|
|
switchBranch,
|
|
unsubscribeMessages,
|
|
retireConnection: (previousClient: GatewayBrowserClient | null) => {
|
|
if (previousClient) {
|
|
// No observer can be acquired before the runtime is installed and its
|
|
// captured connection revalidated, so a pending import needs no reset.
|
|
subscriptionRuntime?.resetGatewaySessionMessageSubscriptionCoordinator(previousClient);
|
|
}
|
|
ownedSubscriptions.clear();
|
|
},
|
|
dispose: () => {
|
|
disposed = true;
|
|
for (const subscription of ownedSubscriptions) {
|
|
void unsubscribeMessages(subscription).catch(() => undefined);
|
|
}
|
|
},
|
|
};
|
|
}
|