mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-26 04:15:48 -06:00
6d7bc062e3
* fix(session-catalog): hide OpenClaw-managed upstream sessions * fix(codex): filter managed paired-node sessions * fix(codex): classify legacy managed sessions * fix(session-catalog): classify managed provider sessions * fix(session-catalog): backfill inter-session ownership * fix(session-catalog): classify Claude internal prompts * fix(session-catalog): retain durable provenance * fix(codex): keep rollout home derivation private * fix(anthropic): declare catalog schema dependency * fix(anthropic): avoid catalog schema dependency * fix(session-catalog): scope managed ownership to Codex * fix(codex): contain catalog provenance reads * fix(codex): bind managed threads to catalog home --------- Co-authored-by: VACInc <3279061+VACInc@users.noreply.github.com> Co-authored-by: Josh Lehman <550978+jalehman@users.noreply.github.com>
435 lines
16 KiB
TypeScript
435 lines
16 KiB
TypeScript
import { resolveAgentDir } from "openclaw/plugin-sdk/agent-runtime";
|
|
import { pruneMapToMaxSize } from "openclaw/plugin-sdk/collection-runtime";
|
|
import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts";
|
|
import { CODEX_CONTROL_METHODS } from "./app-server/capabilities.js";
|
|
import { resolveCodexAppServerClientInstanceId } from "./app-server/client.js";
|
|
import {
|
|
resolveCodexSupervisionAppServerRuntimeOptions,
|
|
type CodexAppServerStartOptions,
|
|
} from "./app-server/config.js";
|
|
import { buildCodexAppServerConnectionFingerprint } from "./app-server/plugin-app-cache-key.js";
|
|
import { assertCodexThreadForkParams } from "./app-server/protocol-validators.js";
|
|
import type {
|
|
CodexAppServerRequestParams,
|
|
CodexAppServerRequestResult,
|
|
CodexThread,
|
|
CodexThreadForkParams,
|
|
CodexThreadForkResponse,
|
|
CodexThreadListParams,
|
|
CodexThreadListResponse,
|
|
CodexThreadTurnsListParams,
|
|
CodexThreadTurnsListResponse,
|
|
} from "./app-server/protocol.js";
|
|
import { requestCodexAppServerClientJson } from "./app-server/request.js";
|
|
import {
|
|
getLeasedSharedCodexAppServerClient,
|
|
releaseLeasedSharedCodexAppServerClient,
|
|
} from "./app-server/shared-client.js";
|
|
import { codexControlRequest } from "./command-rpc.js";
|
|
import { createCodexCatalogHomeResolver, type CodexCatalogHome } from "./session-catalog-homes.js";
|
|
import {
|
|
MAX_TITLE_SEARCH_CATALOG_PAGES,
|
|
normalizeLimit,
|
|
readControlCursor,
|
|
toCatalogSession,
|
|
} from "./session-catalog-parsing.js";
|
|
import { isOpenClawManagedCodexThread } from "./session-catalog-provenance.js";
|
|
import type {
|
|
CodexSessionCatalogControl,
|
|
CodexSessionCatalogControlFactory,
|
|
CodexSessionCatalogPage,
|
|
CodexSessionCatalogPageParams,
|
|
CodexSessionCatalogSession,
|
|
} from "./session-catalog-types.js";
|
|
|
|
const CODEX_SESSION_CATALOG_LIST_TTL_MS = 32_000;
|
|
const CODEX_SESSION_CATALOG_LIST_CACHE_MAX_ENTRIES = 32;
|
|
|
|
type CodexCatalogRequestOptions = {
|
|
agentDir: string;
|
|
config: OpenClawConfig | undefined;
|
|
startOptions: CodexAppServerStartOptions;
|
|
};
|
|
|
|
type CodexCatalogPageCacheEntry = {
|
|
expiresAt: number;
|
|
page: Promise<CodexSessionCatalogPage>;
|
|
value?: CodexSessionCatalogPage;
|
|
};
|
|
|
|
function codexCatalogPageCacheKey(
|
|
params: CodexSessionCatalogPageParams,
|
|
agentId: string,
|
|
source?: CodexCatalogHome,
|
|
): string {
|
|
// Mirror listPage's search/cwd normalization; these trimmed values are what reach app-server.
|
|
return JSON.stringify([
|
|
agentId,
|
|
source?.sourceHomeId ?? null,
|
|
params.cursor ?? null,
|
|
params.limit ?? null,
|
|
params.searchTerm?.trim().toLocaleLowerCase() || null,
|
|
params.cwd?.trim() || null,
|
|
]);
|
|
}
|
|
|
|
type CodexSessionCatalogRequestSnapshot = {
|
|
requestTimeoutMs: number;
|
|
listThreads(params: CodexThreadListParams, timeoutMs: number): Promise<CodexThreadListResponse>;
|
|
listThreadTurns(params: CodexThreadTurnsListParams): Promise<CodexThreadTurnsListResponse>;
|
|
forkThread(params: CodexThreadForkParams): Promise<CodexThreadForkResponse>;
|
|
readThread(threadId: string, includeTurns: boolean): Promise<CodexThread>;
|
|
archiveThread(threadId: string): Promise<void>;
|
|
};
|
|
|
|
type CodexCatalogRequestMethod =
|
|
| typeof CODEX_CONTROL_METHODS.archiveThread
|
|
| typeof CODEX_CONTROL_METHODS.forkThread
|
|
| typeof CODEX_CONTROL_METHODS.listThreads
|
|
| typeof CODEX_CONTROL_METHODS.listThreadTurns
|
|
| typeof CODEX_CONTROL_METHODS.readThread;
|
|
|
|
type CodexCatalogRequest = <M extends CodexCatalogRequestMethod>(
|
|
method: M,
|
|
requestParams: CodexAppServerRequestParams<M>,
|
|
timeoutMs?: number,
|
|
) => Promise<CodexAppServerRequestResult<M>>;
|
|
|
|
function createCodexCatalogRequestSnapshot(
|
|
requestTimeoutMs: number,
|
|
request: CodexCatalogRequest,
|
|
): CodexSessionCatalogRequestSnapshot {
|
|
return {
|
|
requestTimeoutMs,
|
|
listThreads: (params, timeoutMs) =>
|
|
request(CODEX_CONTROL_METHODS.listThreads, params, timeoutMs),
|
|
listThreadTurns: (params) => request(CODEX_CONTROL_METHODS.listThreadTurns, params),
|
|
forkThread: (params) =>
|
|
request(CODEX_CONTROL_METHODS.forkThread, assertCodexThreadForkParams(params)),
|
|
readThread: async (threadId, includeTurns) =>
|
|
(await request(CODEX_CONTROL_METHODS.readThread, { threadId, includeTurns })).thread,
|
|
archiveThread: async (threadId) => {
|
|
await request(CODEX_CONTROL_METHODS.archiveThread, { threadId });
|
|
},
|
|
};
|
|
}
|
|
|
|
function createCodexSessionCatalogControlFromRequests(params: {
|
|
clientId?: string;
|
|
connectionFingerprint?: string;
|
|
createRequestSnapshot: () => CodexSessionCatalogRequestSnapshot;
|
|
localSessionsRoot?: string;
|
|
now: () => number;
|
|
withPinnedConnection: CodexSessionCatalogControl["withPinnedConnection"];
|
|
}): CodexSessionCatalogControl {
|
|
return {
|
|
...(params.clientId ? { clientId: params.clientId } : {}),
|
|
...(params.connectionFingerprint
|
|
? { connectionFingerprint: params.connectionFingerprint }
|
|
: {}),
|
|
withPinnedConnection: params.withPinnedConnection,
|
|
async listPage(pageParams) {
|
|
const limit = normalizeLimit(pageParams.limit, "limit");
|
|
// App Server search also matches transcript previews. Scan native pages
|
|
// without that filter so this catalog remains a title-only surface.
|
|
const search = pageParams.searchTerm?.trim().toLocaleLowerCase() || undefined;
|
|
const cwd = pageParams.cwd?.trim() || undefined;
|
|
const maxPages = search ? MAX_TITLE_SEARCH_CATALOG_PAGES : 1;
|
|
const sessions: CodexSessionCatalogSession[] = [];
|
|
const managedThreads: Array<{ threadId: string; rolloutPath?: string }> = [];
|
|
let cursor = readControlCursor(pageParams.cursor, "request");
|
|
let nextCursor: string | undefined;
|
|
let backwardsCursor: string | undefined;
|
|
const seenCursors = new Set(cursor ? [cursor] : []);
|
|
const requests = params.createRequestSnapshot();
|
|
const deadline = params.now() + requests.requestTimeoutMs;
|
|
|
|
for (let pageIndex = 0; pageIndex < maxPages; pageIndex += 1) {
|
|
const remainingTimeoutMs = Math.ceil(deadline - params.now());
|
|
if (remainingTimeoutMs <= 0) {
|
|
throw new Error("Codex session catalog listing timed out");
|
|
}
|
|
const response = await requests.listThreads(
|
|
{
|
|
archived: false,
|
|
limit: limit - sessions.length,
|
|
modelProviders: [],
|
|
// Match Codex's resume picker/latest-session ordering so a session
|
|
// created outside OpenClaw enters the first catalog page immediately.
|
|
sortKey: "updated_at",
|
|
sortDirection: "desc",
|
|
...(cwd ? { cwd } : {}),
|
|
...(cursor ? { cursor } : {}),
|
|
},
|
|
remainingTimeoutMs,
|
|
);
|
|
if (pageIndex === 0) {
|
|
backwardsCursor = readControlCursor(response.backwardsCursor, "backwards response");
|
|
}
|
|
for (const thread of response.data) {
|
|
if (await isOpenClawManagedCodexThread(thread, params.localSessionsRoot)) {
|
|
const rolloutPath = typeof thread.path === "string" ? thread.path.trim() : "";
|
|
managedThreads.push({
|
|
threadId: thread.id,
|
|
...(rolloutPath ? { rolloutPath } : {}),
|
|
});
|
|
continue;
|
|
}
|
|
const session = toCatalogSession(thread, false);
|
|
if (session && (!search || session.name?.toLocaleLowerCase().includes(search))) {
|
|
sessions.push(session);
|
|
}
|
|
}
|
|
nextCursor = readControlCursor(response.nextCursor, "next response");
|
|
if (!nextCursor || sessions.length >= limit) {
|
|
break;
|
|
}
|
|
if (seenCursors.has(nextCursor)) {
|
|
throw new Error("Codex session catalog returned a repeated search cursor");
|
|
}
|
|
seenCursors.add(nextCursor);
|
|
cursor = nextCursor;
|
|
}
|
|
return {
|
|
sessions,
|
|
...(managedThreads.length > 0 ? { managedThreads } : {}),
|
|
...(nextCursor ? { nextCursor } : {}),
|
|
...(backwardsCursor ? { backwardsCursor } : {}),
|
|
};
|
|
},
|
|
async listDescendantPage(listParams) {
|
|
const requests = params.createRequestSnapshot();
|
|
const response = await requests.listThreads(listParams, requests.requestTimeoutMs);
|
|
return response;
|
|
},
|
|
async readThread(threadId, includeTurns = false) {
|
|
const thread = await params.createRequestSnapshot().readThread(threadId, includeTurns);
|
|
return thread;
|
|
},
|
|
async listTurnPage(listParams) {
|
|
const response = await params.createRequestSnapshot().listThreadTurns(listParams);
|
|
return response;
|
|
},
|
|
async forkThread(forkParams) {
|
|
return await params.createRequestSnapshot().forkThread(forkParams);
|
|
},
|
|
async archiveThread(threadId) {
|
|
await params.createRequestSnapshot().archiveThread(threadId);
|
|
},
|
|
};
|
|
}
|
|
|
|
/** Builds the passive catalog over the Codex plugin's canonical shared client. */
|
|
export function createCodexSessionCatalogControl(params: {
|
|
config?: OpenClawConfig;
|
|
env?: NodeJS.ProcessEnv;
|
|
getPluginConfig: () => unknown;
|
|
getRuntimeConfig: () => OpenClawConfig | undefined;
|
|
now?: () => number;
|
|
}): CodexSessionCatalogControlFactory {
|
|
const now = params.now ?? Date.now;
|
|
const getPluginConfig = () => params.getPluginConfig();
|
|
const homeResolver = createCodexCatalogHomeResolver({
|
|
config: params.getRuntimeConfig() ?? params.config ?? {},
|
|
getRuntimeConfig: params.getRuntimeConfig,
|
|
getPluginConfig: params.getPluginConfig,
|
|
...(params.env ? { env: params.env } : {}),
|
|
});
|
|
const requestOptionsByConfig = new WeakMap<
|
|
OpenClawConfig,
|
|
Map<string, CodexCatalogRequestOptions>
|
|
>();
|
|
const catalogPagesByConfig = new WeakMap<
|
|
OpenClawConfig,
|
|
Map<string, CodexCatalogPageCacheEntry>
|
|
>();
|
|
const resolveRequestOptions = (
|
|
startOptions: CodexAppServerStartOptions,
|
|
agentId: string,
|
|
source?: CodexCatalogHome,
|
|
): CodexCatalogRequestOptions => {
|
|
const runtimeConfig = params.getRuntimeConfig();
|
|
const agentDir = source?.agentDir ?? resolveAgentDir(runtimeConfig ?? {}, agentId);
|
|
const resolvedStartOptions = source?.appServer.start ?? startOptions;
|
|
if (!runtimeConfig) {
|
|
return {
|
|
agentDir,
|
|
config: undefined,
|
|
startOptions: structuredClone(resolvedStartOptions),
|
|
};
|
|
}
|
|
let byAgent = requestOptionsByConfig.get(runtimeConfig);
|
|
const cacheKey = `${agentId ?? ""}\0${source?.sourceHomeId ?? ""}`;
|
|
const cached = byAgent?.get(cacheKey);
|
|
if (cached) {
|
|
// Plugin start options derive from this same immutable config snapshot. Config reload changes
|
|
// object identity; re-cloning on every poll only adds CPU and allocation to the catalog path.
|
|
return cached;
|
|
}
|
|
const resolved = {
|
|
agentDir,
|
|
config: structuredClone(runtimeConfig),
|
|
startOptions: structuredClone(resolvedStartOptions),
|
|
};
|
|
if (!byAgent) {
|
|
byAgent = new Map();
|
|
requestOptionsByConfig.set(runtimeConfig, byAgent);
|
|
}
|
|
byAgent.set(cacheKey, resolved);
|
|
return resolved;
|
|
};
|
|
const createRequestSnapshot = (
|
|
agentId: string,
|
|
source?: CodexCatalogHome,
|
|
): CodexSessionCatalogRequestSnapshot => {
|
|
const pluginConfig = getPluginConfig();
|
|
const runtime =
|
|
source?.appServer ?? resolveCodexSupervisionAppServerRuntimeOptions({ pluginConfig });
|
|
const requestOptions = resolveRequestOptions(runtime.start, agentId, source);
|
|
return createCodexCatalogRequestSnapshot(
|
|
runtime.requestTimeoutMs,
|
|
async (method, requestParams, timeoutMs) =>
|
|
await codexControlRequest(pluginConfig, method, requestParams, {
|
|
...requestOptions,
|
|
...(timeoutMs === undefined ? {} : { timeoutMs }),
|
|
}),
|
|
);
|
|
};
|
|
|
|
const forRequest = (agentId: string, source?: CodexCatalogHome): CodexSessionCatalogControl => {
|
|
const withPinnedConnection: CodexSessionCatalogControl["withPinnedConnection"] = async (
|
|
run,
|
|
) => {
|
|
const pluginConfig = getPluginConfig();
|
|
const runtime =
|
|
source?.appServer ?? resolveCodexSupervisionAppServerRuntimeOptions({ pluginConfig });
|
|
const {
|
|
agentDir,
|
|
config: runtimeConfig,
|
|
startOptions,
|
|
} = resolveRequestOptions(runtime.start, agentId, source);
|
|
const client = await getLeasedSharedCodexAppServerClient({
|
|
agentDir,
|
|
config: runtimeConfig,
|
|
startOptions,
|
|
timeoutMs: runtime.requestTimeoutMs,
|
|
});
|
|
try {
|
|
const requests = createCodexCatalogRequestSnapshot(
|
|
runtime.requestTimeoutMs,
|
|
async <M extends CodexCatalogRequestMethod>(
|
|
method: M,
|
|
requestParams: CodexAppServerRequestParams<M>,
|
|
timeoutMs?: number,
|
|
): Promise<CodexAppServerRequestResult<M>> =>
|
|
await requestCodexAppServerClientJson<CodexAppServerRequestResult<M>>({
|
|
client,
|
|
method,
|
|
requestParams,
|
|
config: runtimeConfig,
|
|
timeoutMs: timeoutMs ?? runtime.requestTimeoutMs,
|
|
}),
|
|
);
|
|
const pinnedControl: CodexSessionCatalogControl =
|
|
createCodexSessionCatalogControlFromRequests({
|
|
clientId: resolveCodexAppServerClientInstanceId(client),
|
|
connectionFingerprint: buildCodexAppServerConnectionFingerprint(runtime, agentDir),
|
|
createRequestSnapshot: () => requests,
|
|
...(source?.localSessionsRoot ? { localSessionsRoot: source.localSessionsRoot } : {}),
|
|
now,
|
|
withPinnedConnection: async (nestedRun) => await nestedRun(pinnedControl),
|
|
});
|
|
return await run(pinnedControl);
|
|
} finally {
|
|
releaseLeasedSharedCodexAppServerClient(client);
|
|
}
|
|
};
|
|
const control = createCodexSessionCatalogControlFromRequests({
|
|
createRequestSnapshot: () => createRequestSnapshot(agentId, source),
|
|
...(source?.localSessionsRoot ? { localSessionsRoot: source.localSessionsRoot } : {}),
|
|
now,
|
|
withPinnedConnection,
|
|
});
|
|
return {
|
|
...control,
|
|
async listPage(pageParams: CodexSessionCatalogPageParams) {
|
|
const runtimeConfig = params.getRuntimeConfig();
|
|
if (!runtimeConfig) {
|
|
return await control.listPage(pageParams);
|
|
}
|
|
let cache = catalogPagesByConfig.get(runtimeConfig);
|
|
if (!cache) {
|
|
cache = new Map();
|
|
catalogPagesByConfig.set(runtimeConfig, cache);
|
|
}
|
|
const key = codexCatalogPageCacheKey(pageParams, agentId, source);
|
|
const cached = cache.get(key);
|
|
if (pageParams.forceRefresh !== true && cached) {
|
|
cache.delete(key);
|
|
cache.set(key, cached);
|
|
if (cached.expiresAt > now()) {
|
|
return cached.value ?? (await cached.page);
|
|
}
|
|
}
|
|
if (cached) {
|
|
cache.delete(key);
|
|
}
|
|
const page = control.listPage(pageParams);
|
|
const staleValue = cached?.value;
|
|
const entry: CodexCatalogPageCacheEntry = {
|
|
expiresAt: Number.POSITIVE_INFINITY,
|
|
page,
|
|
...(staleValue ? { value: staleValue } : {}),
|
|
};
|
|
cache.set(key, entry);
|
|
pruneMapToMaxSize(cache, CODEX_SESSION_CATALOG_LIST_CACHE_MAX_ENTRIES);
|
|
const settle = (value: CodexSessionCatalogPage) => {
|
|
if (cache.get(key) === entry) {
|
|
entry.value = value;
|
|
entry.expiresAt = now() + CODEX_SESSION_CATALOG_LIST_TTL_MS;
|
|
}
|
|
return value;
|
|
};
|
|
const restore = () => {
|
|
if (cache.get(key) !== entry) {
|
|
return;
|
|
}
|
|
if (staleValue) {
|
|
cache.set(key, {
|
|
expiresAt: now(),
|
|
page: Promise.resolve(staleValue),
|
|
value: staleValue,
|
|
});
|
|
} else {
|
|
cache.delete(key);
|
|
}
|
|
};
|
|
// Expiry starts one background refresh. Passive callers keep the last settled page while
|
|
// the next poll publishes success or retries failure; a forced caller still sees failure.
|
|
if (pageParams.forceRefresh !== true && staleValue) {
|
|
void page.then(settle, restore);
|
|
return staleValue;
|
|
}
|
|
try {
|
|
return settle(await page);
|
|
} catch (error) {
|
|
restore();
|
|
throw error;
|
|
}
|
|
},
|
|
};
|
|
};
|
|
const homesForAgent = (agentId: string) => homeResolver.forAgent(agentId);
|
|
const forUpstream = (agentId: string, connectionFingerprint: string) => {
|
|
// A fingerprint is correlation only. A miss must stay fail-closed instead of selecting a
|
|
// different home whose thread namespace could contain the same copied identifier.
|
|
const source = homesForAgent(agentId).find(
|
|
(home) =>
|
|
buildCodexAppServerConnectionFingerprint(home.appServer, home.agentDir) ===
|
|
connectionFingerprint,
|
|
);
|
|
return source ? forRequest(agentId, source) : undefined;
|
|
};
|
|
return { forRequest, forUpstream, homesForAgent };
|
|
}
|