src/watch.js (6699 bytes)
1 // The watcher: the cheapest model triages new ideas and pairs each one with its project, in the 2 // background. No process stays alive between passes, because a process like that can stop (a 3 // crash, a reboot, a full context, a usage limit). Instead, the UserPromptSubmit hook calls kick() 4 // for each prompt, and kick() starts a pass when there is work. Thus the watcher goes on while it 5 // is on, and it costs nothing while no new ideas come in. 6 7 import { spawn } from 'node:child_process'; 8 import fs from 'node:fs'; 9 import path from 'node:path'; 10 import { fileURLToPath } from 'node:url'; 11 import { stamp } from './render.js'; 12 import { home, lane, load } from './store.js'; 13 14 const MODEL = 'haiku'; 15 const RETRY_WAITS_MS = [30 * 1000, 2 * 60 * 1000]; // after a failed call, in the same pass 16 const WAIT_AFTER_ERROR_MS = 10 * 60 * 1000; // after a failed pass, before kick() starts a new one 17 const LOCK_STALE_MS = 30 * 60 * 1000; // a pass is much shorter: an older lock is from a crashed pass 18 const MAX_ROUNDS = 5; // triage calls in one pass, 20 ideas each 19 const LOG_LINES = 200; 20 21 const statePath = () => path.join(home(), 'watch.json'); 22 const logPath = () => path.join(home(), 'watch.log'); 23 const lockPath = () => path.join(home(), '.watch'); 24 const sleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms)); 25 26 export function readState() { 27 try { 28 return JSON.parse(fs.readFileSync(statePath(), 'utf8')); 29 } catch { 30 return { on: false }; 31 } 32 } 33 34 function writeState(patch) { 35 fs.mkdirSync(home(), { recursive: true }); 36 const state = { ...readState(), ...patch }; 37 fs.writeFileSync(statePath(), JSON.stringify(state, null, 2) + '\n'); 38 return state; 39 } 40 41 export function turnOn() { 42 const state = readState(); 43 return state.on ? state : writeState({ on: true, since: new Date().toISOString(), error: null, errorAt: null }); 44 } 45 46 export function turnOff() { 47 return writeState({ on: false }); 48 } 49 50 /** Ids that the watcher must triage: the inbox, and open ideas that were triaged before pairing. */ 51 export function pendingWork(db) { 52 return db.ideas.filter((i) => lane(i) === 'inbox' || (['do', 'maybe'].includes(lane(i)) && !i.triage?.paired)).map((i) => i.id); 53 } 54 55 function isRunning(now) { 56 try { 57 return now - fs.statSync(lockPath()).mtimeMs < LOCK_STALE_MS; 58 } catch { 59 return false; 60 } 61 } 62 63 /** Run a pass in a separate process, which goes on after the caller (the hook) exits. */ 64 function startPassProcess() { 65 const bin = fileURLToPath(new URL('../bin/ideamine.js', import.meta.url)); 66 const child = spawn(process.execPath, [bin, 'watch-pass'], { cwd: home(), detached: true, stdio: 'ignore', windowsHide: true }); 67 // A pass that cannot start must not stop the caller, for example the server of the dashboard. 68 child.on('error', (e) => writeState({ error: `cannot start a pass: ${e.message}`, errorAt: new Date().toISOString() })); 69 child.unref(); 70 } 71 72 /** Start a pass if the watcher is on, there is work, no pass runs, and no pass failed a short time ago. */ 73 export function kick({ startPass = startPassProcess, now = Date.now() } = {}) { 74 const state = readState(); 75 if (!state.on) return false; 76 if (state.errorAt && now - Date.parse(state.errorAt) < WAIT_AFTER_ERROR_MS) return false; 77 if (isRunning(now) || !pendingWork(load()).length) return false; 78 startPass(); 79 return true; 80 } 81 82 function readLog() { 83 try { 84 return fs.readFileSync(logPath(), 'utf8').split('\n').filter(Boolean); 85 } catch { 86 return []; 87 } 88 } 89 90 function log(text) { 91 const lines = [...readLog(), `${stamp(new Date().toISOString())} ${text}`]; 92 fs.writeFileSync(logPath(), lines.slice(-LOG_LINES).join('\n') + '\n'); 93 } 94 95 function describe(out) { 96 const ok = out.results.filter((r) => !r.error); 97 const moved = ok.filter((r) => r.moved).map((r) => `#${r.id} → ${path.basename(r.moved)}`); 98 const tokens = `${out.tokens.input} in / ${out.tokens.output} out tokens`; 99 return `triaged ${ok.map((r) => `#${r.id}`).join(' ')} (${out.model}, ${tokens})${moved.length ? ` · paired ${moved.join(', ')}` : ''}`; 100 } 101 102 async function withRetries(fn, retries) { 103 for (let attempt = 0; ; attempt++) { 104 try { 105 return await fn(); 106 } catch (e) { 107 if (attempt >= retries) throw e; 108 await sleep(RETRY_WAITS_MS[Math.min(attempt, RETRY_WAITS_MS.length - 1)]); 109 } 110 } 111 } 112 113 /** One pass: triage and pair all pending ideas with the cheapest model. Only one pass runs at a time. */ 114 export async function pass({ model = MODEL, retries = RETRY_WAITS_MS.length } = {}) { 115 fs.mkdirSync(home(), { recursive: true }); 116 try { 117 fs.mkdirSync(lockPath()); 118 } catch (e) { 119 if (e.code !== 'EEXIST' || isRunning(Date.now())) return 'busy'; 120 fs.utimesSync(lockPath(), new Date(), new Date()); // take over the lock of a crashed pass 121 } 122 try { 123 // Loaded here, not at the top: the hook imports this module for each prompt. 124 const { headlessTriage } = await import('./claude.js'); 125 for (let round = 0; round < MAX_ROUNDS; round++) { 126 const ids = pendingWork(load()).slice(0, 20); 127 if (!ids.length) break; 128 const out = await withRetries(() => headlessTriage({ model, ids }), retries); 129 if (out.message) break; // the ideas were deleted in the meantime 130 log(describe(out)); 131 // An idea that the triage cannot sort must not start a new pass for each prompt. 132 const stuck = pendingWork(load()).filter((id) => ids.includes(id)); 133 if (stuck.length) throw new Error(`the triage did not sort ${stuck.map((id) => `#${id}`).join(' ')}`); 134 } 135 writeState({ lastPass: new Date().toISOString(), error: null, errorAt: null }); 136 return 'done'; 137 } catch (e) { 138 log(`error: ${e.message}`); 139 writeState({ error: e.message, errorAt: new Date().toISOString() }); 140 return 'error'; 141 } finally { 142 fs.rmSync(lockPath(), { recursive: true, force: true }); 143 } 144 } 145 146 /** On or off, what waits, and the last passes. */ 147 export function status() { 148 const state = readState(); 149 if (!state.on) return 'ideamine watch: off. /ideas-watch turns it on.'; 150 const waiting = pendingWork(load()).length; 151 const lines = [`ideamine watch: on since ${stamp(state.since)}. ${MODEL} triages new ideas and pairs each one with its project.`]; 152 lines.push(waiting ? `${waiting} idea${waiting === 1 ? ' waits' : 's wait'} for the watcher.` : 'No idea waits.'); 153 if (isRunning(Date.now())) lines.push('A pass runs now.'); 154 if (state.error) lines.push(`The last pass failed: ${state.error}. The watcher tries again 10 minutes after the failure.`); 155 const recent = readLog().slice(-5); 156 if (recent.length) lines.push('last passes:', ...recent.map((line) => ` ${line}`)); 157 lines.push('/ideas-watch off turns it off.'); 158 return lines.join('\n'); 159 }