Compare commits
35
Commits
3a6c09931f
...
main
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
3d7372d04f | ||
|
|
c31fb786f4 | ||
|
|
66f6c4534f | ||
|
|
b3765f30fa | ||
|
|
4b2099498c | ||
|
|
fd8867ee96 | ||
|
|
00101ba14a | ||
|
|
6c953a346a | ||
|
|
feadf0a1f7 | ||
|
|
ef3475db28 | ||
|
|
d2a1fb4ca9 | ||
|
|
f6b2ab82fc | ||
|
|
4c00a8d84a | ||
|
|
565286a4d9 | ||
|
|
6aefc391d9 | ||
|
|
4f6ae8492b | ||
|
|
16b1fa2d07 | ||
|
|
eb9fda62f3 | ||
|
|
dce63b6f7b | ||
|
|
e9ceee15fa | ||
|
|
551c7f03ec | ||
|
|
a223b753e5 | ||
|
|
67e2b446e4 | ||
|
|
b9195921dc | ||
|
|
2adefe4df9 | ||
|
|
2bee4b23d6 | ||
|
|
fab8f32395 | ||
|
|
ff1c0ab215 | ||
|
|
7e8aead917 | ||
|
|
1e7c11bad8 | ||
|
|
1df6acc427 | ||
|
|
36c7cf2251 | ||
|
|
0353ce96e0 | ||
|
|
e570bbecf2 | ||
|
|
cbaf83ff17 |
@@ -6,4 +6,5 @@ data/
|
||||
.env.*
|
||||
*.db
|
||||
.claude/settings.local.json
|
||||
.patch
|
||||
EOF
|
||||
@@ -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 |
|
||||
|
||||
@@ -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
@@ -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*
|
||||
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -30,8 +30,6 @@ or inference services — all traffic flows through orchestration.
|
||||
| LLAMA_SERVER_URL | No | http://localhost:8080 | Direct llama-server URL for /models/props |
|
||||
| QDRANT_URL | No | http://localhost:6333 | Qdrant URL for semantic search |
|
||||
| CORS_ORIGIN | No | http://localhost:5173 | Allowed origin for CORS requests |
|
||||
| EXTRACTION_URL | No | http://localhost:11434 | Ollama URL for summarisation |
|
||||
| EXTRACTION_MODEL | No | qwen2.5:3b | Ollama model used for summarisation |
|
||||
|
||||
## Internal Structure
|
||||
|
||||
@@ -75,6 +73,9 @@ via `appSettings.load()` — changes apply immediately without a service restart
|
||||
| `scoreThreshold` | 0.5 | Minimum similarity score for Qdrant semantic results |
|
||||
| `semanticWeight` | 1.0 | RRF weight for Qdrant semantic results |
|
||||
| `keywordWeight` | 0 | RRF weight for FTS5 keyword results (`0` = disabled) |
|
||||
| `contextBudget` | — | Token budget for context assembly (char/4 estimation on stored text) |
|
||||
| `entityWeight` | — | Scoring bonus for entity-linked episodes in the context pool |
|
||||
| `minRecentEpisodes` | — | Guaranteed floor of recent episodes always included in context |
|
||||
| `modelsFolderPath` | `/mnt/nexus-models` | Path to folder containing .gguf files |
|
||||
| `temperature` | 0.7 | Inference temperature |
|
||||
| `repeatPenalty` | 1.1 | Repeat token penalty |
|
||||
@@ -103,35 +104,43 @@ difference is how the inference response is delivered to the client.
|
||||
|
||||
4. **Recent episode retrieval** — fetch most recent episodes (`recentEpisodeLimit`).
|
||||
|
||||
5. **Fused episode retrieval** — runs semantic (Qdrant) and keyword (FTS5)
|
||||
5. **Trivial-turn gate** — greetings/pleasantries (`isTrivialTurn`) skip all
|
||||
retrieval (semantic, keyword, entity); recent history alone is the context.
|
||||
Breaks the greeting → marginal-retrieval → confabulation loop.
|
||||
|
||||
6. **Fused episode retrieval** — runs semantic (Qdrant) and keyword (FTS5)
|
||||
search in parallel, then merges results via Reciprocal Rank Fusion (RRF).
|
||||
Both paths are filtered against `recentIds` before fusion. FTS is scoped
|
||||
to the current session or all project sessions. If `keywordWeight` is `0`,
|
||||
the FTS call is skipped entirely. Non-critical — failures fall back to
|
||||
whichever strategy succeeded.
|
||||
The query is embedded once and shared with entity search. Both paths are
|
||||
filtered against `recentIds` before fusion. FTS is scoped to the current
|
||||
session or all project sessions. If `keywordWeight` is `0`, the FTS call
|
||||
is skipped entirely. Non-critical — failures fall back to whichever
|
||||
strategy succeeded.
|
||||
|
||||
6. **Entity search** — query `entities` Qdrant collection filtered by
|
||||
`projectId`. Returns entity IDs alongside Qdrant payload data (the Qdrant
|
||||
point ID equals the SQLite entity ID). Non-critical.
|
||||
7. **Entity search + graph expansion** — query `entities` Qdrant collection
|
||||
(project-scoped, or session-scoped via `/sessions/:id/entity-ids` for
|
||||
non-project chats). Entity IDs are expanded into a 1-hop subgraph via
|
||||
`POST /graph/neighbors`; on failure, falls back to flat entity list.
|
||||
Non-critical.
|
||||
|
||||
7. **Graph neighborhood expansion** — call `POST /graph/neighbors` on
|
||||
memory-service with the entity IDs from step 6. Returns a 1-hop subgraph
|
||||
`{ nodes, edges }` — entity objects plus the relationships connecting them.
|
||||
If no entities were found or the graph call fails, falls back to flat entity
|
||||
list (no edges). Non-critical.
|
||||
8. **Scored pool + budget selection** — `buildScoredPool` combines RRF scores,
|
||||
recency, and entity-linkage bonus; `selectWithinBudget` fills `contextBudget`
|
||||
(char/4 token estimation on stored text) above a guaranteed floor of
|
||||
`minRecentEpisodes` recent episodes. Selected episode IDs are then reported
|
||||
to `POST /episodes/touch` fire-and-forget (access tracking for the
|
||||
consolidation lifecycle).
|
||||
|
||||
8. **Prompt assembly** — combine system prompt, graph context, fused episodes,
|
||||
recent episodes, and user message.
|
||||
9. **Prompt assembly** — combine system prompt, graph context, selected
|
||||
episodes, guaranteed recent episodes, and user message.
|
||||
|
||||
9. **Inference** — send to inference service. `/chat` awaits full response;
|
||||
10. **Inference** — send to inference service. `/chat` awaits full response;
|
||||
`/chat/stream` pipes SSE chunks to the client.
|
||||
|
||||
10. **Episode write** — write exchange back to memory with `projectId`.
|
||||
11. **Episode write** — write exchange back to memory with `projectId`.
|
||||
|
||||
11. **Summarisation trigger** — `triggerSummary(session, allEpisodes)` called
|
||||
12. **Summarisation trigger** — `triggerSummary(session)` called
|
||||
fire-and-forget. See `summarization.md` for full details.
|
||||
|
||||
12. **Auto-naming** — on first message with no session name, fires a secondary
|
||||
13. **Auto-naming** — on first message with no session name, fires a secondary
|
||||
inference call (max 20 tokens, temperature 0.3) to generate a session name.
|
||||
|
||||
### Prompt Structure
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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 |
|
||||
|
||||
@@ -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,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;
|
||||
@@ -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};
|
||||
@@ -1,5 +1,6 @@
|
||||
const Database = require('better-sqlite3');
|
||||
const schema = require('./schema');
|
||||
const { migrate } = require('./migrations');
|
||||
const { getEnv, SQLITE, logger } = require('@nexusai/shared');
|
||||
|
||||
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('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(`ALTER TABLE sessions ADD COLUMN name TEXT`)
|
||||
} catch {}
|
||||
db.exec(schema); // complete current shape — fresh DBs get everything
|
||||
migrate(db); // carry an older DB forward; no-op on fresh/current DBs
|
||||
|
||||
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`); // ← 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`);
|
||||
// One-time FTS backfill: only when the index was just created on a DB that
|
||||
// already holds episodes (i.e. episodes predate FTS). During normal
|
||||
// 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
|
||||
// external-content FTS5 table proxies the content table, so it can't detect
|
||||
// a desync — the "was it just created" check is what makes this correct.
|
||||
if (!ftsExisted) {
|
||||
const epCount = db.prepare('SELECT COUNT(*) AS c FROM episodes').get().c;
|
||||
if (epCount > 0) {
|
||||
db.exec(`INSERT INTO episodes_fts(episodes_fts) VALUES('rebuild')`);
|
||||
logger.info(`[db] Backfilled FTS index for ${epCount} pre-existing episodes`);
|
||||
}
|
||||
}
|
||||
|
||||
logger.info(`Connected to SQLite database at ${path}`);
|
||||
}
|
||||
|
||||
@@ -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 };
|
||||
@@ -1,10 +1,17 @@
|
||||
// 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 = `
|
||||
CREATE TABLE IF NOT EXISTS sessions (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
external_id TEXT UNIQUE NOT NULL,
|
||||
created_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 (
|
||||
@@ -14,7 +21,9 @@ const schema = `
|
||||
ai_response TEXT NOT NULL,
|
||||
created_at INTEGER NOT NULL DEFAULT (unixepoch()),
|
||||
token_count INTEGER,
|
||||
metadata TEXT
|
||||
metadata TEXT,
|
||||
last_accessed_at INTEGER,
|
||||
access_count INTEGER NOT NULL DEFAULT 0
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS entities (
|
||||
@@ -25,6 +34,10 @@ const schema = `
|
||||
created_at INTEGER NOT NULL DEFAULT (unixepoch()),
|
||||
updated_at INTEGER NOT NULL DEFAULT (unixepoch()),
|
||||
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)
|
||||
);
|
||||
|
||||
@@ -35,6 +48,8 @@ const schema = `
|
||||
label TEXT NOT NULL,
|
||||
created_at INTEGER NOT NULL DEFAULT (unixepoch()),
|
||||
metadata TEXT,
|
||||
mention_count INTEGER NOT NULL DEFAULT 1,
|
||||
notes TEXT,
|
||||
UNIQUE(from_id, to_id, label)
|
||||
);
|
||||
|
||||
@@ -50,15 +65,16 @@ const schema = `
|
||||
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 TABLE IF NOT EXISTS projects (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
name TEXT NOT NULL,
|
||||
description TEXT,
|
||||
colour 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 (
|
||||
@@ -72,12 +88,12 @@ const schema = `
|
||||
metadata TEXT
|
||||
);
|
||||
|
||||
CREATE INDEX IF NOT EXISTS idx_episodes_session
|
||||
ON episodes(session_id);
|
||||
CREATE INDEX IF NOT EXISTS idx_episodes_created
|
||||
ON episodes(created_at);
|
||||
CREATE INDEX IF NOT EXISTS idx_entities_type
|
||||
ON entities(type);
|
||||
CREATE INDEX IF NOT EXISTS idx_episodes_session ON episodes(session_id);
|
||||
CREATE INDEX IF NOT EXISTS idx_episodes_created ON episodes(created_at);
|
||||
CREATE INDEX IF NOT EXISTS idx_entities_type ON entities(type);
|
||||
CREATE INDEX IF NOT EXISTS idx_sessions_project ON sessions(project_id);
|
||||
CREATE INDEX IF NOT EXISTS idx_summaries_project ON summaries(project_id);
|
||||
CREATE INDEX IF NOT EXISTS idx_summaries_session ON summaries(session_id);
|
||||
|
||||
CREATE VIRTUAL TABLE IF NOT EXISTS episodes_fts
|
||||
USING fts5(user_message, ai_response, content=episodes, content_rowid=id);
|
||||
@@ -101,8 +117,6 @@ const schema = `
|
||||
INSERT INTO episodes_fts(rowid, user_message, ai_response)
|
||||
VALUES (new.id, new.user_message, new.ai_response);
|
||||
END;
|
||||
|
||||
|
||||
`;
|
||||
|
||||
module.exports = schema;
|
||||
@@ -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,11 +38,9 @@ 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',
|
||||
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)',
|
||||
@@ -65,9 +60,8 @@ function buildExtractionPrompt(userMessage, aiResponse, knownEntities = []) {
|
||||
`User: ${userMessage}`,
|
||||
`Assistant: ${aiResponse}`,
|
||||
'--- END CONVERSATION ---',
|
||||
'<|im_end|>',
|
||||
'<|im_start|>assistant',
|
||||
].join('\n');
|
||||
].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: {
|
||||
const raw = await utilityInference({
|
||||
system,
|
||||
user,
|
||||
json: true,
|
||||
temperature: ENTITIES.TEMPERATURE,
|
||||
num_predict: ENTITIES.NUM_PREDICT,
|
||||
},
|
||||
}),
|
||||
signal: AbortSignal.timeout(60_000),
|
||||
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');
|
||||
|
||||
@@ -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));
|
||||
|
||||
// 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,
|
||||
};
|
||||
@@ -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 };
|
||||
|
||||
@@ -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,8 +293,13 @@ 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.
|
||||
// 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);
|
||||
@@ -292,11 +307,13 @@ async function assembleContext(externalId, userMessage) {
|
||||
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([
|
||||
[fusedWithScores, entityResults] = await Promise.all([
|
||||
getFusedEpisodes(userMessage, session, recentIds, projectSessionIds, queryVector, { semanticLimit, scoreThreshold, semanticWeight, keywordWeight }),
|
||||
getRelevantEntities(queryVector, session.project_id ?? null),
|
||||
getRelevantEntities(queryVector, { projectId: session.project_id ?? null, sessionId: session.id }),
|
||||
]);
|
||||
} else {
|
||||
logger.debug('[orchestration] Trivial turn — skipping semantic/keyword/entity retrieval');
|
||||
}
|
||||
|
||||
// 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;
|
||||
}
|
||||
|
||||
|
||||
@@ -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 = {
|
||||
@@ -101,10 +108,24 @@ const ENTITIES = {
|
||||
|
||||
const RETRIEVAL = {
|
||||
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
|
||||
}
|
||||
|
||||
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,
|
||||
};
|
||||
@@ -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,
|
||||
};
|
||||
@@ -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 };
|
||||
@@ -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);
|
||||
});
|
||||
@@ -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);
|
||||
});
|
||||
@@ -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);
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
Reference in New Issue
Block a user