Compare commits

...
30 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
31 changed files with 787 additions and 302 deletions
+23
View File
@@ -205,6 +205,9 @@ Returns `503` if llama-server is unreachable.
| `scoreThreshold` | float | 0–1 | Minimum similarity score for Qdrant results |
| `semanticWeight` | float | 0–5 | RRF weight for Qdrant semantic results |
| `keywordWeight` | float | 0–5 | RRF weight for FTS5 keyword results (`0` = disabled) |
| `contextBudget` | integer | — | Token budget for context assembly (char/4 estimation) |
| `entityWeight` | float | — | Scoring bonus for entity-linked episodes in the context pool |
| `minRecentEpisodes` | integer | — | Guaranteed floor of recent episodes always included in context |
| `modelsFolderPath` | string | — | Path to folder containing .gguf files |
| `temperature` | float | 0–2 | Inference randomness |
| `repeatPenalty` | float | 1–2 | Repeat token penalty |
@@ -253,6 +256,7 @@ orchestration.
| GET | /sessions/by-external/:externalId | Get session by external ID |
| PATCH | /sessions/by-external/:externalId | Update session fields |
| DELETE | /sessions/by-external/:externalId | Delete session (cascades to episodes) |
| GET | /sessions/:id/entity-ids | Entity IDs linked to this session's episodes (non-project memory isolation scoping) |
> Route ordering: `by-external/:externalId` must be defined before `/:id`
> to prevent `by-external` being captured as an ID param.
@@ -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/:id | Get episode by ID |
| GET | /sessions/:id/episodes?limit=&offset= | Paginated episodes for a session |
| GET | /sessions/:id/episode-stats | Aggregate: count, total tokens, max id (summarization threshold check) |
| GET | /sessions/:id/episodes/since/:afterId | Episodes newer than :afterId, chronological (un-summarized tail) |
| POST | /episodes/touch | Batch access-tracking bump — increments `access_count`, sets `last_accessed_at` |
| GET | /sessions/:id/consolidation-candidates | Dry-run consolidation scoring (aging score, floors applied) |
| DELETE | /episodes/:id | Delete episode (SQLite + Qdrant cleanup) |
> Route ordering: `/episodes/search` must be defined before `/episodes/:id`.
@@ -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
| 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.
- [x] Reciprocal Rank Fusion (RRF) — merge semantic (Qdrant) + keyword (FTS5) results
- [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
Prevents long-term memory degradation and enables compression.
- [ ] Episode aging — score/weight episodes by recency and access frequency
- [ ] Consolidation pass — merge related low-weight episodes into summary nodes
- [x] Episode aging — `access_count` + `last_accessed_at` columns (v2 migration), batch touch on retrieval selection (`POST /episodes/touch`), aging score `access_count / (1 + days since last access)`
- [x] Dry-run candidates endpoint — `GET /sessions/:id/consolidation-candidates` with age + session-size floors (observe-only phase before destructive pass)
- [ ] Consolidation pass — merge related low-weight episodes into summary nodes (incl. Qdrant vector + entity link cleanup)
- [ ] Orphan cleanup — remove entities no longer referenced by active episodes
### 4. User Preference Model
@@ -88,11 +91,11 @@ Short-circuit simple requests before they reach the LLM.
- [ ] Confidence bands — FAST PATH (memory lookup only) vs FULL (LLM + context)
- [ ] Fast-path handlers — direct memory queries, session lookups, factual recalls
### 6. Smarter Context Assembly *(inspired by acid2lake)*
### 6. Smarter Context Assembly *(inspired by acid2lake)* ✅
Budget-aware context selection instead of dumping all relevant memory into the prompt.
- [ ] Token budget manager in orchestration
- [ ] Priority scoring — recency × relevance × entity weight
- [ ] Configurable context budget via env var
- [x] Token budget manager in orchestration (`selectWithinBudget`, char/4 estimation on stored text)
- [x] Priority scoring — RRF fusion + recency + entity boost (`buildScoredPool`)
- [x] Configurable via settings (`contextBudget`, `entityWeight`, `minRecentEpisodes`) — live, no restart
### 7. Procedural Memory Store *(inspired by acid2lake)*
Learns "how NexusAI has successfully handled this type of request before."
@@ -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`
**Triggered by:** Episode creation (`POST /episodes`)
**Model:** `qwen2.5:3b` via Ollama (configurable via `EXTRACTION_MODEL` env var)
**Model:** the utility model served by the inference service (`/utility/complete`), configurable via `UTILITY_MODEL` on the inference service
## Purpose
@@ -28,7 +28,7 @@ swallowed.
| Setting | Value | Notes |
|---|---|---|
| Model | `qwen2.5:3b` | Ollama, configurable via `EXTRACTION_MODEL` |
| Model | utility model | Served by inference-service `/utility/complete`, set via `UTILITY_MODEL` |
| Temperature | 0.1 | Low for consistent, deterministic output |
| `num_predict` | 1500 | Higher ceiling to accommodate entity + relationship JSON |
| `format` | `'json'` | Ollama constrained decoding — enforces valid JSON output |
@@ -131,6 +131,20 @@ check alone won't reject.
> imperfectly — switch to `/[^\p{L}\p{N}\s]/gu` if the set gains non-ASCII
> 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
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 |
| QDRANT_URL | No | http://localhost:6333 | Qdrant instance URL |
| EMBEDDING_SERVICE_URL | No | http://localhost:3003 | Embedding service URL |
| EXTRACTION_URL | No | http://localhost:11434 | Ollama URL for entity extraction |
| EXTRACTION_MODEL | No | qwen2.5:3b | Ollama model used for entity extraction |
| INFERENCE_SERVICE_URL | No | http://localhost:3001 | Inference service URL — entity extraction routes through its `/utility/complete` endpoint |
## Internal Structure
```
src/
├── db/
│ ├── index.js # SQLite connection + initialization + migrations
│ ├── schema.js # Table definitions, indexes, FTS5, triggers
│ ├── index.js # SQLite connection + init + migrate() + one-time FTS backfill
│ ├── migrations.js # Forward-only versioned migration runner (PRAGMA user_version)
│ ├── schema.js # Complete current shape: tables, indexes, FTS5, triggers
│ ├── projects.js # Project CRUD functions
│ └── summaries.js # Summary CRUD functions
├── episodic/
@@ -64,36 +64,61 @@ Eight core tables:
- **summaries** — condensed episode groups for efficient context retrieval
- **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
idempotent migrations in `db/index.js` at startup:
`schema.js` holds the **complete current shape** — every table, column, index,
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
try { db.exec(`ALTER TABLE sessions ADD COLUMN name TEXT`); } catch {}
try { db.exec(`ALTER TABLE sessions ADD COLUMN project_id INTEGER REFERENCES projects(id)`); } catch {}
try { db.exec(`CREATE INDEX IF NOT EXISTS idx_sessions_project ON sessions(project_id)`); } catch {}
try { db.exec(`ALTER TABLE projects ADD COLUMN isolated INTEGER NOT NULL DEFAULT 0`); } catch {}
try { db.exec(`ALTER TABLE projects ADD COLUMN notes TEXT`); } catch {}
try { db.exec(`ALTER TABLE projects ADD COLUMN system_prompt TEXT`); } catch {}
// 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 {}
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 {}
const migrations = [
(_db) => {}, // v0 → v1: consolidated baseline (historical ALTERs folded into schema.js)
(db) => { // v1 → v2: access tracking for consolidation lifecycle
db.exec(`ALTER TABLE episodes ADD COLUMN last_accessed_at INTEGER`);
db.exec(`ALTER TABLE episodes ADD COLUMN access_count INTEGER NOT NULL DEFAULT 0`);
db.exec(`UPDATE episodes SET last_accessed_at = created_at`); // backfill
},
];
const LATEST_VERSION = migrations.length; // derived, never hand-maintained
```
`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
An `episodes_fts` virtual table enables keyword search across all episodes.
Three triggers (`episodes_fts_insert`, `episodes_fts_update`, `episodes_fts_delete`)
keep the FTS index automatically in sync with the episodes table.
An `episodes_fts` external-content virtual table enables keyword search across
episodes. Three triggers (`episodes_fts_insert`, `episodes_fts_update`,
`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
@@ -101,6 +126,12 @@ keep the FTS index automatically in sync with the episodes table.
- `foreign_keys = ON` — enforces referential integrity and cascade deletes
- PRAGMAs set via `db.pragma()`, not `db.exec()`
> **Copying a live WAL database:** `cp` on the `.db` file alone silently loses
> everything in the un-checkpointed `-wal` file (recent writes, even the
> migration version stamp). Always use
> `sqlite3 nexusai.db "VACUUM INTO './copy.db'"` (or `.backup`) — safe while
> the service is running, produces a complete single-file snapshot.
### Dynamic Updates
Both `updateSession` and `updateProject` build their `SET` clause dynamically
@@ -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,
> and ChatML token stripping, see `summarization.md`.
## Access Tracking & Consolidation (dry-run)
Every episode selected into a chat context window (budget-selected, not the
guaranteed-recency floor) gets an access bump via `POST /episodes/touch` —
`access_count` incremented, `last_accessed_at` set to `Date.now()` (ms).
Called fire-and-forget from orchestration; a failure loses one increment,
nothing more.
`GET /sessions/:id/consolidation-candidates` scores episodes by
`access_count / (1 + days since last access)` — never-accessed episodes fall
back to `created_at` for the recency term and score exactly 0 (most eligible).
Two floors apply: episodes younger than `CONSOLIDATION.MIN_AGE_DAYS` are
excluded in SQL; sessions under `CONSOLIDATION.MIN_SESSION_EPISODES` return
`eligible: false` before scoring runs. The endpoint is observe-only — the
destructive pass (merge → summarize → Qdrant cleanup → orphan sweep) is not
yet built.
> **Unit note:** `created_at` is unix **seconds** (`unixepoch()`);
> `last_accessed_at` is unix **milliseconds** (`Date.now()`). The scoring
> query normalizes with `created_at * 1000`. Keep this in mind for any new
> queries touching both columns.
## Delete Behaviour (SQLite + Qdrant consistency)
SQLite cascades handle relational cleanup, but Qdrant is a separate store and
+31 -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 |
| QDRANT_URL | No | http://localhost:6333 | Qdrant URL for semantic search |
| CORS_ORIGIN | No | http://localhost:5173 | Allowed origin for CORS requests |
| EXTRACTION_URL | No | http://localhost:11434 | Ollama URL for summarisation |
| EXTRACTION_MODEL | No | qwen2.5:3b | Ollama model used for summarisation |
## Internal Structure
@@ -75,6 +73,9 @@ via `appSettings.load()` — changes apply immediately without a service restart
| `scoreThreshold` | 0.5 | Minimum similarity score for Qdrant semantic results |
| `semanticWeight` | 1.0 | RRF weight for Qdrant semantic results |
| `keywordWeight` | 0 | RRF weight for FTS5 keyword results (`0` = disabled) |
| `contextBudget` | — | Token budget for context assembly (char/4 estimation on stored text) |
| `entityWeight` | — | Scoring bonus for entity-linked episodes in the context pool |
| `minRecentEpisodes` | — | Guaranteed floor of recent episodes always included in context |
| `modelsFolderPath` | `/mnt/nexus-models` | Path to folder containing .gguf files |
| `temperature` | 0.7 | Inference temperature |
| `repeatPenalty` | 1.1 | Repeat token penalty |
@@ -103,35 +104,43 @@ difference is how the inference response is delivered to the client.
4. **Recent episode retrieval** — fetch most recent episodes (`recentEpisodeLimit`).
5. **Fused episode retrieval** — runs semantic (Qdrant) and keyword (FTS5)
5. **Trivial-turn gate** — greetings/pleasantries (`isTrivialTurn`) skip all
retrieval (semantic, keyword, entity); recent history alone is the context.
Breaks the greeting → marginal-retrieval → confabulation loop.
6. **Fused episode retrieval** — runs semantic (Qdrant) and keyword (FTS5)
search in parallel, then merges results via Reciprocal Rank Fusion (RRF).
Both paths are filtered against `recentIds` before fusion. FTS is scoped
to the current session or all project sessions. If `keywordWeight` is `0`,
the FTS call is skipped entirely. Non-critical — failures fall back to
whichever strategy succeeded.
The query is embedded once and shared with entity search. Both paths are
filtered against `recentIds` before fusion. FTS is scoped to the current
session or all project sessions. If `keywordWeight` is `0`, the FTS call
is skipped entirely. Non-critical — failures fall back to whichever
strategy succeeded.
6. **Entity search** — query `entities` Qdrant collection filtered by
`projectId`. Returns entity IDs alongside Qdrant payload data (the Qdrant
point ID equals the SQLite entity ID). Non-critical.
7. **Entity search + graph expansion** — query `entities` Qdrant collection
(project-scoped, or session-scoped via `/sessions/:id/entity-ids` for
non-project chats). Entity IDs are expanded into a 1-hop subgraph via
`POST /graph/neighbors`; on failure, falls back to flat entity list.
Non-critical.
7. **Graph neighborhood expansion** — call `POST /graph/neighbors` on
memory-service with the entity IDs from step 6. Returns a 1-hop subgraph
`{ nodes, edges }` — entity objects plus the relationships connecting them.
If no entities were found or the graph call fails, falls back to flat entity
list (no edges). Non-critical.
8. **Scored pool + budget selection** — `buildScoredPool` combines RRF scores,
recency, and entity-linkage bonus; `selectWithinBudget` fills `contextBudget`
(char/4 token estimation on stored text) above a guaranteed floor of
`minRecentEpisodes` recent episodes. Selected episode IDs are then reported
to `POST /episodes/touch` fire-and-forget (access tracking for the
consolidation lifecycle).
8. **Prompt assembly** — combine system prompt, graph context, fused episodes,
recent episodes, and user message.
9. **Prompt assembly** — combine system prompt, graph context, selected
episodes, guaranteed recent episodes, and user message.
9. **Inference** — send to inference service. `/chat` awaits full response;
`/chat/stream` pipes SSE chunks to the client.
10. **Inference** — send to inference service. `/chat` awaits full response;
`/chat/stream` pipes SSE chunks to the client.
10. **Episode write** — write exchange back to memory with `projectId`.
11. **Episode write** — write exchange back to memory with `projectId`.
11. **Summarisation trigger** — `triggerSummary(session, allEpisodes)` called
12. **Summarisation trigger** — `triggerSummary(session)` called
fire-and-forget. See `summarization.md` for full details.
12. **Auto-naming** — on first message with no session name, fires a secondary
13. **Auto-naming** — on first message with no session name, fires a secondary
inference call (max 20 tokens, temperature 0.3) to generate a session name.
### Prompt Structure
+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.
```
getFusedEpisodes()
├── getSemanticEpisodes() — Qdrant embed+search → fetch full rows by ID
│ (existing path, unchanged)
getFusedEpisodes(…, queryVector, …)
├── getSemanticEpisodes(queryVector) — Qdrant search → fetch full rows by ID
│ (query embedded ONCE upstream in assembleContext and shared with entity
│ search — no longer embedded separately here)
└── getFTSResults() — memory-service /episodes/search → full rows directly
(skipped entirely if keywordWeight == 0)
↓
@@ -48,6 +49,10 @@ fuseEpisodeResults() — pure RRF, no I/O
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
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
`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
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
```
#### `CONSOLIDATION`
Controls the memory consolidation lifecycle (currently dry-run only).
| Key | Value | Description |
|---|---|---|
| `MIN_AGE_DAYS` | `7` | Episodes younger than this are never consolidation candidates |
| `MIN_SESSION_EPISODES` | `20` | Sessions with fewer episodes are skipped entirely |
| `CANDIDATE_LIMIT` | `50` | Max candidates returned per scoring query |
#### `SQLITE`
| Key | Value | Description |
+46 -43
View File
@@ -6,13 +6,27 @@ the full context window with raw episodes.
**Location:** `packages/orchestration-service/src/services/summarization.js`
**Triggered by:** `chat/index.js` after every episode write (fire-and-forget)
**Model:** `qwen2.5:3b` via Ollama on Mini PC 1 (192.168.0.81)
**Model:** the utility model served by the inference service (`/utility/complete`), set via `UTILITY_MODEL` (backed by Ollama on Mini PC 1, 192.168.0.81)
---
## 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:
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
{
model: EXTRACTION_MODEL, // qwen2.5:3b (set via EXTRACTION_MODEL env var)
prompt: buildSummaryPrompt(episodesToSummarize, existingSummary),
stream: false,
// No format: 'json' — free-text output required for summaries
options: {
temperature: 0.2,
num_predict: 500,
},
}
const content = await utilityInference({
user: buildSummaryPrompt(episodesToSummarize, existingSummary),
temperature: SUMMARIES.TEMPERATURE, // 0.2
maxTokens: SUMMARIES.SESSION_GEN_MAX_TOKENS, // 500
});
```
`temperature: 0.2` is slightly higher than extraction (0.1) — summaries
benefit from some fluency. `num_predict: 500` gives room for 5 thorough
sentences without risk of runoff.
`TEMPERATURE` (0.2) is slightly higher than extraction (0.1) — summaries benefit
from some fluency. `SESSION_GEN_MAX_TOKENS` (500) gives room for ~5 thorough
sentences without runoff. Both live in `@nexusai/shared` `SUMMARIES` constants.
There is no `json: true` here — summaries are free-text, unlike entity extraction.
---
## Prompt Format
ChatML format — native to qwen2.5:
The prompt is plain text describing the task; the model's own prompt template
(ChatML for qwen, etc.) is applied **server-side** by the inference service via
Ollama's `/api/chat`. No `<|im_start|>` tags belong in this codebase, and the
utility model can be swapped (via `UTILITY_MODEL` on the inference service) with
no prompt changes here.
Fresh summary instruction:
```
<|im_start|>user
Summarize the conversation below in 3-5 sentences.
Write in third person. Do not quote directly — paraphrase only.
Do not include greetings, sign-offs, or filler. Output only the summary text.
Conversation:
{context}
<|im_end|>
<|im_start|>assistant
```
For cumulative updates, the instruction and context change:
Cumulative update instruction:
```
<|im_start|>user
Update the summary below to incorporate the new exchanges.
Write 3-5 sentences in third person. Do not quote directly — paraphrase only.
Do not include greetings, sign-offs, or filler. Output only the updated summary text.
@@ -92,35 +109,22 @@ Previous summary:
New exchanges:
{context}
<|im_end|>
<|im_start|>assistant
```
### Input truncation
Episode context is truncated to `MAX_CHARS = 3000` characters, keeping the
most recent exchanges (sliced from the end). This keeps Qwen focused and
most recent exchanges (sliced from the end). This keeps the model focused and
prevents the prompt from exceeding its effective context window.
---
## ChatML Token Stripping
## Output Handling
Qwen occasionally echoes ChatML tokens back into its response. The raw output
is cleaned before saving:
```js
const raw = data.response?.trim() ?? '';
const content = raw
.replace(/<\|im_start\|>.*?<\|im_end\|>/gs, '')
.replace(/<\|im_start\|>|<\|im_end\|>|<\|im_sep\|>/g, '')
.trim();
return content;
```
Without this, leaked tokens get stored in the summary and then injected
back into the next summarisation prompt — causing the model to append a new
summary after the old one rather than replacing it.
Because `/api/chat` applies and removes the prompt template server-side, the
returned text is already clean — the previous ChatML token-stripping step (and
the class of bug where leaked tokens got stored and re-injected into the next
summarisation prompt) no longer applies.
---
@@ -193,8 +197,7 @@ Set in `packages/orchestration-service/src/.env`:
| Variable | Default | Description |
|---|---|---|
| `EXTRACTION_URL` | `http://localhost:11434` | Ollama instance URL |
| `EXTRACTION_MODEL` | `qwen2.5:3b` | Model for summarisation |
| `INFERENCE_SERVICE_URL` | `http://localhost:3001` | Inference service — summaries route through its `/utility/complete` endpoint (model set via `UTILITY_MODEL` there) |
| `MEMORY_SERVICE_URL` | `http://localhost:3002` | Memory service URL |
| `SUMMARY_THRESHOLD_TOKENS` | `200` | Token threshold before summarisation triggers |
| `SUMMARY_MAX_TOKENS` | `800` | Max summary length before a new row is created |
+2
View File
@@ -2,6 +2,7 @@ require ('dotenv').config();
const express = require('express');
const {getEnv, PORTS, OLLAMA, logger} = require('@nexusai/shared');
const inferenceRouter = require('./routes/inference');
const utilityRouter = require('./routes/utility')
const app = express();
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);
// 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};
@@ -16,6 +16,14 @@ const migrations = [
// used to run (wrapped in try/catch) on every boot are folded into schema.js.
// Nothing to do here — this entry exists to mark v1 as the consolidated baseline.
(_db) => {},
(db) => {
db.exec(`
ALTER TABLE episodes ADD COLUMN last_accessed_at INTEGER;
ALTER TABLE episodes ADD COLUMN access_count INTEGER NOT NULL DEFAULT 0;
`)
db.exec(`UPDATE episodes SET last_accessed_at = created_at WHERE last_accessed_at IS NULL`)
}
];
const LATEST_VERSION = migrations.length;
+9 -7
View File
@@ -15,13 +15,15 @@ const schema = `
);
CREATE TABLE IF NOT EXISTS episodes (
id INTEGER PRIMARY KEY AUTOINCREMENT,
session_id INTEGER NOT NULL REFERENCES sessions(id) ON DELETE CASCADE,
user_message TEXT NOT NULL,
ai_response TEXT NOT NULL,
created_at INTEGER NOT NULL DEFAULT (unixepoch()),
token_count INTEGER,
metadata TEXT
id INTEGER PRIMARY KEY AUTOINCREMENT,
session_id INTEGER NOT NULL REFERENCES sessions(id) ON DELETE CASCADE,
user_message TEXT NOT NULL,
ai_response TEXT NOT NULL,
created_at INTEGER NOT NULL DEFAULT (unixepoch()),
token_count INTEGER,
metadata TEXT,
last_accessed_at INTEGER,
access_count INTEGER NOT NULL DEFAULT 0
);
CREATE TABLE IF NOT EXISTS entities (
@@ -1,9 +1,7 @@
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 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 ENTITY_TYPES = ENTITIES.TYPES;
@@ -28,10 +26,9 @@ function mentionedIn(name, haystack) {
return norm(haystack).includes(norm(name));
}
// NOTE: This prompt uses ChatML format (<|im_start|> / <|im_end|> tags), which is
// specific to qwen-family models. If EXTRACTION_MODEL is changed to a Llama-family
// or other model, this format will need to change — most alternatives use either
// plain text or [INST] / <<SYS>> tags. Silent degradation is likely if mismatched.
// 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 = []) {
const knownBlock = knownEntities.length > 0
? [
@@ -41,33 +38,30 @@ function buildExtractionPrompt(userMessage, aiResponse, knownEntities = []) {
].join('\n')
: '';
return [
'<|im_start|>system',
'You are a named entity and relationship extractor. You output only valid JSON.',
'<|im_end|>',
'<|im_start|>user',
'Read the conversation below and extract all named entities and the relationships between them.',
`Entity types: ${ENTITY_TYPES.join(', ')}`,
'Use "character" for any fictional, game, or media characters (e.g. characters from anime, games, books, TV shows, movies)',
'Use "person" only for real people',
'For each entity provide:',
' "name": short proper noun only (max 4 words)',
' "type": one of the valid types',
' "notes": one specific sentence about this entity based on the conversation',
'For relationships, use snake_case verb labels (e.g. works_on, manages, uses, knows, located_in, part_of, created_by).',
'Only include relationships between entities you have listed above.',
'The known-entities list below is ONLY for consistent spelling and types. Do NOT output an entity unless it actually appears in the conversation.',
'Return this exact JSON structure:',
'{ "entities": [{"name": "...", "type": "...", "notes": "..."}], "relationships": [{"from": "...", "fromType": "...", "to": "...", "toType": "...", "label": "...", "notes": "..."}] }',
'',
knownBlock,
'--- CONVERSATION ---',
`User: ${userMessage}`,
`Assistant: ${aiResponse}`,
'--- END CONVERSATION ---',
'<|im_end|>',
'<|im_start|>assistant',
].join('\n');
return {
system: 'You are a named entity and relationship extractor. You output only valid JSON.',
user: [
'Read the conversation below and extract all named entities and the relationships between them.',
`Entity types: ${ENTITY_TYPES.join(', ')}`,
'Use "character" for any fictional, game, or media characters (e.g. characters from anime, games, books, TV shows, movies)',
'Use "person" only for real people',
'For each entity provide:',
' "name": short proper noun only (max 4 words)',
' "type": one of the valid types',
' "notes": one specific sentence about this entity based on the conversation',
'For relationships, use snake_case verb labels (e.g. works_on, manages, uses, knows, located_in, part_of, created_by).',
'Only include relationships between entities you have listed above.',
'The known-entities list below is ONLY for consistent spelling and types. Do NOT output an entity unless it actually appears in the conversation.',
'Return this exact JSON structure:',
'{ "entities": [{"name": "...", "type": "...", "notes": "..."}], "relationships": [{"from": "...", "fromType": "...", "to": "...", "toType": "...", "label": "...", "notes": "..."}] }',
'',
knownBlock,
'--- CONVERSATION ---',
`User: ${userMessage}`,
`Assistant: ${aiResponse}`,
'--- END CONVERSATION ---',
].join('\n'),
};
}
async function embedEntity(entity) {
@@ -91,30 +85,16 @@ async function extractAndStoreEntities(userMessage, aiResponse, episodeId=null,
// Fetch existing entities to guide the model toward consistent name/type pairs
const db = require('../db').getDB();
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 res = await fetch(`${EXTRACTION_URL}/api/generate`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({
model: EXTRACTION_MODEL,
prompt: prompt,
stream: false,
format: 'json',
options: {
temperature: ENTITIES.TEMPERATURE,
num_predict: ENTITIES.NUM_PREDICT,
},
}),
signal: AbortSignal.timeout(60_000),
const raw = await utilityInference({
system,
user,
json: true,
temperature: ENTITIES.TEMPERATURE,
maxTokens: ENTITIES.NUM_PREDICT,
});
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]*\}/);
if (!jsonMatch) {
logger.warn('[entities] No JSON object found in response');
+42 -3
View File
@@ -1,5 +1,5 @@
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 { 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));
extractAndStoreEntities(userMessage, aiResponse, episode.id, projectId)
.catch(err => logger.error(`Failed to extract entities for episode ${episode.id}:`, err.message));
// Skip entity extraction on contentless social turns (greetings, sign-offs).
// 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;
@@ -258,6 +266,19 @@ function deleteEpisode(id) {
db.prepare(`DELETE FROM episodes WHERE id = ?`).run(id);
}
function touchEpisodes(ids) {
if (!ids.length) return;
const db = getDB(); // <-- missing in your version
const placeholders = ids.map(() => '?').join(',');
db.prepare(`
UPDATE episodes
SET access_count = access_count + 1,
last_accessed_at = ?
WHERE id IN (${placeholders})
`).run(Date.now(), ...ids);
}
/******** Embedding Helper ********/
async function getEpisodeEmbedding(userMessage, aiResponse){
const url = getEnv('EMBEDDING_SERVICE_URL', SERVICES.EMBEDDING_URL);
@@ -290,6 +311,22 @@ function getEpisodesByProject(projectId, limit = SUMMARIES.MAX_PROJECT_EPISODE_L
`).all(projectId, limit).map(parseRow);
}
function getConsolidationCandidates(sessionId, limit = CONSOLIDATION.CANDIDATE_LIMIT) {
const db = getDB();
const now = Date.now();
const cutoff = Math.floor(now / 1000) - CONSOLIDATION.MIN_AGE_DAYS *86400 // seconds, so it matches created_at
return db.prepare(`
SELECT id, access_count, created_at, last_accessed_at,
SUBSTR(user_message, 1, 80) AS preview,
CAST( access_count AS REAL) / (1+(?-COALESCE(last_accessed_at, created_at * 1000)) / 86400000.0) AS aging_score
FROM episodes
WHERE session_id = ? AND created_at < ?
ORDER BY aging_score ASC
LIMIT ?
`).all(now, sessionId, cutoff, limit);
}
module.exports = {
createSession,
getSession,
@@ -307,6 +344,8 @@ module.exports = {
getEpisodesSince,
searchEpisodes,
deleteEpisode,
touchEpisodes,
getEpisodesByProject,
buildFtsQuery,
getConsolidationCandidates,
};
+14 -1
View File
@@ -74,4 +74,17 @@ function getEpisodeIdsByEntities(entityIds) {
).all(...entityIds).map(r => r.episode_id);
}
module.exports = { getNeighborhood, getEntityNeighbors, getEpisodeIdsByEntities };
//Entity IDs linked (via entity_episodes) to any episode in a given session
//Scopes non-project entity search to the session's own entities under the "isolated chats" model.
//Entities globally deduped
function getEntityIdsBySession(sessionId){
const db = getDB();
return db.prepare(`
SELECT DISTINCT ee.entity_id
FROM entity_episodes ee
JOIN episodes e on e.id = ee.episode_id
WHERE e.session_id = ?
`).all(sessionId).map(r => r.entity_id);
}
module.exports = { getNeighborhood, getEntityNeighbors, getEpisodeIdsByEntities, getEntityIdsBySession };
+35 -1
View File
@@ -1,6 +1,6 @@
require ('dotenv').config();
const express = require('express');
const {getEnv, PORTS, EPISODIC, logger} = require('@nexusai/shared');
const {getEnv, PORTS, EPISODIC, logger, CONSOLIDATION} = require('@nexusai/shared');
const { getDB } = require('./db');
const { createProject, getProjects, getProject, updateProject, deleteProject } = require('./db/projects');
const { createSummary, getSummary, getSummariesBySession, getSummariesByProject, updateSummary, deleteSummary } = require('./db/summaries');
@@ -139,6 +139,13 @@ app.get('/episodes/search', (req, res) => {
res.json(episodic.searchEpisodes(q, Number(limit), parsedSessionIds));
});
app.post('/episodes/touch', (req, res) => {
const { ids } = req.body;
if(!Array.isArray(ids)) return res.status(400).json({error: 'ids must be an array'});
episodic.touchEpisodes(ids);
res.json({touched: ids.length});
})
app.get('/episodes/:id', (req, res) => {
const episode = episodic.getEpisode(req.params.id);
if (!episode) return res.status(404).json({ error: 'Episode not found' });
@@ -162,12 +169,39 @@ app.get('/sessions/:id/episode-stats', (req, res) => {
res.json(episodic.getSessionEpisodeStats(Number(req.params.id)));
});
//Entity IDs linked to this session's episodes: sesion-scoped entity search
app.get('/sessions/:id/entity-ids', (req, res) => {
res.json({
entityIds: graph.getEntityIdsBySession(Number(req.params.id))
})
})
// Episodes newer than :afterId, chronological — the un-summarized tail.
app.get('/sessions/:id/episodes/since/:afterId', (req, res) => {
const episodes = episodic.getEpisodesSince(Number(req.params.id), Number(req.params.afterId));
res.json(episodes);
});
app.get('/sessions/:id/consolidation-candidates', (req, res) => {
const sessionId = Number(req.params.id);
const stats = episodic.getSessionEpisodeStats(sessionId);
if(stats.count < CONSOLIDATION.MIN_SESSION_EPISODES) {
return res.json({
eligible: false,
reason: `session has ${stats.count} episodes, floor is ${CONSOLIDATION.MIN_SESSION_EPISODES}`,
candidates: [],
});
}
const candidates = episodic.getConsolidationCandidates(sessionId);
res.json({
eligible:true,
sessionEpisodeCount: stats.count,
candidates
})
})
app.delete('/episodes/:id', (req, res) => {
const id = Number(req.params.id);
episodic.deleteEpisode(id);
@@ -1,4 +1,4 @@
const { SERVICES, getEnv, SUMMARIES } = require('@nexusai/shared');
const { SERVICES, getEnv, SUMMARIES, utilityInference } = require('@nexusai/shared');
const {
getSessionSummariesForProject,
getProjectOverviewSummary,
@@ -9,9 +9,6 @@ const {
const { getEpisodesByProject } = require('../episodic');
const { getProject } = require('../db/projects');
const EXTRACTION_URL = getEnv('EXTRACTION_URL', 'http://localhost:11434');
const EXTRACTION_MODEL = getEnv('EXTRACTION_MODEL', 'qwen2.5:3b');
const MAX_SUMMARY_CHARS = SUMMARIES.MAX_SUMMARY_CHARS; // generous ceiling before we truncate input
function buildProjectSummaryPrompt(projectName, sessionSummaries) {
@@ -24,8 +21,9 @@ function buildProjectSummaryPrompt(projectName, sessionSummaries) {
summaryBlock = summaryBlock.slice(-MAX_SUMMARY_CHARS);
}
// No ChatML wrapper — the model's own prompt template is applied server-side
// by the inference service's /utility/complete route (Ollama /api/chat).
return [
'<|im_start|>user',
`The following are session summaries from a project called "${projectName}".`,
'Write a project overview covering: goals, progress, key decisions, and current state.',
'Scale the length to the material — use multiple paragraphs for complex projects, a few sentences for simple ones.',
@@ -33,8 +31,6 @@ function buildProjectSummaryPrompt(projectName, sessionSummaries) {
'Write in third person. Output only the overview text, no headings or labels.',
'',
summaryBlock,
'<|im_end|>',
'<|im_start|>assistant',
].join('\n');
}
@@ -49,8 +45,9 @@ function buildProjectSummaryFromEpisodesPrompt(projectName, episodes) {
episodeBlock = episodeBlock.slice(-MAX_SUMMARY_CHARS);
}
// No ChatML wrapper — the model's own prompt template is applied server-side
// by the inference service's /utility/complete route (Ollama /api/chat).
return [
'<|im_start|>user',
`The following are conversations from a project called "${projectName}".`,
'Write a project overview covering: goals, progress, key decisions, and current state.',
'Scale the length to the material — use multiple paragraphs for complex projects, a few sentences for simple ones.',
@@ -58,58 +55,25 @@ function buildProjectSummaryFromEpisodesPrompt(projectName, episodes) {
'Write in third person. Output only the overview text, no headings or labels.',
'',
episodeBlock,
'<|im_end|>',
'<|im_start|>assistant',
].join('\n');
}
async function generateProjectSummaryFromEpisodes(projectName, episodes) {
const prompt = buildProjectSummaryFromEpisodesPrompt(projectName, episodes);
const res = await fetch(`${EXTRACTION_URL}/api/generate`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({
model: EXTRACTION_MODEL,
prompt,
stream: false,
options: { temperature: 0.2, num_predict: 1200 },
}),
const user = buildProjectSummaryFromEpisodesPrompt(projectName, episodes);
return utilityInference({
user,
temperature: SUMMARIES.TEMPERATURE,
maxTokens: SUMMARIES.PROJECT_GEN_MAX_TOKENS,
});
if (!res.ok) throw new Error(`Ollama responded ${res.status}`);
const data = await res.json();
const raw = data.response?.trim() ?? '';
return raw
.replace(/<\|im_start\|>.*?<\|im_end\|>/gs, '')
.replace(/<\|im_start\|>|<\|im_end\|>|<\|im_sep\|>/g, '')
.trim();
}
async function generateProjectSummary(projectName, sessionSummaries) {
const prompt = buildProjectSummaryPrompt(projectName, sessionSummaries);
const res = await fetch(`${EXTRACTION_URL}/api/generate`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({
model: EXTRACTION_MODEL,
prompt,
stream: false,
// No format: 'json' — we want free-text narrative, same as session summarization
options: { temperature: 0.2, num_predict: 1200 },
}),
const user = buildProjectSummaryPrompt(projectName, sessionSummaries);
return utilityInference({
user,
temperature: SUMMARIES.TEMPERATURE,
maxTokens: SUMMARIES.PROJECT_GEN_MAX_TOKENS,
});
if (!res.ok) throw new Error(`Ollama responded ${res.status}`);
const data = await res.json();
const raw = data.response?.trim() ?? '';
return raw
.replace(/<\|im_start\|>.*?<\|im_end\|>/gs, '')
.replace(/<\|im_start\|>|<\|im_end\|>|<\|im_sep\|>/g, '')
.trim();
}
// Main entry point — called by the route handler
@@ -2,7 +2,7 @@ const memory = require("../services/memory");
const inference = require("../services/inference");
const embedding = require("../services/embedding");
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 {triggerSummary} = require('../services/summarization')
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 [];
try {
const results = await qdrant.searchEntities(vector, { projectId });
logger.info(
'[orchestration] Entity search results:',
results.map((r) => ({ name: r.payload?.name, score: r.score })),
);
// Include the Qdrant point ID (== SQLite entity ID) for graph traversal
return results.map((r) => r.payload ? { id: r.id, ...r.payload } : null).filter(Boolean);
let allowedIds;
if (projectId === null || projectId === undefined) {
// Non-project chat is its own island — scope to entities linked to
// THIS session. No links yet ⇒ nothing to retrieve, and we return
// early so searchEntities is never called unfiltered.
allowedIds = await memory.getEntityIdsBySession(sessionId);
logger.info(`[orchestration] Non-project chat, session ${sessionId}: ${allowedIds.length} linked entities`);
if (allowedIds.length === 0) return [];
}
const results = await qdrant.searchEntities(vector, { projectId, allowedIds });
logger.info('[orchestration] Entity search results:',
results.map(r => ({ name: r.payload?.name, score: r.score })));
return results.map(r => r.payload ? { id: r.id, ...r.payload } : null).filter(Boolean);
} catch (err) {
logger.debug('[orchestration] Entity search failed, continuing without:', err.message);
return [];
@@ -177,8 +183,11 @@ function fuseEpisodeResults(semanticEps, keywordEps, { semanticWeight, keywordWe
}
function estimateTokens(episode) {
return episode.token_count
?? Math.ceil((episode.user_message.length + episode.ai_response.length) / 4);
//NOTE: episode.token_count is not used here. It stores the
//full inference cost of that turn(prompt + injected memory + response),
//this inflate the episodes apparent size and starving the budget.
//Therefore, char/4 on the actual stored text is an honest selectedl
return Math.ceil((episode.user_message.length + episode.ai_response.length) / 4);
}
function buildScoredPool(fusedWithScores, recentEpisodes, entityBoostedIds, { entityWeight }) {
@@ -252,6 +261,7 @@ async function getFusedEpisodes(userMessage, session, recentIds, projectSessionI
}
async function assembleContext(externalId, userMessage) {
logger.info(`[orchestration] assembleContext ENTERED — msg: "${userMessage}", trivial: ${isTrivialTurn(userMessage)}`);
const settings = appSettings.load();
const { recentEpisodeLimit, semanticLimit, scoreThreshold,
temperature, repeatPenalty, topP, topK, systemPrompt,
@@ -283,20 +293,27 @@ async function assembleContext(externalId, userMessage) {
const isFirstMessage = recentEpisodes.length === 0;
const recentIds = new Set(recentEpisodes.map(e => e.id));
// 4. Embed the query once — the vector is shared by semantic episode search
// and entity search, so embedding it twice was a wasted round-trip + Ollama call.
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);
}
// 4. Retrieval — skipped entirely on contentless social turns (greetings,
// sign-offs). On those, recent history alone is the right context; running
// semantic/keyword/entity retrieval only surfaces marginal noise the model
// then confabulates around. Embed once (shared by episode + entity search).
let fusedWithScores = [];
let entityResults = [];
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)
const [fusedWithScores, entityResults] = await Promise.all([
getFusedEpisodes(userMessage, session, recentIds, projectSessionIds, queryVector, { semanticLimit, scoreThreshold, semanticWeight, keywordWeight }),
getRelevantEntities(queryVector, session.project_id ?? null),
]);
[fusedWithScores, entityResults] = await Promise.all([
getFusedEpisodes(userMessage, session, recentIds, projectSessionIds, queryVector, { semanticLimit, scoreThreshold, semanticWeight, keywordWeight }),
getRelevantEntities(queryVector, { projectId: session.project_id ?? null, sessionId: session.id }),
]);
} else {
logger.debug('[orchestration] Trivial turn — skipping semantic/keyword/entity retrieval');
}
// 5. Entity-linked episode IDs for scoring bonus
const entityIds = entityResults.map(e => e.id);
@@ -314,6 +331,9 @@ async function assembleContext(externalId, userMessage) {
const scoredPool = buildScoredPool(fusedWithScores, recentEpisodes, entityBoostedIds, { entityWeight });
const { guaranteed, selected } = selectWithinBudget(scoredPool, contextBudget, minRecentEpisodes, recentEpisodes);
const selectedIds = selected.map(ep => ep.id);
memory.touchEpisodes(selectedIds);
// 7. Graph neighborhood expansion
let neighborhood = { nodes: [], edges: [] };
if (entityIds.length > 0) {
@@ -1,4 +1,4 @@
const { getEnv, SERVICES, EPISODIC } = require('@nexusai/shared');
const { getEnv, SERVICES, EPISODIC, logger } = require('@nexusai/shared');
const BASE_URL = getEnv('MEMORY_SERVICE_URL', SERVICES.MEMORY_URL);
@@ -216,6 +216,23 @@ async function getEpisodesByEntities(entityIds) {
return res.json(); // { episodeIds: [...] }
}
async function getEntityIdsBySession(sessionId){
const res = await fetch(`${BASE_URL}/sessions/${sessionId}/entity-ids`);
if (!res.ok) throw new Error(`Entity-ids-by-session error: ${res.status}`);
const { entityIds } = await res.json();
return entityIds;
}
// orchestration-service/src/services/memory.js
async function touchEpisodes(ids) {
if (!ids.length) return;
fetch(`${BASE_URL}/episodes/touch`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ ids }),
}).catch(err => logger.warn(`[memory] touch failed: ${err.message}`));
}
module.exports = {
getSessionByExternalId,
createSession,
@@ -242,4 +259,6 @@ module.exports = {
getProjectOverviewSummary,
searchEpisodes,
getEpisodesByEntities,
getEntityIdsBySession,
touchEpisodes,
}
@@ -29,30 +29,28 @@ async function searchEpisodes( vector, {limit = ORCHESTRATION.RECENT_EPISODE_LIM
return data.result;
}
async function searchEntities(vector, { limit = ORCHESTRATION.ENTITIES_LIMIT, scoreThreshold = ORCHESTRATION.ENTITIES_THRESHOLD, projectId = undefined } = {}) {
async function searchEntities(vector, { limit = ORCHESTRATION.ENTITIES_LIMIT, scoreThreshold = ORCHESTRATION.ENTITIES_THRESHOLD, projectId, allowedIds } = {}) {
const body = { vector, limit, score_threshold: scoreThreshold, with_payload: true };
if (projectId !== null && projectId !== undefined) {
body.filter = {
must: [{ key: 'projectId', match: { value: projectId } }]
};
// Project chat: entities shared across the project's sessions.
body.filter = { must: [{ key: 'projectId', match: { value: projectId } }] };
} else if (allowedIds && allowedIds.length > 0) {
// Non-project chat: restrict to this session's own entities (Model 2).
body.filter = { must: [{ has_id: allowedIds }] };
}
// No else: the caller returns early when a non-project session has no linked
// entities, so an unfiltered (leaky) search is never reached.
const res = await fetch(
`${BASE_URL}/collections/${COLLECTIONS.ENTITIES}/points/search`,
{
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify(body),
}
{ method: 'POST', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify(body) }
);
if (!res.ok) {
const body = await res.text();
throw new Error(`Qdrant error: ${res.status} - ${body}`);
const text = await res.text();
throw new Error(`Qdrant error: ${res.status} - ${text}`);
}
const data = await res.json();
return data.result;
return (await res.json()).result;
}
module.exports = { searchEpisodes, searchEntities };
@@ -1,7 +1,5 @@
const { getEnv, SERVICES, SUMMARIES, logger } = require('@nexusai/shared');
const { getEnv, SERVICES, SUMMARIES, logger, utilityInference } = require('@nexusai/shared');
const EXTRACTION_URL = getEnv('EXTRACTION_URL', 'http://localhost:11434');
const EXTRACTION_MODEL = getEnv('EXTRACTION_MODEL', 'qwen2.5:3b');
const MEMORY_URL = getEnv('MEMORY_SERVICE_URL', SERVICES.MEMORY_URL);
const THRESHOLD_TOKENS = parseInt(getEnv('SUMMARY_THRESHOLD_TOKENS', SUMMARIES.THRESHOLD_TOKENS));
@@ -35,41 +33,20 @@ Do not include greetings, sign-offs, or filler. Output only the summary text.
Conversation:
${context}`;
return [
'<|im_start|>user', // ChatML for qwen2.5
instruction,
'<|im_end|>',
'<|im_start|>assistant',
].join('\n');
// No ChatML wrapper — the model's own prompt template is applied server-side
// by the inference service's /utility/complete route (Ollama /api/chat).
return instruction;
}
async function generateSummary(episodes, existingSummary = null) {
const prompt = buildSummaryPrompt(episodes, existingSummary);
const user = buildSummaryPrompt(episodes, existingSummary);
const res = await fetch(`${EXTRACTION_URL}/api/generate`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({
model: EXTRACTION_MODEL,
prompt,
stream: false,
options: {
temperature: 0.2, // slightly higher than entities — summaries benefit from some fluency
num_predict: 500, // generous but bounded — keeps summaries from running long
},
}),
const content = await utilityInference({
user,
temperature: SUMMARIES.TEMPERATURE,
maxTokens: SUMMARIES.SESSION_GEN_MAX_TOKENS,
});
if (!res.ok) throw new Error(`Ollama responded ${res.status}`);
const data = await res.json();
const raw = data.response?.trim() ?? '';
// Strip any leaked ChatML tokens Qwen echoes back
const content = raw
.replace(/<\|im_start\|>.*?<\|im_end\|>/gs, '')
.replace(/<\|im_start\|>|<\|im_end\|>|<\|im_sep\|>/g, '')
.trim();
return content;
}
+23
View File
@@ -78,6 +78,13 @@ const SUMMARIES = {
MIN_EPISODES_SINCE: 5, // don't resummarize until N new episodes since last summary
MAX_SUMMARY_CHARS: 8000, // max chars to include from recent episodes when generating summary (to control prompt size)
MAX_PROJECT_EPISODE_LIMIT: 200, // max number of episodes to consider from the entire project when generating summary (to control prompt size)
// Generation params for the utility model (passed to utilityInference).
// Distinct from MAX_SUMMARY_TOKENS above, which gates STORED summary size;
// these two cap GENERATION length (num_predict) per summary type.
TEMPERATURE: 0.2, // slightly higher than entities (0.1) — summaries benefit from some fluency
SESSION_GEN_MAX_TOKENS: 500, // num_predict for a session summary (3-5 sentences)
PROJECT_GEN_MAX_TOKENS: 1200, // num_predict for a project overview (multi-paragraph)
}
const ENTITIES = {
@@ -105,6 +112,20 @@ const RETRIEVAL = {
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 = {
QDRANT,
COLLECTIONS,
@@ -119,4 +140,6 @@ module.exports = {
SUMMARIES,
ENTITIES,
RETRIEVAL,
UTILITY,
CONSOLIDATION,
};
+24 -2
View File
@@ -1,7 +1,25 @@
const {getEnv} = require('./config/env');
const {QDRANT, COLLECTIONS, EPISODIC, SERVICES, OLLAMA, PORTS, LLAMACPP, INFERENCE_DEFAULTS, SQLITE, ORCHESTRATION, SUMMARIES, ENTITIES, RETRIEVAL } = require('./config/constants');
const {parseRow, formatEpisodeText} = require('./utils')
const {
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 {utilityInference} = require('./utils/utilityInference');
module.exports = {
getEnv,
@@ -17,8 +35,12 @@ module.exports = {
ORCHESTRATION,
parseRow,
formatEpisodeText,
isTrivialTurn,
SUMMARIES,
ENTITIES,
logger,
RETRIEVAL,
UTILITY,
CONSOLIDATION,
utilityInference,
};
+43 -1
View File
@@ -10,4 +10,46 @@ function formatEpisodeText(userMessage, 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 };
+11 -2
View File
@@ -25,14 +25,16 @@ function fakeDb(startVersion = 0) {
test('a fresh DB (user_version 0) is stamped to LATEST_VERSION', () => {
const db = fakeDb(0);
const result = migrate(db);
const stubs = Array.from({ length: LATEST_VERSION }, () => () => {});
const result = migrate(db, stubs); // inject no-op migrations
assert.strictEqual(result, LATEST_VERSION);
assert.strictEqual(db.version, LATEST_VERSION);
});
test('an already-current DB is a no-op (no version writes)', () => {
const db = fakeDb(LATEST_VERSION);
const result = migrate(db);
const stubs = Array.from({ length: LATEST_VERSION }, () => () => {});
const result = migrate(db, stubs); // fixed: db is now the first arg
assert.strictEqual(result, LATEST_VERSION);
assert.deepStrictEqual(db.setCalls, [], 'should not touch user_version when current');
});
@@ -63,3 +65,10 @@ test('only migrations newer than the current version run', () => {
assert.deepStrictEqual(order, ['v2', 'v3'], 'v1 is not re-run');
assert.deepStrictEqual(db.setCalls, [2, 3]);
});
test('LATEST_VERSION reflects the appended v1 → v2 migration', () => {
// Guards that adding the access-tracking migration bumped the target version.
// The real SQL is verified out-of-band (scratch run against a DB copy); here we
// only lock the runner-visible contract: the array length drives LATEST_VERSION.
assert.strictEqual(LATEST_VERSION, 2);
});
+1 -1
View File
@@ -13,7 +13,7 @@ try { ({ DatabaseSync } = require('node:sqlite')); } catch { DatabaseSync = null
// The complete intended shape = original base tables + every historical ALTER.
const EXPECTED = {
sessions: ['id','external_id','created_at','updated_at','metadata','name','project_id'],
episodes: ['id','session_id','user_message','ai_response','created_at','token_count','metadata'],
episodes: ['access_count', 'ai_response', 'created_at', 'id', 'last_accessed_at', 'metadata', 'session_id', 'token_count', 'user_message'],
entities: ['id','name','type','notes','created_at','updated_at','metadata','mention_count','confidence','source','last_seen_at'],
relationships: ['id','from_id','to_id','label','created_at','metadata','mention_count','notes'],
entity_episodes: ['entity_id','episode_id'],
+1 -1
View File
@@ -19,7 +19,7 @@ function mockMemory({ stats, summaries, since }) {
calls.sinceAfterId = Number(u.split('/since/').at(-1));
return json(since.filter(ep => ep.id > calls.sinceAfterId));
}
if (u.includes('/api/generate')) return json({ response: 'A concise third-person summary.' });
if (u.endsWith('/utility/complete')) return json({ text: 'A concise third-person summary.' });
if (u.endsWith('/summaries') && opts.method === 'POST') { calls.posted = JSON.parse(opts.body); return json({ id: 99 }); }
if (/\/summaries\/\d+$/.test(u) && opts.method === 'PATCH') { calls.patched = JSON.parse(opts.body); return json({ ok: true }); }
throw new Error('unexpected fetch: ' + u);
+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);
});