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 }