Compare commits
4 Commits
90035b3568
...
e11122ed14
| Author | SHA1 | Date | |
|---|---|---|---|
| e11122ed14 | |||
|
|
7b47e46852 | ||
|
|
62477d009c | ||
|
|
ee6fa9e576 |
@@ -55,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:
|
||||||
|
|
||||||
|
|||||||
50
backend/db/fixtures/notes.jsonl
Normal file
50
backend/db/fixtures/notes.jsonl
Normal 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"}
|
||||||
@@ -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;
|
||||||
};
|
};
|
||||||
|
|||||||
@@ -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
72
backend/lib/streams.js
Normal 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 ? ']' : '[]');
|
||||||
|
},
|
||||||
|
});
|
||||||
|
};
|
||||||
29
backend/package-lock.json
generated
29
backend/package-lock.json
generated
@@ -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",
|
||||||
|
|||||||
@@ -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": {
|
||||||
|
|||||||
@@ -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' });
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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),
|
||||||
]);
|
]);
|
||||||
|
|
||||||
|
|||||||
@@ -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",
|
}
|
||||||
});
|
|
||||||
|
|
||||||
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;
|
||||||
};
|
};
|
||||||
|
|||||||
@@ -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;
|
|
||||||
};
|
};
|
||||||
|
|||||||
@@ -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');
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -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 () => {
|
||||||
|
|||||||
@@ -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);
|
||||||
|
}
|
||||||
|
});
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|||||||
82
backend/tests/streams.test.js
Normal file
82
backend/tests/streams.test.js
Normal 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();
|
||||||
|
});
|
||||||
|
});
|
||||||
21
package-lock.json
generated
Normal file
21
package-lock.json
generated
Normal 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
5
package.json
Normal file
@@ -0,0 +1,5 @@
|
|||||||
|
{
|
||||||
|
"dependencies": {
|
||||||
|
"split2": "^4.2.0"
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user