Compare commits

...
7 Commits
Author SHA1 Message Date
Storme-bit 3d7372d04f documentation update 2026-08-23 23:27:53 -07:00
Storme-bit c31fb786f4 documentation update 2026-08-23 23:26:51 -07:00
Storme-bit 66f6c4534f consolidation test 2026-08-23 22:55:39 -07:00
Storme-bit b3765f30fa consolidation test 2026-08-23 22:46:53 -07:00
Storme-bit 4b2099498c consolidation and aging 2026-08-23 22:41:21 -07:00
Storme-bit fd8867ee96 episodes touch 2026-08-23 19:24:10 -07:00
Storme-bit 00101ba14a test 2026-08-23 06:57:58 -07:00
15 changed files with 202 additions and 36 deletions
+21
View File
@@ -205,6 +205,9 @@ Returns `503` if llama-server is unreachable.
| `scoreThreshold` | float | 0–1 | Minimum similarity score for Qdrant results |
| `semanticWeight` | float | 0–5 | RRF weight for Qdrant semantic results |
| `keywordWeight` | float | 0–5 | RRF weight for FTS5 keyword results (`0` = disabled) |
| `contextBudget` | integer | — | Token budget for context assembly (char/4 estimation) |
| `entityWeight` | float | — | Scoring bonus for entity-linked episodes in the context pool |
| `minRecentEpisodes` | integer | — | Guaranteed floor of recent episodes always included in context |
| `modelsFolderPath` | string | — | Path to folder containing .gguf files |
| `temperature` | float | 0–2 | Inference randomness |
| `repeatPenalty` | float | 1–2 | Repeat token penalty |
@@ -253,6 +256,7 @@ orchestration.
| GET | /sessions/by-external/:externalId | Get session by external ID |
| PATCH | /sessions/by-external/:externalId | Update session fields |
| DELETE | /sessions/by-external/:externalId | Delete session (cascades to episodes) |
| GET | /sessions/:id/entity-ids | Entity IDs linked to this session's episodes (non-project memory isolation scoping) |
> Route ordering: `by-external/:externalId` must be defined before `/:id`
> to prevent `by-external` being captured as an ID param.
@@ -279,6 +283,8 @@ Both fields are optional. Only provided fields are updated.
| GET | /sessions/:id/episodes?limit=&offset= | Paginated episodes for a session |
| GET | /sessions/:id/episode-stats | Aggregate: count, total tokens, max id (summarization threshold check) |
| GET | /sessions/:id/episodes/since/:afterId | Episodes newer than :afterId, chronological (un-summarized tail) |
| POST | /episodes/touch | Batch access-tracking bump — increments `access_count`, sets `last_accessed_at` |
| GET | /sessions/:id/consolidation-candidates | Dry-run consolidation scoring (aging score, floors applied) |
| DELETE | /episodes/:id | Delete episode (SQLite + Qdrant cleanup) |
> Route ordering: `/episodes/search` must be defined before `/episodes/:id`.
@@ -293,6 +299,21 @@ Both fields are optional. Only provided fields are updated.
}
```
**POST /episodes/touch — body:**
```json
{ "ids": [54, 55, 56] }
```
Returns `{ "touched": 3 }`. Called fire-and-forget by orchestration after
budget selection. Nonexistent IDs are silently skipped; `touched` echoes the
request count, not rows matched.
**GET /sessions/:id/consolidation-candidates** — returns
`{ eligible, reason?, candidates: [...] }` where each candidate is
`{ id, access_count, created_at, last_accessed_at, preview, aging_score }`,
sorted by `aging_score` ascending (most eligible first). Sessions under
`CONSOLIDATION.MIN_SESSION_EPISODES` return `eligible: false` with a `reason`.
Observe-only — nothing is modified.
### Projects
| Method | Path | Description |
+8
View File
@@ -18,6 +18,7 @@ npm test # = node --test (discovers test/*.test.js at the repo root)
| `test/summarization.test.js` | Summarization decision logic | `maybeSummarize` (orchestration) |
| `test/entity-extraction.test.js` | Greeting + regurgitation guards | `mentionedIn`, `isIgnoredName` (memory-service) |
| `test/schema.test.js` | Fresh-DB schema completeness | `schema.js` string (memory-service) |
| `test/trivial-turn.test.js` | Greeting/trivial-turn detection | `isTrivialTurn` (shared) |
| `test/migrations.test.js` | Migration version-stepping | `migrate` (memory-service) |
Tests import the **real** functions rather than reimplementing logic — the
@@ -44,3 +45,10 @@ Keep the pattern: export the real function, import it, mock I/O at the boundary
logic (ranking, tokenizing, version-stepping, decision branches) over wiring.
Several of these tests were written *after* a bug slipped through — each new
class of mistake is worth a case so it can't recur silently.
When a mocked service call changes shape (e.g. the `utilityInference` refactor
moving summarization from Ollama's `/api/generate` to the inference service's
`/utility/complete`), the mock URL router must move with it — an "unexpected
fetch" throw in these tests usually means the code under test evolved, not
broke. Runner-contract tests (migrations) inject stub no-op migration arrays
rather than letting the real SQL-bearing migrations hit the minimal fake db.
+8 -7
View File
@@ -74,8 +74,9 @@ Multi-strategy retrieval merged into a single ranked result set.
### 3. Memory Consolidation Lifecycle
Prevents long-term memory degradation and enables compression.
- [ ] Episode aging — score/weight episodes by recency and access frequency
- [ ] Consolidation pass — merge related low-weight episodes into summary nodes
- [x] Episode aging — `access_count` + `last_accessed_at` columns (v2 migration), batch touch on retrieval selection (`POST /episodes/touch`), aging score `access_count / (1 + days since last access)`
- [x] Dry-run candidates endpoint — `GET /sessions/:id/consolidation-candidates` with age + session-size floors (observe-only phase before destructive pass)
- [ ] Consolidation pass — merge related low-weight episodes into summary nodes (incl. Qdrant vector + entity link cleanup)
- [ ] Orphan cleanup — remove entities no longer referenced by active episodes
### 4. User Preference Model
@@ -90,11 +91,11 @@ Short-circuit simple requests before they reach the LLM.
- [ ] Confidence bands — FAST PATH (memory lookup only) vs FULL (LLM + context)
- [ ] Fast-path handlers — direct memory queries, session lookups, factual recalls
### 6. Smarter Context Assembly *(inspired by acid2lake)*
### 6. Smarter Context Assembly *(inspired by acid2lake)* ✅
Budget-aware context selection instead of dumping all relevant memory into the prompt.
- [ ] Token budget manager in orchestration
- [ ] Priority scoring — recency × relevance × entity weight
- [ ] Configurable context budget via env var
- [x] Token budget manager in orchestration (`selectWithinBudget`, char/4 estimation on stored text)
- [x] Priority scoring — RRF fusion + recency + entity boost (`buildScoredPool`)
- [x] Configurable via settings (`contextBudget`, `entityWeight`, `minRecentEpisodes`) — live, no restart
### 7. Procedural Memory Store *(inspired by acid2lake)*
Learns "how NexusAI has successfully handled this type of request before."
@@ -227,4 +228,4 @@ The JARVIS moment — NexusAI reasons, plans, and acts across multiple steps.
---
*Last updated: April 2026*
*Last updated: August 2026*
+33
View File
@@ -78,6 +78,11 @@ does **not** reconcile columns on old tables — that's what migrations are for)
```js
const migrations = [
(_db) => {}, // v0 → v1: consolidated baseline (historical ALTERs folded into schema.js)
(db) => { // v1 → v2: access tracking for consolidation lifecycle
db.exec(`ALTER TABLE episodes ADD COLUMN last_accessed_at INTEGER`);
db.exec(`ALTER TABLE episodes ADD COLUMN access_count INTEGER NOT NULL DEFAULT 0`);
db.exec(`UPDATE episodes SET last_accessed_at = created_at`); // backfill
},
];
const LATEST_VERSION = migrations.length; // derived, never hand-maintained
```
@@ -121,6 +126,12 @@ previously ran on every startup.
- `foreign_keys = ON` — enforces referential integrity and cascade deletes
- PRAGMAs set via `db.pragma()`, not `db.exec()`
> **Copying a live WAL database:** `cp` on the `.db` file alone silently loses
> everything in the un-checkpointed `-wal` file (recent writes, even the
> migration version stamp). Always use
> `sqlite3 nexusai.db "VACUUM INTO './copy.db'"` (or `.backup`) — safe while
> the service is running, produces a complete single-file snapshot.
### Dynamic Updates
Both `updateSession` and `updateProject` build their `SET` clause dynamically
@@ -204,6 +215,28 @@ service is responsible only for CRUD — generation logic lives in orchestration
> For full details on trigger conditions, prompt format, cumulative updates,
> and ChatML token stripping, see `summarization.md`.
## Access Tracking & Consolidation (dry-run)
Every episode selected into a chat context window (budget-selected, not the
guaranteed-recency floor) gets an access bump via `POST /episodes/touch` —
`access_count` incremented, `last_accessed_at` set to `Date.now()` (ms).
Called fire-and-forget from orchestration; a failure loses one increment,
nothing more.
`GET /sessions/:id/consolidation-candidates` scores episodes by
`access_count / (1 + days since last access)` — never-accessed episodes fall
back to `created_at` for the recency term and score exactly 0 (most eligible).
Two floors apply: episodes younger than `CONSOLIDATION.MIN_AGE_DAYS` are
excluded in SQL; sessions under `CONSOLIDATION.MIN_SESSION_EPISODES` return
`eligible: false` before scoring runs. The endpoint is observe-only — the
destructive pass (merge → summarize → Qdrant cleanup → orphan sweep) is not
yet built.
> **Unit note:** `created_at` is unix **seconds** (`unixepoch()`);
> `last_accessed_at` is unix **milliseconds** (`Date.now()`). The scoring
> query normalizes with `created_at * 1000`. Keep this in mind for any new
> queries touching both columns.
## Delete Behaviour (SQLite + Qdrant consistency)
SQLite cascades handle relational cleanup, but Qdrant is a separate store and
+31 -20
View File
@@ -73,6 +73,9 @@ via `appSettings.load()` — changes apply immediately without a service restart
| `scoreThreshold` | 0.5 | Minimum similarity score for Qdrant semantic results |
| `semanticWeight` | 1.0 | RRF weight for Qdrant semantic results |
| `keywordWeight` | 0 | RRF weight for FTS5 keyword results (`0` = disabled) |
| `contextBudget` | — | Token budget for context assembly (char/4 estimation on stored text) |
| `entityWeight` | — | Scoring bonus for entity-linked episodes in the context pool |
| `minRecentEpisodes` | — | Guaranteed floor of recent episodes always included in context |
| `modelsFolderPath` | `/mnt/nexus-models` | Path to folder containing .gguf files |
| `temperature` | 0.7 | Inference temperature |
| `repeatPenalty` | 1.1 | Repeat token penalty |
@@ -101,35 +104,43 @@ difference is how the inference response is delivered to the client.
4. **Recent episode retrieval** — fetch most recent episodes (`recentEpisodeLimit`).
5. **Fused episode retrieval** — runs semantic (Qdrant) and keyword (FTS5)
5. **Trivial-turn gate** — greetings/pleasantries (`isTrivialTurn`) skip all
retrieval (semantic, keyword, entity); recent history alone is the context.
Breaks the greeting → marginal-retrieval → confabulation loop.
6. **Fused episode retrieval** — runs semantic (Qdrant) and keyword (FTS5)
search in parallel, then merges results via Reciprocal Rank Fusion (RRF).
Both paths are filtered against `recentIds` before fusion. FTS is scoped
to the current session or all project sessions. If `keywordWeight` is `0`,
the FTS call is skipped entirely. Non-critical — failures fall back to
whichever strategy succeeded.
The query is embedded once and shared with entity search. Both paths are
filtered against `recentIds` before fusion. FTS is scoped to the current
session or all project sessions. If `keywordWeight` is `0`, the FTS call
is skipped entirely. Non-critical — failures fall back to whichever
strategy succeeded.
6. **Entity search** — query `entities` Qdrant collection filtered by
`projectId`. Returns entity IDs alongside Qdrant payload data (the Qdrant
point ID equals the SQLite entity ID). Non-critical.
7. **Entity search + graph expansion** — query `entities` Qdrant collection
(project-scoped, or session-scoped via `/sessions/:id/entity-ids` for
non-project chats). Entity IDs are expanded into a 1-hop subgraph via
`POST /graph/neighbors`; on failure, falls back to flat entity list.
Non-critical.
7. **Graph neighborhood expansion** — call `POST /graph/neighbors` on
memory-service with the entity IDs from step 6. Returns a 1-hop subgraph
`{ nodes, edges }` — entity objects plus the relationships connecting them.
If no entities were found or the graph call fails, falls back to flat entity
list (no edges). Non-critical.
8. **Scored pool + budget selection** — `buildScoredPool` combines RRF scores,
recency, and entity-linkage bonus; `selectWithinBudget` fills `contextBudget`
(char/4 token estimation on stored text) above a guaranteed floor of
`minRecentEpisodes` recent episodes. Selected episode IDs are then reported
to `POST /episodes/touch` fire-and-forget (access tracking for the
consolidation lifecycle).
8. **Prompt assembly** — combine system prompt, graph context, fused episodes,
recent episodes, and user message.
9. **Prompt assembly** — combine system prompt, graph context, selected
episodes, guaranteed recent episodes, and user message.
9. **Inference** — send to inference service. `/chat` awaits full response;
`/chat/stream` pipes SSE chunks to the client.
10. **Inference** — send to inference service. `/chat` awaits full response;
`/chat/stream` pipes SSE chunks to the client.
10. **Episode write** — write exchange back to memory with `projectId`.
11. **Episode write** — write exchange back to memory with `projectId`.
11. **Summarisation trigger** — `triggerSummary(session, allEpisodes)` called
12. **Summarisation trigger** — `triggerSummary(session)` called
fire-and-forget. See `summarization.md` for full details.
12. **Auto-naming** — on first message with no session name, fires a secondary
13. **Auto-naming** — on first message with no session name, fires a secondary
inference call (max 20 tokens, temperature 0.3) to generate a session name.
### Prompt Structure
+10
View File
@@ -203,6 +203,16 @@ SUMMARY_MAX_TOKENS=800
SUMMARY_MIN_EPISODES=5
```
#### `CONSOLIDATION`
Controls the memory consolidation lifecycle (currently dry-run only).
| Key | Value | Description |
|---|---|---|
| `MIN_AGE_DAYS` | `7` | Episodes younger than this are never consolidation candidates |
| `MIN_SESSION_EPISODES` | `20` | Sessions with fewer episodes are skipped entirely |
| `CANDIDATE_LIMIT` | `50` | Max candidates returned per scoring query |
#### `SQLITE`
| Key | Value | Description |
+32 -1
View File
@@ -1,5 +1,5 @@
const {getDB} = require('../db');
const { EPISODIC, getEnv, SERVICES, parseRow, formatEpisodeText, SUMMARIES, logger, isTrivialTurn } = require('@nexusai/shared');
const { EPISODIC, getEnv, SERVICES, parseRow, formatEpisodeText, SUMMARIES, logger, isTrivialTurn, CONSOLIDATION } = require('@nexusai/shared');
const semantic = require('../semantic');
const { extractAndStoreEntities } = require('../entities/extraction')
@@ -266,6 +266,19 @@ function deleteEpisode(id) {
db.prepare(`DELETE FROM episodes WHERE id = ?`).run(id);
}
function touchEpisodes(ids) {
if (!ids.length) return;
const db = getDB(); // <-- missing in your version
const placeholders = ids.map(() => '?').join(',');
db.prepare(`
UPDATE episodes
SET access_count = access_count + 1,
last_accessed_at = ?
WHERE id IN (${placeholders})
`).run(Date.now(), ...ids);
}
/******** Embedding Helper ********/
async function getEpisodeEmbedding(userMessage, aiResponse){
const url = getEnv('EMBEDDING_SERVICE_URL', SERVICES.EMBEDDING_URL);
@@ -298,6 +311,22 @@ function getEpisodesByProject(projectId, limit = SUMMARIES.MAX_PROJECT_EPISODE_L
`).all(projectId, limit).map(parseRow);
}
function getConsolidationCandidates(sessionId, limit = CONSOLIDATION.CANDIDATE_LIMIT) {
const db = getDB();
const now = Date.now();
const cutoff = Math.floor(now / 1000) - CONSOLIDATION.MIN_AGE_DAYS *86400 // seconds, so it matches created_at
return db.prepare(`
SELECT id, access_count, created_at, last_accessed_at,
SUBSTR(user_message, 1, 80) AS preview,
CAST( access_count AS REAL) / (1+(?-COALESCE(last_accessed_at, created_at * 1000)) / 86400000.0) AS aging_score
FROM episodes
WHERE session_id = ? AND created_at < ?
ORDER BY aging_score ASC
LIMIT ?
`).all(now, sessionId, cutoff, limit);
}
module.exports = {
createSession,
getSession,
@@ -315,6 +344,8 @@ module.exports = {
getEpisodesSince,
searchEpisodes,
deleteEpisode,
touchEpisodes,
getEpisodesByProject,
buildFtsQuery,
getConsolidationCandidates,
};
+28 -1
View File
@@ -1,6 +1,6 @@
require ('dotenv').config();
const express = require('express');
const {getEnv, PORTS, EPISODIC, logger} = require('@nexusai/shared');
const {getEnv, PORTS, EPISODIC, logger, CONSOLIDATION} = require('@nexusai/shared');
const { getDB } = require('./db');
const { createProject, getProjects, getProject, updateProject, deleteProject } = require('./db/projects');
const { createSummary, getSummary, getSummariesBySession, getSummariesByProject, updateSummary, deleteSummary } = require('./db/summaries');
@@ -139,6 +139,13 @@ app.get('/episodes/search', (req, res) => {
res.json(episodic.searchEpisodes(q, Number(limit), parsedSessionIds));
});
app.post('/episodes/touch', (req, res) => {
const { ids } = req.body;
if(!Array.isArray(ids)) return res.status(400).json({error: 'ids must be an array'});
episodic.touchEpisodes(ids);
res.json({touched: ids.length});
})
app.get('/episodes/:id', (req, res) => {
const episode = episodic.getEpisode(req.params.id);
if (!episode) return res.status(404).json({ error: 'Episode not found' });
@@ -175,6 +182,26 @@ app.get('/sessions/:id/episodes/since/:afterId', (req, res) => {
res.json(episodes);
});
app.get('/sessions/:id/consolidation-candidates', (req, res) => {
const sessionId = Number(req.params.id);
const stats = episodic.getSessionEpisodeStats(sessionId);
if(stats.count < CONSOLIDATION.MIN_SESSION_EPISODES) {
return res.json({
eligible: false,
reason: `session has ${stats.count} episodes, floor is ${CONSOLIDATION.MIN_SESSION_EPISODES}`,
candidates: [],
});
}
const candidates = episodic.getConsolidationCandidates(sessionId);
res.json({
eligible:true,
sessionEpisodeCount: stats.count,
candidates
})
})
app.delete('/episodes/:id', (req, res) => {
const id = Number(req.params.id);
episodic.deleteEpisode(id);
@@ -331,6 +331,9 @@ async function assembleContext(externalId, userMessage) {
const scoredPool = buildScoredPool(fusedWithScores, recentEpisodes, entityBoostedIds, { entityWeight });
const { guaranteed, selected } = selectWithinBudget(scoredPool, contextBudget, minRecentEpisodes, recentEpisodes);
const selectedIds = selected.map(ep => ep.id);
memory.touchEpisodes(selectedIds);
// 7. Graph neighborhood expansion
let neighborhood = { nodes: [], edges: [] };
if (entityIds.length > 0) {
@@ -1,4 +1,4 @@
const { getEnv, SERVICES, EPISODIC } = require('@nexusai/shared');
const { getEnv, SERVICES, EPISODIC, logger } = require('@nexusai/shared');
const BASE_URL = getEnv('MEMORY_SERVICE_URL', SERVICES.MEMORY_URL);
@@ -223,6 +223,16 @@ async function getEntityIdsBySession(sessionId){
return entityIds;
}
// orchestration-service/src/services/memory.js
async function touchEpisodes(ids) {
if (!ids.length) return;
fetch(`${BASE_URL}/episodes/touch`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ ids }),
}).catch(err => logger.warn(`[memory] touch failed: ${err.message}`));
}
module.exports = {
getSessionByExternalId,
createSession,
@@ -250,4 +260,5 @@ module.exports = {
searchEpisodes,
getEpisodesByEntities,
getEntityIdsBySession,
touchEpisodes,
}
+7
View File
@@ -120,6 +120,12 @@ const UTILITY = {
TIMEOUT_MS: 120_000, // extraction previously used 60s; summaries had none — standardized
};
const CONSOLIDATION = {
MIN_AGE_DAYS: 7, // episodes younger than this are never consolidation candidates
MIN_SESSION_EPISODES: 20, // sessions smaller than this are left entirely alone
CANDIDATE_LIMIT: 50,
}
module.exports = {
QDRANT,
COLLECTIONS,
@@ -135,4 +141,5 @@ module.exports = {
ENTITIES,
RETRIEVAL,
UTILITY,
CONSOLIDATION,
};
+3 -1
View File
@@ -13,7 +13,8 @@ const {
SUMMARIES,
ENTITIES,
RETRIEVAL,
UTILITY
UTILITY,
CONSOLIDATION
} = require('./config/constants');
const {parseRow, formatEpisodeText, isTrivialTurn} = require('./utils')
@@ -40,5 +41,6 @@ module.exports = {
logger,
RETRIEVAL,
UTILITY,
CONSOLIDATION,
utilityInference,
};
+4 -3
View File
@@ -4,7 +4,6 @@
// advance user_version correctly, and no-op when the DB is already current.
const { test } = require('node:test');
const assert = require('node:assert');
const Database = require('better-sqlite3');
const { migrate, LATEST_VERSION } = require('../packages/memory-service/src/db/migrations');
// Emulates just the slice of the better-sqlite3 API that migrate() uses.
@@ -26,14 +25,16 @@ function fakeDb(startVersion = 0) {
test('a fresh DB (user_version 0) is stamped to LATEST_VERSION', () => {
const db = fakeDb(0);
const result = migrate(db);
const stubs = Array.from({ length: LATEST_VERSION }, () => () => {});
const result = migrate(db, stubs); // inject no-op migrations
assert.strictEqual(result, LATEST_VERSION);
assert.strictEqual(db.version, LATEST_VERSION);
});
test('an already-current DB is a no-op (no version writes)', () => {
const db = fakeDb(LATEST_VERSION);
const result = migrate(db);
const stubs = Array.from({ length: LATEST_VERSION }, () => () => {});
const result = migrate(db, stubs); // fixed: db is now the first arg
assert.strictEqual(result, LATEST_VERSION);
assert.deepStrictEqual(db.setCalls, [], 'should not touch user_version when current');
});
+1 -1
View File
@@ -13,7 +13,7 @@ try { ({ DatabaseSync } = require('node:sqlite')); } catch { DatabaseSync = null
// The complete intended shape = original base tables + every historical ALTER.
const EXPECTED = {
sessions: ['id','external_id','created_at','updated_at','metadata','name','project_id'],
episodes: ['id','session_id','user_message','ai_response','created_at','token_count','metadata'],
episodes: ['access_count', 'ai_response', 'created_at', 'id', 'last_accessed_at', 'metadata', 'session_id', 'token_count', 'user_message'],
entities: ['id','name','type','notes','created_at','updated_at','metadata','mention_count','confidence','source','last_seen_at'],
relationships: ['id','from_id','to_id','label','created_at','metadata','mention_count','notes'],
entity_episodes: ['entity_id','episode_id'],
+1 -1
View File
@@ -19,7 +19,7 @@ function mockMemory({ stats, summaries, since }) {
calls.sinceAfterId = Number(u.split('/since/').at(-1));
return json(since.filter(ep => ep.id > calls.sinceAfterId));
}
if (u.includes('/api/generate')) return json({ response: 'A concise third-person summary.' });
if (u.endsWith('/utility/complete')) return json({ text: 'A concise third-person summary.' });
if (u.endsWith('/summaries') && opts.method === 'POST') { calls.posted = JSON.parse(opts.body); return json({ id: 99 }); }
if (/\/summaries\/\d+$/.test(u) && opts.method === 'PATCH') { calls.patched = JSON.parse(opts.body); return json({ ok: true }); }
throw new Error('unexpected fetch: ' + u);