antcolony
All repositories: gitoria
42.4 KB
// daemon.hl — `run`: the scheduler as a long-running loop (mission 020, tickets antcolony#1/#2, build order 4// part 2; concept docs/scheduler-agent.md §2 leases, §3 quota / parking / ports). Still ONE host, no WebSocket.//// Every --interval s (and once at the start) ONE iteration:// 1. survey (next.hl) + the lease of every `in progress` ticket (lease.hl): an expired colony lease → the// ticket back to `open` with a comment; a parked session of THIS host whose resume time is over → a resume// candidate; a live lease → skipped (never a second session on a leased ticket, also after a restart);// 2. candidates = due resumes first, then the eligible open tickets whose dev folder is on this host (oldest// first) and not in use by a running session of this run;// 3. free slots = --parallel − running; nothing to start or no slot → no quota probe;// 4. quota: `claude -p /usage` (quota.hl, costs nothing) — at or above --reserve % of the 5 h window → no start;// above --week-limit % (default 84, creator 2026-09-24) of the WEEKLY usage, or weekly unknown → no start;// 100 turns either check off (--reserve 100: never blocks; --week-limit 100: never blocks, also when unknown);// 5. start: a port range per session from --port-pool (--ports-per-session each, disjoint while running) →// `work` (lease, worker, controller, post by the verdict) or `runResume`;// 6. ONE line: `ITER <n> <time> · quota … · running … · started / resumed / expired … · skipped: …`.// Mission 031: step 1 and the resume in 5 act on what the survey READ, seconds (or a librarian pass) earlier — so both// re-read the ticket right before they write (lease.hl expireLease / stillLeased) and write nothing unless it is STILL// `in progress` with the same lease; else `NOT expired …` / `NOT resumed …: re-read right before — <why> → nothing written`.// A lease that has ENDED (the ticket left `in progress`, a `colony-deployed:` line, the session's own report …) is// listed `… the lease of S has ended (…) → never expired`.// A finished session prints `SESSION END <ref> <session> · <outcome> · cost …`.// STOP: the `colony` wrapper turns SIGTERM/SIGINT into the file $COLONY_STOP_FILE ('finish' / 'park'); the run// checks it twice a second. First signal: start nothing more, let the running sessions finish, then exit.// Second signal: kill the running sessions' processes and PARK them (resumed by the next `run`). Never silent.// `--once`: one iteration, wait for what it started, exit. Refuses the LIVE tickets server without `--live`.// Mission 029 (antcolony#21): every iteration STARTS with the LIBRARIAN (librarian.hl), awaited: it reads the creator's new// ticket texts and files decisions (no model call when nothing is new; quiet then); only afterwards the iteration picks and// starts — a pass that failed or left a text unfiled starts nothing. `--librarian off` / COLONY_LIBRARIAN=off switches it// off; `--librarian-model`, `--librarian-max` (texts per pass, 0 = all).// Mission 034 (antcolony#25): with the librarian on, the QUOTA is read before it (`claude -p /usage`, $0): at or above the// start threshold (a reached limit = 100 % always) → no librarian call, no start. A model call that still hits the usage// limit (the librarian's, or a worker / controller parked on it) → `LIMIT:` pause until the reset (15 min when unknown):// no librarian, no probe, no start; one line when it begins, one when it ends (README "Usage limit").import { env } from 'hl:proc'import { exists, readFile, mkDir } from 'hl:fs'import { every, now, timestamp } from 'hl:time'import { survey } from './next.hl'import { baseUrl, readToken, getJson, ticketPath } from './tickets.hl'import { runWork, runResume, hostName, runsDir } from './work.hl'import { setting, isWholeNumber, round4 } from './claude.hl'import { leaseOf, expireLease, stillLeased, leaseSettings } from './lease.hl'import { probeUsage } from './quota.hl'import { portRange, newSession } from './brief.hl'import { librarianPass, librarianSettings } from './librarian.hl'import { keywordPass } from './keywords.hl'import { digestTick } from './status.hl'import { questionWhy } from './relations.hl'import { buildBrief, PORTS_MARKER } from './brief.hl'import { readAgentToken } from './tickets.hl'import { jsonErrorAt } from './jsoncheck.hl'import { parseCapabilities, missingCapabilities } from './util.hl'import AgentServer from './AgentServer.hl'// options `run` hands to every `work` / resume unchangedstatic PASS = ['--model' '--max-turns' '--permission-mode' '--timeout' '--max-budget-usd' '--claude' '--controller-claude' '--controller-model' '--controller-max-turns' '--controller-timeout' '--max-fails' '--lease-ttl' '--heartbeat' '--finalize-max-turns' '--finalize-timeout']static LIVE_HOST = 'tickets.worldapi.org'// mission 034: a usage limit without a known reset time pauses the librarian + starts this long, then tries oncestatic LIMIT_RETRY_MS = 900000// the host part of a URL, lower casestatic hostOf = (url) => {let u = url.toLowerCase()let i = u.indexOf('://')if (i >= 0) { u = u.slice(i + 3) }let j = u.indexOf('/')if (j >= 0) { u = u.slice(0, j) }let k = u.indexOf(':')if (k >= 0) { u = u.slice(0, k) }return u}static passOptions = (argv) => {let out = []let i = 0while (i < argv.length) {if (PASS.includes(argv[i]) && i + 1 < argv.length) {out.push(argv[i])out.push(argv[i + 1])i = i + 2} else { i = i + 1 }}return out}static money = (n) => {if (n == null) { return '–' }return '$' + round4(n)}static ago = (ms) => { return timestamp(ms).slice(0, 19) + 'Z' }// `leases [--expire]`: every `in progress` ticket with its lease (a look at what `run` would do)static runLeases = (argv) => {let expire = argv.includes('--expire')let tk = nullif (expire) {if (hostOf(baseUrl()) == LIVE_HOST && !argv.includes('--live')) {console.log('leases: REFUSED — ' + baseUrl() + ' is the LIVE tickets server; --expire writes, pass --live to do it there')return false}tk = readToken()if (tk.error != null) {console.log('leases: REFUSED — ' + tk.error + ' (--expire writes)')return false}}let sv = survey()if (sv.error != null) {console.log('leases: ' + sv.error)return false}let n = 0let t0 = now()for (e of sv.notEligible) {let t = e.tif (t.state == 'in progress') {n = n + 1let ref = t.project + t.reflet d = getJson(ticketPath(t.project, t.number))if (d.status != 200) { console.log('LEASE ' + ref + ' — cannot read: ' + d.status) } else {let l = leaseOf(d.json.events)if (l.kind == 'ready') { console.log('LEASE ' + ref + ' — built, waiting for the deploy (session ' + l.session + '): never expires — ./colony deployed ' + t.project + ' ' + t.number + ' after the deploy') }else if (l.kind == 'ended') { console.log('LEASE ' + ref + ' — the lease of session ' + l.session + ' has ended (' + l.text + '): never expires') }else if (l.kind != 'colony') { console.log('LEASE ' + ref + ' — no colony lease ("' + l.text + '"): never expires') }else {let st = l.parked ? 'parked (' + l.phase + ') until ' + ago(l.resumeMs) + ', ' : ''if (t0 < l.expiryMs) { console.log('LEASE ' + ref + ' — session ' + l.session + ' on ' + l.host + ': ' + st + 'alive until ' + ago(l.expiryMs)) }else if (!expire) { console.log('LEASE ' + ref + ' — session ' + l.session + ' on ' + l.host + ': ' + st + 'EXPIRED at ' + ago(l.expiryMs) + ' (--expire gives it back)') }else {let r = expireLease(t.project, t.number, l, tk.token)let res = ' → back to open'if (!r.ok && r.skipped == true) { res = ' → NOT expired, re-read right before: ' + r.text + ' (nothing written)' }else if (!r.ok) { res = ' → giving it back FAILED: ' + r.text }console.log('LEASE ' + ref + ' — session ' + l.session + ' on ' + l.host + ': ' + st + 'EXPIRED at ' + ago(l.expiryMs) + res)}}}}}console.log('LEASES: ' + n + ' ticket(s) in progress')return true}static runRun = (argv) => {let home = env('COLONY_HOME')if (home == null || home == '') {console.log('run: REFUSED — COLONY_HOME is not set (run through ./colony, which sets it)')return false}let base = baseUrl()let live = hostOf(base) == LIVE_HOSTif (live && !argv.includes('--live')) {console.log('run: REFUSED — ' + base + ' is the LIVE tickets server (COLONY_TICKETS_URL unset or pointing there); pass --live to run against it')return false}let capText = setting(argv, '--capabilities', 'COLONY_CAPABILITIES', '')let hostCaps = parseCapabilities(capText)let intervalText = setting(argv, '--interval', 'COLONY_RUN_INTERVAL', '60')let parallelText = setting(argv, '--parallel', 'COLONY_PARALLEL', '1')let reserveText = setting(argv, '--reserve', 'COLONY_QUOTA_RESERVE', '70')let weekText = setting(argv, '--week-limit', 'COLONY_WEEK_LIMIT', '84')let poolText = setting(argv, '--port-pool', 'COLONY_PORT_POOL', '8700-8799')let perText = setting(argv, '--ports-per-session', 'COLONY_PORTS_PER_SESSION', '10')let claudeBin = setting(argv, '--claude', 'COLONY_CLAUDE', 'claude')let quotaBin = setting(argv, '--quota-claude', 'COLONY_QUOTA_CLAUDE', claudeBin)if (!isWholeNumber(intervalText) || intervalText == '0') {console.log('run: REFUSED — --interval must be whole seconds ≥ 1, got ' + intervalText)return false}if (!isWholeNumber(parallelText) || parallelText == '0') {console.log('run: REFUSED — --parallel must be a whole number ≥ 1, got ' + parallelText)return false}if (!isWholeNumber(reserveText) || toNumber(reserveText) > 100) {console.log('run: REFUSED — --reserve must be a percentage 0–100, got ' + reserveText)return false}if (!isWholeNumber(weekText) || toNumber(weekText) > 100) {console.log('run: REFUSED — --week-limit must be a percentage 0–100, got ' + weekText)return false}// mission 035 (antcolony#14): `--serve PORT` = the scheduler proper — it starts no worker itself; agents connect to itlet serveText = setting(argv, '--serve', 'COLONY_SERVE', '')let serveHost = setting(argv, '--serve-host', 'COLONY_SERVE_HOST', '127.0.0.1')let ttlText = setting(argv, '--agent-ttl', 'COLONY_AGENT_TTL', '30')let settleText = setting(argv, '--account-settle', 'COLONY_ACCOUNT_SETTLE', '60')let serving = serveText != ''if (serving && (!isWholeNumber(serveText) || serveText == '0' || toNumber(serveText) > 65535)) {console.log('run: REFUSED — --serve must be a port 1–65535, got ' + serveText)return false}if (serving && (!isWholeNumber(ttlText) || ttlText == '0')) {console.log('run: REFUSED — --agent-ttl must be whole seconds ≥ 1, got ' + ttlText)return false}if (serving && (!isWholeNumber(settleText))) {console.log('run: REFUSED — --account-settle must be whole seconds ≥ 0, got ' + settleText)return false}if (serving && argv.includes('--once')) {console.log('run: REFUSED — --serve and --once do not go together (an agent needs the scheduler to stay)')return false}let agentToken = { token = '' }if (serving) {agentToken = readAgentToken()if (agentToken.error != null) {console.log('run: REFUSED — ' + agentToken.error + ' (agents must authenticate)')return false}}let pool = portRange(poolText)if (pool == null) {console.log('run: REFUSED — --port-pool must be A-B (1 ≤ A ≤ B ≤ 65535), got ' + poolText)return false}if (!isWholeNumber(perText) || perText == '0') {console.log('run: REFUSED — --ports-per-session must be a whole number ≥ 1, got ' + perText)return false}let interval = toNumber(intervalText)let parallel = toNumber(parallelText)let reserve = toNumber(reserveText)let weekLimit = toNumber(weekText)let per = toNumber(perText)let size = pool.to - pool.from + 1let slots = (size - size % per) / perif (!serving && slots < parallel) {console.log('run: REFUSED — the port pool ' + poolText + ' has ' + slots + ' range(s) of ' + per + ' ports, --parallel ' + parallel + ' needs ' + parallel)return false}let libOn = setting(argv, '--librarian', 'COLONY_LIBRARIAN', 'on')if (libOn != 'on' && libOn != 'off') {console.log('run: REFUSED — --librarian must be on or off, got ' + libOn)return false}let libSet = librarianSettings(argv, true)if (libSet.error != null) {console.log('run: REFUSED — ' + libSet.error)return false}let ls = leaseSettings(argv)if (ls.error != null) {console.log('run: REFUSED — ' + ls.error)return false}let tk = readToken()if (tk.error != null) {console.log('run: REFUSED — ' + tk.error + ' (run leases and posts)')return false}let host = hostName()let runs = runsDir(home)mkDir(runs, 448)let stopFile = env('COLONY_STOP_FILE')let once = argv.includes('--once')let pass = passOptions(argv)let running = []let live = {}let libBusy = falselet busy = falselet stopping = nulllet finished = falselet iter = 0let loop = nulllet watch = nullconsole.log('RUN: started on ' + host + (serving ? ' · SCHEDULER for agents on ' + serveHost + ':' + serveText + ' (starts no worker itself; agents offline after ' + ttlText + ' s without a sync)' : '') + ' · tickets ' + base + (live ? ' (LIVE)' : '') + ' · ' + (once ? 'once' : 'every ' + interval + ' s') + ' · parallel ' + parallel + ' · quota reserve ' + reserve + '% of the 5 h window, no start above ' + weekLimit + '% of the week (probe: ' + quotaBin + ' -p /usage) · ports ' + pool.from + '-' + pool.to + ', ' + per + ' per session · lease ttl ' + ls.ttl + ' s, heartbeat ' + ls.heartbeat + ' s · runs ' + runs + '/ · librarian ' + (libOn == 'on' ? 'on, first in every iteration (' + libSet.model + (libSet.max > 0 ? ', at most ' + libSet.max + ' text(s) per pass' : '') + ')' : 'off') + ' · stop file ' + (stopFile == null ? '— (no wrapper: SIGTERM kills without parking)' : stopFile))// mission 034 (antcolony#25): the Claude usage limit. `limitUntil` (epoch ms) = "limited until": set when a model call// of this run hit the limit — the librarian's (its pass reports `limited`), or a worker / controller that was parked on// it. Until then: no librarian pass, no quota probe, no start (resumes too); leases are still looked after; the// iterations are silent. ONE line when the pause begins, ONE when it ends. Unknown reset time → 15 min, then try once// (the quota is read first — below).let limitUntil = nulllet limitReset = (resetMs) => { return resetMs == null ? now() + LIMIT_RETRY_MS : resetMs }let pauseForLimit = (resetMs, why) => {let until = limitReset(resetMs)if (limitUntil == null) {limitUntil = untilconsole.log('LIMIT: usage limit reached — librarian and starts paused until ' + ago(until) + (resetMs == null ? ' (no reset time known → ' + (LIMIT_RETRY_MS / 60000) + ' min, then try once)' : '') + ' · ' + why)} else if (until > limitUntil) {limitUntil = untilconsole.log('LIMIT: the pause is extended until ' + ago(until) + ' · ' + why)}return null}// the quota reading at or above the start threshold → why (text), else null. A REACHED limit (100 %, 5 h or week)// always counts, also with --reserve 100 / --week-limit 100. An unknown week or a failed probe → null (the librarian// runs; starts are still held by iterate as before).let quotaBlock = (q) => {if (q == null || q.error != null || q.percent == null) { return null }if (q.percent >= 100) { return 'the 5 h limit is reached (100%)' }if (q.weekly != null && q.weekly >= 100) { return 'the weekly limit is reached (100%)' }if (reserve < 100 && q.percent >= reserve) { return '5 h ' + q.percent + '% ≥ reserve ' + reserve + '%' }if (weekLimit < 100 && q.weekly != null && q.weekly > weekLimit) { return 'week ' + q.weekly + '% > week limit ' + weekLimit + '%' }return null}let quotaText = (q) => { return 'quota ' + q.percent + '% of 5 h' + (q.resetsMs == null ? '' : ' (resets ' + ago(q.resetsMs) + ')') + (q.weekly == null ? '' : ', week ' + q.weekly + '%') + (q.cost != null && q.cost > 0 ? ', PROBE COST $' + q.cost : '') }// mission 035 (antcolony#14): the agents that are connected (in memory only — an agent's next sync rebuilds all of it).// agent = { host, seen, parallel, accepting, why, online }; `running` holds the sessions of the agents (entry.remote =// the agent's host): what an agent reports as running, plus a command queued / just sent that the agent has not shown yet.let agents = []let seenEvents = []let kwSeen = {}let digestState = { day = '' }let agentTtl = toNumber(ttlText) * 1000// antcolony#16: the quota belongs to the Claude ACCOUNT. Agents on one account (`--account`) report the same reading; the// scheduler pools it: freshest reading of the account, a limit pause of any of them, and one new start per account at a// time until a reading taken AFTER that start (--account-settle s later) shows what it costs. accountStart = { account: ms }.let settleMs = toNumber(settleText) * 1000let accountStart = {}let accountHosts = (acct) => {let hs = []for (a of agents) { if (agentOnline(a) && a.account == acct) { hs.push(a.host) } }return hs}let accountReading = (acct) => {let best = nullfor (a of agents) {if (agentOnline(a) && a.account == acct && a.quota != null && a.quota.at != null && (best == null || a.quota.at > best.at)) { best = a.quota }}return best}let accountBlock = (acct) => {let hs = accountHosts(acct)for (a of agents) {if (agentOnline(a) && a.account == acct && a.limitUntil != null && a.limitUntil > now()) { return 'usage limit reached until ' + ago(a.limitUntil) + ' (hit on ' + a.host + ')' }}let q = accountReading(acct)if (q == null) { return 'the quota is not read yet' }if (now() - q.at > 300000) { return 'the newest quota reading is older than 5 min' }if (q.percent >= 100) { return '5 h limit reached (100%)' }if (q.weekly != null && q.weekly >= 100) { return 'weekly limit reached (100%)' }if (reserve < 100 && q.percent >= reserve) { return 'quota 5 h ' + q.percent + '% ≥ reserve ' + reserve + '%' }if (weekLimit < 100 && (q.weekly == null || q.weekly > weekLimit)) { return q.weekly == null ? 'week UNKNOWN' : 'week ' + q.weekly + '% > week limit ' + weekLimit + '%' }let st = accountStart[acct]if (st != null && q.at < st + settleMs) { return 'a start on it (' + ago(st) + ') is not in the quota reading yet' + (hs.length > 1 ? ' — shared by ' + hs.join(', ') : '') }return null}let agentOf = (h) => {if (h == null) { return null }for (a of agents) { if (a.host.toLowerCase() == h.toLowerCase()) { return a } }return null}let agentOnline = (a) => { return a != null && now() - a.seen <= agentTtl }let hostRunning = (h) => {let n = 0for (r of running) { if (r.remote != null && r.remote.toLowerCase() == h.toLowerCase()) { n = n + 1 } }return n}// the sum of what the connected agents can run at once (the "/N" of the ITER line)let capacity = () => {if (!serving) { return parallel }let n = 0for (a of agents) { if (agentOnline(a)) { n = n + a.parallel } }return n}let dropHost = (h) => {let out = []for (x of running) { if (x.remote == null || x.remote.toLowerCase() != h.toLowerCase()) { out.push(x) } }running = outreturn null}let answer = (status, body) => { return { status = status body = body } }let handleAgent = (method, path, body, auth) => {if (path != '/agent/sync' || method != 'POST') { return answer(404, { error = 'POST /agent/sync only' }) }if (auth == null || auth != 'Bearer ' + agentToken.token) { return answer(401, { error = 'the agent token is wrong or missing' }) }if (body == null || body == '' || jsonErrorAt(body) >= 0) { return answer(400, { error = 'the body is not valid JSON' }) }let j = JSON.parse(body)if (j.host == null || hlTypeName(j.host) != 'String' || j.host == '' || j.running == null || j.events == null || j.parallel == null) { return answer(400, { error = 'host, parallel, running and events are required' }) }let a = agentOf(j.host)if (a == null) {agents.push({ host = j.host account = j.host quota = null limitUntil = null caps = [] seen = now() parallel = j.parallel accepting = j.accepting why = j.why announced = false wasAccepting = false })a = agentOf(j.host)}let wasOnline = agentOnline(a) && a.announced == truea.seen = now()a.caps = j.capabilities == null ? [] : parseCapabilities(j.capabilities.join(','))a.parallel = j.parallela.account = j.account == null || j.account == '' ? j.host : j.accounta.quota = j.quotaa.limitUntil = j.limitUntila.accepting = j.accepting == truea.why = j.whyif (!wasOnline) {a.announced = trueconsole.log('AGENT: ' + j.host + ' connected · offers ' + (a.caps.length == 0 ? '—' : a.caps.join(', ')) + ' · parallel ' + j.parallel + ' · Claude account ' + a.account + ' · ' + (a.accepting ? 'accepting work' : 'not accepting work (' + j.why + ')') + ' · running ' + j.running.length + ' · ' + j.events.length + ' buffered event(s)')} else if (a.accepting != a.wasAccepting) {console.log('AGENT: ' + j.host + (a.accepting ? ' accepts work again' : ' does not accept work: ' + j.why))}a.wasAccepting = a.acceptinglet ack = []for (e of j.events) {ack.push(e.id)if (seenEvents.includes(e.id)) { }else {seenEvents.push(e.id)if (seenEvents.length > 300) { seenEvents = seenEvents.slice(100) }if (e.kind == 'ended') { console.log('SESSION END ' + e.ref + ' ' + e.session + ' · ' + e.outcome + ' · cost ' + money(e.workerCost + e.controllerCost) + ' (worker ' + money(e.workerCost) + ' + controller ' + money(e.controllerCost) + ') · on ' + j.host + ' · ports ' + e.ports) }else if (e.kind == 'refused') { console.log('AGENT: ' + j.host + ' refused ' + e.ref + ' (' + e.session + ') — ' + e.why) }}}// rebuild this host's entries: what it runs + a command it has not shown yet (sent one round ago and still missing = lost / refused)let shown = []let commands = []let mine = []for (x of running) { if (x.remote != null && x.remote.toLowerCase() == j.host.toLowerCase()) { mine.push(x) } }dropHost(j.host)for (r of j.running) {shown.push(r.session)running.push({ ref = r.ref session = r.session folder = r.folder slot = -1 ports = r.ports phase = r.phase remote = j.host proc = null sentAt = null cmd = null })}for (x of mine) {if (shown.includes(x.session)) { }else if (x.sentAt != null || x.cmd == null) { }else {x.sentAt = now()commands.push(x.cmd)running.push(x)}}return answer(200, { ack = ack commands = commands })}let server = nullif (serving) { server = new AgentServer(port = toNumber(serveText), host = serveHost, handler = handleAgent) }let slotFree = (i) => {for (r of running) { if (r.slot == i) { return false } }return true}let freeSlot = () => {let i = 0while (i < slots) {if (slotFree(i)) { return i }i = i + 1}return -1}let rangeOf = (i) => {let a = pool.from + i * perreturn a + '-' + (a + per - 1)}let runningSession = (s) => {for (r of running) { if (r.session == s) { return true } }return false}let folderBusy = (f) => {for (r of running) { if (r.folder == f) { return true } }return false}let checkEnd = () => {if (finished || busy || libBusy || (running.length > 0 && !serving)) { return null }if (stopping == null && !(once && iter >= 1)) { return null }finished = trueif (loop != null) { loop.stop() }if (watch != null) { watch.stop() }console.log('RUN STOPPED — ' + (stopping != null ? 'stop requested (' + stopping + ')' : '--once') + ', ' + iter + ' iteration(s), ' + (serving ? running.length + ' session(s) keep running on their agents (leases in tickets)' : 'no session running'))return null}let sessionEnded = (entry, r) => {let out = []for (x of running) { if (x.session != entry.session) { out.push(x) } }running = outlet total = (r.workerCost == null ? 0 : r.workerCost) + (r.controllerCost == null ? 0 : r.controllerCost)console.log('SESSION END ' + entry.ref + ' ' + entry.session + ' · ' + r.outcome + ' · cost ' + money(total) + ' (worker ' + money(r.workerCost) + ' + controller ' + money(r.controllerCost) + ') · ports ' + entry.ports + ' freed · running ' + running.length + '/' + capacity())// mission 034: parked on the usage limit → no librarian call and no start into it until the resetif (r.limited == true) { pauseForLimit(r.resetMs, entry.ref + ' ' + entry.session + ' was parked on it (' + r.limitText + ')') }checkEnd()return null}let startOne = (c) => {let slot = freeSlot()if (slot < 0) { return 'no port range free for ' + c.ref }let ports = rangeOf(slot)let entry = { ref = c.ref session = c.session folder = c.folder slot = slot ports = ports proc = null phase = 'worker' parkReason = null }// antcolony#28: process, phase and park reason live in `live[session]`, NOT on `entry` — `running` is rebuilt by// copying its entries (a Hybriel push copies), so a park reason set on a `running` entry never reached the hooklive[c.session] = { proc = null phase = 'worker' parkReason = null }let hooks = {onProc = (p, phase) => {live[c.session].proc = plive[c.session].phase = phasereturn null}parkReason = () => { return live[c.session].parkReason }}// mission 031: a resume acts on what the survey saw — re-read first; still the same parked lease, or nothing is writtenif (c.kind == 'resume') {let chk = stillLeased('resume', c.project, c.number, c.lease)if (!chk.ok) {console.log('lease: NOT resumed ' + c.ref + ' (session ' + c.session + '): re-read right before — ' + chk.why + ' → nothing written')return 'NOT resumed ' + c.ref + ' (' + c.session + '): re-read right before — ' + chk.why + ' → nothing written'}}running.push(entry)let ok = falseif (c.kind == 'resume') {ok = runResume(runs + '/' + c.session, pass, (r) => {sessionEnded(entry, r)return null}, hooks, ports)} else {let wargv = ['work' c.project '' + c.number '--post' '--session' c.session '--ports' ports]for (o of pass) { wargv.push(o) }ok = runWork(wargv, (r) => {sessionEnded(entry, r)return null}, hooks)}if (!ok) {let out = []for (x of running) { if (x.session != entry.session) { out.push(x) } }running = outreturn c.ref + ' not started (see above)'}return (c.kind == 'resume' ? 'resumed ' : 'started ') + c.ref + ' (' + c.session + ', ports ' + ports + ')'}let iterate = nulllet runLibrarian = nulllet tick = () => {if (busy || finished || stopping != null) { return null }if (once && iter >= 1) { return null }busy = trueiter = iter + 1let head = 'ITER ' + iter + ' ' + ago(now())// antcolony#17: the creator's /confirm /reject /prio /hold in comments — code only, no quota, before anything is pickedlet kw = keywordPass(true, kwSeen)if (kw == null || kw.error != null) { console.log(head + ' · keywords: ' + (kw == null ? 'failed' : kw.error)) }// antcolony#19: the daily summary of what waits for the creator — code only, once a day after COLONY_DIGEST_HOURlet dg = digestTick(digestState, now())if (dg != null) { console.log(head + ' · ' + dg) }// mission 034: a known limit → no model call (no librarian), no probe, no start until it resetsif (limitUntil != null) {if (now() < limitUntil) {iterate(head, 'usage limit', null, true, null)return null}console.log('LIMIT: resumed — the pause until ' + ago(limitUntil) + ' is over; the quota is read first, then the librarian and starts')limitUntil = null}if (libOn != 'on') {iterate(head, null, null, false, null)return null}// mission 034 (creator: "it cant recognize when it gets launches when we already hit the limit"): the quota FIRST// (`claude -p /usage`, local, $0) — at or above the start threshold → no librarian call and no start this iterationprobeUsage(quotaBin, runs, (q) => {let qb = quotaBlock(q)if (qb != null) {iterate(head, null, q, false, qb)return null}runLibrarian(head, q)return null})return null}// mission 029 — ORDER: the librarian pass comes FIRST and is awaited; only then the iteration looks at tickets and// starts anything, so every brief built in this iteration has every creator decision written before it began.// A pass that failed or left a creator text unfiled → nothing starts in this iteration (`hold`).runLibrarian = (head, q) => {libBusy = truelibrarianPass({ argv = argv post = true auto = true stop = () => { return stopping != null } }, (s) => {libBusy = false// mission 034: the librarian's call hit the usage limit → pause (librarian + starts) until the resetif (s != null && s.limited == true) {pauseForLimit(s.resetMs, 'the librarian\'s model call hit it (' + s.limitText + '); ' + s.left + ' creator text(s) not filed yet')iterate(head, 'usage limit', null, true, null)return null}let hold = nullif (s == null || s.error != null) { hold = 'librarian failed (' + (s == null ? '?' : s.error) + ') → no start this iteration' }else if (s.left != null && s.left > 0) { hold = 'librarian: ' + s.left + ' creator text(s) not filed yet → no start this iteration' }// the reading is re-taken before a start when the pass made model calls (it may be minutes old then)iterate(head, hold, s != null && s.calls != null && s.calls > 0 ? null : q, false, null)return null})return null}// the rest of one iteration (after the librarian): leases, candidates, quota, starts.// q = the quota reading of this iteration (null → probed here when something could start); paused = a known usage// limit (mission 034): leases only, no start, no line unless a lease was expired; libSkipped = why the librarian did// not run (the quota at or above the start threshold)iterate = (head, hold, q0, paused, libSkipped) => {if (stopping != null) {console.log(head + ' · stop requested — nothing started · running ' + running.length + '/' + capacity())busy = falsecheckEnd()return null}let notes = []let did = []// mission 035: an agent that has not synced for --agent-ttl s is offline — its entries go; its sessions keep their// leases in tickets (the heartbeat is theirs) and are handled like any lease: renewed, or expiredfor (a of agents) {if (a.announced == true && !agentOnline(a)) {a.announced = falselet had = hostRunning(a.host)dropHost(a.host)console.log('AGENT: ' + a.host + ' is OFFLINE (no sync for ' + ttlText + ' s) — ' + had + ' session(s) were on it; their leases in tickets decide (renewed by the workers, else expired)')}}let sv = survey()if (sv.error != null) {console.log(head + ' · ERROR ' + sv.error + ' · running ' + running.length + '/' + capacity())busy = falsecheckEnd()return null}let t0 = now()let resumes = []for (e of sv.notEligible) {let t = e.tif (t.state == 'in progress') {let ref = t.project + t.reflet d = getJson(ticketPath(t.project, t.number))if (d.status != 200) { notes.push(ref + ' unreadable (' + d.status + ')') } else {let l = leaseOf(d.json.events)if (l.kind == 'ready') { notes.push(ref + ' built, waiting for the deploy') }else if (l.kind == 'ended') { notes.push(ref + ' in progress, but the lease of ' + l.session + ' has ended (' + l.text + ') → never expired') }else if (l.kind != 'colony') { notes.push(ref + ' in progress without a colony lease') }else if (runningSession(l.session)) { }else if (t0 >= l.expiryMs) {// mission 031: expireLease re-reads the ticket first and writes nothing unless it is STILL this dead leaselet r = expireLease(t.project, t.number, l, tk.token)if (r.ok) { did.push('expired ' + ref + ' (session ' + l.session + ' on ' + l.host + ', lease ended ' + ago(l.expiryMs) + ') → open') }else if (r.skipped == true) {did.push('NOT expired ' + ref + ' (session ' + l.session + ', lease ended ' + ago(l.expiryMs) + '): re-read right before — ' + r.text + ' → nothing written')console.log('lease: NOT expired ' + ref + ' (session ' + l.session + '): re-read right before — ' + r.text + ' → nothing written')}else { did.push('EXPIRING ' + ref + ' FAILED: ' + r.text) }} else if (l.parked) {// mission 030: a question ticket is never work — a session parked on one (pre-030) is not resumedif (questionWhy(t) != null) { notes.push(ref + ' parked, but ' + questionWhy(t) + ' → not resumed') }else if (!serving && l.host.toLowerCase() != host.toLowerCase()) { notes.push(ref + ' parked on ' + l.host) }else if (t0 < l.resumeMs) { notes.push(ref + ' parked until ' + ago(l.resumeMs)) }else if (serving && !agentOnline(agentOf(l.host))) { notes.push(ref + ' parked on ' + l.host + ' (no agent online there)') }else if (!serving && !exists(runs + '/' + l.session + '/session.json')) { notes.push(ref + ' parked, but ' + runs + '/' + l.session + ' is missing') }else {let sj = serving ? null : runs + '/' + l.session + '/session.json'resumes.push({ kind = 'resume' ref = ref project = t.project number = t.number lease = l session = l.session folder = null sj = sj host = l.host })}} else { notes.push(ref + ' leased by ' + l.session + ' on ' + l.host + ' until ' + ago(l.expiryMs)) }}}}let cands = []for (r of resumes) {if (serving) { cands.push(r) }else {let m = JSON.parse(readFile(r.sj))r.folder = m.folderif (folderBusy(r.folder)) { notes.push(r.ref + ' (resume) dev folder busy') } else { cands.push(r) }}}let other = 0for (t of sv.eligible) {let m = sv.reg.projects[t.project]let ref = t.project + t.refif (m.dev == null || m.dev.host == null) { other = other + 1 }else if (serving && !agentOnline(agentOf(m.dev.host))) { notes.push(ref + ' waits — no agent online on ' + m.dev.host) }else if (!serving && m.dev.host.toLowerCase() != host.toLowerCase()) { other = other + 1 }else if (missingCapabilities(m.needList, serving ? agentOf(m.dev.host).caps : hostCaps).length > 0) { notes.push(ref + ' waits — ' + m.dev.host + ' does not offer: ' + missingCapabilities(m.needList, serving ? agentOf(m.dev.host).caps : hostCaps).join(', ')) }else if (m.dev.folder == null || folderBusy(m.dev.folder)) { notes.push(ref + ' dev folder busy') }else {let dup = nullfor (c of cands) { if (dup == null && c.folder == m.dev.folder) { dup = c.ref } }if (dup != null) { notes.push(ref + ' waits (same dev folder as ' + dup + ')') }else { cands.push({ kind = 'new' ref = ref project = t.project number = t.number folder = m.dev.folder session = null host = m.dev.host }) }}}if (other > 0) { notes.push(other + ' eligible on other hosts') }let free = serving ? capacity() - running.length : parallel - running.lengthlet libNote = libSkipped == null ? [] : ['librarian not run: quota at or above the start threshold (' + libSkipped + ')']let line = (quota, list) => {let parts = [head, quota, 'running ' + running.length + '/' + capacity()]for (x of did) { parts.push(x) }for (x of list) { parts.push(x) }for (x of libNote) { parts.push(x) }if (notes.length > 0) { parts.push('skipped: ' + notes.join(', ')) }console.log(parts.join(' · '))return null}// mission 034: paused on a known usage limit — nothing starts; the iteration is silent unless a lease was acted onif (paused == true) {if (did.length > 0) { line('usage limit — librarian and starts paused until ' + ago(limitUntil), []) }busy = falsecheckEnd()return null}if (hold != null && cands.length > 0) {line(hold, [])busy = falsecheckEnd()return null}// mission 035 (antcolony#14): the scheduler starts nothing itself — a candidate becomes a command for the agent of its// dev host (one with a free slot that says it is accepting work; the quota is the agent's), sent with the agent's next syncif (serving) {let queued = []for (c of cands) {let a = agentOf(c.host)if (!agentOnline(a)) { notes.push(c.ref + ' waits — no agent online on ' + c.host) }else if (!a.accepting) { notes.push(c.ref + ' waits — agent ' + a.host + ' is not accepting work (' + a.why + ')') }else if (hostRunning(a.host) >= a.parallel) { notes.push(c.ref + ' waits — agent ' + a.host + ' has all ' + a.parallel + ' slot(s) busy') }else if (accountBlock(a.account) != null) { notes.push(c.ref + ' waits — Claude account ' + a.account + ': ' + accountBlock(a.account)) }else if (c.kind == 'resume') {// mission 031: a resume acts on what the survey saw — re-read firstlet chk = stillLeased('resume', c.project, c.number, c.lease)if (!chk.ok) {console.log('lease: NOT resumed ' + c.ref + ' (session ' + c.session + '): re-read right before — ' + chk.why + ' → nothing written')notes.push('NOT resumed ' + c.ref + ' (' + c.session + '): re-read right before — ' + chk.why + ' → nothing written')} else {running.push({ ref = c.ref session = c.session folder = null slot = -1 ports = '' phase = 'worker' remote = a.host proc = null sentAt = null cmd = { kind = 'resume' project = c.project number = c.number session = c.session folder = null brief = null } })accountStart[a.account] = now()queued.push('resume ' + c.ref + ' (' + c.session + ') → ' + a.host)}} else {let session = newSession()let b = buildBrief(c.project, '' + c.number, ['--session' session '--ports' PORTS_MARKER])if (b.error != null) { notes.push(c.ref + ' not started: ' + b.error) }else {running.push({ ref = c.ref session = session folder = c.folder slot = -1 ports = '' phase = 'worker' remote = a.host proc = null sentAt = null cmd = { kind = 'start' project = c.project number = c.number session = session folder = c.folder brief = b.text } })accountStart[a.account] = now()queued.push('start ' + c.ref + ' (' + session + ') → ' + a.host)}}}let names = []for (a of agents) { if (agentOnline(a)) { names.push(a.host + ' ' + hostRunning(a.host) + '/' + a.parallel + (a.accepting ? '' : ' (not accepting: ' + a.why + ')')) } }let accts = []for (a of agents) {if (agentOnline(a) && !accts.includes(a.account)) {accts.push(a.account)let q = accountReading(a.account)names.push('account ' + a.account + (q == null ? ' quota unknown' : ' ' + q.percent + '%' + (q.weekly == null ? '' : ' week ' + q.weekly + '%')) + ' [' + accountHosts(a.account).join('+') + ']')}}line('agents: ' + (names.length == 0 ? 'none online' : names.join(', ')), queued)busy = falsecheckEnd()return null}if (cands.length == 0 || free <= 0) {let why = cands.length == 0 ? 'nothing to start' : 'all ' + parallel + ' slot(s) busy, waiting: ' + cands.lengthif (q0 == null) { line('quota not read (' + why + ')', []) }else if (q0.error != null) { line('quota UNKNOWN (' + q0.error + ') · ' + why, []) }else { line(quotaText(q0) + (libSkipped == null ? '' : ' → librarian and starts paused') + ' · ' + why, []) }busy = falsecheckEnd()return null}let withQuota = (fn) => {if (q0 != null) { return fn(q0) }return probeUsage(quotaBin, runs, fn)}withQuota((q) => {if (stopping != null) {line('stop requested — nothing started', [])busy = falsecheckEnd()return null}if (q.error != null) {line('quota UNKNOWN (' + q.error + ') → no start', [])busy = falsecheckEnd()return null}let qt = quotaText(q)// --reserve 100 = the 5 h window never blocks (creator 2026-09-24: first live runs up to 100 %)if (reserve < 100 && q.percent >= reserve) {let names = []for (c of cands) { names.push(c.ref) }line(qt + ' ≥ reserve ' + reserve + '% → no start (waiting: ' + names.join(', ') + ')', [])busy = falsecheckEnd()return null}// mission 023 (creator): no new start (resumes too) while the WEEKLY usage is above the week limit (84 %);// a probe without a weekly number starts nothing either — better wait than overspend the week;// --week-limit 100 = the week never blocks (also when unknown)if (weekLimit < 100 && (q.weekly == null || q.weekly > weekLimit)) {let wn = []for (c of cands) { wn.push(c.ref) }line(qt + (q.weekly == null ? ', week UNKNOWN' : ' > week limit ' + weekLimit + '%') + ' → no start (waiting: ' + wn.join(', ') + ')', [])busy = falsecheckEnd()return null}// mission 034: a REACHED limit (100 %) holds every start, also with --reserve 100 / --week-limit 100let reached = quotaBlock(q)if (reached != null) {let rn = []for (c of cands) { rn.push(c.ref) }line(qt + ' → ' + reached + ' → no start (waiting: ' + rn.join(', ') + ')', [])busy = falsecheckEnd()return null}let started = []let i = 0while (i < cands.length && running.length < parallel) {let c = cands[i]if (c.session == null) { c.session = newSession() }started.push(startOne(c))i = i + 1}while (i < cands.length) {notes.push(cands[i].ref + ' waits for a slot')i = i + 1}line(qt + (reserve >= 100 ? ' · 5 h reserve off (100%)' : ' < reserve ' + reserve + '%') + (weekLimit >= 100 ? ' · week limit off (100%)' : ''), started)busy = falsecheckEnd()return null})return null}let checkStop = () => {if (finished || stopFile == null || stopFile == '' || !exists(stopFile)) { return null }let mode = readFile(stopFile).trim()if (stopping == null) {stopping = 'finish'console.log('RUN: stop requested — starting nothing more; ' + running.length + ' running session(s) finish (a second SIGTERM/SIGINT parks them instead)')}if (mode == 'park' && stopping != 'park') {stopping = 'park'console.log('RUN: second stop signal — parking ' + running.length + ' running session(s)')for (r of running) {if (live[r.session] != null) {live[r.session].parkReason = 'the scheduler run was stopped (second SIGTERM/SIGINT)'if (live[r.session].proc != null) {console.log('RUN: killing ' + live[r.session].phase + ' pid ' + live[r.session].proc.pid + ' of ' + r.ref + ' ' + r.session + ' → parked')live[r.session].proc.kill()}}}}checkEnd()return null}watch = every(0.5)on watch.tick(x) {checkStop()return null}loop = every(interval)on loop.tick(x) {tick()return null}tick()return true}
Branches
- mainmain branch