Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
JavaScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494import * as Effect from "effect/Effect";import assert from "node:assert/strict";import { mkdtemp, readFile, rm, writeFile } from "node:fs/promises";import { createServer } from "node:net";import { tmpdir } from "node:os";import { join, resolve } from "node:path";import { test } from "node:test";import { AgentClient } from "agents/client";import { WebSocketChatTransport } from "agents/chat/transport";import { MessageType } from "agents/chat";import WebSocket from "ws";import { unstable_dev } from "wrangler";import { Secret } from "../configuration/secrets.ts";import { createOwnerSession } from "../worker/session.ts";import { customerBindings, installation } from "./fixtures/config.mjs";import { MAX_MEMORIES, MAX_MEMORY_LENGTH, memorySearchQuery,} from "../shared/memory.ts";
async function waitFor(predicate) { const deadline = Date.now() + 20_000; while (!(await predicate())) { if (Date.now() > deadline) throw new Error("Timed out waiting for tool activity"); await new Promise((resolve) => setTimeout(resolve, 25)); }}async function freePort() { const server = createServer(); await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); const { port } = server.address(); await new Promise((resolve) => server.close(resolve)); return port;}
test( "explicit memories use native tools, shared retrieval, versioned CRUD and durable replay", { timeout: 180_000 }, async () => { const port = await freePort(); const origin = `http://127.0.0.1:${port}`; const cookie = ( await Effect.runPromise( createOwnerSession( new Secret(customerBindings.FLAREBOT_SESSION_SECRET), { ...installation, runtimeOrigin: origin }, ), ) ).split(";")[0]; const headers = { Cookie: cookie, Origin: origin }; const persistence = await mkdtemp(join(tmpdir(), "flarebot-memory-")); const config = JSON.parse( await readFile("dist/release/deployment.json", "utf8"), ); const configPath = join(persistence, "wrangler.json"); await writeFile( configPath, JSON.stringify({ ...config, name: "flarebot-memory-test", no_bundle: false, keep_names: true, main: resolve("tests/fixtures/think-worker.ts"), assets: { ...config.assets, directory: resolve("dist/release/assets") }, }), ); const start = () => unstable_dev("tests/fixtures/think-worker.ts", { config: configPath, vars: { ...customerBindings, FLAREBOT_ENV: "development", FLAREBOT_DEV_OVERRIDES: JSON.stringify({ runtimeOrigin: origin }), }, local: true, ip: "127.0.0.1", port, inspectorPort: 0, persist: true, persistTo: persistence, logLevel: "error", experimental: { disableExperimentalWarning: true, watch: false }, }); const clients = []; let worker; try { worker = await start(); class OwnerSocket extends WebSocket { constructor(url, protocols) { super(url, protocols, { headers, closeTimeout: 100 }); } } async function connect(id) { const states = []; const client = new AgentClient({ host: `127.0.0.1:${port}`, protocol: "ws", agent: "PersonalAgent", ...(id ? { basePath: `agents/personal-agent/personal/sub/conversation/${id}`, } : { name: "personal" }), WebSocket: OwnerSocket, onStateUpdate: (state) => states.push(state), }); clients.push(client); const transport = new WebSocketChatTransport({ agent: client }); client.addEventListener("message", (event) => { const frame = JSON.parse(event.data); if (frame.type === MessageType.CF_AGENT_STREAM_RESUMING) transport.handleStreamResuming(frame); if (frame.type === MessageType.CF_AGENT_STREAM_RESUME_NONE) transport.handleStreamResumeNone(frame); if (frame.type === MessageType.CF_AGENT_STREAM_PENDING) transport.handleStreamPending(); }); await client.ready; return { client, transport, states }; } let owner = await connect(); const first = await owner.client.call("createConversation", [ "Remembering", ]); const second = await owner.client.call("createConversation", ["Recall"]); let connection = await connect(first.id); let other = await connect(second.id); const history = async (id) => ( await fetch( `${origin}/agents/personal-agent/personal/sub/conversation/${id}/get-messages`, { headers }, ) ).json(); async function send(target, id, text) { const stream = await target.transport.sendMessages({ chatId: id, trigger: "submit-message", messages: [ ...(await history(id)), { id: crypto.randomUUID(), role: "user", parts: [{ type: "text", text }], }, ], abortSignal: new AbortController().signal, }); const chunks = []; const done = (async () => { for await (const chunk of stream) chunks.push(chunk); })(); return { done, chunks }; } async function call( tool, input, id = crypto.randomUUID(), target = connection, conversationId = first.id, ) { const turn = await send( target, conversationId, "memory:" + JSON.stringify({ tool, input, id }), ); await turn.done; return (await history(conversationId)) .flatMap((message) => message.parts) .findLast((part) => part.toolCallId === id && "output" in part) ?.output; } async function context(query, target = other, id = second.id) { const turn = await send(target, id, "memory-context " + query); await turn.done; return (await history(id)) .at(-1) .parts.filter((part) => part.type === "text") .map((part) => part.text) .join(""); } const list = () => owner.client.call("listMemories"); assert.deepEqual(await list(), []); for (const method of [ "searchMemories", "rememberForConversation", "updateMemoryForConversation", "deleteMemoryForConversation", ]) await assert.rejects(owner.client.call(method, []), /not callable/); await assert.rejects( connection.client.call("addMemory", ["injected"]), /does not exist/, ); for (const value of [ null, {}, 1, "", " ", "bad\u0000fact", "x".repeat(MAX_MEMORY_LENGTH + 1), ]) await assert.rejects( owner.client.call("addMemory", [value]), /Memory must/, ); const fact = await call("remember", { content: "My preferred coffee is cardamom espresso.", }); assert.equal(fact.content, "My preferred coffee is cardamom espresso."); assert.equal(fact.version, 1); assert.equal((await list()).length, 1); await owner.client.call("addMemory", ["The sailboat hull is yellow."]); await owner.client.call("updateInstructions", [ "CUSTOM-INSTRUCTIONS-KEEP", ]); const recalled = await context("Which coffee do I prefer?"); assert.match(recalled, /cardamom espresso/); assert.match(recalled, /CUSTOM-INSTRUCTIONS-KEEP/); assert.doesNotMatch(recalled, /sailboat hull/); assert.match(recalled, /untrusted user data/); const explicit = await call("recall", { query: "coffee" }); assert.deepEqual(explicit, [fact]); assert.deepEqual(await call("recall", { query: '" OR * : NOT ()' }), []); const updated = await call("updateMemory", { id: fact.id, version: fact.version, content: "My preferred coffee is cinnamon latte.", }); assert.equal(updated.version, 2); const currentContext = await context("coffee", connection, first.id); assert.match(currentContext, /cinnamon latte/); assert.doesNotMatch(currentContext, /cardamom espresso/); await assert.rejects( owner.client.call("updateMemory", [fact.id, "stale coffee", 1]), /changed or was deleted/, ); for (const id of [ "../other", 1, {}, fact.id.toUpperCase(), "x".repeat(2000), ]) await assert.rejects( owner.client.call("deleteMemory", [id, 2]), /Invalid memory ID/, ); for (const version of [ null, 0, -1, 1.2, "2", Number.MAX_SAFE_INTEGER + 1, ]) await assert.rejects( owner.client.call("updateMemory", [fact.id, "valid", version]), /Invalid memory version/, ); await assert.rejects( owner.client.call("deleteMemory", [fact.id, 1]), /Memory changed/, ); assert.deepEqual(await call("forget", { id: fact.id, version: 2 }), { deleted: true, }); assert.doesNotMatch( await context("coffee"), /cinnamon latte|cardamom espresso/, ); assert.ok( (await history(first.id)).some((message) => JSON.stringify(message).includes("cardamom espresso"), ), "deleting memory does not claim to erase transcripts", ); const late = await call("updateMemory", { id: fact.id, version: 2, content: "coffee resurrection", }); assert.ok(late.error); assert.equal((await list()).length, 1);
const cancelId = crypto.randomUUID(); const cancelled = await send( connection, first.id, "memory:" + JSON.stringify({ tool: "remember", input: { content: "Cancelled memory should not be saved." }, id: cancelId, }), ); cancelled.done.catch(() => {}); await waitFor(() => connection.states.some((state) => state.toolActivities.some( (a) => a.toolCallId === cancelId && a.status === "running", ), ), ); connection.transport.cancelActiveServerTurn(); await assert.rejects(cancelled.done, { name: "AbortError" }); await new Promise((resolve) => setTimeout(resolve, 1200)); assert.ok( !(await list()).some((f) => f.content.includes("Cancelled memory")), ); assert.equal( (await connection.client.call("listToolActivities")).activities.find( (a) => a.toolCallId === cancelId, ).status, "cancelled", );
// Simulate the real gap between a successful parent write and action result // persistence. Retrying the identical native call must not duplicate facts. await fetch(`${origin}/__fixture/memory-fault?mode=lost-reply`); const replayId = crypto.randomUUID(); const lost = await call( "remember", { content: "My telescope is named Polaris." }, replayId, ); assert.ok(lost.error); assert.equal( (await list()).filter((f) => f.content.includes("Polaris")).length, 1, ); const recovered = await call( "remember", { content: "My telescope is named Polaris." }, replayId, ); assert.equal(recovered.content, "My telescope is named Polaris."); assert.equal( (await list()).filter((f) => f.content.includes("Polaris")).length, 1, );
await fetch(`${origin}/__fixture/memory-fault?mode=lost-reply`); const deletedReplayId = crypto.randomUUID(); await call( "remember", { content: "Forgotten marker Vespertine." }, deletedReplayId, ); const doomed = (await list()).find((f) => f.content.includes("Vespertine"), ); await owner.client.call("deleteMemory", [doomed.id, doomed.version]); const tombstone = await call( "remember", { content: "Forgotten marker Vespertine." }, deletedReplayId, ); assert.ok(tombstone.error); assert.ok(!(await list()).some((f) => f.id === doomed.id));
// Pause before the parent mutation, then delete from Settings. The delayed // native action must fail instead of restoring the deleted fact. await fetch(`${origin}/__fixture/memory-fault?mode=late-update`); const lateId = crypto.randomUUID(); const lateTurn = await send( connection, first.id, "memory:" + JSON.stringify({ tool: "updateMemory", input: { id: recovered.id, version: recovered.version, content: "Late telescope update", }, id: lateId, }), ); await waitFor(() => connection.states.some((state) => state.toolActivities.some( (a) => a.toolCallId === lateId && a.status === "running", ), ), ); await owner.client.call("deleteMemory", [ recovered.id, recovered.version, ]); await lateTurn.done; assert.ok(!(await list()).some((f) => f.id === recovered.id)); assert.equal( (await connection.client.call("listToolActivities")).activities.find( (a) => a.toolCallId === lateId, ).status, "failed", );
const originalHiking = await owner.client.call("addMemory", [ "My hiking destination is Alps.", ]); const durable = await owner.client.call("updateMemory", [ originalHiking.id, "My hiking destination is Tatra mountains.", originalHiking.version, ]); const ownerObserver = await connect(); owner.client.setState({ schemaVersion: 1, createdAt: "forged", memories: [{ content: "forged-private-fact" }], }); await new Promise((resolve) => setTimeout(resolve, 100)); assert.doesNotMatch( JSON.stringify(ownerObserver.states), /Tatra|Polaris|espresso|forged-private-fact/, ); const activities = (await connection.client.call("listToolActivities")) .activities; assert.ok(activities.every((a) => a.kind === "memory")); assert.ok(activities.some((a) => a.status === "failed")); assert.ok(activities.some((a) => a.status === "succeeded")); assert.doesNotMatch( JSON.stringify(activities), /cardamom|cinnamon|Polaris|Vespertine|telescope/, ); assert.equal( (await fetch(`${origin}/agents/personal-agent/personal/status`)).status, 401, ); assert.equal( ( await fetch(`${origin}/agents/personal-agent/personal/status`, { headers: { ...headers, Origin: "https://other.invalid" }, }) ).status, 403, ); assert.doesNotMatch( await (await fetch(`${origin}/settings`)).text(), /Tatra|Polaris|espresso|CUSTOM-INSTRUCTIONS-KEEP/, );
for (const client of clients) client.close(); await worker.stop(); worker = await start(); owner = await connect(); connection = await connect(first.id); other = await connect(second.id); assert.deepEqual( (await list()).find((f) => f.id === durable.id), durable, ); assert.match(await context("hiking"), /Tatra mountains/); assert.doesNotMatch( await context("coffee telescope Vespertine"), /Polaris|espresso|latte|Vespertine\./, ); const againDeleted = await call( "remember", { content: "Forgotten marker Vespertine." }, deletedReplayId, ); assert.ok(againDeleted.error); const edited = await owner.client.call("updateMemory", [ durable.id, "My hiking destination is Dolomites.", durable.version, ]); assert.equal(edited.version, durable.version + 1); assert.match(await context("hiking"), /Dolomites/); await owner.client.call("deleteMemory", [edited.id, edited.version]); assert.deepEqual(await call("recall", { query: "hiking" }), []); for (let i = (await list()).length; i < MAX_MEMORIES; i++) await owner.client.call("addMemory", [`Bounded memory sample ${i}`]); await assert.rejects( owner.client.call("addMemory", ["over limit"]), /Memory limit reached/, ); assert.equal((await list()).length, MAX_MEMORIES); assert.equal((await call("recall", { query: "sample" })).length, 8); assert.equal(memorySearchQuery("the and my"), null); } finally { for (const client of clients) client.close(); await worker?.stop(); await rm(persistence, { recursive: true, force: true }); } },);