import path from "node:path";
import { definePluginEntry, type OpenClawPluginApi } from "./api.js";
import { ensureSecafsMountForSession } from "./src/auto-mount-hook.js";
import { resolveSecafsConfig } from "./src/config.js";
import { registerSecafsCli } from "./src/cli.js";
import { createDaemonSupervisor } from "./src/daemon-supervisor.js";
import { registerGatewayMethods } from "./src/gateway-methods.js";
import { startIdleScanner } from "./src/idle-scanner.js";
import { truncateJsonlAfterMessage } from "./src/jsonl-truncate.js";
import { reconcile } from "./src/reconcile.js";
import { createSnapshotOnTurnEnd, findLastAssistantMessageId } from "./src/rollback-hook.js";
import { createHandleRestore } from "./src/rollback-orchestration.js";
import type { SecafsRollbackState } from "./src/rollback-types.js";
import { createSpawnedWorkspaceRedirect } from "./src/spawned-workspace.js";
import { createSecafsRpcClient } from "./src/rpc-client.js";
import { handleSessionEnd } from "./src/session-bindings.js";
import { foldSecafsExt, hoistSecafsExt, hoistSecafsStore } from "./src/session-ext.js";
import {
loadSessionStore,
resolveStorePath,
updateSessionStore,
} from "openclaw/plugin-sdk/session-store-runtime";
import { resolveAgentWorkspaceDir } from "openclaw/plugin-sdk/health";
* The plugin persists three custom fields on session entries that upstream's
* `SessionEntry` type does not declare: `kind` (marks a conversation as
* SecAFS-owned), `mountState`, and `secafsRollback` (per-volume rollback
* state). On disk these live under the sanctioned
* `pluginExtensions["secafs-chat"]` slot (see `session-ext.ts`); the load/save
* wrappers below hoist them to top level in memory so the rest of the plugin
* keeps reading `entry.kind` / `entry.mountState` / `entry.secafsRollback`.
* This type captures that in-memory top-level shape for the typed write sites.
*/
type SecafsSessionFields = { kind?: "secafs"; secafsRollback?: SecafsRollbackState };
export default definePluginEntry({
id: "secafs-chat",
name: "SecAFS Chat",
description: "Per-conversation FUSE-mounted workspace backed by SecAFS.",
register(api: OpenClawPluginApi) {
const cfg = resolveSecafsConfig({ pluginConfig: api.pluginConfig ?? {} }, process.env);
const supervisor = createDaemonSupervisor({
manageDaemon: cfg.manageDaemon,
binary: "secafs",
args: [
"serve",
"api",
"--socket",
cfg.socketPath,
...(cfg.postgresUrl ? ["--pg-url", cfg.postgresUrl] : []),
"--mount-root",
cfg.mountRoot,
],
onExit: (code) => {
api.logger.warn?.(`[secafs-chat] daemon exited with code ${String(code)}`);
},
});
const rpc = createSecafsRpcClient({ socketPath: cfg.socketPath });
const mainKey = (api.config.session?.mainKey?.trim() || "main").toLowerCase();
const agentsList = api.config.agents?.list ?? [];
const defaultAgentId =
agentsList.find((a) => a?.default)?.id?.trim() ?? agentsList[0]?.id?.trim() ?? "main";
const storePath = resolveStorePath(api.config.session?.store, { agentId: defaultAgentId });
const loadStore = () => hoistSecafsStore(loadSessionStore(storePath));
type StoreShape = ReturnType<typeof loadStore>;
const updateStore = async <T>(fn: (store: StoreShape) => T): Promise<T> => {
return await updateSessionStore(storePath, (raw) => {
const store = raw as StoreShape;
for (const k of Object.keys(store)) store[k] = hoistSecafsExt(store[k]);
const result = fn(store);
for (const k of Object.keys(store)) store[k] = foldSecafsExt(store[k]);
return result;
});
};
const sessionStore = {
async load(sessionKey: string) {
const store = loadStore();
return store[sessionKey] ?? null;
},
async patch(sessionKey: string, patch: Record<string, unknown>) {
await updateStore((store) => {
const existing = store[sessionKey];
if (!existing) return;
store[sessionKey] = { ...existing, ...patch } as StoreShape[string];
});
},
};
const sessions = {
async create(_opts: { kind: "secafs" }): Promise<{ sessionKey: string }> {
const { randomUUID } = await import("node:crypto");
const uuid = randomUUID();
const sessionKey = `${mainKey}:${uuid}`;
await updateStore((store) => {
store[sessionKey] = {
sessionId: uuid,
updatedAt: Date.now(),
kind: "secafs",
} as (typeof store)[string];
});
return { sessionKey };
},
async patch(key: string, patch: Record<string, unknown>): Promise<void> {
await updateStore((store) => {
const existing = store[key];
if (!existing) {
return;
}
store[key] = { ...existing, ...patch } as (typeof store)[string];
});
},
async load(key: string): Promise<{ sessionId: string } | null> {
const store = loadStore();
const entry = store[key];
if (!entry?.sessionId) {
return null;
}
return { sessionId: entry.sessionId };
},
async keys(): Promise<string[]> {
const store = loadStore();
return Object.keys(store);
},
async entries(): Promise<Record<string, Record<string, unknown>>> {
return loadStore() as unknown as Record<string, Record<string, unknown>>;
},
async delete(key: string): Promise<void> {
await updateStore((store) => {
delete store[key];
});
},
};
let defaultWorkspaceDir: string | undefined;
try {
defaultWorkspaceDir = resolveAgentWorkspaceDir(api.config, defaultAgentId);
} catch (err) {
api.logger.warn?.(
`[secafs-chat] could not resolve default workspace dir for seeding: ${
err instanceof Error ? err.message : String(err)
}`,
);
}
const sessionFileFor = (sessionKey: string): string => {
const store = loadStore();
const entry = store[sessionKey];
if (entry?.sessionFile) {
return entry.sessionFile;
}
if (entry?.sessionId) {
const dir = path.dirname(storePath);
return path.join(dir, `${entry.sessionId}.jsonl`);
}
const sid = sessionKey.split(":").findLast(Boolean) ?? sessionKey;
return path.join(path.dirname(storePath), `${sid}.jsonl`);
};
const trajectoryFor = (sessionKey: string): string =>
sessionFileFor(sessionKey).replace(/\.jsonl$/, ".trajectory.jsonl");
const workspaceFor = (sessionKey: string): string => {
const sid = sessionKey.split(":").findLast(Boolean) ?? sessionKey;
return path.join(cfg.mountRoot, sid);
};
const extractConvIdFn = (sessionKey: string): string => {
const parts = sessionKey.split(":").filter((p) => p.length > 0);
return parts.length >= 2 ? parts[parts.length - 1] : sessionKey;
};
const events = (event: { type: string; sessionKey: string; [k: string]: unknown }) => {
const broadcast = (
api as unknown as { broadcast?: (channel: string, payload: unknown) => void }
).broadcast;
if (typeof broadcast === "function") {
broadcast.call(api, "secafs.rollback", event);
} else {
api.logger.info?.(`[secafs-chat] event: ${JSON.stringify(event)}`);
}
};
const collectKeyForms = (sessionKey: string): string[] => {
const sid = extractConvIdFn(sessionKey);
const bare = `${mainKey}:${sid}`;
const canonical = `agent:${defaultAgentId}:${mainKey}:${sid}`;
return [...new Set([bare, sessionKey, canonical])];
};
const deleteSessionArtifacts = async (sid: string): Promise<void> => {
const fs = await import("node:fs/promises");
const store = loadStore();
const dir = path.dirname(storePath);
const transcripts = new Set<string>([path.join(dir, `${sid}.jsonl`)]);
for (const key of collectKeyForms(`${mainKey}:${sid}`)) {
const sessionFile = (store[key] as { sessionFile?: string } | undefined)?.sessionFile;
if (sessionFile) transcripts.add(sessionFile);
}
for (const file of transcripts) {
for (const target of [file, file.replace(/\.jsonl$/, ".trajectory.jsonl")]) {
try {
await fs.rm(target, { force: true });
} catch (err) {
api.logger.warn?.(
`[secafs-chat] failed to delete transcript ${target}: ${
err instanceof Error ? err.message : String(err)
}`,
);
}
}
}
};
const readTranscripts = async (
sid: string,
): Promise<{ sessionJsonl?: string; trajectoryJsonl?: string }> => {
const fs = await import("node:fs/promises");
const store = loadStore();
const dir = path.dirname(storePath);
const candidates = [path.join(dir, `${sid}.jsonl`)];
for (const key of collectKeyForms(`${mainKey}:${sid}`)) {
const sessionFile = (store[key] as { sessionFile?: string } | undefined)?.sessionFile;
if (sessionFile && !candidates.includes(sessionFile)) candidates.push(sessionFile);
}
const result: { sessionJsonl?: string; trajectoryJsonl?: string } = {};
for (const file of candidates) {
try {
if (!result.sessionJsonl) {
result.sessionJsonl = (await fs.readFile(file)).toString("base64");
}
} catch {
}
try {
if (!result.trajectoryJsonl) {
const traj = file.replace(/\.jsonl$/, ".trajectory.jsonl");
result.trajectoryJsonl = (await fs.readFile(traj)).toString("base64");
}
} catch {
}
}
return result;
};
const writeTranscripts = async (
sid: string,
chat: { sessionJsonl?: string; trajectoryJsonl?: string },
): Promise<void> => {
const fs = await import("node:fs/promises");
const dir = path.dirname(storePath);
if (chat.sessionJsonl) {
await fs.writeFile(path.join(dir, `${sid}.jsonl`), Buffer.from(chat.sessionJsonl, "base64"));
}
if (chat.trajectoryJsonl) {
await fs.writeFile(
path.join(dir, `${sid}.trajectory.jsonl`),
Buffer.from(chat.trajectoryJsonl, "base64"),
);
}
};
const rollbackSessionStore = {
load: async (key: string): Promise<{ secafsRollback?: SecafsRollbackState } | null> => {
const store = loadStore();
for (const k of collectKeyForms(key)) {
const entry = store[k] as (SecafsSessionFields & { sessionId?: string }) | undefined;
if (entry?.secafsRollback !== undefined) {
return entry as { secafsRollback?: SecafsRollbackState };
}
}
const sid = extractConvIdFn(key);
const bare = store[`${mainKey}:${sid}`];
return bare ? (bare as { secafsRollback?: SecafsRollbackState }) : null;
},
patch: async (key: string, patch: Record<string, unknown>) => {
for (const k of collectKeyForms(key)) {
await sessionStore.patch(k, patch);
}
},
};
const workspace = createSpawnedWorkspaceRedirect(
{ patch: (key, patch) => rollbackSessionStore.patch(key, patch) },
"secafs-chat",
);
const snapshotOnTurnEnd = createSnapshotOnTurnEnd({
rpc,
sessionStore: rollbackSessionStore,
findLastAssistantMessageId,
events,
sessionFileFor,
extractConvId: extractConvIdFn,
logger: api.logger,
});
const snapshotNow = async ({ sessionKey }: { sessionKey: string }) => {
const entry = await rollbackSessionStore.load(sessionKey);
const rb: SecafsRollbackState = entry?.secafsRollback ?? { enabled: false };
if (!rb.enabled || rb.inProgress) {
return { enabled: rb.enabled === true, committed: false };
}
const messageId = await findLastAssistantMessageId(sessionFileFor(sessionKey));
if (!messageId) {
return { enabled: true, committed: false };
}
const conversationId = extractConvIdFn(sessionKey);
if (rb.lastSnapshotMessageId === messageId) {
const list = await rpc.snapshotList({ conversationId });
const hit = [...list.snapshots].reverse().find((s) => s.label === messageId);
return { enabled: true, committed: false, snapId: hit?.snapId, messageId };
}
const r = await rpc.snapshotCommit({ conversationId, label: messageId });
await rollbackSessionStore.patch(sessionKey, {
secafsRollback: { ...rb, lastSnapshotMessageId: messageId },
});
return { enabled: true, committed: true, snapId: r.snapId, messageId };
};
const handleRestoreFn = createHandleRestore({
rpc,
sessionStore: rollbackSessionStore,
truncate: truncateJsonlAfterMessage,
events,
sessionFileFor,
trajectoryFor,
workspaceFor,
extractConvId: extractConvIdFn,
logger: api.logger,
});
registerGatewayMethods({
registerGatewayMethod: (name, handler) => {
api.registerGatewayMethod(name, async ({ params, respond }) => {
try {
const result = await handler(params);
respond(true, result);
} catch (error) {
const message = error instanceof Error ? error.message : String(error);
respond(false, undefined, { code: "UNAVAILABLE", message });
}
});
},
rpc,
sessions,
workspace,
mountRoot: cfg.mountRoot,
defaultWorkspaceDir,
mainKey,
defaultAgentId,
logger: api.logger,
handleRestore: handleRestoreFn,
snapshotNow,
sessionStore: rollbackSessionStore,
enableRollbackUI: cfg.enableRollbackUI,
sessionFileFor: (sid: string) => path.join(path.dirname(storePath), `${sid}.jsonl`),
deleteSessionArtifacts,
readTranscripts,
writeTranscripts,
ensureSessionEntry: async (key, entry) => {
await updateStore((store) => {
if (store[key]) return;
store[key] = { ...entry } as (typeof store)[string];
});
},
});
api.registerCli((ctx) =>
registerSecafsCli(ctx as unknown as Parameters<typeof registerSecafsCli>[0]),
);
api.on("session_end", async (event) => {
const key = event.sessionKey;
if (!key) {
return;
}
const store = loadStore();
const entry = store[key] as { kind?: string; spawnedBy?: string } | undefined;
const managed = entry?.kind === "secafs" || entry?.spawnedBy === "secafs-chat";
await handleSessionEnd(
{ sessionKey: key, workspace: managed ? { managedBy: "secafs-chat" } : undefined },
{ rpc, sessions, logger: api.logger },
);
});
api.on("before_prompt_build", async (_event, ctx) => {
await ensureSecafsMountForSession(
{
loadStore: () => loadStore(),
patchSession: (key, patch) => sessionStore.patch(key, patch),
workspace: { set: (key, spec) => workspace.set(key, spec) },
rpc,
mountRoot: cfg.mountRoot,
logger: api.logger,
},
ctx?.sessionKey,
);
if (cfg.enableRollbackUI) {
await snapshotOnTurnEnd({ sessionKey: ctx?.sessionKey });
}
});
const idleScanner = startIdleScanner(
{
loadStore: () => loadStore(),
patchSession: (key, patch) => sessionStore.patch(key, patch),
rpc,
workspace,
mountRoot: cfg.mountRoot,
logger: api.logger,
},
{
idleUnmountSeconds: cfg.idleUnmountSeconds,
idleScanSeconds: cfg.idleScanSeconds,
},
);
api.on("gateway_stop", async () => {
idleScanner.stop();
try {
rpc.close();
} catch {
}
await supervisor.stop();
});
void (async () => {
try {
await supervisor.start();
if (cfg.manageDaemon) {
await new Promise((resolve) => setTimeout(resolve, 500));
}
const result = await reconcile({
rpc,
sessions: {
load: async (key: string) => {
const store = loadStore();
for (const k of collectKeyForms(key)) {
if (store[k]) return store[k];
}
return null;
},
keys: async () => {
const store = loadStore();
return Object.keys(store);
},
patch: async (key: string, p: Record<string, unknown>) => sessionStore.patch(key, p),
},
sessionKeyFor: (id: string) => `${mainKey}:${id}`,
logger: api.logger,
truncateJsonlAfterMessage,
sessionFileFor,
trajectoryFor,
workspaceFor,
extractConvId: extractConvIdFn,
});
if (result.unmounted > 0) {
api.logger.info?.(`[secafs-chat] reconcile unmounted ${result.unmounted} orphan(s)`);
}
if (result.rollbacksResumed > 0) {
api.logger.info?.(
`[secafs-chat] reconcile resumed ${result.rollbacksResumed} in-progress rollback(s)`,
);
}
} catch (e) {
api.logger.warn?.(`[secafs-chat] async startup failed: ${String(e)}`);
}
})();
api.logger.info?.("[secafs-chat] plugin registered");
},
});