antcolony
All repositories: gitoria
17.3 KB
// agent.hl — `agent`: the program on ONE host that starts the workers (mission 035, ticket antcolony#14; concept// docs/scheduler-agent.md §1 + §3). The scheduler (`run --serve`) decides WHAT runs; this agent runs it here, on the host// that holds the project's dev folder, and reports back. It holds no data of its own: tickets stay the only state — a// worker's lease, heartbeat and report go to tickets directly (work.hl), the agent only keeps in-flight session bookkeeping.//// Connects OUT: every --poll s one `POST <scheduler>/agent/sync` (HTTPS, `Authorization: Bearer <agent token>`):// → { host, capabilities: [what this computer offers, e.g. chrome, gpu, deploy], parallel, accepting, why, running: [{ref, session, folder, phase, ports}], events: [buffered, not yet acked] }// ← { ack: [event ids], commands: [{kind: start|resume, project, number, session, folder, brief}] }// Hybriel has no WebSocket client (hybriel#109), so this is request/answer instead of one held-open socket; the effects the// concept asks for are the same: the agent reconnects on its own, sessions keep running while the scheduler is not// reachable, `SESSION END` events are buffered (until acked) and delivered on reconnect, a resent command is ignored.// Quota is the agent's (concept §3): `claude -p /usage` is read here (at most every --quota-every s, only while a slot is// free); at/above the reserve, above the weekly limit, or after a usage limit was hit → `accepting: false` + why, and a// command that arrives anyway is refused (event `refused`). Ports: the agent's own pool, one disjoint range per session.// Stop: the `colony` wrapper writes $COLONY_STOP_FILE: first signal → accept nothing, running sessions finish;// second → they are parked (resumed by the next agent / scheduler round). Never silent: every line starts with AGENT / SESSION END / LIMIT.import { env } from 'hl:proc'import { exists, readFile, writeFile, mkDir } from 'hl:fs'import { fetch } from 'hl:fetch'import { every, now } from 'hl:time'import { readToken, readAgentToken, userAgent } from './tickets.hl'import { runWork, runResume, hostName, runsDir } from './work.hl'import { setting, isWholeNumber } from './claude.hl'import { parseCapabilities } from './util.hl'import { probeUsage } from './quota.hl'import { portRange } from './brief.hl'import { passOptions, money, ago } from './daemon.hl'static LIMIT_RETRY_MS = 900000static runAgent = (argv) => {let home = env('COLONY_HOME')if (home == null || home == '') {console.log('agent: REFUSED — COLONY_HOME is not set (run through ./colony, which sets it)')return false}let schedText = setting(argv, '--scheduler', 'COLONY_SCHEDULER_URL', '')if (schedText == '') {console.log('agent: REFUSED — --scheduler URL (COLONY_SCHEDULER_URL) is not set: where is the scheduler?')return false}while (schedText.endsWith('/')) { schedText = schedText.slice(0, schedText.length - 1) }let pollText = setting(argv, '--poll', 'COLONY_AGENT_POLL', '5')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 everyText = setting(argv, '--quota-every', 'COLONY_AGENT_QUOTA_EVERY', '60')let capText = setting(argv, '--capabilities', 'COLONY_CAPABILITIES', '')let accountText = setting(argv, '--account', 'COLONY_CLAUDE_ACCOUNT', '')let claudeBin = setting(argv, '--claude', 'COLONY_CLAUDE', 'claude')let quotaBin = setting(argv, '--quota-claude', 'COLONY_QUOTA_CLAUDE', claudeBin)if (!isWholeNumber(pollText) || pollText == '0') {console.log('agent: REFUSED — --poll must be whole seconds ≥ 1, got ' + pollText)return false}if (!isWholeNumber(parallelText) || parallelText == '0') {console.log('agent: REFUSED — --parallel must be a whole number ≥ 1, got ' + parallelText)return false}if (!isWholeNumber(reserveText) || toNumber(reserveText) > 100) {console.log('agent: REFUSED — --reserve must be a percentage 0–100, got ' + reserveText)return false}if (!isWholeNumber(weekText) || toNumber(weekText) > 100) {console.log('agent: REFUSED — --week-limit must be a percentage 0–100, got ' + weekText)return false}if (!isWholeNumber(everyText) || everyText == '0') {console.log('agent: REFUSED — --quota-every must be whole seconds ≥ 1, got ' + everyText)return false}let pool = portRange(poolText)if (pool == null) {console.log('agent: REFUSED — --port-pool must be A-B (1 ≤ A ≤ B ≤ 65535), got ' + poolText)return false}if (!isWholeNumber(perText) || perText == '0') {console.log('agent: REFUSED — --ports-per-session must be a whole number ≥ 1, got ' + perText)return false}let caps = parseCapabilities(capText)let poll = toNumber(pollText)let parallel = toNumber(parallelText)let reserve = toNumber(reserveText)let weekLimit = toNumber(weekText)let per = toNumber(perText)let quotaEvery = toNumber(everyText)let size = pool.to - pool.from + 1let slots = (size - size % per) / perif (slots < parallel) {console.log('agent: REFUSED — the port pool ' + poolText + ' has ' + slots + ' range(s) of ' + per + ' ports, --parallel ' + parallel + ' needs ' + parallel)return false}let at = readAgentToken()if (at.error != null) {console.log('agent: REFUSED — ' + at.error)return false}let tk = readToken()if (tk.error != null) {console.log('agent: REFUSED — ' + tk.error + ' (the workers of this host lease and post to tickets)')return false}let host = hostName()// antcolony#16: the name of the Claude account this host is logged in to; hosts sharing one account give it the same name// (default: the host's own name = its own account). The scheduler pools the quota of all agents with the same name.let account = accountText.trim() == '' ? host : accountText.trim()let runs = runsDir(home)mkDir(runs, 448)mkDir(runs + '/.briefs', 448)let stopFile = env('COLONY_STOP_FILE')let pass = passOptions(argv)let running = []let live = {}let events = []let handled = []let stopping = nulllet finished = falselet syncing = falselet connected = nulllet probing = falselet lastProbe = nulllet reading = nulllet quotaOk = falselet quotaWhy = 'the quota is not read yet'let limitUntil = nulllet loop = nulllet watch = nulllet n = 0console.log('AGENT: started on ' + host + ' → scheduler ' + schedText + ' · poll ' + poll + ' s · parallel ' + parallel + ' · offers ' + (caps.length == 0 ? '— (nothing announced)' : caps.join(', ')) + ' · quota reserve ' + reserve + '% of the 5 h window, no start above ' + weekLimit + '% of the week (probe: ' + quotaBin + ' -p /usage, every ' + quotaEvery + ' s while a slot is free) · ports ' + pool.from + '-' + pool.to + ', ' + per + ' per session · runs ' + runs + '/ · stop file ' + (stopFile == null ? '— (no wrapper: SIGTERM kills without parking)' : stopFile) + ' · Claude account ' + account + (accountText.trim() == '' ? ' (own — none named: not shared)' : ' (quota shared with the agents of the same name)'))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 known = (s) => {for (r of running) { if (r.session == s) { return true } }return handled.includes(s)}let folderBusy = (f) => {for (r of running) { if (r.folder == f) { return true } }return false}let pushEvent = (ev) => {for (e of events) { if (e.id == ev.id) { return null } }events.push(ev)return null}// the same rules as the single-host `run` (daemon.hl quotaBlock + the reserve / week checks of iterate)let quotaVerdict = (q) => {if (q == null || q.error != null || q.percent == null) { return 'the quota is UNKNOWN (' + (q == null ? 'no reading' : q.error) + ')' }let qt = 'quota ' + q.percent + '% of 5 h' + (q.resetsMs == null ? '' : ' (resets ' + ago(q.resetsMs) + ')') + (q.weekly == null ? '' : ', week ' + q.weekly + '%')if (q.percent >= 100) { return qt + ' → the 5 h limit is reached (100%)' }if (q.weekly != null && q.weekly >= 100) { return qt + ' → the weekly limit is reached (100%)' }if (reserve < 100 && q.percent >= reserve) { return qt + ' ≥ reserve ' + reserve + '%' }if (weekLimit < 100 && (q.weekly == null || q.weekly > weekLimit)) { return qt + (q.weekly == null ? ', week UNKNOWN' : ' > week limit ' + weekLimit + '%') }return null}let pauseForLimit = (resetMs, why) => {let until = resetMs == null ? now() + LIMIT_RETRY_MS : resetMsif (limitUntil == null) {limitUntil = untilconsole.log('LIMIT: usage limit reached — no start 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}let accepting = () => {if (stopping != null) { return { ok = false why = 'stop requested' } }if (limitUntil != null) {if (now() < limitUntil) { return { ok = false why = 'usage limit reached until ' + ago(limitUntil) } }console.log('LIMIT: resumed — the pause until ' + ago(limitUntil) + ' is over; the quota is read again')limitUntil = nulllastProbe = null}if (!quotaOk) { return { ok = false why = quotaWhy } }return { ok = true why = null }}let flushed = falselet sync = nulllet checkEnd = () => {if (finished || stopping == null || running.length > 0) { return null }// one last try to hand over what the sessions reported while stoppingif (events.length > 0 && !flushed) {flushed = truesync()return null}finished = trueif (loop != null) { loop.stop() }if (watch != null) { watch.stop() }console.log('AGENT STOPPED — stop requested (' + stopping + '), no session running' + (events.length > 0 ? ' · ' + events.length + ' event(s) not delivered (the scheduler was not reachable)' : ''))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 + '/' + parallel)pushEvent({ id = entry.session + ':ended' kind = 'ended' session = entry.session ref = entry.ref outcome = r.outcome workerCost = r.workerCost == null ? 0 : r.workerCost controllerCost = r.controllerCost == null ? 0 : r.controllerCost ports = entry.ports limited = r.limited == true })if (r.limited == true) { pauseForLimit(r.resetMs, entry.ref + ' ' + entry.session + ' was parked on it (' + r.limitText + ')') }checkEnd()return null}let refuse = (c, why) => {console.log('AGENT: refused ' + c.kind + ' ' + c.project + '#' + c.number + ' (' + c.session + ') — ' + why)pushEvent({ id = c.session + ':refused' kind = 'refused' session = c.session ref = c.project + '#' + c.number why = why })return null}let startCommand = (c) => {if (known(c.session)) { return null }let ref = c.project + '#' + c.numberlet a = accepting()if (!a.ok) { return refuse(c, 'this agent is not accepting: ' + a.why) }let slot = freeSlot()if (slot < 0) { return refuse(c, 'no slot / port range free (' + running.length + '/' + parallel + ' running)') }let ports = rangeOf(slot)let entry = { ref = 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` (a push into `running` copies it)live[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 }}handled.push(c.session)let ok = falseif (c.kind == 'resume') {if (!exists(runs + '/' + c.session + '/session.json')) { return refuse(c, 'no run folder ' + runs + '/' + c.session + ' on ' + host + ' (parked on another host?)') }let sj = JSON.parse(readFile(runs + '/' + c.session + '/session.json'))entry.folder = sj.folderif (folderBusy(entry.folder)) { return refuse(c, 'dev folder ' + entry.folder + ' is busy') }running.push(entry)ok = runResume(runs + '/' + c.session, pass, (r) => {sessionEnded(entry, r)return null}, hooks, ports)} else {if (c.folder != null && folderBusy(c.folder)) { return refuse(c, 'dev folder ' + c.folder + ' is busy') }let briefFile = runs + '/.briefs/' + c.session + '.md'writeFile(briefFile, c.brief, 384)running.push(entry)let wargv = ['work' c.project '' + c.number '--post' '--session' c.session '--ports' ports '--brief-file' briefFile]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 refuse(c, 'the worker did not start (see the lines above)')}console.log('AGENT: ' + (c.kind == 'resume' ? 'resumed ' : 'started ') + ref + ' (' + c.session + ', ports ' + ports + ') · running ' + running.length + '/' + parallel)return null}// quota: read at most every --quota-every s, only while a slot is free and nothing is paused / stoppinglet maybeProbe = () => {if (probing || stopping != null || running.length >= parallel) { return null }if (limitUntil != null && now() < limitUntil) { return null }if (lastProbe != null && now() - lastProbe < quotaEvery * 1000) { return null }probing = truelastProbe = now()probeUsage(quotaBin, runs, (q) => {probing = falsereading = q == null || q.error != null || q.percent == null ? null : { percent = q.percent weekly = q.weekly resetsMs = q.resetsMs at = now() }let v = quotaVerdict(q)let was = quotaOkquotaOk = v == nullquotaWhy = v == null ? null : vif (quotaOk != was || n == 0) { console.log('AGENT: ' + (quotaOk ? 'quota ok (' + q.percent + '% of 5 h' + (q.weekly == null ? '' : ', week ' + q.weekly + '%') + ') → accepting work' : 'not accepting work — ' + v)) }return null})return null}sync = () => {if (syncing || finished) { return null }syncing = truen = n + 1maybeProbe()let a = accepting()let list = []for (r of running) { list.push({ ref = r.ref session = r.session folder = r.folder phase = r.phase ports = r.ports }) }let payload = { host = host account = account quota = reading limitUntil = limitUntil capabilities = caps parallel = parallel accepting = a.ok why = a.why running = list events = events }let res = fetch(schedText + '/agent/sync', { method = 'POST' headers = { 'user-agent' = userAgent() authorization = 'Bearer ' + at.token } json = payload timeoutMs = 15000 })if (res == null || res.status != 200) {let why = res == null ? 'no answer' : 'answered ' + res.statusif (connected != false) { console.log('AGENT: scheduler ' + schedText + ' not reachable (' + why + ') — ' + running.length + ' session(s) keep running, ' + events.length + ' event(s) buffered, trying again every ' + poll + ' s') }connected = falsesyncing = falsereturn null}if (connected != true) { console.log('AGENT: connected to ' + schedText + (connected == false ? ' again' : '') + ' · ' + events.length + ' buffered event(s) delivered') }connected = truelet j = res.json()let keep = []for (e of events) { if (!j.ack.includes(e.id)) { keep.push(e) } }events = keepfor (c of j.commands) { startCommand(c) }syncing = falsecheckEnd()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('AGENT: stop requested — accepting nothing more; ' + running.length + ' running session(s) finish (a second SIGTERM/SIGINT parks them instead)')}if (mode == 'park' && stopping != 'park') {stopping = 'park'console.log('AGENT: second stop signal — parking ' + running.length + ' running session(s)')for (r of running) {if (live[r.session] != null) {live[r.session].parkReason = 'the agent was stopped (second SIGTERM/SIGINT)'if (live[r.session].proc != null) {console.log('AGENT: 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(poll)on loop.tick(x) {sync()return null}sync()return true}
Branches
- mainmain branch
Latest commits
- 3a4d0324antcolony#37: a too-long report gets up to 3 fix tries, finished work is never thrown away for lengthmre
- a6af7883tracker: worker box sees calendar.worldapi.org (login to copy)mre
- c613d26btemplates: bridges to external components (login.js for ident's selector) are allowed (creator 2026-09-27)mre
- 9062978ctracker: worker box sees /media/STORAGE/projects/old-tracker read-only (tracker#2 source data)mre
- 7f9660eeState of 2026-09-27, before the move to gitoriamre