Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 30 additions & 0 deletions src/app.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import express from "express";
import {
ANALYTICS_DB_PATH,
FOXMEMORY_REGISTRY_DB_PATH,
effectiveLlmModel,
effectiveGraphLlmModel,
EMBED_MODEL,
Expand All @@ -26,18 +27,25 @@ import {
} from "./config/env.js";
import { MODEL_CATALOG_SEED } from "./config/defaults.js";
import { initAnalyticsDb, analyticsDb } from "./analytics/db.js";
import { initRegistry, registry } from "./registry/db.js";
import { migrateToRegistry } from "./registry/migrate.js";
import { recreateMemory } from "./memory/factory.js";
import { agentResolver } from "./middleware/agentResolver.js";
import { createHealthRouter } from "./routes/health.js";
import { createMemoriesRouter } from "./routes/memories.js";
import { createConfigRouter } from "./routes/config.js";
import { createGraphRouter } from "./routes/graph.js";
import { createStatsRouter } from "./routes/stats.js";
import { createJobsRouter } from "./routes/jobs.js";
import { createAdminRouter } from "./routes/admin.js";

const AGENT_PREFIX = "/v2/agents/:agentId";

export const createApp = () => {
const app = express();
app.use(express.json({ limit: "1mb" }));

/* ── Analytics DB ────────────────────────────────────── */
const db = initAnalyticsDb(ANALYTICS_DB_PATH);
if (db) {
const persisted = db.getConfig("custom_prompt");
Expand Down Expand Up @@ -90,12 +98,34 @@ export const createApp = () => {
console.log("[config] memory instance recreated with restored DB config");
}

/* ── Registry DB + migration ─────────────────────────── */
const reg = initRegistry(FOXMEMORY_REGISTRY_DB_PATH);
if (reg) {
console.log("[registry] initialized");
migrateToRegistry(reg, analyticsDb);
}

/* ── Legacy routes (no prefix change) ────────────────── */
app.use(createHealthRouter());
app.use(createMemoriesRouter());
app.use(createConfigRouter());
app.use(createGraphRouter());
app.use(createStatsRouter());
app.use(createJobsRouter());

/* ── Admin routes ────────────────────────────────────── */
app.use(createAdminRouter());

/* ── Agent-scoped routes ─────────────────────────────── */
// The agentResolver middleware resolves req.agent + req.agentMemory from :agentId.
// Each agent-scoped router uses mergeParams and the routes embed the full path
// including :agentId, so Express populates req.params.agentId automatically.
// We apply agentResolver as a param-aware middleware on the agent path prefix.
app.use("/v2/agents/:agentId", agentResolver);
app.use(createMemoriesRouter(AGENT_PREFIX));
app.use(createConfigRouter(AGENT_PREFIX));
app.use(createGraphRouter(AGENT_PREFIX));
app.use(createStatsRouter(AGENT_PREFIX));

return app;
};
3 changes: 3 additions & 0 deletions src/config/env.ts
Original file line number Diff line number Diff line change
Expand Up @@ -91,6 +91,9 @@ export const IDEM_TTL_MS = Math.max(60_000, Number(process.env.IDEMPOTENCY_TTL_M

export const ANALYTICS_DB_PATH = process.env.FOXMEMORY_ANALYTICS_DB_PATH || "/data/foxmemory-analytics.db";

export const FOXMEMORY_REGISTRY_DB_PATH = process.env.FOXMEMORY_REGISTRY_DB_PATH || "/data/foxmemory-registry.db";
export const DEFAULT_AGENT = process.env.DEFAULT_AGENT || "";

export type RuntimeStats = {
startedAt: string;
writesByMode: { infer: number; raw: number };
Expand Down
115 changes: 115 additions & 0 deletions src/memory/pool.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,115 @@
import { Memory } from "@foxlight-foundation/mem0ai/oss";
import type { AgentRecord } from "../registry/types.js";
import {
OPENAI_API_KEY,
OPENAI_BASE_URL,
effectiveLlmModel,
EMBED_MODEL,
GRAPH_ENABLED,
NEO4J_URL,
NEO4J_USERNAME,
NEO4J_PASSWORD,
effectiveGraphLlmModel,
GRAPH_SEARCH_THRESHOLD,
GRAPH_NODE_DEDUP_THRESHOLD,
GRAPH_BM25_TOPK,
roleUserName,
roleAssistantName,
} from "../config/env.js";
import { hardenGraphJsonContract } from "../utils/json.js";

const pool = new Map<string, { memory: Memory; lastUsed: number }>();
const MAX_POOL_SIZE = 50;

const createAgentMemory = (agent: AgentRecord): Memory => {
const mem = new Memory({
version: "v1.1",
historyDbPath: process.env.MEM0_HISTORY_DB_PATH || "/tmp/history.db",
roleNames: { user: roleUserName, assistant: roleAssistantName },
llm: {
provider: "openai",
config: {
apiKey: OPENAI_API_KEY,
model: effectiveLlmModel,
...(OPENAI_BASE_URL ? { baseURL: OPENAI_BASE_URL } : {}),
},
},
embedder: {
provider: "openai",
config: {
apiKey: OPENAI_API_KEY,
model: EMBED_MODEL,
...(OPENAI_BASE_URL ? { baseURL: OPENAI_BASE_URL } : {}),
},
},
...(process.env.QDRANT_HOST
? {
vectorStore: {
provider: "qdrant",
config: {
host: process.env.QDRANT_HOST,
port: Number(process.env.QDRANT_PORT || 6333),
apiKey: process.env.QDRANT_API_KEY,
collectionName: agent.qdrant_collection,
},
},
}
: {}),
...(GRAPH_ENABLED
? {
enableGraph: true,
graphStore: {
provider: "neo4j",
config: {
url: NEO4J_URL!,
username: NEO4J_USERNAME,
password: NEO4J_PASSWORD!,
...(agent.neo4j_database !== "neo4j" ? { database: agent.neo4j_database } : {}),
},
...(GRAPH_SEARCH_THRESHOLD !== undefined ? { searchThreshold: GRAPH_SEARCH_THRESHOLD } : {}),
...(GRAPH_NODE_DEDUP_THRESHOLD !== undefined ? { nodeDeduplicationThreshold: GRAPH_NODE_DEDUP_THRESHOLD } : {}),
...(GRAPH_BM25_TOPK !== undefined ? { bm25TopK: GRAPH_BM25_TOPK } : {}),
llm: {
provider: "openai",
config: {
apiKey: OPENAI_API_KEY,
model: effectiveGraphLlmModel,
...(OPENAI_BASE_URL ? { baseURL: OPENAI_BASE_URL } : {}),
},
},
},
}
: {}),
});

return hardenGraphJsonContract(mem);
};

export const getOrCreateMemory = (agentId: string, agent: AgentRecord): Memory => {
const entry = pool.get(agentId);
if (entry) {
entry.lastUsed = Date.now();
return entry.memory;
}

// LRU eviction if at capacity
if (pool.size >= MAX_POOL_SIZE) {
let oldestKey: string | null = null;
let oldestTime = Infinity;
for (const [key, val] of pool) {
if (val.lastUsed < oldestTime) {
oldestTime = val.lastUsed;
oldestKey = key;
}
}
if (oldestKey) pool.delete(oldestKey);
}

const memory = createAgentMemory(agent);
pool.set(agentId, { memory, lastUsed: Date.now() });
return memory;
};

export const evictMemory = (agentId: string): void => {
pool.delete(agentId);
};
58 changes: 58 additions & 0 deletions src/middleware/agentResolver.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,58 @@
import type { Request, Response, NextFunction } from "express";
import type { Memory } from "@foxlight-foundation/mem0ai/oss";
import type { AgentRecord } from "../registry/types.js";
import { registry } from "../registry/db.js";
import { getOrCreateMemory } from "../memory/pool.js";
import { getMemory } from "../memory/factory.js";
import { DEFAULT_AGENT } from "../config/env.js";
import { v2Err } from "../utils/response.js";

declare global {
namespace Express {
interface Request {
agent?: AgentRecord;
agentMemory?: Memory;
}
}
}

/**
* Middleware for agent-scoped routes (`:agentId` in params).
* Looks up the agent in the registry, attaches `req.agent` and `req.agentMemory`.
*/
export const agentResolver = (req: Request, res: Response, next: NextFunction): void => {
const agentId = req.params.agentId;

if (!agentId) {
// Legacy route — resolve via DEFAULT_AGENT or fall back to singleton
if (DEFAULT_AGENT && registry?.ready) {
const agent = registry.getAgent(DEFAULT_AGENT);
if (agent) {
req.agent = agent;
req.agentMemory = getOrCreateMemory(agent.id, agent);
}
}
next();
return;
}

if (!registry?.ready) {
v2Err(res, 503, "SERVICE_UNAVAILABLE", "Registry not available");
return;
}

const agent = registry.getAgent(agentId);
if (!agent) {
v2Err(res, 404, "NOT_FOUND", `Agent ${agentId} not found`);
return;
}

if (agent.status !== "active") {
v2Err(res, 503, "SERVICE_UNAVAILABLE", `Agent ${agentId} is ${agent.status}, not active`);
return;
}

req.agent = agent;
req.agentMemory = getOrCreateMemory(agent.id, agent);
next();
};
6 changes: 4 additions & 2 deletions src/pipeline/retry.ts
Original file line number Diff line number Diff line change
@@ -1,11 +1,13 @@
import type { Memory } from "@foxlight-foundation/mem0ai/oss";
import { ADD_RETRIES, ADD_RETRY_DELAY_MS } from "../config/env.js";
import { getMemory } from "../memory/factory.js";

export const addWithRetries = async (
messages: Array<{ role: string; content: string }>,
opts: { userId?: string; runId?: string; metadata?: Record<string, unknown> }
opts: { userId?: string; runId?: string; metadata?: Record<string, unknown> },
memoryOverride?: Memory,
) => {
const memory = getMemory();
const memory = memoryOverride ?? getMemory();
let last: any = { results: [] };
for (let attempt = 1; attempt <= Math.max(1, ADD_RETRIES); attempt++) {
try {
Expand Down
15 changes: 9 additions & 6 deletions src/pipeline/write.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import { z } from "zod";
import express from "express";
import { randomUUID } from "node:crypto";
import type { Memory } from "@foxlight-foundation/mem0ai/oss";
import {
runtimeStats,
ADD_RETRIES,
Expand Down Expand Up @@ -48,8 +49,8 @@ export const captureGraphLinks = (result: any, userId?: string) => {
}
};

export const v2Write = async (body: z.infer<typeof v2WriteSchema>) => {
const memory = getMemory();
export const v2Write = async (body: z.infer<typeof v2WriteSchema>, memoryOverride?: Memory) => {
const memory = memoryOverride ?? getMemory();
const userId = body.user_id;
const runId = body.run_id;
const metadata = body.metadata;
Expand Down Expand Up @@ -79,7 +80,7 @@ export const v2Write = async (body: z.infer<typeof v2WriteSchema>) => {
userId,
runId,
metadata
});
}, memoryOverride);
const hasResults = Array.isArray(inferResult?.results) && inferResult.results.length > 0;
if (hasResults) {
trackAddResult("infer", inferResult);
Expand Down Expand Up @@ -127,9 +128,10 @@ export const v2Write = async (body: z.infer<typeof v2WriteSchema>) => {
export const executeWriteAndRecord = async (
parsed: z.infer<typeof v2WriteSchema>,
idem: ReturnType<typeof idempotencyPrecheck>,
memoryOverride?: Memory,
): Promise<{ status: number; body: any }> => {
const t0 = Date.now();
const out = await v2Write(parsed);
const out = await v2Write(parsed, memoryOverride);
const latencyMs = Date.now() - t0;
analyticsDb?.recordWriteResults({
results: out.result?.results || [],
Expand Down Expand Up @@ -159,6 +161,7 @@ export const handleV2Write = async (
req: express.Request,
res: express.Response,
route: string,
memoryOverride?: Memory,
) => {
try {
const parsed = v2WriteSchema.safeParse(req.body);
Expand Down Expand Up @@ -198,7 +201,7 @@ export const handleV2Write = async (
job.status = "running";
try {
const noopIdem = { type: "none" as const };
const { body } = await executeWriteAndRecord(parsed.data, noopIdem);
const { body } = await executeWriteAndRecord(parsed.data, noopIdem, memoryOverride);
job.status = "completed";
job.completed_at = new Date().toISOString();
job.result = body.data;
Expand All @@ -213,7 +216,7 @@ export const handleV2Write = async (
return res.status(202).json(acceptedBody);
}

const { status, body } = await executeWriteAndRecord(parsed.data, idem);
const { status, body } = await executeWriteAndRecord(parsed.data, idem, memoryOverride);
return res.status(status).json(body);
} catch (err: any) {
return v2Err(res, 500, "INTERNAL_ERROR", String(err?.message || err));
Expand Down
Loading
Loading