129 lines
5.0 KiB
JavaScript
129 lines
5.0 KiB
JavaScript
const { getEnv, SERVICES, SUMMARIES, logger, utilityInference } = require('@nexusai/shared');
|
|
|
|
const MEMORY_URL = getEnv('MEMORY_SERVICE_URL', SERVICES.MEMORY_URL);
|
|
|
|
const THRESHOLD_TOKENS = parseInt(getEnv('SUMMARY_THRESHOLD_TOKENS', SUMMARIES.THRESHOLD_TOKENS));
|
|
const MAX_SUMMARY_TOKENS = parseInt(getEnv('SUMMARY_MAX_TOKENS', SUMMARIES.MAX_SUMMARY_TOKENS));
|
|
const MIN_EPISODES_SINCE = parseInt(getEnv('SUMMARY_MIN_EPISODES', SUMMARIES.MIN_EPISODES_SINCE));
|
|
|
|
function buildSummaryPrompt(episodes, existingSummary = null) {
|
|
const MAX_CHARS = 3000;
|
|
let context = episodes
|
|
.map(ep => `User: ${ep.user_message}\nAssistant: ${ep.ai_response}`)
|
|
.join('\n\n');
|
|
|
|
if (context.length > MAX_CHARS) {
|
|
context = context.slice(-MAX_CHARS);
|
|
}
|
|
|
|
const instruction = existingSummary
|
|
? `Update the summary below to incorporate the new exchanges.
|
|
Write 3-5 sentences in third person. Do not quote directly — paraphrase only.
|
|
Do not include greetings, sign-offs, or filler. Output only the updated summary text.
|
|
|
|
Previous summary:
|
|
${existingSummary}
|
|
|
|
New exchanges:
|
|
${context}`
|
|
: `Summarize the conversation below in 3-5 sentences.
|
|
Write in third person. Do not quote directly — paraphrase only.
|
|
Do not include greetings, sign-offs, or filler. Output only the summary text.
|
|
|
|
Conversation:
|
|
${context}`;
|
|
|
|
// No ChatML wrapper — the model's own prompt template is applied server-side
|
|
// by the inference service's /utility/complete route (Ollama /api/chat).
|
|
return instruction;
|
|
}
|
|
|
|
async function generateSummary(episodes, existingSummary = null) {
|
|
const user = buildSummaryPrompt(episodes, existingSummary);
|
|
|
|
const content = await utilityInference({
|
|
user,
|
|
temperature: SUMMARIES.TEMPERATURE,
|
|
maxTokens: SUMMARIES.SESSION_GEN_MAX_TOKENS,
|
|
});
|
|
|
|
return content;
|
|
}
|
|
|
|
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
|
|
const summariesRes = await fetch(`${MEMORY_URL}/sessions/${session.id}/summaries`);
|
|
if (!summariesRes.ok) return;
|
|
const summaries = await summariesRes.json();
|
|
|
|
const latest = summaries.at(-1) ?? null;
|
|
const lastCoveredId = latest
|
|
? parseInt(latest.episode_range?.split('-').at(-1)) || 0
|
|
: 0;
|
|
|
|
// 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 episodeRange = `${summarizedIds.at(0)}-${summarizedIds.at(-1)}`;
|
|
|
|
logger.debug('[summarization] episodes to summarize:', episodesToSummarize.length);
|
|
|
|
const content = await generateSummary(
|
|
episodesToSummarize,
|
|
latest && latest.content.length < MAX_SUMMARY_TOKENS ? latest.content : null
|
|
// if existing summary is already large, treat as fresh rather than appending to a huge blob
|
|
);
|
|
|
|
if (!content) return;
|
|
|
|
// 6. Create new row or update existing
|
|
if (!latest || latest.content.length >= MAX_SUMMARY_TOKENS) {
|
|
await fetch(`${MEMORY_URL}/summaries`, {
|
|
method: 'POST',
|
|
headers: { 'Content-Type': 'application/json' },
|
|
body: JSON.stringify({
|
|
sessionId: session.id,
|
|
content,
|
|
tokenCount: totalTokens,
|
|
episodeRange,
|
|
}),
|
|
});
|
|
logger.debug(`[summarization] Created new summary for session ${session.id}`);
|
|
} else {
|
|
await fetch(`${MEMORY_URL}/summaries/${latest.id}`, {
|
|
method: 'PATCH',
|
|
headers: { 'Content-Type': 'application/json' },
|
|
body: JSON.stringify({
|
|
content,
|
|
tokenCount: totalTokens,
|
|
episodeRange,
|
|
}),
|
|
});
|
|
logger.debug(`[summarization] Updated summary ${latest.id} for session ${session.id}`);
|
|
}
|
|
}
|
|
|
|
async function triggerSummary(session) {
|
|
// Intentionally fire-and-forget — caller doesn't await this
|
|
maybeSummarize(session).catch(err =>
|
|
logger.warn('[summarization] Summary failed (non-critical):', err.message)
|
|
);
|
|
}
|
|
|
|
module.exports = { triggerSummary, maybeSummarize };
|