diff --git a/README.md b/README.md index aa0a0c0..dcc4019 100644 --- a/README.md +++ b/README.md @@ -55,7 +55,13 @@ EOF 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: diff --git a/backend/db/fixtures/notes.jsonl b/backend/db/fixtures/notes.jsonl new file mode 100644 index 0000000..7494c54 --- /dev/null +++ b/backend/db/fixtures/notes.jsonl @@ -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"} diff --git a/backend/db/notes.dao.js b/backend/db/notes.dao.js index 1818e83..5f79f64 100644 --- a/backend/db/notes.dao.js +++ b/backend/db/notes.dao.js @@ -1,12 +1,49 @@ -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, so stay well under it. +const INSERT_BATCH_SIZE = 1000; export const getAllNotes = async () => { - const { rows } = await query( - 'SELECT id, text, x, y, author, color FROM notes ORDER BY id' - ); + const { rows } = await query(SELECT_NOTES); return rows; }; +/** + * Streams every note as an object-mode Readable. The pooled client is checked + * out for the life of the stream and released once it ends, errors, or is + * destroyed early by a consumer. + * + * @returns {Promise} + */ +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) => { const { rows } = await query( 'SELECT id, text, x, y, author, color FROM notes WHERE id = $1', @@ -25,7 +62,7 @@ export const createNote = async (note) => { return rows[0]; }; -export const createNotes = async (notes) => { +const insertNoteBatch = async (notes) => { const values = []; const placeholders = []; @@ -52,3 +89,23 @@ export const createNotes = async (notes) => { ); return rows; }; + +export const createNotes = async (notes) => { + if (!notes || notes.length === 0) { + return []; + } + + const inserted = []; + + await pipeline( + Readable.from(notes, { objectMode: true }), + batch(INSERT_BATCH_SIZE), + async (batches) => { + for await (const chunk of batches) { + inserted.push(...await insertNoteBatch(chunk)); + } + } + ); + + return inserted; +}; diff --git a/backend/db/seed.js b/backend/db/seed.js index c7cf7e4..8c8071d 100644 --- a/backend/db/seed.js +++ b/backend/db/seed.js @@ -1,85 +1,78 @@ 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 = [ - { 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 FIXTURE = fileURLToPath(new URL('./fixtures/notes.jsonl', import.meta.url)); -const run = async () => { - try { - const insertQuery = ` - INSERT INTO notes (id, text, x, y, author, color) - VALUES ($1, $2, $3, $4, $5, $6) - ON CONFLICT (id) DO NOTHING - `; +const COLUMNS = ['id', 'text', 'x', 'y', 'author', 'color']; - let inserted = 0; - for (const note of notes) { - const result = await query(insertQuery, [ - note.id, - note.text, - note.x, - note.y, - note.author, - note.color, - ]); - inserted += result.rowCount; +const DEFAULTS = { x: 0, y: 0, color: 'yellow' }; + +const csvField = (value) => { + if (value === null || value === undefined) return ''; + return `"${String(value).replaceAll('"', '""')}"`; +}; + +const toCsvRows = () => new Transform({ + writableObjectMode: true, + transform(line, _encoding, callback) { + 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) { + await client.query('ROLLBACK').catch(() => {}); console.error('Seed failed:', err.message); - process.exit(1); + process.exitCode = 1; } finally { + client.release(); await close(); } }; diff --git a/backend/lib/streams.js b/backend/lib/streams.js new file mode 100644 index 0000000..5cf35f5 --- /dev/null +++ b/backend/lib/streams.js @@ -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 ? ']' : '[]'); + }, + }); +}; diff --git a/backend/package-lock.json b/backend/package-lock.json index 08b4700..e4d5be4 100644 --- a/backend/package-lock.json +++ b/backend/package-lock.json @@ -13,6 +13,8 @@ "dotenv": "^16.4.7", "express": "^4.21.2", "pg": "^8.18.0", + "pg-copy-streams": "^7.0.0", + "pg-query-stream": "^4.16.0", "voyageai": "^0.1.0" }, "devDependencies": { @@ -2339,6 +2341,21 @@ "integrity": "sha512-kecgoJwhOpxYU21rZjULrmrBJ698U2RxXofKVzOn5UDj61BPj/qMb7diYUR1nLScCDbrztQFl1TaQZT0t1EtzQ==", "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": { "version": "1.0.1", "resolved": "https://registry.npmjs.org/pg-int8/-/pg-int8-1.0.1.tgz", @@ -2363,6 +2380,18 @@ "integrity": "sha512-pfsxk2M9M3BuGgDOfuy37VNRRX3jmKgMjcvAcWqNDpZSf4cUmv8HSOl5ViRQFsfARFn0KuUQTgLxVMbNq5NW3g==", "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": { "version": "2.2.0", "resolved": "https://registry.npmjs.org/pg-types/-/pg-types-2.2.0.tgz", diff --git a/backend/package.json b/backend/package.json index 481200d..4b43856 100644 --- a/backend/package.json +++ b/backend/package.json @@ -17,6 +17,8 @@ "dotenv": "^16.4.7", "express": "^4.21.2", "pg": "^8.18.0", + "pg-copy-streams": "^7.0.0", + "pg-query-stream": "^4.16.0", "voyageai": "^0.1.0" }, "devDependencies": { diff --git a/backend/routes/notes.routes.js b/backend/routes/notes.routes.js index b2c7232..069ebb9 100644 --- a/backend/routes/notes.routes.js +++ b/backend/routes/notes.routes.js @@ -1,25 +1,38 @@ 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 { jsonArray } from '../lib/streams.js'; const router = Router(); router.get('/', async (req, res) => { try { - const notes = await getAllNotes(); - res.json(notes); + const rows = await streamAllNotes(); + res.type('application/json'); + await pipeline(rows, jsonArray(), res); } catch (err) { console.error(`Error loading notes: ${err}`); + if (res.headersSent) { + res.destroy(err); + return; + } res.status(500).json({ error: 'Failed to load notes' }); } }); router.post('/cluster', async (req, res) => { + const controller = new AbortController(); + res.on('close', () => { + if (!res.writableEnded) controller.abort(); + }); + try { const notes = await getAllNotes(); - const result = await clusterNotes(notes); + const result = await clusterNotes(notes, { signal: controller.signal }); res.json(result); } catch (err) { + if (controller.signal.aborted) return; console.error(`Clustering failed: ${err}`); res.status(500).json({ error: 'Clustering failed' }); } diff --git a/backend/services/clustering.service.js b/backend/services/clustering.service.js index cdc0af5..8064d19 100644 --- a/backend/services/clustering.service.js +++ b/backend/services/clustering.service.js @@ -1,4 +1,6 @@ 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 { validateStructure, computeCohesionScore } from "./validation.service.js"; @@ -33,31 +35,71 @@ Here are the notes: ${notesJson}`; }; -const requestClusters = async (notes) => { - const response = await client.messages.create({ - model: "claude-sonnet-4-20250514", - max_tokens: 4096, - messages: [ - { role: "user", content: buildPrompt(notes) }, - ], - }); +const textDeltas = () => new Transform({ + objectMode: true, + transform(event, _encoding, callback) { + if (event?.type === 'content_block_delta' && event.delta?.type === 'text_delta') { + callback(null, event.delta.text); + 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-4-20250514", + 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'); } try { - return JSON.parse(textBlock.text); + return JSON.parse(text); } catch { 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([ - requestClusters(notes), + requestClusters(notes, signal), embedNotes(notes), ]); diff --git a/backend/services/embedding.service.js b/backend/services/embedding.service.js index 29c3318..ffb21e1 100644 --- a/backend/services/embedding.service.js +++ b/backend/services/embedding.service.js @@ -1,25 +1,42 @@ 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({ apiKey: process.env.VOYAGEAI_API_KEY, }); +// Well under Voyage's per-request input and token ceilings. +const EMBED_BATCH_SIZE = 128; + /** * @param {Array<{id: string, text: string}>} notes * @returns {Promise>} noteId → embedding vector */ export const embedNotes = async (notes) => { - const texts = notes.map((n) => n.text); - - const response = await client.embed({ - input: texts, - model: "voyage-3", - }); - const embeddingMap = new Map(); - response.data.forEach((item, i) => { - embeddingMap.set(notes[i].id, item.embedding); - }); + + if (!notes || notes.length === 0) { + return embeddingMap; + } + + await pipeline( + Readable.from(notes, { objectMode: true }), + 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", + }); + + response.data.forEach((item, i) => { + embeddingMap.set(chunk[i].id, item.embedding); + }); + } + } + ); return embeddingMap; }; diff --git a/backend/tests/Clustering.service.test.js b/backend/tests/Clustering.service.test.js index 139cd3e..fbbf454 100644 --- a/backend/tests/Clustering.service.test.js +++ b/backend/tests/Clustering.service.test.js @@ -1,18 +1,18 @@ import { describe, it, expect, vi, beforeEach } from 'vitest'; -const { createMock, mockEmbeddings } = vi.hoisted(() => { +const { streamMock, mockEmbeddings } = vi.hoisted(() => { const embeddings = new Map([ ['note_001', [1.0, 0.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', () => { return { default: class MockAnthropic { constructor() { - this.messages = { create: createMock }; + this.messages = { stream: streamMock }; } }, }; @@ -34,6 +34,40 @@ const MOCK_CLUSTERS = [ { 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', () => { @@ -41,27 +75,32 @@ describe('clusterNotes service', () => { vi.clearAllMocks(); }); - it('should call Anthropic messages.create with the correct model', async () => { - createMock.mockResolvedValue({ - content: [{ type: 'text', text: JSON.stringify(MOCK_CLUSTERS) }], - }); + it('should call Anthropic messages.stream with the correct model', async () => { + mockStreamOf(JSON.stringify(MOCK_CLUSTERS)); await clusterNotes(MOCK_NOTES); - expect(createMock).toHaveBeenCalledOnce(); - const callArgs = createMock.mock.calls[0][0]; + expect(streamMock).toHaveBeenCalledOnce(); + const callArgs = streamMock.mock.calls[0][0]; expect(callArgs.model).toBe('claude-sonnet-4-20250514'); 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 () => { - createMock.mockResolvedValue({ - content: [{ type: 'text', text: JSON.stringify(MOCK_CLUSTERS) }], - }); + mockStreamOf(JSON.stringify(MOCK_CLUSTERS)); 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('Login is broken'); expect(prompt).toContain('note_002'); @@ -69,9 +108,7 @@ describe('clusterNotes service', () => { }); it('should return clusters and a cohesion score', async () => { - createMock.mockResolvedValue({ - content: [{ type: 'text', text: JSON.stringify(MOCK_CLUSTERS) }], - }); + mockStreamOf(JSON.stringify(MOCK_CLUSTERS)); const result = await clusterNotes(MOCK_NOTES); @@ -81,22 +118,35 @@ describe('clusterNotes service', () => { 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 () => { - createMock.mockResolvedValue({ - content: [{ type: 'text', text: 'An unknown error occured when generting structured response.' }], - }); + mockStreamOf('An unknown error occured when generting structured response.'); await expect(clusterNotes(MOCK_NOTES)).rejects.toThrow('non-JSON response'); }); 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'); }); 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'); }); @@ -105,9 +155,7 @@ describe('clusterNotes service', () => { const incompleteClusters = [ { label: 'Auth Issues', noteIds: ['note_001'] }, ]; - createMock.mockResolvedValue({ - content: [{ type: 'text', text: JSON.stringify(incompleteClusters) }], - }); + mockStreamOf(JSON.stringify(incompleteClusters)); await expect(clusterNotes(MOCK_NOTES)).rejects.toThrow('Cluster validation failed'); }); diff --git a/backend/tests/Notes.test.js b/backend/tests/Notes.test.js index 47a1bfd..d8fe67b 100644 --- a/backend/tests/Notes.test.js +++ b/backend/tests/Notes.test.js @@ -1,16 +1,18 @@ import { describe, it, expect, vi, beforeEach } from 'vitest'; import request from 'supertest'; +import { Readable } from 'node:stream'; import app from '../app.js'; vi.mock('../db/notes.dao.js', () => ({ getAllNotes: vi.fn(), + streamAllNotes: vi.fn(), })); vi.mock('../services/clustering.service.js', () => ({ 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'; const MOCK_NOTES = [ @@ -24,6 +26,8 @@ const MOCK_CLUSTERS = [ { label: 'Export Problems', noteIds: ['note_003'] }, ]; +const rowStream = (rows) => Readable.from(rows, { objectMode: true }); + describe('GET /v1/notes', () => { beforeEach(() => { @@ -31,7 +35,7 @@ describe('GET /v1/notes', () => { }); 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'); @@ -40,8 +44,26 @@ describe('GET /v1/notes', () => { 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 () => { - getAllNotes.mockResolvedValue(MOCK_NOTES); + streamAllNotes.mockResolvedValue(rowStream(MOCK_NOTES)); const res = await request(app).get('/v1/notes'); const note = res.body[0]; @@ -55,7 +77,7 @@ describe('GET /v1/notes', () => { }); 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'); @@ -63,6 +85,19 @@ describe('GET /v1/notes', () => { expect(res.body).toHaveProperty('error'); 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', () => { @@ -101,7 +136,22 @@ describe('POST /v1/notes/cluster', () => { await request(app).post('/v1/notes/cluster'); 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 () => { diff --git a/backend/tests/notes.dao.test.js b/backend/tests/notes.dao.test.js index 8bda4ba..10e708c 100644 --- a/backend/tests/notes.dao.test.js +++ b/backend/tests/notes.dao.test.js @@ -1,14 +1,23 @@ 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(), + mockConnect: vi.fn(), })); vi.mock('../db/index.js', () => ({ 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 = [ { 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', () => { it('should return a single note when found', async () => { mockQuery.mockResolvedValue({ rows: [MOCK_ROWS[0]] }); @@ -109,5 +170,31 @@ describe('notes.dao', () => { expect(sql).toContain('INSERT INTO notes'); 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); + } + }); }); }); diff --git a/backend/tests/streams.test.js b/backend/tests/streams.test.js new file mode 100644 index 0000000..5edb932 --- /dev/null +++ b/backend/tests/streams.test.js @@ -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(); + }); +}); diff --git a/package-lock.json b/package-lock.json new file mode 100644 index 0000000..9b63f7d --- /dev/null +++ b/package-lock.json @@ -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" + } + } + } +} diff --git a/package.json b/package.json new file mode 100644 index 0000000..e9ac23b --- /dev/null +++ b/package.json @@ -0,0 +1,5 @@ +{ + "dependencies": { + "split2": "^4.2.0" + } +}