gitoriaLog in with ident

antcolony

All repositories: gitoria

ReadmeCodePull requestsReleasesTicketsSettings
Commit7f9660ee7f9660eeState of 2026-09-27, before the move to gitoriamre7f9660ee/lib/agent.hl

17.3 KB

  1. // agent.hl — `agent`: the program on ONE host that starts the workers (mission 035, ticket antcolony#14; concept
  2. // docs/scheduler-agent.md §1 + §3). The scheduler (`run --serve`) decides WHAT runs; this agent runs it here, on the host
  3. // that holds the project's dev folder, and reports back. It holds no data of its own: tickets stay the only state — a
  4. // worker's lease, heartbeat and report go to tickets directly (work.hl), the agent only keeps in-flight session bookkeeping.
  5. //
  6. // Connects OUT: every --poll s one `POST <scheduler>/agent/sync` (HTTPS, `Authorization: Bearer <agent token>`):
  7. // → { 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] }
  8. // ← { ack: [event ids], commands: [{kind: start|resume, project, number, session, folder, brief}] }
  9. // Hybriel has no WebSocket client (hybriel#109), so this is request/answer instead of one held-open socket; the effects the
  10. // concept asks for are the same: the agent reconnects on its own, sessions keep running while the scheduler is not
  11. // reachable, `SESSION END` events are buffered (until acked) and delivered on reconnect, a resent command is ignored.
  12. // Quota is the agent's (concept §3): `claude -p /usage` is read here (at most every --quota-every s, only while a slot is
  13. // free); at/above the reserve, above the weekly limit, or after a usage limit was hit → `accepting: false` + why, and a
  14. // command that arrives anyway is refused (event `refused`). Ports: the agent's own pool, one disjoint range per session.
  15. // Stop: the `colony` wrapper writes $COLONY_STOP_FILE: first signal → accept nothing, running sessions finish;
  16. // second → they are parked (resumed by the next agent / scheduler round). Never silent: every line starts with AGENT / SESSION END / LIMIT.
  17. import { env } from 'hl:proc'
  18. import { exists, readFile, writeFile, mkDir } from 'hl:fs'
  19. import { fetch } from 'hl:fetch'
  20. import { every, now } from 'hl:time'
  21. import { readToken, readAgentToken, userAgent } from './tickets.hl'
  22. import { runWork, runResume, hostName, runsDir } from './work.hl'
  23. import { setting, isWholeNumber } from './claude.hl'
  24. import { parseCapabilities } from './util.hl'
  25. import { probeUsage } from './quota.hl'
  26. import { portRange } from './brief.hl'
  27. import { passOptions, money, ago } from './daemon.hl'
  28. static LIMIT_RETRY_MS = 900000
  29. static runAgent = (argv) => {
  30. let home = env('COLONY_HOME')
  31. if (home == null || home == '') {
  32. console.log('agent: REFUSED — COLONY_HOME is not set (run through ./colony, which sets it)')
  33. return false
  34. }
  35. let schedText = setting(argv, '--scheduler', 'COLONY_SCHEDULER_URL', '')
  36. if (schedText == '') {
  37. console.log('agent: REFUSED — --scheduler URL (COLONY_SCHEDULER_URL) is not set: where is the scheduler?')
  38. return false
  39. }
  40. while (schedText.endsWith('/')) { schedText = schedText.slice(0, schedText.length - 1) }
  41. let pollText = setting(argv, '--poll', 'COLONY_AGENT_POLL', '5')
  42. let parallelText = setting(argv, '--parallel', 'COLONY_PARALLEL', '1')
  43. let reserveText = setting(argv, '--reserve', 'COLONY_QUOTA_RESERVE', '70')
  44. let weekText = setting(argv, '--week-limit', 'COLONY_WEEK_LIMIT', '84')
  45. let poolText = setting(argv, '--port-pool', 'COLONY_PORT_POOL', '8700-8799')
  46. let perText = setting(argv, '--ports-per-session', 'COLONY_PORTS_PER_SESSION', '10')
  47. let everyText = setting(argv, '--quota-every', 'COLONY_AGENT_QUOTA_EVERY', '60')
  48. let capText = setting(argv, '--capabilities', 'COLONY_CAPABILITIES', '')
  49. let accountText = setting(argv, '--account', 'COLONY_CLAUDE_ACCOUNT', '')
  50. let claudeBin = setting(argv, '--claude', 'COLONY_CLAUDE', 'claude')
  51. let quotaBin = setting(argv, '--quota-claude', 'COLONY_QUOTA_CLAUDE', claudeBin)
  52. if (!isWholeNumber(pollText) || pollText == '0') {
  53. console.log('agent: REFUSED — --poll must be whole seconds ≥ 1, got ' + pollText)
  54. return false
  55. }
  56. if (!isWholeNumber(parallelText) || parallelText == '0') {
  57. console.log('agent: REFUSED — --parallel must be a whole number ≥ 1, got ' + parallelText)
  58. return false
  59. }
  60. if (!isWholeNumber(reserveText) || toNumber(reserveText) > 100) {
  61. console.log('agent: REFUSED — --reserve must be a percentage 0–100, got ' + reserveText)
  62. return false
  63. }
  64. if (!isWholeNumber(weekText) || toNumber(weekText) > 100) {
  65. console.log('agent: REFUSED — --week-limit must be a percentage 0–100, got ' + weekText)
  66. return false
  67. }
  68. if (!isWholeNumber(everyText) || everyText == '0') {
  69. console.log('agent: REFUSED — --quota-every must be whole seconds ≥ 1, got ' + everyText)
  70. return false
  71. }
  72. let pool = portRange(poolText)
  73. if (pool == null) {
  74. console.log('agent: REFUSED — --port-pool must be A-B (1 ≤ A ≤ B ≤ 65535), got ' + poolText)
  75. return false
  76. }
  77. if (!isWholeNumber(perText) || perText == '0') {
  78. console.log('agent: REFUSED — --ports-per-session must be a whole number ≥ 1, got ' + perText)
  79. return false
  80. }
  81. let caps = parseCapabilities(capText)
  82. let poll = toNumber(pollText)
  83. let parallel = toNumber(parallelText)
  84. let reserve = toNumber(reserveText)
  85. let weekLimit = toNumber(weekText)
  86. let per = toNumber(perText)
  87. let quotaEvery = toNumber(everyText)
  88. let size = pool.to - pool.from + 1
  89. let slots = (size - size % per) / per
  90. if (slots < parallel) {
  91. console.log('agent: REFUSED — the port pool ' + poolText + ' has ' + slots + ' range(s) of ' + per + ' ports, --parallel ' + parallel + ' needs ' + parallel)
  92. return false
  93. }
  94. let at = readAgentToken()
  95. if (at.error != null) {
  96. console.log('agent: REFUSED — ' + at.error)
  97. return false
  98. }
  99. let tk = readToken()
  100. if (tk.error != null) {
  101. console.log('agent: REFUSED — ' + tk.error + ' (the workers of this host lease and post to tickets)')
  102. return false
  103. }
  104. let host = hostName()
  105. // antcolony#16: the name of the Claude account this host is logged in to; hosts sharing one account give it the same name
  106. // (default: the host's own name = its own account). The scheduler pools the quota of all agents with the same name.
  107. let account = accountText.trim() == '' ? host : accountText.trim()
  108. let runs = runsDir(home)
  109. mkDir(runs, 448)
  110. mkDir(runs + '/.briefs', 448)
  111. let stopFile = env('COLONY_STOP_FILE')
  112. let pass = passOptions(argv)
  113. let running = []
  114. let live = {}
  115. let events = []
  116. let handled = []
  117. let stopping = null
  118. let finished = false
  119. let syncing = false
  120. let connected = null
  121. let probing = false
  122. let lastProbe = null
  123. let reading = null
  124. let quotaOk = false
  125. let quotaWhy = 'the quota is not read yet'
  126. let limitUntil = null
  127. let loop = null
  128. let watch = null
  129. let n = 0
  130. console.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)'))
  131. let slotFree = (i) => {
  132. for (r of running) { if (r.slot == i) { return false } }
  133. return true
  134. }
  135. let freeSlot = () => {
  136. let i = 0
  137. while (i < slots) {
  138. if (slotFree(i)) { return i }
  139. i = i + 1
  140. }
  141. return -1
  142. }
  143. let rangeOf = (i) => {
  144. let a = pool.from + i * per
  145. return a + '-' + (a + per - 1)
  146. }
  147. let known = (s) => {
  148. for (r of running) { if (r.session == s) { return true } }
  149. return handled.includes(s)
  150. }
  151. let folderBusy = (f) => {
  152. for (r of running) { if (r.folder == f) { return true } }
  153. return false
  154. }
  155. let pushEvent = (ev) => {
  156. for (e of events) { if (e.id == ev.id) { return null } }
  157. events.push(ev)
  158. return null
  159. }
  160. // the same rules as the single-host `run` (daemon.hl quotaBlock + the reserve / week checks of iterate)
  161. let quotaVerdict = (q) => {
  162. if (q == null || q.error != null || q.percent == null) { return 'the quota is UNKNOWN (' + (q == null ? 'no reading' : q.error) + ')' }
  163. let qt = 'quota ' + q.percent + '% of 5 h' + (q.resetsMs == null ? '' : ' (resets ' + ago(q.resetsMs) + ')') + (q.weekly == null ? '' : ', week ' + q.weekly + '%')
  164. if (q.percent >= 100) { return qt + ' → the 5 h limit is reached (100%)' }
  165. if (q.weekly != null && q.weekly >= 100) { return qt + ' → the weekly limit is reached (100%)' }
  166. if (reserve < 100 && q.percent >= reserve) { return qt + ' ≥ reserve ' + reserve + '%' }
  167. if (weekLimit < 100 && (q.weekly == null || q.weekly > weekLimit)) { return qt + (q.weekly == null ? ', week UNKNOWN' : ' > week limit ' + weekLimit + '%') }
  168. return null
  169. }
  170. let pauseForLimit = (resetMs, why) => {
  171. let until = resetMs == null ? now() + LIMIT_RETRY_MS : resetMs
  172. if (limitUntil == null) {
  173. limitUntil = until
  174. console.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)
  175. } else if (until > limitUntil) {
  176. limitUntil = until
  177. console.log('LIMIT: the pause is extended until ' + ago(until) + ' · ' + why)
  178. }
  179. return null
  180. }
  181. let accepting = () => {
  182. if (stopping != null) { return { ok = false why = 'stop requested' } }
  183. if (limitUntil != null) {
  184. if (now() < limitUntil) { return { ok = false why = 'usage limit reached until ' + ago(limitUntil) } }
  185. console.log('LIMIT: resumed — the pause until ' + ago(limitUntil) + ' is over; the quota is read again')
  186. limitUntil = null
  187. lastProbe = null
  188. }
  189. if (!quotaOk) { return { ok = false why = quotaWhy } }
  190. return { ok = true why = null }
  191. }
  192. let flushed = false
  193. let sync = null
  194. let checkEnd = () => {
  195. if (finished || stopping == null || running.length > 0) { return null }
  196. // one last try to hand over what the sessions reported while stopping
  197. if (events.length > 0 && !flushed) {
  198. flushed = true
  199. sync()
  200. return null
  201. }
  202. finished = true
  203. if (loop != null) { loop.stop() }
  204. if (watch != null) { watch.stop() }
  205. console.log('AGENT STOPPED — stop requested (' + stopping + '), no session running' + (events.length > 0 ? ' · ' + events.length + ' event(s) not delivered (the scheduler was not reachable)' : ''))
  206. return null
  207. }
  208. let sessionEnded = (entry, r) => {
  209. let out = []
  210. for (x of running) { if (x.session != entry.session) { out.push(x) } }
  211. running = out
  212. let total = (r.workerCost == null ? 0 : r.workerCost) + (r.controllerCost == null ? 0 : r.controllerCost)
  213. 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)
  214. 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 })
  215. if (r.limited == true) { pauseForLimit(r.resetMs, entry.ref + ' ' + entry.session + ' was parked on it (' + r.limitText + ')') }
  216. checkEnd()
  217. return null
  218. }
  219. let refuse = (c, why) => {
  220. console.log('AGENT: refused ' + c.kind + ' ' + c.project + '#' + c.number + ' (' + c.session + ') — ' + why)
  221. pushEvent({ id = c.session + ':refused' kind = 'refused' session = c.session ref = c.project + '#' + c.number why = why })
  222. return null
  223. }
  224. let startCommand = (c) => {
  225. if (known(c.session)) { return null }
  226. let ref = c.project + '#' + c.number
  227. let a = accepting()
  228. if (!a.ok) { return refuse(c, 'this agent is not accepting: ' + a.why) }
  229. let slot = freeSlot()
  230. if (slot < 0) { return refuse(c, 'no slot / port range free (' + running.length + '/' + parallel + ' running)') }
  231. let ports = rangeOf(slot)
  232. let entry = { ref = ref session = c.session folder = c.folder slot = slot ports = ports proc = null phase = 'worker' parkReason = null }
  233. // antcolony#28: process, phase and park reason live in `live[session]`, not on `entry` (a push into `running` copies it)
  234. live[c.session] = { proc = null phase = 'worker' parkReason = null }
  235. let hooks = {
  236. onProc = (p, phase) => {
  237. live[c.session].proc = p
  238. live[c.session].phase = phase
  239. return null
  240. }
  241. parkReason = () => { return live[c.session].parkReason }
  242. }
  243. handled.push(c.session)
  244. let ok = false
  245. if (c.kind == 'resume') {
  246. if (!exists(runs + '/' + c.session + '/session.json')) { return refuse(c, 'no run folder ' + runs + '/' + c.session + ' on ' + host + ' (parked on another host?)') }
  247. let sj = JSON.parse(readFile(runs + '/' + c.session + '/session.json'))
  248. entry.folder = sj.folder
  249. if (folderBusy(entry.folder)) { return refuse(c, 'dev folder ' + entry.folder + ' is busy') }
  250. running.push(entry)
  251. ok = runResume(runs + '/' + c.session, pass, (r) => {
  252. sessionEnded(entry, r)
  253. return null
  254. }, hooks, ports)
  255. } else {
  256. if (c.folder != null && folderBusy(c.folder)) { return refuse(c, 'dev folder ' + c.folder + ' is busy') }
  257. let briefFile = runs + '/.briefs/' + c.session + '.md'
  258. writeFile(briefFile, c.brief, 384)
  259. running.push(entry)
  260. let wargv = ['work' c.project '' + c.number '--post' '--session' c.session '--ports' ports '--brief-file' briefFile]
  261. for (o of pass) { wargv.push(o) }
  262. ok = runWork(wargv, (r) => {
  263. sessionEnded(entry, r)
  264. return null
  265. }, hooks)
  266. }
  267. if (!ok) {
  268. let out = []
  269. for (x of running) { if (x.session != entry.session) { out.push(x) } }
  270. running = out
  271. return refuse(c, 'the worker did not start (see the lines above)')
  272. }
  273. console.log('AGENT: ' + (c.kind == 'resume' ? 'resumed ' : 'started ') + ref + ' (' + c.session + ', ports ' + ports + ') · running ' + running.length + '/' + parallel)
  274. return null
  275. }
  276. // quota: read at most every --quota-every s, only while a slot is free and nothing is paused / stopping
  277. let maybeProbe = () => {
  278. if (probing || stopping != null || running.length >= parallel) { return null }
  279. if (limitUntil != null && now() < limitUntil) { return null }
  280. if (lastProbe != null && now() - lastProbe < quotaEvery * 1000) { return null }
  281. probing = true
  282. lastProbe = now()
  283. probeUsage(quotaBin, runs, (q) => {
  284. probing = false
  285. reading = q == null || q.error != null || q.percent == null ? null : { percent = q.percent weekly = q.weekly resetsMs = q.resetsMs at = now() }
  286. let v = quotaVerdict(q)
  287. let was = quotaOk
  288. quotaOk = v == null
  289. quotaWhy = v == null ? null : v
  290. if (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)) }
  291. return null
  292. })
  293. return null
  294. }
  295. sync = () => {
  296. if (syncing || finished) { return null }
  297. syncing = true
  298. n = n + 1
  299. maybeProbe()
  300. let a = accepting()
  301. let list = []
  302. for (r of running) { list.push({ ref = r.ref session = r.session folder = r.folder phase = r.phase ports = r.ports }) }
  303. let payload = { host = host account = account quota = reading limitUntil = limitUntil capabilities = caps parallel = parallel accepting = a.ok why = a.why running = list events = events }
  304. let res = fetch(schedText + '/agent/sync', { method = 'POST' headers = { 'user-agent' = userAgent() authorization = 'Bearer ' + at.token } json = payload timeoutMs = 15000 })
  305. if (res == null || res.status != 200) {
  306. let why = res == null ? 'no answer' : 'answered ' + res.status
  307. if (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') }
  308. connected = false
  309. syncing = false
  310. return null
  311. }
  312. if (connected != true) { console.log('AGENT: connected to ' + schedText + (connected == false ? ' again' : '') + ' · ' + events.length + ' buffered event(s) delivered') }
  313. connected = true
  314. let j = res.json()
  315. let keep = []
  316. for (e of events) { if (!j.ack.includes(e.id)) { keep.push(e) } }
  317. events = keep
  318. for (c of j.commands) { startCommand(c) }
  319. syncing = false
  320. checkEnd()
  321. return null
  322. }
  323. let checkStop = () => {
  324. if (finished || stopFile == null || stopFile == '' || !exists(stopFile)) { return null }
  325. let mode = readFile(stopFile).trim()
  326. if (stopping == null) {
  327. stopping = 'finish'
  328. console.log('AGENT: stop requested — accepting nothing more; ' + running.length + ' running session(s) finish (a second SIGTERM/SIGINT parks them instead)')
  329. }
  330. if (mode == 'park' && stopping != 'park') {
  331. stopping = 'park'
  332. console.log('AGENT: second stop signal — parking ' + running.length + ' running session(s)')
  333. for (r of running) {
  334. if (live[r.session] != null) {
  335. live[r.session].parkReason = 'the agent was stopped (second SIGTERM/SIGINT)'
  336. if (live[r.session].proc != null) {
  337. console.log('AGENT: killing ' + live[r.session].phase + ' pid ' + live[r.session].proc.pid + ' of ' + r.ref + ' ' + r.session + ' → parked')
  338. live[r.session].proc.kill()
  339. }
  340. }
  341. }
  342. }
  343. checkEnd()
  344. return null
  345. }
  346. watch = every(0.5)
  347. on watch.tick(x) {
  348. checkStop()
  349. return null
  350. }
  351. loop = every(poll)
  352. on loop.tick(x) {
  353. sync()
  354. return null
  355. }
  356. sync()
  357. return true
  358. }

Branches

Latest commits

  • 7f9660eeState of 2026-09-27, before the move to gitoriamre