Compare commits

..
2 Commits
Author SHA1 Message Date
Storme-bit 33536982bc refinement fixes 2026-08-17 00:37:52 -07:00
Storme-bit 7757de76b7 refinement fixes 2026-08-17 00:37:17 -07:00
4 changed files with 79 additions and 34 deletions
@@ -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 // Searches episodes using FTS5 full-text search, ordered by relevance, with a limit
function searchEpisodes(query, limit = EPISODIC.DEFAULT_SEARCH_LIMIT, sessionIds = null) { function searchEpisodes(query, limit = EPISODIC.DEFAULT_SEARCH_LIMIT, sessionIds = null) {
const db = getDB(); const db = getDB();
@@ -247,6 +271,8 @@ module.exports = {
getEpisode, getEpisode,
getEpisodesBySession, getEpisodesBySession,
getRecentEpisodes, getRecentEpisodes,
getSessionEpisodeStats,
getEpisodesSince,
searchEpisodes, searchEpisodes,
deleteEpisode, deleteEpisode,
getEpisodesByProject getEpisodesByProject
+12
View File
@@ -156,6 +156,18 @@ app.get('/sessions/:id/episodes', (req, res) => {
res.json(episodes); 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) => { app.delete('/episodes/:id', (req, res) => {
const id = Number(req.params.id); const id = Number(req.params.id);
episodic.deleteEpisode(id); episodic.deleteEpisode(id);
@@ -97,14 +97,14 @@ async function autoNameSession(externalId, userMessage, aiResponse) {
} }
async function getSemanticEpisodes( async function getSemanticEpisodes(
userMessage, vector,
sessionId, sessionId,
recentIds, recentIds,
projectSessionIds = null, projectSessionIds = null,
{ semanticLimit, scoreThreshold } = {}, { semanticLimit, scoreThreshold } = {},
) { ) {
if (!vector) return [];
try { try {
const vector = await embedding.embed(userMessage);
const results = await qdrant.searchEpisodes(vector, { const results = await qdrant.searchEpisodes(vector, {
limit: semanticLimit, limit: semanticLimit,
scoreThreshold: scoreThreshold, 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 { try {
const vector = await embedding.embed(userMessage);
const results = await qdrant.searchEntities(vector, { projectId }); const results = await qdrant.searchEntities(vector, { projectId });
logger.info( logger.info(
'[orchestration] Entity search results:', '[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 { semanticLimit, scoreThreshold, semanticWeight, keywordWeight } = settings;
const ftsSessionIds = projectSessionIds ?? [session.id]; const ftsSessionIds = projectSessionIds ?? [session.id];
@@ -243,7 +243,7 @@ async function getFusedEpisodes(userMessage, session, recentIds, projectSessionI
: Promise.resolve([]); : Promise.resolve([]);
const [semanticEps, rawKeywordEps] = await Promise.all([ const [semanticEps, rawKeywordEps] = await Promise.all([
getSemanticEpisodes(userMessage, session.id, recentIds, projectSessionIds, { semanticLimit, scoreThreshold }), getSemanticEpisodes(vector, session.id, recentIds, projectSessionIds, { semanticLimit, scoreThreshold }),
ftsPromise, ftsPromise,
]); ]);
@@ -283,10 +283,19 @@ async function assembleContext(externalId, userMessage) {
const isFirstMessage = recentEpisodes.length === 0; const isFirstMessage = recentEpisodes.length === 0;
const recentIds = new Set(recentEpisodes.map(e => e.id)); 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([ const [fusedWithScores, entityResults] = await Promise.all([
getFusedEpisodes(userMessage, session, recentIds, projectSessionIds, { semanticLimit, scoreThreshold, semanticWeight, keywordWeight }), getFusedEpisodes(userMessage, session, recentIds, projectSessionIds, queryVector, { semanticLimit, scoreThreshold, semanticWeight, keywordWeight }),
getRelevantEntities(userMessage, session.project_id ?? null), getRelevantEntities(queryVector, session.project_id ?? null),
]); ]);
// 5. Entity-linked episode IDs for scoring bonus // 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); logger.error('[orchestration] Failed to save episode:', err.message);
} }
const allEpisodes = await memory.getRecentEpisodes(session.id, 9999); triggerSummary(session);
triggerSummary(session, allEpisodes);
if (isFirstMessage && !session.name) { if (isFirstMessage && !session.name) {
autoNameSession(externalId, userMessage, result.text).catch(() => {}); autoNameSession(externalId, userMessage, result.text).catch(() => {});
@@ -393,8 +401,7 @@ async function chatStream(externalId, userMessage, onChunk, options = {}) {
if (fullText.trim()) { if (fullText.trim()) {
await memory.createEpisode(session.id, userMessage, fullText, tokenCount, session.project_id ?? null); await memory.createEpisode(session.id, userMessage, fullText, tokenCount, session.project_id ?? null);
const allEpisodes = await memory.getRecentEpisodes(session.id, 9999); triggerSummary(session);
triggerSummary(session, allEpisodes);
} else { } else {
logger.warn('[orchestration] Stream finished with no assistant text; episode not saved'); logger.warn('[orchestration] Stream finished with no assistant text; episode not saved');
} }
@@ -73,9 +73,12 @@ async function generateSummary(episodes, existingSummary = null) {
return content; return content;
} }
async function maybeSummarize(session, allEpisodes) { async function maybeSummarize(session) {
// 1. Sum total tokens for this session // 1. Cheap aggregate — is this session even over the token threshold?
const totalTokens = allEpisodes.reduce((sum, ep) => sum + (ep.token_count || 0), 0); // 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 if (totalTokens < THRESHOLD_TOKENS) return; // under threshold — nothing to do
// 2. Fetch existing summaries for session // 2. Fetch existing summaries for session
@@ -87,23 +90,20 @@ async function maybeSummarize(session, allEpisodes) {
const lastCoveredId = latest const lastCoveredId = latest
? parseInt(latest.episode_range?.split('-').at(-1)) || 0 ? parseInt(latest.episode_range?.split('-').at(-1)) || 0
: 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 // 3. Fetch only the un-summarized tail (full text), not the whole session.
const episodesToSummarize = latest // With no prior summary, lastCoveredId is 0 → this returns all episodes.
? allEpisodes.filter(ep => ep.id > lastCoveredId) const episodesToSummarize = await fetch(`${MEMORY_URL}/sessions/${session.id}/episodes/since/${lastCoveredId}`)
: allEpisodes; .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 // 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 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); logger.debug('[summarization] episodes to summarize:', episodesToSummarize.length);
const content = await generateSummary( const content = await generateSummary(
@@ -122,7 +122,7 @@ async function maybeSummarize(session, allEpisodes) {
body: JSON.stringify({ body: JSON.stringify({
sessionId: session.id, sessionId: session.id,
content, content,
tokenCount: totalEpisodeTokens, tokenCount: totalTokens,
episodeRange, episodeRange,
}), }),
}); });
@@ -133,7 +133,7 @@ async function maybeSummarize(session, allEpisodes) {
headers: { 'Content-Type': 'application/json' }, headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ body: JSON.stringify({
content, content,
tokenCount: totalEpisodeTokens, tokenCount: totalTokens,
episodeRange, 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 // 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) logger.warn('[summarization] Summary failed (non-critical):', err.message)
); );
} }