Recently Written · git

ideamine

An idea inbox for Claude Code: /idea saves ideas at zero tokens; Claude triages them and routes each to the cheapest model that can build it.

git clone https://github.com/equwal/ideamine

Log | Files | Refs


src/sync.js (14986 bytes)

1 // Sync with an ideamine server, so that every Claude on every machine sees one archive. The server
2 // (`ideamine serve`, for example behind nginx on a WireGuard address) holds the archive, and
3 // ideas.json on this machine is a copy of it. A change goes into the outbox first and then to the
4 // server. When the server does not answer, the change waits in the outbox, so no idea is lost. The
5 // prompt log takes the same way. No process stays alive: the hook calls kick() for each prompt, and
6 // kick() starts a background sync when work waits or the copy is old.
7 
8 import { spawn } from 'node:child_process';
9 import crypto from 'node:crypto';
10 import fs from 'node:fs';
11 import os from 'node:os';
12 import path from 'node:path';
13 import { fileURLToPath } from 'node:url';
14 import * as config from './config.js';
15 import { stamp } from './render.js';
16 import { dbPath, home, withLock, writeAtomic } from './store.js';
17 import { clip } from './text.js';
18 
19 const PULL_EVERY_MS = 60 * 1000; // a copy older than this is pulled again in the background
20 const PROMPTS_EVERY_MS = 20 * 1000; // prompts wait at least this long, so that a prompt starts no process each time
21 const WAIT_AFTER_ERROR_MS = 2 * 60 * 1000; // after a failed sync, before kick() tries again
22 const LOCK_STALE_MS = 10 * 60 * 1000;
23 const OPS_PER_REQUEST = 100;
24 const PROMPT_BYTES_PER_REQUEST = 256 * 1024;
25 const MAX_PROMPT_CHARS = 100 * 1000; // a longer paste is cut, so that each prompt fits in a request
26 
27 const outboxPath = () => path.join(home(), 'outbox.jsonl');
28 const promptsPath = () => path.join(home(), 'prompts-outbox.jsonl');
29 const statePath = () => path.join(home(), 'sync.json');
30 const lockPath = () => path.join(home(), '.sync');
31 
32 /** The server does not answer, or it refuses the request. The change stays in the outbox. */
33 export class SyncError extends Error {}
34 
35 /** The address of the server, with a slash at the end, or '' when sync is off. */
36 export function serverUrl() {
37   const url = config.get('sync_url');
38   return url ? (url.endsWith('/') ? url : `${url}/`) : '';
39 }
40 
41 export const enabled = () => serverUrl() !== '';
42 export const logsPrompts = () => enabled() && /^(on|true|yes|1)$/i.test(config.get('prompt_log'));
43 
44 function readLines(file) {
45   let raw;
46   try {
47     raw = fs.readFileSync(file, 'utf8');
48   } catch {
49     return [];
50   }
51   const out = [];
52   for (const line of raw.split('\n')) {
53     if (!line.trim()) continue;
54     try {
55       out.push(JSON.parse(line));
56     } catch {
57       // A line that a crash cut in half. The rest of the file is still good.
58     }
59   }
60   return out;
61 }
62 
63 const writeLines = (file, items) => writeAtomic(file, items.map((i) => `${JSON.stringify(i)}\n`).join(''));
64 const append = (file, items) => withLock(() => fs.appendFileSync(file, items.map((i) => `${JSON.stringify(i)}\n`).join('')));
65 const size = (file) => fs.statSync(file, { throwIfNoEntry: false })?.size || 0;
66 
67 export function readState() {
68   try {
69     return JSON.parse(fs.readFileSync(statePath(), 'utf8'));
70   } catch {
71     return {};
72   }
73 }
74 
75 // Call it while you hold the lock.
76 const saveState = (patch) => writeAtomic(statePath(), `${JSON.stringify({ ...readState(), ...patch }, null, 2)}\n`);
77 const writeState = (patch) => withLock(() => saveState(patch));
78 const failed = (e) => writeState({ error: e.message, errorAt: new Date().toISOString() });
79 
80 async function request(route, { body, timeoutMs }) {
81   const target = new URL(route, serverUrl());
82   const init = { signal: AbortSignal.timeout(timeoutMs) };
83   if (body !== undefined) Object.assign(init, { method: 'POST', headers: { 'content-type': 'application/json' }, body: JSON.stringify(body) });
84   let res;
85   try {
86     res = await fetch(target, init);
87   } catch (e) {
88     const why = e.name === 'TimeoutError' ? `no answer in ${timeoutMs / 1000} s` : e.cause?.code || e.cause?.message || e.message;
89     throw new SyncError(`cannot reach ${target.origin} (${why})`);
90   }
91   const text = await res.text();
92   let json = null;
93   try {
94     json = JSON.parse(text);
95   } catch {
96     // Not JSON: the error below shows the text.
97   }
98   if (!res.ok || !json?.ok) throw new SyncError(`${target.origin} answered ${res.status}: ${json?.error || clip(text, 160)}`);
99   return json;
100 }
101 
102 /** Keep the archive of the server as the copy here, unless the copy is newer: a slow answer must not undo a later one. */
103 function keep(db) {
104   if (!db || !Array.isArray(db.ideas)) throw new SyncError('the server sent no archive');
105   withLock(() => {
106     const state = readState();
107     if (db.uid === state.uid && (Number(db.rev) || 0) < (Number(state.rev) || 0)) return;
108     writeAtomic(dbPath(), `${JSON.stringify(db, null, 2)}\n`, { backup: true });
109     saveState({ uid: db.uid, rev: db.rev, pulled: new Date().toISOString(), error: null, errorAt: null });
110   });
111 }
112 
113 /**
114  * Send the changes of the outbox, oldest first. `extra` must go even when another process took it
115  * from the outbox already: the server applies a change only once, and it answers with the result.
116  * Resolves to the results by oid, and to the archive of the server after the last change.
117  */
118 async function sendOps({ extra = null, timeoutMs }) {
119   const ops = withLock(() => readLines(outboxPath()));
120   if (extra && !ops.some((op) => op.oid === extra.oid)) ops.push(extra);
121   const results = new Map();
122   let db = null;
123   for (let i = 0; i < ops.length; i += OPS_PER_REQUEST) {
124     const part = ops.slice(i, i + OPS_PER_REQUEST);
125     const out = await request('api/ops', { body: { ops: part }, timeoutMs });
126     part.forEach((op, k) => results.set(op.oid, out.results[k]));
127     withLock(() => writeLines(outboxPath(), readLines(outboxPath()).filter((op) => !results.has(op.oid))));
128     db = out.db;
129   }
130   return { results, db };
131 }
132 
133 async function sendPrompts({ timeoutMs }) {
134   const prompts = withLock(() => readLines(promptsPath()));
135   let part = [];
136   let bytes = 0;
137   const sendPart = async () => {
138     if (!part.length) return;
139     await request('api/prompts', { body: { prompts: part }, timeoutMs });
140     const sent = new Set(part.map((p) => p.id));
141     withLock(() => writeLines(promptsPath(), readLines(promptsPath()).filter((p) => !sent.has(p.id))));
142     part = [];
143     bytes = 0;
144   };
145   for (const p of prompts) {
146     const n = Buffer.byteLength(JSON.stringify(p));
147     if (bytes + n > PROMPT_BYTES_PER_REQUEST) await sendPart();
148     part.push(p);
149     bytes += n;
150   }
151   await sendPart();
152 }
153 
154 /**
155  * Bring the copy here up to date: send the changes that wait, then take the archive of the server.
156  * With `prompts`, send the prompts that wait too. Throws SyncError when the server does not answer.
157  */
158 export async function pull({ timeoutMs = 10000, prompts = false } = {}) {
159   let { db } = await sendOps({ timeoutMs });
160   if (prompts) await sendPrompts({ timeoutMs });
161   db ||= (await request('api/db', { timeoutMs })).db;
162   keep(db);
163 }
164 
165 /**
166  * Change the archive on the server. The change goes into the outbox, then to the server. Resolves
167  * to { value }, the result of the change, or to { queued, reason } when the server does not answer:
168  * the change then waits in the outbox. A change that the server refuses throws.
169  */
170 export async function change(op, { timeoutMs = 3000 } = {}) {
171   const item = { ...op, oid: crypto.randomUUID(), at: new Date().toISOString() };
172   append(outboxPath(), [item]);
173   let sent;
174   try {
175     sent = await sendOps({ extra: item, timeoutMs });
176   } catch (e) {
177     if (!(e instanceof SyncError)) throw e;
178     failed(e);
179     return { queued: true, reason: e.message };
180   }
181   keep(sent.db);
182   const r = sent.results.get(item.oid);
183   if (!r?.ok) throw new Error(r?.error || 'the server sent no result');
184   return { value: r.value };
185 }
186 
187 /** How many changes wait in the outbox. */
188 export const waiting = () => readLines(outboxPath()).length;
189 
190 /**
191  * The text of a prompt for the log, or null when nothing is left. The notes that Claude Code and the
192  * desktop app put into a prompt (<system-reminder> blocks) go, and a huge paste is cut.
193  */
194 function promptText(raw) {
195   const text = String(raw ?? '').replace(/<system-reminder>[\s\S]*?<\/system-reminder>/g, '').trim();
196   if (!text) return null;
197   return text.length > MAX_PROMPT_CHARS ? `${text.slice(0, MAX_PROMPT_CHARS)}\n… (cut: ${text.length - MAX_PROMPT_CHARS} more characters)` : text;
198 }
199 
200 /** Put a prompt into the prompt outbox. The next background sync sends it. */
201 export function logPrompt({ prompt, session = null, cwd = null, at = new Date().toISOString() }) {
202   const text = promptText(prompt);
203   if (!text) return;
204   append(promptsPath(), [{ id: crypto.randomUUID(), at, host: os.hostname(), session, cwd, prompt: text }]);
205   if (!readState().promptsSince) writeState({ promptsSince: at });
206 }
207 
208 // The text that the user typed, from one line of a Claude Code transcript. Null for everything
209 // else: tool results, subagent prompts, command output, and notes that Claude Code adds itself.
210 function typedText(entry) {
211   if (entry?.type !== 'user' || entry.isMeta || entry.isSidechain || entry.toolUseResult) return null;
212   const content = entry.message?.content;
213   let text = null;
214   if (typeof content === 'string') text = content;
215   else if (Array.isArray(content) && content.every((b) => b?.type === 'text' || b?.type === 'image')) {
216     text = content.filter((b) => b.type === 'text').map((b) => b.text).join('\n');
217   }
218   if (!text?.trim()) return null;
219   const name = text.match(/<command-name>([^<]*)<\/command-name>/);
220   if (name) {
221     const args = text.match(/<command-args>([\s\S]*?)<\/command-args>/);
222     return `${name[1]} ${args ? args[1] : ''}`.trim();
223   }
224   if (/^(<(local-command-stdout|local-command-stderr|command-message|bash-input|bash-stdout|bash-stderr)>|\[Request interrupted)/.test(text.trim())) return null;
225   return text;
226 }
227 
228 /** The transcripts folder of Claude Code on this machine. */
229 export function transcriptsDir() {
230   return path.join(process.env.CLAUDE_CONFIG_DIR || path.join(os.homedir(), '.claude'), 'projects');
231 }
232 
233 /**
234  * Prompts from the Claude Code transcripts of this machine, for the prompt log. Only prompts from
235  * before the live prompt log started are taken, so that no prompt comes twice. An id made from the
236  * session, the time, and the text lets the server skip a prompt that an earlier import sent.
237  */
238 export function importPrompts({ dir = transcriptsDir() } = {}) {
239   const before = readState().promptsSince || null;
240   const found = [];
241   let files = [];
242   try {
243     files = fs.readdirSync(dir, { recursive: true }).filter((f) => String(f).endsWith('.jsonl'));
244   } catch {
245     return 0;
246   }
247   for (const file of files) {
248     let raw;
249     try {
250       raw = fs.readFileSync(path.join(dir, String(file)), 'utf8');
251     } catch {
252       continue;
253     }
254     for (const line of raw.split('\n')) {
255       if (!line.includes('"user"')) continue;
256       let entry;
257       try {
258         entry = JSON.parse(line);
259       } catch {
260         continue;
261       }
262       const prompt = promptText(typedText(entry));
263       if (!prompt || !entry.timestamp || (before && entry.timestamp >= before)) continue;
264       const id = crypto.createHash('sha256').update(`${entry.sessionId}\n${entry.timestamp}\n${prompt}`).digest('hex').slice(0, 32);
265       found.push({ id, at: entry.timestamp, host: os.hostname(), session: entry.sessionId || null, cwd: entry.cwd || null, prompt });
266     }
267   }
268   if (found.length) append(promptsPath(), found);
269   return found.length;
270 }
271 
272 function isRunning(now) {
273   try {
274     return now - fs.statSync(lockPath()).mtimeMs < LOCK_STALE_MS;
275   } catch {
276     return false;
277   }
278 }
279 
280 function startProcess() {
281   const bin = fileURLToPath(new URL('../bin/ideamine.js', import.meta.url));
282   const child = spawn(process.execPath, [bin, 'sync', '--background'], { cwd: home(), detached: true, stdio: 'ignore', windowsHide: true });
283   // A sync that cannot start must not stop the caller.
284   child.on('error', (e) => failed(new Error(`cannot start a sync: ${e.message}`)));
285   child.unref();
286 }
287 
288 /**
289  * Start a background sync when sync is on, no sync runs, no sync failed a short time ago, and
290  * changes wait, prompts wait since the last sync 20 seconds ago, or the copy is older than a minute.
291  */
292 export function kick({ start = startProcess, now = Date.now() } = {}) {
293   if (!enabled()) return false;
294   const state = readState();
295   if (state.errorAt && now - Date.parse(state.errorAt) < WAIT_AFTER_ERROR_MS) return false;
296   if (isRunning(now)) return false;
297   const age = state.pulled ? now - Date.parse(state.pulled) : Infinity;
298   const due = age > PULL_EVERY_MS || size(outboxPath()) > 0 || (size(promptsPath()) > 0 && age > PROMPTS_EVERY_MS);
299   if (!due) return false;
300   start();
301   return true;
302 }
303 
304 /** One background sync, from kick(). Only one runs at a time. */
305 export async function backgroundSync() {
306   fs.mkdirSync(home(), { recursive: true });
307   try {
308     fs.mkdirSync(lockPath());
309   } catch (e) {
310     if (e.code !== 'EEXIST' || isRunning(Date.now())) return 'busy';
311     fs.utimesSync(lockPath(), new Date(), new Date()); // take over the lock of a crashed sync
312   }
313   try {
314     await pull({ timeoutMs: 30000, prompts: true });
315     return 'done';
316   } catch (e) {
317     failed(e);
318     return 'error';
319   } finally {
320     fs.rmSync(lockPath(), { recursive: true, force: true });
321   }
322 }
323 
324 /**
325  * Turn sync on with the server at `url`, or off with ''. Before the first sync with a server, the
326  * archive here goes to a backup file, because the archive of the server takes its place.
327  */
328 export function setServer(url) {
329   const clean = String(url || '').trim();
330   if (clean && !/^https?:\/\//i.test(clean)) throw new Error(`"${clean}" is not an http:// or https:// address`);
331   let backup = null;
332   if (clean && fs.existsSync(dbPath()) && !readState().uid) {
333     backup = path.join(home(), `ideas.before-sync-${new Date().toISOString().slice(0, 10)}.json`);
334     fs.copyFileSync(dbPath(), backup);
335   }
336   config.set('sync_url', clean);
337   if (clean) config.set('publish_url', ''); // the server shows the live archive: no uploads
338   else writeState({ uid: null, rev: null, pulled: null, error: null, errorAt: null });
339   return backup;
340 }
341 
342 /** On or off, the last sync, and what waits. */
343 export function status() {
344   const url = serverUrl();
345   if (!url) return 'ideamine sync: off. The archive is on this machine. /ideas-sync <url> shares it through an ideamine server.';
346   const s = readState();
347   const lines = [`ideamine sync: ${url}`];
348   lines.push(s.pulled ? `last sync ${stamp(s.pulled)}${s.rev ? `, revision ${s.rev}` : ''}.` : 'Not synced yet.');
349   const changes = waiting();
350   if (changes) lines.push(`${changes} change${changes === 1 ? ' waits' : 's wait'} for the server.`);
351   const prompts = readLines(promptsPath()).length;
352   if (prompts) lines.push(`${prompts} prompt${prompts === 1 ? ' waits' : 's wait'} for the server.`);
353   lines.push(logsPrompts() ? 'The prompt log is on.' : 'The prompt log is off. `ideamine config prompt_log on` turns it on.');
354   if (s.error) lines.push(`The last sync failed at ${stamp(s.errorAt)}: ${s.error}`);
355   return lines.join('\n');
356 }