13 Commits

Author SHA1 Message Date
e11122ed14 Merge pull request 'FEAT-nonblocking-backend-io-imporvements' (#1) from FEAT-nonblocking-backend-io-imporvements into master
Reviewed-on: #1
2026-08-01 05:43:24 +00:00
KS Jannette
7b47e46852 Cleanup 2026-08-01 01:42:43 -04:00
KS Jannette
62477d009c Updated Claude model 2026-08-01 01:35:27 -04:00
KS Jannette
ee6fa9e576 Add nonblocking/asyn I/O operations 2026-08-01 01:06:19 -04:00
90035b3568 Update README.ms
hotfix
2026-07-31 13:45:35 +00:00
S Jannette
9636be78c9 Fix typos and enhance README content
Improved clarity in the README.
2026-03-06 03:52:47 -05:00
S Jannette
a50ca5c171 Add MIT License to the project 2026-02-25 18:50:59 -05:00
S Jannette
e271702bca Fix typo in README.md regarding implementation planning
Corrected 'incoporated in' to 'incoporated into' for clarity.
2026-02-24 21:11:24 -05:00
S Jannette
d3d76b1d0d Refine README.md content for clarity and accuracy
Updated descriptions for clarity and corrected typos.
2026-02-24 21:10:51 -05:00
S Jannette
8cae0665db Merge pull request #16 from kjannette/refinements
Refinements
2026-02-24 21:04:56 -05:00
KS Jannette
21453aab7d UI development 2026-02-24 21:02:43 -05:00
KS Jannette
691186e02c upgraded scoring algorithm model 2026-02-24 20:55:19 -05:00
S Jannette
fd48aaab4d Merge pull request #15 from kjannette/test-data-xfer
further DB infra buildout
2026-02-24 19:41:31 -05:00
20 changed files with 881 additions and 269 deletions

21
LICENSE Normal file
View File

@@ -0,0 +1,21 @@
MIT License
Copyright (c) 2026 Steven Jannette
Permission is hereby granted, free of charge, to any person obtaining a copy
of this software and associated documentation files (the "Software"), to deal
in the Software without restriction, including without limitation the rights
to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
copies of the Software, and to permit persons to whom the Software is
furnished to do so, subject to the following conditions:
The above copyright notice and this permission notice shall be included in all
copies or substantial portions of the Software.
THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
SOFTWARE.

View File

@@ -1,10 +1,12 @@
# kongruity # kongruity
kongruity pulls in the unstructured artifacts of the creative-engineering process -- to-dos, action items, agile tickets, Jira comment threads, Slack threads, retrospective notes -- and synthesizes them into semantically coherent, prioritized clusters ready for implementation planning. kongruity pulls in unstructured artifacts of the creative-engineering process -- to-dos, action items, agile tickets, Jira thread comments, Slack thread comments, retrospective notes -- and synthesizes them into semantically coherent, prioritized clusters that can be incorporated into implementation planning.
In kongruity world, these artifacts are "sticky notes." A board full of them looks chaotic. With a click, an LLM analyzes their semantic meaning and groups them into thematic clusters, each with a descriptive header. In kongruity world, the artifacts become "sticky notes." A board full of them looks chaotic. With a click, the LLM analyzes their semantic meaning and groups them into thematic clusters, each with a descriptive header.
An independent embedding-based evaluation scores clustering quality, so the output is data-backed. From there, teams can simply drag-and-rank clusters by implementation priority, turning noise into an actionable workflow. ## Result evalutation/validation:
An independent RAG embedding-based model scores evaluates quality using cosine similiarty, so the output is qualitatively refined. From there, teams can drag-and-rank related task clusters by implementation priority - turning noise into an actionable workflow.
## How it works ## How it works
@@ -14,9 +16,14 @@ An independent embedding-based evaluation scores clustering quality, so the outp
4. **Validate** — Structural checks confirm every note is assigned to exactly one cluster, no clusters are empty, and labels are present. 4. **Validate** — Structural checks confirm every note is assigned to exactly one cluster, no clusters are empty, and labels are present.
5. **Prioritize** — Clusters appear ranked and are drag-reorderable. Teams set implementation priority by dragging clusters into position. 5. **Prioritize** — Clusters appear ranked and are drag-reorderable. Teams set implementation priority by dragging clusters into position.
## Dev implementation note ## Dev implementation notes
Developers may swap in other LLM SDKs/APIs and alter prompt syntax in `backend/services/clustering.service.js` to experiment with any model or platform. As of 03.04.2026, two of the above-described features are in the planning and implementation phase:
1. **Ingestion** - The implementation goal is a system for easily tagging items/issues mentioned in Slack, Jira comments, etc. (similar to hashtagging), and running batch "pulls" of these tagged items into kongruity via REST API interface.
2. **Prioritization** -- The final prioritization and planning stage will ultimately result in "pushing" these ordered items back **out** into Jira, Asana, Rally (etc.), for incorporation into Epic/Sprint workflows.
Also, developers may swap in other LLM SDKs/APIs and alter prompt syntax in `backend/services/clustering.service.js` to experiment with any model or platform of their choice.
## Prerequisites ## Prerequisites
@@ -48,7 +55,13 @@ EOF
Replace placeholder values with your actual keys and database credentials. Replace placeholder values with your actual keys and database credentials.
### 3. Set up the database ### 3. Set up/run the database
Start DB for local development (assumes local dev env MacOS and Homebrew installed)
```bash
brew services start postgresql@15
```
Create a PostgreSQL database for the project: Create a PostgreSQL database for the project:
@@ -144,3 +157,5 @@ This runs Vitest with jsdom. For watch mode during development:
```bash ```bash
npm run test:watch npm run test:watch
``` ```

View File

@@ -0,0 +1,50 @@
{"id":"note_001","text":"Login flow feels confusing","x":193,"y":191,"author":"user_5","color":"yellow"}
{"id":"note_002","text":"Login flow is broken on mobile","x":214,"y":281,"author":"user_9","color":"yellow"}
{"id":"note_003","text":"When I enter my username and password I get an unknown error","x":189,"y":193,"author":"user_2","color":"yellow"}
{"id":"note_004","text":"Password reset email never arrives","x":207,"y":267,"author":"user_9","color":"yellow"}
{"id":"note_005","text":"SSO login loops back to the sign-in page","x":184,"y":124,"author":"user_7","color":"yellow"}
{"id":"note_006","text":"Two-factor code is rejected even when correct","x":162,"y":213,"author":"user_1","color":"yellow"}
{"id":"note_007","text":"I'm logged out unexpectedly after a few minutes","x":185,"y":157,"author":"user_3","color":"yellow"}
{"id":"note_008","text":"Cannot change my password — save button does nothing","x":244,"y":188,"author":"user_2","color":"orange"}
{"id":"note_009","text":"Account gets locked too quickly after one failed attempt","x":280,"y":255,"author":"user_10","color":"orange"}
{"id":"note_010","text":"OAuth consent screen appears every time I sign in","x":228,"y":124,"author":"user_9","color":"yellow"}
{"id":"note_011","text":"Need better export options (PDF quality is too low)","x":748,"y":212,"author":"user_9","color":"green"}
{"id":"note_012","text":"Exported PDF cuts off content near the edges","x":733,"y":159,"author":"user_6","color":"blue"}
{"id":"note_013","text":"Cannot export selected area — only full board exports","x":696,"y":205,"author":"user_4","color":"green"}
{"id":"note_014","text":"Export takes too long and sometimes never finishes","x":798,"y":211,"author":"user_2","color":"green"}
{"id":"note_015","text":"Sharing link permissions are confusing","x":688,"y":290,"author":"user_6","color":"blue"}
{"id":"note_016","text":"Downloaded image is blurry compared to the canvas","x":677,"y":245,"author":"user_5","color":"blue"}
{"id":"note_017","text":"Exported file names are inconsistent and hard to track","x":676,"y":201,"author":"user_4","color":"blue"}
{"id":"note_018","text":"No way to schedule recurring exports for stakeholders","x":661,"y":229,"author":"user_9","color":"blue"}
{"id":"note_019","text":"Embedded exports don't update when the board changes","x":662,"y":132,"author":"user_1","color":"blue"}
{"id":"note_020","text":"Can't export comments and reactions with the content","x":739,"y":139,"author":"user_7","color":"green"}
{"id":"note_021","text":"Board feels slow when there are many sticky notes","x":341,"y":695,"author":"user_10","color":"purple"}
{"id":"note_022","text":"Canvas freezes for a few seconds when zooming","x":254,"y":707,"author":"user_8","color":"pink"}
{"id":"note_023","text":"Undo/redo sometimes lags and applies late","x":236,"y":687,"author":"user_9","color":"purple"}
{"id":"note_024","text":"App crashes when opening a large board","x":239,"y":597,"author":"user_10","color":"purple"}
{"id":"note_025","text":"Saving indicator spins but changes aren't saved","x":129,"y":781,"author":"user_3","color":"purple"}
{"id":"note_026","text":"Search is slow on boards with lots of content","x":253,"y":658,"author":"user_2","color":"pink"}
{"id":"note_027","text":"Scrolling stutters on older laptops","x":178,"y":586,"author":"user_7","color":"pink"}
{"id":"note_028","text":"High CPU usage even when idle on a board","x":190,"y":695,"author":"user_8","color":"purple"}
{"id":"note_029","text":"Offline mode loses edits when reconnecting","x":338,"y":632,"author":"user_1","color":"pink"}
{"id":"note_030","text":"I see random 'something went wrong' banners with no details","x":214,"y":594,"author":"user_5","color":"purple"}
{"id":"note_031","text":"Hard to tell who is editing what in real time","x":761,"y":704,"author":"user_8","color":"yellow"}
{"id":"note_032","text":"Comments get lost — no clear thread view","x":818,"y":641,"author":"user_5","color":"yellow"}
{"id":"note_033","text":"Mentions (@) don't notify the right people","x":696,"y":669,"author":"user_5","color":"yellow"}
{"id":"note_034","text":"Too many notification emails for minor edits","x":769,"y":739,"author":"user_9","color":"yellow"}
{"id":"note_035","text":"No notification when someone resolves my comment","x":673,"y":636,"author":"user_2","color":"blue"}
{"id":"note_036","text":"Cursor presence is distracting and overlaps content","x":788,"y":605,"author":"user_5","color":"yellow"}
{"id":"note_037","text":"Can't easily hand off facilitation to another user","x":816,"y":707,"author":"user_2","color":"yellow"}
{"id":"note_038","text":"Live follow mode frequently breaks","x":710,"y":579,"author":"user_9","color":"yellow"}
{"id":"note_039","text":"Guest collaborators can't see updates without refreshing","x":759,"y":711,"author":"user_9","color":"yellow"}
{"id":"note_040","text":"Activity feed lacks context about what changed","x":710,"y":771,"author":"user_7","color":"yellow"}
{"id":"note_041","text":"Hard to find the right template quickly","x":536,"y":394,"author":"user_4","color":"orange"}
{"id":"note_042","text":"Template search results feel irrelevant","x":400,"y":474,"author":"user_6","color":"orange"}
{"id":"note_043","text":"Can't organize boards into folders the way I need","x":504,"y":398,"author":"user_4","color":"green"}
{"id":"note_044","text":"Naming conventions aren't enforced and things get messy","x":469,"y":434,"author":"user_9","color":"green"}
{"id":"note_045","text":"No bulk rename for multiple sticky notes","x":455,"y":427,"author":"user_1","color":"green"}
{"id":"note_046","text":"Tags are missing — I need better categorization","x":472,"y":435,"author":"user_6","color":"green"}
{"id":"note_047","text":"Hard to keep consistent styles across boards","x":420,"y":426,"author":"user_8","color":"green"}
{"id":"note_048","text":"Duplicating a board loses some formatting","x":382,"y":410,"author":"user_10","color":"orange"}
{"id":"note_049","text":"Need better controls for aligning and distributing notes","x":462,"y":487,"author":"user_7","color":"green"}
{"id":"note_050","text":"Can't lock sections to prevent accidental edits during workshops","x":521,"y":471,"author":"user_6","color":"orange"}

View File

@@ -1,54 +1,110 @@
import { query } from './index.js'; import QueryStream from 'pg-query-stream';
import { pipeline } from 'node:stream/promises';
import { Readable } from 'node:stream';
import { query, getPool } from './index.js';
import { batch } from '../lib/streams.js';
const SELECT_NOTES = 'SELECT id, text, x, y, author, color FROM notes ORDER BY id';
// Postgres caps a statement at 65535 bind parameters; six columns per note
// leaves 10922 as the hard ceiling.
const INSERT_BATCH_SIZE = 1000;
export const getAllNotes = async () => { export const getAllNotes = async () => {
const { rows } = await query( const { rows } = await query(SELECT_NOTES);
'SELECT id, text, x, y, author, color FROM notes ORDER BY id' return rows;
); };
return rows;
/**
* Streams every note as an object-mode Readable. The pooled client is released on end, error, or
* destruction by consumer.
*
* @returns {Promise<import('node:stream').Readable>}
*/
export const streamAllNotes = async () => {
const client = await getPool().connect();
let released = false;
const release = () => {
if (released) return;
released = true;
client.release();
};
try {
const rows = client.query(new QueryStream(SELECT_NOTES));
rows.once('end', release);
rows.once('error', release);
rows.once('close', release);
return rows;
} catch (err) {
release();
throw err;
}
}; };
export const getNoteById = async (id) => { export const getNoteById = async (id) => {
const { rows } = await query( const { rows } = await query(
'SELECT id, text, x, y, author, color FROM notes WHERE id = $1', 'SELECT id, text, x, y, author, color FROM notes WHERE id = $1',
[id] [id]
); );
return rows[0] || null; return rows[0] || null;
}; };
export const createNote = async (note) => { export const createNote = async (note) => {
const { rows } = await query( const { rows } = await query(
`INSERT INTO notes (id, text, x, y, author, color) `INSERT INTO notes (id, text, x, y, author, color)
VALUES ($1, $2, $3, $4, $5, $6) VALUES ($1, $2, $3, $4, $5, $6)
RETURNING id, text, x, y, author, color`, RETURNING id, text, x, y, author, color`,
[note.id, note.text, note.x ?? 0, note.y ?? 0, note.author, note.color ?? 'yellow'] [note.id, note.text, note.x ?? 0, note.y ?? 0, note.author, note.color ?? 'yellow']
); );
return rows[0]; return rows[0];
};
const insertNoteBatch = async (notes) => {
const values = [];
const placeholders = [];
notes.forEach((note, i) => {
const offset = i * 6;
placeholders.push(
`($${offset + 1}, $${offset + 2}, $${offset + 3}, $${offset + 4}, $${offset + 5}, $${offset + 6})`
);
values.push(
note.id,
note.text,
note.x ?? 0,
note.y ?? 0,
note.author,
note.color ?? 'yellow'
);
});
const { rows } = await query(
`INSERT INTO notes (id, text, x, y, author, color)
VALUES ${placeholders.join(', ')}
RETURNING id, text, x, y, author, color`,
values
);
return rows;
}; };
export const createNotes = async (notes) => { export const createNotes = async (notes) => {
const values = []; if (!notes || notes.length === 0) {
const placeholders = []; return [];
}
notes.forEach((note, i) => { const inserted = [];
const offset = i * 6;
placeholders.push(
`($${offset + 1}, $${offset + 2}, $${offset + 3}, $${offset + 4}, $${offset + 5}, $${offset + 6})`
);
values.push(
note.id,
note.text,
note.x ?? 0,
note.y ?? 0,
note.author,
note.color ?? 'yellow'
);
});
const { rows } = await query( await pipeline(
`INSERT INTO notes (id, text, x, y, author, color) Readable.from(notes, { objectMode: true }),
VALUES ${placeholders.join(', ')} batch(INSERT_BATCH_SIZE),
RETURNING id, text, x, y, author, color`, async (batches) => {
values for await (const chunk of batches) {
); inserted.push(...await insertNoteBatch(chunk));
return rows; }
}
);
return inserted;
}; };

View File

@@ -1,85 +1,78 @@
import 'dotenv/config'; import 'dotenv/config';
import { query, close } from './index.js'; import split2 from 'split2';
import { from as copyFrom } from 'pg-copy-streams';
import { createReadStream } from 'node:fs';
import { Transform } from 'node:stream';
import { pipeline } from 'node:stream/promises';
import { fileURLToPath } from 'node:url';
import { getPool, close } from './index.js';
const notes = [ const FIXTURE = fileURLToPath(new URL('./fixtures/notes.jsonl', import.meta.url));
{ id: 'note_001', text: 'Login flow feels confusing', x: 193, y: 191, author: 'user_5', color: 'yellow' },
{ id: 'note_002', text: 'Login flow is broken on mobile', x: 214, y: 281, author: 'user_9', color: 'yellow' },
{ id: 'note_003', text: 'When I enter my username and password I get an unknown error', x: 189, y: 193, author: 'user_2', color: 'yellow' },
{ id: 'note_004', text: 'Password reset email never arrives', x: 207, y: 267, author: 'user_9', color: 'yellow' },
{ id: 'note_005', text: 'SSO login loops back to the sign-in page', x: 184, y: 124, author: 'user_7', color: 'yellow' },
{ id: 'note_006', text: 'Two-factor code is rejected even when correct', x: 162, y: 213, author: 'user_1', color: 'yellow' },
{ id: 'note_007', text: "I'm logged out unexpectedly after a few minutes", x: 185, y: 157, author: 'user_3', color: 'yellow' },
{ id: 'note_008', text: 'Cannot change my password — save button does nothing', x: 244, y: 188, author: 'user_2', color: 'orange' },
{ id: 'note_009', text: 'Account gets locked too quickly after one failed attempt', x: 280, y: 255, author: 'user_10', color: 'orange' },
{ id: 'note_010', text: 'OAuth consent screen appears every time I sign in', x: 228, y: 124, author: 'user_9', color: 'yellow' },
{ id: 'note_011', text: 'Need better export options (PDF quality is too low)', x: 748, y: 212, author: 'user_9', color: 'green' },
{ id: 'note_012', text: 'Exported PDF cuts off content near the edges', x: 733, y: 159, author: 'user_6', color: 'blue' },
{ id: 'note_013', text: 'Cannot export selected area — only full board exports', x: 696, y: 205, author: 'user_4', color: 'green' },
{ id: 'note_014', text: 'Export takes too long and sometimes never finishes', x: 798, y: 211, author: 'user_2', color: 'green' },
{ id: 'note_015', text: 'Sharing link permissions are confusing', x: 688, y: 290, author: 'user_6', color: 'blue' },
{ id: 'note_016', text: 'Downloaded image is blurry compared to the canvas', x: 677, y: 245, author: 'user_5', color: 'blue' },
{ id: 'note_017', text: 'Exported file names are inconsistent and hard to track', x: 676, y: 201, author: 'user_4', color: 'blue' },
{ id: 'note_018', text: 'No way to schedule recurring exports for stakeholders', x: 661, y: 229, author: 'user_9', color: 'blue' },
{ id: 'note_019', text: "Embedded exports don't update when the board changes", x: 662, y: 132, author: 'user_1', color: 'blue' },
{ id: 'note_020', text: "Can't export comments and reactions with the content", x: 739, y: 139, author: 'user_7', color: 'green' },
{ id: 'note_021', text: 'Board feels slow when there are many sticky notes', x: 341, y: 695, author: 'user_10', color: 'purple' },
{ id: 'note_022', text: 'Canvas freezes for a few seconds when zooming', x: 254, y: 707, author: 'user_8', color: 'pink' },
{ id: 'note_023', text: 'Undo/redo sometimes lags and applies late', x: 236, y: 687, author: 'user_9', color: 'purple' },
{ id: 'note_024', text: 'App crashes when opening a large board', x: 239, y: 597, author: 'user_10', color: 'purple' },
{ id: 'note_025', text: "Saving indicator spins but changes aren't saved", x: 129, y: 781, author: 'user_3', color: 'purple' },
{ id: 'note_026', text: 'Search is slow on boards with lots of content', x: 253, y: 658, author: 'user_2', color: 'pink' },
{ id: 'note_027', text: 'Scrolling stutters on older laptops', x: 178, y: 586, author: 'user_7', color: 'pink' },
{ id: 'note_028', text: 'High CPU usage even when idle on a board', x: 190, y: 695, author: 'user_8', color: 'purple' },
{ id: 'note_029', text: 'Offline mode loses edits when reconnecting', x: 338, y: 632, author: 'user_1', color: 'pink' },
{ id: 'note_030', text: "I see random 'something went wrong' banners with no details", x: 214, y: 594, author: 'user_5', color: 'purple' },
{ id: 'note_031', text: 'Hard to tell who is editing what in real time', x: 761, y: 704, author: 'user_8', color: 'yellow' },
{ id: 'note_032', text: 'Comments get lost — no clear thread view', x: 818, y: 641, author: 'user_5', color: 'yellow' },
{ id: 'note_033', text: "Mentions (@) don't notify the right people", x: 696, y: 669, author: 'user_5', color: 'yellow' },
{ id: 'note_034', text: 'Too many notification emails for minor edits', x: 769, y: 739, author: 'user_9', color: 'yellow' },
{ id: 'note_035', text: 'No notification when someone resolves my comment', x: 673, y: 636, author: 'user_2', color: 'blue' },
{ id: 'note_036', text: 'Cursor presence is distracting and overlaps content', x: 788, y: 605, author: 'user_5', color: 'yellow' },
{ id: 'note_037', text: "Can't easily hand off facilitation to another user", x: 816, y: 707, author: 'user_2', color: 'yellow' },
{ id: 'note_038', text: 'Live follow mode frequently breaks', x: 710, y: 579, author: 'user_9', color: 'yellow' },
{ id: 'note_039', text: "Guest collaborators can't see updates without refreshing", x: 759, y: 711, author: 'user_9', color: 'yellow' },
{ id: 'note_040', text: 'Activity feed lacks context about what changed', x: 710, y: 771, author: 'user_7', color: 'yellow' },
{ id: 'note_041', text: 'Hard to find the right template quickly', x: 536, y: 394, author: 'user_4', color: 'orange' },
{ id: 'note_042', text: 'Template search results feel irrelevant', x: 400, y: 474, author: 'user_6', color: 'orange' },
{ id: 'note_043', text: "Can't organize boards into folders the way I need", x: 504, y: 398, author: 'user_4', color: 'green' },
{ id: 'note_044', text: "Naming conventions aren't enforced and things get messy", x: 469, y: 434, author: 'user_9', color: 'green' },
{ id: 'note_045', text: 'No bulk rename for multiple sticky notes', x: 455, y: 427, author: 'user_1', color: 'green' },
{ id: 'note_046', text: 'Tags are missing — I need better categorization', x: 472, y: 435, author: 'user_6', color: 'green' },
{ id: 'note_047', text: 'Hard to keep consistent styles across boards', x: 420, y: 426, author: 'user_8', color: 'green' },
{ id: 'note_048', text: 'Duplicating a board loses some formatting', x: 382, y: 410, author: 'user_10', color: 'orange' },
{ id: 'note_049', text: 'Need better controls for aligning and distributing notes', x: 462, y: 487, author: 'user_7', color: 'green' },
{ id: 'note_050', text: "Can't lock sections to prevent accidental edits during workshops", x: 521, y: 471, author: 'user_6', color: 'orange' },
];
const run = async () => { const COLUMNS = ['id', 'text', 'x', 'y', 'author', 'color'];
try {
const insertQuery = `
INSERT INTO notes (id, text, x, y, author, color)
VALUES ($1, $2, $3, $4, $5, $6)
ON CONFLICT (id) DO NOTHING
`;
let inserted = 0; const DEFAULTS = { x: 0, y: 0, color: 'yellow' };
for (const note of notes) {
const result = await query(insertQuery, [ const csvField = (value) => {
note.id, if (value === null || value === undefined) return '';
note.text, return `"${String(value).replaceAll('"', '""')}"`;
note.x, };
note.y,
note.author, const toCsvRows = () => new Transform({
note.color, writableObjectMode: true,
]); transform(line, _encoding, callback) {
inserted += result.rowCount; if (line.trim().length === 0) {
callback();
return;
} }
console.log(`Seed complete — ${inserted} notes inserted (${notes.length - inserted} already existed).`); let note;
try {
note = JSON.parse(line);
} catch {
callback(new Error(`Fixture contains a malformed JSON line: ${line.slice(0, 80)}`));
return;
}
const row = COLUMNS.map((col) => csvField(note[col] ?? DEFAULTS[col]));
callback(null, `${row.join(',')}\n`);
},
});
const run = async () => {
const client = await getPool().connect();
try {
await client.query('BEGIN');
// COPY has no ON CONFLICT, so land the fixture in a temp table first and
// let a single INSERT ... SELECT apply the existing idempotency.
await client.query(
'CREATE TEMP TABLE notes_import (LIKE notes INCLUDING DEFAULTS) ON COMMIT DROP'
);
const copy = client.query(
copyFrom(`COPY notes_import (${COLUMNS.join(', ')}) FROM STDIN WITH (FORMAT csv)`)
);
await pipeline(createReadStream(FIXTURE), split2(), toCsvRows(), copy);
const { rowCount: staged } = await client.query('SELECT 1 FROM notes_import');
const { rowCount: inserted } = await client.query(
`INSERT INTO notes (${COLUMNS.join(', ')})
SELECT ${COLUMNS.join(', ')} FROM notes_import
ON CONFLICT (id) DO NOTHING`
);
await client.query('COMMIT');
console.log(`Seed complete — ${inserted} notes inserted (${staged - inserted} already existed).`);
} catch (err) { } catch (err) {
await client.query('ROLLBACK').catch(() => {});
console.error('Seed failed:', err.message); console.error('Seed failed:', err.message);
process.exit(1); process.exitCode = 1;
} finally { } finally {
client.release();
await close(); await close();
} }
}; };

72
backend/lib/streams.js Normal file
View File

@@ -0,0 +1,72 @@
import { Transform } from 'node:stream';
/**
* Groups an object-mode stream into fixed-size arrays. Backpressure on the
* readable side is what limits how many batches are ever in flight.
*
* @param {number} size - maximum items per emitted batch
* @returns {Transform}
*/
export const batch = (size) => {
if (!Number.isInteger(size) || size < 1) {
throw new TypeError('batch(size) requires a positive integer size');
}
let pending = [];
return new Transform({
objectMode: true,
transform(item, _encoding, callback) {
pending.push(item);
if (pending.length < size) {
callback();
return;
}
const full = pending;
pending = [];
callback(null, full);
},
flush(callback) {
if (pending.length === 0) {
callback();
return;
}
const remainder = pending;
pending = [];
callback(null, remainder);
},
});
};
/**
* Serializes an object-mode stream into a JSON array, one element at a time,
* so no complete copy of the payload is ever held in memory.
*
* @returns {Transform}
*/
export const jsonArray = () => {
let wroteFirst = false;
return new Transform({
writableObjectMode: true,
transform(item, _encoding, callback) {
let serialized;
try {
serialized = JSON.stringify(item);
} catch (err) {
callback(err);
return;
}
const prefix = wroteFirst ? ',' : '[';
wroteFirst = true;
callback(null, prefix + serialized);
},
flush(callback) {
callback(null, wroteFirst ? ']' : '[]');
},
});
};

View File

@@ -13,6 +13,8 @@
"dotenv": "^16.4.7", "dotenv": "^16.4.7",
"express": "^4.21.2", "express": "^4.21.2",
"pg": "^8.18.0", "pg": "^8.18.0",
"pg-copy-streams": "^7.0.0",
"pg-query-stream": "^4.16.0",
"voyageai": "^0.1.0" "voyageai": "^0.1.0"
}, },
"devDependencies": { "devDependencies": {
@@ -2339,6 +2341,21 @@
"integrity": "sha512-kecgoJwhOpxYU21rZjULrmrBJ698U2RxXofKVzOn5UDj61BPj/qMb7diYUR1nLScCDbrztQFl1TaQZT0t1EtzQ==", "integrity": "sha512-kecgoJwhOpxYU21rZjULrmrBJ698U2RxXofKVzOn5UDj61BPj/qMb7diYUR1nLScCDbrztQFl1TaQZT0t1EtzQ==",
"license": "MIT" "license": "MIT"
}, },
"node_modules/pg-copy-streams": {
"version": "7.0.0",
"resolved": "https://registry.npmjs.org/pg-copy-streams/-/pg-copy-streams-7.0.0.tgz",
"integrity": "sha512-zBvnY6wtaBRE2ae2xXWOOGMaNVPkXh1vhypAkNSKgMdciJeTyIQAHZaEeRAxUjs/p1El5jgzYmwG5u871Zj3dQ==",
"license": "MIT"
},
"node_modules/pg-cursor": {
"version": "2.21.0",
"resolved": "https://registry.npmjs.org/pg-cursor/-/pg-cursor-2.21.0.tgz",
"integrity": "sha512-IYvk/j+Suhtbo/C3uOf4JLsLK/gWxOTUOmYbDsbKnLaVJDq+KwhwK6ngpRfiCk8eDMS3AmGQABZCv0cREEzHQw==",
"license": "MIT",
"peerDependencies": {
"pg": "^8"
}
},
"node_modules/pg-int8": { "node_modules/pg-int8": {
"version": "1.0.1", "version": "1.0.1",
"resolved": "https://registry.npmjs.org/pg-int8/-/pg-int8-1.0.1.tgz", "resolved": "https://registry.npmjs.org/pg-int8/-/pg-int8-1.0.1.tgz",
@@ -2363,6 +2380,18 @@
"integrity": "sha512-pfsxk2M9M3BuGgDOfuy37VNRRX3jmKgMjcvAcWqNDpZSf4cUmv8HSOl5ViRQFsfARFn0KuUQTgLxVMbNq5NW3g==", "integrity": "sha512-pfsxk2M9M3BuGgDOfuy37VNRRX3jmKgMjcvAcWqNDpZSf4cUmv8HSOl5ViRQFsfARFn0KuUQTgLxVMbNq5NW3g==",
"license": "MIT" "license": "MIT"
}, },
"node_modules/pg-query-stream": {
"version": "4.16.0",
"resolved": "https://registry.npmjs.org/pg-query-stream/-/pg-query-stream-4.16.0.tgz",
"integrity": "sha512-vyqxAG4YVax43BCSfqcKxHSbh8YQshhVR0CjLNdo7MIM2UOqI+XsIYk0bVZJ0aKjQfjJQcGiGo9xK3hz0JcFyg==",
"license": "MIT",
"dependencies": {
"pg-cursor": "^2.21.0"
},
"peerDependencies": {
"pg": "^8"
}
},
"node_modules/pg-types": { "node_modules/pg-types": {
"version": "2.2.0", "version": "2.2.0",
"resolved": "https://registry.npmjs.org/pg-types/-/pg-types-2.2.0.tgz", "resolved": "https://registry.npmjs.org/pg-types/-/pg-types-2.2.0.tgz",

View File

@@ -17,6 +17,8 @@
"dotenv": "^16.4.7", "dotenv": "^16.4.7",
"express": "^4.21.2", "express": "^4.21.2",
"pg": "^8.18.0", "pg": "^8.18.0",
"pg-copy-streams": "^7.0.0",
"pg-query-stream": "^4.16.0",
"voyageai": "^0.1.0" "voyageai": "^0.1.0"
}, },
"devDependencies": { "devDependencies": {

View File

@@ -1,25 +1,38 @@
import { Router } from 'express'; import { Router } from 'express';
import { getAllNotes } from '../db/notes.dao.js'; import { pipeline } from 'node:stream/promises';
import { getAllNotes, streamAllNotes } from '../db/notes.dao.js';
import { clusterNotes } from '../services/clustering.service.js'; import { clusterNotes } from '../services/clustering.service.js';
import { jsonArray } from '../lib/streams.js';
const router = Router(); const router = Router();
router.get('/', async (req, res) => { router.get('/', async (req, res) => {
try { try {
const notes = await getAllNotes(); const rows = await streamAllNotes();
res.json(notes); res.type('application/json');
await pipeline(rows, jsonArray(), res);
} catch (err) { } catch (err) {
console.error(`Error loading notes: ${err}`); console.error(`Error loading notes: ${err}`);
if (res.headersSent) {
res.destroy(err);
return;
}
res.status(500).json({ error: 'Failed to load notes' }); res.status(500).json({ error: 'Failed to load notes' });
} }
}); });
router.post('/cluster', async (req, res) => { router.post('/cluster', async (req, res) => {
const controller = new AbortController();
res.on('close', () => {
if (!res.writableEnded) controller.abort();
});
try { try {
const notes = await getAllNotes(); const notes = await getAllNotes();
const result = await clusterNotes(notes); const result = await clusterNotes(notes, { signal: controller.signal });
res.json(result); res.json(result);
} catch (err) { } catch (err) {
if (controller.signal.aborted) return;
console.error(`Clustering failed: ${err}`); console.error(`Clustering failed: ${err}`);
res.status(500).json({ error: 'Clustering failed' }); res.status(500).json({ error: 'Clustering failed' });
} }

View File

@@ -1,4 +1,6 @@
import Anthropic from "@anthropic-ai/sdk"; import Anthropic from "@anthropic-ai/sdk";
import { pipeline } from "node:stream/promises";
import { Transform, Writable } from "node:stream";
import { embedNotes } from "./embedding.service.js"; import { embedNotes } from "./embedding.service.js";
import { validateStructure, computeCohesionScore } from "./validation.service.js"; import { validateStructure, computeCohesionScore } from "./validation.service.js";
@@ -33,31 +35,71 @@ Here are the notes:
${notesJson}`; ${notesJson}`;
}; };
const requestClusters = async (notes) => { const textDeltas = () => new Transform({
const response = await client.messages.create({ objectMode: true,
model: "claude-sonnet-4-20250514", transform(event, _encoding, callback) {
max_tokens: 4096, if (event?.type === 'content_block_delta' && event.delta?.type === 'text_delta') {
messages: [ callback(null, event.delta.text);
{ role: "user", content: buildPrompt(notes) }, return;
], }
}); callback();
},
});
const textBlock = response?.content?.[0]; // Rejects as soon as the first non-whitespace character proves the response
// is not the JSON array we asked for, rather than after the full generation.
const collectClusterJson = (sink) => new Writable({
objectMode: true,
write(text, _encoding, callback) {
if (!sink.sawOpeningBracket) {
const leading = (sink.parts.join('') + text).trimStart();
if (leading.length > 0) {
if (!leading.startsWith('[')) {
callback(new Error('LLM API returned non-JSON response'));
return;
}
sink.sawOpeningBracket = true;
}
}
sink.parts.push(text);
callback();
},
});
if (!textBlock || textBlock.type !== 'text' || typeof textBlock.text !== 'string') { const requestClusters = async (notes, signal) => {
const sink = { parts: [], sawOpeningBracket: false };
const options = signal ? [{ signal }] : [];
const events = client.messages.stream(
{
model: "claude-sonnet-5",
max_tokens: 4096,
messages: [
{ role: "user", content: buildPrompt(notes) },
],
},
...options
);
await pipeline(events, textDeltas(), collectClusterJson(sink), ...options);
const text = sink.parts.join('');
if (text.trim().length === 0) {
throw new Error('Unexpected response from LLM API: no text content returned'); throw new Error('Unexpected response from LLM API: no text content returned');
} }
try { try {
return JSON.parse(textBlock.text); return JSON.parse(text);
} catch { } catch {
throw new Error('LLM API returned non-JSON response'); throw new Error('LLM API returned non-JSON response');
} }
}; };
export const clusterNotes = async (notes) => { export const clusterNotes = async (notes, { signal } = {}) => {
const [clusters, embeddingMap] = await Promise.all([ const [clusters, embeddingMap] = await Promise.all([
requestClusters(notes), requestClusters(notes, signal),
embedNotes(notes), embedNotes(notes),
]); ]);

View File

@@ -1,25 +1,41 @@
import { VoyageAIClient } from "voyageai"; import { VoyageAIClient } from "voyageai";
import { pipeline } from "node:stream/promises";
import { Readable } from "node:stream";
import { batch } from "../lib/streams.js";
const client = new VoyageAIClient({ const client = new VoyageAIClient({
apiKey: process.env.VOYAGEAI_API_KEY, apiKey: process.env.VOYAGEAI_API_KEY,
}); });
const EMBED_BATCH_SIZE = 128;
/** /**
* @param {Array<{id: string, text: string}>} notes * @param {Array<{id: string, text: string}>} notes
* @returns {Promise<Map<string, number[]>>} noteId → embedding vector * @returns {Promise<Map<string, number[]>>} noteId → embedding vector
*/ */
export const embedNotes = async (notes) => { export const embedNotes = async (notes) => {
const texts = notes.map((n) => n.text); const embeddingMap = new Map();
const response = await client.embed({ if (!notes || notes.length === 0) {
input: texts, return embeddingMap;
model: "voyage-3-lite", }
});
const embeddingMap = new Map(); await pipeline(
response.data.forEach((item, i) => { Readable.from(notes, { objectMode: true }),
embeddingMap.set(notes[i].id, item.embedding); batch(EMBED_BATCH_SIZE),
}); async (batches) => {
for await (const chunk of batches) {
const response = await client.embed({
input: chunk.map((n) => n.text),
model: "voyage-3",
});
return embeddingMap; response.data.forEach((item, i) => {
embeddingMap.set(chunk[i].id, item.embedding);
});
}
}
);
return embeddingMap;
}; };

View File

@@ -1,68 +1,67 @@
/** /**
* Structural validation: confirms the LLM output is well-formed * Structural validation for LLM output
* before it reaches the frontend.
* *
* @param {Array<{label: string, noteIds: string[]}>} clusters * @param {Array<{label: string, noteIds: string[]}>} clusters
* @param {string[]} inputNoteIds - the original note IDs that were sent to the LLM * @param {string[]} inputNoteIds - the original note IDs that were sent to the LLM
* @returns {{valid: boolean, reasons: string[]}} * @returns {{valid: boolean, reasons: string[]}}
*/ */
export const validateStructure = (clusters, inputNoteIds) => { export const validateStructure = (clusters, inputNoteIds) => {
const reasons = []; const reasons = [];
if (!Array.isArray(clusters) || clusters.length === 0) { if (!Array.isArray(clusters) || clusters.length === 0) {
return { valid: false, reasons: ['Response is not a non-empty array'] }; return { valid: false, reasons: ['Response is not a non-empty array'] };
}
const assignedIds = [];
for (const cluster of clusters) {
if (!cluster.label || typeof cluster.label !== 'string') {
reasons.push(`Cluster missing a valid label`);
} }
if (!Array.isArray(cluster.noteIds) || cluster.noteIds.length === 0) {
reasons.push(`Cluster "${cluster.label ?? '(unlabeled)'}" has no noteIds`); const assignedIds = [];
for (const cluster of clusters) {
if (!cluster.label || typeof cluster.label !== 'string') {
reasons.push(`Cluster missing a valid label`);
}
if (!Array.isArray(cluster.noteIds) || cluster.noteIds.length === 0) {
reasons.push(`Cluster "${cluster.label ?? '(unlabeled)'}" has no noteIds`);
}
assignedIds.push(...(cluster.noteIds ?? []));
} }
assignedIds.push(...(cluster.noteIds ?? []));
}
const inputSet = new Set(inputNoteIds); const inputSet = new Set(inputNoteIds);
const assignedSet = new Set(assignedIds); const assignedSet = new Set(assignedIds);
if (assignedIds.length !== assignedSet.size) { if (assignedIds.length !== assignedSet.size) {
reasons.push('One or more notes appear in multiple clusters'); reasons.push('One or more notes appear in multiple clusters');
} }
const missing = inputNoteIds.filter((id) => !assignedSet.has(id)); const missing = inputNoteIds.filter((id) => !assignedSet.has(id));
if (missing.length > 0) { if (missing.length > 0) {
reasons.push(`Notes missing from clusters: ${missing.join(', ')}`); reasons.push(`Notes missing from clusters: ${missing.join(', ')}`);
} }
const extra = assignedIds.filter((id) => !inputSet.has(id)); const extra = assignedIds.filter((id) => !inputSet.has(id));
if (extra.length > 0) { if (extra.length > 0) {
reasons.push(`Unknown noteIds in clusters: ${[...new Set(extra)].join(', ')}`); reasons.push(`Unknown noteIds in clusters: ${[...new Set(extra)].join(', ')}`);
} }
if (clusters.length > inputNoteIds.length) { if (clusters.length > inputNoteIds.length) {
reasons.push(`More clusters (${clusters.length}) than notes (${inputNoteIds.length})`); reasons.push(`More clusters (${clusters.length}) than notes (${inputNoteIds.length})`);
} }
return { valid: reasons.length === 0, reasons }; return { valid: reasons.length === 0, reasons };
}; };
const cosineSimilarity = (a, b) => { const cosineSimilarity = (a, b) => {
let dot = 0; let dot = 0;
let magA = 0; let magA = 0;
let magB = 0; let magB = 0;
for (let i = 0; i < a.length; i++) { for (let i = 0; i < a.length; i++) {
dot += a[i] * b[i]; dot += a[i] * b[i];
magA += a[i] * a[i]; magA += a[i] * a[i];
magB += b[i] * b[i]; magB += b[i] * b[i];
} }
const denom = Math.sqrt(magA) * Math.sqrt(magB); const denom = Math.sqrt(magA) * Math.sqrt(magB);
return denom === 0 ? 0 : dot / denom; return denom === 0 ? 0 : dot / denom;
}; };
/** /**
* Computes a silhouette-style cohesion score for the clustering. * Computes silhouette-style cohesion score for the clustering.
* *
* For each note, measures how much more similar it is to its own cluster * For each note, measures how much more similar it is to its own cluster
* versus the nearest neighboring cluster. Returns a score in [-1, 1] * versus the nearest neighboring cluster. Returns a score in [-1, 1]
@@ -73,57 +72,57 @@ const cosineSimilarity = (a, b) => {
* @returns {number} average silhouette score * @returns {number} average silhouette score
*/ */
export const computeCohesionScore = (clusters, embeddingMap) => { export const computeCohesionScore = (clusters, embeddingMap) => {
if (clusters.length <= 1) return 1.0; if (clusters.length <= 1) return 1.0;
const scores = []; const scores = [];
for (let ci = 0; ci < clusters.length; ci++) { for (let ci = 0; ci < clusters.length; ci++) {
const clusterIds = clusters[ci].noteIds; const clusterIds = clusters[ci].noteIds;
if (clusterIds.length <= 1) { if (clusterIds.length <= 1) {
scores.push(0); scores.push(0);
continue; continue;
}
for (const noteId of clusterIds) {
const vec = embeddingMap.get(noteId);
if (!vec) continue;
// a(i): avg distance to other notes in same cluster
let intraSum = 0;
let intraCount = 0;
for (const otherId of clusterIds) {
if (otherId === noteId) continue;
const otherVec = embeddingMap.get(otherId);
if (!otherVec) continue;
intraSum += 1 - cosineSimilarity(vec, otherVec);
intraCount++;
}
const a = intraCount > 0 ? intraSum / intraCount : 0;
// b(i): min avg distance to notes in any other cluster
let b = Infinity;
for (let oi = 0; oi < clusters.length; oi++) {
if (oi === ci) continue;
const otherClusterIds = clusters[oi].noteIds;
let interSum = 0;
let interCount = 0;
for (const otherId of otherClusterIds) {
const otherVec = embeddingMap.get(otherId);
if (!otherVec) continue;
interSum += 1 - cosineSimilarity(vec, otherVec);
interCount++;
}
if (interCount > 0) {
b = Math.min(b, interSum / interCount);
}
}
if (b === Infinity) b = 0;
const max = Math.max(a, b);
scores.push(max === 0 ? 0 : (b - a) / max);
}
} }
for (const noteId of clusterIds) { if (scores.length === 0) return 0;
const vec = embeddingMap.get(noteId); return scores.reduce((sum, s) => sum + s, 0) / scores.length;
if (!vec) continue;
// a(i): avg distance to other notes in same cluster
let intraSum = 0;
let intraCount = 0;
for (const otherId of clusterIds) {
if (otherId === noteId) continue;
const otherVec = embeddingMap.get(otherId);
if (!otherVec) continue;
intraSum += 1 - cosineSimilarity(vec, otherVec);
intraCount++;
}
const a = intraCount > 0 ? intraSum / intraCount : 0;
// b(i): min avg distance to notes in any other cluster
let b = Infinity;
for (let oi = 0; oi < clusters.length; oi++) {
if (oi === ci) continue;
const otherClusterIds = clusters[oi].noteIds;
let interSum = 0;
let interCount = 0;
for (const otherId of otherClusterIds) {
const otherVec = embeddingMap.get(otherId);
if (!otherVec) continue;
interSum += 1 - cosineSimilarity(vec, otherVec);
interCount++;
}
if (interCount > 0) {
b = Math.min(b, interSum / interCount);
}
}
if (b === Infinity) b = 0;
const max = Math.max(a, b);
scores.push(max === 0 ? 0 : (b - a) / max);
}
}
if (scores.length === 0) return 0;
return scores.reduce((sum, s) => sum + s, 0) / scores.length;
}; };

View File

@@ -1,18 +1,18 @@
import { describe, it, expect, vi, beforeEach } from 'vitest'; import { describe, it, expect, vi, beforeEach } from 'vitest';
const { createMock, mockEmbeddings } = vi.hoisted(() => { const { streamMock, mockEmbeddings } = vi.hoisted(() => {
const embeddings = new Map([ const embeddings = new Map([
['note_001', [1.0, 0.0, 0.0]], ['note_001', [1.0, 0.0, 0.0]],
['note_002', [0.0, 1.0, 0.0]], ['note_002', [0.0, 1.0, 0.0]],
]); ]);
return { createMock: vi.fn(), mockEmbeddings: embeddings }; return { streamMock: vi.fn(), mockEmbeddings: embeddings };
}); });
vi.mock('@anthropic-ai/sdk', () => { vi.mock('@anthropic-ai/sdk', () => {
return { return {
default: class MockAnthropic { default: class MockAnthropic {
constructor() { constructor() {
this.messages = { create: createMock }; this.messages = { stream: streamMock };
} }
}, },
}; };
@@ -34,6 +34,40 @@ const MOCK_CLUSTERS = [
{ label: 'Export Issues', noteIds: ['note_002'] }, { label: 'Export Issues', noteIds: ['note_002'] },
]; ];
// Splits text into several text_delta events so the service is exercised
// against a genuinely incremental stream rather than one whole payload.
const textEvents = (text, pieces = 4) => {
const size = Math.max(1, Math.ceil(text.length / pieces));
const events = [];
for (let i = 0; i < text.length; i += size) {
events.push({
type: 'content_block_delta',
delta: { type: 'text_delta', text: text.slice(i, i + size) },
});
}
return events;
};
const mockStreamOf = (text) => {
const events = [
{ type: 'message_start' },
...textEvents(text),
{ type: 'message_stop' },
];
streamMock.mockImplementation(() => ({
async *[Symbol.asyncIterator]() {
for (const event of events) yield event;
},
}));
};
const mockStreamThrowing = (err) => {
streamMock.mockImplementation(() => ({
async *[Symbol.asyncIterator]() {
throw err;
},
}));
};
describe('clusterNotes service', () => { describe('clusterNotes service', () => {
@@ -41,27 +75,32 @@ describe('clusterNotes service', () => {
vi.clearAllMocks(); vi.clearAllMocks();
}); });
it('should call Anthropic messages.create with the correct model', async () => { it('should call Anthropic messages.stream with the correct model', async () => {
createMock.mockResolvedValue({ mockStreamOf(JSON.stringify(MOCK_CLUSTERS));
content: [{ type: 'text', text: JSON.stringify(MOCK_CLUSTERS) }],
});
await clusterNotes(MOCK_NOTES); await clusterNotes(MOCK_NOTES);
expect(createMock).toHaveBeenCalledOnce(); expect(streamMock).toHaveBeenCalledOnce();
const callArgs = createMock.mock.calls[0][0]; const callArgs = streamMock.mock.calls[0][0];
expect(callArgs.model).toBe('claude-sonnet-4-20250514'); expect(callArgs.model).toBe('claude-sonnet-5');
expect(callArgs.max_tokens).toBe(4096); expect(callArgs.max_tokens).toBe(4096);
}); });
it('should forward an abort signal to the LLM request', async () => {
mockStreamOf(JSON.stringify(MOCK_CLUSTERS));
const controller = new AbortController();
await clusterNotes(MOCK_NOTES, { signal: controller.signal });
expect(streamMock.mock.calls[0][1]).toEqual({ signal: controller.signal });
});
it('should include all note texts in prompt sent to the LLM API', async () => { it('should include all note texts in prompt sent to the LLM API', async () => {
createMock.mockResolvedValue({ mockStreamOf(JSON.stringify(MOCK_CLUSTERS));
content: [{ type: 'text', text: JSON.stringify(MOCK_CLUSTERS) }],
});
await clusterNotes(MOCK_NOTES); await clusterNotes(MOCK_NOTES);
const prompt = createMock.mock.calls[0][0].messages[0].content; const prompt = streamMock.mock.calls[0][0].messages[0].content;
expect(prompt).toContain('note_001'); expect(prompt).toContain('note_001');
expect(prompt).toContain('Login is broken'); expect(prompt).toContain('Login is broken');
expect(prompt).toContain('note_002'); expect(prompt).toContain('note_002');
@@ -69,9 +108,7 @@ describe('clusterNotes service', () => {
}); });
it('should return clusters and a cohesion score', async () => { it('should return clusters and a cohesion score', async () => {
createMock.mockResolvedValue({ mockStreamOf(JSON.stringify(MOCK_CLUSTERS));
content: [{ type: 'text', text: JSON.stringify(MOCK_CLUSTERS) }],
});
const result = await clusterNotes(MOCK_NOTES); const result = await clusterNotes(MOCK_NOTES);
@@ -81,22 +118,35 @@ describe('clusterNotes service', () => {
expect(result.score).toBeLessThanOrEqual(1); expect(result.score).toBeLessThanOrEqual(1);
}); });
it('should reassemble clusters split across many stream deltas', async () => {
const json = JSON.stringify(MOCK_CLUSTERS);
streamMock.mockImplementation(() => ({
async *[Symbol.asyncIterator]() {
for (const char of json) {
yield { type: 'content_block_delta', delta: { type: 'text_delta', text: char } };
}
},
}));
const result = await clusterNotes(MOCK_NOTES);
expect(result.clusters).toEqual(MOCK_CLUSTERS);
});
it('should throw error when the API returns non-JSON', async () => { it('should throw error when the API returns non-JSON', async () => {
createMock.mockResolvedValue({ mockStreamOf('An unknown error occured when generting structured response.');
content: [{ type: 'text', text: 'An unknown error occured when generting structured response.' }],
});
await expect(clusterNotes(MOCK_NOTES)).rejects.toThrow('non-JSON response'); await expect(clusterNotes(MOCK_NOTES)).rejects.toThrow('non-JSON response');
}); });
it('should throw error when the API response has no text content', async () => { it('should throw error when the API response has no text content', async () => {
createMock.mockResolvedValue({ content: [] }); mockStreamOf('');
await expect(clusterNotes(MOCK_NOTES)).rejects.toThrow('no text content returned'); await expect(clusterNotes(MOCK_NOTES)).rejects.toThrow('no text content returned');
}); });
it('should throw error when Anthropic API authentication fails', async () => { it('should throw error when Anthropic API authentication fails', async () => {
createMock.mockRejectedValue(new Error('401 Unauthorized')); mockStreamThrowing(new Error('401 Unauthorized'));
await expect(clusterNotes(MOCK_NOTES)).rejects.toThrow('401 Unauthorized'); await expect(clusterNotes(MOCK_NOTES)).rejects.toThrow('401 Unauthorized');
}); });
@@ -105,9 +155,7 @@ describe('clusterNotes service', () => {
const incompleteClusters = [ const incompleteClusters = [
{ label: 'Auth Issues', noteIds: ['note_001'] }, { label: 'Auth Issues', noteIds: ['note_001'] },
]; ];
createMock.mockResolvedValue({ mockStreamOf(JSON.stringify(incompleteClusters));
content: [{ type: 'text', text: JSON.stringify(incompleteClusters) }],
});
await expect(clusterNotes(MOCK_NOTES)).rejects.toThrow('Cluster validation failed'); await expect(clusterNotes(MOCK_NOTES)).rejects.toThrow('Cluster validation failed');
}); });

View File

@@ -1,16 +1,18 @@
import { describe, it, expect, vi, beforeEach } from 'vitest'; import { describe, it, expect, vi, beforeEach } from 'vitest';
import request from 'supertest'; import request from 'supertest';
import { Readable } from 'node:stream';
import app from '../app.js'; import app from '../app.js';
vi.mock('../db/notes.dao.js', () => ({ vi.mock('../db/notes.dao.js', () => ({
getAllNotes: vi.fn(), getAllNotes: vi.fn(),
streamAllNotes: vi.fn(),
})); }));
vi.mock('../services/clustering.service.js', () => ({ vi.mock('../services/clustering.service.js', () => ({
clusterNotes: vi.fn(), clusterNotes: vi.fn(),
})); }));
import { getAllNotes } from '../db/notes.dao.js'; import { getAllNotes, streamAllNotes } from '../db/notes.dao.js';
import { clusterNotes } from '../services/clustering.service.js'; import { clusterNotes } from '../services/clustering.service.js';
const MOCK_NOTES = [ const MOCK_NOTES = [
@@ -24,6 +26,8 @@ const MOCK_CLUSTERS = [
{ label: 'Export Problems', noteIds: ['note_003'] }, { label: 'Export Problems', noteIds: ['note_003'] },
]; ];
const rowStream = (rows) => Readable.from(rows, { objectMode: true });
describe('GET /v1/notes', () => { describe('GET /v1/notes', () => {
beforeEach(() => { beforeEach(() => {
@@ -31,7 +35,7 @@ describe('GET /v1/notes', () => {
}); });
it('should return 200 and an array of notes', async () => { it('should return 200 and an array of notes', async () => {
getAllNotes.mockResolvedValue(MOCK_NOTES); streamAllNotes.mockResolvedValue(rowStream(MOCK_NOTES));
const res = await request(app).get('/v1/notes'); const res = await request(app).get('/v1/notes');
@@ -40,8 +44,26 @@ describe('GET /v1/notes', () => {
expect(Array.isArray(res.body)).toBe(true); expect(Array.isArray(res.body)).toBe(true);
}); });
it('should send JSON incrementally rather than buffering the row set', async () => {
streamAllNotes.mockResolvedValue(rowStream(MOCK_NOTES));
const res = await request(app).get('/v1/notes');
expect(res.headers['content-type']).toMatch(/application\/json/);
expect(res.headers['content-length']).toBeUndefined();
});
it('should return an empty array when there are no notes', async () => {
streamAllNotes.mockResolvedValue(rowStream([]));
const res = await request(app).get('/v1/notes');
expect(res.status).toBe(200);
expect(res.body).toEqual([]);
});
it('should return notes with expected properties', async () => { it('should return notes with expected properties', async () => {
getAllNotes.mockResolvedValue(MOCK_NOTES); streamAllNotes.mockResolvedValue(rowStream(MOCK_NOTES));
const res = await request(app).get('/v1/notes'); const res = await request(app).get('/v1/notes');
const note = res.body[0]; const note = res.body[0];
@@ -55,7 +77,7 @@ describe('GET /v1/notes', () => {
}); });
it('should return 500 when the database query fails', async () => { it('should return 500 when the database query fails', async () => {
getAllNotes.mockRejectedValue(new Error('connection refused')); streamAllNotes.mockRejectedValue(new Error('connection refused'));
const res = await request(app).get('/v1/notes'); const res = await request(app).get('/v1/notes');
@@ -63,6 +85,19 @@ describe('GET /v1/notes', () => {
expect(res.body).toHaveProperty('error'); expect(res.body).toHaveProperty('error');
expect(res.body.error).toBe('Failed to load notes'); expect(res.body.error).toBe('Failed to load notes');
}); });
it('should abort the response when the row stream fails mid-flight', async () => {
const failing = new Readable({
objectMode: true,
read() {
this.push(MOCK_NOTES[0]);
this.destroy(new Error('connection lost'));
},
});
streamAllNotes.mockResolvedValue(failing);
await expect(request(app).get('/v1/notes')).rejects.toThrow();
});
}); });
describe('POST /v1/notes/cluster', () => { describe('POST /v1/notes/cluster', () => {
@@ -101,7 +136,22 @@ describe('POST /v1/notes/cluster', () => {
await request(app).post('/v1/notes/cluster'); await request(app).post('/v1/notes/cluster');
expect(clusterNotes).toHaveBeenCalledOnce(); expect(clusterNotes).toHaveBeenCalledOnce();
expect(clusterNotes).toHaveBeenCalledWith(MOCK_NOTES); expect(clusterNotes).toHaveBeenCalledWith(
MOCK_NOTES,
expect.objectContaining({ signal: expect.any(AbortSignal) })
);
});
it('should pass a signal that is not aborted while the request is open', async () => {
getAllNotes.mockResolvedValue(MOCK_NOTES);
clusterNotes.mockImplementation(async (_notes, { signal }) => {
expect(signal.aborted).toBe(false);
return MOCK_CLUSTERS;
});
const res = await request(app).post('/v1/notes/cluster');
expect(res.status).toBe(200);
}); });
it('should return 500 when clusterNotes (API call) fails', async () => { it('should return 500 when clusterNotes (API call) fails', async () => {

View File

@@ -1,14 +1,23 @@
import { describe, it, expect, vi, beforeEach } from 'vitest'; import { describe, it, expect, vi, beforeEach } from 'vitest';
import { Readable } from 'node:stream';
const { mockQuery } = vi.hoisted(() => ({ const { mockQuery, mockConnect } = vi.hoisted(() => ({
mockQuery: vi.fn(), mockQuery: vi.fn(),
mockConnect: vi.fn(),
})); }));
vi.mock('../db/index.js', () => ({ vi.mock('../db/index.js', () => ({
query: mockQuery, query: mockQuery,
getPool: () => ({ connect: mockConnect }),
})); }));
import { getAllNotes, getNoteById, createNote, createNotes } from '../db/notes.dao.js'; import {
getAllNotes,
streamAllNotes,
getNoteById,
createNote,
createNotes,
} from '../db/notes.dao.js';
const MOCK_ROWS = [ const MOCK_ROWS = [
{ id: 'note_001', text: 'Login flow feels confusing', x: 193, y: 191, author: 'user_5', color: 'yellow' }, { id: 'note_001', text: 'Login flow feels confusing', x: 193, y: 191, author: 'user_5', color: 'yellow' },
@@ -48,6 +57,58 @@ describe('notes.dao', () => {
}); });
}); });
describe('streamAllNotes', () => {
const mockClient = (rows) => {
const release = vi.fn();
const client = {
release,
query: vi.fn(() => Readable.from(rows, { objectMode: true })),
};
mockConnect.mockResolvedValue(client);
return { client, release };
};
it('should stream rows without buffering them into an array', async () => {
mockClient(MOCK_ROWS);
const stream = await streamAllNotes();
const received = [];
for await (const row of stream) received.push(row);
expect(received).toEqual(MOCK_ROWS);
});
it('should release the pooled client once the stream ends', async () => {
const { release } = mockClient(MOCK_ROWS);
const stream = await streamAllNotes();
for await (const _row of stream) { /* drain */ }
expect(release).toHaveBeenCalledOnce();
});
it('should release the pooled client when a consumer destroys the stream early', async () => {
const { release } = mockClient(MOCK_ROWS);
const stream = await streamAllNotes();
stream.destroy();
await new Promise((resolve) => stream.once('close', resolve));
expect(release).toHaveBeenCalledOnce();
});
it('should release the pooled client when starting the query throws', async () => {
const release = vi.fn();
mockConnect.mockResolvedValue({
release,
query: vi.fn(() => { throw new Error('cursor failed'); }),
});
await expect(streamAllNotes()).rejects.toThrow('cursor failed');
expect(release).toHaveBeenCalledOnce();
});
});
describe('getNoteById', () => { describe('getNoteById', () => {
it('should return a single note when found', async () => { it('should return a single note when found', async () => {
mockQuery.mockResolvedValue({ rows: [MOCK_ROWS[0]] }); mockQuery.mockResolvedValue({ rows: [MOCK_ROWS[0]] });
@@ -109,5 +170,31 @@ describe('notes.dao', () => {
expect(sql).toContain('INSERT INTO notes'); expect(sql).toContain('INSERT INTO notes');
expect(params).toHaveLength(12); expect(params).toHaveLength(12);
}); });
it('should return an empty array without querying when given no notes', async () => {
const result = await createNotes([]);
expect(result).toEqual([]);
expect(mockQuery).not.toHaveBeenCalled();
});
it('should split large inputs into multiple statements under the bind-parameter limit', async () => {
const many = Array.from({ length: 2500 }, (_, i) => ({
id: `note_${i}`,
text: `text ${i}`,
author: 'user_1',
}));
mockQuery.mockImplementation(async (_sql, params) => ({
rows: new Array(params.length / 6).fill(null).map((_, i) => ({ i })),
}));
const result = await createNotes(many);
expect(mockQuery).toHaveBeenCalledTimes(3);
expect(result).toHaveLength(2500);
for (const [, params] of mockQuery.mock.calls) {
expect(params.length).toBeLessThan(65535);
}
});
}); });
}); });

View File

@@ -0,0 +1,82 @@
import { describe, it, expect } from 'vitest';
import { Readable } from 'node:stream';
import { pipeline } from 'node:stream/promises';
import { batch, jsonArray } from '../lib/streams.js';
const collect = async (source, transform) => {
const out = [];
await pipeline(source, transform, async (results) => {
for await (const item of results) out.push(item);
});
return out;
};
describe('batch', () => {
it('should group items into fixed-size arrays', async () => {
const source = Readable.from([1, 2, 3, 4], { objectMode: true });
const result = await collect(source, batch(2));
expect(result).toEqual([[1, 2], [3, 4]]);
});
it('should flush a partial trailing batch', async () => {
const source = Readable.from([1, 2, 3, 4, 5], { objectMode: true });
const result = await collect(source, batch(2));
expect(result).toEqual([[1, 2], [3, 4], [5]]);
});
it('should emit nothing for an empty source', async () => {
const source = Readable.from([], { objectMode: true });
const result = await collect(source, batch(3));
expect(result).toEqual([]);
});
it('should reject a non-positive size', () => {
expect(() => batch(0)).toThrow(TypeError);
expect(() => batch(1.5)).toThrow(TypeError);
});
});
describe('jsonArray', () => {
const serialize = async (items) => {
const chunks = await collect(
Readable.from(items, { objectMode: true }),
jsonArray()
);
return chunks.map(String).join('');
};
it('should serialize objects into a JSON array', async () => {
const items = [{ id: 'a' }, { id: 'b' }];
const output = await serialize(items);
expect(output).toBe('[{"id":"a"},{"id":"b"}]');
expect(JSON.parse(output)).toEqual(items);
});
it('should emit an empty array when the source yields nothing', async () => {
const output = await serialize([]);
expect(output).toBe('[]');
expect(JSON.parse(output)).toEqual([]);
});
it('should emit a valid single-element array', async () => {
const output = await serialize([{ id: 'only' }]);
expect(JSON.parse(output)).toEqual([{ id: 'only' }]);
});
it('should propagate serialization errors', async () => {
const circular = {};
circular.self = circular;
await expect(serialize([circular])).rejects.toThrow();
});
});

View File

@@ -109,6 +109,9 @@ const Stickies = () => {
<span className="cluster-rank" aria-label={`Priority ${group.rank}`}> <span className="cluster-rank" aria-label={`Priority ${group.rank}`}>
{group.rank} {group.rank}
</span> </span>
{group.rank === 1 && (
<span className="cluster-reorder-hint">Drag and drop to reorganize cluster priority</span>
)}
<h3 className="cluster-label">{group.label}</h3> <h3 className="cluster-label">{group.label}</h3>
<span className="cluster-drag-handle" aria-hidden="true">⠿</span> <span className="cluster-drag-handle" aria-hidden="true">⠿</span>
</div> </div>

View File

@@ -60,6 +60,14 @@
flex-shrink: 0; flex-shrink: 0;
} }
.cluster-reorder-hint {
font-size: 0.8em;
color: #9ca3af;
font-style: italic;
white-space: nowrap;
flex-shrink: 0;
}
.cluster-label { .cluster-label {
margin: 0; margin: 0;
font-size: 1.2em; font-size: 1.2em;

21
package-lock.json generated Normal file
View File

@@ -0,0 +1,21 @@
{
"name": "kongruity",
"lockfileVersion": 3,
"requires": true,
"packages": {
"": {
"dependencies": {
"split2": "^4.2.0"
}
},
"node_modules/split2": {
"version": "4.2.0",
"resolved": "https://registry.npmjs.org/split2/-/split2-4.2.0.tgz",
"integrity": "sha512-UcjcJOWknrNkF6PLX83qcHM6KHgVKNkV62Y8a5uYDVv9ydGQVwAHMKqHdJje1VTWpljG0WYpCDhrCdAOYH4TWg==",
"license": "ISC",
"engines": {
"node": ">= 10.x"
}
}
}
}

5
package.json Normal file
View File

@@ -0,0 +1,5 @@
{
"dependencies": {
"split2": "^4.2.0"
}
}