fixes for summary prompt, chat route logger, QDrant orphan cleanup, llamacpp stream buffering
This commit is contained in:
@@ -66,18 +66,28 @@ async function* completeStream(prompt, options = {}) {
|
|||||||
if (!res.ok)
|
if (!res.ok)
|
||||||
throw new Error(`llama.cpp error: ${res.status} ${res.statusText}`);
|
throw new Error(`llama.cpp error: ${res.status} ${res.statusText}`);
|
||||||
|
|
||||||
for await (const chunk of res.body) {
|
//SSE lines can be split across network chunks, so we need to buffer them until we have a complete line
|
||||||
const lines = Buffer.from(chunk)
|
let buffer = '';
|
||||||
.toString("utf8")
|
|
||||||
.split("\n")
|
|
||||||
.filter((l) => l.startsWith("data: ") && l !== "data: [DONE]");
|
|
||||||
|
|
||||||
for (const line of lines) {
|
function* processLine(line){
|
||||||
const json = JSON.parse(line.slice(6));
|
if (!line.startsWith('data: ') || line === 'data: [DONE]') return;
|
||||||
const delta = json.choices?.[0]?.delta?.content;
|
|
||||||
|
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' ) {
|
if (json.choices?.[0]?.finish_reason === 'stop' ) {
|
||||||
finalModel = json.model ?? finalModel;
|
finalModel = json.model ?? finalModel;
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// usage arrives in a separate final chunk with empty choices array
|
// usage arrives in a separate final chunk with empty choices array
|
||||||
@@ -85,10 +95,23 @@ async function* completeStream(prompt, options = {}) {
|
|||||||
finalTokenCount = (json.usage.completion_tokens ?? 0) + (json.usage.prompt_tokens ?? 0);
|
finalTokenCount = (json.usage.completion_tokens ?? 0) + (json.usage.prompt_tokens ?? 0);
|
||||||
}
|
}
|
||||||
|
|
||||||
if (delta) yield { response: delta, done: false };
|
for (const line of lines) {
|
||||||
|
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);
|
logger.info('[llamacpp] finalTokenCount:', finalTokenCount);
|
||||||
|
|
||||||
yield { response: '', done: true, model: finalModel, tokenCount: finalTokenCount };
|
yield { response: '', done: true, model: finalModel, tokenCount: finalTokenCount };
|
||||||
|
|||||||
@@ -62,6 +62,11 @@ function touchSession(id) {
|
|||||||
function deleteSession(id) {
|
function deleteSession(id) {
|
||||||
const db = getDB();
|
const db = getDB();
|
||||||
db.prepare(`DELETE FROM sessions WHERE id = ?`).run(id);
|
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 } = {}) {
|
function updateSession(id, { name, projectId } = {}) {
|
||||||
|
|||||||
@@ -192,10 +192,14 @@ app.get('/entities/:id', (req, res) => {
|
|||||||
});
|
});
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
// Delete an entity by ID
|
// Delete an entity by ID
|
||||||
app.delete('/entities/:id', (req, res) => {
|
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();
|
res.status(204).send();
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|||||||
@@ -99,6 +99,23 @@ async function deleteEpisode(id) {
|
|||||||
return deleteVector(COLLECTIONS.EPISODES, 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 = {
|
module.exports = {
|
||||||
initCollections,
|
initCollections,
|
||||||
@@ -109,5 +126,7 @@ module.exports = {
|
|||||||
searchEntities,
|
searchEntities,
|
||||||
searchSummaries,
|
searchSummaries,
|
||||||
deleteVector,
|
deleteVector,
|
||||||
deleteEpisode
|
deleteEpisode,
|
||||||
|
deleteEntity,
|
||||||
|
deleteEpisodesBySession
|
||||||
};
|
};
|
||||||
@@ -32,6 +32,9 @@ function buildProjectSummaryPrompt(projectName, sessionSummaries) {
|
|||||||
'Be comprehensive but avoid padding. Do not repeat the same point twice.',
|
'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.',
|
'Write in third person. Output only the overview text, no headings or labels.',
|
||||||
'',
|
'',
|
||||||
|
summaryBlock,
|
||||||
|
'<|im_end|>',
|
||||||
|
'<|im_start|>assistant',
|
||||||
].join('\n');
|
].join('\n');
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -1,7 +1,7 @@
|
|||||||
const { Router } = require('express')
|
const { Router } = require('express')
|
||||||
const { chat, chatStream } = require('../chat/index');
|
const { chat, chatStream } = require('../chat/index');
|
||||||
const memory = require('../services/memory')
|
const memory = require('../services/memory')
|
||||||
const logger = require('@nexusai/shared');
|
const { logger } = require('@nexusai/shared');
|
||||||
|
|
||||||
|
|
||||||
const router = Router();
|
const router = Router();
|
||||||
|
|||||||
Reference in New Issue
Block a user