Skip to content

memory.ts ​

Source: docs/worker-api/examples/memory.ts · Download source · Example guide

ts
import { createDitto, graph, loadRuntimeConfigFile } from "@codesoul-co/ditto";
import {
  createMemory, createMemoryWorker, MemoryError, memoryGetNode,
  type MemoryResources, type MemoryStore, type MemorySearchProvider,
} from "@codesoul-co/ditto/worker/memory";

// example: setup
export function setupMemory(resources: MemoryResources) {
  const config = loadRuntimeConfigFile("ditto.yaml", process.env);
  const worker = createMemoryWorker({ ...resources, concurrency: 16 });
  const runtime = createDitto({ config, workers: [worker] });
  const memory = createMemory({ ...resources, defaults: config.memory });
  return { runtime, memory };
}

// example: get
export async function getMemory(resources: MemoryResources) {
  const memory = createMemory(resources);
  const result = await memory.get({ ids: ["m1", "m1"], keys: ["preference"] });
  if (result.status !== "success" || !result.output) throw new Error(result.error?.code ?? result.status);
  return result.output.map(item => ({ id: item.id, content: item.content }));
}

// example: query
export async function queryMemory(resources: MemoryResources) {
  const memory = createMemory(resources);
  const first = await memory.query({ limit: 20, orderBy: [{ field: "id", direction: "asc" }] });
  if (first.status !== "success" || !first.output) throw new Error(first.error?.code ?? first.status);
  const items = [...first.output.items];
  if (first.output.nextCursor) {
    const next = await memory.query({ limit: 20, orderBy: [{ field: "id", direction: "asc" }], cursor: first.output.nextCursor });
    if (next.status !== "success" || !next.output) throw new Error(next.error?.code ?? next.status);
    items.push(...next.output.items);
  }
  return items;
}

// example: search
export async function searchMemory(resources: MemoryResources) {
  const memory = createMemory(resources);
  const result = await memory.search({ query: "preferred language", limit: 5 });
  if (result.status !== "success" || !result.output) throw new Error(result.error?.code ?? result.status);
  return result.output.map(hit => ({ id: hit.memory.id, content: hit.memory.content, score: hit.score }));
}

// example: write
export async function writeMemory(resources: MemoryResources) {
  const memory = createMemory(resources);
  const result = await memory.write({ memories: [{ key: "preference", content: { language: "zh-CN" }, metadata: { source: "user" } }] });
  if (result.status !== "success" || !result.output) throw new Error(result.error?.code ?? result.status);
  return result.output[0]!.id; // Allocated by the storage plugin.
}

// example: update
export async function updateMemory(resources: MemoryResources, id: string) {
  const memory = createMemory(resources);
  const result = await memory.update({ memories: [{ id, content: { language: "en" }, metadata: {} }] });
  if (result.status !== "success" || !result.output) throw new Error(result.error?.code ?? result.status);
  return result.output[0]!; // metadata is replaced, not merged; key is unchanged.
}

// example: delete
export async function deleteMemory(resources: MemoryResources, id: string) {
  const memory = createMemory(resources);
  const result = await memory.delete({ ids: [id, id] });
  if (result.status !== "success" || !result.output) throw new Error(result.error?.code ?? result.status);
  return result.output.deleted; // A missing id is not reported as deleted.
}

// example: execute
export async function executeMemory(resources: MemoryResources) {
  const memory = createMemory(resources);
  return memory.execute("MEMORY.QUERY", { limit: 10 });
}

// example: plugin
export function adaptDatabase(database: MemoryStore & Partial<MemorySearchProvider>): MemoryResources {
  // These methods are the application's SDK adapter, not raw SQL/Milvus SDK methods.
  // Explicit calls retain the SDK adapter's receiver and connection pool.
  const store: MemoryStore = {
    get: (input, options) => database.get(input, options),
    query: (input, options) => database.query(input, options),
    write: (input, options) => database.write(input, options),
    update: (input, options) => database.update(input, options),
    delete: (input, options) => database.delete(input, options),
  };
  const search = database.search ? { search: (input: Parameters<MemorySearchProvider["search"]>[0], options?: Parameters<MemorySearchProvider["search"]>[1]) => database.search!(input, options) } : undefined;
  return { store, ...(search ? { search } : {}) };
}

// example: errors
export async function memoryErrors(store: MemoryStore) {
  const search: MemorySearchProvider = {
    async search(input) {
      if (input.strategy !== "keyword") throw new MemoryError("UNSUPPORTED_STRATEGY", "Only keyword search is supported");
      throw new MemoryError("SEARCH_UNAVAILABLE", "Search is temporarily unavailable");
    },
  };
  const memory = createMemory({ store, search });
  const invalid = await memory.query({ limit: 0 }); // failed / INVALID_INPUT; store.query is not called.
  const unavailable = await memory.search({ query: "x", strategy: "keyword" });
  return { invalid, unavailable };
}

// example: graph
export async function memoryGraph(resources: MemoryResources) {
  const runtime = createDitto({ workers: [createMemoryWorker(resources)] });
  const plan = graph<string>("read-memory")
    .node("memory", "MEMORY.GET", [], id => ({ ids: [id] }));
  try { return await runtime.run(plan, "m1"); }
  finally { await runtime.close(); } // Close the application's database pool afterwards.
}

// example: scaffold
export const customGet = memoryGetNode.define("MEMORY", async () => ({
  executionId: "example-call", node: "MEMORY.GET", status: "failed",
  error: { code: "NOT_CONFIGURED", message: "Configure the application storage adapter" },
}));

Ditto · @codesoul-co/ditto · Node.js 24+