From 0248fcb58ca73b4e2b0ebf7f10bd4a019c869d6d Mon Sep 17 00:00:00 2001 From: Storme-bit Date: Sun, 16 Aug 2026 06:05:33 -0700 Subject: [PATCH] fixes for summary prompt, chat route logger, QDrant orphan cleanup, llamacpp stream buffering --- .../src/providers/llamacpp.js | 59 +++++++++++++------ packages/memory-service/src/episodic/index.js | 5 ++ packages/memory-service/src/index.js | 8 ++- packages/memory-service/src/semantic/index.js | 21 ++++++- .../src/summarization/project.js | 3 + .../orchestration-service/src/routes/chat.js | 2 +- 6 files changed, 76 insertions(+), 22 deletions(-) diff --git a/packages/inference-service/src/providers/llamacpp.js b/packages/inference-service/src/providers/llamacpp.js index d4dec70..6477343 100644 --- a/packages/inference-service/src/providers/llamacpp.js +++ b/packages/inference-service/src/providers/llamacpp.js @@ -66,29 +66,52 @@ async function* completeStream(prompt, options = {}) { if (!res.ok) throw new Error(`llama.cpp error: ${res.status} ${res.statusText}`); - for await (const chunk of res.body) { - const lines = Buffer.from(chunk) - .toString("utf8") - .split("\n") - .filter((l) => l.startsWith("data: ") && l !== "data: [DONE]"); + //SSE lines can be split across network chunks, so we need to buffer them until we have a complete line + let buffer = ''; + + function* processLine(line){ + if (!line.startsWith('data: ') || line === 'data: [DONE]') return; + + let json; + + try { + json = JSON.parse(line.slice(6)); + } catch (err) { + logger.error('[llamacpp] Skipping unparseable SSE line:', line.slice(0,120)); + return; + } + + + + const delta = json.choices?.[0]?.delta?.content ?? ''; + + if (json.choices?.[0]?.finish_reason === 'stop' ) { + finalModel = json.model ?? finalModel; + + } + + // usage arrives in a separate final chunk with empty choices array + if (json.usage){ + finalTokenCount = (json.usage.completion_tokens ?? 0) + (json.usage.prompt_tokens ?? 0); + } for (const line of lines) { - const json = JSON.parse(line.slice(6)); - const delta = json.choices?.[0]?.delta?.content; - - if (json.choices?.[0]?.finish_reason === 'stop') { - finalModel = json.model ?? finalModel; - } - - // usage arrives in a separate final chunk with empty choices array - if (json.usage) { - finalTokenCount = (json.usage.completion_tokens ?? 0) + (json.usage.prompt_tokens ?? 0); - } - - if (delta) yield { response: delta, done: false }; + yield* processLine(line.trim()); } } + for await (const chunk of res.body) { + buffer += Buffer.from(chunk).toString('utf-8'); + + const lines = buffer.split('\n'); + buffer = lines.pop() ?? ""; // keep the last line in the buffer, it may be incomplete + } + + //Flush anything left in buffer after stream closes + if (buffer.trim()){ + yield* processLine(buffer.trim()); + } + logger.info('[llamacpp] finalTokenCount:', finalTokenCount); yield { response: '', done: true, model: finalModel, tokenCount: finalTokenCount }; diff --git a/packages/memory-service/src/episodic/index.js b/packages/memory-service/src/episodic/index.js index 0d73f80..570bc49 100644 --- a/packages/memory-service/src/episodic/index.js +++ b/packages/memory-service/src/episodic/index.js @@ -62,6 +62,11 @@ function touchSession(id) { function deleteSession(id) { const db = getDB(); db.prepare(`DELETE FROM sessions WHERE id = ?`).run(id); + + //SQLite cascade removed episode rows, clean up QDrant vectors too + // Fire-and-forget - payload-filter delete, so it doesn't depend on the rows existing + semantic.deleteEpisodesBySession(id) + .catch(err => logger.error(`[Memory] QDrant cleanup failed for session ${id}:`, err.message)); } function updateSession(id, { name, projectId } = {}) { diff --git a/packages/memory-service/src/index.js b/packages/memory-service/src/index.js index f572900..1284c5d 100644 --- a/packages/memory-service/src/index.js +++ b/packages/memory-service/src/index.js @@ -192,10 +192,14 @@ app.get('/entities/:id', (req, res) => { }); - // Delete an entity by ID app.delete('/entities/:id', (req, res) => { - entities.deleteEntity(req.params.id); + const id = Number(req.params.id); + entities.deleteEntity(id); + + semantic.deleteEntity(id) //fire-and-forget, same pattern as episode delete + .catch(err => logger.error(`[Memory] Qdrant delete failed for entity ${id};`, err.message)); + res.status(204).send(); }); diff --git a/packages/memory-service/src/semantic/index.js b/packages/memory-service/src/semantic/index.js index fea41c8..89beedd 100644 --- a/packages/memory-service/src/semantic/index.js +++ b/packages/memory-service/src/semantic/index.js @@ -99,6 +99,23 @@ async function deleteEpisode(id) { return deleteVector(COLLECTIONS.EPISODES, id); } +async function deleteEntity(id) { + return deleteVector(COLLECTIONS.ENTITIES, id); +} + +//Delete ALL episode vectors associated with a session via payload filter +//Used on session delete: SQLite cascade removes the rows, +//Filter-based, +async function deleteEpisodesBySession(sessionId) { + const client = getClient(); + await client.delete(CoLLECTIONS.EPISODES, { + wait: true, + filter: { + must: [{ key: 'sessionId', match: { value: sessionId } }] + } + }); +} + module.exports = { initCollections, @@ -109,5 +126,7 @@ module.exports = { searchEntities, searchSummaries, deleteVector, - deleteEpisode + deleteEpisode, + deleteEntity, + deleteEpisodesBySession }; \ No newline at end of file diff --git a/packages/memory-service/src/summarization/project.js b/packages/memory-service/src/summarization/project.js index 76d8a48..1443c92 100644 --- a/packages/memory-service/src/summarization/project.js +++ b/packages/memory-service/src/summarization/project.js @@ -32,6 +32,9 @@ function buildProjectSummaryPrompt(projectName, sessionSummaries) { 'Be comprehensive but avoid padding. Do not repeat the same point twice.', 'Write in third person. Output only the overview text, no headings or labels.', '', + summaryBlock, + '<|im_end|>', + '<|im_start|>assistant', ].join('\n'); } diff --git a/packages/orchestration-service/src/routes/chat.js b/packages/orchestration-service/src/routes/chat.js index ebb6cf5..6f83264 100644 --- a/packages/orchestration-service/src/routes/chat.js +++ b/packages/orchestration-service/src/routes/chat.js @@ -1,7 +1,7 @@ const { Router } = require('express') const { chat, chatStream } = require('../chat/index'); const memory = require('../services/memory') -const logger = require('@nexusai/shared'); +const { logger } = require('@nexusai/shared'); const router = Router();