mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 09:39:25 +00:00
refactor(sessions): move conversation reads into workers (#163208)
* refactor(sessions): move conversation reads into workers Read conversation lists, standalone send/turn addresses, and manual-compaction primary bindings through the existing history worker. Retain source custody through discovery and cleanup, current route checks, and native incognito reads. Keep getConversationSession synchronous for the public plugin SDK. Transaction and delivery authority guards remain native pending the S10 entry-patch cutover. No schema, stored data, migration, retention, or permission changes. * test(sessions): await conversation reads in account history proof
This commit is contained in:
parent
c79b1cfca4
commit
745421393a
23 changed files with 478 additions and 111 deletions
|
|
@ -57,7 +57,7 @@ describe("dashboard turns retain external conversation identity", () => {
|
|||
);
|
||||
await initSessionState({ cfg, ctx: externalContext(), commandAuthorized: true });
|
||||
const originalEntry = loadSessionEntry({ ...scope, sessionKey });
|
||||
const originalConversations = listConversations(scope, { channel });
|
||||
const originalConversations = await listConversations(scope, { channel });
|
||||
expect(originalConversations).toHaveLength(1);
|
||||
expect(originalConversations[0]).toMatchObject({ kind, role: "primary", sessionKey });
|
||||
|
||||
|
|
@ -76,7 +76,7 @@ describe("dashboard turns retain external conversation identity", () => {
|
|||
chatType: kind,
|
||||
delivery: originalEntry?.delivery,
|
||||
});
|
||||
expect(listConversations(scope, { channel })).toEqual([
|
||||
expect(await listConversations(scope, { channel })).toEqual([
|
||||
expect.objectContaining({
|
||||
conversationRef: originalConversations[0]?.conversationRef,
|
||||
kind,
|
||||
|
|
@ -86,7 +86,7 @@ describe("dashboard turns retain external conversation identity", () => {
|
|||
]);
|
||||
|
||||
await initSessionState({ cfg, ctx: externalContext(), commandAuthorized: true });
|
||||
expect(listConversations(scope, { channel })).toEqual([
|
||||
expect(await listConversations(scope, { channel })).toEqual([
|
||||
expect.objectContaining({
|
||||
conversationRef: originalConversations[0]?.conversationRef,
|
||||
kind,
|
||||
|
|
@ -110,6 +110,6 @@ describe("dashboard turns retain external conversation identity", () => {
|
|||
chatType: "direct",
|
||||
delivery: { kind: "internal" },
|
||||
});
|
||||
expect(listConversations({ agentId: "main", storePath })).toEqual([]);
|
||||
expect(await listConversations({ agentId: "main", storePath })).toEqual([]);
|
||||
});
|
||||
});
|
||||
|
|
|
|||
|
|
@ -66,7 +66,7 @@ describe("doctor Telegram General-topic conversation repair", () => {
|
|||
delivery: delivery("telegram:-1001234567890"),
|
||||
});
|
||||
|
||||
const before = listConversations(scope, { channel: "telegram" });
|
||||
const before = await listConversations(scope, { channel: "telegram" });
|
||||
expect(
|
||||
before
|
||||
.map(({ target, role }) => ({ target, role }))
|
||||
|
|
@ -122,7 +122,7 @@ describe("doctor Telegram General-topic conversation repair", () => {
|
|||
"Merged 1 stale Telegram General-topic conversation identity row(s).",
|
||||
]);
|
||||
|
||||
expect(listConversations(scope, { channel: "telegram" })).toEqual([
|
||||
expect(await listConversations(scope, { channel: "telegram" })).toEqual([
|
||||
expect.objectContaining({
|
||||
conversationRef: canonical!.conversationRef,
|
||||
target: "telegram:-1001234567890",
|
||||
|
|
@ -150,7 +150,7 @@ describe("doctor Telegram General-topic conversation repair", () => {
|
|||
);
|
||||
expect(repeated.findings).toEqual([]);
|
||||
expect(repeated.changes).toEqual([]);
|
||||
expect(listConversations(scope, { channel: "telegram" })).toHaveLength(1);
|
||||
expect(await listConversations(scope, { channel: "telegram" })).toHaveLength(1);
|
||||
});
|
||||
|
||||
it("canonicalizes a legacy-only current entry before later session writes", async () => {
|
||||
|
|
@ -180,7 +180,7 @@ describe("doctor Telegram General-topic conversation repair", () => {
|
|||
{ checks: [check!] },
|
||||
);
|
||||
expect(repaired.remainingFindings).toEqual([]);
|
||||
expect(listConversations(scope, { channel: "telegram" })).toEqual([
|
||||
expect(await listConversations(scope, { channel: "telegram" })).toEqual([
|
||||
expect.objectContaining({
|
||||
target: "telegram:-1002223334444",
|
||||
threadId: "1",
|
||||
|
|
@ -190,7 +190,7 @@ describe("doctor Telegram General-topic conversation repair", () => {
|
|||
]);
|
||||
|
||||
await patchSessionEntryCore(scope, () => ({ displayName: "harmless later write" }));
|
||||
expect(listConversations(scope, { channel: "telegram" })).toEqual([
|
||||
expect(await listConversations(scope, { channel: "telegram" })).toEqual([
|
||||
expect.objectContaining({
|
||||
target: "telegram:-1002223334444",
|
||||
role: "primary",
|
||||
|
|
|
|||
|
|
@ -75,7 +75,10 @@ describe("conversation registry", () => {
|
|||
skillsSnapshot: { prompt: "Saved instructions. ".repeat(4096), skills: [] },
|
||||
});
|
||||
|
||||
const conversations = listConversations({ agentId: "main", storePath }, { channel: "reef" });
|
||||
const conversations = await listConversations(
|
||||
{ agentId: "main", storePath },
|
||||
{ channel: "reef" },
|
||||
);
|
||||
expect(conversations.map((entry) => entry.target).toSorted()).toEqual([
|
||||
"reef:peer-a",
|
||||
"reef:peer-b",
|
||||
|
|
@ -83,7 +86,10 @@ describe("conversation registry", () => {
|
|||
expect(conversations.every((entry) => entry.role === "participant")).toBe(true);
|
||||
expect(conversations.every((entry) => entry.sessionKey === scope.sessionKey)).toBe(true);
|
||||
expect(
|
||||
resolveCurrentSessionPrimaryConversation({ ...scope, sessionId: "shared-main-session" }),
|
||||
await resolveCurrentSessionPrimaryConversation({
|
||||
...scope,
|
||||
sessionId: "shared-main-session",
|
||||
}),
|
||||
).toBeUndefined();
|
||||
|
||||
const peerA = conversations.find((entry) => entry.target === "reef:peer-a");
|
||||
|
|
@ -93,7 +99,7 @@ describe("conversation registry", () => {
|
|||
);
|
||||
});
|
||||
|
||||
it("catalogs a directory address without inventing a model-context session", () => {
|
||||
it("catalogs a directory address without inventing a model-context session", async () => {
|
||||
const identity = buildConversationIdentity({
|
||||
channel: "reef",
|
||||
accountId: "default",
|
||||
|
|
@ -106,7 +112,10 @@ describe("conversation registry", () => {
|
|||
expect(identity).toBeDefined();
|
||||
registerConversationAddresses({ agentId: "main", storePath }, [identity!], 100);
|
||||
|
||||
const [conversation] = listConversations({ agentId: "main", storePath }, { channel: "reef" });
|
||||
const [conversation] = await listConversations(
|
||||
{ agentId: "main", storePath },
|
||||
{ channel: "reef" },
|
||||
);
|
||||
expect(conversation).toMatchObject({
|
||||
conversationRef: identity?.conversationRef,
|
||||
target: "reef:peer-a",
|
||||
|
|
@ -169,19 +178,19 @@ describe("conversation registry", () => {
|
|||
});
|
||||
|
||||
expect(committed.ok).toBe(true);
|
||||
const canonicalConversation = listConversations(scope).find(
|
||||
const canonicalConversation = (await listConversations(scope)).find(
|
||||
(conversation) => conversation.peerId === "canonical-ops",
|
||||
);
|
||||
expect(canonicalConversation).toBeDefined();
|
||||
const conversationRef = canonicalConversation!.conversationRef;
|
||||
expect(
|
||||
resolveCurrentSessionPrimaryConversation({ ...scope, sessionId: "ops-session" }),
|
||||
await resolveCurrentSessionPrimaryConversation({ ...scope, sessionId: "ops-session" }),
|
||||
).toEqual(canonicalConversation);
|
||||
expect(
|
||||
resolveCurrentSessionPrimaryConversation({ ...scope, sessionId: "another-session" }),
|
||||
await resolveCurrentSessionPrimaryConversation({ ...scope, sessionId: "another-session" }),
|
||||
).toBeUndefined();
|
||||
expect(
|
||||
resolveCurrentSessionPrimaryConversation({
|
||||
await resolveCurrentSessionPrimaryConversation({
|
||||
...scope,
|
||||
sessionId: "ops-session",
|
||||
sessionKey: "agent:main:discord:channel:another",
|
||||
|
|
@ -201,7 +210,7 @@ describe("conversation registry", () => {
|
|||
|
||||
await upsertCanonicalSessionEntry(scope, { label: "generic current write", updatedAt: 200 });
|
||||
expect(
|
||||
listConversations(scope).filter((conversation) => conversation.role === "primary"),
|
||||
(await listConversations(scope)).filter((conversation) => conversation.role === "primary"),
|
||||
).toEqual([expect.objectContaining({ conversationRef, peerId: "canonical-ops" })]);
|
||||
const afterCurrentWrite = resolveConversation({ agentId: "main", storePath }, conversationRef);
|
||||
expect(afterCurrentWrite).toMatchObject({
|
||||
|
|
@ -233,7 +242,7 @@ describe("conversation registry", () => {
|
|||
routeContextObserved: true,
|
||||
});
|
||||
expect(
|
||||
resolveCurrentSessionPrimaryConversation({ ...scope, sessionId: "ops-session" })
|
||||
(await resolveCurrentSessionPrimaryConversation({ ...scope, sessionId: "ops-session" }))
|
||||
?.routeContext,
|
||||
).toBeUndefined();
|
||||
await upsertCanonicalSessionEntry(scope, {
|
||||
|
|
@ -273,7 +282,7 @@ describe("conversation registry", () => {
|
|||
await writeRoute("channel:beta", "guild-beta", 200);
|
||||
|
||||
expect(
|
||||
listConversations(scope, { channel: "discord" })
|
||||
(await listConversations(scope, { channel: "discord" }))
|
||||
.map(({ target, routeContext }) => ({ target, routeContext }))
|
||||
.toSorted((left, right) => left.target.localeCompare(right.target)),
|
||||
).toEqual([
|
||||
|
|
@ -291,7 +300,7 @@ describe("conversation registry", () => {
|
|||
.where("role", "=", "related"),
|
||||
);
|
||||
expect(
|
||||
resolveCurrentSessionPrimaryConversation({ ...scope, sessionId: "shared-session" }),
|
||||
await resolveCurrentSessionPrimaryConversation({ ...scope, sessionId: "shared-session" }),
|
||||
).toMatchObject({ target: "channel:beta", routeContext: { guildId: "guild-beta" } });
|
||||
});
|
||||
|
||||
|
|
@ -327,16 +336,16 @@ describe("conversation registry", () => {
|
|||
storePath,
|
||||
});
|
||||
expect(rollover.ok).toBe(true);
|
||||
expect(listConversations(scope)[0]).toMatchObject({
|
||||
expect((await listConversations(scope))[0]).toMatchObject({
|
||||
sessionId: "after-rollover",
|
||||
routeContextObserved: true,
|
||||
routeContext: { guildId: "guild-a", memberRoleIds: ["support"] },
|
||||
});
|
||||
expect(
|
||||
resolveCurrentSessionPrimaryConversation({ ...scope, sessionId: "before-rollover" }),
|
||||
await resolveCurrentSessionPrimaryConversation({ ...scope, sessionId: "before-rollover" }),
|
||||
).toBeUndefined();
|
||||
expect(
|
||||
resolveCurrentSessionPrimaryConversation({ ...scope, sessionId: "after-rollover" }),
|
||||
await resolveCurrentSessionPrimaryConversation({ ...scope, sessionId: "after-rollover" }),
|
||||
).toMatchObject({ routeContext: { guildId: "guild-a" } });
|
||||
|
||||
snapshot = loadReplySessionInitializationSnapshot(scope);
|
||||
|
|
@ -350,13 +359,13 @@ describe("conversation registry", () => {
|
|||
snapshotEntry: snapshot.currentEntry,
|
||||
storePath,
|
||||
});
|
||||
expect(listConversations(scope)[0]).toMatchObject({
|
||||
expect((await listConversations(scope))[0]).toMatchObject({
|
||||
sessionId: "after-rollover",
|
||||
routeContextObserved: true,
|
||||
});
|
||||
expect(listConversations(scope)[0]?.routeContext).toBeUndefined();
|
||||
expect((await listConversations(scope))[0]?.routeContext).toBeUndefined();
|
||||
expect(
|
||||
resolveCurrentSessionPrimaryConversation({ ...scope, sessionId: "after-rollover" })
|
||||
(await resolveCurrentSessionPrimaryConversation({ ...scope, sessionId: "after-rollover" }))
|
||||
?.routeContext,
|
||||
).toBeUndefined();
|
||||
});
|
||||
|
|
@ -386,10 +395,10 @@ describe("conversation registry", () => {
|
|||
.where("session_key", "=", scope.sessionKey),
|
||||
);
|
||||
expect(
|
||||
resolveCurrentSessionPrimaryConversation({ ...scope, sessionId: "peer-a-session" }),
|
||||
await resolveCurrentSessionPrimaryConversation({ ...scope, sessionId: "peer-a-session" }),
|
||||
).toBeUndefined();
|
||||
if (invalid.entry_json !== undefined) {
|
||||
const [conversation] = listConversations(scope);
|
||||
const [conversation] = await listConversations(scope);
|
||||
expect(conversation).toMatchObject({ target: "reef:peer-a" });
|
||||
expect(conversation?.sessionId).toBeUndefined();
|
||||
expect(conversation?.sessionKey).toBeUndefined();
|
||||
|
|
@ -418,7 +427,7 @@ describe("conversation registry", () => {
|
|||
registerConversationAddresses({ agentId: "main", storePath }, [freshIdentity!], freshAt);
|
||||
|
||||
expect(
|
||||
listConversations({ agentId: "main", storePath }, { channel: "reef", limit: 1 }),
|
||||
await listConversations({ agentId: "main", storePath }, { channel: "reef", limit: 1 }),
|
||||
).toEqual([
|
||||
expect.objectContaining({
|
||||
conversationRef: freshIdentity?.conversationRef,
|
||||
|
|
@ -468,7 +477,7 @@ describe("conversation registry", () => {
|
|||
);
|
||||
|
||||
expect(
|
||||
listConversations({ agentId: "main", storePath }, { channel: "reef", limit: 1 })[0],
|
||||
(await listConversations({ agentId: "main", storePath }, { channel: "reef", limit: 1 }))[0],
|
||||
).toMatchObject({
|
||||
target: "reef:peer-a",
|
||||
sessionId: "live-session",
|
||||
|
|
@ -487,7 +496,10 @@ describe("conversation registry", () => {
|
|||
deliveryContext: { channel: "reef", accountId: "default", to: "reef:peer-a" },
|
||||
origin: { provider: "reef", accountId: "default", nativeDirectUserId: "peer-a" },
|
||||
});
|
||||
const [historical] = listConversations({ agentId: "main", storePath }, { channel: "reef" });
|
||||
const [historical] = await listConversations(
|
||||
{ agentId: "main", storePath },
|
||||
{ channel: "reef" },
|
||||
);
|
||||
expect(historical?.sessionId).toBe("old-session");
|
||||
|
||||
await upsertSessionEntry(scope, {
|
||||
|
|
@ -515,7 +527,7 @@ describe("conversation registry", () => {
|
|||
chatType: "direct",
|
||||
deliveryContext: { channel: "reef", accountId: "default", to: "reef:peer-a" },
|
||||
});
|
||||
const [linked] = listConversations({ agentId: "main", storePath }, { channel: "reef" });
|
||||
const [linked] = await listConversations({ agentId: "main", storePath }, { channel: "reef" });
|
||||
expect(linked?.sessionId).toBe("deleted-session");
|
||||
|
||||
await deleteSessionEntryLifecycle({
|
||||
|
|
@ -525,7 +537,7 @@ describe("conversation registry", () => {
|
|||
archiveTranscript: false,
|
||||
});
|
||||
expect(
|
||||
resolveCurrentSessionPrimaryConversation({ ...scope, sessionId: "deleted-session" }),
|
||||
await resolveCurrentSessionPrimaryConversation({ ...scope, sessionId: "deleted-session" }),
|
||||
).toBeUndefined();
|
||||
|
||||
expect(
|
||||
|
|
|
|||
|
|
@ -1,3 +1,4 @@
|
|||
import path from "node:path";
|
||||
import {
|
||||
withOpenClawAgentDatabaseReadOnly,
|
||||
type OpenClawAgentReadOnlyDatabase,
|
||||
|
|
@ -7,19 +8,25 @@ import {
|
|||
getOpenClawAgentDatabaseIfOpen,
|
||||
openOpenClawAgentDatabase,
|
||||
} from "../../state/openclaw-agent-db.js";
|
||||
import { resolveOpenClawAgentSqlitePath } from "../../state/openclaw-agent-db.paths.js";
|
||||
import {
|
||||
createOpenClawAgentDatabasePathMatcher,
|
||||
isIncognitoOpenClawAgentSqlitePath,
|
||||
resolveOpenClawAgentSqlitePath,
|
||||
} from "../../state/openclaw-agent-db.paths.js";
|
||||
import { captureOpenClawStateReadWorkerContext } from "../../state/openclaw-state-worker-context.js";
|
||||
import { resolveStateDir } from "../state-dir.js";
|
||||
import type { OpenClawConfig } from "../types.openclaw.js";
|
||||
import type { ConversationIdentity } from "./conversation-identity.js";
|
||||
import type { ConversationReadQuery, ConversationRecord } from "./conversation-registry.types.js";
|
||||
import { resolveSessionStorePathCore } from "./paths.js";
|
||||
import {
|
||||
selectConversationRowsFromDatabase,
|
||||
type ConversationRecord,
|
||||
} from "./session-accessor.sqlite-conversation-read.js";
|
||||
import { selectConversationRowsFromDatabase } from "./session-accessor.sqlite-conversation-read.js";
|
||||
import { upsertConversationIdentity } from "./session-accessor.sqlite-conversation.js";
|
||||
import { resolveSqliteReadScope, toDatabaseOptions } from "./session-accessor.sqlite-scope.js";
|
||||
import { withSessionStoreReaderInWorker } from "./session-entry-read-runtime.js";
|
||||
import { captureSessionStoreReadCandidates } from "./session-store-target-inventory.js";
|
||||
import { captureSessionTranscriptStorageEnvironment } from "./transcript-target-binding.js";
|
||||
|
||||
export type { ConversationRecord } from "./session-accessor.sqlite-conversation-read.js";
|
||||
export type { ConversationRecord } from "./conversation-registry.types.js";
|
||||
|
||||
export type ConversationRegistryScope = {
|
||||
agentId: string;
|
||||
|
|
@ -49,6 +56,81 @@ export function resolveConversationRegistryScope(params: {
|
|||
return pinConversationDatabaseScope(scope).scope;
|
||||
}
|
||||
|
||||
export async function prepareConversationRegistryScope(params: {
|
||||
agentId: string;
|
||||
config: OpenClawConfig;
|
||||
}): Promise<PreparedConversationRegistryScope> {
|
||||
const input = {
|
||||
agentId: params.agentId,
|
||||
storePath: resolveSessionStorePathCore(params.config.session?.store, {
|
||||
agentId: params.agentId,
|
||||
}),
|
||||
};
|
||||
if (isIncognitoOpenClawAgentSqlitePath(input.storePath, input)) {
|
||||
return pinConversationDatabaseScope(input).scope;
|
||||
}
|
||||
return withConversationRead(input, async ({ database, logicalAgentId }) => ({
|
||||
agentId: logicalAgentId,
|
||||
databaseAgentId: database.agentId,
|
||||
storePath: database.path,
|
||||
env: database.env,
|
||||
}));
|
||||
}
|
||||
|
||||
function withConversationRead<T>(
|
||||
input: ConversationRegistryScope,
|
||||
read: Parameters<typeof withSessionStoreReaderInWorker<T>>[1],
|
||||
): Promise<T> {
|
||||
const env = captureSessionTranscriptStorageEnvironment(input.env ?? process.env);
|
||||
const storePath = path.resolve(
|
||||
input.storePath ?? resolveSessionStorePathCore(undefined, { agentId: input.agentId, env }),
|
||||
);
|
||||
const context = captureOpenClawStateReadWorkerContext({ env });
|
||||
const source = createOpenClawAgentDatabasePathMatcher();
|
||||
for (const candidate of captureSessionStoreReadCandidates(storePath)) {
|
||||
source(candidate.path, candidate.path);
|
||||
}
|
||||
return withSessionStoreReaderInWorker(
|
||||
{ agentId: input.agentId, storePath, env },
|
||||
async (owner) => {
|
||||
if (input.databaseAgentId && owner.database.agentId !== input.databaseAgentId) {
|
||||
throw new Error("Conversation database owner changed. Retry the request.");
|
||||
}
|
||||
const result = await read(owner);
|
||||
owner.assertCurrent();
|
||||
return result;
|
||||
},
|
||||
{
|
||||
dataOnly: true,
|
||||
logical: {
|
||||
assertCurrent() {
|
||||
context.maintenanceScope?.assertAdmission();
|
||||
context.admission.assertCurrent();
|
||||
if (!source.isCurrent()) {
|
||||
throw new Error(
|
||||
"Session store changed while reading conversations. Retry the request.",
|
||||
);
|
||||
}
|
||||
},
|
||||
},
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
function selectConversationRowsInWorker(
|
||||
scope: ConversationRegistryScope,
|
||||
query: ConversationReadQuery,
|
||||
): Promise<ConversationRecord[]> {
|
||||
const capturedQuery = structuredClone(query);
|
||||
if (scope.storePath && isIncognitoOpenClawAgentSqlitePath(scope.storePath, scope)) {
|
||||
// Process-held databases retain their native owner until the incognito cutover.
|
||||
return Promise.resolve(selectConversationRows(scope, capturedQuery));
|
||||
}
|
||||
return withConversationRead(scope, ({ reader, database }) =>
|
||||
reader.readConversations({ query: capturedQuery, env: database.env }),
|
||||
);
|
||||
}
|
||||
|
||||
export function pinConversationDatabaseScope(input: ConversationRegistryScope) {
|
||||
const env = { ...(input.env ?? process.env) };
|
||||
env.OPENCLAW_STATE_DIR = resolveStateDir(env);
|
||||
|
|
@ -118,8 +200,15 @@ export function registerConversationAddresses(
|
|||
export function listConversations(
|
||||
scope: ConversationRegistryScope,
|
||||
options: { channel?: string; limit?: number } = {},
|
||||
): ConversationRecord[] {
|
||||
return selectConversationRows(scope, options);
|
||||
): Promise<ConversationRecord[]> {
|
||||
return selectConversationRowsInWorker(scope, options);
|
||||
}
|
||||
|
||||
export async function readConversation(
|
||||
scope: ConversationRegistryScope,
|
||||
conversationRef: string,
|
||||
): Promise<ConversationRecord | undefined> {
|
||||
return (await selectConversationRowsInWorker(scope, { conversationRef, limit: 1 }))[0];
|
||||
}
|
||||
|
||||
/** Resolves an opaque address to one exact channel target and its context binding, when present. */
|
||||
|
|
@ -151,10 +240,12 @@ export function resolveCurrentConversationSession(
|
|||
}
|
||||
|
||||
/** Reads only the primary address bound to this exact current session window. */
|
||||
export function resolveCurrentSessionPrimaryConversation(
|
||||
export async function resolveCurrentSessionPrimaryConversation(
|
||||
scope: ConversationRegistryScope & { sessionId: string; sessionKey: string },
|
||||
): ConversationRecord | undefined {
|
||||
const [conversation] = selectConversationRows(scope, { primarySession: scope });
|
||||
): Promise<ConversationRecord | undefined> {
|
||||
const [conversation] = await selectConversationRowsInWorker(scope, {
|
||||
primarySession: { sessionId: scope.sessionId, sessionKey: scope.sessionKey },
|
||||
});
|
||||
return conversation?.sessionId === scope.sessionId && conversation.sessionKey === scope.sessionKey
|
||||
? conversation
|
||||
: undefined;
|
||||
|
|
|
|||
43
src/config/sessions/conversation-registry.types.ts
Normal file
43
src/config/sessions/conversation-registry.types.ts
Normal file
|
|
@ -0,0 +1,43 @@
|
|||
import type { ConversationKind } from "./conversation-identity.js";
|
||||
import type { ConversationRouteContext } from "./conversation-route-context.js";
|
||||
|
||||
export type ConversationRecord = {
|
||||
conversationRef: string;
|
||||
channel: string;
|
||||
accountId: string;
|
||||
kind: ConversationKind;
|
||||
peerId: string;
|
||||
target: string;
|
||||
parentConversationRef?: string;
|
||||
threadId?: string;
|
||||
nativeChannelId?: string;
|
||||
nativeDirectUserId?: string;
|
||||
label?: string;
|
||||
sessionId?: string;
|
||||
sessionKey?: string;
|
||||
role?: "participant" | "primary" | "related";
|
||||
/** True when this address has been linked to a session in this agent store. */
|
||||
observedFromSession?: true;
|
||||
/** Exact contextual facts from the authoritative inbound route. */
|
||||
routeContext?: ConversationRouteContext;
|
||||
/** True when authoritative ingress observed empty or populated route context. */
|
||||
routeContextObserved?: true;
|
||||
firstSeenAt: number;
|
||||
lastSeenAt: number;
|
||||
};
|
||||
|
||||
export type ConversationReadQuery = {
|
||||
channel?: string;
|
||||
conversationRef?: string;
|
||||
limit?: number;
|
||||
primarySession?: { sessionId: string; sessionKey: string };
|
||||
currentBindingOnly?: boolean;
|
||||
currentSession?: { sessionKey: string; sessionId: string };
|
||||
};
|
||||
|
||||
export type ConversationRowsWorkerInput = {
|
||||
kind: "conversation-rows";
|
||||
database: { agentId: string; path: string };
|
||||
query: ConversationReadQuery;
|
||||
env: NodeJS.ProcessEnv;
|
||||
};
|
||||
|
|
@ -1,5 +1,5 @@
|
|||
import { sha256Hex } from "@openclaw/normalization-core/node-crypto";
|
||||
import type { ConversationRecord } from "./conversation-registry.js";
|
||||
import type { ConversationRecord } from "./conversation-registry.types.js";
|
||||
import {
|
||||
parseConversationRouteContext,
|
||||
type ConversationRouteContext,
|
||||
|
|
|
|||
|
|
@ -2,40 +2,12 @@ import { normalizeOptionalLowercaseString } from "@openclaw/normalization-core/s
|
|||
import { executeSqliteQuerySync, getNodeSqliteKysely } from "../../infra/kysely-sync.js";
|
||||
import type { DB as OpenClawAgentKyselyDatabase } from "../../state/openclaw-agent-db.generated.js";
|
||||
import type { OpenClawAgentDatabase } from "../../state/openclaw-agent-db.js";
|
||||
import type { ConversationKind } from "./conversation-identity.js";
|
||||
import {
|
||||
parseStoredConversationRouteContext,
|
||||
type ConversationRouteContext,
|
||||
} from "./conversation-route-context.js";
|
||||
import type { ConversationReadQuery, ConversationRecord } from "./conversation-registry.types.js";
|
||||
import { parseStoredConversationRouteContext } from "./conversation-route-context.js";
|
||||
import { parseSessionEntryJson } from "./session-accessor.sqlite-status.js";
|
||||
|
||||
const CONVERSATION_REF_PATTERN = /^conv_[a-f0-9]{32}$/u;
|
||||
|
||||
export type ConversationRecord = {
|
||||
conversationRef: string;
|
||||
channel: string;
|
||||
accountId: string;
|
||||
kind: ConversationKind;
|
||||
peerId: string;
|
||||
target: string;
|
||||
parentConversationRef?: string;
|
||||
threadId?: string;
|
||||
nativeChannelId?: string;
|
||||
nativeDirectUserId?: string;
|
||||
label?: string;
|
||||
sessionId?: string;
|
||||
sessionKey?: string;
|
||||
role?: "participant" | "primary" | "related";
|
||||
/** True when this address has been linked to a session in this agent store. */
|
||||
observedFromSession?: true;
|
||||
/** Exact contextual facts from the authoritative inbound route. */
|
||||
routeContext?: ConversationRouteContext;
|
||||
/** True when authoritative ingress observed empty or populated route context. */
|
||||
routeContextObserved?: true;
|
||||
firstSeenAt: number;
|
||||
lastSeenAt: number;
|
||||
};
|
||||
|
||||
function normalizeConversationRef(value: string): string {
|
||||
const normalized = value.trim().toLowerCase();
|
||||
if (!CONVERSATION_REF_PATTERN.test(normalized)) {
|
||||
|
|
@ -123,14 +95,7 @@ function mapConversationRow(row: {
|
|||
|
||||
export function selectConversationRowsFromDatabase(
|
||||
database: Pick<OpenClawAgentDatabase, "db">,
|
||||
options: {
|
||||
channel?: string;
|
||||
conversationRef?: string;
|
||||
limit?: number;
|
||||
primarySession?: { sessionId: string; sessionKey: string };
|
||||
currentBindingOnly?: boolean;
|
||||
currentSession?: { sessionKey: string; sessionId: string };
|
||||
} = {},
|
||||
options: ConversationReadQuery = {},
|
||||
): ConversationRecord[] {
|
||||
const db = getNodeSqliteKysely<
|
||||
Pick<
|
||||
|
|
|
|||
|
|
@ -329,7 +329,7 @@ test.each(["directory discovery", "Gateway send", "durable completion"] as const
|
|||
expect(result).toMatchObject({ kind: "lifecycle-artifacts", value: { removedEntries: 1 } });
|
||||
expect(loadSessionEntryReadOnly(scopes[0]!)).toBeUndefined();
|
||||
if (operation === "directory discovery") {
|
||||
expect(listConversations(scope)).toEqual([
|
||||
expect(await listConversations(scope)).toEqual([
|
||||
expect.objectContaining({
|
||||
conversationRef: conversation.conversationRef,
|
||||
target: conversation.target,
|
||||
|
|
|
|||
|
|
@ -14,7 +14,7 @@ import type {
|
|||
SessionTranscriptSummaryResult,
|
||||
} from "../../gateway/session-transcript-summary.js";
|
||||
import type { AgentHistoryActivity } from "../../infra/agent-activity-events.js";
|
||||
import type { ConversationRecord } from "./conversation-registry.js";
|
||||
import type { ConversationRecord } from "./conversation-registry.types.js";
|
||||
import type {
|
||||
SessionTranscriptDisplayDeltaResult,
|
||||
SessionTranscriptMessageByIdOptions,
|
||||
|
|
|
|||
|
|
@ -53,6 +53,12 @@ export function createSessionHistoryWorkerReaders(
|
|||
);
|
||||
}
|
||||
return {
|
||||
readConversations: reader(
|
||||
"conversation-rows",
|
||||
"conversations",
|
||||
(input) => ({ kind: "conversation-rows", ...input }),
|
||||
(value) => value.rows,
|
||||
),
|
||||
prewarm: reader(
|
||||
"prewarm",
|
||||
"prewarm acknowledgement",
|
||||
|
|
|
|||
|
|
@ -22,6 +22,10 @@ import type {
|
|||
SessionActivitySummaryBatchResult,
|
||||
} from "./activity-summary-source.types.js";
|
||||
import type { ConversationDeliveryRecord } from "./conversation-delivery-store.types.js";
|
||||
import type {
|
||||
ConversationRowsWorkerInput,
|
||||
ConversationRecord,
|
||||
} from "./conversation-registry.types.js";
|
||||
import type {
|
||||
ArchivedSessionEvictionBatch,
|
||||
ArchivedSessionEvictionQuery,
|
||||
|
|
@ -499,6 +503,7 @@ export type SessionHistoryWorkerInput =
|
|||
| SessionProgressCardWorkerInput
|
||||
| SessionPendingInputReceiptsWorkerInput
|
||||
| SessionGoalOperationReceiptWorkerInput
|
||||
| ConversationRowsWorkerInput
|
||||
| ConversationDeliveryWorkerInput
|
||||
| SessionEntryListWorkerInput
|
||||
| SessionEntryReadWorkerInput
|
||||
|
|
@ -529,6 +534,7 @@ export type SessionHistoryWorkerPreparedInput =
|
|||
PreparedHistoryInput<SessionHistoryDatabaseWorkerInput>;
|
||||
|
||||
export type SessionTranscriptWorkerValues = SessionTranscriptInventoryWorkerValues & {
|
||||
"conversation-rows": { kind: "conversation-rows"; rows: ConversationRecord[] };
|
||||
"conversation-delivery": { kind: "conversation-delivery"; record?: ConversationDeliveryRecord };
|
||||
prewarm: { kind: "prewarm" };
|
||||
"session-pending-archives": { kind: "session-pending-archives"; pending: boolean };
|
||||
|
|
@ -646,6 +652,7 @@ type CancellableSessionHistoryReader<
|
|||
> = (input: Omit<Input, "kind" | "database">, signal?: AbortSignal) => Promise<Value>;
|
||||
|
||||
export type SessionHistoryWorkerDatabase = SessionTranscriptInventoryReaders & {
|
||||
readConversations: SessionHistoryReader<ConversationRowsWorkerInput, ConversationRecord[]>;
|
||||
prewarm: (input: { env: NodeJS.ProcessEnv }) => Promise<void>;
|
||||
readPendingArchives: CancellableSessionHistoryReader<SessionPendingArchivesWorkerInput, boolean>;
|
||||
findTranscriptEvent: (
|
||||
|
|
|
|||
|
|
@ -454,6 +454,17 @@ serveOwnedWorkerTasks(
|
|||
),
|
||||
};
|
||||
}
|
||||
if (request.kind === "conversation-rows") {
|
||||
const { withOpenClawAgentDatabaseReadOnly } =
|
||||
await import("../../state/openclaw-agent-db-readonly.js");
|
||||
const { selectConversationRowsFromDatabase } =
|
||||
await import("./session-accessor.sqlite-conversation-read.js");
|
||||
const read = withOpenClawAgentDatabaseReadOnly(
|
||||
(database) => selectConversationRowsFromDatabase(database, request.query),
|
||||
{ ...request.database, env: cloneEnvWithPlatformSemantics(request.env) },
|
||||
);
|
||||
return { kind: "conversation-rows", rows: read.found ? read.value : [] };
|
||||
}
|
||||
if (request.kind === "conversation-delivery") {
|
||||
const { withOpenClawAgentDatabaseReadOnly } =
|
||||
await import("../../state/openclaw-agent-db-readonly.js");
|
||||
|
|
|
|||
|
|
@ -127,7 +127,7 @@ describe("conversation listings after removing accounts", () => {
|
|||
),
|
||||
);
|
||||
registerConversationAddresses(scope, identities);
|
||||
const before = listConversations(scope);
|
||||
const before = await listConversations(scope);
|
||||
expect(before).toHaveLength(identities.length);
|
||||
|
||||
const result = await runGatewayConversationList({ config, agentId: "main", limit: 10 }, deps);
|
||||
|
|
@ -138,7 +138,7 @@ describe("conversation listings after removing accounts", () => {
|
|||
accountId: conversation.accountId,
|
||||
})),
|
||||
).toEqual(LIVE_CHANNELS.map(({ channel }) => ({ channel, accountId: LIVE_ACCOUNT_ID })));
|
||||
expect(listConversations(scope)).toEqual(before);
|
||||
expect(await listConversations(scope)).toEqual(before);
|
||||
});
|
||||
|
||||
it("searches the active Discord account and preserves inactive history", async () => {
|
||||
|
|
@ -159,7 +159,7 @@ describe("conversation listings after removing accounts", () => {
|
|||
),
|
||||
);
|
||||
registerConversationAddresses(scope, identities);
|
||||
const before = listConversations(scope);
|
||||
const before = await listConversations(scope);
|
||||
|
||||
const result = await runGatewayConversationList(
|
||||
{ config, agentId: "main", channel: "discord", query: "123456789012345678", limit: 10 },
|
||||
|
|
@ -173,7 +173,7 @@ describe("conversation listings after removing accounts", () => {
|
|||
target: "channel:123456789012345678",
|
||||
}),
|
||||
]);
|
||||
expect(listConversations(scope)).toEqual(before);
|
||||
expect(await listConversations(scope)).toEqual(before);
|
||||
expect(before).toHaveLength(3);
|
||||
});
|
||||
});
|
||||
|
|
|
|||
|
|
@ -101,7 +101,7 @@ describe("conversation directory write admission", () => {
|
|||
await fixture.routed.promise;
|
||||
await setImmediate();
|
||||
expect(settled).toBe(false);
|
||||
expect(listConversations(fixture.scope)).toEqual([]);
|
||||
expect(await listConversations(fixture.scope)).toEqual([]);
|
||||
if (eligibility === "denied") {
|
||||
currentConfig = {
|
||||
...fixture.config,
|
||||
|
|
@ -129,7 +129,9 @@ describe("conversation directory write admission", () => {
|
|||
},
|
||||
});
|
||||
}
|
||||
expect(listConversations(fixture.scope)).toHaveLength(eligibility === "eligible" ? 1 : 0);
|
||||
expect(await listConversations(fixture.scope)).toHaveLength(
|
||||
eligibility === "eligible" ? 1 : 0,
|
||||
);
|
||||
expect(listSessionEntriesCore(fixture.scope)).toEqual([]);
|
||||
} finally {
|
||||
release.resolve();
|
||||
|
|
@ -173,7 +175,7 @@ describe("conversation directory write admission", () => {
|
|||
await expect(result).resolves.toMatchObject({
|
||||
conversations: [expect.objectContaining({ target: "reef:peer" })],
|
||||
});
|
||||
expect(listConversations(fixture.scope)).toHaveLength(1);
|
||||
expect(await listConversations(fixture.scope)).toHaveLength(1);
|
||||
expect(listSessionEntriesCore(fixture.scope)).toEqual([]);
|
||||
} finally {
|
||||
releaseRoute.resolve();
|
||||
|
|
|
|||
|
|
@ -11,7 +11,7 @@ import {
|
|||
import {
|
||||
listConversations,
|
||||
registerConversationAddresses,
|
||||
resolveConversationRegistryScope,
|
||||
prepareConversationRegistryScope,
|
||||
runConversationDatabaseWrite,
|
||||
type ConversationRecord,
|
||||
type ConversationRegistryScope,
|
||||
|
|
@ -239,7 +239,7 @@ export async function runGatewayConversationList(
|
|||
},
|
||||
deps: ConversationListDeps = defaultDeps,
|
||||
): Promise<ConversationListResult> {
|
||||
const scope = resolveConversationRegistryScope(params);
|
||||
const scope = await prepareConversationRegistryScope(params);
|
||||
const query = params.query?.trim() || undefined;
|
||||
const discovery = params.channel
|
||||
? await discoverChannelAddresses({
|
||||
|
|
@ -253,7 +253,7 @@ export async function runGatewayConversationList(
|
|||
...(params.readCurrentConfig ? { readCurrentConfig: params.readCurrentConfig } : {}),
|
||||
})
|
||||
: undefined;
|
||||
const conversations = deps.listConversations(
|
||||
const conversations = await deps.listConversations(
|
||||
scope,
|
||||
discovery ? { channel: discovery.channel } : {},
|
||||
);
|
||||
|
|
|
|||
219
src/gateway/conversation-list.worker-read.test.ts
Normal file
219
src/gateway/conversation-list.worker-read.test.ts
Normal file
|
|
@ -0,0 +1,219 @@
|
|||
import fs from "node:fs/promises";
|
||||
import { afterAll, afterEach, beforeAll, expect, it, vi } from "vitest";
|
||||
import { awaitGateBeforeSettlement } from "../../test/helpers/promise.js";
|
||||
import {
|
||||
emptySqliteCounts,
|
||||
observeParentSqlite,
|
||||
} from "../../test/helpers/sqlite-parent-observer.js";
|
||||
import {
|
||||
listConversations,
|
||||
prepareConversationRegistryScope,
|
||||
readConversation,
|
||||
resolveCurrentSessionPrimaryConversation,
|
||||
} from "../config/sessions/conversation-registry.js";
|
||||
import { replaceSessionEntrySync } from "../config/sessions/session-accessor.sqlite-entry.js";
|
||||
import { historyLane } from "../config/sessions/session-transcript-worker-resources.js";
|
||||
import type { OpenClawConfig } from "../config/types.openclaw.js";
|
||||
import { createDeferredCore } from "../shared/deferred.js";
|
||||
import { openOpenClawAgentDatabase } from "../state/openclaw-agent-db.js";
|
||||
import { resolveIncognitoOpenClawAgentSqlitePath } from "../state/openclaw-agent-db.paths.js";
|
||||
import {
|
||||
createOpenClawTestState,
|
||||
type OpenClawTestState,
|
||||
} from "../test-utils/openclaw-test-state.js";
|
||||
import { runGatewayConversationList } from "./conversation-list.js";
|
||||
|
||||
let state: OpenClawTestState;
|
||||
let storePath: string;
|
||||
const sessionKey = "agent:main:reef:channel:room";
|
||||
const sessionId = "conversation-read-session";
|
||||
const config = (): OpenClawConfig => ({
|
||||
agents: { entries: { main: { default: true }, other: {} } },
|
||||
session: { store: storePath },
|
||||
});
|
||||
const scope = () => ({ agentId: "main", storePath, sessionKey, sessionId });
|
||||
|
||||
beforeAll(async () => {
|
||||
state = await createOpenClawTestState({ scenario: "minimal" });
|
||||
storePath = state.statePath("shared-conversations.sqlite");
|
||||
openOpenClawAgentDatabase({ agentId: "physical-owner", path: storePath });
|
||||
replaceSessionEntrySync(scope(), {
|
||||
sessionId,
|
||||
updatedAt: 100,
|
||||
chatType: "channel",
|
||||
delivery: {
|
||||
kind: "external",
|
||||
route: {
|
||||
channel: "reef",
|
||||
accountId: "default",
|
||||
target: { to: "reef:room", chatType: "channel" },
|
||||
},
|
||||
context: { channel: "reef", accountId: "default", to: "reef:room" },
|
||||
origin: { provider: "reef", accountId: "default" },
|
||||
},
|
||||
});
|
||||
});
|
||||
afterEach(() => vi.restoreAllMocks());
|
||||
afterAll(async () => state.cleanup());
|
||||
|
||||
it("reads a shared store's list, exact address and primary binding without caller-thread SQLite", async () => {
|
||||
const observer = observeParentSqlite();
|
||||
try {
|
||||
const result = await runGatewayConversationList({
|
||||
config: config(),
|
||||
agentId: "main",
|
||||
limit: 10,
|
||||
});
|
||||
expect(result.conversations).toEqual([
|
||||
expect.objectContaining({ channel: "reef", target: "reef:room" }),
|
||||
]);
|
||||
const ref = result.conversations[0]!.conversationRef;
|
||||
const exact = await readConversation(scope(), ref);
|
||||
expect(exact).toMatchObject({ conversationRef: ref, sessionId, sessionKey, role: "primary" });
|
||||
expect(await resolveCurrentSessionPrimaryConversation(scope())).toEqual(exact);
|
||||
const missing = state.statePath("absent", "openclaw-agent.sqlite");
|
||||
expect(await listConversations({ agentId: "main", storePath: missing })).toEqual([]);
|
||||
expect(observer.counts).toEqual(emptySqliteCounts());
|
||||
await expect(fs.stat(missing)).rejects.toMatchObject({ code: "ENOENT" });
|
||||
} finally {
|
||||
observer.restore();
|
||||
}
|
||||
});
|
||||
|
||||
it("rechecks current route ownership after the worker reply is delayed", async () => {
|
||||
const entered = createDeferredCore();
|
||||
const release = createDeferredCore();
|
||||
const run = historyLane.pool.run.bind(historyLane.pool);
|
||||
vi.spyOn(historyLane.pool, "run").mockImplementation(async (...args) => {
|
||||
const reply = await run(...args);
|
||||
if (
|
||||
reply.ok &&
|
||||
typeof reply.value === "object" &&
|
||||
!Array.isArray(reply.value) &&
|
||||
"kind" in reply.value &&
|
||||
reply.value.kind === "conversation-rows"
|
||||
) {
|
||||
entered.resolve();
|
||||
await release.promise;
|
||||
}
|
||||
return reply;
|
||||
});
|
||||
let current = config();
|
||||
const pending = runGatewayConversationList({
|
||||
config: current,
|
||||
readCurrentConfig: () => current,
|
||||
agentId: "main",
|
||||
limit: 10,
|
||||
});
|
||||
const outcome = pending.catch((error: unknown) => error);
|
||||
try {
|
||||
await awaitGateBeforeSettlement(
|
||||
entered.promise,
|
||||
pending,
|
||||
"Conversation read was not dispatched",
|
||||
);
|
||||
current = {
|
||||
...current,
|
||||
bindings: [{ type: "route", agentId: "other", match: { channel: "reef" } }],
|
||||
};
|
||||
release.resolve();
|
||||
await expect(pending).resolves.toEqual({ conversations: [] });
|
||||
} finally {
|
||||
release.resolve();
|
||||
await outcome;
|
||||
}
|
||||
});
|
||||
|
||||
it("propagates worker rejection without reading SQLite on the caller", async () => {
|
||||
const failure = new Error("conversation worker refused");
|
||||
vi.spyOn(historyLane.pool, "run").mockRejectedValueOnce(failure);
|
||||
const observer = observeParentSqlite();
|
||||
try {
|
||||
await expect(listConversations({ ...scope(), databaseAgentId: "physical-owner" })).rejects.toBe(
|
||||
failure,
|
||||
);
|
||||
expect(observer.counts).toEqual(emptySqliteCounts());
|
||||
} finally {
|
||||
observer.restore();
|
||||
}
|
||||
});
|
||||
|
||||
it.each(["conversation-rows", "session-store-target"] as const)(
|
||||
"refuses a replaced source while %s is pending",
|
||||
async (phase) => {
|
||||
const directory = state.statePath(phase);
|
||||
const missing = state.statePath(phase, "openclaw-agent.sqlite");
|
||||
const locator =
|
||||
phase === "session-store-target" ? state.statePath(phase, "sessions.json") : missing;
|
||||
await fs.mkdir(directory, { recursive: true });
|
||||
const entered = createDeferredCore();
|
||||
const release = createDeferredCore();
|
||||
const run = historyLane.pool.run.bind(historyLane.pool);
|
||||
vi.spyOn(historyLane.pool, "run").mockImplementation(async (...args) => {
|
||||
const reply = await run(...args);
|
||||
if (
|
||||
reply.ok &&
|
||||
typeof reply.value === "object" &&
|
||||
!Array.isArray(reply.value) &&
|
||||
"kind" in reply.value &&
|
||||
reply.value.kind === phase
|
||||
) {
|
||||
entered.resolve();
|
||||
await release.promise;
|
||||
}
|
||||
return reply;
|
||||
});
|
||||
const pending = listConversations({ agentId: "main", storePath: locator });
|
||||
const outcome = pending.catch((error: unknown) => error);
|
||||
try {
|
||||
await awaitGateBeforeSettlement(
|
||||
entered.promise,
|
||||
pending,
|
||||
"Conversation read was not dispatched",
|
||||
);
|
||||
await fs.writeFile(missing, "replacement source");
|
||||
release.resolve();
|
||||
await expect(pending).rejects.toThrow(/Session store changed/);
|
||||
} finally {
|
||||
release.resolve();
|
||||
await outcome;
|
||||
await fs.unlink(missing).catch(() => {});
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
it("keeps process-held incognito conversations with their native owner", async () => {
|
||||
const native = {
|
||||
agentId: "main",
|
||||
storePath: resolveIncognitoOpenClawAgentSqlitePath({ agentId: "main", env: state.env }),
|
||||
sessionKey: "agent:main:dashboard:incognito-conversation-read",
|
||||
sessionId: "incognito-conversation-read",
|
||||
};
|
||||
replaceSessionEntrySync(native, {
|
||||
sessionId: native.sessionId,
|
||||
updatedAt: 100,
|
||||
chatType: "channel",
|
||||
delivery: {
|
||||
kind: "external",
|
||||
route: {
|
||||
channel: "reef",
|
||||
accountId: "default",
|
||||
target: { to: "reef:room", chatType: "channel" },
|
||||
},
|
||||
context: { channel: "reef", accountId: "default", to: "reef:room" },
|
||||
origin: { provider: "reef", accountId: "default" },
|
||||
},
|
||||
});
|
||||
const worker = vi.spyOn(historyLane.pool, "run");
|
||||
const prepared = await prepareConversationRegistryScope({
|
||||
agentId: "main",
|
||||
config: { session: { store: native.storePath } },
|
||||
});
|
||||
const rows = await listConversations(prepared);
|
||||
expect(rows).toEqual([
|
||||
expect.objectContaining({ sessionId: native.sessionId, sessionKey: native.sessionKey }),
|
||||
]);
|
||||
expect(await resolveCurrentSessionPrimaryConversation(native)).toEqual(rows[0]);
|
||||
expect(worker).not.toHaveBeenCalled();
|
||||
await expect(fs.stat(native.storePath)).rejects.toMatchObject({ code: "ENOENT" });
|
||||
});
|
||||
|
|
@ -54,6 +54,9 @@ function sentResult() {
|
|||
|
||||
function createDeps(agentId = "main") {
|
||||
const store = createConversationDeliveryTestStore(agentId);
|
||||
vi.spyOn(conversationRegistry, "readConversation").mockImplementation(async (scope, ref) =>
|
||||
conversationRegistry.resolveConversation(scope, ref),
|
||||
);
|
||||
return {
|
||||
...store,
|
||||
beginOperation: vi.spyOn(deliveryStore, "beginConversationDeliveryOperation"),
|
||||
|
|
|
|||
|
|
@ -5,8 +5,8 @@ import {
|
|||
type ConversationDeliveryRecord,
|
||||
} from "../config/sessions/conversation-delivery-store.js";
|
||||
import {
|
||||
resolveConversation,
|
||||
resolveConversationRegistryScope,
|
||||
readConversation,
|
||||
prepareConversationRegistryScope,
|
||||
} from "../config/sessions/conversation-registry.js";
|
||||
import { resolveConversationRouteFingerprint } from "../config/sessions/conversation-route-fingerprint.js";
|
||||
import type { OpenClawConfig } from "../config/types.openclaw.js";
|
||||
|
|
@ -36,7 +36,8 @@ export async function runGatewayConversationSend(params: {
|
|||
message: string;
|
||||
signal?: AbortSignal;
|
||||
}): Promise<ConversationSendResult> {
|
||||
const scope = resolveConversationRegistryScope(params);
|
||||
const scope = await prepareConversationRegistryScope(params);
|
||||
params.signal?.throwIfAborted();
|
||||
try {
|
||||
const operation: ConversationDeliveryRecord | undefined =
|
||||
await getConversationDeliveryOperation(scope, params.operationId, {
|
||||
|
|
@ -46,7 +47,8 @@ export async function runGatewayConversationSend(params: {
|
|||
message: params.message,
|
||||
});
|
||||
|
||||
const conversation = resolveConversation(scope, params.conversationRef);
|
||||
const conversation = await readConversation(scope, params.conversationRef);
|
||||
params.signal?.throwIfAborted();
|
||||
if (!conversation) {
|
||||
throw new ConversationInputError(
|
||||
`Conversation not found: ${params.conversationRef} (use conversations_list)`,
|
||||
|
|
|
|||
|
|
@ -57,6 +57,9 @@ function sentResult(messageId = "reef-outbound-1") {
|
|||
|
||||
function createDeps() {
|
||||
const store = createConversationDeliveryTestStore();
|
||||
vi.spyOn(conversationRegistry, "readConversation").mockImplementation(async (scope, ref) =>
|
||||
conversationRegistry.resolveConversation(scope, ref),
|
||||
);
|
||||
return {
|
||||
...store,
|
||||
beginOperation: vi.spyOn(deliveryStore, "beginConversationDeliveryOperation"),
|
||||
|
|
|
|||
|
|
@ -6,8 +6,8 @@ import {
|
|||
getConversationDeliveryOperation,
|
||||
} from "../config/sessions/conversation-delivery-store.js";
|
||||
import {
|
||||
resolveConversation,
|
||||
resolveConversationRegistryScope,
|
||||
readConversation,
|
||||
prepareConversationRegistryScope,
|
||||
type ConversationRecord,
|
||||
type ConversationRegistryScope,
|
||||
} from "../config/sessions/conversation-registry.js";
|
||||
|
|
@ -201,7 +201,7 @@ async function ensureConversationContextBinding(params: {
|
|||
},
|
||||
binding,
|
||||
);
|
||||
const bound = resolveConversation(params.scope, params.conversation.conversationRef);
|
||||
const bound = await readConversation(params.scope, params.conversation.conversationRef);
|
||||
if (!bound || !hasConversationSessionBinding(bound)) {
|
||||
throw new Error(
|
||||
`Conversation ${params.conversation.conversationRef} could not create its local context binding`,
|
||||
|
|
@ -222,7 +222,7 @@ export async function runGatewayConversationTurn(params: {
|
|||
message: string;
|
||||
timeoutMs: number;
|
||||
}): Promise<ConversationTurnResult> {
|
||||
const scope = resolveConversationRegistryScope(params);
|
||||
const scope = await prepareConversationRegistryScope(params);
|
||||
const binding = captureOutboundSessionBinding({
|
||||
cfg: params.config,
|
||||
scope,
|
||||
|
|
@ -244,7 +244,7 @@ export async function runGatewayConversationTurn(params: {
|
|||
throw error;
|
||||
}
|
||||
|
||||
const discoveredConversation = resolveConversation(scope, params.conversationRef);
|
||||
const discoveredConversation = await readConversation(scope, params.conversationRef);
|
||||
if (!discoveredConversation) {
|
||||
throw new ConversationInputError(
|
||||
`Conversation not found: ${params.conversationRef} (use conversations_list)`,
|
||||
|
|
|
|||
|
|
@ -98,7 +98,8 @@ export async function runGatewaySessionCompaction(
|
|||
entry: params.entry,
|
||||
cfg: params.cfg,
|
||||
});
|
||||
const primaryConversation = resolveCurrentSessionPrimaryConversation(transcriptTarget);
|
||||
const primaryConversation = await resolveCurrentSessionPrimaryConversation(transcriptTarget);
|
||||
params.abortSignal?.throwIfAborted();
|
||||
return await compactEmbeddedAgentSession(
|
||||
{
|
||||
abortSignal: params.abortSignal,
|
||||
|
|
|
|||
|
|
@ -399,7 +399,7 @@ describe("outbound session persistence", () => {
|
|||
established!.entry.updatedAt + 1,
|
||||
);
|
||||
|
||||
const discovered = listConversations({ agentId: "main", storePath }, { channel: "reef" });
|
||||
const discovered = await listConversations({ agentId: "main", storePath }, { channel: "reef" });
|
||||
expect(discovered[0]?.conversationRef).toBe(threadlessIdentity!.conversationRef);
|
||||
expect(discovered[0]).not.toMatchObject({ sessionId: expect.any(String) });
|
||||
|
||||
|
|
|
|||
|
|
@ -252,10 +252,12 @@ async function runAccountHistoryProof(compactionMode: "client" | "server-endpoin
|
|||
);
|
||||
const session = listed.sessions.find((entry) => entry.key === key);
|
||||
expect(session).toBeDefined();
|
||||
const conversation = listConversations({
|
||||
agentId: "main",
|
||||
storePath: path.join(instance.state.agentDir("main"), "openclaw-agent.sqlite"),
|
||||
}).find(
|
||||
const conversation = (
|
||||
await listConversations({
|
||||
agentId: "main",
|
||||
storePath: path.join(instance.state.agentDir("main"), "openclaw-agent.sqlite"),
|
||||
})
|
||||
).find(
|
||||
(entry) =>
|
||||
entry.sessionKey === key &&
|
||||
entry.sessionId === session!.sessionId &&
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue