Files
openclaw/extensions/codex/src/session-catalog-control.ts
Vito Cappello 6d7bc062e3 fix(session-catalog): hide OpenClaw-managed provider sessions (#125424)
* 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>
2026-08-20 19:51:55 -07:00

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