Compare commits
33
Commits
e570bbecf2
..
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 |
@@ -6,4 +6,5 @@ data/
|
|||||||
.env.*
|
.env.*
|
||||||
*.db
|
*.db
|
||||||
.claude/settings.local.json
|
.claude/settings.local.json
|
||||||
|
.patch
|
||||||
EOF
|
EOF
|
||||||
@@ -205,6 +205,9 @@ Returns `503` if llama-server is unreachable.
|
|||||||
| `scoreThreshold` | float | 0–1 | Minimum similarity score for Qdrant results |
|
| `scoreThreshold` | float | 0–1 | Minimum similarity score for Qdrant results |
|
||||||
| `semanticWeight` | float | 0–5 | RRF weight for Qdrant semantic results |
|
| `semanticWeight` | float | 0–5 | RRF weight for Qdrant semantic results |
|
||||||
| `keywordWeight` | float | 0–5 | RRF weight for FTS5 keyword results (`0` = disabled) |
|
| `keywordWeight` | float | 0–5 | RRF weight for FTS5 keyword results (`0` = disabled) |
|
||||||
|
| `contextBudget` | integer | — | Token budget for context assembly (char/4 estimation) |
|
||||||
|
| `entityWeight` | float | — | Scoring bonus for entity-linked episodes in the context pool |
|
||||||
|
| `minRecentEpisodes` | integer | — | Guaranteed floor of recent episodes always included in context |
|
||||||
| `modelsFolderPath` | string | — | Path to folder containing .gguf files |
|
| `modelsFolderPath` | string | — | Path to folder containing .gguf files |
|
||||||
| `temperature` | float | 0–2 | Inference randomness |
|
| `temperature` | float | 0–2 | Inference randomness |
|
||||||
| `repeatPenalty` | float | 1–2 | Repeat token penalty |
|
| `repeatPenalty` | float | 1–2 | Repeat token penalty |
|
||||||
@@ -253,6 +256,7 @@ orchestration.
|
|||||||
| GET | /sessions/by-external/:externalId | Get session by external ID |
|
| GET | /sessions/by-external/:externalId | Get session by external ID |
|
||||||
| PATCH | /sessions/by-external/:externalId | Update session fields |
|
| PATCH | /sessions/by-external/:externalId | Update session fields |
|
||||||
| DELETE | /sessions/by-external/:externalId | Delete session (cascades to episodes) |
|
| DELETE | /sessions/by-external/:externalId | Delete session (cascades to episodes) |
|
||||||
|
| GET | /sessions/:id/entity-ids | Entity IDs linked to this session's episodes (non-project memory isolation scoping) |
|
||||||
|
|
||||||
> Route ordering: `by-external/:externalId` must be defined before `/:id`
|
> Route ordering: `by-external/:externalId` must be defined before `/:id`
|
||||||
> to prevent `by-external` being captured as an ID param.
|
> to prevent `by-external` being captured as an ID param.
|
||||||
@@ -277,6 +281,10 @@ Both fields are optional. Only provided fields are updated.
|
|||||||
| GET | /episodes/search?q=&limit= | FTS keyword search across all episodes |
|
| GET | /episodes/search?q=&limit= | FTS keyword search across all episodes |
|
||||||
| GET | /episodes/:id | Get episode by ID |
|
| GET | /episodes/:id | Get episode by ID |
|
||||||
| GET | /sessions/:id/episodes?limit=&offset= | Paginated episodes for a session |
|
| GET | /sessions/:id/episodes?limit=&offset= | Paginated episodes for a session |
|
||||||
|
| GET | /sessions/:id/episode-stats | Aggregate: count, total tokens, max id (summarization threshold check) |
|
||||||
|
| GET | /sessions/:id/episodes/since/:afterId | Episodes newer than :afterId, chronological (un-summarized tail) |
|
||||||
|
| POST | /episodes/touch | Batch access-tracking bump — increments `access_count`, sets `last_accessed_at` |
|
||||||
|
| GET | /sessions/:id/consolidation-candidates | Dry-run consolidation scoring (aging score, floors applied) |
|
||||||
| DELETE | /episodes/:id | Delete episode (SQLite + Qdrant cleanup) |
|
| DELETE | /episodes/:id | Delete episode (SQLite + Qdrant cleanup) |
|
||||||
|
|
||||||
> Route ordering: `/episodes/search` must be defined before `/episodes/:id`.
|
> Route ordering: `/episodes/search` must be defined before `/episodes/:id`.
|
||||||
@@ -291,6 +299,21 @@ Both fields are optional. Only provided fields are updated.
|
|||||||
}
|
}
|
||||||
```
|
```
|
||||||
|
|
||||||
|
**POST /episodes/touch — body:**
|
||||||
|
```json
|
||||||
|
{ "ids": [54, 55, 56] }
|
||||||
|
```
|
||||||
|
Returns `{ "touched": 3 }`. Called fire-and-forget by orchestration after
|
||||||
|
budget selection. Nonexistent IDs are silently skipped; `touched` echoes the
|
||||||
|
request count, not rows matched.
|
||||||
|
|
||||||
|
**GET /sessions/:id/consolidation-candidates** — returns
|
||||||
|
`{ eligible, reason?, candidates: [...] }` where each candidate is
|
||||||
|
`{ id, access_count, created_at, last_accessed_at, preview, aging_score }`,
|
||||||
|
sorted by `aging_score` ascending (most eligible first). Sessions under
|
||||||
|
`CONSOLIDATION.MIN_SESSION_EPISODES` return `eligible: false` with a `reason`.
|
||||||
|
Observe-only — nothing is modified.
|
||||||
|
|
||||||
### Projects
|
### Projects
|
||||||
|
|
||||||
| Method | Path | Description |
|
| Method | Path | Description |
|
||||||
|
|||||||
@@ -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.
|
Multi-strategy retrieval merged into a single ranked result set.
|
||||||
- [x] Reciprocal Rank Fusion (RRF) — merge semantic (Qdrant) + keyword (FTS5) results
|
- [x] Reciprocal Rank Fusion (RRF) — merge semantic (Qdrant) + keyword (FTS5) results
|
||||||
- [x] Configurable weights per retrieval strategy (`semanticWeight`, `keywordWeight` via `PATCH /settings`)
|
- [x] Configurable weights per retrieval strategy (`semanticWeight`, `keywordWeight` via `PATCH /settings`)
|
||||||
- [x] Score threshold retained per-strategy; FTS scoped to session/project sessions; `keywordWeight: 0` default (disabled until tuned)
|
- [x] Score threshold retained per-strategy; FTS scoped to session/project sessions
|
||||||
|
- [x] FTS query tokenization (`buildFtsQuery`) — tokenize + stopword-filter + quoted `OR` terms; replaced whole-phrase matching that made keyword recall nil
|
||||||
|
- [x] Fusion tuned and verified live (keyword path + project scoping confirmed via weight inversion); running `keywordWeight: 0.5` / `semanticWeight: 1.0`
|
||||||
|
|
||||||
### 3. Memory Consolidation Lifecycle
|
### 3. Memory Consolidation Lifecycle
|
||||||
Prevents long-term memory degradation and enables compression.
|
Prevents long-term memory degradation and enables compression.
|
||||||
- [ ] Episode aging — score/weight episodes by recency and access frequency
|
- [x] Episode aging — `access_count` + `last_accessed_at` columns (v2 migration), batch touch on retrieval selection (`POST /episodes/touch`), aging score `access_count / (1 + days since last access)`
|
||||||
- [ ] Consolidation pass — merge related low-weight episodes into summary nodes
|
- [x] Dry-run candidates endpoint — `GET /sessions/:id/consolidation-candidates` with age + session-size floors (observe-only phase before destructive pass)
|
||||||
|
- [ ] Consolidation pass — merge related low-weight episodes into summary nodes (incl. Qdrant vector + entity link cleanup)
|
||||||
- [ ] Orphan cleanup — remove entities no longer referenced by active episodes
|
- [ ] Orphan cleanup — remove entities no longer referenced by active episodes
|
||||||
|
|
||||||
### 4. User Preference Model
|
### 4. User Preference Model
|
||||||
@@ -88,11 +91,11 @@ Short-circuit simple requests before they reach the LLM.
|
|||||||
- [ ] Confidence bands — FAST PATH (memory lookup only) vs FULL (LLM + context)
|
- [ ] Confidence bands — FAST PATH (memory lookup only) vs FULL (LLM + context)
|
||||||
- [ ] Fast-path handlers — direct memory queries, session lookups, factual recalls
|
- [ ] Fast-path handlers — direct memory queries, session lookups, factual recalls
|
||||||
|
|
||||||
### 6. Smarter Context Assembly *(inspired by acid2lake)*
|
### 6. Smarter Context Assembly *(inspired by acid2lake)* ✅
|
||||||
Budget-aware context selection instead of dumping all relevant memory into the prompt.
|
Budget-aware context selection instead of dumping all relevant memory into the prompt.
|
||||||
- [ ] Token budget manager in orchestration
|
- [x] Token budget manager in orchestration (`selectWithinBudget`, char/4 estimation on stored text)
|
||||||
- [ ] Priority scoring — recency × relevance × entity weight
|
- [x] Priority scoring — RRF fusion + recency + entity boost (`buildScoredPool`)
|
||||||
- [ ] Configurable context budget via env var
|
- [x] Configurable via settings (`contextBudget`, `entityWeight`, `minRecentEpisodes`) — live, no restart
|
||||||
|
|
||||||
### 7. Procedural Memory Store *(inspired by acid2lake)*
|
### 7. Procedural Memory Store *(inspired by acid2lake)*
|
||||||
Learns "how NexusAI has successfully handled this type of request before."
|
Learns "how NexusAI has successfully handled this type of request before."
|
||||||
@@ -225,4 +228,4 @@ The JARVIS moment — NexusAI reasons, plans, and acts across multiple steps.
|
|||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
*Last updated: April 2026*
|
*Last updated: August 2026*
|
||||||
@@ -2,7 +2,7 @@
|
|||||||
|
|
||||||
**Location:** `packages/memory-service/src/entities/extraction.js`
|
**Location:** `packages/memory-service/src/entities/extraction.js`
|
||||||
**Triggered by:** Episode creation (`POST /episodes`)
|
**Triggered by:** Episode creation (`POST /episodes`)
|
||||||
**Model:** `qwen2.5:3b` via Ollama (configurable via `EXTRACTION_MODEL` env var)
|
**Model:** the utility model served by the inference service (`/utility/complete`), configurable via `UTILITY_MODEL` on the inference service
|
||||||
|
|
||||||
## Purpose
|
## Purpose
|
||||||
|
|
||||||
@@ -28,7 +28,7 @@ swallowed.
|
|||||||
|
|
||||||
| Setting | Value | Notes |
|
| Setting | Value | Notes |
|
||||||
|---|---|---|
|
|---|---|---|
|
||||||
| Model | `qwen2.5:3b` | Ollama, configurable via `EXTRACTION_MODEL` |
|
| Model | utility model | Served by inference-service `/utility/complete`, set via `UTILITY_MODEL` |
|
||||||
| Temperature | 0.1 | Low for consistent, deterministic output |
|
| Temperature | 0.1 | Low for consistent, deterministic output |
|
||||||
| `num_predict` | 1500 | Higher ceiling to accommodate entity + relationship JSON |
|
| `num_predict` | 1500 | Higher ceiling to accommodate entity + relationship JSON |
|
||||||
| `format` | `'json'` | Ollama constrained decoding — enforces valid JSON output |
|
| `format` | `'json'` | Ollama constrained decoding — enforces valid JSON output |
|
||||||
@@ -131,6 +131,20 @@ check alone won't reject.
|
|||||||
> imperfectly — switch to `/[^\p{L}\p{N}\s]/gu` if the set gains non-ASCII
|
> imperfectly — switch to `/[^\p{L}\p{N}\s]/gu` if the set gains non-ASCII
|
||||||
> entries.
|
> entries.
|
||||||
|
|
||||||
|
**Regurgitation guard (`mentionedIn`):** the extraction prompt feeds the model
|
||||||
|
a "known entities" hint block (the 20 most-recent entities) for spelling/type
|
||||||
|
consistency. The small model (qwen2.5:3b) will sometimes echo that list back as
|
||||||
|
if those entities appeared in the conversation — most visibly on contentless
|
||||||
|
turns (a greeting produced fake extractions of unrelated authors, game titles,
|
||||||
|
etc.). After parsing, each extracted name is checked against the actual
|
||||||
|
`userMessage + aiResponse` text (case- and whitespace-normalized substring); any
|
||||||
|
name not present is dropped before upsert. Since the prompt constrains names to
|
||||||
|
short proper nouns, a genuinely-discussed entity appears verbatim while a
|
||||||
|
regurgitated hint does not. Relationships referencing a dropped entity fall away
|
||||||
|
automatically (they resolve against the surviving `entityMap`). A prompt line
|
||||||
|
also tells the model the hint list is spelling-only — a backstop, with
|
||||||
|
`mentionedIn` as the deterministic guarantee.
|
||||||
|
|
||||||
## Relationship Processing
|
## Relationship Processing
|
||||||
|
|
||||||
After all entities are saved, relationships are processed:
|
After all entities are saved, relationships are processed:
|
||||||
|
|||||||
@@ -28,16 +28,16 @@ relationship extraction and embeds results into Qdrant.
|
|||||||
| SQLITE_PATH | Yes | — | Path to SQLite database file |
|
| SQLITE_PATH | Yes | — | Path to SQLite database file |
|
||||||
| QDRANT_URL | No | http://localhost:6333 | Qdrant instance URL |
|
| QDRANT_URL | No | http://localhost:6333 | Qdrant instance URL |
|
||||||
| EMBEDDING_SERVICE_URL | No | http://localhost:3003 | Embedding service URL |
|
| EMBEDDING_SERVICE_URL | No | http://localhost:3003 | Embedding service URL |
|
||||||
| EXTRACTION_URL | No | http://localhost:11434 | Ollama URL for entity extraction |
|
| INFERENCE_SERVICE_URL | No | http://localhost:3001 | Inference service URL — entity extraction routes through its `/utility/complete` endpoint |
|
||||||
| EXTRACTION_MODEL | No | qwen2.5:3b | Ollama model used for entity extraction |
|
|
||||||
|
|
||||||
## Internal Structure
|
## Internal Structure
|
||||||
|
|
||||||
```
|
```
|
||||||
src/
|
src/
|
||||||
├── db/
|
├── db/
|
||||||
│ ├── index.js # SQLite connection + initialization + migrations
|
│ ├── index.js # SQLite connection + init + migrate() + one-time FTS backfill
|
||||||
│ ├── schema.js # Table definitions, indexes, FTS5, triggers
|
│ ├── migrations.js # Forward-only versioned migration runner (PRAGMA user_version)
|
||||||
|
│ ├── schema.js # Complete current shape: tables, indexes, FTS5, triggers
|
||||||
│ ├── projects.js # Project CRUD functions
|
│ ├── projects.js # Project CRUD functions
|
||||||
│ └── summaries.js # Summary CRUD functions
|
│ └── summaries.js # Summary CRUD functions
|
||||||
├── episodic/
|
├── episodic/
|
||||||
@@ -64,36 +64,61 @@ Eight core tables:
|
|||||||
- **summaries** — condensed episode groups for efficient context retrieval
|
- **summaries** — condensed episode groups for efficient context retrieval
|
||||||
- **projects** — named groupings of sessions with `name`, `description`, `colour`, `icon`, `isolated`, `notes`, `system_prompt`
|
- **projects** — named groupings of sessions with `name`, `description`, `colour`, `icon`, `isolated`, `notes`, `system_prompt`
|
||||||
|
|
||||||
### Migrations
|
### Schema & Migrations
|
||||||
|
|
||||||
Schema changes that cannot use `CREATE TABLE IF NOT EXISTS` are applied as
|
`schema.js` holds the **complete current shape** — every table, column, index,
|
||||||
idempotent migrations in `db/index.js` at startup:
|
the FTS5 virtual table, and its triggers — as the single source of truth for a
|
||||||
|
fresh database. It uses `CREATE TABLE IF NOT EXISTS`, so on a fresh DB it builds
|
||||||
|
everything; on an existing DB it skips tables that already exist (and therefore
|
||||||
|
does **not** reconcile columns on old tables — that's what migrations are for).
|
||||||
|
|
||||||
|
`db/migrations.js` is a forward-only versioned runner keyed on
|
||||||
|
`PRAGMA user_version`:
|
||||||
|
|
||||||
```js
|
```js
|
||||||
try { db.exec(`ALTER TABLE sessions ADD COLUMN name TEXT`); } catch {}
|
const migrations = [
|
||||||
try { db.exec(`ALTER TABLE sessions ADD COLUMN project_id INTEGER REFERENCES projects(id)`); } catch {}
|
(_db) => {}, // v0 → v1: consolidated baseline (historical ALTERs folded into schema.js)
|
||||||
try { db.exec(`CREATE INDEX IF NOT EXISTS idx_sessions_project ON sessions(project_id)`); } catch {}
|
(db) => { // v1 → v2: access tracking for consolidation lifecycle
|
||||||
try { db.exec(`ALTER TABLE projects ADD COLUMN isolated INTEGER NOT NULL DEFAULT 0`); } catch {}
|
db.exec(`ALTER TABLE episodes ADD COLUMN last_accessed_at INTEGER`);
|
||||||
try { db.exec(`ALTER TABLE projects ADD COLUMN notes TEXT`); } catch {}
|
db.exec(`ALTER TABLE episodes ADD COLUMN access_count INTEGER NOT NULL DEFAULT 0`);
|
||||||
try { db.exec(`ALTER TABLE projects ADD COLUMN system_prompt TEXT`); } catch {}
|
db.exec(`UPDATE episodes SET last_accessed_at = created_at`); // backfill
|
||||||
// Knowledge graph columns:
|
},
|
||||||
try { db.exec(`ALTER TABLE entities ADD COLUMN mention_count INTEGER NOT NULL DEFAULT 1`) } catch {}
|
];
|
||||||
try { db.exec(`ALTER TABLE entities ADD COLUMN confidence REAL NOT NULL DEFAULT 1.0`) } catch {}
|
const LATEST_VERSION = migrations.length; // derived, never hand-maintained
|
||||||
try { db.exec(`ALTER TABLE entities ADD COLUMN source TEXT NOT NULL DEFAULT 'extraction'`) } catch {}
|
|
||||||
try { db.exec(`ALTER TABLE entities ADD COLUMN last_seen_at INTEGER`) } catch {}
|
|
||||||
try { db.exec(`ALTER TABLE relationships ADD COLUMN mention_count INTEGER NOT NULL DEFAULT 1`) } catch {}
|
|
||||||
try { db.exec(`ALTER TABLE relationships ADD COLUMN notes TEXT`) } catch {}
|
|
||||||
```
|
```
|
||||||
|
|
||||||
`entity_episodes` is defined in `schema.js` itself (not a migration) since it is a new table.
|
`migrate(db)` reads `user_version`, applies every entry newer than it (each in a
|
||||||
|
transaction alongside its version bump), and stamps the result. A fresh DB is
|
||||||
|
built whole by `schema.js` and simply stamped to `LATEST_VERSION`; the baseline
|
||||||
|
entry is a no-op.
|
||||||
|
|
||||||
New migrations are always appended — never modify the schema file for existing tables since `ALTER TABLE` cannot use `IF NOT EXISTS`.
|
**Adding a schema change:** append a new function to the `migrations` array
|
||||||
|
(which bumps `LATEST_VERSION` automatically). Never edit an already-shipped
|
||||||
|
entry, and never edit a table in `schema.js` expecting existing DBs to pick it
|
||||||
|
up — they won't. This replaces the previous pattern of stacking silent
|
||||||
|
`try/catch ALTER TABLE` statements in `db/index.js` on every boot.
|
||||||
|
|
||||||
|
> **Consolidation note:** the historical ALTERs were folded into `schema.js`
|
||||||
|
> rather than preserved as replayable migrations, so this assumes a fresh
|
||||||
|
> database (which is the case post-wipe). An older, pre-consolidation database
|
||||||
|
> would **not** auto-upgrade — `schema.js` skips its existing tables and the
|
||||||
|
> baseline migration is a no-op. To support upgrading old DBs, the v1 baseline
|
||||||
|
> would instead perform guarded (`ADD COLUMN if missing`) catch-up.
|
||||||
|
|
||||||
### FTS5 Full-Text Search
|
### FTS5 Full-Text Search
|
||||||
|
|
||||||
An `episodes_fts` virtual table enables keyword search across all episodes.
|
An `episodes_fts` external-content virtual table enables keyword search across
|
||||||
Three triggers (`episodes_fts_insert`, `episodes_fts_update`, `episodes_fts_delete`)
|
episodes. Three triggers (`episodes_fts_insert`, `episodes_fts_update`,
|
||||||
keep the FTS index automatically in sync with the episodes table.
|
`episodes_fts_delete`) keep the index in sync with the `episodes` table
|
||||||
|
automatically during normal operation.
|
||||||
|
|
||||||
|
A one-time backfill in `db/index.js` handles the case where the FTS table is
|
||||||
|
created on a DB that already holds episodes (e.g. episodes predating FTS). It is
|
||||||
|
gated on "did `episodes_fts` not exist before this boot," checked via
|
||||||
|
`sqlite_master` **before** running the schema — not on a row-count comparison,
|
||||||
|
because `COUNT(*)` on an external-content FTS5 table proxies the content table
|
||||||
|
and cannot detect a desync. This replaced an unconditional full FTS rebuild that
|
||||||
|
previously ran on every startup.
|
||||||
|
|
||||||
### SQLite Configuration
|
### SQLite Configuration
|
||||||
|
|
||||||
@@ -101,6 +126,12 @@ keep the FTS index automatically in sync with the episodes table.
|
|||||||
- `foreign_keys = ON` — enforces referential integrity and cascade deletes
|
- `foreign_keys = ON` — enforces referential integrity and cascade deletes
|
||||||
- PRAGMAs set via `db.pragma()`, not `db.exec()`
|
- PRAGMAs set via `db.pragma()`, not `db.exec()`
|
||||||
|
|
||||||
|
> **Copying a live WAL database:** `cp` on the `.db` file alone silently loses
|
||||||
|
> everything in the un-checkpointed `-wal` file (recent writes, even the
|
||||||
|
> migration version stamp). Always use
|
||||||
|
> `sqlite3 nexusai.db "VACUUM INTO './copy.db'"` (or `.backup`) — safe while
|
||||||
|
> the service is running, produces a complete single-file snapshot.
|
||||||
|
|
||||||
### Dynamic Updates
|
### Dynamic Updates
|
||||||
|
|
||||||
Both `updateSession` and `updateProject` build their `SET` clause dynamically
|
Both `updateSession` and `updateProject` build their `SET` clause dynamically
|
||||||
@@ -184,6 +215,28 @@ service is responsible only for CRUD — generation logic lives in orchestration
|
|||||||
> For full details on trigger conditions, prompt format, cumulative updates,
|
> For full details on trigger conditions, prompt format, cumulative updates,
|
||||||
> and ChatML token stripping, see `summarization.md`.
|
> and ChatML token stripping, see `summarization.md`.
|
||||||
|
|
||||||
|
## Access Tracking & Consolidation (dry-run)
|
||||||
|
|
||||||
|
Every episode selected into a chat context window (budget-selected, not the
|
||||||
|
guaranteed-recency floor) gets an access bump via `POST /episodes/touch` —
|
||||||
|
`access_count` incremented, `last_accessed_at` set to `Date.now()` (ms).
|
||||||
|
Called fire-and-forget from orchestration; a failure loses one increment,
|
||||||
|
nothing more.
|
||||||
|
|
||||||
|
`GET /sessions/:id/consolidation-candidates` scores episodes by
|
||||||
|
`access_count / (1 + days since last access)` — never-accessed episodes fall
|
||||||
|
back to `created_at` for the recency term and score exactly 0 (most eligible).
|
||||||
|
Two floors apply: episodes younger than `CONSOLIDATION.MIN_AGE_DAYS` are
|
||||||
|
excluded in SQL; sessions under `CONSOLIDATION.MIN_SESSION_EPISODES` return
|
||||||
|
`eligible: false` before scoring runs. The endpoint is observe-only — the
|
||||||
|
destructive pass (merge → summarize → Qdrant cleanup → orphan sweep) is not
|
||||||
|
yet built.
|
||||||
|
|
||||||
|
> **Unit note:** `created_at` is unix **seconds** (`unixepoch()`);
|
||||||
|
> `last_accessed_at` is unix **milliseconds** (`Date.now()`). The scoring
|
||||||
|
> query normalizes with `created_at * 1000`. Keep this in mind for any new
|
||||||
|
> queries touching both columns.
|
||||||
|
|
||||||
## Delete Behaviour (SQLite + Qdrant consistency)
|
## Delete Behaviour (SQLite + Qdrant consistency)
|
||||||
|
|
||||||
SQLite cascades handle relational cleanup, but Qdrant is a separate store and
|
SQLite cascades handle relational cleanup, but Qdrant is a separate store and
|
||||||
|
|||||||
@@ -30,8 +30,6 @@ or inference services — all traffic flows through orchestration.
|
|||||||
| LLAMA_SERVER_URL | No | http://localhost:8080 | Direct llama-server URL for /models/props |
|
| LLAMA_SERVER_URL | No | http://localhost:8080 | Direct llama-server URL for /models/props |
|
||||||
| QDRANT_URL | No | http://localhost:6333 | Qdrant URL for semantic search |
|
| QDRANT_URL | No | http://localhost:6333 | Qdrant URL for semantic search |
|
||||||
| CORS_ORIGIN | No | http://localhost:5173 | Allowed origin for CORS requests |
|
| CORS_ORIGIN | No | http://localhost:5173 | Allowed origin for CORS requests |
|
||||||
| EXTRACTION_URL | No | http://localhost:11434 | Ollama URL for summarisation |
|
|
||||||
| EXTRACTION_MODEL | No | qwen2.5:3b | Ollama model used for summarisation |
|
|
||||||
|
|
||||||
## Internal Structure
|
## Internal Structure
|
||||||
|
|
||||||
@@ -75,6 +73,9 @@ via `appSettings.load()` — changes apply immediately without a service restart
|
|||||||
| `scoreThreshold` | 0.5 | Minimum similarity score for Qdrant semantic results |
|
| `scoreThreshold` | 0.5 | Minimum similarity score for Qdrant semantic results |
|
||||||
| `semanticWeight` | 1.0 | RRF weight for Qdrant semantic results |
|
| `semanticWeight` | 1.0 | RRF weight for Qdrant semantic results |
|
||||||
| `keywordWeight` | 0 | RRF weight for FTS5 keyword results (`0` = disabled) |
|
| `keywordWeight` | 0 | RRF weight for FTS5 keyword results (`0` = disabled) |
|
||||||
|
| `contextBudget` | — | Token budget for context assembly (char/4 estimation on stored text) |
|
||||||
|
| `entityWeight` | — | Scoring bonus for entity-linked episodes in the context pool |
|
||||||
|
| `minRecentEpisodes` | — | Guaranteed floor of recent episodes always included in context |
|
||||||
| `modelsFolderPath` | `/mnt/nexus-models` | Path to folder containing .gguf files |
|
| `modelsFolderPath` | `/mnt/nexus-models` | Path to folder containing .gguf files |
|
||||||
| `temperature` | 0.7 | Inference temperature |
|
| `temperature` | 0.7 | Inference temperature |
|
||||||
| `repeatPenalty` | 1.1 | Repeat token penalty |
|
| `repeatPenalty` | 1.1 | Repeat token penalty |
|
||||||
@@ -103,35 +104,43 @@ difference is how the inference response is delivered to the client.
|
|||||||
|
|
||||||
4. **Recent episode retrieval** — fetch most recent episodes (`recentEpisodeLimit`).
|
4. **Recent episode retrieval** — fetch most recent episodes (`recentEpisodeLimit`).
|
||||||
|
|
||||||
5. **Fused episode retrieval** — runs semantic (Qdrant) and keyword (FTS5)
|
5. **Trivial-turn gate** — greetings/pleasantries (`isTrivialTurn`) skip all
|
||||||
|
retrieval (semantic, keyword, entity); recent history alone is the context.
|
||||||
|
Breaks the greeting → marginal-retrieval → confabulation loop.
|
||||||
|
|
||||||
|
6. **Fused episode retrieval** — runs semantic (Qdrant) and keyword (FTS5)
|
||||||
search in parallel, then merges results via Reciprocal Rank Fusion (RRF).
|
search in parallel, then merges results via Reciprocal Rank Fusion (RRF).
|
||||||
Both paths are filtered against `recentIds` before fusion. FTS is scoped
|
The query is embedded once and shared with entity search. Both paths are
|
||||||
to the current session or all project sessions. If `keywordWeight` is `0`,
|
filtered against `recentIds` before fusion. FTS is scoped to the current
|
||||||
the FTS call is skipped entirely. Non-critical — failures fall back to
|
session or all project sessions. If `keywordWeight` is `0`, the FTS call
|
||||||
whichever strategy succeeded.
|
is skipped entirely. Non-critical — failures fall back to whichever
|
||||||
|
strategy succeeded.
|
||||||
|
|
||||||
6. **Entity search** — query `entities` Qdrant collection filtered by
|
7. **Entity search + graph expansion** — query `entities` Qdrant collection
|
||||||
`projectId`. Returns entity IDs alongside Qdrant payload data (the Qdrant
|
(project-scoped, or session-scoped via `/sessions/:id/entity-ids` for
|
||||||
point ID equals the SQLite entity ID). Non-critical.
|
non-project chats). Entity IDs are expanded into a 1-hop subgraph via
|
||||||
|
`POST /graph/neighbors`; on failure, falls back to flat entity list.
|
||||||
|
Non-critical.
|
||||||
|
|
||||||
7. **Graph neighborhood expansion** — call `POST /graph/neighbors` on
|
8. **Scored pool + budget selection** — `buildScoredPool` combines RRF scores,
|
||||||
memory-service with the entity IDs from step 6. Returns a 1-hop subgraph
|
recency, and entity-linkage bonus; `selectWithinBudget` fills `contextBudget`
|
||||||
`{ nodes, edges }` — entity objects plus the relationships connecting them.
|
(char/4 token estimation on stored text) above a guaranteed floor of
|
||||||
If no entities were found or the graph call fails, falls back to flat entity
|
`minRecentEpisodes` recent episodes. Selected episode IDs are then reported
|
||||||
list (no edges). Non-critical.
|
to `POST /episodes/touch` fire-and-forget (access tracking for the
|
||||||
|
consolidation lifecycle).
|
||||||
|
|
||||||
8. **Prompt assembly** — combine system prompt, graph context, fused episodes,
|
9. **Prompt assembly** — combine system prompt, graph context, selected
|
||||||
recent episodes, and user message.
|
episodes, guaranteed recent episodes, and user message.
|
||||||
|
|
||||||
9. **Inference** — send to inference service. `/chat` awaits full response;
|
10. **Inference** — send to inference service. `/chat` awaits full response;
|
||||||
`/chat/stream` pipes SSE chunks to the client.
|
`/chat/stream` pipes SSE chunks to the client.
|
||||||
|
|
||||||
10. **Episode write** — write exchange back to memory with `projectId`.
|
11. **Episode write** — write exchange back to memory with `projectId`.
|
||||||
|
|
||||||
11. **Summarisation trigger** — `triggerSummary(session, allEpisodes)` called
|
12. **Summarisation trigger** — `triggerSummary(session)` called
|
||||||
fire-and-forget. See `summarization.md` for full details.
|
fire-and-forget. See `summarization.md` for full details.
|
||||||
|
|
||||||
12. **Auto-naming** — on first message with no session name, fires a secondary
|
13. **Auto-naming** — on first message with no session name, fires a secondary
|
||||||
inference call (max 20 tokens, temperature 0.3) to generate a session name.
|
inference call (max 20 tokens, temperature 0.3) to generate a session name.
|
||||||
|
|
||||||
### Prompt Structure
|
### Prompt Structure
|
||||||
|
|||||||
@@ -37,9 +37,10 @@ Fusion lives in orchestration — the service already coordinates multiple data
|
|||||||
sources, and fusion is a retrieval strategy, not a storage concern.
|
sources, and fusion is a retrieval strategy, not a storage concern.
|
||||||
|
|
||||||
```
|
```
|
||||||
getFusedEpisodes()
|
getFusedEpisodes(…, queryVector, …)
|
||||||
├── getSemanticEpisodes() — Qdrant embed+search → fetch full rows by ID
|
├── getSemanticEpisodes(queryVector) — Qdrant search → fetch full rows by ID
|
||||||
│ (existing path, unchanged)
|
│ (query embedded ONCE upstream in assembleContext and shared with entity
|
||||||
|
│ search — no longer embedded separately here)
|
||||||
└── getFTSResults() — memory-service /episodes/search → full rows directly
|
└── getFTSResults() — memory-service /episodes/search → full rows directly
|
||||||
(skipped entirely if keywordWeight == 0)
|
(skipped entirely if keywordWeight == 0)
|
||||||
↓
|
↓
|
||||||
@@ -48,6 +49,10 @@ fuseEpisodeResults() — pure RRF, no I/O
|
|||||||
fusedEpisodes[] — top semanticLimit episodes by RRF score
|
fusedEpisodes[] — top semanticLimit episodes by RRF score
|
||||||
```
|
```
|
||||||
|
|
||||||
|
The query embedding is computed once per turn in `assembleContext` and passed
|
||||||
|
into both fused retrieval and entity search; if embedding fails, both receive
|
||||||
|
`null` and degrade to empty results rather than erroring.
|
||||||
|
|
||||||
### Data Shape Consistency
|
### Data Shape Consistency
|
||||||
|
|
||||||
Both sides must enter fusion as `Episode[]` — full SQLite row objects with
|
Both sides must enter fusion as `Episode[]` — full SQLite row objects with
|
||||||
@@ -59,6 +64,38 @@ the same shape — and both must be filtered against `recentIds` first:
|
|||||||
FTS requests `semanticLimit * 2` results to provide headroom for the
|
FTS requests `semanticLimit * 2` results to provide headroom for the
|
||||||
`recentIds` filter without under-serving the fusion.
|
`recentIds` filter without under-serving the fusion.
|
||||||
|
|
||||||
|
## Query Tokenization
|
||||||
|
|
||||||
|
Before FTS5 sees the query, `buildFtsQuery(query)` (in
|
||||||
|
`memory-service/src/episodic/index.js`) turns the raw message into a MATCH
|
||||||
|
expression:
|
||||||
|
|
||||||
|
1. Lowercase and split on any non-letter/number (`/[^\p{L}\p{N}]+/u`, unicode-aware)
|
||||||
|
2. Drop stopwords (a small `FTS_STOPWORDS` set of common function words) and single-character tokens
|
||||||
|
3. Wrap each surviving token in double quotes and join with ` OR `
|
||||||
|
|
||||||
|
So `"How do I configure the Qdrant collection?"` becomes
|
||||||
|
`"configure" OR "qdrant" OR "collection"`. If nothing survives (an
|
||||||
|
all-stopword message like `"how do I do it?"`), it returns `null` and
|
||||||
|
`searchEpisodes` returns `[]` — keyword search sits out that turn and
|
||||||
|
semantic retrieval carries it.
|
||||||
|
|
||||||
|
**Why this matters:** the earlier implementation quoted the *entire* message
|
||||||
|
as one FTS5 phrase, which required the whole string to appear verbatim in an
|
||||||
|
episode — so keyword recall was effectively nil for conversational queries.
|
||||||
|
Tokenizing into OR-joined terms is what makes `keywordWeight > 0` actually
|
||||||
|
contribute anything.
|
||||||
|
|
||||||
|
**Injection safety:** quoting each token individually means any
|
||||||
|
FTS5-significant token inside the user's message (a literal `OR`, `*`, `"`,
|
||||||
|
etc.) is matched as a search term rather than parsed as an operator. This
|
||||||
|
replaces the safety the old whole-phrase quoting provided.
|
||||||
|
|
||||||
|
The stopword set is deliberately conservative and tuned iteratively — high
|
||||||
|
frequency filler (`the`, `is`, `one`, `there`, …) is dropped, but borderline
|
||||||
|
words that can carry signal (`time`, `good`, `way`) are kept. Add to the set
|
||||||
|
when a common word is observed producing noisy matches.
|
||||||
|
|
||||||
## FTS Session Scoping
|
## FTS Session Scoping
|
||||||
|
|
||||||
Without scoping, FTS5 searches across all episodes in the database. For
|
Without scoping, FTS5 searches across all episodes in the database. For
|
||||||
|
|||||||
@@ -203,6 +203,16 @@ SUMMARY_MAX_TOKENS=800
|
|||||||
SUMMARY_MIN_EPISODES=5
|
SUMMARY_MIN_EPISODES=5
|
||||||
```
|
```
|
||||||
|
|
||||||
|
#### `CONSOLIDATION`
|
||||||
|
|
||||||
|
Controls the memory consolidation lifecycle (currently dry-run only).
|
||||||
|
|
||||||
|
| Key | Value | Description |
|
||||||
|
|---|---|---|
|
||||||
|
| `MIN_AGE_DAYS` | `7` | Episodes younger than this are never consolidation candidates |
|
||||||
|
| `MIN_SESSION_EPISODES` | `20` | Sessions with fewer episodes are skipped entirely |
|
||||||
|
| `CANDIDATE_LIMIT` | `50` | Max candidates returned per scoring query |
|
||||||
|
|
||||||
#### `SQLITE`
|
#### `SQLITE`
|
||||||
|
|
||||||
| Key | Value | Description |
|
| Key | Value | Description |
|
||||||
|
|||||||
@@ -6,13 +6,27 @@ the full context window with raw episodes.
|
|||||||
|
|
||||||
**Location:** `packages/orchestration-service/src/services/summarization.js`
|
**Location:** `packages/orchestration-service/src/services/summarization.js`
|
||||||
**Triggered by:** `chat/index.js` after every episode write (fire-and-forget)
|
**Triggered by:** `chat/index.js` after every episode write (fire-and-forget)
|
||||||
**Model:** `qwen2.5:3b` via Ollama on Mini PC 1 (192.168.0.81)
|
**Model:** the utility model served by the inference service (`/utility/complete`), set via `UTILITY_MODEL` (backed by Ollama on Mini PC 1, 192.168.0.81)
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
## Trigger Conditions
|
## Trigger Conditions
|
||||||
|
|
||||||
`triggerSummary(session, allEpisodes)` calls `maybeSummarize` fire-and-forget.
|
`triggerSummary(session)` calls `maybeSummarize` fire-and-forget. It takes only
|
||||||
|
the `session` — it no longer receives the full episode list. Instead
|
||||||
|
`maybeSummarize` fetches exactly what it needs:
|
||||||
|
|
||||||
|
1. `GET /sessions/:id/episode-stats` — a cheap aggregate (`COUNT`, `SUM(token_count)`,
|
||||||
|
`MAX(id)`) that gates the token threshold **without** pulling every episode row
|
||||||
|
2. Only if over threshold: `GET /sessions/:id/episodes/since/:afterId` — the
|
||||||
|
un-summarized tail (episodes newer than the last summary's range), fetched in
|
||||||
|
full for the summary prompt
|
||||||
|
|
||||||
|
This replaced an earlier approach where `chat/index.js` fetched the entire
|
||||||
|
session (`getRecentEpisodes(session.id, 9999)`) on every message just to hand it
|
||||||
|
over — an O(session length) cost per turn. The stats query now runs on every
|
||||||
|
message; full episode text is fetched only when a summary actually fires.
|
||||||
|
|
||||||
`maybeSummarize` proceeds only when both conditions are met:
|
`maybeSummarize` proceeds only when both conditions are met:
|
||||||
|
|
||||||
1. Total session token count exceeds `SUMMARIES.THRESHOLD_TOKENS` (default 200)
|
1. Total session token count exceeds `SUMMARIES.THRESHOLD_TOKENS` (default 200)
|
||||||
@@ -42,47 +56,50 @@ not all episodes in the session.
|
|||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
## Ollama Request
|
## Utility Inference Request
|
||||||
|
|
||||||
|
Summaries are generated through the shared `utilityInference()` helper, which
|
||||||
|
POSTs to the inference service's `/utility/complete` endpoint. `buildSummaryPrompt`
|
||||||
|
returns a plain instruction string (no template tags) passed as the `user` message:
|
||||||
|
|
||||||
```js
|
```js
|
||||||
{
|
const content = await utilityInference({
|
||||||
model: EXTRACTION_MODEL, // qwen2.5:3b (set via EXTRACTION_MODEL env var)
|
user: buildSummaryPrompt(episodesToSummarize, existingSummary),
|
||||||
prompt: buildSummaryPrompt(episodesToSummarize, existingSummary),
|
temperature: SUMMARIES.TEMPERATURE, // 0.2
|
||||||
stream: false,
|
maxTokens: SUMMARIES.SESSION_GEN_MAX_TOKENS, // 500
|
||||||
// No format: 'json' — free-text output required for summaries
|
});
|
||||||
options: {
|
|
||||||
temperature: 0.2,
|
|
||||||
num_predict: 500,
|
|
||||||
},
|
|
||||||
}
|
|
||||||
```
|
```
|
||||||
|
|
||||||
`temperature: 0.2` is slightly higher than extraction (0.1) — summaries
|
`TEMPERATURE` (0.2) is slightly higher than extraction (0.1) — summaries benefit
|
||||||
benefit from some fluency. `num_predict: 500` gives room for 5 thorough
|
from some fluency. `SESSION_GEN_MAX_TOKENS` (500) gives room for ~5 thorough
|
||||||
sentences without risk of runoff.
|
sentences without runoff. Both live in `@nexusai/shared` `SUMMARIES` constants.
|
||||||
|
|
||||||
|
There is no `json: true` here — summaries are free-text, unlike entity extraction.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
## Prompt Format
|
## Prompt Format
|
||||||
|
|
||||||
ChatML format — native to qwen2.5:
|
The prompt is plain text describing the task; the model's own prompt template
|
||||||
|
(ChatML for qwen, etc.) is applied **server-side** by the inference service via
|
||||||
|
Ollama's `/api/chat`. No `<|im_start|>` tags belong in this codebase, and the
|
||||||
|
utility model can be swapped (via `UTILITY_MODEL` on the inference service) with
|
||||||
|
no prompt changes here.
|
||||||
|
|
||||||
|
Fresh summary instruction:
|
||||||
|
|
||||||
```
|
```
|
||||||
<|im_start|>user
|
|
||||||
Summarize the conversation below in 3-5 sentences.
|
Summarize the conversation below in 3-5 sentences.
|
||||||
Write in third person. Do not quote directly — paraphrase only.
|
Write in third person. Do not quote directly — paraphrase only.
|
||||||
Do not include greetings, sign-offs, or filler. Output only the summary text.
|
Do not include greetings, sign-offs, or filler. Output only the summary text.
|
||||||
|
|
||||||
Conversation:
|
Conversation:
|
||||||
{context}
|
{context}
|
||||||
<|im_end|>
|
|
||||||
<|im_start|>assistant
|
|
||||||
```
|
```
|
||||||
|
|
||||||
For cumulative updates, the instruction and context change:
|
Cumulative update instruction:
|
||||||
|
|
||||||
```
|
```
|
||||||
<|im_start|>user
|
|
||||||
Update the summary below to incorporate the new exchanges.
|
Update the summary below to incorporate the new exchanges.
|
||||||
Write 3-5 sentences in third person. Do not quote directly — paraphrase only.
|
Write 3-5 sentences in third person. Do not quote directly — paraphrase only.
|
||||||
Do not include greetings, sign-offs, or filler. Output only the updated summary text.
|
Do not include greetings, sign-offs, or filler. Output only the updated summary text.
|
||||||
@@ -92,35 +109,22 @@ Previous summary:
|
|||||||
|
|
||||||
New exchanges:
|
New exchanges:
|
||||||
{context}
|
{context}
|
||||||
<|im_end|>
|
|
||||||
<|im_start|>assistant
|
|
||||||
```
|
```
|
||||||
|
|
||||||
### Input truncation
|
### Input truncation
|
||||||
|
|
||||||
Episode context is truncated to `MAX_CHARS = 3000` characters, keeping the
|
Episode context is truncated to `MAX_CHARS = 3000` characters, keeping the
|
||||||
most recent exchanges (sliced from the end). This keeps Qwen focused and
|
most recent exchanges (sliced from the end). This keeps the model focused and
|
||||||
prevents the prompt from exceeding its effective context window.
|
prevents the prompt from exceeding its effective context window.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
## ChatML Token Stripping
|
## Output Handling
|
||||||
|
|
||||||
Qwen occasionally echoes ChatML tokens back into its response. The raw output
|
Because `/api/chat` applies and removes the prompt template server-side, the
|
||||||
is cleaned before saving:
|
returned text is already clean — the previous ChatML token-stripping step (and
|
||||||
|
the class of bug where leaked tokens got stored and re-injected into the next
|
||||||
```js
|
summarisation prompt) no longer applies.
|
||||||
const raw = data.response?.trim() ?? '';
|
|
||||||
const content = raw
|
|
||||||
.replace(/<\|im_start\|>.*?<\|im_end\|>/gs, '')
|
|
||||||
.replace(/<\|im_start\|>|<\|im_end\|>|<\|im_sep\|>/g, '')
|
|
||||||
.trim();
|
|
||||||
return content;
|
|
||||||
```
|
|
||||||
|
|
||||||
Without this, leaked tokens get stored in the summary and then injected
|
|
||||||
back into the next summarisation prompt — causing the model to append a new
|
|
||||||
summary after the old one rather than replacing it.
|
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
@@ -193,8 +197,7 @@ Set in `packages/orchestration-service/src/.env`:
|
|||||||
|
|
||||||
| Variable | Default | Description |
|
| Variable | Default | Description |
|
||||||
|---|---|---|
|
|---|---|---|
|
||||||
| `EXTRACTION_URL` | `http://localhost:11434` | Ollama instance URL |
|
| `INFERENCE_SERVICE_URL` | `http://localhost:3001` | Inference service — summaries route through its `/utility/complete` endpoint (model set via `UTILITY_MODEL` there) |
|
||||||
| `EXTRACTION_MODEL` | `qwen2.5:3b` | Model for summarisation |
|
|
||||||
| `MEMORY_SERVICE_URL` | `http://localhost:3002` | Memory service URL |
|
| `MEMORY_SERVICE_URL` | `http://localhost:3002` | Memory service URL |
|
||||||
| `SUMMARY_THRESHOLD_TOKENS` | `200` | Token threshold before summarisation triggers |
|
| `SUMMARY_THRESHOLD_TOKENS` | `200` | Token threshold before summarisation triggers |
|
||||||
| `SUMMARY_MAX_TOKENS` | `800` | Max summary length before a new row is created |
|
| `SUMMARY_MAX_TOKENS` | `800` | Max summary length before a new row is created |
|
||||||
|
|||||||
@@ -2,6 +2,7 @@ require ('dotenv').config();
|
|||||||
const express = require('express');
|
const express = require('express');
|
||||||
const {getEnv, PORTS, OLLAMA, logger} = require('@nexusai/shared');
|
const {getEnv, PORTS, OLLAMA, logger} = require('@nexusai/shared');
|
||||||
const inferenceRouter = require('./routes/inference');
|
const inferenceRouter = require('./routes/inference');
|
||||||
|
const utilityRouter = require('./routes/utility')
|
||||||
|
|
||||||
const app = express();
|
const app = express();
|
||||||
app.use(express.json({ limit: '8mb' })); // prompts include full context window
|
app.use(express.json({ limit: '8mb' })); // prompts include full context window
|
||||||
@@ -20,6 +21,7 @@ app.get('/health', (req, res) => {
|
|||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
|
app.use('/', utilityRouter)
|
||||||
app.use('/', inferenceRouter);
|
app.use('/', inferenceRouter);
|
||||||
|
|
||||||
// Start the server
|
// Start the server
|
||||||
|
|||||||
@@ -0,0 +1,21 @@
|
|||||||
|
const { Router } = require('express');
|
||||||
|
const { logger } = require('@nexusai/shared');
|
||||||
|
const { utilityComplete } = require('../utility');
|
||||||
|
|
||||||
|
const router = Router();
|
||||||
|
|
||||||
|
router.post('/utility/complete', async (req, res) => {
|
||||||
|
const { system, user, json, temperature, maxTokens } = req.body;
|
||||||
|
|
||||||
|
if (!user) return res.status(400).json({ error: 'user message is required' });
|
||||||
|
|
||||||
|
try {
|
||||||
|
const result = await utilityComplete({ system, user, json, temperature, maxTokens });
|
||||||
|
res.json(result);
|
||||||
|
} catch (error) {
|
||||||
|
logger.error('[Utility] Completion error:', error.message);
|
||||||
|
res.status(500).json({ error: 'Utility inference failed', detail: error.message });
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
module.exports = router;
|
||||||
@@ -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,6 +1,7 @@
|
|||||||
const Database = require('better-sqlite3');
|
const Database = require('better-sqlite3');
|
||||||
const schema = require('./schema');
|
const schema = require('./schema');
|
||||||
const {getEnv, SQLITE, logger } = require('@nexusai/shared');
|
const { migrate } = require('./migrations');
|
||||||
|
const { getEnv, SQLITE, logger } = require('@nexusai/shared');
|
||||||
|
|
||||||
let db; // Declare db variable in a scope accessible to all functions
|
let db; // Declare db variable in a scope accessible to all functions
|
||||||
|
|
||||||
@@ -12,60 +13,29 @@ function getDB() {
|
|||||||
db.pragma('journal_mode = WAL');
|
db.pragma('journal_mode = WAL');
|
||||||
db.pragma('foreign_keys = ON');
|
db.pragma('foreign_keys = ON');
|
||||||
|
|
||||||
db.exec(schema);
|
// Was the FTS index absent before this boot? (True for a fresh DB, and for
|
||||||
|
// an older DB from before FTS existed.) Checked BEFORE schema runs so we can
|
||||||
|
// decide whether a one-time backfill is needed below.
|
||||||
|
const ftsExisted = db.prepare(
|
||||||
|
`SELECT 1 FROM sqlite_master WHERE type='table' AND name='episodes_fts'`
|
||||||
|
).get() !== undefined;
|
||||||
|
|
||||||
try{
|
db.exec(schema); // complete current shape — fresh DBs get everything
|
||||||
db.exec(`ALTER TABLE sessions ADD COLUMN name TEXT`)
|
migrate(db); // carry an older DB forward; no-op on fresh/current DBs
|
||||||
} catch {}
|
|
||||||
|
|
||||||
try {
|
// One-time FTS backfill: only when the index was just created on a DB that
|
||||||
db.exec(`ALTER TABLE sessions ADD COLUMN project_id INTEGER REFERENCES projects(id)`);
|
// already holds episodes (i.e. episodes predate FTS). During normal
|
||||||
} catch {}
|
// operation the insert/delete/update triggers keep it in sync, so this no
|
||||||
|
// longer rebuilds the whole index on every boot. NOTE: COUNT(*) on an
|
||||||
try {
|
// external-content FTS5 table proxies the content table, so it can't detect
|
||||||
db.exec(`CREATE INDEX IF NOT EXISTS idx_sessions_project ON sessions(project_id)`);
|
// a desync — the "was it just created" check is what makes this correct.
|
||||||
} catch {}
|
if (!ftsExisted) {
|
||||||
|
const epCount = db.prepare('SELECT COUNT(*) AS c FROM episodes').get().c;
|
||||||
try {
|
if (epCount > 0) {
|
||||||
db.exec(`ALTER TABLE projects ADD COLUMN isolated INTEGER NOT NULL DEFAULT 0`);
|
db.exec(`INSERT INTO episodes_fts(episodes_fts) VALUES('rebuild')`);
|
||||||
} catch {}
|
logger.info(`[db] Backfilled FTS index for ${epCount} pre-existing episodes`);
|
||||||
|
}
|
||||||
try {
|
}
|
||||||
db.exec(`ALTER TABLE projects ADD COLUMN notes TEXT`); // ← add this
|
|
||||||
} catch {}
|
|
||||||
|
|
||||||
try {
|
|
||||||
db.exec(`ALTER TABLE projects ADD COLUMN system_prompt TEXT`);
|
|
||||||
} catch {}
|
|
||||||
|
|
||||||
try {
|
|
||||||
db.exec(`ALTER TABLE summaries ADD COLUMN project_id INTEGER REFERENCES projects(id) ON DELETE CASCADE`);
|
|
||||||
} catch {}
|
|
||||||
|
|
||||||
try {
|
|
||||||
db.exec(`ALTER TABLE summaries ADD COLUMN token_count INTEGER`);
|
|
||||||
} catch {}
|
|
||||||
|
|
||||||
try {
|
|
||||||
db.exec(`CREATE INDEX IF NOT EXISTS idx_summaries_project ON summaries(project_id)`);
|
|
||||||
} catch {}
|
|
||||||
|
|
||||||
try {
|
|
||||||
db.exec(`CREATE INDEX IF NOT EXISTS idx_summaries_session ON summaries(session_id)`);
|
|
||||||
} catch {}
|
|
||||||
|
|
||||||
try { db.exec(`ALTER TABLE entities ADD COLUMN mention_count INTEGER NOT NULL DEFAULT 1`) } catch {}
|
|
||||||
try { db.exec(`ALTER TABLE entities ADD COLUMN confidence REAL NOT NULL DEFAULT 1.0`) } catch {}
|
|
||||||
try { db.exec(`ALTER TABLE entities ADD COLUMN source TEXT NOT NULL DEFAULT 'extraction'`) } catch {}
|
|
||||||
try { db.exec(`ALTER TABLE entities ADD COLUMN last_seen_at INTEGER`) } catch {}
|
|
||||||
|
|
||||||
try { db.exec(`ALTER TABLE relationships ADD COLUMN mention_count INTEGER NOT NULL DEFAULT 1`) } catch {}
|
|
||||||
try { db.exec(`ALTER TABLE relationships ADD COLUMN notes TEXT`) } catch {}
|
|
||||||
|
|
||||||
|
|
||||||
// Sync FTS index with any existing episodes data
|
|
||||||
db.exec(`INSERT OR REPLACE INTO episodes_fts(rowid, user_message, ai_response)
|
|
||||||
SELECT id, user_message, ai_response FROM episodes`);
|
|
||||||
|
|
||||||
logger.info(`Connected to SQLite database at ${path}`);
|
logger.info(`Connected to SQLite database at ${path}`);
|
||||||
}
|
}
|
||||||
@@ -74,4 +44,4 @@ function getDB() {
|
|||||||
|
|
||||||
module.exports = {
|
module.exports = {
|
||||||
getDB
|
getDB
|
||||||
};
|
};
|
||||||
|
|||||||
@@ -0,0 +1,51 @@
|
|||||||
|
const { logger } = require('@nexusai/shared');
|
||||||
|
|
||||||
|
// Forward-only schema migrations. Entry i takes the database from user_version i
|
||||||
|
// to i+1, so migrations[0] is the v0→v1 step, migrations[1] the v1→v2 step, etc.
|
||||||
|
//
|
||||||
|
// schema.js already holds the COMPLETE current shape, so a fresh database is
|
||||||
|
// created whole and stamped straight to LATEST_VERSION — these run only to carry
|
||||||
|
// an OLDER database forward. Add a new schema change by appending a function here
|
||||||
|
// (which bumps LATEST_VERSION by one); never edit an existing entry once shipped,
|
||||||
|
// since databases already stamped past it will not re-run it.
|
||||||
|
//
|
||||||
|
// Each migration receives the better-sqlite3 db handle and runs inside a
|
||||||
|
// transaction together with its version bump, so a failure rolls back cleanly.
|
||||||
|
const migrations = [
|
||||||
|
// v0 → v1: baseline. The historical ALTER TABLE / CREATE INDEX statements that
|
||||||
|
// used to run (wrapped in try/catch) on every boot are folded into schema.js.
|
||||||
|
// Nothing to do here — this entry exists to mark v1 as the consolidated baseline.
|
||||||
|
(_db) => {},
|
||||||
|
|
||||||
|
(db) => {
|
||||||
|
db.exec(`
|
||||||
|
ALTER TABLE episodes ADD COLUMN last_accessed_at INTEGER;
|
||||||
|
ALTER TABLE episodes ADD COLUMN access_count INTEGER NOT NULL DEFAULT 0;
|
||||||
|
`)
|
||||||
|
db.exec(`UPDATE episodes SET last_accessed_at = created_at WHERE last_accessed_at IS NULL`)
|
||||||
|
}
|
||||||
|
];
|
||||||
|
|
||||||
|
const LATEST_VERSION = migrations.length;
|
||||||
|
|
||||||
|
// Applies any migrations newer than the database's current user_version, one at a
|
||||||
|
// time, each with its version bump, inside a transaction. Returns the resulting
|
||||||
|
// version. A no-op when the database is already current. The migrations list and
|
||||||
|
// target version are injectable for testing; production callers pass just `db`.
|
||||||
|
function migrate(db, migs = migrations, latest = migs.length) {
|
||||||
|
const current = db.pragma('user_version', { simple: true });
|
||||||
|
if (current >= latest) return current;
|
||||||
|
|
||||||
|
for (let v = current; v < latest; v++) {
|
||||||
|
const step = db.transaction(() => {
|
||||||
|
migs[v](db);
|
||||||
|
db.pragma(`user_version = ${v + 1}`);
|
||||||
|
});
|
||||||
|
step();
|
||||||
|
}
|
||||||
|
|
||||||
|
logger.info(`[db] schema migrated ${current} -> ${latest}`);
|
||||||
|
return latest;
|
||||||
|
}
|
||||||
|
|
||||||
|
module.exports = { migrate, migrations, LATEST_VERSION };
|
||||||
@@ -1,40 +1,55 @@
|
|||||||
|
// Complete current schema — the single source of truth for a fresh database.
|
||||||
|
// Historical ALTER TABLE statements that used to run on every boot are folded in
|
||||||
|
// here. Any change that must reach EXISTING databases goes in migrations.js as a
|
||||||
|
// new numbered migration, NOT by editing a table below (CREATE ... IF NOT EXISTS
|
||||||
|
// silently skips tables that already exist, so column edits here never reach them).
|
||||||
const schema = `
|
const schema = `
|
||||||
CREATE TABLE IF NOT EXISTS sessions (
|
CREATE TABLE IF NOT EXISTS sessions (
|
||||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||||
external_id TEXT UNIQUE NOT NULL,
|
external_id TEXT UNIQUE NOT NULL,
|
||||||
created_at INTEGER NOT NULL DEFAULT (unixepoch()),
|
created_at INTEGER NOT NULL DEFAULT (unixepoch()),
|
||||||
updated_at INTEGER NOT NULL DEFAULT (unixepoch()),
|
updated_at INTEGER NOT NULL DEFAULT (unixepoch()),
|
||||||
metadata TEXT
|
metadata TEXT,
|
||||||
|
name TEXT,
|
||||||
|
project_id INTEGER REFERENCES projects(id)
|
||||||
);
|
);
|
||||||
|
|
||||||
CREATE TABLE IF NOT EXISTS episodes (
|
CREATE TABLE IF NOT EXISTS episodes (
|
||||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||||
session_id INTEGER NOT NULL REFERENCES sessions(id) ON DELETE CASCADE,
|
session_id INTEGER NOT NULL REFERENCES sessions(id) ON DELETE CASCADE,
|
||||||
user_message TEXT NOT NULL,
|
user_message TEXT NOT NULL,
|
||||||
ai_response TEXT NOT NULL,
|
ai_response TEXT NOT NULL,
|
||||||
created_at INTEGER NOT NULL DEFAULT (unixepoch()),
|
created_at INTEGER NOT NULL DEFAULT (unixepoch()),
|
||||||
token_count INTEGER,
|
token_count INTEGER,
|
||||||
metadata TEXT
|
metadata TEXT,
|
||||||
|
last_accessed_at INTEGER,
|
||||||
|
access_count INTEGER NOT NULL DEFAULT 0
|
||||||
);
|
);
|
||||||
|
|
||||||
CREATE TABLE IF NOT EXISTS entities (
|
CREATE TABLE IF NOT EXISTS entities (
|
||||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||||
name TEXT NOT NULL,
|
name TEXT NOT NULL,
|
||||||
type TEXT NOT NULL,
|
type TEXT NOT NULL,
|
||||||
notes TEXT,
|
notes TEXT,
|
||||||
created_at INTEGER NOT NULL DEFAULT (unixepoch()),
|
created_at INTEGER NOT NULL DEFAULT (unixepoch()),
|
||||||
updated_at INTEGER NOT NULL DEFAULT (unixepoch()),
|
updated_at INTEGER NOT NULL DEFAULT (unixepoch()),
|
||||||
metadata TEXT,
|
metadata TEXT,
|
||||||
|
mention_count INTEGER NOT NULL DEFAULT 1,
|
||||||
|
confidence REAL NOT NULL DEFAULT 1.0,
|
||||||
|
source TEXT NOT NULL DEFAULT 'extraction',
|
||||||
|
last_seen_at INTEGER,
|
||||||
UNIQUE(name, type)
|
UNIQUE(name, type)
|
||||||
);
|
);
|
||||||
|
|
||||||
CREATE TABLE IF NOT EXISTS relationships (
|
CREATE TABLE IF NOT EXISTS relationships (
|
||||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||||
from_id INTEGER NOT NULL REFERENCES entities(id) ON DELETE CASCADE,
|
from_id INTEGER NOT NULL REFERENCES entities(id) ON DELETE CASCADE,
|
||||||
to_id INTEGER NOT NULL REFERENCES entities(id) ON DELETE CASCADE,
|
to_id INTEGER NOT NULL REFERENCES entities(id) ON DELETE CASCADE,
|
||||||
label TEXT NOT NULL,
|
label TEXT NOT NULL,
|
||||||
created_at INTEGER NOT NULL DEFAULT (unixepoch()),
|
created_at INTEGER NOT NULL DEFAULT (unixepoch()),
|
||||||
metadata TEXT,
|
metadata TEXT,
|
||||||
|
mention_count INTEGER NOT NULL DEFAULT 1,
|
||||||
|
notes TEXT,
|
||||||
UNIQUE(from_id, to_id, label)
|
UNIQUE(from_id, to_id, label)
|
||||||
);
|
);
|
||||||
|
|
||||||
@@ -49,16 +64,17 @@ const schema = `
|
|||||||
|
|
||||||
CREATE INDEX IF NOT EXISTS idx_entity_episodes_entity ON entity_episodes(entity_id);
|
CREATE INDEX IF NOT EXISTS idx_entity_episodes_entity ON entity_episodes(entity_id);
|
||||||
CREATE INDEX IF NOT EXISTS idx_entity_episodes_episode ON entity_episodes(episode_id);
|
CREATE INDEX IF NOT EXISTS idx_entity_episodes_episode ON entity_episodes(episode_id);
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
CREATE TABLE IF NOT EXISTS projects (
|
CREATE TABLE IF NOT EXISTS projects (
|
||||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||||
name TEXT NOT NULL,
|
name TEXT NOT NULL,
|
||||||
description TEXT,
|
description TEXT,
|
||||||
colour TEXT,
|
colour TEXT,
|
||||||
icon TEXT,
|
icon TEXT,
|
||||||
created_at INTEGER NOT NULL DEFAULT (unixepoch())
|
created_at INTEGER NOT NULL DEFAULT (unixepoch()),
|
||||||
|
isolated INTEGER NOT NULL DEFAULT 0,
|
||||||
|
notes TEXT,
|
||||||
|
system_prompt TEXT
|
||||||
);
|
);
|
||||||
|
|
||||||
CREATE TABLE IF NOT EXISTS summaries (
|
CREATE TABLE IF NOT EXISTS summaries (
|
||||||
@@ -72,17 +88,17 @@ const schema = `
|
|||||||
metadata TEXT
|
metadata TEXT
|
||||||
);
|
);
|
||||||
|
|
||||||
CREATE INDEX IF NOT EXISTS idx_episodes_session
|
CREATE INDEX IF NOT EXISTS idx_episodes_session ON episodes(session_id);
|
||||||
ON episodes(session_id);
|
CREATE INDEX IF NOT EXISTS idx_episodes_created ON episodes(created_at);
|
||||||
CREATE INDEX IF NOT EXISTS idx_episodes_created
|
CREATE INDEX IF NOT EXISTS idx_entities_type ON entities(type);
|
||||||
ON episodes(created_at);
|
CREATE INDEX IF NOT EXISTS idx_sessions_project ON sessions(project_id);
|
||||||
CREATE INDEX IF NOT EXISTS idx_entities_type
|
CREATE INDEX IF NOT EXISTS idx_summaries_project ON summaries(project_id);
|
||||||
ON entities(type);
|
CREATE INDEX IF NOT EXISTS idx_summaries_session ON summaries(session_id);
|
||||||
|
|
||||||
CREATE VIRTUAL TABLE IF NOT EXISTS episodes_fts
|
CREATE VIRTUAL TABLE IF NOT EXISTS episodes_fts
|
||||||
USING fts5(user_message, ai_response, content=episodes, content_rowid=id);
|
USING fts5(user_message, ai_response, content=episodes, content_rowid=id);
|
||||||
|
|
||||||
CREATE TRIGGER IF NOT EXISTS episodes_fts_insert
|
CREATE TRIGGER IF NOT EXISTS episodes_fts_insert
|
||||||
AFTER INSERT ON episodes BEGIN
|
AFTER INSERT ON episodes BEGIN
|
||||||
INSERT INTO episodes_fts(rowid, user_message, ai_response)
|
INSERT INTO episodes_fts(rowid, user_message, ai_response)
|
||||||
VALUES (new.id, new.user_message, new.ai_response);
|
VALUES (new.id, new.user_message, new.ai_response);
|
||||||
@@ -101,8 +117,6 @@ const schema = `
|
|||||||
INSERT INTO episodes_fts(rowid, user_message, ai_response)
|
INSERT INTO episodes_fts(rowid, user_message, ai_response)
|
||||||
VALUES (new.id, new.user_message, new.ai_response);
|
VALUES (new.id, new.user_message, new.ai_response);
|
||||||
END;
|
END;
|
||||||
|
|
||||||
|
|
||||||
`;
|
`;
|
||||||
|
|
||||||
module.exports = schema;
|
module.exports = schema;
|
||||||
|
|||||||
@@ -1,9 +1,7 @@
|
|||||||
const semantic = require('../semantic')
|
const semantic = require('../semantic')
|
||||||
const { getEnv, SERVICES, formatEpisodeText, ENTITIES, logger } = require('@nexusai/shared');
|
const { getEnv, SERVICES, formatEpisodeText, ENTITIES, logger, utilityInference } = require('@nexusai/shared');
|
||||||
const { upsertEntity, upsertRelationship, linkEntityToEpisode } = require('./index');
|
const { upsertEntity, upsertRelationship, linkEntityToEpisode } = require('./index');
|
||||||
|
|
||||||
const EXTRACTION_URL = getEnv('EXTRACTION_URL', 'http://localhost:11434');
|
|
||||||
const EXTRACTION_MODEL = getEnv('EXTRACTION_MODEL', 'qwen2.5:3b'); // ChatML format — see buildExtractionPrompt
|
|
||||||
const EMBEDDING_SERVICE_URL = getEnv('EMBEDDING_SERVICE_URL', SERVICES.EMBEDDING_URL);
|
const EMBEDDING_SERVICE_URL = getEnv('EMBEDDING_SERVICE_URL', SERVICES.EMBEDDING_URL);
|
||||||
|
|
||||||
const ENTITY_TYPES = ENTITIES.TYPES;
|
const ENTITY_TYPES = ENTITIES.TYPES;
|
||||||
@@ -28,10 +26,9 @@ function mentionedIn(name, haystack) {
|
|||||||
return norm(haystack).includes(norm(name));
|
return norm(haystack).includes(norm(name));
|
||||||
}
|
}
|
||||||
|
|
||||||
// NOTE: This prompt uses ChatML format (<|im_start|> / <|im_end|> tags), which is
|
// Returns { system, user } for utilityInference. The model's prompt template
|
||||||
// specific to qwen-family models. If EXTRACTION_MODEL is changed to a Llama-family
|
// (ChatML for qwen, etc.) is applied by Ollama via the inference service's
|
||||||
// or other model, this format will need to change — most alternatives use either
|
// /utility/complete route — no template tags belong in this file.
|
||||||
// plain text or [INST] / <<SYS>> tags. Silent degradation is likely if mismatched.
|
|
||||||
function buildExtractionPrompt(userMessage, aiResponse, knownEntities = []) {
|
function buildExtractionPrompt(userMessage, aiResponse, knownEntities = []) {
|
||||||
const knownBlock = knownEntities.length > 0
|
const knownBlock = knownEntities.length > 0
|
||||||
? [
|
? [
|
||||||
@@ -41,33 +38,30 @@ function buildExtractionPrompt(userMessage, aiResponse, knownEntities = []) {
|
|||||||
].join('\n')
|
].join('\n')
|
||||||
: '';
|
: '';
|
||||||
|
|
||||||
return [
|
return {
|
||||||
'<|im_start|>system',
|
system: 'You are a named entity and relationship extractor. You output only valid JSON.',
|
||||||
'You are a named entity and relationship extractor. You output only valid JSON.',
|
user: [
|
||||||
'<|im_end|>',
|
'Read the conversation below and extract all named entities and the relationships between them.',
|
||||||
'<|im_start|>user',
|
`Entity types: ${ENTITY_TYPES.join(', ')}`,
|
||||||
'Read the conversation below and extract all named entities and the relationships between them.',
|
'Use "character" for any fictional, game, or media characters (e.g. characters from anime, games, books, TV shows, movies)',
|
||||||
`Entity types: ${ENTITY_TYPES.join(', ')}`,
|
'Use "person" only for real people',
|
||||||
'Use "character" for any fictional, game, or media characters (e.g. characters from anime, games, books, TV shows, movies)',
|
'For each entity provide:',
|
||||||
'Use "person" only for real people',
|
' "name": short proper noun only (max 4 words)',
|
||||||
'For each entity provide:',
|
' "type": one of the valid types',
|
||||||
' "name": short proper noun only (max 4 words)',
|
' "notes": one specific sentence about this entity based on the conversation',
|
||||||
' "type": one of the valid types',
|
'For relationships, use snake_case verb labels (e.g. works_on, manages, uses, knows, located_in, part_of, created_by).',
|
||||||
' "notes": one specific sentence about this entity based on the conversation',
|
'Only include relationships between entities you have listed above.',
|
||||||
'For relationships, use snake_case verb labels (e.g. works_on, manages, uses, knows, located_in, part_of, created_by).',
|
'The known-entities list below is ONLY for consistent spelling and types. Do NOT output an entity unless it actually appears in the conversation.',
|
||||||
'Only include relationships between entities you have listed above.',
|
'Return this exact JSON structure:',
|
||||||
'The known-entities list below is ONLY for consistent spelling and types. Do NOT output an entity unless it actually appears in the conversation.',
|
'{ "entities": [{"name": "...", "type": "...", "notes": "..."}], "relationships": [{"from": "...", "fromType": "...", "to": "...", "toType": "...", "label": "...", "notes": "..."}] }',
|
||||||
'Return this exact JSON structure:',
|
'',
|
||||||
'{ "entities": [{"name": "...", "type": "...", "notes": "..."}], "relationships": [{"from": "...", "fromType": "...", "to": "...", "toType": "...", "label": "...", "notes": "..."}] }',
|
knownBlock,
|
||||||
'',
|
'--- CONVERSATION ---',
|
||||||
knownBlock,
|
`User: ${userMessage}`,
|
||||||
'--- CONVERSATION ---',
|
`Assistant: ${aiResponse}`,
|
||||||
`User: ${userMessage}`,
|
'--- END CONVERSATION ---',
|
||||||
`Assistant: ${aiResponse}`,
|
].join('\n'),
|
||||||
'--- END CONVERSATION ---',
|
};
|
||||||
'<|im_end|>',
|
|
||||||
'<|im_start|>assistant',
|
|
||||||
].join('\n');
|
|
||||||
}
|
}
|
||||||
|
|
||||||
async function embedEntity(entity) {
|
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
|
// Fetch existing entities to guide the model toward consistent name/type pairs
|
||||||
const db = require('../db').getDB();
|
const db = require('../db').getDB();
|
||||||
const knownEntities = db.prepare(`SELECT name, type FROM entities ORDER BY rowid DESC LIMIT 20`).all();
|
const knownEntities = db.prepare(`SELECT name, type FROM entities ORDER BY rowid DESC LIMIT 20`).all();
|
||||||
const prompt = buildExtractionPrompt(userMessage, aiResponse, knownEntities);
|
const { system, user } = buildExtractionPrompt(userMessage, aiResponse, knownEntities);
|
||||||
|
|
||||||
|
const raw = await utilityInference({
|
||||||
const res = await fetch(`${EXTRACTION_URL}/api/generate`, {
|
system,
|
||||||
method: 'POST',
|
user,
|
||||||
headers: { 'Content-Type': 'application/json' },
|
json: true,
|
||||||
body: JSON.stringify({
|
temperature: ENTITIES.TEMPERATURE,
|
||||||
model: EXTRACTION_MODEL,
|
maxTokens: ENTITIES.NUM_PREDICT,
|
||||||
prompt: prompt,
|
|
||||||
stream: false,
|
|
||||||
format: 'json',
|
|
||||||
options: {
|
|
||||||
temperature: ENTITIES.TEMPERATURE,
|
|
||||||
num_predict: ENTITIES.NUM_PREDICT,
|
|
||||||
},
|
|
||||||
}),
|
|
||||||
signal: AbortSignal.timeout(60_000),
|
|
||||||
});
|
});
|
||||||
|
|
||||||
if (!res.ok) throw new Error(`Ollama responded ${res.status}`);
|
|
||||||
|
|
||||||
const data = await res.json();
|
|
||||||
const raw = data.response?.trim() ?? '';
|
|
||||||
|
|
||||||
const jsonMatch = raw.match(/\{[\s\S]*\}/);
|
const jsonMatch = raw.match(/\{[\s\S]*\}/);
|
||||||
if (!jsonMatch) {
|
if (!jsonMatch) {
|
||||||
logger.warn('[entities] No JSON object found in response');
|
logger.warn('[entities] No JSON object found in response');
|
||||||
|
|||||||
@@ -1,5 +1,5 @@
|
|||||||
const {getDB} = require('../db');
|
const {getDB} = require('../db');
|
||||||
const { EPISODIC, getEnv, SERVICES, parseRow, formatEpisodeText, SUMMARIES, logger } = require('@nexusai/shared');
|
const { EPISODIC, getEnv, SERVICES, parseRow, formatEpisodeText, SUMMARIES, logger, isTrivialTurn, CONSOLIDATION } = require('@nexusai/shared');
|
||||||
const semantic = require('../semantic');
|
const semantic = require('../semantic');
|
||||||
const { extractAndStoreEntities } = require('../entities/extraction')
|
const { extractAndStoreEntities } = require('../entities/extraction')
|
||||||
|
|
||||||
@@ -162,8 +162,16 @@ async function createEpisode(sessionId, userMessage, aiResponse, tokenCount = nu
|
|||||||
}))
|
}))
|
||||||
.catch(err => logger.error(`Failed to embed episode ${episode.id}:`, err.message));
|
.catch(err => logger.error(`Failed to embed episode ${episode.id}:`, err.message));
|
||||||
|
|
||||||
extractAndStoreEntities(userMessage, aiResponse, episode.id, projectId)
|
// Skip entity extraction on contentless social turns (greetings, sign-offs).
|
||||||
.catch(err => logger.error(`Failed to extract entities for episode ${episode.id}:`, err.message));
|
// They carry nothing worth storing, and running extraction on them was a
|
||||||
|
// source of junk/confabulated entities. The informational content lives in
|
||||||
|
// substantive turns, which still extract normally.
|
||||||
|
if (isTrivialTurn(userMessage)) {
|
||||||
|
logger.debug(`[entities] Skipping extraction for episode ${episode.id} — trivial turn`);
|
||||||
|
} else {
|
||||||
|
extractAndStoreEntities(userMessage, aiResponse, episode.id, projectId)
|
||||||
|
.catch(err => logger.error(`Failed to extract entities for episode ${episode.id}:`, err.message));
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
return episode;
|
return episode;
|
||||||
@@ -258,6 +266,19 @@ function deleteEpisode(id) {
|
|||||||
db.prepare(`DELETE FROM episodes WHERE id = ?`).run(id);
|
db.prepare(`DELETE FROM episodes WHERE id = ?`).run(id);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
function touchEpisodes(ids) {
|
||||||
|
if (!ids.length) return;
|
||||||
|
|
||||||
|
const db = getDB(); // <-- missing in your version
|
||||||
|
const placeholders = ids.map(() => '?').join(',');
|
||||||
|
db.prepare(`
|
||||||
|
UPDATE episodes
|
||||||
|
SET access_count = access_count + 1,
|
||||||
|
last_accessed_at = ?
|
||||||
|
WHERE id IN (${placeholders})
|
||||||
|
`).run(Date.now(), ...ids);
|
||||||
|
}
|
||||||
|
|
||||||
/******** Embedding Helper ********/
|
/******** Embedding Helper ********/
|
||||||
async function getEpisodeEmbedding(userMessage, aiResponse){
|
async function getEpisodeEmbedding(userMessage, aiResponse){
|
||||||
const url = getEnv('EMBEDDING_SERVICE_URL', SERVICES.EMBEDDING_URL);
|
const url = getEnv('EMBEDDING_SERVICE_URL', SERVICES.EMBEDDING_URL);
|
||||||
@@ -290,6 +311,22 @@ function getEpisodesByProject(projectId, limit = SUMMARIES.MAX_PROJECT_EPISODE_L
|
|||||||
`).all(projectId, limit).map(parseRow);
|
`).all(projectId, limit).map(parseRow);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
function getConsolidationCandidates(sessionId, limit = CONSOLIDATION.CANDIDATE_LIMIT) {
|
||||||
|
const db = getDB();
|
||||||
|
const now = Date.now();
|
||||||
|
const cutoff = Math.floor(now / 1000) - CONSOLIDATION.MIN_AGE_DAYS *86400 // seconds, so it matches created_at
|
||||||
|
|
||||||
|
return db.prepare(`
|
||||||
|
SELECT id, access_count, created_at, last_accessed_at,
|
||||||
|
SUBSTR(user_message, 1, 80) AS preview,
|
||||||
|
CAST( access_count AS REAL) / (1+(?-COALESCE(last_accessed_at, created_at * 1000)) / 86400000.0) AS aging_score
|
||||||
|
FROM episodes
|
||||||
|
WHERE session_id = ? AND created_at < ?
|
||||||
|
ORDER BY aging_score ASC
|
||||||
|
LIMIT ?
|
||||||
|
`).all(now, sessionId, cutoff, limit);
|
||||||
|
}
|
||||||
|
|
||||||
module.exports = {
|
module.exports = {
|
||||||
createSession,
|
createSession,
|
||||||
getSession,
|
getSession,
|
||||||
@@ -307,6 +344,8 @@ module.exports = {
|
|||||||
getEpisodesSince,
|
getEpisodesSince,
|
||||||
searchEpisodes,
|
searchEpisodes,
|
||||||
deleteEpisode,
|
deleteEpisode,
|
||||||
|
touchEpisodes,
|
||||||
getEpisodesByProject,
|
getEpisodesByProject,
|
||||||
buildFtsQuery,
|
buildFtsQuery,
|
||||||
|
getConsolidationCandidates,
|
||||||
};
|
};
|
||||||
@@ -74,4 +74,17 @@ function getEpisodeIdsByEntities(entityIds) {
|
|||||||
).all(...entityIds).map(r => r.episode_id);
|
).all(...entityIds).map(r => r.episode_id);
|
||||||
}
|
}
|
||||||
|
|
||||||
module.exports = { getNeighborhood, getEntityNeighbors, getEpisodeIdsByEntities };
|
//Entity IDs linked (via entity_episodes) to any episode in a given session
|
||||||
|
//Scopes non-project entity search to the session's own entities under the "isolated chats" model.
|
||||||
|
//Entities globally deduped
|
||||||
|
function getEntityIdsBySession(sessionId){
|
||||||
|
const db = getDB();
|
||||||
|
return db.prepare(`
|
||||||
|
SELECT DISTINCT ee.entity_id
|
||||||
|
FROM entity_episodes ee
|
||||||
|
JOIN episodes e on e.id = ee.episode_id
|
||||||
|
WHERE e.session_id = ?
|
||||||
|
`).all(sessionId).map(r => r.entity_id);
|
||||||
|
}
|
||||||
|
|
||||||
|
module.exports = { getNeighborhood, getEntityNeighbors, getEpisodeIdsByEntities, getEntityIdsBySession };
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
require ('dotenv').config();
|
require ('dotenv').config();
|
||||||
const express = require('express');
|
const express = require('express');
|
||||||
const {getEnv, PORTS, EPISODIC, logger} = require('@nexusai/shared');
|
const {getEnv, PORTS, EPISODIC, logger, CONSOLIDATION} = require('@nexusai/shared');
|
||||||
const { getDB } = require('./db');
|
const { getDB } = require('./db');
|
||||||
const { createProject, getProjects, getProject, updateProject, deleteProject } = require('./db/projects');
|
const { createProject, getProjects, getProject, updateProject, deleteProject } = require('./db/projects');
|
||||||
const { createSummary, getSummary, getSummariesBySession, getSummariesByProject, updateSummary, deleteSummary } = require('./db/summaries');
|
const { createSummary, getSummary, getSummariesBySession, getSummariesByProject, updateSummary, deleteSummary } = require('./db/summaries');
|
||||||
@@ -139,6 +139,13 @@ app.get('/episodes/search', (req, res) => {
|
|||||||
res.json(episodic.searchEpisodes(q, Number(limit), parsedSessionIds));
|
res.json(episodic.searchEpisodes(q, Number(limit), parsedSessionIds));
|
||||||
});
|
});
|
||||||
|
|
||||||
|
app.post('/episodes/touch', (req, res) => {
|
||||||
|
const { ids } = req.body;
|
||||||
|
if(!Array.isArray(ids)) return res.status(400).json({error: 'ids must be an array'});
|
||||||
|
episodic.touchEpisodes(ids);
|
||||||
|
res.json({touched: ids.length});
|
||||||
|
})
|
||||||
|
|
||||||
app.get('/episodes/:id', (req, res) => {
|
app.get('/episodes/:id', (req, res) => {
|
||||||
const episode = episodic.getEpisode(req.params.id);
|
const episode = episodic.getEpisode(req.params.id);
|
||||||
if (!episode) return res.status(404).json({ error: 'Episode not found' });
|
if (!episode) return res.status(404).json({ error: 'Episode not found' });
|
||||||
@@ -162,12 +169,39 @@ app.get('/sessions/:id/episode-stats', (req, res) => {
|
|||||||
res.json(episodic.getSessionEpisodeStats(Number(req.params.id)));
|
res.json(episodic.getSessionEpisodeStats(Number(req.params.id)));
|
||||||
});
|
});
|
||||||
|
|
||||||
|
//Entity IDs linked to this session's episodes: sesion-scoped entity search
|
||||||
|
app.get('/sessions/:id/entity-ids', (req, res) => {
|
||||||
|
res.json({
|
||||||
|
entityIds: graph.getEntityIdsBySession(Number(req.params.id))
|
||||||
|
})
|
||||||
|
})
|
||||||
|
|
||||||
// Episodes newer than :afterId, chronological — the un-summarized tail.
|
// Episodes newer than :afterId, chronological — the un-summarized tail.
|
||||||
app.get('/sessions/:id/episodes/since/:afterId', (req, res) => {
|
app.get('/sessions/:id/episodes/since/:afterId', (req, res) => {
|
||||||
const episodes = episodic.getEpisodesSince(Number(req.params.id), Number(req.params.afterId));
|
const episodes = episodic.getEpisodesSince(Number(req.params.id), Number(req.params.afterId));
|
||||||
res.json(episodes);
|
res.json(episodes);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
app.get('/sessions/:id/consolidation-candidates', (req, res) => {
|
||||||
|
const sessionId = Number(req.params.id);
|
||||||
|
const stats = episodic.getSessionEpisodeStats(sessionId);
|
||||||
|
|
||||||
|
if(stats.count < CONSOLIDATION.MIN_SESSION_EPISODES) {
|
||||||
|
return res.json({
|
||||||
|
eligible: false,
|
||||||
|
reason: `session has ${stats.count} episodes, floor is ${CONSOLIDATION.MIN_SESSION_EPISODES}`,
|
||||||
|
candidates: [],
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
const candidates = episodic.getConsolidationCandidates(sessionId);
|
||||||
|
res.json({
|
||||||
|
eligible:true,
|
||||||
|
sessionEpisodeCount: stats.count,
|
||||||
|
candidates
|
||||||
|
})
|
||||||
|
})
|
||||||
|
|
||||||
app.delete('/episodes/:id', (req, res) => {
|
app.delete('/episodes/:id', (req, res) => {
|
||||||
const id = Number(req.params.id);
|
const id = Number(req.params.id);
|
||||||
episodic.deleteEpisode(id);
|
episodic.deleteEpisode(id);
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
const { SERVICES, getEnv, SUMMARIES } = require('@nexusai/shared');
|
const { SERVICES, getEnv, SUMMARIES, utilityInference } = require('@nexusai/shared');
|
||||||
const {
|
const {
|
||||||
getSessionSummariesForProject,
|
getSessionSummariesForProject,
|
||||||
getProjectOverviewSummary,
|
getProjectOverviewSummary,
|
||||||
@@ -9,9 +9,6 @@ const {
|
|||||||
const { getEpisodesByProject } = require('../episodic');
|
const { getEpisodesByProject } = require('../episodic');
|
||||||
const { getProject } = require('../db/projects');
|
const { getProject } = require('../db/projects');
|
||||||
|
|
||||||
const EXTRACTION_URL = getEnv('EXTRACTION_URL', 'http://localhost:11434');
|
|
||||||
const EXTRACTION_MODEL = getEnv('EXTRACTION_MODEL', 'qwen2.5:3b');
|
|
||||||
|
|
||||||
const MAX_SUMMARY_CHARS = SUMMARIES.MAX_SUMMARY_CHARS; // generous ceiling before we truncate input
|
const MAX_SUMMARY_CHARS = SUMMARIES.MAX_SUMMARY_CHARS; // generous ceiling before we truncate input
|
||||||
|
|
||||||
function buildProjectSummaryPrompt(projectName, sessionSummaries) {
|
function buildProjectSummaryPrompt(projectName, sessionSummaries) {
|
||||||
@@ -24,8 +21,9 @@ function buildProjectSummaryPrompt(projectName, sessionSummaries) {
|
|||||||
summaryBlock = summaryBlock.slice(-MAX_SUMMARY_CHARS);
|
summaryBlock = summaryBlock.slice(-MAX_SUMMARY_CHARS);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// No ChatML wrapper — the model's own prompt template is applied server-side
|
||||||
|
// by the inference service's /utility/complete route (Ollama /api/chat).
|
||||||
return [
|
return [
|
||||||
'<|im_start|>user',
|
|
||||||
`The following are session summaries from a project called "${projectName}".`,
|
`The following are session summaries from a project called "${projectName}".`,
|
||||||
'Write a project overview covering: goals, progress, key decisions, and current state.',
|
'Write a project overview covering: goals, progress, key decisions, and current state.',
|
||||||
'Scale the length to the material — use multiple paragraphs for complex projects, a few sentences for simple ones.',
|
'Scale the length to the material — use multiple paragraphs for complex projects, a few sentences for simple ones.',
|
||||||
@@ -33,8 +31,6 @@ function buildProjectSummaryPrompt(projectName, sessionSummaries) {
|
|||||||
'Write in third person. Output only the overview text, no headings or labels.',
|
'Write in third person. Output only the overview text, no headings or labels.',
|
||||||
'',
|
'',
|
||||||
summaryBlock,
|
summaryBlock,
|
||||||
'<|im_end|>',
|
|
||||||
'<|im_start|>assistant',
|
|
||||||
].join('\n');
|
].join('\n');
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -49,8 +45,9 @@ function buildProjectSummaryFromEpisodesPrompt(projectName, episodes) {
|
|||||||
episodeBlock = episodeBlock.slice(-MAX_SUMMARY_CHARS);
|
episodeBlock = episodeBlock.slice(-MAX_SUMMARY_CHARS);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// No ChatML wrapper — the model's own prompt template is applied server-side
|
||||||
|
// by the inference service's /utility/complete route (Ollama /api/chat).
|
||||||
return [
|
return [
|
||||||
'<|im_start|>user',
|
|
||||||
`The following are conversations from a project called "${projectName}".`,
|
`The following are conversations from a project called "${projectName}".`,
|
||||||
'Write a project overview covering: goals, progress, key decisions, and current state.',
|
'Write a project overview covering: goals, progress, key decisions, and current state.',
|
||||||
'Scale the length to the material — use multiple paragraphs for complex projects, a few sentences for simple ones.',
|
'Scale the length to the material — use multiple paragraphs for complex projects, a few sentences for simple ones.',
|
||||||
@@ -58,58 +55,25 @@ function buildProjectSummaryFromEpisodesPrompt(projectName, episodes) {
|
|||||||
'Write in third person. Output only the overview text, no headings or labels.',
|
'Write in third person. Output only the overview text, no headings or labels.',
|
||||||
'',
|
'',
|
||||||
episodeBlock,
|
episodeBlock,
|
||||||
'<|im_end|>',
|
|
||||||
'<|im_start|>assistant',
|
|
||||||
].join('\n');
|
].join('\n');
|
||||||
}
|
}
|
||||||
|
|
||||||
async function generateProjectSummaryFromEpisodes(projectName, episodes) {
|
async function generateProjectSummaryFromEpisodes(projectName, episodes) {
|
||||||
const prompt = buildProjectSummaryFromEpisodesPrompt(projectName, episodes);
|
const user = buildProjectSummaryFromEpisodesPrompt(projectName, episodes);
|
||||||
|
return utilityInference({
|
||||||
const res = await fetch(`${EXTRACTION_URL}/api/generate`, {
|
user,
|
||||||
method: 'POST',
|
temperature: SUMMARIES.TEMPERATURE,
|
||||||
headers: { 'Content-Type': 'application/json' },
|
maxTokens: SUMMARIES.PROJECT_GEN_MAX_TOKENS,
|
||||||
body: JSON.stringify({
|
|
||||||
model: EXTRACTION_MODEL,
|
|
||||||
prompt,
|
|
||||||
stream: false,
|
|
||||||
options: { temperature: 0.2, num_predict: 1200 },
|
|
||||||
}),
|
|
||||||
});
|
});
|
||||||
|
|
||||||
if (!res.ok) throw new Error(`Ollama responded ${res.status}`);
|
|
||||||
const data = await res.json();
|
|
||||||
|
|
||||||
const raw = data.response?.trim() ?? '';
|
|
||||||
return raw
|
|
||||||
.replace(/<\|im_start\|>.*?<\|im_end\|>/gs, '')
|
|
||||||
.replace(/<\|im_start\|>|<\|im_end\|>|<\|im_sep\|>/g, '')
|
|
||||||
.trim();
|
|
||||||
}
|
}
|
||||||
|
|
||||||
async function generateProjectSummary(projectName, sessionSummaries) {
|
async function generateProjectSummary(projectName, sessionSummaries) {
|
||||||
const prompt = buildProjectSummaryPrompt(projectName, sessionSummaries);
|
const user = buildProjectSummaryPrompt(projectName, sessionSummaries);
|
||||||
|
return utilityInference({
|
||||||
const res = await fetch(`${EXTRACTION_URL}/api/generate`, {
|
user,
|
||||||
method: 'POST',
|
temperature: SUMMARIES.TEMPERATURE,
|
||||||
headers: { 'Content-Type': 'application/json' },
|
maxTokens: SUMMARIES.PROJECT_GEN_MAX_TOKENS,
|
||||||
body: JSON.stringify({
|
|
||||||
model: EXTRACTION_MODEL,
|
|
||||||
prompt,
|
|
||||||
stream: false,
|
|
||||||
// No format: 'json' — we want free-text narrative, same as session summarization
|
|
||||||
options: { temperature: 0.2, num_predict: 1200 },
|
|
||||||
}),
|
|
||||||
});
|
});
|
||||||
|
|
||||||
if (!res.ok) throw new Error(`Ollama responded ${res.status}`);
|
|
||||||
const data = await res.json();
|
|
||||||
|
|
||||||
const raw = data.response?.trim() ?? '';
|
|
||||||
return raw
|
|
||||||
.replace(/<\|im_start\|>.*?<\|im_end\|>/gs, '')
|
|
||||||
.replace(/<\|im_start\|>|<\|im_end\|>|<\|im_sep\|>/g, '')
|
|
||||||
.trim();
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// Main entry point — called by the route handler
|
// Main entry point — called by the route handler
|
||||||
@@ -142,4 +106,4 @@ async function generateAndStoreProjectSummary(projectId) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
module.exports = { generateAndStoreProjectSummary };
|
module.exports = { generateAndStoreProjectSummary };
|
||||||
|
|||||||
@@ -2,7 +2,7 @@ const memory = require("../services/memory");
|
|||||||
const inference = require("../services/inference");
|
const inference = require("../services/inference");
|
||||||
const embedding = require("../services/embedding");
|
const embedding = require("../services/embedding");
|
||||||
const qdrant = require("../services/qdrant");
|
const qdrant = require("../services/qdrant");
|
||||||
const { ORCHESTRATION, RETRIEVAL, logger } = require("@nexusai/shared");
|
const { ORCHESTRATION, RETRIEVAL, logger, isTrivialTurn } = require("@nexusai/shared");
|
||||||
const appSettings = require("../config/settings");
|
const appSettings = require("../config/settings");
|
||||||
const {triggerSummary} = require('../services/summarization')
|
const {triggerSummary} = require('../services/summarization')
|
||||||
const graph = require('../services/graph');
|
const graph = require('../services/graph');
|
||||||
@@ -127,16 +127,22 @@ async function getSemanticEpisodes(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
async function getRelevantEntities(vector, projectId = null) {
|
async function getRelevantEntities(vector, { projectId = null, sessionId } = {}) {
|
||||||
if (!vector) return [];
|
if (!vector) return [];
|
||||||
try {
|
try {
|
||||||
const results = await qdrant.searchEntities(vector, { projectId });
|
let allowedIds;
|
||||||
logger.info(
|
if (projectId === null || projectId === undefined) {
|
||||||
'[orchestration] Entity search results:',
|
// Non-project chat is its own island — scope to entities linked to
|
||||||
results.map((r) => ({ name: r.payload?.name, score: r.score })),
|
// THIS session. No links yet ⇒ nothing to retrieve, and we return
|
||||||
);
|
// early so searchEntities is never called unfiltered.
|
||||||
// Include the Qdrant point ID (== SQLite entity ID) for graph traversal
|
allowedIds = await memory.getEntityIdsBySession(sessionId);
|
||||||
return results.map((r) => r.payload ? { id: r.id, ...r.payload } : null).filter(Boolean);
|
logger.info(`[orchestration] Non-project chat, session ${sessionId}: ${allowedIds.length} linked entities`);
|
||||||
|
if (allowedIds.length === 0) return [];
|
||||||
|
}
|
||||||
|
const results = await qdrant.searchEntities(vector, { projectId, allowedIds });
|
||||||
|
logger.info('[orchestration] Entity search results:',
|
||||||
|
results.map(r => ({ name: r.payload?.name, score: r.score })));
|
||||||
|
return results.map(r => r.payload ? { id: r.id, ...r.payload } : null).filter(Boolean);
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
logger.debug('[orchestration] Entity search failed, continuing without:', err.message);
|
logger.debug('[orchestration] Entity search failed, continuing without:', err.message);
|
||||||
return [];
|
return [];
|
||||||
@@ -177,8 +183,11 @@ function fuseEpisodeResults(semanticEps, keywordEps, { semanticWeight, keywordWe
|
|||||||
}
|
}
|
||||||
|
|
||||||
function estimateTokens(episode) {
|
function estimateTokens(episode) {
|
||||||
return episode.token_count
|
//NOTE: episode.token_count is not used here. It stores the
|
||||||
?? Math.ceil((episode.user_message.length + episode.ai_response.length) / 4);
|
//full inference cost of that turn(prompt + injected memory + response),
|
||||||
|
//this inflate the episodes apparent size and starving the budget.
|
||||||
|
//Therefore, char/4 on the actual stored text is an honest selectedl
|
||||||
|
return Math.ceil((episode.user_message.length + episode.ai_response.length) / 4);
|
||||||
}
|
}
|
||||||
|
|
||||||
function buildScoredPool(fusedWithScores, recentEpisodes, entityBoostedIds, { entityWeight }) {
|
function buildScoredPool(fusedWithScores, recentEpisodes, entityBoostedIds, { entityWeight }) {
|
||||||
@@ -252,6 +261,7 @@ async function getFusedEpisodes(userMessage, session, recentIds, projectSessionI
|
|||||||
}
|
}
|
||||||
|
|
||||||
async function assembleContext(externalId, userMessage) {
|
async function assembleContext(externalId, userMessage) {
|
||||||
|
logger.info(`[orchestration] assembleContext ENTERED — msg: "${userMessage}", trivial: ${isTrivialTurn(userMessage)}`);
|
||||||
const settings = appSettings.load();
|
const settings = appSettings.load();
|
||||||
const { recentEpisodeLimit, semanticLimit, scoreThreshold,
|
const { recentEpisodeLimit, semanticLimit, scoreThreshold,
|
||||||
temperature, repeatPenalty, topP, topK, systemPrompt,
|
temperature, repeatPenalty, topP, topK, systemPrompt,
|
||||||
@@ -283,20 +293,27 @@ async function assembleContext(externalId, userMessage) {
|
|||||||
const isFirstMessage = recentEpisodes.length === 0;
|
const isFirstMessage = recentEpisodes.length === 0;
|
||||||
const recentIds = new Set(recentEpisodes.map(e => e.id));
|
const recentIds = new Set(recentEpisodes.map(e => e.id));
|
||||||
|
|
||||||
// 4. Embed the query once — the vector is shared by semantic episode search
|
// 4. Retrieval — skipped entirely on contentless social turns (greetings,
|
||||||
// and entity search, so embedding it twice was a wasted round-trip + Ollama call.
|
// sign-offs). On those, recent history alone is the right context; running
|
||||||
let queryVector = null;
|
// semantic/keyword/entity retrieval only surfaces marginal noise the model
|
||||||
try {
|
// then confabulates around. Embed once (shared by episode + entity search).
|
||||||
queryVector = await embedding.embed(userMessage);
|
let fusedWithScores = [];
|
||||||
} catch (err) {
|
let entityResults = [];
|
||||||
logger.warn('[orchestration] Query embedding failed; semantic + entity search disabled this turn:', err.message);
|
if (!isTrivialTurn(userMessage)) {
|
||||||
}
|
let queryVector = null;
|
||||||
|
try {
|
||||||
|
queryVector = await embedding.embed(userMessage);
|
||||||
|
} catch (err) {
|
||||||
|
logger.warn('[orchestration] Query embedding failed; semantic + entity search disabled this turn:', err.message);
|
||||||
|
}
|
||||||
|
|
||||||
// 4b. Fused retrieval + entity search in parallel (both are independent)
|
[fusedWithScores, entityResults] = await Promise.all([
|
||||||
const [fusedWithScores, entityResults] = await Promise.all([
|
getFusedEpisodes(userMessage, session, recentIds, projectSessionIds, queryVector, { semanticLimit, scoreThreshold, semanticWeight, keywordWeight }),
|
||||||
getFusedEpisodes(userMessage, session, recentIds, projectSessionIds, queryVector, { semanticLimit, scoreThreshold, semanticWeight, keywordWeight }),
|
getRelevantEntities(queryVector, { projectId: session.project_id ?? null, sessionId: session.id }),
|
||||||
getRelevantEntities(queryVector, session.project_id ?? null),
|
]);
|
||||||
]);
|
} else {
|
||||||
|
logger.debug('[orchestration] Trivial turn — skipping semantic/keyword/entity retrieval');
|
||||||
|
}
|
||||||
|
|
||||||
// 5. Entity-linked episode IDs for scoring bonus
|
// 5. Entity-linked episode IDs for scoring bonus
|
||||||
const entityIds = entityResults.map(e => e.id);
|
const entityIds = entityResults.map(e => e.id);
|
||||||
@@ -314,6 +331,9 @@ async function assembleContext(externalId, userMessage) {
|
|||||||
const scoredPool = buildScoredPool(fusedWithScores, recentEpisodes, entityBoostedIds, { entityWeight });
|
const scoredPool = buildScoredPool(fusedWithScores, recentEpisodes, entityBoostedIds, { entityWeight });
|
||||||
const { guaranteed, selected } = selectWithinBudget(scoredPool, contextBudget, minRecentEpisodes, recentEpisodes);
|
const { guaranteed, selected } = selectWithinBudget(scoredPool, contextBudget, minRecentEpisodes, recentEpisodes);
|
||||||
|
|
||||||
|
const selectedIds = selected.map(ep => ep.id);
|
||||||
|
memory.touchEpisodes(selectedIds);
|
||||||
|
|
||||||
// 7. Graph neighborhood expansion
|
// 7. Graph neighborhood expansion
|
||||||
let neighborhood = { nodes: [], edges: [] };
|
let neighborhood = { nodes: [], edges: [] };
|
||||||
if (entityIds.length > 0) {
|
if (entityIds.length > 0) {
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
const { getEnv, SERVICES, EPISODIC } = require('@nexusai/shared');
|
const { getEnv, SERVICES, EPISODIC, logger } = require('@nexusai/shared');
|
||||||
|
|
||||||
const BASE_URL = getEnv('MEMORY_SERVICE_URL', SERVICES.MEMORY_URL);
|
const BASE_URL = getEnv('MEMORY_SERVICE_URL', SERVICES.MEMORY_URL);
|
||||||
|
|
||||||
@@ -216,6 +216,23 @@ async function getEpisodesByEntities(entityIds) {
|
|||||||
return res.json(); // { episodeIds: [...] }
|
return res.json(); // { episodeIds: [...] }
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async function getEntityIdsBySession(sessionId){
|
||||||
|
const res = await fetch(`${BASE_URL}/sessions/${sessionId}/entity-ids`);
|
||||||
|
if (!res.ok) throw new Error(`Entity-ids-by-session error: ${res.status}`);
|
||||||
|
const { entityIds } = await res.json();
|
||||||
|
return entityIds;
|
||||||
|
}
|
||||||
|
|
||||||
|
// orchestration-service/src/services/memory.js
|
||||||
|
async function touchEpisodes(ids) {
|
||||||
|
if (!ids.length) return;
|
||||||
|
fetch(`${BASE_URL}/episodes/touch`, {
|
||||||
|
method: 'POST',
|
||||||
|
headers: { 'Content-Type': 'application/json' },
|
||||||
|
body: JSON.stringify({ ids }),
|
||||||
|
}).catch(err => logger.warn(`[memory] touch failed: ${err.message}`));
|
||||||
|
}
|
||||||
|
|
||||||
module.exports = {
|
module.exports = {
|
||||||
getSessionByExternalId,
|
getSessionByExternalId,
|
||||||
createSession,
|
createSession,
|
||||||
@@ -242,4 +259,6 @@ module.exports = {
|
|||||||
getProjectOverviewSummary,
|
getProjectOverviewSummary,
|
||||||
searchEpisodes,
|
searchEpisodes,
|
||||||
getEpisodesByEntities,
|
getEpisodesByEntities,
|
||||||
|
getEntityIdsBySession,
|
||||||
|
touchEpisodes,
|
||||||
}
|
}
|
||||||
@@ -29,30 +29,28 @@ async function searchEpisodes( vector, {limit = ORCHESTRATION.RECENT_EPISODE_LIM
|
|||||||
return data.result;
|
return data.result;
|
||||||
}
|
}
|
||||||
|
|
||||||
async function searchEntities(vector, { limit = ORCHESTRATION.ENTITIES_LIMIT, scoreThreshold = ORCHESTRATION.ENTITIES_THRESHOLD, projectId = undefined } = {}) {
|
async function searchEntities(vector, { limit = ORCHESTRATION.ENTITIES_LIMIT, scoreThreshold = ORCHESTRATION.ENTITIES_THRESHOLD, projectId, allowedIds } = {}) {
|
||||||
const body = { vector, limit, score_threshold: scoreThreshold, with_payload: true };
|
const body = { vector, limit, score_threshold: scoreThreshold, with_payload: true };
|
||||||
|
|
||||||
if (projectId !== null && projectId !== undefined) {
|
if (projectId !== null && projectId !== undefined) {
|
||||||
body.filter = {
|
// Project chat: entities shared across the project's sessions.
|
||||||
must: [{ key: 'projectId', match: { value: projectId } }]
|
body.filter = { must: [{ key: 'projectId', match: { value: projectId } }] };
|
||||||
};
|
} else if (allowedIds && allowedIds.length > 0) {
|
||||||
|
// Non-project chat: restrict to this session's own entities (Model 2).
|
||||||
|
body.filter = { must: [{ has_id: allowedIds }] };
|
||||||
}
|
}
|
||||||
|
// No else: the caller returns early when a non-project session has no linked
|
||||||
|
// entities, so an unfiltered (leaky) search is never reached.
|
||||||
|
|
||||||
const res = await fetch(
|
const res = await fetch(
|
||||||
`${BASE_URL}/collections/${COLLECTIONS.ENTITIES}/points/search`,
|
`${BASE_URL}/collections/${COLLECTIONS.ENTITIES}/points/search`,
|
||||||
{
|
{ method: 'POST', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify(body) }
|
||||||
method: 'POST',
|
|
||||||
headers: { 'Content-Type': 'application/json' },
|
|
||||||
body: JSON.stringify(body),
|
|
||||||
}
|
|
||||||
);
|
);
|
||||||
|
|
||||||
if (!res.ok) {
|
if (!res.ok) {
|
||||||
const body = await res.text();
|
const text = await res.text();
|
||||||
throw new Error(`Qdrant error: ${res.status} - ${body}`);
|
throw new Error(`Qdrant error: ${res.status} - ${text}`);
|
||||||
}
|
}
|
||||||
|
return (await res.json()).result;
|
||||||
const data = await res.json();
|
|
||||||
return data.result;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
module.exports = { searchEpisodes, searchEntities };
|
module.exports = { searchEpisodes, searchEntities };
|
||||||
@@ -1,7 +1,5 @@
|
|||||||
const { getEnv, SERVICES, SUMMARIES, logger } = require('@nexusai/shared');
|
const { getEnv, SERVICES, SUMMARIES, logger, utilityInference } = require('@nexusai/shared');
|
||||||
|
|
||||||
const EXTRACTION_URL = getEnv('EXTRACTION_URL', 'http://localhost:11434');
|
|
||||||
const EXTRACTION_MODEL = getEnv('EXTRACTION_MODEL', 'qwen2.5:3b');
|
|
||||||
const MEMORY_URL = getEnv('MEMORY_SERVICE_URL', SERVICES.MEMORY_URL);
|
const MEMORY_URL = getEnv('MEMORY_SERVICE_URL', SERVICES.MEMORY_URL);
|
||||||
|
|
||||||
const THRESHOLD_TOKENS = parseInt(getEnv('SUMMARY_THRESHOLD_TOKENS', SUMMARIES.THRESHOLD_TOKENS));
|
const THRESHOLD_TOKENS = parseInt(getEnv('SUMMARY_THRESHOLD_TOKENS', SUMMARIES.THRESHOLD_TOKENS));
|
||||||
@@ -35,41 +33,20 @@ Do not include greetings, sign-offs, or filler. Output only the summary text.
|
|||||||
Conversation:
|
Conversation:
|
||||||
${context}`;
|
${context}`;
|
||||||
|
|
||||||
return [
|
// No ChatML wrapper — the model's own prompt template is applied server-side
|
||||||
'<|im_start|>user', // ChatML for qwen2.5
|
// by the inference service's /utility/complete route (Ollama /api/chat).
|
||||||
instruction,
|
return instruction;
|
||||||
'<|im_end|>',
|
|
||||||
'<|im_start|>assistant',
|
|
||||||
].join('\n');
|
|
||||||
}
|
}
|
||||||
|
|
||||||
async function generateSummary(episodes, existingSummary = null) {
|
async function generateSummary(episodes, existingSummary = null) {
|
||||||
const prompt = buildSummaryPrompt(episodes, existingSummary);
|
const user = buildSummaryPrompt(episodes, existingSummary);
|
||||||
|
|
||||||
const res = await fetch(`${EXTRACTION_URL}/api/generate`, {
|
const content = await utilityInference({
|
||||||
method: 'POST',
|
user,
|
||||||
headers: { 'Content-Type': 'application/json' },
|
temperature: SUMMARIES.TEMPERATURE,
|
||||||
body: JSON.stringify({
|
maxTokens: SUMMARIES.SESSION_GEN_MAX_TOKENS,
|
||||||
model: EXTRACTION_MODEL,
|
|
||||||
prompt,
|
|
||||||
stream: false,
|
|
||||||
options: {
|
|
||||||
temperature: 0.2, // slightly higher than entities — summaries benefit from some fluency
|
|
||||||
num_predict: 500, // generous but bounded — keeps summaries from running long
|
|
||||||
},
|
|
||||||
}),
|
|
||||||
});
|
});
|
||||||
|
|
||||||
if (!res.ok) throw new Error(`Ollama responded ${res.status}`);
|
|
||||||
const data = await res.json();
|
|
||||||
|
|
||||||
|
|
||||||
const raw = data.response?.trim() ?? '';
|
|
||||||
// Strip any leaked ChatML tokens Qwen echoes back
|
|
||||||
const content = raw
|
|
||||||
.replace(/<\|im_start\|>.*?<\|im_end\|>/gs, '')
|
|
||||||
.replace(/<\|im_start\|>|<\|im_end\|>|<\|im_sep\|>/g, '')
|
|
||||||
.trim();
|
|
||||||
return content;
|
return content;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -148,4 +125,4 @@ async function triggerSummary(session) {
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
module.exports = { triggerSummary, maybeSummarize };
|
module.exports = { triggerSummary, maybeSummarize };
|
||||||
|
|||||||
@@ -78,6 +78,13 @@ const SUMMARIES = {
|
|||||||
MIN_EPISODES_SINCE: 5, // don't resummarize until N new episodes since last summary
|
MIN_EPISODES_SINCE: 5, // don't resummarize until N new episodes since last summary
|
||||||
MAX_SUMMARY_CHARS: 8000, // max chars to include from recent episodes when generating summary (to control prompt size)
|
MAX_SUMMARY_CHARS: 8000, // max chars to include from recent episodes when generating summary (to control prompt size)
|
||||||
MAX_PROJECT_EPISODE_LIMIT: 200, // max number of episodes to consider from the entire project when generating summary (to control prompt size)
|
MAX_PROJECT_EPISODE_LIMIT: 200, // max number of episodes to consider from the entire project when generating summary (to control prompt size)
|
||||||
|
|
||||||
|
// Generation params for the utility model (passed to utilityInference).
|
||||||
|
// Distinct from MAX_SUMMARY_TOKENS above, which gates STORED summary size;
|
||||||
|
// these two cap GENERATION length (num_predict) per summary type.
|
||||||
|
TEMPERATURE: 0.2, // slightly higher than entities (0.1) — summaries benefit from some fluency
|
||||||
|
SESSION_GEN_MAX_TOKENS: 500, // num_predict for a session summary (3-5 sentences)
|
||||||
|
PROJECT_GEN_MAX_TOKENS: 1200, // num_predict for a project overview (multi-paragraph)
|
||||||
}
|
}
|
||||||
|
|
||||||
const ENTITIES = {
|
const ENTITIES = {
|
||||||
@@ -105,6 +112,20 @@ const RETRIEVAL = {
|
|||||||
KEYWORD_WEIGHT: 0.5, // Weight applied to keyword (SQLite) results, 0 = disables, set >0 to enable and tune balance between semantic vs keyword matches
|
KEYWORD_WEIGHT: 0.5, // Weight applied to keyword (SQLite) results, 0 = disables, set >0 to enable and tune balance between semantic vs keyword matches
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const UTILITY = {
|
||||||
|
DEFAULT_URL: 'http://localhost:11434', // Ollama host for background/utility tasks
|
||||||
|
DEFAULT_MODEL: 'qwen2.5:3b',
|
||||||
|
TEMPERATURE: 0.2, // precise, low-creativity default for extraction/summaries
|
||||||
|
MAX_TOKENS: 1200,
|
||||||
|
TIMEOUT_MS: 120_000, // extraction previously used 60s; summaries had none — standardized
|
||||||
|
};
|
||||||
|
|
||||||
|
const CONSOLIDATION = {
|
||||||
|
MIN_AGE_DAYS: 7, // episodes younger than this are never consolidation candidates
|
||||||
|
MIN_SESSION_EPISODES: 20, // sessions smaller than this are left entirely alone
|
||||||
|
CANDIDATE_LIMIT: 50,
|
||||||
|
}
|
||||||
|
|
||||||
module.exports = {
|
module.exports = {
|
||||||
QDRANT,
|
QDRANT,
|
||||||
COLLECTIONS,
|
COLLECTIONS,
|
||||||
@@ -119,4 +140,6 @@ module.exports = {
|
|||||||
SUMMARIES,
|
SUMMARIES,
|
||||||
ENTITIES,
|
ENTITIES,
|
||||||
RETRIEVAL,
|
RETRIEVAL,
|
||||||
|
UTILITY,
|
||||||
|
CONSOLIDATION,
|
||||||
};
|
};
|
||||||
@@ -1,7 +1,25 @@
|
|||||||
const {getEnv} = require('./config/env');
|
const {getEnv} = require('./config/env');
|
||||||
const {QDRANT, COLLECTIONS, EPISODIC, SERVICES, OLLAMA, PORTS, LLAMACPP, INFERENCE_DEFAULTS, SQLITE, ORCHESTRATION, SUMMARIES, ENTITIES, RETRIEVAL } = require('./config/constants');
|
const {
|
||||||
const {parseRow, formatEpisodeText} = require('./utils')
|
QDRANT,
|
||||||
|
COLLECTIONS,
|
||||||
|
EPISODIC,
|
||||||
|
SERVICES,
|
||||||
|
OLLAMA,
|
||||||
|
PORTS,
|
||||||
|
LLAMACPP,
|
||||||
|
INFERENCE_DEFAULTS,
|
||||||
|
SQLITE,
|
||||||
|
ORCHESTRATION,
|
||||||
|
SUMMARIES,
|
||||||
|
ENTITIES,
|
||||||
|
RETRIEVAL,
|
||||||
|
UTILITY,
|
||||||
|
CONSOLIDATION
|
||||||
|
} = require('./config/constants');
|
||||||
|
const {parseRow, formatEpisodeText, isTrivialTurn} = require('./utils')
|
||||||
|
|
||||||
const logger = require('./utils/logger');
|
const logger = require('./utils/logger');
|
||||||
|
const {utilityInference} = require('./utils/utilityInference');
|
||||||
|
|
||||||
module.exports = {
|
module.exports = {
|
||||||
getEnv,
|
getEnv,
|
||||||
@@ -17,8 +35,12 @@ module.exports = {
|
|||||||
ORCHESTRATION,
|
ORCHESTRATION,
|
||||||
parseRow,
|
parseRow,
|
||||||
formatEpisodeText,
|
formatEpisodeText,
|
||||||
|
isTrivialTurn,
|
||||||
SUMMARIES,
|
SUMMARIES,
|
||||||
ENTITIES,
|
ENTITIES,
|
||||||
logger,
|
logger,
|
||||||
RETRIEVAL,
|
RETRIEVAL,
|
||||||
|
UTILITY,
|
||||||
|
CONSOLIDATION,
|
||||||
|
utilityInference,
|
||||||
};
|
};
|
||||||
@@ -10,4 +10,46 @@ function formatEpisodeText(userMessage, aiResponse) {
|
|||||||
return `User: ${userMessage}\nAssistant: ${aiResponse}`;
|
return `User: ${userMessage}\nAssistant: ${aiResponse}`;
|
||||||
}
|
}
|
||||||
|
|
||||||
module.exports = { parseRow, formatEpisodeText };
|
// Contentless "social" turns — greetings, sign-offs, acknowledgements — that
|
||||||
|
// carry no information to store or recall. Used to skip entity extraction and
|
||||||
|
// noisy retrieval on such turns.
|
||||||
|
const TRIVIAL_TURNS = new Set([
|
||||||
|
'good morning', 'good night', 'good evening', 'good afternoon', 'morning', 'evening',
|
||||||
|
'hello', 'hi', 'hey', 'hey there', 'yo', 'sup', 'whats up', 'hiya', 'howdy',
|
||||||
|
'goodbye', 'bye', 'see you', 'see ya', 'see you later', 'talk later', 'talk soon',
|
||||||
|
'later', 'catch you later', 'gtg', 'gotta go', 'im off', 'heading out',
|
||||||
|
'thanks', 'thank you', 'thanks again', 'thank you so much', 'ty', 'thx', 'cheers',
|
||||||
|
'no worries', 'no problem', 'np', 'youre welcome', 'my pleasure',
|
||||||
|
'ok', 'okay', 'k', 'kk', 'alright', 'sure', 'sounds good', 'got it', 'gotcha',
|
||||||
|
'cool', 'nice', 'great', 'awesome', 'perfect', 'lol', 'haha',
|
||||||
|
'just saying hi', 'just saying hello', 'just checking in', 'just dropping in',
|
||||||
|
'stopping by', 'just stopping by', 'just wanted to say hi',
|
||||||
|
]);
|
||||||
|
|
||||||
|
// Trailing filler words that don't change a phrase's social nature, so
|
||||||
|
// "good morning again" / "thanks everyone" collapse to a base phrase.
|
||||||
|
const TRIVIAL_FILLER = new Set(['again', 'there', 'everyone', 'all', 'yall', 'folks', 'man', 'dude', 'friend', 'buddy']);
|
||||||
|
|
||||||
|
// Conservative, HIGH-PRECISION check: is this a contentless social turn?
|
||||||
|
// It exists to skip retrieval/extraction on turns with nothing to remember or
|
||||||
|
// recall, and deliberately errs toward "substantive" — a real short query like
|
||||||
|
// "capital of France" must NOT be treated as trivial. Genuine intent
|
||||||
|
// classification is a separate, later concern (confidence-based routing), not
|
||||||
|
// this heuristic. Extend TRIVIAL_TURNS as new pure-social phrases show up.
|
||||||
|
function isTrivialTurn(message) {
|
||||||
|
if (!message) return true;
|
||||||
|
const norm = String(message)
|
||||||
|
.toLowerCase()
|
||||||
|
.replace(/[^\p{L}\p{N}\s]/gu, ' ') // punctuation → space (unicode-aware)
|
||||||
|
.replace(/\s+/g, ' ')
|
||||||
|
.trim();
|
||||||
|
if (!norm) return true; // empty / punctuation-only
|
||||||
|
if (TRIVIAL_TURNS.has(norm)) return true;
|
||||||
|
|
||||||
|
// Drop trailing filler and re-check ("good morning again" → "good morning")
|
||||||
|
const words = norm.split(' ');
|
||||||
|
while (words.length > 1 && TRIVIAL_FILLER.has(words[words.length - 1])) words.pop();
|
||||||
|
return TRIVIAL_TURNS.has(words.join(' '));
|
||||||
|
}
|
||||||
|
|
||||||
|
module.exports = { parseRow, formatEpisodeText, isTrivialTurn };
|
||||||
@@ -0,0 +1,20 @@
|
|||||||
|
const { getEnv } = require('../config/env');
|
||||||
|
const { SERVICES } = require('../config/constants');
|
||||||
|
|
||||||
|
// Client for the inference service's /utility/complete route.
|
||||||
|
// Resolved at call time (not module load) so each service's .env applies.
|
||||||
|
async function utilityInference({ system, user, json = false, temperature, maxTokens }) {
|
||||||
|
const base = getEnv('INFERENCE_SERVICE_URL', SERVICES.INFERENCE_URL);
|
||||||
|
|
||||||
|
const res = await fetch(`${base}/utility/complete`, {
|
||||||
|
method: 'POST',
|
||||||
|
headers: { 'Content-Type': 'application/json' },
|
||||||
|
body: JSON.stringify({ system, user, json, temperature, maxTokens }),
|
||||||
|
});
|
||||||
|
|
||||||
|
if (!res.ok) throw new Error(`Utility inference error: ${res.status}`);
|
||||||
|
const data = await res.json();
|
||||||
|
return data.text;
|
||||||
|
}
|
||||||
|
|
||||||
|
module.exports = { utilityInference };
|
||||||
@@ -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));
|
calls.sinceAfterId = Number(u.split('/since/').at(-1));
|
||||||
return json(since.filter(ep => ep.id > calls.sinceAfterId));
|
return json(since.filter(ep => ep.id > calls.sinceAfterId));
|
||||||
}
|
}
|
||||||
if (u.includes('/api/generate')) return json({ response: 'A concise third-person summary.' });
|
if (u.endsWith('/utility/complete')) return json({ text: 'A concise third-person summary.' });
|
||||||
if (u.endsWith('/summaries') && opts.method === 'POST') { calls.posted = JSON.parse(opts.body); return json({ id: 99 }); }
|
if (u.endsWith('/summaries') && opts.method === 'POST') { calls.posted = JSON.parse(opts.body); return json({ id: 99 }); }
|
||||||
if (/\/summaries\/\d+$/.test(u) && opts.method === 'PATCH') { calls.patched = JSON.parse(opts.body); return json({ ok: true }); }
|
if (/\/summaries\/\d+$/.test(u) && opts.method === 'PATCH') { calls.patched = JSON.parse(opts.body); return json({ ok: true }); }
|
||||||
throw new Error('unexpected fetch: ' + u);
|
throw new Error('unexpected fetch: ' + u);
|
||||||
|
|||||||
@@ -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