Compare commits

...
37 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
Storme-bit b9195921dc utilityInference layer 2 2026-08-17 07:52:44 -07:00
Storme-bit 2adefe4df9 utilityInference layer 2 2026-08-17 07:48:04 -07:00
Storme-bit 2bee4b23d6 utilityInference layer 2 2026-08-17 07:41:41 -07:00
Storme-bit fab8f32395 utilityInference 2026-08-17 07:18:51 -07:00
Storme-bit ff1c0ab215 trivial turn patch 2026-08-17 06:23:00 -07:00
Storme-bit 7e8aead917 documentation updates 2026-08-17 05:31:54 -07:00
Storme-bit 1e7c11bad8 documentation updates 2026-08-17 05:31:19 -07:00
Storme-bit 1df6acc427 migration consolidation patch 2026-08-17 04:54:30 -07:00
Storme-bit 36c7cf2251 migration consolidation patch 2026-08-17 04:54:14 -07:00
Storme-bit 0353ce96e0 migration consolidation patch 2026-08-17 04:53:12 -07:00
Storme-bit e570bbecf2 testing 2026-08-17 03:43:06 -07:00
Storme-bit cbaf83ff17 testing 2026-08-17 03:29:00 -07:00
Storme-bit 3a6c09931f entities extraction fix 2026-08-17 02:37:03 -07:00
Storme-bit a4854cbfde entities extraction fix 2026-08-17 02:36:48 -07:00
34 changed files with 1090 additions and 388 deletions
+1
View File
@@ -6,4 +6,5 @@ data/
.env.* .env.*
*.db *.db
.claude/settings.local.json .claude/settings.local.json
.patch
EOF EOF
+23
View File
@@ -205,6 +205,9 @@ Returns `503` if llama-server is unreachable.
| `scoreThreshold` | float | 0–1 | Minimum similarity score for Qdrant results | | `scoreThreshold` | float | 0–1 | Minimum similarity score for Qdrant results |
| `semanticWeight` | float | 0–5 | RRF weight for Qdrant semantic results | | `semanticWeight` | float | 0–5 | RRF weight for Qdrant semantic results |
| `keywordWeight` | float | 0–5 | RRF weight for FTS5 keyword results (`0` = disabled) | | `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 | | `modelsFolderPath` | string | — | Path to folder containing .gguf files |
| `temperature` | float | 0–2 | Inference randomness | | `temperature` | float | 0–2 | Inference randomness |
| `repeatPenalty` | float | 1–2 | Repeat token penalty | | `repeatPenalty` | float | 1–2 | Repeat token penalty |
@@ -253,6 +256,7 @@ orchestration.
| GET | /sessions/by-external/:externalId | Get session by external ID | | GET | /sessions/by-external/:externalId | Get session by external ID |
| PATCH | /sessions/by-external/:externalId | Update session fields | | PATCH | /sessions/by-external/:externalId | Update session fields |
| DELETE | /sessions/by-external/:externalId | Delete session (cascades to episodes) | | 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` > Route ordering: `by-external/:externalId` must be defined before `/:id`
> to prevent `by-external` being captured as an ID param. > to prevent `by-external` being captured as an ID param.
@@ -277,6 +281,10 @@ Both fields are optional. Only provided fields are updated.
| GET | /episodes/search?q=&limit= | FTS keyword search across all episodes | | GET | /episodes/search?q=&limit= | FTS keyword search across all episodes |
| GET | /episodes/:id | Get episode by ID | | GET | /episodes/:id | Get episode by ID |
| GET | /sessions/:id/episodes?limit=&offset= | Paginated episodes for a session | | 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) | | DELETE | /episodes/:id | Delete episode (SQLite + Qdrant cleanup) |
> Route ordering: `/episodes/search` must be defined before `/episodes/:id`. > Route ordering: `/episodes/search` must be defined before `/episodes/:id`.
@@ -291,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 ### Projects
| Method | Path | Description | | Method | Path | Description |
+54
View File
@@ -0,0 +1,54 @@
# Testing
NexusAI has a lightweight regression suite built on Node's **built-in** test
runner (`node:test`) and assertions (`node:assert`) — no external test
dependencies. It runs offline on any node in the homelab.
```bash
npm test # = node --test (discovers test/*.test.js at the repo root)
```
## What's covered
| File | Subject | Imports the real… |
|---|---|---|
| `test/fusion.test.js` | RRF ranking math | `fuseEpisodeResults` (orchestration) |
| `test/fts-query.test.js` | Keyword tokenizer + FTS5 matching | `buildFtsQuery` (memory-service) |
| `test/llamacpp-stream.test.js` | SSE stream reassembly across chunks | `completeStream` (inference) |
| `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
functions are exported for this purpose. External calls (Ollama, Qdrant, the
memory service) are mocked via `global.fetch`; SQLite-backed tests use the
built-in `node:sqlite` module against a throwaway in-memory database.
## `node:sqlite` and skips
`test/schema.test.js` and three cases in `test/fts-query.test.js` need
`node:sqlite`, which requires **Node ≥ 22.5**. On older Node they `skip`
themselves cleanly (guarded by `{ skip: !DatabaseSync }`) rather than failing,
so the suite stays green everywhere — but those checks only *verify* anything on
a node new enough to run them. A dev machine on current Node is the source of
truth for the schema and FTS matching tests.
`node:sqlite` is still marked experimental, so runs print a one-line
`ExperimentalWarning`. It's harmless; `node --test --no-warnings` suppresses it.
## Adding tests
Keep the pattern: export the real function, import it, mock I/O at the boundary
(`global.fetch`) or use `node:sqlite` for DB behaviour. Prefer testing pure
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.
+11 -8
View File
@@ -68,12 +68,15 @@ The highest-leverage memory upgrade. Transforms NexusAI from "remembers conversa
Multi-strategy retrieval merged into a single ranked result set. Multi-strategy retrieval merged into a single ranked result set.
- [x] Reciprocal Rank Fusion (RRF) — merge semantic (Qdrant) + keyword (FTS5) results - [x] Reciprocal Rank Fusion (RRF) — merge semantic (Qdrant) + keyword (FTS5) results
- [x] Configurable weights per retrieval strategy (`semanticWeight`, `keywordWeight` via `PATCH /settings`) - [x] Configurable weights per retrieval strategy (`semanticWeight`, `keywordWeight` via `PATCH /settings`)
- [x] Score threshold retained per-strategy; FTS scoped to session/project sessions; `keywordWeight: 0` default (disabled until tuned) - [x] Score threshold retained per-strategy; FTS scoped to session/project sessions
- [x] FTS query tokenization (`buildFtsQuery`) — tokenize + stopword-filter + quoted `OR` terms; replaced whole-phrase matching that made keyword recall nil
- [x] Fusion tuned and verified live (keyword path + project scoping confirmed via weight inversion); running `keywordWeight: 0.5` / `semanticWeight: 1.0`
### 3. Memory Consolidation Lifecycle ### 3. Memory Consolidation Lifecycle
Prevents long-term memory degradation and enables compression. Prevents long-term memory degradation and enables compression.
- [ ] Episode aging — score/weight episodes by recency and access frequency - [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)`
- [ ] Consolidation pass — merge related low-weight episodes into summary nodes - [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 - [ ] Orphan cleanup — remove entities no longer referenced by active episodes
### 4. User Preference Model ### 4. User Preference Model
@@ -88,11 +91,11 @@ Short-circuit simple requests before they reach the LLM.
- [ ] Confidence bands — FAST PATH (memory lookup only) vs FULL (LLM + context) - [ ] Confidence bands — FAST PATH (memory lookup only) vs FULL (LLM + context)
- [ ] Fast-path handlers — direct memory queries, session lookups, factual recalls - [ ] 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. Budget-aware context selection instead of dumping all relevant memory into the prompt.
- [ ] Token budget manager in orchestration - [x] Token budget manager in orchestration (`selectWithinBudget`, char/4 estimation on stored text)
- [ ] Priority scoring — recency × relevance × entity weight - [x] Priority scoring — RRF fusion + recency + entity boost (`buildScoredPool`)
- [ ] Configurable context budget via env var - [x] Configurable via settings (`contextBudget`, `entityWeight`, `minRecentEpisodes`) — live, no restart
### 7. Procedural Memory Store *(inspired by acid2lake)* ### 7. Procedural Memory Store *(inspired by acid2lake)*
Learns "how NexusAI has successfully handled this type of request before." Learns "how NexusAI has successfully handled this type of request before."
@@ -225,4 +228,4 @@ The JARVIS moment — NexusAI reasons, plans, and acts across multiple steps.
--- ---
*Last updated: April 2026* *Last updated: August 2026*
+16 -2
View File
@@ -2,7 +2,7 @@
**Location:** `packages/memory-service/src/entities/extraction.js` **Location:** `packages/memory-service/src/entities/extraction.js`
**Triggered by:** Episode creation (`POST /episodes`) **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 ## Purpose
@@ -28,7 +28,7 @@ swallowed.
| Setting | Value | Notes | | 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 | | Temperature | 0.1 | Low for consistent, deterministic output |
| `num_predict` | 1500 | Higher ceiling to accommodate entity + relationship JSON | | `num_predict` | 1500 | Higher ceiling to accommodate entity + relationship JSON |
| `format` | `'json'` | Ollama constrained decoding — enforces valid JSON output | | `format` | `'json'` | Ollama constrained decoding — enforces valid JSON output |
@@ -131,6 +131,20 @@ check alone won't reject.
> imperfectly — switch to `/[^\p{L}\p{N}\s]/gu` if the set gains non-ASCII > imperfectly — switch to `/[^\p{L}\p{N}\s]/gu` if the set gains non-ASCII
> entries. > entries.
**Regurgitation guard (`mentionedIn`):** the extraction prompt feeds the model
a "known entities" hint block (the 20 most-recent entities) for spelling/type
consistency. The small model (qwen2.5:3b) will sometimes echo that list back as
if those entities appeared in the conversation — most visibly on contentless
turns (a greeting produced fake extractions of unrelated authors, game titles,
etc.). After parsing, each extracted name is checked against the actual
`userMessage + aiResponse` text (case- and whitespace-normalized substring); any
name not present is dropped before upsert. Since the prompt constrains names to
short proper nouns, a genuinely-discussed entity appears verbatim while a
regurgitated hint does not. Relationships referencing a dropped entity fall away
automatically (they resolve against the surviving `entityMap`). A prompt line
also tells the model the hint list is spelling-only — a backstop, with
`mentionedIn` as the deterministic guarantee.
## Relationship Processing ## Relationship Processing
After all entities are saved, relationships are processed: After all entities are saved, relationships are processed:
+78 -25
View File
@@ -28,16 +28,16 @@ relationship extraction and embeds results into Qdrant.
| SQLITE_PATH | Yes | — | Path to SQLite database file | | SQLITE_PATH | Yes | — | Path to SQLite database file |
| QDRANT_URL | No | http://localhost:6333 | Qdrant instance URL | | QDRANT_URL | No | http://localhost:6333 | Qdrant instance URL |
| EMBEDDING_SERVICE_URL | No | http://localhost:3003 | Embedding service URL | | EMBEDDING_SERVICE_URL | No | http://localhost:3003 | Embedding service URL |
| EXTRACTION_URL | No | http://localhost:11434 | Ollama URL for entity extraction | | INFERENCE_SERVICE_URL | No | http://localhost:3001 | Inference service URL — entity extraction routes through its `/utility/complete` endpoint |
| EXTRACTION_MODEL | No | qwen2.5:3b | Ollama model used for entity extraction |
## Internal Structure ## Internal Structure
``` ```
src/ src/
├── db/ ├── db/
│ ├── index.js # SQLite connection + initialization + migrations │ ├── index.js # SQLite connection + init + migrate() + one-time FTS backfill
│ ├── schema.js # Table definitions, indexes, FTS5, triggers │ ├── migrations.js # Forward-only versioned migration runner (PRAGMA user_version)
│ ├── schema.js # Complete current shape: tables, indexes, FTS5, triggers
│ ├── projects.js # Project CRUD functions │ ├── projects.js # Project CRUD functions
│ └── summaries.js # Summary CRUD functions │ └── summaries.js # Summary CRUD functions
├── episodic/ ├── episodic/
@@ -64,36 +64,61 @@ Eight core tables:
- **summaries** — condensed episode groups for efficient context retrieval - **summaries** — condensed episode groups for efficient context retrieval
- **projects** — named groupings of sessions with `name`, `description`, `colour`, `icon`, `isolated`, `notes`, `system_prompt` - **projects** — named groupings of sessions with `name`, `description`, `colour`, `icon`, `isolated`, `notes`, `system_prompt`
### Migrations ### Schema & Migrations
Schema changes that cannot use `CREATE TABLE IF NOT EXISTS` are applied as `schema.js` holds the **complete current shape** — every table, column, index,
idempotent migrations in `db/index.js` at startup: the FTS5 virtual table, and its triggers — as the single source of truth for a
fresh database. It uses `CREATE TABLE IF NOT EXISTS`, so on a fresh DB it builds
everything; on an existing DB it skips tables that already exist (and therefore
does **not** reconcile columns on old tables — that's what migrations are for).
`db/migrations.js` is a forward-only versioned runner keyed on
`PRAGMA user_version`:
```js ```js
try { db.exec(`ALTER TABLE sessions ADD COLUMN name TEXT`); } catch {} const migrations = [
try { db.exec(`ALTER TABLE sessions ADD COLUMN project_id INTEGER REFERENCES projects(id)`); } catch {} (_db) => {}, // v0 → v1: consolidated baseline (historical ALTERs folded into schema.js)
try { db.exec(`CREATE INDEX IF NOT EXISTS idx_sessions_project ON sessions(project_id)`); } catch {} (db) => { // v1 → v2: access tracking for consolidation lifecycle
try { db.exec(`ALTER TABLE projects ADD COLUMN isolated INTEGER NOT NULL DEFAULT 0`); } catch {} db.exec(`ALTER TABLE episodes ADD COLUMN last_accessed_at INTEGER`);
try { db.exec(`ALTER TABLE projects ADD COLUMN notes TEXT`); } catch {} db.exec(`ALTER TABLE episodes ADD COLUMN access_count INTEGER NOT NULL DEFAULT 0`);
try { db.exec(`ALTER TABLE projects ADD COLUMN system_prompt TEXT`); } catch {} db.exec(`UPDATE episodes SET last_accessed_at = created_at`); // backfill
// Knowledge graph columns: },
try { db.exec(`ALTER TABLE entities ADD COLUMN mention_count INTEGER NOT NULL DEFAULT 1`) } catch {} ];
try { db.exec(`ALTER TABLE entities ADD COLUMN confidence REAL NOT NULL DEFAULT 1.0`) } catch {} const LATEST_VERSION = migrations.length; // derived, never hand-maintained
try { db.exec(`ALTER TABLE entities ADD COLUMN source TEXT NOT NULL DEFAULT 'extraction'`) } catch {}
try { db.exec(`ALTER TABLE entities ADD COLUMN last_seen_at INTEGER`) } catch {}
try { db.exec(`ALTER TABLE relationships ADD COLUMN mention_count INTEGER NOT NULL DEFAULT 1`) } catch {}
try { db.exec(`ALTER TABLE relationships ADD COLUMN notes TEXT`) } catch {}
``` ```
`entity_episodes` is defined in `schema.js` itself (not a migration) since it is a new table. `migrate(db)` reads `user_version`, applies every entry newer than it (each in a
transaction alongside its version bump), and stamps the result. A fresh DB is
built whole by `schema.js` and simply stamped to `LATEST_VERSION`; the baseline
entry is a no-op.
New migrations are always appended — never modify the schema file for existing tables since `ALTER TABLE` cannot use `IF NOT EXISTS`. **Adding a schema change:** append a new function to the `migrations` array
(which bumps `LATEST_VERSION` automatically). Never edit an already-shipped
entry, and never edit a table in `schema.js` expecting existing DBs to pick it
up — they won't. This replaces the previous pattern of stacking silent
`try/catch ALTER TABLE` statements in `db/index.js` on every boot.
> **Consolidation note:** the historical ALTERs were folded into `schema.js`
> rather than preserved as replayable migrations, so this assumes a fresh
> database (which is the case post-wipe). An older, pre-consolidation database
> would **not** auto-upgrade — `schema.js` skips its existing tables and the
> baseline migration is a no-op. To support upgrading old DBs, the v1 baseline
> would instead perform guarded (`ADD COLUMN if missing`) catch-up.
### FTS5 Full-Text Search ### FTS5 Full-Text Search
An `episodes_fts` virtual table enables keyword search across all episodes. An `episodes_fts` external-content virtual table enables keyword search across
Three triggers (`episodes_fts_insert`, `episodes_fts_update`, `episodes_fts_delete`) episodes. Three triggers (`episodes_fts_insert`, `episodes_fts_update`,
keep the FTS index automatically in sync with the episodes table. `episodes_fts_delete`) keep the index in sync with the `episodes` table
automatically during normal operation.
A one-time backfill in `db/index.js` handles the case where the FTS table is
created on a DB that already holds episodes (e.g. episodes predating FTS). It is
gated on "did `episodes_fts` not exist before this boot," checked via
`sqlite_master` **before** running the schema — not on a row-count comparison,
because `COUNT(*)` on an external-content FTS5 table proxies the content table
and cannot detect a desync. This replaced an unconditional full FTS rebuild that
previously ran on every startup.
### SQLite Configuration ### SQLite Configuration
@@ -101,6 +126,12 @@ keep the FTS index automatically in sync with the episodes table.
- `foreign_keys = ON` — enforces referential integrity and cascade deletes - `foreign_keys = ON` — enforces referential integrity and cascade deletes
- PRAGMAs set via `db.pragma()`, not `db.exec()` - 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 ### Dynamic Updates
Both `updateSession` and `updateProject` build their `SET` clause dynamically Both `updateSession` and `updateProject` build their `SET` clause dynamically
@@ -184,6 +215,28 @@ service is responsible only for CRUD — generation logic lives in orchestration
> For full details on trigger conditions, prompt format, cumulative updates, > For full details on trigger conditions, prompt format, cumulative updates,
> and ChatML token stripping, see `summarization.md`. > 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) ## Delete Behaviour (SQLite + Qdrant consistency)
SQLite cascades handle relational cleanup, but Qdrant is a separate store and SQLite cascades handle relational cleanup, but Qdrant is a separate store and
+31 -22
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 | | 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 | | QDRANT_URL | No | http://localhost:6333 | Qdrant URL for semantic search |
| CORS_ORIGIN | No | http://localhost:5173 | Allowed origin for CORS requests | | 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 ## 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 | | `scoreThreshold` | 0.5 | Minimum similarity score for Qdrant semantic results |
| `semanticWeight` | 1.0 | RRF weight for Qdrant semantic results | | `semanticWeight` | 1.0 | RRF weight for Qdrant semantic results |
| `keywordWeight` | 0 | RRF weight for FTS5 keyword results (`0` = disabled) | | `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 | | `modelsFolderPath` | `/mnt/nexus-models` | Path to folder containing .gguf files |
| `temperature` | 0.7 | Inference temperature | | `temperature` | 0.7 | Inference temperature |
| `repeatPenalty` | 1.1 | Repeat token penalty | | `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`). 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). search in parallel, then merges results via Reciprocal Rank Fusion (RRF).
Both paths are filtered against `recentIds` before fusion. FTS is scoped The query is embedded once and shared with entity search. Both paths are
to the current session or all project sessions. If `keywordWeight` is `0`, filtered against `recentIds` before fusion. FTS is scoped to the current
the FTS call is skipped entirely. Non-critical — failures fall back to session or all project sessions. If `keywordWeight` is `0`, the FTS call
whichever strategy succeeded. is skipped entirely. Non-critical — failures fall back to whichever
strategy succeeded.
6. **Entity search** — query `entities` Qdrant collection filtered by 7. **Entity search + graph expansion** — query `entities` Qdrant collection
`projectId`. Returns entity IDs alongside Qdrant payload data (the Qdrant (project-scoped, or session-scoped via `/sessions/:id/entity-ids` for
point ID equals the SQLite entity ID). Non-critical. 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 8. **Scored pool + budget selection** — `buildScoredPool` combines RRF scores,
memory-service with the entity IDs from step 6. Returns a 1-hop subgraph recency, and entity-linkage bonus; `selectWithinBudget` fills `contextBudget`
`{ nodes, edges }` — entity objects plus the relationships connecting them. (char/4 token estimation on stored text) above a guaranteed floor of
If no entities were found or the graph call fails, falls back to flat entity `minRecentEpisodes` recent episodes. Selected episode IDs are then reported
list (no edges). Non-critical. to `POST /episodes/touch` fire-and-forget (access tracking for the
consolidation lifecycle).
8. **Prompt assembly** — combine system prompt, graph context, fused episodes, 9. **Prompt assembly** — combine system prompt, graph context, selected
recent episodes, and user message. 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. `/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. 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. inference call (max 20 tokens, temperature 0.3) to generate a session name.
### Prompt Structure ### Prompt Structure
+40 -3
View File
@@ -37,9 +37,10 @@ Fusion lives in orchestration — the service already coordinates multiple data
sources, and fusion is a retrieval strategy, not a storage concern. sources, and fusion is a retrieval strategy, not a storage concern.
``` ```
getFusedEpisodes() getFusedEpisodes(…, queryVector, …)
├── getSemanticEpisodes() — Qdrant embed+search → fetch full rows by ID ├── getSemanticEpisodes(queryVector) — Qdrant search → fetch full rows by ID
│ (existing path, unchanged) │ (query embedded ONCE upstream in assembleContext and shared with entity
│ search — no longer embedded separately here)
└── getFTSResults() — memory-service /episodes/search → full rows directly └── getFTSResults() — memory-service /episodes/search → full rows directly
(skipped entirely if keywordWeight == 0) (skipped entirely if keywordWeight == 0)
↓ ↓
@@ -48,6 +49,10 @@ fuseEpisodeResults() — pure RRF, no I/O
fusedEpisodes[] — top semanticLimit episodes by RRF score fusedEpisodes[] — top semanticLimit episodes by RRF score
``` ```
The query embedding is computed once per turn in `assembleContext` and passed
into both fused retrieval and entity search; if embedding fails, both receive
`null` and degrade to empty results rather than erroring.
### Data Shape Consistency ### Data Shape Consistency
Both sides must enter fusion as `Episode[]` — full SQLite row objects with Both sides must enter fusion as `Episode[]` — full SQLite row objects with
@@ -59,6 +64,38 @@ the same shape — and both must be filtered against `recentIds` first:
FTS requests `semanticLimit * 2` results to provide headroom for the FTS requests `semanticLimit * 2` results to provide headroom for the
`recentIds` filter without under-serving the fusion. `recentIds` filter without under-serving the fusion.
## Query Tokenization
Before FTS5 sees the query, `buildFtsQuery(query)` (in
`memory-service/src/episodic/index.js`) turns the raw message into a MATCH
expression:
1. Lowercase and split on any non-letter/number (`/[^\p{L}\p{N}]+/u`, unicode-aware)
2. Drop stopwords (a small `FTS_STOPWORDS` set of common function words) and single-character tokens
3. Wrap each surviving token in double quotes and join with ` OR `
So `"How do I configure the Qdrant collection?"` becomes
`"configure" OR "qdrant" OR "collection"`. If nothing survives (an
all-stopword message like `"how do I do it?"`), it returns `null` and
`searchEpisodes` returns `[]` — keyword search sits out that turn and
semantic retrieval carries it.
**Why this matters:** the earlier implementation quoted the *entire* message
as one FTS5 phrase, which required the whole string to appear verbatim in an
episode — so keyword recall was effectively nil for conversational queries.
Tokenizing into OR-joined terms is what makes `keywordWeight > 0` actually
contribute anything.
**Injection safety:** quoting each token individually means any
FTS5-significant token inside the user's message (a literal `OR`, `*`, `"`,
etc.) is matched as a search term rather than parsed as an operator. This
replaces the safety the old whole-phrase quoting provided.
The stopword set is deliberately conservative and tuned iteratively — high
frequency filler (`the`, `is`, `one`, `there`, …) is dropped, but borderline
words that can carry signal (`time`, `good`, `way`) are kept. Add to the set
when a common word is observed producing noisy matches.
## FTS Session Scoping ## FTS Session Scoping
Without scoping, FTS5 searches across all episodes in the database. For Without scoping, FTS5 searches across all episodes in the database. For
+10
View File
@@ -203,6 +203,16 @@ SUMMARY_MAX_TOKENS=800
SUMMARY_MIN_EPISODES=5 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` #### `SQLITE`
| Key | Value | Description | | Key | Value | Description |
+46 -43
View File
@@ -6,13 +6,27 @@ the full context window with raw episodes.
**Location:** `packages/orchestration-service/src/services/summarization.js` **Location:** `packages/orchestration-service/src/services/summarization.js`
**Triggered by:** `chat/index.js` after every episode write (fire-and-forget) **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)
--- ---
## Trigger Conditions ## Trigger Conditions
`triggerSummary(session, allEpisodes)` calls `maybeSummarize` fire-and-forget. `triggerSummary(session)` calls `maybeSummarize` fire-and-forget. It takes only
the `session` — it no longer receives the full episode list. Instead
`maybeSummarize` fetches exactly what it needs:
1. `GET /sessions/:id/episode-stats` — a cheap aggregate (`COUNT`, `SUM(token_count)`,
`MAX(id)`) that gates the token threshold **without** pulling every episode row
2. Only if over threshold: `GET /sessions/:id/episodes/since/:afterId` — the
un-summarized tail (episodes newer than the last summary's range), fetched in
full for the summary prompt
This replaced an earlier approach where `chat/index.js` fetched the entire
session (`getRecentEpisodes(session.id, 9999)`) on every message just to hand it
over — an O(session length) cost per turn. The stats query now runs on every
message; full episode text is fetched only when a summary actually fires.
`maybeSummarize` proceeds only when both conditions are met: `maybeSummarize` proceeds only when both conditions are met:
1. Total session token count exceeds `SUMMARIES.THRESHOLD_TOKENS` (default 200) 1. Total session token count exceeds `SUMMARIES.THRESHOLD_TOKENS` (default 200)
@@ -42,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 ```js
{ const content = await utilityInference({
model: EXTRACTION_MODEL, // qwen2.5:3b (set via EXTRACTION_MODEL env var) user: buildSummaryPrompt(episodesToSummarize, existingSummary),
prompt: buildSummaryPrompt(episodesToSummarize, existingSummary), temperature: SUMMARIES.TEMPERATURE, // 0.2
stream: false, maxTokens: SUMMARIES.SESSION_GEN_MAX_TOKENS, // 500
// No format: 'json' — free-text output required for summaries });
options: {
temperature: 0.2,
num_predict: 500,
},
}
``` ```
`temperature: 0.2` is slightly higher than extraction (0.1) — summaries `TEMPERATURE` (0.2) is slightly higher than extraction (0.1) — summaries benefit
benefit from some fluency. `num_predict: 500` gives room for 5 thorough from some fluency. `SESSION_GEN_MAX_TOKENS` (500) gives room for ~5 thorough
sentences without risk of runoff. 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 ## 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. Summarize the conversation below in 3-5 sentences.
Write in third person. Do not quote directly — paraphrase only. Write in third person. Do not quote directly — paraphrase only.
Do not include greetings, sign-offs, or filler. Output only the summary text. Do not include greetings, sign-offs, or filler. Output only the summary text.
Conversation: Conversation:
{context} {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. Update the summary below to incorporate the new exchanges.
Write 3-5 sentences in third person. Do not quote directly — paraphrase only. 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. Do not include greetings, sign-offs, or filler. Output only the updated summary text.
@@ -92,35 +109,22 @@ Previous summary:
New exchanges: New exchanges:
{context} {context}
<|im_end|>
<|im_start|>assistant
``` ```
### Input truncation ### Input truncation
Episode context is truncated to `MAX_CHARS = 3000` characters, keeping the 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. 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 Because `/api/chat` applies and removes the prompt template server-side, the
is cleaned before saving: 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
```js summarisation prompt) no longer applies.
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.
--- ---
@@ -193,8 +197,7 @@ Set in `packages/orchestration-service/src/.env`:
| Variable | Default | Description | | Variable | Default | Description |
|---|---|---| |---|---|---|
| `EXTRACTION_URL` | `http://localhost:11434` | Ollama instance URL | | `INFERENCE_SERVICE_URL` | `http://localhost:3001` | Inference service — summaries route through its `/utility/complete` endpoint (model set via `UTILITY_MODEL` there) |
| `EXTRACTION_MODEL` | `qwen2.5:3b` | Model for summarisation |
| `MEMORY_SERVICE_URL` | `http://localhost:3002` | Memory service URL | | `MEMORY_SERVICE_URL` | `http://localhost:3002` | Memory service URL |
| `SUMMARY_THRESHOLD_TOKENS` | `200` | Token threshold before summarisation triggers | | `SUMMARY_THRESHOLD_TOKENS` | `200` | Token threshold before summarisation triggers |
| `SUMMARY_MAX_TOKENS` | `800` | Max summary length before a new row is created | | `SUMMARY_MAX_TOKENS` | `800` | Max summary length before a new row is created |
+2
View File
@@ -2,6 +2,7 @@ require ('dotenv').config();
const express = require('express'); const express = require('express');
const {getEnv, PORTS, OLLAMA, logger} = require('@nexusai/shared'); const {getEnv, PORTS, OLLAMA, logger} = require('@nexusai/shared');
const inferenceRouter = require('./routes/inference'); const inferenceRouter = require('./routes/inference');
const utilityRouter = require('./routes/utility')
const app = express(); const app = express();
app.use(express.json({ limit: '8mb' })); // prompts include full context window app.use(express.json({ limit: '8mb' })); // prompts include full context window
@@ -20,6 +21,7 @@ app.get('/health', (req, res) => {
}); });
}); });
app.use('/', utilityRouter)
app.use('/', inferenceRouter); app.use('/', inferenceRouter);
// Start the server // Start the server
@@ -0,0 +1,21 @@
const { Router } = require('express');
const { logger } = require('@nexusai/shared');
const { utilityComplete } = require('../utility');
const router = Router();
router.post('/utility/complete', async (req, res) => {
const { system, user, json, temperature, maxTokens } = req.body;
if (!user) return res.status(400).json({ error: 'user message is required' });
try {
const result = await utilityComplete({ system, user, json, temperature, maxTokens });
res.json(result);
} catch (error) {
logger.error('[Utility] Completion error:', error.message);
res.status(500).json({ error: 'Utility inference failed', detail: error.message });
}
});
module.exports = router;
+44
View File
@@ -0,0 +1,44 @@
const { getEnv, UTILITY } = require ('@nexusai/shared')
const UTILITY_URL = getEnv('UTILITY_URL', UTILITY.DEFAULT_URL);
const UTILITY_MODEL = getEnv('UTILITY_MODEL', UTILITY.DEFAULT_MODEL);
// Background task inference (extraction/summarization). Uses ollama's /api/chat
// so the model's own prompt template is server-side, no need for ChatML
async function utilityComplete({
system,
user,
json = false,
temperature,
maxTokens
}) {
const messages = [];
if(system) messages.push({ role: 'system', content: system});
messages.push ({role: 'user', content: user});
const res = await fetch (`${UTILITY_URL}/api/chat`, {
method: 'POST',
headers: { 'Content-Type': 'application/json'},
body: JSON.stringify({
model: UTILITY_MODEL,
messages,
stream: false,
...(json && {format: 'json' }),
options: {
temperature: temperature ?? UTILITY.TEMPERATURE,
num_predict: maxTokens ?? UTILITY.MAX_TOKENS,
},
}),
signal: AbortSignal.timeout(UTILITY.TIMEOUT_MS),
});
if (!res.ok) throw new Error(`Utility backend responded ${res.status}`);
const data = await res.json();
return {
text: (data.message?.content ?? '').trim(),
model: data.model,
}
}
module.exports = { utilityComplete};
+24 -54
View File
@@ -1,6 +1,7 @@
const Database = require('better-sqlite3'); const Database = require('better-sqlite3');
const schema = require('./schema'); const schema = require('./schema');
const {getEnv, SQLITE, logger } = require('@nexusai/shared'); const { migrate } = require('./migrations');
const { getEnv, SQLITE, logger } = require('@nexusai/shared');
let db; // Declare db variable in a scope accessible to all functions let db; // Declare db variable in a scope accessible to all functions
@@ -12,60 +13,29 @@ function getDB() {
db.pragma('journal_mode = WAL'); db.pragma('journal_mode = WAL');
db.pragma('foreign_keys = ON'); db.pragma('foreign_keys = ON');
db.exec(schema); // Was the FTS index absent before this boot? (True for a fresh DB, and for
// an older DB from before FTS existed.) Checked BEFORE schema runs so we can
// decide whether a one-time backfill is needed below.
const ftsExisted = db.prepare(
`SELECT 1 FROM sqlite_master WHERE type='table' AND name='episodes_fts'`
).get() !== undefined;
try{ db.exec(schema); // complete current shape — fresh DBs get everything
db.exec(`ALTER TABLE sessions ADD COLUMN name TEXT`) migrate(db); // carry an older DB forward; no-op on fresh/current DBs
} catch {}
try { // One-time FTS backfill: only when the index was just created on a DB that
db.exec(`ALTER TABLE sessions ADD COLUMN project_id INTEGER REFERENCES projects(id)`); // already holds episodes (i.e. episodes predate FTS). During normal
} catch {} // operation the insert/delete/update triggers keep it in sync, so this no
// longer rebuilds the whole index on every boot. NOTE: COUNT(*) on an
try { // external-content FTS5 table proxies the content table, so it can't detect
db.exec(`CREATE INDEX IF NOT EXISTS idx_sessions_project ON sessions(project_id)`); // a desync — the "was it just created" check is what makes this correct.
} catch {} if (!ftsExisted) {
const epCount = db.prepare('SELECT COUNT(*) AS c FROM episodes').get().c;
try { if (epCount > 0) {
db.exec(`ALTER TABLE projects ADD COLUMN isolated INTEGER NOT NULL DEFAULT 0`); db.exec(`INSERT INTO episodes_fts(episodes_fts) VALUES('rebuild')`);
} catch {} logger.info(`[db] Backfilled FTS index for ${epCount} pre-existing episodes`);
}
try { }
db.exec(`ALTER TABLE projects ADD COLUMN notes TEXT`); // ← add this
} catch {}
try {
db.exec(`ALTER TABLE projects ADD COLUMN system_prompt TEXT`);
} catch {}
try {
db.exec(`ALTER TABLE summaries ADD COLUMN project_id INTEGER REFERENCES projects(id) ON DELETE CASCADE`);
} catch {}
try {
db.exec(`ALTER TABLE summaries ADD COLUMN token_count INTEGER`);
} catch {}
try {
db.exec(`CREATE INDEX IF NOT EXISTS idx_summaries_project ON summaries(project_id)`);
} catch {}
try {
db.exec(`CREATE INDEX IF NOT EXISTS idx_summaries_session ON summaries(session_id)`);
} catch {}
try { db.exec(`ALTER TABLE entities ADD COLUMN mention_count INTEGER NOT NULL DEFAULT 1`) } catch {}
try { db.exec(`ALTER TABLE entities ADD COLUMN confidence REAL NOT NULL DEFAULT 1.0`) } catch {}
try { db.exec(`ALTER TABLE entities ADD COLUMN source TEXT NOT NULL DEFAULT 'extraction'`) } catch {}
try { db.exec(`ALTER TABLE entities ADD COLUMN last_seen_at INTEGER`) } catch {}
try { db.exec(`ALTER TABLE relationships ADD COLUMN mention_count INTEGER NOT NULL DEFAULT 1`) } catch {}
try { db.exec(`ALTER TABLE relationships ADD COLUMN notes TEXT`) } catch {}
// Sync FTS index with any existing episodes data
db.exec(`INSERT OR REPLACE INTO episodes_fts(rowid, user_message, ai_response)
SELECT id, user_message, ai_response FROM episodes`);
logger.info(`Connected to SQLite database at ${path}`); logger.info(`Connected to SQLite database at ${path}`);
} }
@@ -74,4 +44,4 @@ function getDB() {
module.exports = { module.exports = {
getDB getDB
}; };
@@ -0,0 +1,51 @@
const { logger } = require('@nexusai/shared');
// Forward-only schema migrations. Entry i takes the database from user_version i
// to i+1, so migrations[0] is the v0→v1 step, migrations[1] the v1→v2 step, etc.
//
// schema.js already holds the COMPLETE current shape, so a fresh database is
// created whole and stamped straight to LATEST_VERSION — these run only to carry
// an OLDER database forward. Add a new schema change by appending a function here
// (which bumps LATEST_VERSION by one); never edit an existing entry once shipped,
// since databases already stamped past it will not re-run it.
//
// Each migration receives the better-sqlite3 db handle and runs inside a
// transaction together with its version bump, so a failure rolls back cleanly.
const migrations = [
// v0 → v1: baseline. The historical ALTER TABLE / CREATE INDEX statements that
// 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;
// Applies any migrations newer than the database's current user_version, one at a
// time, each with its version bump, inside a transaction. Returns the resulting
// version. A no-op when the database is already current. The migrations list and
// target version are injectable for testing; production callers pass just `db`.
function migrate(db, migs = migrations, latest = migs.length) {
const current = db.pragma('user_version', { simple: true });
if (current >= latest) return current;
for (let v = current; v < latest; v++) {
const step = db.transaction(() => {
migs[v](db);
db.pragma(`user_version = ${v + 1}`);
});
step();
}
logger.info(`[db] schema migrated ${current} -> ${latest}`);
return latest;
}
module.exports = { migrate, migrations, LATEST_VERSION };
+54 -40
View File
@@ -1,40 +1,55 @@
// Complete current schema — the single source of truth for a fresh database.
// Historical ALTER TABLE statements that used to run on every boot are folded in
// here. Any change that must reach EXISTING databases goes in migrations.js as a
// new numbered migration, NOT by editing a table below (CREATE ... IF NOT EXISTS
// silently skips tables that already exist, so column edits here never reach them).
const schema = ` const schema = `
CREATE TABLE IF NOT EXISTS sessions ( CREATE TABLE IF NOT EXISTS sessions (
id INTEGER PRIMARY KEY AUTOINCREMENT, id INTEGER PRIMARY KEY AUTOINCREMENT,
external_id TEXT UNIQUE NOT NULL, external_id TEXT UNIQUE NOT NULL,
created_at INTEGER NOT NULL DEFAULT (unixepoch()), created_at INTEGER NOT NULL DEFAULT (unixepoch()),
updated_at INTEGER NOT NULL DEFAULT (unixepoch()), updated_at INTEGER NOT NULL DEFAULT (unixepoch()),
metadata TEXT metadata TEXT,
name TEXT,
project_id INTEGER REFERENCES projects(id)
); );
CREATE TABLE IF NOT EXISTS episodes ( CREATE TABLE IF NOT EXISTS episodes (
id INTEGER PRIMARY KEY AUTOINCREMENT, id INTEGER PRIMARY KEY AUTOINCREMENT,
session_id INTEGER NOT NULL REFERENCES sessions(id) ON DELETE CASCADE, session_id INTEGER NOT NULL REFERENCES sessions(id) ON DELETE CASCADE,
user_message TEXT NOT NULL, user_message TEXT NOT NULL,
ai_response TEXT NOT NULL, ai_response TEXT NOT NULL,
created_at INTEGER NOT NULL DEFAULT (unixepoch()), created_at INTEGER NOT NULL DEFAULT (unixepoch()),
token_count INTEGER, token_count INTEGER,
metadata TEXT metadata TEXT,
last_accessed_at INTEGER,
access_count INTEGER NOT NULL DEFAULT 0
); );
CREATE TABLE IF NOT EXISTS entities ( CREATE TABLE IF NOT EXISTS entities (
id INTEGER PRIMARY KEY AUTOINCREMENT, id INTEGER PRIMARY KEY AUTOINCREMENT,
name TEXT NOT NULL, name TEXT NOT NULL,
type TEXT NOT NULL, type TEXT NOT NULL,
notes TEXT, notes TEXT,
created_at INTEGER NOT NULL DEFAULT (unixepoch()), created_at INTEGER NOT NULL DEFAULT (unixepoch()),
updated_at INTEGER NOT NULL DEFAULT (unixepoch()), updated_at INTEGER NOT NULL DEFAULT (unixepoch()),
metadata TEXT, metadata TEXT,
mention_count INTEGER NOT NULL DEFAULT 1,
confidence REAL NOT NULL DEFAULT 1.0,
source TEXT NOT NULL DEFAULT 'extraction',
last_seen_at INTEGER,
UNIQUE(name, type) UNIQUE(name, type)
); );
CREATE TABLE IF NOT EXISTS relationships ( CREATE TABLE IF NOT EXISTS relationships (
id INTEGER PRIMARY KEY AUTOINCREMENT, id INTEGER PRIMARY KEY AUTOINCREMENT,
from_id INTEGER NOT NULL REFERENCES entities(id) ON DELETE CASCADE, from_id INTEGER NOT NULL REFERENCES entities(id) ON DELETE CASCADE,
to_id INTEGER NOT NULL REFERENCES entities(id) ON DELETE CASCADE, to_id INTEGER NOT NULL REFERENCES entities(id) ON DELETE CASCADE,
label TEXT NOT NULL, label TEXT NOT NULL,
created_at INTEGER NOT NULL DEFAULT (unixepoch()), created_at INTEGER NOT NULL DEFAULT (unixepoch()),
metadata TEXT, metadata TEXT,
mention_count INTEGER NOT NULL DEFAULT 1,
notes TEXT,
UNIQUE(from_id, to_id, label) UNIQUE(from_id, to_id, label)
); );
@@ -49,16 +64,17 @@ const schema = `
CREATE INDEX IF NOT EXISTS idx_entity_episodes_entity ON entity_episodes(entity_id); CREATE INDEX IF NOT EXISTS idx_entity_episodes_entity ON entity_episodes(entity_id);
CREATE INDEX IF NOT EXISTS idx_entity_episodes_episode ON entity_episodes(episode_id); CREATE INDEX IF NOT EXISTS idx_entity_episodes_episode ON entity_episodes(episode_id);
CREATE TABLE IF NOT EXISTS projects ( CREATE TABLE IF NOT EXISTS projects (
id INTEGER PRIMARY KEY AUTOINCREMENT, id INTEGER PRIMARY KEY AUTOINCREMENT,
name TEXT NOT NULL, name TEXT NOT NULL,
description TEXT, description TEXT,
colour TEXT, colour TEXT,
icon TEXT, icon TEXT,
created_at INTEGER NOT NULL DEFAULT (unixepoch()) created_at INTEGER NOT NULL DEFAULT (unixepoch()),
isolated INTEGER NOT NULL DEFAULT 0,
notes TEXT,
system_prompt TEXT
); );
CREATE TABLE IF NOT EXISTS summaries ( CREATE TABLE IF NOT EXISTS summaries (
@@ -72,17 +88,17 @@ const schema = `
metadata TEXT metadata TEXT
); );
CREATE INDEX IF NOT EXISTS idx_episodes_session CREATE INDEX IF NOT EXISTS idx_episodes_session ON episodes(session_id);
ON episodes(session_id); CREATE INDEX IF NOT EXISTS idx_episodes_created ON episodes(created_at);
CREATE INDEX IF NOT EXISTS idx_episodes_created CREATE INDEX IF NOT EXISTS idx_entities_type ON entities(type);
ON episodes(created_at); CREATE INDEX IF NOT EXISTS idx_sessions_project ON sessions(project_id);
CREATE INDEX IF NOT EXISTS idx_entities_type CREATE INDEX IF NOT EXISTS idx_summaries_project ON summaries(project_id);
ON entities(type); CREATE INDEX IF NOT EXISTS idx_summaries_session ON summaries(session_id);
CREATE VIRTUAL TABLE IF NOT EXISTS episodes_fts CREATE VIRTUAL TABLE IF NOT EXISTS episodes_fts
USING fts5(user_message, ai_response, content=episodes, content_rowid=id); USING fts5(user_message, ai_response, content=episodes, content_rowid=id);
CREATE TRIGGER IF NOT EXISTS episodes_fts_insert CREATE TRIGGER IF NOT EXISTS episodes_fts_insert
AFTER INSERT ON episodes BEGIN AFTER INSERT ON episodes BEGIN
INSERT INTO episodes_fts(rowid, user_message, ai_response) INSERT INTO episodes_fts(rowid, user_message, ai_response)
VALUES (new.id, new.user_message, new.ai_response); VALUES (new.id, new.user_message, new.ai_response);
@@ -101,8 +117,6 @@ const schema = `
INSERT INTO episodes_fts(rowid, user_message, ai_response) INSERT INTO episodes_fts(rowid, user_message, ai_response)
VALUES (new.id, new.user_message, new.ai_response); VALUES (new.id, new.user_message, new.ai_response);
END; END;
`; `;
module.exports = schema; module.exports = schema;
@@ -1,9 +1,7 @@
const semantic = require('../semantic') const semantic = require('../semantic')
const { getEnv, SERVICES, formatEpisodeText, ENTITIES, logger } = require('@nexusai/shared'); const { getEnv, SERVICES, formatEpisodeText, ENTITIES, logger, utilityInference } = require('@nexusai/shared');
const { upsertEntity, upsertRelationship, linkEntityToEpisode } = require('./index'); const { upsertEntity, upsertRelationship, linkEntityToEpisode } = require('./index');
const EXTRACTION_URL = getEnv('EXTRACTION_URL', 'http://localhost:11434');
const EXTRACTION_MODEL = getEnv('EXTRACTION_MODEL', 'qwen2.5:3b'); // ChatML format — see buildExtractionPrompt
const EMBEDDING_SERVICE_URL = getEnv('EMBEDDING_SERVICE_URL', SERVICES.EMBEDDING_URL); const EMBEDDING_SERVICE_URL = getEnv('EMBEDDING_SERVICE_URL', SERVICES.EMBEDDING_URL);
const ENTITY_TYPES = ENTITIES.TYPES; const ENTITY_TYPES = ENTITIES.TYPES;
@@ -18,10 +16,19 @@ function isIgnoredName(name) {
return IGNORED_NAMES.has(normalized); return IGNORED_NAMES.has(normalized);
} }
// NOTE: This prompt uses ChatML format (<|im_start|> / <|im_end|> tags), which is // Guards against the extraction model regurgitating the "known entities" hint block
// specific to qwen-family models. If EXTRACTION_MODEL is changed to a Llama-family // as if those entities appeared in the conversation. The prompt constrains names to
// or other model, this format will need to change — most alternatives use either // short proper nouns, so a genuinely-discussed entity appears verbatim in the
// plain text or [INST] / <<SYS>> tags. Silent degradation is likely if mismatched. // exchange; a regurgitated one generally does not. Whitespace/case-normalized
// substring match.
function mentionedIn(name, haystack) {
const norm = s => s.toLowerCase().replace(/\s+/g, ' ').trim();
return norm(haystack).includes(norm(name));
}
// Returns { system, user } for utilityInference. The model's prompt template
// (ChatML for qwen, etc.) is applied by Ollama via the inference service's
// /utility/complete route — no template tags belong in this file.
function buildExtractionPrompt(userMessage, aiResponse, knownEntities = []) { function buildExtractionPrompt(userMessage, aiResponse, knownEntities = []) {
const knownBlock = knownEntities.length > 0 const knownBlock = knownEntities.length > 0
? [ ? [
@@ -31,32 +38,30 @@ function buildExtractionPrompt(userMessage, aiResponse, knownEntities = []) {
].join('\n') ].join('\n')
: ''; : '';
return [ return {
'<|im_start|>system', system: 'You are a named entity and relationship extractor. You output only valid JSON.',
'You are a named entity and relationship extractor. You output only valid JSON.', user: [
'<|im_end|>', 'Read the conversation below and extract all named entities and the relationships between them.',
'<|im_start|>user', `Entity types: ${ENTITY_TYPES.join(', ')}`,
'Read the conversation below and extract all named entities and the relationships between them.', 'Use "character" for any fictional, game, or media characters (e.g. characters from anime, games, books, TV shows, movies)',
`Entity types: ${ENTITY_TYPES.join(', ')}`, 'Use "person" only for real people',
'Use "character" for any fictional, game, or media characters (e.g. characters from anime, games, books, TV shows, movies)', 'For each entity provide:',
'Use "person" only for real people', ' "name": short proper noun only (max 4 words)',
'For each entity provide:', ' "type": one of the valid types',
' "name": short proper noun only (max 4 words)', ' "notes": one specific sentence about this entity based on the conversation',
' "type": one of the valid types', 'For relationships, use snake_case verb labels (e.g. works_on, manages, uses, knows, located_in, part_of, created_by).',
' "notes": one specific sentence about this entity based on the conversation', 'Only include relationships between entities you have listed above.',
'For relationships, use snake_case verb labels (e.g. works_on, manages, uses, knows, located_in, part_of, created_by).', 'The known-entities list below is ONLY for consistent spelling and types. Do NOT output an entity unless it actually appears in the conversation.',
'Only include relationships between entities you have listed above.', 'Return this exact JSON structure:',
'Return this exact JSON structure:', '{ "entities": [{"name": "...", "type": "...", "notes": "..."}], "relationships": [{"from": "...", "fromType": "...", "to": "...", "toType": "...", "label": "...", "notes": "..."}] }',
'{ "entities": [{"name": "...", "type": "...", "notes": "..."}], "relationships": [{"from": "...", "fromType": "...", "to": "...", "toType": "...", "label": "...", "notes": "..."}] }', '',
'', knownBlock,
knownBlock, '--- CONVERSATION ---',
'--- CONVERSATION ---', `User: ${userMessage}`,
`User: ${userMessage}`, `Assistant: ${aiResponse}`,
`Assistant: ${aiResponse}`, '--- END CONVERSATION ---',
'--- END CONVERSATION ---', ].join('\n'),
'<|im_end|>', };
'<|im_start|>assistant',
].join('\n');
} }
async function embedEntity(entity) { async function embedEntity(entity) {
@@ -80,30 +85,16 @@ async function extractAndStoreEntities(userMessage, aiResponse, episodeId=null,
// Fetch existing entities to guide the model toward consistent name/type pairs // Fetch existing entities to guide the model toward consistent name/type pairs
const db = require('../db').getDB(); const db = require('../db').getDB();
const knownEntities = db.prepare(`SELECT name, type FROM entities ORDER BY rowid DESC LIMIT 20`).all(); const knownEntities = db.prepare(`SELECT name, type FROM entities ORDER BY rowid DESC LIMIT 20`).all();
const prompt = buildExtractionPrompt(userMessage, aiResponse, knownEntities); const { system, user } = buildExtractionPrompt(userMessage, aiResponse, knownEntities);
const raw = await utilityInference({
const res = await fetch(`${EXTRACTION_URL}/api/generate`, { system,
method: 'POST', user,
headers: { 'Content-Type': 'application/json' }, json: true,
body: JSON.stringify({ temperature: ENTITIES.TEMPERATURE,
model: EXTRACTION_MODEL, maxTokens: ENTITIES.NUM_PREDICT,
prompt: prompt,
stream: false,
format: 'json',
options: {
temperature: ENTITIES.TEMPERATURE,
num_predict: ENTITIES.NUM_PREDICT,
},
}),
signal: AbortSignal.timeout(60_000),
}); });
if (!res.ok) throw new Error(`Ollama responded ${res.status}`);
const data = await res.json();
const raw = data.response?.trim() ?? '';
const jsonMatch = raw.match(/\{[\s\S]*\}/); const jsonMatch = raw.match(/\{[\s\S]*\}/);
if (!jsonMatch) { if (!jsonMatch) {
logger.warn('[entities] No JSON object found in response'); logger.warn('[entities] No JSON object found in response');
@@ -129,9 +120,17 @@ async function extractAndStoreEntities(userMessage, aiResponse, episodeId=null,
const entityMap = new Map(); const entityMap = new Map();
let saved = 0; let saved = 0;
// Everything the model was actually shown from THIS exchange. Extracted names
// are verified against this to reject regurgitated known-entity hints.
const conversationText = `${userMessage} ${aiResponse}`;
for (const { name, type, notes } of entities) { for (const { name, type, notes } of entities) {
if (!name || !type || !ENTITY_TYPES.includes(type)) continue; if (!name || !type || !ENTITY_TYPES.includes(type)) continue;
if (isIgnoredName(name)) continue; if (isIgnoredName(name)) continue;
if (!mentionedIn(name, conversationText)) {
logger.debug(`[entities] Skipping "${name}" — not present in the exchange (likely hint regurgitation)`);
continue;
}
const entity = upsertEntity(name, type, notes ?? null); const entity = upsertEntity(name, type, notes ?? null);
entityMap.set(`${name}::${type}`, entity); entityMap.set(`${name}::${type}`, entity);
@@ -178,4 +177,4 @@ async function extractAndStoreEntities(userMessage, aiResponse, episodeId=null,
} }
} }
module.exports = { extractAndStoreEntities }; module.exports = { extractAndStoreEntities, mentionedIn, isIgnoredName };
+42 -3
View File
@@ -1,5 +1,5 @@
const {getDB} = require('../db'); const {getDB} = require('../db');
const { EPISODIC, getEnv, SERVICES, parseRow, formatEpisodeText, SUMMARIES, logger } = require('@nexusai/shared'); const { EPISODIC, getEnv, SERVICES, parseRow, formatEpisodeText, SUMMARIES, logger, isTrivialTurn, CONSOLIDATION } = require('@nexusai/shared');
const semantic = require('../semantic'); const semantic = require('../semantic');
const { extractAndStoreEntities } = require('../entities/extraction') const { extractAndStoreEntities } = require('../entities/extraction')
@@ -162,8 +162,16 @@ async function createEpisode(sessionId, userMessage, aiResponse, tokenCount = nu
})) }))
.catch(err => logger.error(`Failed to embed episode ${episode.id}:`, err.message)); .catch(err => logger.error(`Failed to embed episode ${episode.id}:`, err.message));
extractAndStoreEntities(userMessage, aiResponse, episode.id, projectId) // Skip entity extraction on contentless social turns (greetings, sign-offs).
.catch(err => logger.error(`Failed to extract entities for episode ${episode.id}:`, err.message)); // They carry nothing worth storing, and running extraction on them was a
// source of junk/confabulated entities. The informational content lives in
// substantive turns, which still extract normally.
if (isTrivialTurn(userMessage)) {
logger.debug(`[entities] Skipping extraction for episode ${episode.id} — trivial turn`);
} else {
extractAndStoreEntities(userMessage, aiResponse, episode.id, projectId)
.catch(err => logger.error(`Failed to extract entities for episode ${episode.id}:`, err.message));
}
return episode; return episode;
@@ -258,6 +266,19 @@ function deleteEpisode(id) {
db.prepare(`DELETE FROM episodes WHERE id = ?`).run(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 ********/ /******** Embedding Helper ********/
async function getEpisodeEmbedding(userMessage, aiResponse){ async function getEpisodeEmbedding(userMessage, aiResponse){
const url = getEnv('EMBEDDING_SERVICE_URL', SERVICES.EMBEDDING_URL); const url = getEnv('EMBEDDING_SERVICE_URL', SERVICES.EMBEDDING_URL);
@@ -290,6 +311,22 @@ function getEpisodesByProject(projectId, limit = SUMMARIES.MAX_PROJECT_EPISODE_L
`).all(projectId, limit).map(parseRow); `).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 = { module.exports = {
createSession, createSession,
getSession, getSession,
@@ -307,6 +344,8 @@ module.exports = {
getEpisodesSince, getEpisodesSince,
searchEpisodes, searchEpisodes,
deleteEpisode, deleteEpisode,
touchEpisodes,
getEpisodesByProject, getEpisodesByProject,
buildFtsQuery, buildFtsQuery,
getConsolidationCandidates,
}; };
+14 -1
View File
@@ -74,4 +74,17 @@ function getEpisodeIdsByEntities(entityIds) {
).all(...entityIds).map(r => r.episode_id); ).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(); require ('dotenv').config();
const express = require('express'); 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 { getDB } = require('./db');
const { createProject, getProjects, getProject, updateProject, deleteProject } = require('./db/projects'); const { createProject, getProjects, getProject, updateProject, deleteProject } = require('./db/projects');
const { createSummary, getSummary, getSummariesBySession, getSummariesByProject, updateSummary, deleteSummary } = require('./db/summaries'); 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)); 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) => { app.get('/episodes/:id', (req, res) => {
const episode = episodic.getEpisode(req.params.id); const episode = episodic.getEpisode(req.params.id);
if (!episode) return res.status(404).json({ error: 'Episode not found' }); 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))); 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. // Episodes newer than :afterId, chronological — the un-summarized tail.
app.get('/sessions/:id/episodes/since/:afterId', (req, res) => { app.get('/sessions/:id/episodes/since/:afterId', (req, res) => {
const episodes = episodic.getEpisodesSince(Number(req.params.id), Number(req.params.afterId)); const episodes = episodic.getEpisodesSince(Number(req.params.id), Number(req.params.afterId));
res.json(episodes); 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) => { app.delete('/episodes/:id', (req, res) => {
const id = Number(req.params.id); const id = Number(req.params.id);
episodic.deleteEpisode(id); episodic.deleteEpisode(id);
@@ -1,4 +1,4 @@
const { SERVICES, getEnv, SUMMARIES } = require('@nexusai/shared'); const { SERVICES, getEnv, SUMMARIES, utilityInference } = require('@nexusai/shared');
const { const {
getSessionSummariesForProject, getSessionSummariesForProject,
getProjectOverviewSummary, getProjectOverviewSummary,
@@ -9,9 +9,6 @@ const {
const { getEpisodesByProject } = require('../episodic'); const { getEpisodesByProject } = require('../episodic');
const { getProject } = require('../db/projects'); 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 const MAX_SUMMARY_CHARS = SUMMARIES.MAX_SUMMARY_CHARS; // generous ceiling before we truncate input
function buildProjectSummaryPrompt(projectName, sessionSummaries) { function buildProjectSummaryPrompt(projectName, sessionSummaries) {
@@ -24,8 +21,9 @@ function buildProjectSummaryPrompt(projectName, sessionSummaries) {
summaryBlock = summaryBlock.slice(-MAX_SUMMARY_CHARS); 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 [ return [
'<|im_start|>user',
`The following are session summaries from a project called "${projectName}".`, `The following are session summaries from a project called "${projectName}".`,
'Write a project overview covering: goals, progress, key decisions, and current state.', '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.', '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.', 'Write in third person. Output only the overview text, no headings or labels.',
'', '',
summaryBlock, summaryBlock,
'<|im_end|>',
'<|im_start|>assistant',
].join('\n'); ].join('\n');
} }
@@ -49,8 +45,9 @@ function buildProjectSummaryFromEpisodesPrompt(projectName, episodes) {
episodeBlock = episodeBlock.slice(-MAX_SUMMARY_CHARS); 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 [ return [
'<|im_start|>user',
`The following are conversations from a project called "${projectName}".`, `The following are conversations from a project called "${projectName}".`,
'Write a project overview covering: goals, progress, key decisions, and current state.', '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.', '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.', 'Write in third person. Output only the overview text, no headings or labels.',
'', '',
episodeBlock, episodeBlock,
'<|im_end|>',
'<|im_start|>assistant',
].join('\n'); ].join('\n');
} }
async function generateProjectSummaryFromEpisodes(projectName, episodes) { async function generateProjectSummaryFromEpisodes(projectName, episodes) {
const prompt = buildProjectSummaryFromEpisodesPrompt(projectName, episodes); const user = buildProjectSummaryFromEpisodesPrompt(projectName, episodes);
return utilityInference({
const res = await fetch(`${EXTRACTION_URL}/api/generate`, { user,
method: 'POST', temperature: SUMMARIES.TEMPERATURE,
headers: { 'Content-Type': 'application/json' }, maxTokens: SUMMARIES.PROJECT_GEN_MAX_TOKENS,
body: JSON.stringify({
model: EXTRACTION_MODEL,
prompt,
stream: false,
options: { temperature: 0.2, num_predict: 1200 },
}),
}); });
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) { async function generateProjectSummary(projectName, sessionSummaries) {
const prompt = buildProjectSummaryPrompt(projectName, sessionSummaries); const user = buildProjectSummaryPrompt(projectName, sessionSummaries);
return utilityInference({
const res = await fetch(`${EXTRACTION_URL}/api/generate`, { user,
method: 'POST', temperature: SUMMARIES.TEMPERATURE,
headers: { 'Content-Type': 'application/json' }, maxTokens: SUMMARIES.PROJECT_GEN_MAX_TOKENS,
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 },
}),
}); });
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 // Main entry point — called by the route handler
@@ -142,4 +106,4 @@ async function generateAndStoreProjectSummary(projectId) {
} }
} }
module.exports = { generateAndStoreProjectSummary }; module.exports = { generateAndStoreProjectSummary };
@@ -2,7 +2,7 @@ const memory = require("../services/memory");
const inference = require("../services/inference"); const inference = require("../services/inference");
const embedding = require("../services/embedding"); const embedding = require("../services/embedding");
const qdrant = require("../services/qdrant"); const qdrant = require("../services/qdrant");
const { ORCHESTRATION, RETRIEVAL, logger } = require("@nexusai/shared"); const { ORCHESTRATION, RETRIEVAL, logger, isTrivialTurn } = require("@nexusai/shared");
const appSettings = require("../config/settings"); const appSettings = require("../config/settings");
const {triggerSummary} = require('../services/summarization') const {triggerSummary} = require('../services/summarization')
const graph = require('../services/graph'); const graph = require('../services/graph');
@@ -127,16 +127,22 @@ async function getSemanticEpisodes(
} }
} }
async function getRelevantEntities(vector, projectId = null) { async function getRelevantEntities(vector, { projectId = null, sessionId } = {}) {
if (!vector) return []; if (!vector) return [];
try { try {
const results = await qdrant.searchEntities(vector, { projectId }); let allowedIds;
logger.info( if (projectId === null || projectId === undefined) {
'[orchestration] Entity search results:', // Non-project chat is its own island — scope to entities linked to
results.map((r) => ({ name: r.payload?.name, score: r.score })), // THIS session. No links yet ⇒ nothing to retrieve, and we return
); // early so searchEntities is never called unfiltered.
// Include the Qdrant point ID (== SQLite entity ID) for graph traversal allowedIds = await memory.getEntityIdsBySession(sessionId);
return results.map((r) => r.payload ? { id: r.id, ...r.payload } : null).filter(Boolean); 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) { } catch (err) {
logger.debug('[orchestration] Entity search failed, continuing without:', err.message); logger.debug('[orchestration] Entity search failed, continuing without:', err.message);
return []; return [];
@@ -177,8 +183,11 @@ function fuseEpisodeResults(semanticEps, keywordEps, { semanticWeight, keywordWe
} }
function estimateTokens(episode) { function estimateTokens(episode) {
return episode.token_count //NOTE: episode.token_count is not used here. It stores the
?? Math.ceil((episode.user_message.length + episode.ai_response.length) / 4); //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 }) { function buildScoredPool(fusedWithScores, recentEpisodes, entityBoostedIds, { entityWeight }) {
@@ -252,6 +261,7 @@ async function getFusedEpisodes(userMessage, session, recentIds, projectSessionI
} }
async function assembleContext(externalId, userMessage) { async function assembleContext(externalId, userMessage) {
logger.info(`[orchestration] assembleContext ENTERED — msg: "${userMessage}", trivial: ${isTrivialTurn(userMessage)}`);
const settings = appSettings.load(); const settings = appSettings.load();
const { recentEpisodeLimit, semanticLimit, scoreThreshold, const { recentEpisodeLimit, semanticLimit, scoreThreshold,
temperature, repeatPenalty, topP, topK, systemPrompt, temperature, repeatPenalty, topP, topK, systemPrompt,
@@ -283,20 +293,27 @@ async function assembleContext(externalId, userMessage) {
const isFirstMessage = recentEpisodes.length === 0; const isFirstMessage = recentEpisodes.length === 0;
const recentIds = new Set(recentEpisodes.map(e => e.id)); const recentIds = new Set(recentEpisodes.map(e => e.id));
// 4. Embed the query once — the vector is shared by semantic episode search // 4. Retrieval — skipped entirely on contentless social turns (greetings,
// and entity search, so embedding it twice was a wasted round-trip + Ollama call. // sign-offs). On those, recent history alone is the right context; running
let queryVector = null; // semantic/keyword/entity retrieval only surfaces marginal noise the model
try { // then confabulates around. Embed once (shared by episode + entity search).
queryVector = await embedding.embed(userMessage); let fusedWithScores = [];
} catch (err) { let entityResults = [];
logger.warn('[orchestration] Query embedding failed; semantic + entity search disabled this turn:', err.message); if (!isTrivialTurn(userMessage)) {
} let queryVector = null;
try {
queryVector = await embedding.embed(userMessage);
} catch (err) {
logger.warn('[orchestration] Query embedding failed; semantic + entity search disabled this turn:', err.message);
}
// 4b. Fused retrieval + entity search in parallel (both are independent) [fusedWithScores, entityResults] = await Promise.all([
const [fusedWithScores, entityResults] = await Promise.all([ getFusedEpisodes(userMessage, session, recentIds, projectSessionIds, queryVector, { semanticLimit, scoreThreshold, semanticWeight, keywordWeight }),
getFusedEpisodes(userMessage, session, recentIds, projectSessionIds, queryVector, { semanticLimit, scoreThreshold, semanticWeight, keywordWeight }), getRelevantEntities(queryVector, { projectId: session.project_id ?? null, sessionId: session.id }),
getRelevantEntities(queryVector, session.project_id ?? null), ]);
]); } else {
logger.debug('[orchestration] Trivial turn — skipping semantic/keyword/entity retrieval');
}
// 5. Entity-linked episode IDs for scoring bonus // 5. Entity-linked episode IDs for scoring bonus
const entityIds = entityResults.map(e => e.id); const entityIds = entityResults.map(e => e.id);
@@ -314,6 +331,9 @@ async function assembleContext(externalId, userMessage) {
const scoredPool = buildScoredPool(fusedWithScores, recentEpisodes, entityBoostedIds, { entityWeight }); const scoredPool = buildScoredPool(fusedWithScores, recentEpisodes, entityBoostedIds, { entityWeight });
const { guaranteed, selected } = selectWithinBudget(scoredPool, contextBudget, minRecentEpisodes, recentEpisodes); const { guaranteed, selected } = selectWithinBudget(scoredPool, contextBudget, minRecentEpisodes, recentEpisodes);
const selectedIds = selected.map(ep => ep.id);
memory.touchEpisodes(selectedIds);
// 7. Graph neighborhood expansion // 7. Graph neighborhood expansion
let neighborhood = { nodes: [], edges: [] }; let neighborhood = { nodes: [], edges: [] };
if (entityIds.length > 0) { 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); const BASE_URL = getEnv('MEMORY_SERVICE_URL', SERVICES.MEMORY_URL);
@@ -216,6 +216,23 @@ async function getEpisodesByEntities(entityIds) {
return res.json(); // { episodeIds: [...] } 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 = { module.exports = {
getSessionByExternalId, getSessionByExternalId,
createSession, createSession,
@@ -242,4 +259,6 @@ module.exports = {
getProjectOverviewSummary, getProjectOverviewSummary,
searchEpisodes, searchEpisodes,
getEpisodesByEntities, getEpisodesByEntities,
getEntityIdsBySession,
touchEpisodes,
} }
@@ -29,30 +29,28 @@ async function searchEpisodes( vector, {limit = ORCHESTRATION.RECENT_EPISODE_LIM
return data.result; 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 }; const body = { vector, limit, score_threshold: scoreThreshold, with_payload: true };
if (projectId !== null && projectId !== undefined) { if (projectId !== null && projectId !== undefined) {
body.filter = { // Project chat: entities shared across the project's sessions.
must: [{ key: 'projectId', match: { value: projectId } }] 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( const res = await fetch(
`${BASE_URL}/collections/${COLLECTIONS.ENTITIES}/points/search`, `${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) { if (!res.ok) {
const body = await res.text(); const text = await res.text();
throw new Error(`Qdrant error: ${res.status} - ${body}`); throw new Error(`Qdrant error: ${res.status} - ${text}`);
} }
return (await res.json()).result;
const data = await res.json();
return data.result;
} }
module.exports = { searchEpisodes, searchEntities }; 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 MEMORY_URL = getEnv('MEMORY_SERVICE_URL', SERVICES.MEMORY_URL);
const THRESHOLD_TOKENS = parseInt(getEnv('SUMMARY_THRESHOLD_TOKENS', SUMMARIES.THRESHOLD_TOKENS)); 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: Conversation:
${context}`; ${context}`;
return [ // No ChatML wrapper — the model's own prompt template is applied server-side
'<|im_start|>user', // ChatML for qwen2.5 // by the inference service's /utility/complete route (Ollama /api/chat).
instruction, return instruction;
'<|im_end|>',
'<|im_start|>assistant',
].join('\n');
} }
async function generateSummary(episodes, existingSummary = null) { async function generateSummary(episodes, existingSummary = null) {
const prompt = buildSummaryPrompt(episodes, existingSummary); const user = buildSummaryPrompt(episodes, existingSummary);
const res = await fetch(`${EXTRACTION_URL}/api/generate`, { const content = await utilityInference({
method: 'POST', user,
headers: { 'Content-Type': 'application/json' }, temperature: SUMMARIES.TEMPERATURE,
body: JSON.stringify({ maxTokens: SUMMARIES.SESSION_GEN_MAX_TOKENS,
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
},
}),
}); });
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; return content;
} }
@@ -148,4 +125,4 @@ async function triggerSummary(session) {
); );
} }
module.exports = { triggerSummary, maybeSummarize }; module.exports = { triggerSummary, maybeSummarize };
+25 -2
View File
@@ -26,7 +26,7 @@ const ORCHESTRATION = {
SEMANTIC_LIMIT: 5, SEMANTIC_LIMIT: 5,
SCORE_THRESHOLD: 0.5, SCORE_THRESHOLD: 0.5,
ENTITIES_LIMIT: 5, ENTITIES_LIMIT: 5,
ENTITIES_THRESHOLD: 0.55, ENTITIES_THRESHOLD: 0.65,
TEMPERATURE: 0.7, TEMPERATURE: 0.7,
CONTEXT_BUDGET: 4096, CONTEXT_BUDGET: 4096,
ENTITY_WEIGHT: 0.5, ENTITY_WEIGHT: 0.5,
@@ -78,6 +78,13 @@ const SUMMARIES = {
MIN_EPISODES_SINCE: 5, // don't resummarize until N new episodes since last summary 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_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) 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 = { const ENTITIES = {
@@ -101,10 +108,24 @@ const ENTITIES = {
const RETRIEVAL = { const RETRIEVAL = {
RRF_K: 60, // Reciprocal Rank Fusion smoothing constant, softens rank-1 advantage, not exposed in settings RRF_K: 60, // Reciprocal Rank Fusion smoothing constant, softens rank-1 advantage, not exposed in settings
SEMANTIC_WEIGHT: 1.0, // Weight applied to semantic (QDrant) results SEMANTIC_WEIGHT: 1, // Weight applied to semantic (QDrant) results
KEYWORD_WEIGHT: 0.5, // Weight applied to keyword (SQLite) results, 0 = disables, set >0 to enable and tune balance between semantic vs keyword matches KEYWORD_WEIGHT: 0.5, // Weight applied to keyword (SQLite) results, 0 = disables, set >0 to enable and tune balance between semantic vs keyword matches
} }
const UTILITY = {
DEFAULT_URL: 'http://localhost:11434', // Ollama host for background/utility tasks
DEFAULT_MODEL: 'qwen2.5:3b',
TEMPERATURE: 0.2, // precise, low-creativity default for extraction/summaries
MAX_TOKENS: 1200,
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 = { module.exports = {
QDRANT, QDRANT,
COLLECTIONS, COLLECTIONS,
@@ -119,4 +140,6 @@ module.exports = {
SUMMARIES, SUMMARIES,
ENTITIES, ENTITIES,
RETRIEVAL, RETRIEVAL,
UTILITY,
CONSOLIDATION,
}; };
+24 -2
View File
@@ -1,7 +1,25 @@
const {getEnv} = require('./config/env'); const {getEnv} = require('./config/env');
const {QDRANT, COLLECTIONS, EPISODIC, SERVICES, OLLAMA, PORTS, LLAMACPP, INFERENCE_DEFAULTS, SQLITE, ORCHESTRATION, SUMMARIES, ENTITIES, RETRIEVAL } = require('./config/constants'); const {
const {parseRow, formatEpisodeText} = require('./utils') QDRANT,
COLLECTIONS,
EPISODIC,
SERVICES,
OLLAMA,
PORTS,
LLAMACPP,
INFERENCE_DEFAULTS,
SQLITE,
ORCHESTRATION,
SUMMARIES,
ENTITIES,
RETRIEVAL,
UTILITY,
CONSOLIDATION
} = require('./config/constants');
const {parseRow, formatEpisodeText, isTrivialTurn} = require('./utils')
const logger = require('./utils/logger'); const logger = require('./utils/logger');
const {utilityInference} = require('./utils/utilityInference');
module.exports = { module.exports = {
getEnv, getEnv,
@@ -17,8 +35,12 @@ module.exports = {
ORCHESTRATION, ORCHESTRATION,
parseRow, parseRow,
formatEpisodeText, formatEpisodeText,
isTrivialTurn,
SUMMARIES, SUMMARIES,
ENTITIES, ENTITIES,
logger, logger,
RETRIEVAL, RETRIEVAL,
UTILITY,
CONSOLIDATION,
utilityInference,
}; };
+43 -1
View File
@@ -10,4 +10,46 @@ function formatEpisodeText(userMessage, aiResponse) {
return `User: ${userMessage}\nAssistant: ${aiResponse}`; return `User: ${userMessage}\nAssistant: ${aiResponse}`;
} }
module.exports = { parseRow, formatEpisodeText }; // Contentless "social" turns — greetings, sign-offs, acknowledgements — that
// carry no information to store or recall. Used to skip entity extraction and
// noisy retrieval on such turns.
const TRIVIAL_TURNS = new Set([
'good morning', 'good night', 'good evening', 'good afternoon', 'morning', 'evening',
'hello', 'hi', 'hey', 'hey there', 'yo', 'sup', 'whats up', 'hiya', 'howdy',
'goodbye', 'bye', 'see you', 'see ya', 'see you later', 'talk later', 'talk soon',
'later', 'catch you later', 'gtg', 'gotta go', 'im off', 'heading out',
'thanks', 'thank you', 'thanks again', 'thank you so much', 'ty', 'thx', 'cheers',
'no worries', 'no problem', 'np', 'youre welcome', 'my pleasure',
'ok', 'okay', 'k', 'kk', 'alright', 'sure', 'sounds good', 'got it', 'gotcha',
'cool', 'nice', 'great', 'awesome', 'perfect', 'lol', 'haha',
'just saying hi', 'just saying hello', 'just checking in', 'just dropping in',
'stopping by', 'just stopping by', 'just wanted to say hi',
]);
// Trailing filler words that don't change a phrase's social nature, so
// "good morning again" / "thanks everyone" collapse to a base phrase.
const TRIVIAL_FILLER = new Set(['again', 'there', 'everyone', 'all', 'yall', 'folks', 'man', 'dude', 'friend', 'buddy']);
// Conservative, HIGH-PRECISION check: is this a contentless social turn?
// It exists to skip retrieval/extraction on turns with nothing to remember or
// recall, and deliberately errs toward "substantive" — a real short query like
// "capital of France" must NOT be treated as trivial. Genuine intent
// classification is a separate, later concern (confidence-based routing), not
// this heuristic. Extend TRIVIAL_TURNS as new pure-social phrases show up.
function isTrivialTurn(message) {
if (!message) return true;
const norm = String(message)
.toLowerCase()
.replace(/[^\p{L}\p{N}\s]/gu, ' ') // punctuation → space (unicode-aware)
.replace(/\s+/g, ' ')
.trim();
if (!norm) return true; // empty / punctuation-only
if (TRIVIAL_TURNS.has(norm)) return true;
// Drop trailing filler and re-check ("good morning again" → "good morning")
const words = norm.split(' ');
while (words.length > 1 && TRIVIAL_FILLER.has(words[words.length - 1])) words.pop();
return TRIVIAL_TURNS.has(words.join(' '));
}
module.exports = { parseRow, formatEpisodeText, isTrivialTurn };
@@ -0,0 +1,20 @@
const { getEnv } = require('../config/env');
const { SERVICES } = require('../config/constants');
// Client for the inference service's /utility/complete route.
// Resolved at call time (not module load) so each service's .env applies.
async function utilityInference({ system, user, json = false, temperature, maxTokens }) {
const base = getEnv('INFERENCE_SERVICE_URL', SERVICES.INFERENCE_URL);
const res = await fetch(`${base}/utility/complete`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ system, user, json, temperature, maxTokens }),
});
if (!res.ok) throw new Error(`Utility inference error: ${res.status}`);
const data = await res.json();
return data.text;
}
module.exports = { utilityInference };
+49
View File
@@ -0,0 +1,49 @@
// Entity extraction guards — tests the REAL mentionedIn and isIgnoredName filters
// that keep the extractor from storing (a) greetings and (b) "known entities" hint
// regurgitation. The mentionedIn cases are built from an actual observed failure:
// a "good morning" turn where qwen echoed the hint list back as fake extractions.
const { test } = require('node:test');
const assert = require('node:assert');
const { mentionedIn, isIgnoredName } = require('../packages/memory-service/src/entities/extraction');
test('mentionedIn accepts names that actually appear in the exchange', () => {
const text = 'User: I love how Sanderson writes magic systems. Assistant: Brandon Sanderson is known for hard magic.';
assert.ok(mentionedIn('Sanderson', text));
assert.ok(mentionedIn('Brandon Sanderson', text));
});
test('mentionedIn is case- and whitespace-insensitive', () => {
assert.ok(mentionedIn('digimon world next order', 'talking about Digimon World Next Order today'));
assert.ok(mentionedIn('Digimon World', 'the Digimon World game')); // collapsed whitespace
});
test('mentionedIn matches a name inside a possessive', () => {
assert.ok(mentionedIn('Sanderson', "Sanderson's latest book"));
});
test('regurgitated hint entities are rejected on a greeting turn (the bug)', () => {
// The exchange that triggered the bug — a bare greeting, no entities in it.
const exchange = 'good morning again Morning! Good timing — anything on your mind today?';
// The names qwen wrongly "extracted" from the hint block, verbatim from the logs.
const regurgitated = [
'Sanderson', 'Baldree', 'Riordan',
'Digimon World Next Order',
'Digimon Adventure: Last Evolution Kizuna',
'Digimon Adventure (2020 reboot)',
];
for (const name of regurgitated) {
assert.strictEqual(mentionedIn(name, exchange), false, `"${name}" should be rejected — not in the exchange`);
}
});
test('a real entity survives even when hint regurgitation is happening', () => {
const exchange = 'good morning — did you finish the Qdrant migration? Yes, the Qdrant collection is rebuilt.';
assert.ok(mentionedIn('Qdrant', exchange), 'genuinely-discussed entity is kept');
});
test('isIgnoredName catches greetings regardless of punctuation/case', () => {
assert.ok(isIgnoredName('Good morning'));
assert.ok(isIgnoredName('good morning!'));
assert.ok(isIgnoredName(' HELLO '));
assert.strictEqual(isIgnoredName('Qdrant'), false);
});
+74
View File
@@ -0,0 +1,74 @@
// Migration runner — tests the REAL migrate() version-stepping logic against a
// minimal fake db (emulating the better-sqlite3 pragma getter/setter + transaction
// API), so it runs with zero dependencies. Guards that migrations apply in order,
// advance user_version correctly, and no-op when the DB is already current.
const { test } = require('node:test');
const assert = require('node:assert');
const { migrate, LATEST_VERSION } = require('../packages/memory-service/src/db/migrations');
// Emulates just the slice of the better-sqlite3 API that migrate() uses.
function fakeDb(startVersion = 0) {
let version = startVersion;
return {
ran: [],
setCalls: [],
pragma(str, opts) {
const set = str.match(/^user_version\s*=\s*(\d+)$/);
if (set) { version = Number(set[1]); this.setCalls.push(version); return; }
if (str === 'user_version') return opts && opts.simple ? version : [{ user_version: version }];
throw new Error('unexpected pragma: ' + str);
},
transaction(fn) { return (...args) => fn(...args); }, // execute immediately
get version() { return version; },
};
}
test('a fresh DB (user_version 0) is stamped to LATEST_VERSION', () => {
const db = fakeDb(0);
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 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');
});
test('injected migrations run in order, each bumping the version by one', () => {
const db = fakeDb(0);
const order = [];
const migs = [
() => order.push('v1'),
() => order.push('v2'),
() => order.push('v3'),
];
const result = migrate(db, migs);
assert.strictEqual(result, 3);
assert.deepStrictEqual(order, ['v1', 'v2', 'v3'], 'migrations apply in sequence');
assert.deepStrictEqual(db.setCalls, [1, 2, 3], 'user_version advances one step at a time');
});
test('only migrations newer than the current version run', () => {
const db = fakeDb(1); // already at v1
const order = [];
const migs = [
() => order.push('v1'), // should be skipped
() => order.push('v2'),
() => order.push('v3'),
];
migrate(db, migs);
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);
});
+58
View File
@@ -0,0 +1,58 @@
// Schema completeness — execs the REAL schema.js into a throwaway SQLite DB (via
// built-in node:sqlite) and asserts a fresh database gets every column the old
// per-boot ALTER statements used to add. This is the guard for the migration
// consolidation: if someone drops a column from schema.js, a fresh install would
// silently lose it, and this test goes red.
const { test, before } = require('node:test');
const assert = require('node:assert');
const schema = require('../packages/memory-service/src/db/schema');
let DatabaseSync;
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: ['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'],
projects: ['id','name','description','colour','icon','created_at','isolated','notes','system_prompt'],
summaries: ['id','session_id','project_id','content','token_count','episode_range','created_at','metadata'],
};
const EXPECTED_OBJECTS = [
'idx_relationships_from','idx_relationships_to','idx_entity_episodes_entity','idx_entity_episodes_episode',
'idx_episodes_session','idx_episodes_created','idx_entities_type','idx_sessions_project',
'idx_summaries_project','idx_summaries_session',
'episodes_fts','episodes_fts_insert','episodes_fts_delete','episodes_fts_update',
];
let db;
before(() => {
if (!DatabaseSync) return;
db = new DatabaseSync(':memory:');
db.exec('PRAGMA foreign_keys = ON');
db.exec(schema); // throws here if schema.js has a syntax/ordering error
});
for (const [table, cols] of Object.entries(EXPECTED)) {
test(`${table} has exactly its full column set`, { skip: !DatabaseSync }, () => {
const got = db.prepare(`PRAGMA table_info(${table})`).all().map(r => r.name);
assert.deepStrictEqual([...got].sort(), [...cols].sort(), `${table} columns differ from the intended shape`);
});
}
test('all indexes, FTS table, and triggers are present', { skip: !DatabaseSync }, () => {
const names = db.prepare(`SELECT name FROM sqlite_master WHERE name LIKE 'idx_%' OR name LIKE 'episodes_fts%'`)
.all().map(r => r.name);
const missing = EXPECTED_OBJECTS.filter(o => !names.includes(o));
assert.deepStrictEqual(missing, [], `missing schema objects: ${missing}`);
});
test('FTS trigger populates the index on episode insert', { skip: !DatabaseSync }, () => {
db.prepare(`INSERT INTO sessions(external_id) VALUES ('s-test')`).run();
db.prepare(`INSERT INTO episodes(session_id, user_message, ai_response) VALUES (1,'find me qdrant','ok')`).run();
const hit = db.prepare(`SELECT rowid FROM episodes_fts WHERE episodes_fts MATCH 'qdrant'`).all();
assert.strictEqual(hit.length, 1);
});
+1 -1
View File
@@ -19,7 +19,7 @@ function mockMemory({ stats, summaries, since }) {
calls.sinceAfterId = Number(u.split('/since/').at(-1)); calls.sinceAfterId = Number(u.split('/since/').at(-1));
return json(since.filter(ep => ep.id > calls.sinceAfterId)); 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 (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 }); } if (/\/summaries\/\d+$/.test(u) && opts.method === 'PATCH') { calls.patched = JSON.parse(opts.body); return json({ ok: true }); }
throw new Error('unexpected fetch: ' + u); throw new Error('unexpected fetch: ' + u);
+42
View File
@@ -0,0 +1,42 @@
// Trivial-turn guard — tests the REAL isTrivialTurn from @nexusai/shared, which
// gates whether a turn skips retrieval + entity extraction. The critical property
// is high precision: greetings/pleasantries are caught, but real short queries
// must NOT be (a false positive would suppress retrieval on a genuine question).
const { test } = require('node:test');
const assert = require('node:assert');
const { isTrivialTurn } = require('@nexusai/shared');
test('plain greetings and pleasantries are trivial', () => {
for (const m of ['good morning', 'Good morning!', 'hello', 'hey there', 'hi',
'thanks', 'thank you so much', 'bye', 'ok', 'sounds good',
'just dropping in', 'just checking in', 'cheers']) {
assert.strictEqual(isTrivialTurn(m), true, `"${m}" should be trivial`);
}
});
test('trailing filler collapses to a base greeting', () => {
assert.strictEqual(isTrivialTurn('good morning again'), true);
assert.strictEqual(isTrivialTurn('thanks everyone'), true);
assert.strictEqual(isTrivialTurn('hello there'), true);
});
test('empty or punctuation-only input is trivial', () => {
assert.strictEqual(isTrivialTurn(''), true);
assert.strictEqual(isTrivialTurn(' '), true);
assert.strictEqual(isTrivialTurn('!!!'), true);
assert.strictEqual(isTrivialTurn(null), true);
});
test('real queries are NOT trivial — even short ones (the precision guarantee)', () => {
for (const m of ['what is the capital of France', 'France capital', 'capital of France?',
'how do I configure Qdrant', 'One Piece', 'help me debug this',
'what did we decide about the schema', 'morning routine ideas']) {
assert.strictEqual(isTrivialTurn(m), false, `"${m}" should be substantive`);
}
});
test('a greeting prefix does not make a substantive message trivial', () => {
// "good morning, can you help with X" carries a real request — must not be skipped
assert.strictEqual(isTrivialTurn('good morning, can you help me with the migration'), false);
assert.strictEqual(isTrivialTurn('hey, what is the One Piece'), false);
});