Compare commits

..
23 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
Storme-bit 6c953a346a test 2026-08-20 05:07:53 -07:00
Storme-bit feadf0a1f7 test 2026-08-20 04:36:29 -07:00
Storme-bit ef3475db28 test 2026-08-20 04:35:04 -07:00
Storme-bit d2a1fb4ca9 test 2026-08-20 04:28:50 -07:00
Storme-bit f6b2ab82fc test 2026-08-20 04:28:21 -07:00
Storme-bit 4c00a8d84a test 2026-08-20 04:26:58 -07:00
Storme-bit 565286a4d9 isolation logging 2026-08-18 04:16:13 -07:00
Storme-bit 6aefc391d9 isolation logging 2026-08-18 04:08:57 -07:00
Storme-bit 4f6ae8492b fixing typoes 2026-08-18 03:38:43 -07:00
Storme-bit 16b1fa2d07 entity isolation and estimateTokenFix 2026-08-18 03:21:41 -07:00
Storme-bit eb9fda62f3 entity isolation and estimateTokenFix 2026-08-18 03:18:57 -07:00
Storme-bit dce63b6f7b entity isolation and estimateTokenFix 2026-08-18 02:57:05 -07:00
Storme-bit e9ceee15fa entity isolation and estimateTokenFix 2026-08-18 01:31:11 -07:00
Storme-bit 551c7f03ec utility inference cleanup 2026-08-17 23:48:15 -07:00
Storme-bit a223b753e5 utility inference layer 3 2026-08-17 10:11:21 -07:00
Storme-bit 67e2b446e4 utility inference layer 3 2026-08-17 10:11:08 -07:00
23 changed files with 356 additions and 202 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*
+2 -2
View File
@@ -2,7 +2,7 @@
**Location:** `packages/memory-service/src/entities/extraction.js`
**Triggered by:** Episode creation (`POST /episodes`)
**Model:** `qwen2.5:3b` via Ollama (configurable via `EXTRACTION_MODEL` env var)
**Model:** the utility model served by the inference service (`/utility/complete`), configurable via `UTILITY_MODEL` on the inference service
## Purpose
@@ -28,7 +28,7 @@ swallowed.
| Setting | Value | Notes |
|---|---|---|
| Model | `qwen2.5:3b` | Ollama, configurable via `EXTRACTION_MODEL` |
| Model | utility model | Served by inference-service `/utility/complete`, set via `UTILITY_MODEL` |
| Temperature | 0.1 | Low for consistent, deterministic output |
| `num_predict` | 1500 | Higher ceiling to accommodate entity + relationship JSON |
| `format` | `'json'` | Ollama constrained decoding — enforces valid JSON output |
+34 -2
View File
@@ -28,8 +28,7 @@ relationship extraction and embeds results into Qdrant.
| SQLITE_PATH | Yes | — | Path to SQLite database file |
| QDRANT_URL | No | http://localhost:6333 | Qdrant instance URL |
| EMBEDDING_SERVICE_URL | No | http://localhost:3003 | Embedding service URL |
| EXTRACTION_URL | No | http://localhost:11434 | Ollama URL for entity extraction |
| EXTRACTION_MODEL | No | qwen2.5:3b | Ollama model used for entity extraction |
| INFERENCE_SERVICE_URL | No | http://localhost:3001 | Inference service URL — entity extraction routes through its `/utility/complete` endpoint |
## Internal Structure
@@ -79,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
```
@@ -122,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
@@ -205,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
+30 -21
View File
@@ -30,8 +30,6 @@ or inference services — all traffic flows through orchestration.
| LLAMA_SERVER_URL | No | http://localhost:8080 | Direct llama-server URL for /models/props |
| QDRANT_URL | No | http://localhost:6333 | Qdrant URL for semantic search |
| CORS_ORIGIN | No | http://localhost:5173 | Allowed origin for CORS requests |
| EXTRACTION_URL | No | http://localhost:11434 | Ollama URL for summarisation |
| EXTRACTION_MODEL | No | qwen2.5:3b | Ollama model used for summarisation |
## Internal Structure
@@ -75,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 |
@@ -103,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;
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 |
+31 -42
View File
@@ -6,7 +6,7 @@ the full context window with raw episodes.
**Location:** `packages/orchestration-service/src/services/summarization.js`
**Triggered by:** `chat/index.js` after every episode write (fire-and-forget)
**Model:** `qwen2.5:3b` via Ollama on Mini PC 1 (192.168.0.81)
**Model:** the utility model served by the inference service (`/utility/complete`), set via `UTILITY_MODEL` (backed by Ollama on Mini PC 1, 192.168.0.81)
---
@@ -56,47 +56,50 @@ not all episodes in the session.
---
## Ollama Request
## Utility Inference Request
Summaries are generated through the shared `utilityInference()` helper, which
POSTs to the inference service's `/utility/complete` endpoint. `buildSummaryPrompt`
returns a plain instruction string (no template tags) passed as the `user` message:
```js
{
model: EXTRACTION_MODEL, // qwen2.5:3b (set via EXTRACTION_MODEL env var)
prompt: buildSummaryPrompt(episodesToSummarize, existingSummary),
stream: false,
// No format: 'json' — free-text output required for summaries
options: {
temperature: 0.2,
num_predict: 500,
},
}
const content = await utilityInference({
user: buildSummaryPrompt(episodesToSummarize, existingSummary),
temperature: SUMMARIES.TEMPERATURE, // 0.2
maxTokens: SUMMARIES.SESSION_GEN_MAX_TOKENS, // 500
});
```
`temperature: 0.2` is slightly higher than extraction (0.1) — summaries
benefit from some fluency. `num_predict: 500` gives room for 5 thorough
sentences without risk of runoff.
`TEMPERATURE` (0.2) is slightly higher than extraction (0.1) — summaries benefit
from some fluency. `SESSION_GEN_MAX_TOKENS` (500) gives room for ~5 thorough
sentences without runoff. Both live in `@nexusai/shared` `SUMMARIES` constants.
There is no `json: true` here — summaries are free-text, unlike entity extraction.
---
## Prompt Format
ChatML format — native to qwen2.5:
The prompt is plain text describing the task; the model's own prompt template
(ChatML for qwen, etc.) is applied **server-side** by the inference service via
Ollama's `/api/chat`. No `<|im_start|>` tags belong in this codebase, and the
utility model can be swapped (via `UTILITY_MODEL` on the inference service) with
no prompt changes here.
Fresh summary instruction:
```
<|im_start|>user
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}
<|im_end|>
<|im_start|>assistant
```
For cumulative updates, the instruction and context change:
Cumulative update instruction:
```
<|im_start|>user
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.
@@ -106,35 +109,22 @@ Previous summary:
New exchanges:
{context}
<|im_end|>
<|im_start|>assistant
```
### Input truncation
Episode context is truncated to `MAX_CHARS = 3000` characters, keeping the
most recent exchanges (sliced from the end). This keeps Qwen focused and
most recent exchanges (sliced from the end). This keeps the model focused and
prevents the prompt from exceeding its effective context window.
---
## ChatML Token Stripping
## Output Handling
Qwen occasionally echoes ChatML tokens back into its response. The raw output
is cleaned before saving:
```js
const raw = data.response?.trim() ?? '';
const content = raw
.replace(/<\|im_start\|>.*?<\|im_end\|>/gs, '')
.replace(/<\|im_start\|>|<\|im_end\|>|<\|im_sep\|>/g, '')
.trim();
return content;
```
Without this, leaked tokens get stored in the summary and then injected
back into the next summarisation prompt — causing the model to append a new
summary after the old one rather than replacing it.
Because `/api/chat` applies and removes the prompt template server-side, the
returned text is already clean — the previous ChatML token-stripping step (and
the class of bug where leaked tokens got stored and re-injected into the next
summarisation prompt) no longer applies.
---
@@ -207,8 +197,7 @@ Set in `packages/orchestration-service/src/.env`:
| Variable | Default | Description |
|---|---|---|
| `EXTRACTION_URL` | `http://localhost:11434` | Ollama instance URL |
| `EXTRACTION_MODEL` | `qwen2.5:3b` | Model for summarisation |
| `INFERENCE_SERVICE_URL` | `http://localhost:3001` | Inference service — summaries route through its `/utility/complete` endpoint (model set via `UTILITY_MODEL` there) |
| `MEMORY_SERVICE_URL` | `http://localhost:3002` | Memory service URL |
| `SUMMARY_THRESHOLD_TOKENS` | `200` | Token threshold before summarisation triggers |
| `SUMMARY_MAX_TOKENS` | `800` | Max summary length before a new row is created |
@@ -16,6 +16,14 @@ const migrations = [
// used to run (wrapped in try/catch) on every boot are folded into schema.js.
// Nothing to do here — this entry exists to mark v1 as the consolidated baseline.
(_db) => {},
(db) => {
db.exec(`
ALTER TABLE episodes ADD COLUMN last_accessed_at INTEGER;
ALTER TABLE episodes ADD COLUMN access_count INTEGER NOT NULL DEFAULT 0;
`)
db.exec(`UPDATE episodes SET last_accessed_at = created_at WHERE last_accessed_at IS NULL`)
}
];
const LATEST_VERSION = migrations.length;
+3 -1
View File
@@ -21,7 +21,9 @@ const schema = `
ai_response TEXT NOT NULL,
created_at INTEGER NOT NULL DEFAULT (unixepoch()),
token_count INTEGER,
metadata TEXT
metadata TEXT,
last_accessed_at INTEGER,
access_count INTEGER NOT NULL DEFAULT 0
);
CREATE TABLE IF NOT EXISTS entities (
+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,
};
+14 -1
View File
@@ -74,4 +74,17 @@ function getEpisodeIdsByEntities(entityIds) {
).all(...entityIds).map(r => r.episode_id);
}
module.exports = { getNeighborhood, getEntityNeighbors, getEpisodeIdsByEntities };
//Entity IDs linked (via entity_episodes) to any episode in a given session
//Scopes non-project entity search to the session's own entities under the "isolated chats" model.
//Entities globally deduped
function getEntityIdsBySession(sessionId){
const db = getDB();
return db.prepare(`
SELECT DISTINCT ee.entity_id
FROM entity_episodes ee
JOIN episodes e on e.id = ee.episode_id
WHERE e.session_id = ?
`).all(sessionId).map(r => r.entity_id);
}
module.exports = { getNeighborhood, getEntityNeighbors, getEpisodeIdsByEntities, getEntityIdsBySession };
+35 -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' });
@@ -162,12 +169,39 @@ app.get('/sessions/:id/episode-stats', (req, res) => {
res.json(episodic.getSessionEpisodeStats(Number(req.params.id)));
});
//Entity IDs linked to this session's episodes: sesion-scoped entity search
app.get('/sessions/:id/entity-ids', (req, res) => {
res.json({
entityIds: graph.getEntityIdsBySession(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.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);
@@ -1,4 +1,4 @@
const { SERVICES, getEnv, SUMMARIES } = require('@nexusai/shared');
const { SERVICES, getEnv, SUMMARIES, utilityInference } = require('@nexusai/shared');
const {
getSessionSummariesForProject,
getProjectOverviewSummary,
@@ -9,9 +9,6 @@ const {
const { getEpisodesByProject } = require('../episodic');
const { getProject } = require('../db/projects');
const EXTRACTION_URL = getEnv('EXTRACTION_URL', 'http://localhost:11434');
const EXTRACTION_MODEL = getEnv('EXTRACTION_MODEL', 'qwen2.5:3b');
const MAX_SUMMARY_CHARS = SUMMARIES.MAX_SUMMARY_CHARS; // generous ceiling before we truncate input
function buildProjectSummaryPrompt(projectName, sessionSummaries) {
@@ -24,8 +21,9 @@ function buildProjectSummaryPrompt(projectName, sessionSummaries) {
summaryBlock = summaryBlock.slice(-MAX_SUMMARY_CHARS);
}
// 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 [
'<|im_start|>user',
`The following are session summaries from a project called "${projectName}".`,
'Write a project overview covering: goals, progress, key decisions, and current state.',
'Scale the length to the material — use multiple paragraphs for complex projects, a few sentences for simple ones.',
@@ -33,8 +31,6 @@ function buildProjectSummaryPrompt(projectName, sessionSummaries) {
'Write in third person. Output only the overview text, no headings or labels.',
'',
summaryBlock,
'<|im_end|>',
'<|im_start|>assistant',
].join('\n');
}
@@ -49,8 +45,9 @@ function buildProjectSummaryFromEpisodesPrompt(projectName, episodes) {
episodeBlock = episodeBlock.slice(-MAX_SUMMARY_CHARS);
}
// 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 [
'<|im_start|>user',
`The following are conversations from a project called "${projectName}".`,
'Write a project overview covering: goals, progress, key decisions, and current state.',
'Scale the length to the material — use multiple paragraphs for complex projects, a few sentences for simple ones.',
@@ -58,58 +55,25 @@ function buildProjectSummaryFromEpisodesPrompt(projectName, episodes) {
'Write in third person. Output only the overview text, no headings or labels.',
'',
episodeBlock,
'<|im_end|>',
'<|im_start|>assistant',
].join('\n');
}
async function generateProjectSummaryFromEpisodes(projectName, episodes) {
const prompt = buildProjectSummaryFromEpisodesPrompt(projectName, episodes);
const res = await fetch(`${EXTRACTION_URL}/api/generate`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({
model: EXTRACTION_MODEL,
prompt,
stream: false,
options: { temperature: 0.2, num_predict: 1200 },
}),
const user = buildProjectSummaryFromEpisodesPrompt(projectName, episodes);
return utilityInference({
user,
temperature: SUMMARIES.TEMPERATURE,
maxTokens: SUMMARIES.PROJECT_GEN_MAX_TOKENS,
});
if (!res.ok) throw new Error(`Ollama responded ${res.status}`);
const data = await res.json();
const raw = data.response?.trim() ?? '';
return raw
.replace(/<\|im_start\|>.*?<\|im_end\|>/gs, '')
.replace(/<\|im_start\|>|<\|im_end\|>|<\|im_sep\|>/g, '')
.trim();
}
async function generateProjectSummary(projectName, sessionSummaries) {
const prompt = buildProjectSummaryPrompt(projectName, sessionSummaries);
const res = await fetch(`${EXTRACTION_URL}/api/generate`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({
model: EXTRACTION_MODEL,
prompt,
stream: false,
// No format: 'json' — we want free-text narrative, same as session summarization
options: { temperature: 0.2, num_predict: 1200 },
}),
const user = buildProjectSummaryPrompt(projectName, sessionSummaries);
return utilityInference({
user,
temperature: SUMMARIES.TEMPERATURE,
maxTokens: SUMMARIES.PROJECT_GEN_MAX_TOKENS,
});
if (!res.ok) throw new Error(`Ollama responded ${res.status}`);
const data = await res.json();
const raw = data.response?.trim() ?? '';
return raw
.replace(/<\|im_start\|>.*?<\|im_end\|>/gs, '')
.replace(/<\|im_start\|>|<\|im_end\|>|<\|im_sep\|>/g, '')
.trim();
}
// Main entry point — called by the route handler
@@ -127,16 +127,22 @@ async function getSemanticEpisodes(
}
}
async function getRelevantEntities(vector, projectId = null) {
async function getRelevantEntities(vector, { projectId = null, sessionId } = {}) {
if (!vector) return [];
try {
const results = await qdrant.searchEntities(vector, { projectId });
logger.info(
'[orchestration] Entity search results:',
results.map((r) => ({ name: r.payload?.name, score: r.score })),
);
// Include the Qdrant point ID (== SQLite entity ID) for graph traversal
return results.map((r) => r.payload ? { id: r.id, ...r.payload } : null).filter(Boolean);
let allowedIds;
if (projectId === null || projectId === undefined) {
// Non-project chat is its own island — scope to entities linked to
// THIS session. No links yet ⇒ nothing to retrieve, and we return
// early so searchEntities is never called unfiltered.
allowedIds = await memory.getEntityIdsBySession(sessionId);
logger.info(`[orchestration] Non-project chat, session ${sessionId}: ${allowedIds.length} linked entities`);
if (allowedIds.length === 0) return [];
}
const results = await qdrant.searchEntities(vector, { projectId, allowedIds });
logger.info('[orchestration] Entity search results:',
results.map(r => ({ name: r.payload?.name, score: r.score })));
return results.map(r => r.payload ? { id: r.id, ...r.payload } : null).filter(Boolean);
} catch (err) {
logger.debug('[orchestration] Entity search failed, continuing without:', err.message);
return [];
@@ -177,8 +183,11 @@ function fuseEpisodeResults(semanticEps, keywordEps, { semanticWeight, keywordWe
}
function estimateTokens(episode) {
return episode.token_count
?? Math.ceil((episode.user_message.length + episode.ai_response.length) / 4);
//NOTE: episode.token_count is not used here. It stores the
//full inference cost of that turn(prompt + injected memory + response),
//this inflate the episodes apparent size and starving the budget.
//Therefore, char/4 on the actual stored text is an honest selectedl
return Math.ceil((episode.user_message.length + episode.ai_response.length) / 4);
}
function buildScoredPool(fusedWithScores, recentEpisodes, entityBoostedIds, { entityWeight }) {
@@ -252,6 +261,7 @@ async function getFusedEpisodes(userMessage, session, recentIds, projectSessionI
}
async function assembleContext(externalId, userMessage) {
logger.info(`[orchestration] assembleContext ENTERED — msg: "${userMessage}", trivial: ${isTrivialTurn(userMessage)}`);
const settings = appSettings.load();
const { recentEpisodeLimit, semanticLimit, scoreThreshold,
temperature, repeatPenalty, topP, topK, systemPrompt,
@@ -299,7 +309,7 @@ async function assembleContext(externalId, userMessage) {
[fusedWithScores, entityResults] = await Promise.all([
getFusedEpisodes(userMessage, session, recentIds, projectSessionIds, queryVector, { semanticLimit, scoreThreshold, semanticWeight, keywordWeight }),
getRelevantEntities(queryVector, session.project_id ?? null),
getRelevantEntities(queryVector, { projectId: session.project_id ?? null, sessionId: session.id }),
]);
} else {
logger.debug('[orchestration] Trivial turn — skipping semantic/keyword/entity retrieval');
@@ -321,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);
@@ -216,6 +216,23 @@ async function getEpisodesByEntities(entityIds) {
return res.json(); // { episodeIds: [...] }
}
async function getEntityIdsBySession(sessionId){
const res = await fetch(`${BASE_URL}/sessions/${sessionId}/entity-ids`);
if (!res.ok) throw new Error(`Entity-ids-by-session error: ${res.status}`);
const { entityIds } = await res.json();
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,
@@ -242,4 +259,6 @@ module.exports = {
getProjectOverviewSummary,
searchEpisodes,
getEpisodesByEntities,
getEntityIdsBySession,
touchEpisodes,
}
@@ -29,30 +29,28 @@ async function searchEpisodes( vector, {limit = ORCHESTRATION.RECENT_EPISODE_LIM
return data.result;
}
async function searchEntities(vector, { limit = ORCHESTRATION.ENTITIES_LIMIT, scoreThreshold = ORCHESTRATION.ENTITIES_THRESHOLD, projectId = undefined } = {}) {
async function searchEntities(vector, { limit = ORCHESTRATION.ENTITIES_LIMIT, scoreThreshold = ORCHESTRATION.ENTITIES_THRESHOLD, projectId, allowedIds } = {}) {
const body = { vector, limit, score_threshold: scoreThreshold, with_payload: true };
if (projectId !== null && projectId !== undefined) {
body.filter = {
must: [{ key: 'projectId', match: { value: projectId } }]
};
// Project chat: entities shared across the project's sessions.
body.filter = { must: [{ key: 'projectId', match: { value: projectId } }] };
} else if (allowedIds && allowedIds.length > 0) {
// Non-project chat: restrict to this session's own entities (Model 2).
body.filter = { must: [{ has_id: allowedIds }] };
}
// No else: the caller returns early when a non-project session has no linked
// entities, so an unfiltered (leaky) search is never reached.
const res = await fetch(
`${BASE_URL}/collections/${COLLECTIONS.ENTITIES}/points/search`,
{
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify(body),
}
{ method: 'POST', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify(body) }
);
if (!res.ok) {
const body = await res.text();
throw new Error(`Qdrant error: ${res.status} - ${body}`);
const text = await res.text();
throw new Error(`Qdrant error: ${res.status} - ${text}`);
}
const data = await res.json();
return data.result;
return (await res.json()).result;
}
module.exports = { searchEpisodes, searchEntities };
@@ -1,7 +1,5 @@
const { getEnv, SERVICES, SUMMARIES, logger } = require('@nexusai/shared');
const { getEnv, SERVICES, SUMMARIES, logger, utilityInference } = require('@nexusai/shared');
const EXTRACTION_URL = getEnv('EXTRACTION_URL', 'http://localhost:11434');
const EXTRACTION_MODEL = getEnv('EXTRACTION_MODEL', 'qwen2.5:3b');
const MEMORY_URL = getEnv('MEMORY_SERVICE_URL', SERVICES.MEMORY_URL);
const THRESHOLD_TOKENS = parseInt(getEnv('SUMMARY_THRESHOLD_TOKENS', SUMMARIES.THRESHOLD_TOKENS));
@@ -35,41 +33,20 @@ Do not include greetings, sign-offs, or filler. Output only the summary text.
Conversation:
${context}`;
return [
'<|im_start|>user', // ChatML for qwen2.5
instruction,
'<|im_end|>',
'<|im_start|>assistant',
].join('\n');
// 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 prompt = buildSummaryPrompt(episodes, existingSummary);
const user = buildSummaryPrompt(episodes, existingSummary);
const res = await fetch(`${EXTRACTION_URL}/api/generate`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({
model: EXTRACTION_MODEL,
prompt,
stream: false,
options: {
temperature: 0.2, // slightly higher than entities — summaries benefit from some fluency
num_predict: 500, // generous but bounded — keeps summaries from running long
},
}),
const content = await utilityInference({
user,
temperature: SUMMARIES.TEMPERATURE,
maxTokens: SUMMARIES.SESSION_GEN_MAX_TOKENS,
});
if (!res.ok) throw new Error(`Ollama responded ${res.status}`);
const data = await res.json();
const raw = data.response?.trim() ?? '';
// Strip any leaked ChatML tokens Qwen echoes back
const content = raw
.replace(/<\|im_start\|>.*?<\|im_end\|>/gs, '')
.replace(/<\|im_start\|>|<\|im_end\|>|<\|im_sep\|>/g, '')
.trim();
return content;
}
+14
View File
@@ -78,6 +78,13 @@ const SUMMARIES = {
MIN_EPISODES_SINCE: 5, // don't resummarize until N new episodes since last summary
MAX_SUMMARY_CHARS: 8000, // max chars to include from recent episodes when generating summary (to control prompt size)
MAX_PROJECT_EPISODE_LIMIT: 200, // max number of episodes to consider from the entire project when generating summary (to control prompt size)
// Generation params for the utility model (passed to utilityInference).
// Distinct from MAX_SUMMARY_TOKENS above, which gates STORED summary size;
// these two cap GENERATION length (num_predict) per summary type.
TEMPERATURE: 0.2, // slightly higher than entities (0.1) — summaries benefit from some fluency
SESSION_GEN_MAX_TOKENS: 500, // num_predict for a session summary (3-5 sentences)
PROJECT_GEN_MAX_TOKENS: 1200, // num_predict for a project overview (multi-paragraph)
}
const ENTITIES = {
@@ -113,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,
@@ -128,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,
};
+11 -2
View File
@@ -25,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');
});
@@ -63,3 +65,10 @@ test('only migrations newer than the current version run', () => {
assert.deepStrictEqual(order, ['v2', 'v3'], 'v1 is not re-run');
assert.deepStrictEqual(db.setCalls, [2, 3]);
});
test('LATEST_VERSION reflects the appended v1 → v2 migration', () => {
// Guards that adding the access-tracking migration bumped the target version.
// The real SQL is verified out-of-band (scratch run against a DB copy); here we
// only lock the runner-visible contract: the array length drives LATEST_VERSION.
assert.strictEqual(LATEST_VERSION, 2);
});
+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);