From 7757de76b7406366e8d41918a25e683ae2de2ddf Mon Sep 17 00:00:00 2001 From: Storme-bit Date: Mon, 17 Aug 2026 00:37:17 -0700 Subject: [PATCH] refinement fixes --- nexusai-refinements.patch | 245 ++++++++++++++++++ packages/memory-service/src/episodic/index.js | 26 ++ packages/memory-service/src/index.js | 12 + .../orchestration-service/src/chat/index.js | 33 ++- .../src/services/summarization.js | 42 +-- 5 files changed, 324 insertions(+), 34 deletions(-) create mode 100644 nexusai-refinements.patch diff --git a/nexusai-refinements.patch b/nexusai-refinements.patch new file mode 100644 index 0000000..6caebb3 --- /dev/null +++ b/nexusai-refinements.patch @@ -0,0 +1,245 @@ +diff -ruN baseline/packages/memory-service/src/episodic/index.js nexusai/packages/memory-service/src/episodic/index.js +--- baseline/packages/memory-service/src/episodic/index.js 2026-08-17 06:30:54.537312871 +0000 ++++ nexusai/packages/memory-service/src/episodic/index.js 2026-08-17 06:28:47.529976802 +0000 +@@ -172,6 +172,30 @@ + } + + ++// Aggregate stats for a session — count, total tokens, and highest episode id. ++// Lets summarization check the token threshold without pulling every episode row. ++function getSessionEpisodeStats(sessionId) { ++ const db = getDB(); ++ return db.prepare(` ++ SELECT COUNT(*) AS count, ++ COALESCE(SUM(token_count), 0) AS totalTokens, ++ COALESCE(MAX(id), 0) AS maxId ++ FROM episodes ++ WHERE session_id = ? ++ `).get(sessionId); ++} ++ ++// Episodes newer than a given id, in chronological order. Summarization uses ++// this to fetch only the un-summarized tail rather than the whole session. ++function getEpisodesSince(sessionId, afterId) { ++ const db = getDB(); ++ return db.prepare(` ++ SELECT * FROM episodes ++ WHERE session_id = ? AND id > ? ++ ORDER BY id ASC ++ `).all(sessionId, afterId).map(parseRow); ++} ++ + // Searches episodes using FTS5 full-text search, ordered by relevance, with a limit + function searchEpisodes(query, limit = EPISODIC.DEFAULT_SEARCH_LIMIT, sessionIds = null) { + const db = getDB(); +@@ -247,6 +271,8 @@ + getEpisode, + getEpisodesBySession, + getRecentEpisodes, ++ getSessionEpisodeStats, ++ getEpisodesSince, + searchEpisodes, + deleteEpisode, + getEpisodesByProject +diff -ruN baseline/packages/memory-service/src/index.js nexusai/packages/memory-service/src/index.js +--- baseline/packages/memory-service/src/index.js 2026-08-17 06:30:54.537745979 +0000 ++++ nexusai/packages/memory-service/src/index.js 2026-08-17 06:28:57.862794733 +0000 +@@ -156,6 +156,18 @@ + res.json(episodes); + }); + ++// Aggregate episode stats for a session (count, total tokens, max id). ++// Used by summarization to check the token threshold cheaply. ++app.get('/sessions/:id/episode-stats', (req, res) => { ++ res.json(episodic.getSessionEpisodeStats(Number(req.params.id))); ++}); ++ ++// Episodes newer than :afterId, chronological — the un-summarized tail. ++app.get('/sessions/:id/episodes/since/:afterId', (req, res) => { ++ const episodes = episodic.getEpisodesSince(Number(req.params.id), Number(req.params.afterId)); ++ res.json(episodes); ++}); ++ + app.delete('/episodes/:id', (req, res) => { + const id = Number(req.params.id); + episodic.deleteEpisode(id); +diff -ruN baseline/packages/orchestration-service/src/chat/index.js nexusai/packages/orchestration-service/src/chat/index.js +--- baseline/packages/orchestration-service/src/chat/index.js 2026-08-17 06:30:54.391455389 +0000 ++++ nexusai/packages/orchestration-service/src/chat/index.js 2026-08-17 06:28:36.122209942 +0000 +@@ -97,14 +97,14 @@ + } + + async function getSemanticEpisodes( +- userMessage, ++ vector, + sessionId, + recentIds, + projectSessionIds = null, + { semanticLimit, scoreThreshold } = {}, + ) { ++ if (!vector) return []; + try { +- const vector = await embedding.embed(userMessage); + const results = await qdrant.searchEpisodes(vector, { + limit: semanticLimit, + scoreThreshold: scoreThreshold, +@@ -127,9 +127,9 @@ + } + } + +-async function getRelevantEntities(userMessage, projectId = null) { ++async function getRelevantEntities(vector, projectId = null) { ++ if (!vector) return []; + try { +- const vector = await embedding.embed(userMessage); + const results = await qdrant.searchEntities(vector, { projectId }); + logger.info( + '[orchestration] Entity search results:', +@@ -233,7 +233,7 @@ + } + + +-async function getFusedEpisodes(userMessage, session, recentIds, projectSessionIds, settings) { ++async function getFusedEpisodes(userMessage, session, recentIds, projectSessionIds, vector, settings) { + const { semanticLimit, scoreThreshold, semanticWeight, keywordWeight } = settings; + const ftsSessionIds = projectSessionIds ?? [session.id]; + +@@ -243,7 +243,7 @@ + : Promise.resolve([]); + + const [semanticEps, rawKeywordEps] = await Promise.all([ +- getSemanticEpisodes(userMessage, session.id, recentIds, projectSessionIds, { semanticLimit, scoreThreshold }), ++ getSemanticEpisodes(vector, session.id, recentIds, projectSessionIds, { semanticLimit, scoreThreshold }), + ftsPromise, + ]); + +@@ -283,10 +283,19 @@ + const isFirstMessage = recentEpisodes.length === 0; + const recentIds = new Set(recentEpisodes.map(e => e.id)); + +- // 4. Fused retrieval + entity search in parallel (both are independent) ++ // 4. Embed the query once — the vector is shared by semantic episode search ++ // and entity search, so embedding it twice was a wasted round-trip + Ollama call. ++ let queryVector = null; ++ try { ++ queryVector = await embedding.embed(userMessage); ++ } catch (err) { ++ logger.warn('[orchestration] Query embedding failed; semantic + entity search disabled this turn:', err.message); ++ } ++ ++ // 4b. Fused retrieval + entity search in parallel (both are independent) + const [fusedWithScores, entityResults] = await Promise.all([ +- getFusedEpisodes(userMessage, session, recentIds, projectSessionIds, { semanticLimit, scoreThreshold, semanticWeight, keywordWeight }), +- getRelevantEntities(userMessage, session.project_id ?? null), ++ getFusedEpisodes(userMessage, session, recentIds, projectSessionIds, queryVector, { semanticLimit, scoreThreshold, semanticWeight, keywordWeight }), ++ getRelevantEntities(queryVector, session.project_id ?? null), + ]); + + // 5. Entity-linked episode IDs for scoring bonus +@@ -342,8 +351,7 @@ + logger.error('[orchestration] Failed to save episode:', err.message); + } + +- const allEpisodes = await memory.getRecentEpisodes(session.id, 9999); +- triggerSummary(session, allEpisodes); ++ triggerSummary(session); + + if (isFirstMessage && !session.name) { + autoNameSession(externalId, userMessage, result.text).catch(() => {}); +@@ -393,8 +401,7 @@ + + if (fullText.trim()) { + await memory.createEpisode(session.id, userMessage, fullText, tokenCount, session.project_id ?? null); +- const allEpisodes = await memory.getRecentEpisodes(session.id, 9999); +- triggerSummary(session, allEpisodes); ++ triggerSummary(session); + } else { + logger.warn('[orchestration] Stream finished with no assistant text; episode not saved'); + } +diff -ruN baseline/packages/orchestration-service/src/services/summarization.js nexusai/packages/orchestration-service/src/services/summarization.js +--- baseline/packages/orchestration-service/src/services/summarization.js 2026-08-17 06:30:54.397535212 +0000 ++++ nexusai/packages/orchestration-service/src/services/summarization.js 2026-08-17 06:29:22.666268446 +0000 +@@ -73,9 +73,12 @@ + return content; + } + +-async function maybeSummarize(session, allEpisodes) { +- // 1. Sum total tokens for this session +- const totalTokens = allEpisodes.reduce((sum, ep) => sum + (ep.token_count || 0), 0); ++async function maybeSummarize(session) { ++ // 1. Cheap aggregate — is this session even over the token threshold? ++ // Avoids pulling every episode row on messages that won't summarize. ++ const statsRes = await fetch(`${MEMORY_URL}/sessions/${session.id}/episode-stats`); ++ if (!statsRes.ok) return; ++ const { totalTokens } = await statsRes.json(); + if (totalTokens < THRESHOLD_TOKENS) return; // under threshold — nothing to do + + // 2. Fetch existing summaries for session +@@ -84,26 +87,23 @@ + const summaries = await summariesRes.json(); + + const latest = summaries.at(-1) ?? null; +- const lastCoveredId = latest +- ? parseInt(latest.episode_range?.split('-').at(-1)) || 0 ++ const lastCoveredId = latest ++ ? parseInt(latest.episode_range?.split('-').at(-1)) || 0 + : 0; +- // 3. Guard — don't re-summarize until MIN_EPISODES_SINCE new episodes have accumulated +- if (latest) { +- const newEpisodes = allEpisodes.filter(ep => ep.id > lastCoveredId); +- if (newEpisodes.length < MIN_EPISODES_SINCE) return; +- } + +- // 4. Determine episodes to summarize +- const episodesToSummarize = latest +- ? allEpisodes.filter(ep => ep.id > lastCoveredId) +- : allEpisodes; ++ // 3. Fetch only the un-summarized tail (full text), not the whole session. ++ // With no prior summary, lastCoveredId is 0 → this returns all episodes. ++ const episodesToSummarize = await fetch(`${MEMORY_URL}/sessions/${session.id}/episodes/since/${lastCoveredId}`) ++ .then(r => r.ok ? r.json() : []); ++ ++ // 4. Guard — don't re-summarize until MIN_EPISODES_SINCE new episodes have accumulated ++ if (latest && episodesToSummarize.length < MIN_EPISODES_SINCE) return; ++ if (episodesToSummarize.length === 0) return; + + // 5. Determine episode range from the episodes actually being summarized +- const summarizedIds = episodesToSummarize.map(ep => ep.id).sort((a,b) => a - b); ++ const summarizedIds = episodesToSummarize.map(ep => ep.id).sort((a, b) => a - b); + const episodeRange = `${summarizedIds.at(0)}-${summarizedIds.at(-1)}`; +- const totalEpisodeTokens = allEpisodes.reduce((sum, ep) => sum + (ep.token_count || 0), 0); + +- // add temporarily before the generateSummary call + logger.debug('[summarization] episodes to summarize:', episodesToSummarize.length); + + const content = await generateSummary( +@@ -122,7 +122,7 @@ + body: JSON.stringify({ + sessionId: session.id, + content, +- tokenCount: totalEpisodeTokens, ++ tokenCount: totalTokens, + episodeRange, + }), + }); +@@ -133,7 +133,7 @@ + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ + content, +- tokenCount: totalEpisodeTokens, ++ tokenCount: totalTokens, + episodeRange, + }), + }); +@@ -141,9 +141,9 @@ + } + } + +-async function triggerSummary(session, allEpisodes) { ++async function triggerSummary(session) { + // Intentionally fire-and-forget — caller doesn't await this +- maybeSummarize(session, allEpisodes).catch(err => ++ maybeSummarize(session).catch(err => + logger.warn('[summarization] Summary failed (non-critical):', err.message) + ); + } diff --git a/packages/memory-service/src/episodic/index.js b/packages/memory-service/src/episodic/index.js index 570bc49..a8c9caf 100644 --- a/packages/memory-service/src/episodic/index.js +++ b/packages/memory-service/src/episodic/index.js @@ -172,6 +172,30 @@ function getRecentEpisodes(sessionId, limit = EPISODIC.DEFAULT_RECENT_LIMIT) { } +// Aggregate stats for a session — count, total tokens, and highest episode id. +// Lets summarization check the token threshold without pulling every episode row. +function getSessionEpisodeStats(sessionId) { + const db = getDB(); + return db.prepare(` + SELECT COUNT(*) AS count, + COALESCE(SUM(token_count), 0) AS totalTokens, + COALESCE(MAX(id), 0) AS maxId + FROM episodes + WHERE session_id = ? + `).get(sessionId); +} + +// Episodes newer than a given id, in chronological order. Summarization uses +// this to fetch only the un-summarized tail rather than the whole session. +function getEpisodesSince(sessionId, afterId) { + const db = getDB(); + return db.prepare(` + SELECT * FROM episodes + WHERE session_id = ? AND id > ? + ORDER BY id ASC + `).all(sessionId, afterId).map(parseRow); +} + // Searches episodes using FTS5 full-text search, ordered by relevance, with a limit function searchEpisodes(query, limit = EPISODIC.DEFAULT_SEARCH_LIMIT, sessionIds = null) { const db = getDB(); @@ -247,6 +271,8 @@ module.exports = { getEpisode, getEpisodesBySession, getRecentEpisodes, + getSessionEpisodeStats, + getEpisodesSince, searchEpisodes, deleteEpisode, getEpisodesByProject diff --git a/packages/memory-service/src/index.js b/packages/memory-service/src/index.js index 1284c5d..2fa0f8f 100644 --- a/packages/memory-service/src/index.js +++ b/packages/memory-service/src/index.js @@ -156,6 +156,18 @@ app.get('/sessions/:id/episodes', (req, res) => { res.json(episodes); }); +// Aggregate episode stats for a session (count, total tokens, max id). +// Used by summarization to check the token threshold cheaply. +app.get('/sessions/:id/episode-stats', (req, res) => { + res.json(episodic.getSessionEpisodeStats(Number(req.params.id))); +}); + +// Episodes newer than :afterId, chronological — the un-summarized tail. +app.get('/sessions/:id/episodes/since/:afterId', (req, res) => { + const episodes = episodic.getEpisodesSince(Number(req.params.id), Number(req.params.afterId)); + res.json(episodes); +}); + app.delete('/episodes/:id', (req, res) => { const id = Number(req.params.id); episodic.deleteEpisode(id); diff --git a/packages/orchestration-service/src/chat/index.js b/packages/orchestration-service/src/chat/index.js index d3096de..905c85c 100644 --- a/packages/orchestration-service/src/chat/index.js +++ b/packages/orchestration-service/src/chat/index.js @@ -97,14 +97,14 @@ async function autoNameSession(externalId, userMessage, aiResponse) { } async function getSemanticEpisodes( - userMessage, + vector, sessionId, recentIds, projectSessionIds = null, { semanticLimit, scoreThreshold } = {}, ) { + if (!vector) return []; try { - const vector = await embedding.embed(userMessage); const results = await qdrant.searchEpisodes(vector, { limit: semanticLimit, scoreThreshold: scoreThreshold, @@ -127,9 +127,9 @@ async function getSemanticEpisodes( } } -async function getRelevantEntities(userMessage, projectId = null) { +async function getRelevantEntities(vector, projectId = null) { + if (!vector) return []; try { - const vector = await embedding.embed(userMessage); const results = await qdrant.searchEntities(vector, { projectId }); logger.info( '[orchestration] Entity search results:', @@ -233,7 +233,7 @@ function selectWithinBudget(scoredPool, contextBudget, minRecentEpisodes, recent } -async function getFusedEpisodes(userMessage, session, recentIds, projectSessionIds, settings) { +async function getFusedEpisodes(userMessage, session, recentIds, projectSessionIds, vector, settings) { const { semanticLimit, scoreThreshold, semanticWeight, keywordWeight } = settings; const ftsSessionIds = projectSessionIds ?? [session.id]; @@ -243,7 +243,7 @@ async function getFusedEpisodes(userMessage, session, recentIds, projectSessionI : Promise.resolve([]); const [semanticEps, rawKeywordEps] = await Promise.all([ - getSemanticEpisodes(userMessage, session.id, recentIds, projectSessionIds, { semanticLimit, scoreThreshold }), + getSemanticEpisodes(vector, session.id, recentIds, projectSessionIds, { semanticLimit, scoreThreshold }), ftsPromise, ]); @@ -283,10 +283,19 @@ async function assembleContext(externalId, userMessage) { const isFirstMessage = recentEpisodes.length === 0; const recentIds = new Set(recentEpisodes.map(e => e.id)); - // 4. Fused retrieval + entity search in parallel (both are independent) + // 4. Embed the query once — the vector is shared by semantic episode search + // and entity search, so embedding it twice was a wasted round-trip + Ollama call. + let queryVector = null; + try { + queryVector = await embedding.embed(userMessage); + } catch (err) { + logger.warn('[orchestration] Query embedding failed; semantic + entity search disabled this turn:', err.message); + } + + // 4b. Fused retrieval + entity search in parallel (both are independent) const [fusedWithScores, entityResults] = await Promise.all([ - getFusedEpisodes(userMessage, session, recentIds, projectSessionIds, { semanticLimit, scoreThreshold, semanticWeight, keywordWeight }), - getRelevantEntities(userMessage, session.project_id ?? null), + getFusedEpisodes(userMessage, session, recentIds, projectSessionIds, queryVector, { semanticLimit, scoreThreshold, semanticWeight, keywordWeight }), + getRelevantEntities(queryVector, session.project_id ?? null), ]); // 5. Entity-linked episode IDs for scoring bonus @@ -342,8 +351,7 @@ async function chat(externalId, userMessage, options = {}) { logger.error('[orchestration] Failed to save episode:', err.message); } - const allEpisodes = await memory.getRecentEpisodes(session.id, 9999); - triggerSummary(session, allEpisodes); + triggerSummary(session); if (isFirstMessage && !session.name) { autoNameSession(externalId, userMessage, result.text).catch(() => {}); @@ -393,8 +401,7 @@ async function chatStream(externalId, userMessage, onChunk, options = {}) { if (fullText.trim()) { await memory.createEpisode(session.id, userMessage, fullText, tokenCount, session.project_id ?? null); - const allEpisodes = await memory.getRecentEpisodes(session.id, 9999); - triggerSummary(session, allEpisodes); + triggerSummary(session); } else { logger.warn('[orchestration] Stream finished with no assistant text; episode not saved'); } diff --git a/packages/orchestration-service/src/services/summarization.js b/packages/orchestration-service/src/services/summarization.js index 13c8c08..0a2f52d 100644 --- a/packages/orchestration-service/src/services/summarization.js +++ b/packages/orchestration-service/src/services/summarization.js @@ -73,9 +73,12 @@ async function generateSummary(episodes, existingSummary = null) { return content; } -async function maybeSummarize(session, allEpisodes) { - // 1. Sum total tokens for this session - const totalTokens = allEpisodes.reduce((sum, ep) => sum + (ep.token_count || 0), 0); +async function maybeSummarize(session) { + // 1. Cheap aggregate — is this session even over the token threshold? + // Avoids pulling every episode row on messages that won't summarize. + const statsRes = await fetch(`${MEMORY_URL}/sessions/${session.id}/episode-stats`); + if (!statsRes.ok) return; + const { totalTokens } = await statsRes.json(); if (totalTokens < THRESHOLD_TOKENS) return; // under threshold — nothing to do // 2. Fetch existing summaries for session @@ -84,26 +87,23 @@ async function maybeSummarize(session, allEpisodes) { const summaries = await summariesRes.json(); const latest = summaries.at(-1) ?? null; - const lastCoveredId = latest - ? parseInt(latest.episode_range?.split('-').at(-1)) || 0 + const lastCoveredId = latest + ? parseInt(latest.episode_range?.split('-').at(-1)) || 0 : 0; - // 3. Guard — don't re-summarize until MIN_EPISODES_SINCE new episodes have accumulated - if (latest) { - const newEpisodes = allEpisodes.filter(ep => ep.id > lastCoveredId); - if (newEpisodes.length < MIN_EPISODES_SINCE) return; - } - // 4. Determine episodes to summarize - const episodesToSummarize = latest - ? allEpisodes.filter(ep => ep.id > lastCoveredId) - : allEpisodes; + // 3. Fetch only the un-summarized tail (full text), not the whole session. + // With no prior summary, lastCoveredId is 0 → this returns all episodes. + const episodesToSummarize = await fetch(`${MEMORY_URL}/sessions/${session.id}/episodes/since/${lastCoveredId}`) + .then(r => r.ok ? r.json() : []); + + // 4. Guard — don't re-summarize until MIN_EPISODES_SINCE new episodes have accumulated + if (latest && episodesToSummarize.length < MIN_EPISODES_SINCE) return; + if (episodesToSummarize.length === 0) return; // 5. Determine episode range from the episodes actually being summarized - const summarizedIds = episodesToSummarize.map(ep => ep.id).sort((a,b) => a - b); + const summarizedIds = episodesToSummarize.map(ep => ep.id).sort((a, b) => a - b); const episodeRange = `${summarizedIds.at(0)}-${summarizedIds.at(-1)}`; - const totalEpisodeTokens = allEpisodes.reduce((sum, ep) => sum + (ep.token_count || 0), 0); - // add temporarily before the generateSummary call logger.debug('[summarization] episodes to summarize:', episodesToSummarize.length); const content = await generateSummary( @@ -122,7 +122,7 @@ async function maybeSummarize(session, allEpisodes) { body: JSON.stringify({ sessionId: session.id, content, - tokenCount: totalEpisodeTokens, + tokenCount: totalTokens, episodeRange, }), }); @@ -133,7 +133,7 @@ async function maybeSummarize(session, allEpisodes) { headers: { 'Content-Type': 'application/json' }, body: JSON.stringify({ content, - tokenCount: totalEpisodeTokens, + tokenCount: totalTokens, episodeRange, }), }); @@ -141,9 +141,9 @@ async function maybeSummarize(session, allEpisodes) { } } -async function triggerSummary(session, allEpisodes) { +async function triggerSummary(session) { // Intentionally fire-and-forget — caller doesn't await this - maybeSummarize(session, allEpisodes).catch(err => + maybeSummarize(session).catch(err => logger.warn('[summarization] Summary failed (non-critical):', err.message) ); }