Compare commits
7
Commits
6c953a346a
..
main
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
3d7372d04f | ||
|
|
c31fb786f4 | ||
|
|
66f6c4534f | ||
|
|
b3765f30fa | ||
|
|
4b2099498c | ||
|
|
fd8867ee96 | ||
|
|
00101ba14a |
@@ -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.
|
||||||
@@ -279,6 +283,8 @@ Both fields are optional. Only provided fields are updated.
|
|||||||
| GET | /sessions/:id/episodes?limit=&offset= | Paginated episodes for a session |
|
| GET | /sessions/:id/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/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) |
|
| 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`.
|
||||||
@@ -293,6 +299,21 @@ Both fields are optional. Only provided fields are updated.
|
|||||||
}
|
}
|
||||||
```
|
```
|
||||||
|
|
||||||
|
**POST /episodes/touch — body:**
|
||||||
|
```json
|
||||||
|
{ "ids": [54, 55, 56] }
|
||||||
|
```
|
||||||
|
Returns `{ "touched": 3 }`. Called fire-and-forget by orchestration after
|
||||||
|
budget selection. Nonexistent IDs are silently skipped; `touched` echoes the
|
||||||
|
request count, not rows matched.
|
||||||
|
|
||||||
|
**GET /sessions/:id/consolidation-candidates** — returns
|
||||||
|
`{ eligible, reason?, candidates: [...] }` where each candidate is
|
||||||
|
`{ id, access_count, created_at, last_accessed_at, preview, aging_score }`,
|
||||||
|
sorted by `aging_score` ascending (most eligible first). Sessions under
|
||||||
|
`CONSOLIDATION.MIN_SESSION_EPISODES` return `eligible: false` with a `reason`.
|
||||||
|
Observe-only — nothing is modified.
|
||||||
|
|
||||||
### Projects
|
### Projects
|
||||||
|
|
||||||
| Method | Path | Description |
|
| Method | Path | Description |
|
||||||
|
|||||||
@@ -18,6 +18,7 @@ npm test # = node --test (discovers test/*.test.js at the repo root)
|
|||||||
| `test/summarization.test.js` | Summarization decision logic | `maybeSummarize` (orchestration) |
|
| `test/summarization.test.js` | Summarization decision logic | `maybeSummarize` (orchestration) |
|
||||||
| `test/entity-extraction.test.js` | Greeting + regurgitation guards | `mentionedIn`, `isIgnoredName` (memory-service) |
|
| `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/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) |
|
| `test/migrations.test.js` | Migration version-stepping | `migrate` (memory-service) |
|
||||||
|
|
||||||
Tests import the **real** functions rather than reimplementing logic — the
|
Tests import the **real** functions rather than reimplementing logic — the
|
||||||
@@ -44,3 +45,10 @@ Keep the pattern: export the real function, import it, mock I/O at the boundary
|
|||||||
logic (ranking, tokenizing, version-stepping, decision branches) over wiring.
|
logic (ranking, tokenizing, version-stepping, decision branches) over wiring.
|
||||||
Several of these tests were written *after* a bug slipped through — each new
|
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.
|
class of mistake is worth a case so it can't recur silently.
|
||||||
|
|
||||||
|
When a mocked service call changes shape (e.g. the `utilityInference` refactor
|
||||||
|
moving summarization from Ollama's `/api/generate` to the inference service's
|
||||||
|
`/utility/complete`), the mock URL router must move with it — an "unexpected
|
||||||
|
fetch" throw in these tests usually means the code under test evolved, not
|
||||||
|
broke. Runner-contract tests (migrations) inject stub no-op migration arrays
|
||||||
|
rather than letting the real SQL-bearing migrations hit the minimal fake db.
|
||||||
|
|||||||
+8
-7
@@ -74,8 +74,9 @@ Multi-strategy retrieval merged into a single ranked result set.
|
|||||||
|
|
||||||
### 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
|
||||||
@@ -90,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."
|
||||||
@@ -227,4 +228,4 @@ The JARVIS moment — NexusAI reasons, plans, and acts across multiple steps.
|
|||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
*Last updated: April 2026*
|
*Last updated: August 2026*
|
||||||
@@ -78,6 +78,11 @@ does **not** reconcile columns on old tables — that's what migrations are for)
|
|||||||
```js
|
```js
|
||||||
const migrations = [
|
const migrations = [
|
||||||
(_db) => {}, // v0 → v1: consolidated baseline (historical ALTERs folded into schema.js)
|
(_db) => {}, // v0 → v1: consolidated baseline (historical ALTERs folded into schema.js)
|
||||||
|
(db) => { // v1 → v2: access tracking for consolidation lifecycle
|
||||||
|
db.exec(`ALTER TABLE episodes ADD COLUMN last_accessed_at INTEGER`);
|
||||||
|
db.exec(`ALTER TABLE episodes ADD COLUMN access_count INTEGER NOT NULL DEFAULT 0`);
|
||||||
|
db.exec(`UPDATE episodes SET last_accessed_at = created_at`); // backfill
|
||||||
|
},
|
||||||
];
|
];
|
||||||
const LATEST_VERSION = migrations.length; // derived, never hand-maintained
|
const LATEST_VERSION = migrations.length; // derived, never hand-maintained
|
||||||
```
|
```
|
||||||
@@ -121,6 +126,12 @@ previously ran on every startup.
|
|||||||
- `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
|
||||||
@@ -204,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
|
||||||
|
|||||||
@@ -73,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 |
|
||||||
@@ -101,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
|
||||||
|
|||||||
@@ -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 |
|
||||||
|
|||||||
@@ -1,5 +1,5 @@
|
|||||||
const {getDB} = require('../db');
|
const {getDB} = require('../db');
|
||||||
const { EPISODIC, getEnv, SERVICES, parseRow, formatEpisodeText, SUMMARIES, logger, isTrivialTurn } = require('@nexusai/shared');
|
const { EPISODIC, getEnv, SERVICES, parseRow, formatEpisodeText, SUMMARIES, logger, isTrivialTurn, CONSOLIDATION } = require('@nexusai/shared');
|
||||||
const semantic = require('../semantic');
|
const semantic = require('../semantic');
|
||||||
const { extractAndStoreEntities } = require('../entities/extraction')
|
const { extractAndStoreEntities } = require('../entities/extraction')
|
||||||
|
|
||||||
@@ -266,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);
|
||||||
@@ -298,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,
|
||||||
@@ -315,6 +344,8 @@ module.exports = {
|
|||||||
getEpisodesSince,
|
getEpisodesSince,
|
||||||
searchEpisodes,
|
searchEpisodes,
|
||||||
deleteEpisode,
|
deleteEpisode,
|
||||||
|
touchEpisodes,
|
||||||
getEpisodesByProject,
|
getEpisodesByProject,
|
||||||
buildFtsQuery,
|
buildFtsQuery,
|
||||||
|
getConsolidationCandidates,
|
||||||
};
|
};
|
||||||
@@ -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' });
|
||||||
@@ -175,6 +182,26 @@ app.get('/sessions/:id/episodes/since/:afterId', (req, res) => {
|
|||||||
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);
|
||||||
|
|||||||
@@ -331,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);
|
||||||
|
|
||||||
@@ -223,6 +223,16 @@ async function getEntityIdsBySession(sessionId){
|
|||||||
return entityIds;
|
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,
|
||||||
@@ -250,4 +260,5 @@ module.exports = {
|
|||||||
searchEpisodes,
|
searchEpisodes,
|
||||||
getEpisodesByEntities,
|
getEpisodesByEntities,
|
||||||
getEntityIdsBySession,
|
getEntityIdsBySession,
|
||||||
|
touchEpisodes,
|
||||||
}
|
}
|
||||||
@@ -120,6 +120,12 @@ const UTILITY = {
|
|||||||
TIMEOUT_MS: 120_000, // extraction previously used 60s; summaries had none — standardized
|
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,
|
||||||
@@ -135,4 +141,5 @@ module.exports = {
|
|||||||
ENTITIES,
|
ENTITIES,
|
||||||
RETRIEVAL,
|
RETRIEVAL,
|
||||||
UTILITY,
|
UTILITY,
|
||||||
|
CONSOLIDATION,
|
||||||
};
|
};
|
||||||
@@ -13,7 +13,8 @@ const {
|
|||||||
SUMMARIES,
|
SUMMARIES,
|
||||||
ENTITIES,
|
ENTITIES,
|
||||||
RETRIEVAL,
|
RETRIEVAL,
|
||||||
UTILITY
|
UTILITY,
|
||||||
|
CONSOLIDATION
|
||||||
} = require('./config/constants');
|
} = require('./config/constants');
|
||||||
const {parseRow, formatEpisodeText, isTrivialTurn} = require('./utils')
|
const {parseRow, formatEpisodeText, isTrivialTurn} = require('./utils')
|
||||||
|
|
||||||
@@ -40,5 +41,6 @@ module.exports = {
|
|||||||
logger,
|
logger,
|
||||||
RETRIEVAL,
|
RETRIEVAL,
|
||||||
UTILITY,
|
UTILITY,
|
||||||
|
CONSOLIDATION,
|
||||||
utilityInference,
|
utilityInference,
|
||||||
};
|
};
|
||||||
@@ -4,7 +4,6 @@
|
|||||||
// advance user_version correctly, and no-op when the DB is already current.
|
// advance user_version correctly, and no-op when the DB is already current.
|
||||||
const { test } = require('node:test');
|
const { test } = require('node:test');
|
||||||
const assert = require('node:assert');
|
const assert = require('node:assert');
|
||||||
const Database = require('better-sqlite3');
|
|
||||||
const { migrate, LATEST_VERSION } = require('../packages/memory-service/src/db/migrations');
|
const { migrate, LATEST_VERSION } = require('../packages/memory-service/src/db/migrations');
|
||||||
|
|
||||||
// Emulates just the slice of the better-sqlite3 API that migrate() uses.
|
// Emulates just the slice of the better-sqlite3 API that migrate() uses.
|
||||||
@@ -26,14 +25,16 @@ function fakeDb(startVersion = 0) {
|
|||||||
|
|
||||||
test('a fresh DB (user_version 0) is stamped to LATEST_VERSION', () => {
|
test('a fresh DB (user_version 0) is stamped to LATEST_VERSION', () => {
|
||||||
const db = fakeDb(0);
|
const db = fakeDb(0);
|
||||||
const result = migrate(db);
|
const stubs = Array.from({ length: LATEST_VERSION }, () => () => {});
|
||||||
|
const result = migrate(db, stubs); // inject no-op migrations
|
||||||
assert.strictEqual(result, LATEST_VERSION);
|
assert.strictEqual(result, LATEST_VERSION);
|
||||||
assert.strictEqual(db.version, LATEST_VERSION);
|
assert.strictEqual(db.version, LATEST_VERSION);
|
||||||
});
|
});
|
||||||
|
|
||||||
test('an already-current DB is a no-op (no version writes)', () => {
|
test('an already-current DB is a no-op (no version writes)', () => {
|
||||||
const db = fakeDb(LATEST_VERSION);
|
const db = fakeDb(LATEST_VERSION);
|
||||||
const result = migrate(db);
|
const stubs = Array.from({ length: LATEST_VERSION }, () => () => {});
|
||||||
|
const result = migrate(db, stubs); // fixed: db is now the first arg
|
||||||
assert.strictEqual(result, LATEST_VERSION);
|
assert.strictEqual(result, LATEST_VERSION);
|
||||||
assert.deepStrictEqual(db.setCalls, [], 'should not touch user_version when current');
|
assert.deepStrictEqual(db.setCalls, [], 'should not touch user_version when current');
|
||||||
});
|
});
|
||||||
|
|||||||
+1
-1
@@ -13,7 +13,7 @@ try { ({ DatabaseSync } = require('node:sqlite')); } catch { DatabaseSync = null
|
|||||||
// The complete intended shape = original base tables + every historical ALTER.
|
// The complete intended shape = original base tables + every historical ALTER.
|
||||||
const EXPECTED = {
|
const EXPECTED = {
|
||||||
sessions: ['id','external_id','created_at','updated_at','metadata','name','project_id'],
|
sessions: ['id','external_id','created_at','updated_at','metadata','name','project_id'],
|
||||||
episodes: ['id','session_id','user_message','ai_response','created_at','token_count','metadata'],
|
episodes: ['access_count', 'ai_response', 'created_at', 'id', 'last_accessed_at', 'metadata', 'session_id', 'token_count', 'user_message'],
|
||||||
entities: ['id','name','type','notes','created_at','updated_at','metadata','mention_count','confidence','source','last_seen_at'],
|
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'],
|
relationships: ['id','from_id','to_id','label','created_at','metadata','mention_count','notes'],
|
||||||
entity_episodes: ['entity_id','episode_id'],
|
entity_episodes: ['entity_id','episode_id'],
|
||||||
|
|||||||
@@ -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);
|
||||||
|
|||||||
Reference in New Issue
Block a user