antcolony
All repositories: gitoria
34.2 KB
// tests/e2e-agent.mjs — mission 035 (antcolony#14): the scheduler on its own host + an AGENT per host, against the SAME kind// of own tickets + ident as tests/e2e.mjs (never the live ones), FAKE claude only.// COLONY_E2E_PORT_BASE=8710 node tests/e2e-agent.mjs (ident = base+1, tickets = base+2, scheduler = base+3, agent pool base+4 … base+9)// Checks: no agent → nothing starts; a wrong token is refused; agent connects → a ticket is worked THERE and posted; the brief// carries the agent's port range and the scheduler's assembly; the scheduler dies mid-session → the session goes on, its end// is delivered after the reconnect; the quota is the agent's (no start above the reserve); an agent that stops is OFFLINE.import { spawn, spawnSync } from 'node:child_process';import { cpSync, mkdirSync, rmSync, writeFileSync, readFileSync, existsSync, readdirSync } from 'node:fs';import { dirname, join, resolve } from 'node:path';import { fileURLToPath } from 'node:url';import { createServer } from 'node:http';process.env.COLONY_SANDBOX = process.env.COLONY_SANDBOX || 'off'; // antcolony#18: the fake claude lives outside any box (tests/e2e-sandbox.mjs tests the box)const HERE = dirname(fileURLToPath(import.meta.url));const APP = resolve(HERE, '..');const WORK = join(APP, '.scratch', 'e2e-agent');const PB = Number(process.env.COLONY_E2E_PORT_BASE || 8710);const IDENT_PORT = PB + 1, TICKETS_PORT = PB + 2;if (PB < 8700 || TICKETS_PORT > 8799) throw new Error('ports must stay in 8700-8799');const TICKETS_DIR = process.env.COLONY_E2E_TICKETS_DIR || '/media/STORAGE/projects/tickets.worldapi.org';const IDENT_DIR = process.env.COLONY_E2E_IDENT_DIR || '/media/STORAGE/projects/ident.worldapi.org';const BASE = `http://127.0.0.1:${TICKETS_PORT}`;const J = JSON.stringify;const sleep = (ms) => new Promise(r => setTimeout(r, ms));let passes = 0, failures = 0;function check(label, ok, detail = '') {console.log(`${ok ? 'ok ' : 'FAIL'} ${label}${ok ? '' : '\n ' + String(detail).split('\n').join('\n ')}`);if (ok) passes++; else failures++;}// mission 023 (README "Ticket texts are written for the creator", hard limit: at most 5 short lines): the lines// the creator reads = the non-empty lines of a text without the machine lines (`colony-…: …`)// mission 027: the Test steps are a numbered Markdown list ("1. …") under ONE "**Test:**" line — the header counts as one// line (as the one-line "**Test:** 1) … 2) …" did before), the 2–4 step lines are not countedconst readableLines = (t) => String(t || '').split('\n').filter(l => l.trim() !== '' && !/^colony-[a-z-]+: /.test(l) && !/^\d+\. /.test(l)).length;// ---- copies ----------------------------------------------------------------------------------rmSync(WORK, { recursive: true, force: true });mkdirSync(WORK, { recursive: true });// mission 030 (ticket antcolony#1): TEST ISOLATION. The Hybriel interpreter loads the `.env` BESIDE THE ENTRY SCRIPT// (scheduler.hl — not the cwd; the real environment wins, an EMPTY variable counts as set): run from this folder, every// `./colony` call of the e2e got the LIVE settings (max turns, week limit, token file, CLAUDE_CONFIG_DIR …) → 6 failures.// So the scheduler under test runs from a COPY of this folder WITHOUT `.env` (and without the live runs / librarian /// registry / logs): whatever the live `.env` holds — today or tomorrow — cannot reach a test. On top, the environment// every call gets is DEFINED: no inherited COLONY_* / FAKE_* variable of the caller's shell (except the COLONY_E2E_*// controls above), and — unless COLONY_E2E_REAL=1 — a `claude` guard first on PATH that refuses and logs (fake only).const CODE = join(WORK, 'app');const CODE_SKIP = new Set(['.env', '.scratch', 'runs', 'sessions', 'briefs', 'logs', 'librarian', 'projects', 'mcp-needs-auth-cache.json', '.git']);// entry by entry (node refuses to copy a folder into its own subfolder; .scratch is skipped anyway)mkdirSync(CODE, { recursive: true });for (const e of readdirSync(APP)) if (!CODE_SKIP.has(e)) cpSync(join(APP, e), join(CODE, e), { recursive: true });if (existsSync(join(CODE, '.env'))) throw new Error('a .env was copied into the scheduler under test — refusing to start');const REAL = process.env.COLONY_E2E_REAL === '1';const GUARD_DIR = join(WORK, 'claude-guard'), GUARD_LOG = join(WORK, 'claude-guard.log');mkdirSync(GUARD_DIR, { recursive: true });writeFileSync(join(GUARD_DIR, 'claude'), `#!/bin/sh\necho "$(date -u +%FT%TZ) cwd=$(pwd) args=$*" >> '${GUARD_LOG}'\necho 'e2e guard: the REAL claude was called — refused (fake claude only; COLONY_E2E_REAL=1 allows it)' >&2\nexit 97\n`, { mode: 0o755 });const BASE_ENV = { ...Object.fromEntries(Object.entries(process.env).filter(([k]) => !(/^(COLONY_|FAKE_)/.test(k) && !k.startsWith('COLONY_E2E_')))), COLONY_SANDBOX: 'off' };if (!REAL) BASE_ENV.PATH = GUARD_DIR + ':' + (process.env.PATH || '');const SKIP = new Set(['.env', 'storage', '.sessions', '.scratch', 'server.log', 'server.pid', 'testapp', '.git']);const TCODE = join(WORK, 'tickets-code');cpSync(TICKETS_DIR, TCODE, { recursive: true, filter: (src) => {const rel = src.slice(TICKETS_DIR.length).replace(/^\/+/, '');return rel === '' || !SKIP.has(rel.split('/')[0]);} });if (existsSync(join(TCODE, '.env'))) throw new Error('a .env was copied — refusing to start');const { startIdent } = await import(join(TCODE, 'tests', 'identkit.mjs'));let ident = null, tickets = null, tlog = '';const daemons = []; // mission 020: `colony run` processes (wrapper + scheduler pid) — killed at the end if still thereconst fakePids = () => { try { return readdirSync(join(WORK, 'fake-calls')).map(f => Number(f.split('-')[1])); } catch { return []; } };// a pid is killed only while its command line is still ours (pids get reused)const isOurs = (pid, needle) => { try { return readFileSync(`/proc/${pid}/cmdline`, 'utf8').includes(needle); } catch { return false; } };const killOurs = (pid, needle) => { if (pid && isOurs(pid, needle)) { try { process.kill(pid, 'SIGKILL'); return true; } catch {} } return false; };// mission 026: servers the FAKE sessions leave behind (FAKE_SERVER_LOG) + the e2e's own "pre-existing" server — the ones the// scheduler must NOT stop are killed here by the test (only while their command line is still ours)const SERVER_LOG = join(WORK, 'fake-servers.jsonl');const ownServerPids = [];const serverRows = () => { try { return readFileSync(SERVER_LOG, 'utf8').trim().split('\n').filter(Boolean).map(l => JSON.parse(l)); } catch { return []; } };const killTestServers = () => {let n = 0;for (const r of serverRows()) if (killOurs(r.pid, 'fake-left-server')) n++;for (const pid of ownServerPids) if (killOurs(pid, 'e2e-preexisting-server')) n++;return n;};const cleanup = async () => {const ks = killTestServers(); if (ks) console.log('cleanup: killed ' + ks + ' test server(s) left over');for (const d of daemons) { if (d.proc.exitCode === null && d.proc.signalCode === null) { try { d.proc.kill('SIGKILL'); } catch {} } killOurs(d.schedPid, 'scheduler.hl'); }for (const pid of fakePids()) { if (killOurs(pid, 'fake-claude.mjs')) console.log('cleanup: killed a leftover fake claude pid ' + pid); }if (tickets) { try { tickets.kill('SIGTERM'); } catch {} await new Promise(r => { if (tickets.exitCode !== null || tickets.signalCode !== null) return r(); tickets.once('exit', r); setTimeout(r, 3000); }); }if (ident) await ident.stop();writeFileSync(join(WORK, 'tickets.log'), tlog);if (ident) writeFileSync(join(WORK, 'ident.log'), ident.log());};// ---- helpers ---------------------------------------------------------------------------------let TOKEN = null;// antcolony#31: tickets speaks the new states (progress / review / pending / done / reopened); this suite's expectations are// written in the scheduler's own vocabulary (in progress / awaiting creator / on hold / confirmed / rejected — lib/tickets.hl// translates the same way), so every JSON read is translated back: a "Question…" ticket that is pending = awaiting creator.const OLD = { progress: 'in progress', review: 'awaiting creator', pending: 'on hold', done: 'confirmed', reopened: 'rejected' };const oldNames = (v) => {if (Array.isArray(v)) { v.forEach(oldNames); return v; }if (v === null || typeof v !== 'object') return v;if (typeof v.state === 'string' && OLD[v.state]) v.state = (v.state === 'pending' && /^Question\b/.test(v.subject || '')) ? 'awaiting creator' : OLD[v.state];// the server adds an `assign` event when a ticket goes to review / pending with nobody assigned (tickets#20): not part of these checksif (Array.isArray(v.events)) v.events = v.events.filter(e => e.kind !== 'assign');if (v.kind === 'state') { if (OLD[v.from]) v.from = OLD[v.from]; if (OLD[v.to]) v.to = OLD[v.to]; }for (const k of Object.keys(v)) if (v[k] && typeof v[k] === 'object') oldNames(v[k]);return v;};const api = async (method, path, body, accept) => {const headers = { 'user-agent': 'colony-e2e' };if (body) headers['content-type'] = 'application/json';if (accept) headers.accept = accept;if (method !== 'GET') headers.authorization = 'Bearer ' + TOKEN;const r = await fetch(BASE + path, { method, headers, body: body ? J(body) : undefined });const text = await r.text();let json = null; try { json = oldNames(JSON.parse(text)); } catch {}if (method !== 'GET' && !(r.status === 200 || r.status === 201)) throw new Error(`${method} ${path} → ${r.status} ${text}`);return { status: r.status, json, text, ct: r.headers.get('content-type') };};let emitI = 0;const temit = async (event, payload, cookie) => {const r = await fetch(BASE + '/__hl/emit', { method: 'POST', headers: { 'content-type': 'application/json', cookie }, body: J({ t: 'emit', i: ++emitI, event, payload }) });const raw = await r.text(); let j = null; try { j = JSON.parse(raw); } catch {}return { value: j && j.value, raw };};const REG = join(WORK, 'projects');const TOKFILE = join(WORK, 'colony-token');// mission 029: the librarian's folder and the global decisions file always point into .scratch/e2e — no e2e call (`run`,// `brief`, `work`, `cycle` run the librarian first) ever writes the real librarian/ or templates/const colonyEnv = (extra = {}) => ({ ...BASE_ENV, COLONY_TICKETS_URL: BASE, COLONY_PROJECTS_DIR: REG, COLONY_TOKEN_FILE: TOKFILE, COLONY_USER_AGENT: 'colony-e2e',COLONY_LIBRARIAN_DIR: join(WORK, 'librarian'), COLONY_DECISIONS_GLOBAL: join(WORK, 'decisions-global.md'), ...extra });const colony = (args, extra, timeout = 60000) => {const r = spawnSync(join(CODE, 'colony'), args, { cwd: WORK, env: colonyEnv(extra), encoding: 'utf8', timeout });return (r.stdout || '') + (r.stderr || '');};const snapshot = async () => {const all = (await api('GET', '/api/tickets')).json.tickets;const out = [];for (const t of all) {const d = (await api('GET', `/api/projects/${t.project}/tickets/${t.number}`)).json;out.push({ ref: t.project + t.ref, state: t.state, updatedMs: t.updatedMs, events: d.events.map(e => [e.seq, e.kind, e.author, e.text]) });}out.sort((a, b) => a.ref < b.ref ? -1 : 1);return J(out);};const ticket = async (p, n) => (await api('GET', `/api/projects/${p}/tickets/${n}`)).json;let CTOKEN = null; // mission 022: the creator's API tokenconst capi = async (method, path, body) => { const save = TOKEN; TOKEN = CTOKEN; try { return await api(method, path, body); } finally { TOKEN = save; } };try {// ---- ident + tickets (own instances) --------------------------------------------------------ident = await startIdent({ identDir: IDENT_DIR, workDir: join(WORK, 'ident'), port: IDENT_PORT });const colonyAcct = await ident.signIn('[email protected]');const app = await ident.registerApp(colonyAcct, 'tickets (colony e2e)', [BASE]);// mission 022: a CREATOR (only the creator may reject / confirm) — its per-app id is TICKETS_CREATOR_IDENTITYconst creatorAcct = await ident.signIn('[email protected]');const CREATOR_ID = await ident.exchange(app, await ident.selectorCode(creatorAcct, app, BASE));check('own ident runs from a copy without .env, tickets registered in it', /^pk_/.test(app.key) && !existsSync(join(WORK, 'ident', 'ident-code', '.env')));const store = join(WORK, 'tickets-data');tickets = spawn(join(TCODE, 'bin/hybriel'), ['project.hl'], { cwd: TCODE, stdio: ['ignore', 'pipe', 'pipe'], env: { ...process.env,TICKETS_PORT: String(TICKETS_PORT), TICKETS_STORAGE: join(store, 'mpackdb'), TICKETS_SESSIONS: join(store, 'sessions'), HL_HOST: '127.0.0.1', TICKETS_WATCH: '0',IDENT_URL: ident.base, IDENT_EXCHANGE_URL: ident.base, TICKETS_PUBLIC_URL: BASE, IDENT_API_KEY: app.key, IDENT_API_SECRET: app.secret, TICKETS_CREATOR_IDENTITY: CREATOR_ID } });tickets.stdout.on('data', d => tlog += d); tickets.stderr.on('data', d => tlog += d);let up = false;for (let i = 0; i < 80 && !up; i++) { try { up = (await fetch(BASE + '/')).ok; } catch {} if (!up) await sleep(250); }if (!up) throw new Error('tickets did not come up\n' + tlog);// the user "Colony": login button's server half → display name → API token (as the creator will do on /you)const first = await fetch(BASE + '/');const cookie = (first.headers.get('set-cookie') || '').split(';')[0];const cb = await fetch(BASE + '/login/callback?ident_code=' + await ident.selectorCode(colonyAcct, app, BASE), { headers: { cookie }, redirect: 'manual' });await temit('saveDisplayName', ['Colony'], cookie);const tk = await temit('tokenCreate', ['scheduler e2e'], cookie);TOKEN = tk.value && tk.value.token;check('tickets (own copy) runs; user "Colony" logged in via ident and made an API token', cb.status === 302 && /^tkt_[0-9a-f]{48}$/.test(TOKEN || ''), cb.status + ' ' + tk.raw);writeFileSync(TOKFILE, TOKEN + '\n', { mode: 0o600 });// antcolony#31: projects are records with members — the creator (admin) opens a project the first time a ticket is made in it// and makes Colony a member with the role `edit`const ccookie = (await fetch(BASE + '/')).headers.get('set-cookie').split(';')[0];await fetch(BASE + '/login/callback?ident_code=' + await ident.selectorCode(creatorAcct, app, BASE), { headers: { cookie: ccookie }, redirect: 'manual' });await temit('saveDisplayName', ['Creator'], ccookie);const CTOKEN = (await temit('tokenCreate', ['creator e2e'], ccookie)).value.token;const opened = new Set();const openProject = async (slug) => {if (opened.has(slug)) return;opened.add(slug);for (const [path, body] of [['/api/projects', { title: slug, slug }], [`/api/projects/${slug}/members`, { user: 'Colony', role: 'edit' }]]) {const r = await fetch(BASE + path, { method: 'POST', headers: { 'user-agent': 'colony-e2e', 'content-type': 'application/json', authorization: 'Bearer ' + CTOKEN }, body: J(body) });if (r.status > 201) throw new Error(`${path} → ${r.status} ${await r.text()}`);}};// ---- registry + the dev folder of the agent's host ------------------------------------------------------------const REGA = join(WORK, 'projects-agent');const AGDEV = join(WORK, 'agdev');mkdirSync(REGA, { recursive: true }); mkdirSync(AGDEV, { recursive: true });writeFileSync(join(AGDEV, 'CONCEPT.md'), '# ag\nA test project of the agent split: every ticket is a tiny file task in this folder.\n');writeFileSync(join(AGDEV, 'README.md'), '# ag\nTest project. Nothing else to read.\n');writeFileSync(join(REGA, 'ag.json'), J({ name: 'ag', dev: { host: 'e2e-host', folder: AGDEV }, concept: AGDEV + '/CONCEPT.md', dependsOn: [], needs: ['Chrome'] }, null, 1));writeFileSync(join(REGA, 'gp.json'), J({ name: 'gp', dev: { host: 'e2e-host', folder: join(WORK, 'gpdev') }, concept: AGDEV + '/CONCEPT.md', dependsOn: [], needs: ['chrome', 'gpu'] }, null, 1));writeFileSync(join(REGA, 'sh.json'), J({ name: 'sh', dev: { host: 'e2e-b', folder: join(WORK, 'shdev') }, concept: AGDEV + '/CONCEPT.md', dependsOn: [] }, null, 1));writeFileSync(join(REGA, 'sc.json'), J({ name: 'sc', dev: { host: 'e2e-c', folder: join(WORK, 'scdev') }, concept: AGDEV + '/CONCEPT.md', dependsOn: [] }, null, 1));mkdirSync(join(WORK, 'shdev'), { recursive: true }); mkdirSync(join(WORK, 'scdev'), { recursive: true });writeFileSync(join(REGA, 'sd.json'), J({ name: 'sd', dev: { host: 'e2e-c', folder: join(WORK, 'sddev') }, concept: AGDEV + '/CONCEPT.md', dependsOn: [] }, null, 1));mkdirSync(join(WORK, 'sddev'), { recursive: true });await openProject('antcolony'); // the digest / question inbox of the schedulerconst mk = async (project, subject, summary) => { await openProject(project); const t = (await api('POST', `/api/projects/${project}/tickets`, { subject, summary })).json.ticket; await sleep(30); return t; };const FAKE = join(CODE, 'tests', 'fake-claude.mjs');const RUNS_S = join(WORK, 'runs-scheduler'), RUNS_A = join(WORK, 'runs-agent');const FLD = join(WORK, 'fake-calls'), FLOG = join(WORK, 'fake-claude.json');const UF = join(WORK, 'usage-percent'), WF = join(WORK, 'week-percent');writeFileSync(UF, '5\n'); writeFileSync(WF, '41\n');const AGTOKEN = join(WORK, 'agent-token');writeFileSync(AGTOKEN, 'agent-secret-e2e\n', { mode: 0o600 });const SPORT = PB + 3, SURL = `http://127.0.0.1:${SPORT}`;const sEnv = (extra = {}) => ({ COLONY_HOST: 'sched-host', COLONY_RUNS_DIR: RUNS_S, COLONY_PROJECTS_DIR: REGA, COLONY_AGENT_TOKEN_FILE: AGTOKEN, COLONY_LIBRARIAN: 'off',COLONY_CLAUDE: FAKE, FAKE_USAGE_FILE: UF, FAKE_WEEK_FILE: WF, FAKE_LOG_DIR: join(WORK, 'fake-calls-scheduler'), ...extra });const aEnv = (extra = {}) => ({ COLONY_HOST: 'e2e-host', COLONY_CAPABILITIES: 'chrome, Deploy', COLONY_RUNS_DIR: RUNS_A, COLONY_PROJECTS_DIR: REGA, COLONY_AGENT_TOKEN_FILE: AGTOKEN, COLONY_SCHEDULER_URL: SURL,COLONY_CLAUDE: FAKE, FAKE_CLAUDE_MODE: 'good', FAKE_CONTROLLER_MODE: 'pass', FAKE_USAGE_FILE: UF, FAKE_WEEK_FILE: WF, FAKE_LOG_DIR: FLD, FAKE_CLAUDE_LOG: FLOG, FAKE_CONTROLLER_LOG: '',COLONY_PORT_POOL: `${PB + 4}-${PB + 9}`, COLONY_PORTS_PER_SESSION: '3', COLONY_AGENT_POLL: '1', FAKE_CLAUDE_SLEEP: '4', COLONY_AGENT_QUOTA_EVERY: '2', ...extra });const start = (args, env) => {const proc = spawn(join(CODE, 'colony'), args, { cwd: WORK, env: colonyEnv(env), stdio: ['ignore', 'pipe', 'pipe'] });const d = { proc, text: '', schedPid: null, exited: null };d.exited = new Promise(r => proc.once('exit', (code, sig) => r({ code, sig })));const add = (b) => { d.text += b; const m = d.text.match(/scheduler pid (\d+)/); if (m) d.schedPid = Number(m[1]); };proc.stdout.on('data', add); proc.stderr.on('data', add);d.waitFor = async (re, ms = 30000) => { const t0 = Date.now(); while (Date.now() - t0 < ms) { const m = d.text.match(re); if (m) return m; await sleep(100); } return null; };d.kill = () => { try { proc.kill('SIGKILL'); } catch {} killOurs(d.schedPid, 'scheduler.hl'); };d.term = () => { try { proc.kill('SIGTERM'); } catch {} };daemons.push(d);return d;};const waitState = async (p, n, st, ms = 60000) => { const t0 = Date.now(); while (Date.now() - t0 < ms) { if ((await ticket(p, n)).ticket.state === st) return true; await sleep(300); } return false; };const post = async (path, token, method = 'POST', body = '{"host":"x","parallel":1,"running":[],"events":[]}') => { const r = await fetch(SURL + path, { method, headers: token ? { authorization: 'Bearer ' + token } : {}, body: method === 'POST' ? body : undefined }); return { status: r.status, text: await r.text() }; };const workerCalls = (dir) => { try { return readdirSync(dir).filter(f => !/-usage\.json$/.test(f)).length; } catch { return 0; } };// ---- 1. the scheduler alone: no agent → nothing starts, the ticket waits --------------------------------------------const t1 = await mk('ag', 'agent task one', 'Create a file `one.txt` with the content `one`. Verify it.');const tg = await mk('gp', 'needs a gpu', 'Create a file `gpu.txt`.');const sched = start(['run', '--serve', String(SPORT), '--interval', '2', '--agent-ttl', '4'], sEnv());check('scheduler: RUN line names the agent door and that it starts no worker itself', !!(await sched.waitFor(/RUN: started on sched-host · SCHEDULER for agents on 127\.0\.0\.1:\d+ \(starts no worker itself/, 15000)), sched.text);check('scheduler alone: the ticket waits — "no agent online on e2e-host"; still open; no Claude call', !!(await sched.waitFor(new RegExp(`ag#${t1.number} waits — no agent online on e2e-host`, 'm'), 15000)) && (await ticket('ag', t1.number)).ticket.state === 'open' && workerCalls(join(WORK, 'fake-calls-scheduler')) === 0, sched.text);// ---- 2. the door: wrong / missing token, wrong path -------------------------------------------------------------------let r = await post('/agent/sync', 'wrong');check('the door refuses a wrong token (401)', r.status === 401, J(r));r = await post('/agent/sync', null);check('the door refuses a request without a token (401)', r.status === 401, J(r));r = await post('/nope', 'agent-secret-e2e');check('the door answers 404 for any other path', r.status === 404, J(r));r = await post('/agent/sync', 'agent-secret-e2e', 'POST', '{"host":');check('the door answers 400 for a body that is not JSON', r.status === 400, J(r));const noRefuse = !/AGENT: x connected/.test(sched.text);check('a refused request registered no agent', noRefuse, sched.text);// ---- 3. an agent connects → the ticket is worked ON ITS HOST and posted ---------------------------------------------------const agent = start(['agent', '--poll', '1'], aEnv());check('agent: AGENT line (host, scheduler, poll, quota rule, ports)', !!(await agent.waitFor(/AGENT: started on e2e-host → scheduler http:\/\/127\.0\.0\.1:\d+ · poll 1 s · parallel 1 · offers chrome, deploy · quota reserve 70%/, 15000)), agent.text);check('agent: connects, and reads the quota → accepting work', !!(await agent.waitFor(/AGENT: connected to /, 20000)) && !!(await agent.waitFor(/AGENT: quota ok \(5% of 5 h, week 41%\) → accepting work/, 20000)), agent.text);check('capabilities: the scheduler shows what the agent offers, and gp (needs gpu) waits naming what is missing', !!(await sched.waitFor(/AGENT: e2e-host connected · offers chrome, deploy · /, 15000)) && !!(await sched.waitFor(new RegExp(`gp#${tg.number} waits — e2e-host does not offer: gpu`), 15000)) && (await ticket('gp', tg.number)).ticket.state === 'open', sched.text);check('scheduler: "AGENT: e2e-host connected"', !!(await sched.waitFor(/AGENT: e2e-host connected · offers chrome, deploy · parallel 1/, 15000)), sched.text);check('the ticket is worked by the agent: → in progress → awaiting creator (report + controller pass)', await waitState('ag', t1.number, 'awaiting creator', 90000), agent.text + '\n----\n' + sched.text);const d1 = await ticket('ag', t1.number);const rep1 = d1.events.filter(e => e.kind === 'comment').map(e => e.text).join('\n');check('the comment carries the report marker + "controller pass"; the lease was posted by the agent\'s host', /colony-report: s-\d{8}T\d{4}-[0-9a-f]{6} · sha256 [0-9a-f]{12} · controller pass/.test(rep1) && d1.events.some(e => /colony-lease: s-\S+ · host e2e-host · started/.test(e.text || '')), J(d1.events.map(e => e.text)));const s1 = (rep1.match(/colony-report: (s-\d{8}T\d{4}-[0-9a-f]{6})/) || [])[1];check('the worker ran in the AGENT\'s runs folder, not the scheduler\'s', !!s1 && existsSync(join(RUNS_A, s1, 'report.json')) && !existsSync(join(RUNS_S, s1)), s1);const brief1 = s1 && existsSync(join(RUNS_A, s1, 'brief.md')) ? readFileSync(join(RUNS_A, s1, 'brief.md'), 'utf8') : '';check('the brief the scheduler assembled carries THIS agent\'s port range (no marker left), the ticket, the conventions', new RegExp(`Ports for everything you start: ${PB + 4}–${PB + 6} only`).test(brief1) && !/@PORT_(FROM|TO)@|@AGENT@/.test(brief1) && brief1.includes('agent task one') && brief1.includes('## Conventions for all apps') && brief1.includes(`session \`${s1}\``), brief1.slice(0, 600));check('the SCHEDULER started no Claude (worker / controller) itself', workerCalls(join(WORK, 'fake-calls-scheduler')) === 0, '');check('the FAKE worker + controller of the agent ran with the agent\'s port env', existsSync(FLD) && readdirSync(FLD).some(f => f.endsWith('.json') && !/-usage\.json$/.test(f)), '');check('scheduler + agent both print SESSION END (posted); the scheduler names the host', !!(await agent.waitFor(new RegExp(`SESSION END ag#${t1.number} ${s1} · posted`))) && !!(await sched.waitFor(new RegExp(`SESSION END ag#${t1.number} ${s1} · posted[^\\n]* · on e2e-host`))), agent.text + '\n----\n' + sched.text);check('the scheduler holds no state of its own: no run folder for the session, no file written beyond its probe folder', !existsSync(join(RUNS_S, s1 || 'none')), '');// ---- 4. the scheduler dies mid-session: the work goes on, the end is delivered after the reconnect -------------------------const t2 = await mk('ag', 'agent task two', 'Create a file `two.txt` with the content `two`. Verify it.');agent.text += '';const mark2 = agent.text.length;// a slow fake worker (5 s): started by the NEXT command — restart nothing, the agent's env is fixed, so this ticket runs at the normal speed;// the scheduler is killed right after the start and comes back only after the session endedconst t2started = await waitState('ag', t2.number, 'in progress', 60000);check('second ticket: the agent started it (in progress)', t2started, agent.text);sched.kill();const t2done = await waitState('ag', t2.number, 'awaiting creator', 90000);check('scheduler killed while the session ran: the worker finished and posted ON ITS OWN (→ awaiting creator)', t2done, agent.text);check('agent: "scheduler … not reachable — 1 session(s) keep running / N event(s) buffered"', !!(await agent.waitFor(/AGENT: scheduler http:\/\/127\.0\.0\.1:\d+ not reachable \(no answer\) — \d session\(s\) keep running, \d event\(s\) buffered/, 20000)), agent.text);check('agent: the SESSION END event is buffered while the scheduler is away', !!(await agent.waitFor(/SESSION END ag#\d+ s-\S+ · posted/, 30000)) && /1 event\(s\) buffered|\d event\(s\) buffered/.test(agent.text.slice(mark2)), agent.text.slice(mark2));const sched2 = start(['run', '--serve', String(SPORT), '--interval', '2', '--agent-ttl', '4', '--account-settle', '3'], sEnv());check('scheduler restarted: the agent reconnects on its own ("connected … again · 1 buffered event(s) delivered")', !!(await agent.waitFor(/AGENT: connected to http:\/\/127\.0\.0\.1:\d+ again · [1-9]\d* buffered event\(s\) delivered/, 30000)), agent.text.slice(mark2));check('scheduler restarted: it learns the agent + prints the session end it missed, exactly once', !!(await sched2.waitFor(/AGENT: e2e-host connected/, 15000)) && !!(await sched2.waitFor(new RegExp(`SESSION END ag#${t2.number} s-\\S+ · posted`), 15000)) && (sched2.text.match(new RegExp(`SESSION END ag#${t2.number} `, 'g')) || []).length === 1, sched2.text);check('nothing was run twice: ag#2 has exactly one lease start and one report', (await ticket('ag', t2.number)).events.filter(e => /colony-lease: \S+ · host e2e-host · started/.test(e.text || '')).length === 1 && (await ticket('ag', t2.number)).events.filter(e => /colony-report:/.test(e.text || '')).length === 1, '');// ---- 5. the quota is the agent's: above the reserve nothing is accepted, a start waits ------------------------------------writeFileSync(UF, '80\n');check('quota 80% ≥ reserve 70% (agent side, read every 2 s here): "not accepting work — …"', !!(await agent.waitFor(/AGENT: not accepting work — quota 80% of 5 h.* ≥ reserve 70%/, 30000)), agent.text.slice(-500));const t3 = await mk('ag', 'agent task three', 'Create a file `three.txt` with the content `three`. Verify it.');check('…the scheduler says the start waits (the quota is the agent\'s call)', !!(await sched2.waitFor(new RegExp(`ag#${t3.number} waits — agent e2e-host is not accepting work \\(quota 80%`), 30000)), agent.text.slice(-600) + '\n----\n' + sched2.text.slice(-800));check('…and the ticket stays open, no worker started', (await ticket('ag', t3.number)).ticket.state === 'open', '');writeFileSync(UF, '5\n');check('quota back to 5%: the agent accepts again and the ticket is worked', !!(await agent.waitFor(/AGENT: quota ok \(5%/, 120000)) && await waitState('ag', t3.number, 'awaiting creator', 90000), agent.text.slice(-600));// ---- 5b. antcolony#16: the quota belongs to the ACCOUNT — two hosts (e2e-b, e2e-c: hand-made syncs) share the account "shared" --------const fakeRunning = { 'e2e-b': [], 'e2e-c': [] };const beat = async (host, extra) => {const r = await post('/agent/sync', 'agent-secret-e2e', 'POST', J({ host, account: 'shared', parallel: 3, accepting: true, why: null, capabilities: [], running: fakeRunning[host], events: [], ...extra }));const cmds = r.status === 200 ? JSON.parse(r.text).commands : [];for (const c of cmds) fakeRunning[host].push({ ref: c.project + '#' + c.number, session: c.session, folder: c.folder, phase: 'worker', ports: '1-2' });return cmds;};const rd = (percent, at, weekly = 20) => ({ percent, weekly, resetsMs: null, at });// both hosts sync every second for ms; bx / cx = what each reports (a function, evaluated at every beat); returns { b, c } commands receivedconst rounds = async (ms, bx, cx, until) => { const got = { b: [], c: [] }; const t0 = Date.now(); while (Date.now() - t0 < ms) { got.b.push(...await beat('e2e-b', bx())); got.c.push(...await beat('e2e-c', cx())); if (until && until(got)) break; await sleep(1000); } return got; };const tb = await mk('sh', 'share task b', 'Create a file `b.txt`.');const tc = await mk('sc', 'share task c', 'Create a file `c.txt`.');const at0 = Date.now();let g = await rounds(20000, () => ({ quota: rd(10, at0), limitUntil: null }), () => ({ quota: null, limitUntil: null }), (x) => x.b.length > 0);check('account share: host b starts its ticket (reading 10% of the account, own account name "shared")', g.b.length === 1 && g.b[0].kind === 'start' && g.b[0].project === 'sh' && !!(await sched2.waitFor(/AGENT: e2e-b connected · .* · Claude account shared · /, 5000)), J(g) + sched2.text.slice(-600));g = await rounds(6000, () => ({ quota: rd(10, at0), limitUntil: null }), () => ({ quota: null, limitUntil: null }));check('account share: host c (never read the quota itself) waits — the start on host b is not in a reading yet', g.c.length === 0 && new RegExp(`sc#${tc.number} waits — Claude account shared: a start on it \\(\\S+\\) is not in the quota reading yet — shared by e2e-b, e2e-c`).test(sched2.text), J(g) + sched2.text.slice(-800));g = await rounds(20000, () => ({ quota: rd(12, Date.now()), limitUntil: null }), () => ({ quota: null, limitUntil: null }), (x) => x.c.length > 0);check('account share: a reading taken after the start (12%) → host c starts, on host b\'s reading (its own was empty)', g.c.length === 1 && g.c[0].kind === 'start' && g.c[0].project === 'sc', J(g) + sched2.text.slice(-600));const tc2 = await mk('sd', 'share task c2', 'Create a file `c2.txt`.');g = await rounds(6000, () => ({ quota: rd(12, Date.now()), limitUntil: Date.now() + 60000 }), () => ({ quota: null, limitUntil: null }));check('account share: a usage limit hit on host b stops starts on host c too', g.c.length === 0 && new RegExp(`sd#${tc2.number} waits — Claude account shared: usage limit reached until \\S+ \\(hit on e2e-b\\)`).test(sched2.text), J(g) + sched2.text.slice(-800));g = await rounds(6000, () => ({ quota: rd(80, Date.now()), limitUntil: null }), () => ({ quota: null, limitUntil: null }));check('account share: 80% on the account (read by host b) ≥ reserve 70% holds host c as well', g.c.length === 0 && new RegExp(`sd#${tc2.number} waits — Claude account shared: quota 5 h 80% ≥ reserve 70%`).test(sched2.text), J(g) + sched2.text.slice(-800));check('the ITER line names the account with its reading and its hosts', /agents: .*account shared 80% week 20% \[e2e-b\+e2e-c\]/.test(sched2.text), sched2.text.slice(-800));// ---- 6. the agent stops (first signal: finish) → the scheduler sees it OFFLINE --------------------------------------------agent.term();check('agent: stop requested → nothing more accepted, "AGENT STOPPED"', !!(await agent.waitFor(/AGENT: stop requested/, 20000)) && !!(await agent.waitFor(/AGENT STOPPED — stop requested \(finish\), no session running/, 20000)), agent.text.slice(-500));check('scheduler: "AGENT: e2e-host is OFFLINE (no sync for 4 s)"', !!(await sched2.waitFor(/AGENT: e2e-host is OFFLINE \(no sync for 4 s\)/, 30000)), sched2.text.slice(-500));sched2.term();check('scheduler: SIGTERM → RUN STOPPED, nothing left running', !!(await sched2.waitFor(/RUN STOPPED — stop requested \(finish\)/, 20000)), sched2.text.slice(-400));// ---- refusals of the agent -----------------------------------------------------------------------------------------------const bad = (args, env) => spawnSync(join(CODE, 'colony'), args, { cwd: WORK, env: colonyEnv(env), encoding: 'utf8', timeout: 30000 });let o = bad(['agent'], aEnv({ COLONY_SCHEDULER_URL: '' }));check('agent without a scheduler URL: REFUSED', /^agent: REFUSED — --scheduler URL/m.test(o.stdout), o.stdout + o.stderr);o = bad(['agent'], aEnv({ COLONY_AGENT_TOKEN_FILE: join(WORK, 'nope') }));check('agent without its token file: REFUSED', /^agent: REFUSED — COLONY_AGENT_TOKEN_FILE: no such file/m.test(o.stdout), o.stdout + o.stderr);o = bad(['run', '--serve', '8713', '--once'], sEnv());check('run --serve --once: REFUSED', /^run: REFUSED — --serve and --once do not go together/m.test(o.stdout), o.stdout + o.stderr);o = bad(['run', '--serve', '8713'], sEnv({ COLONY_AGENT_TOKEN_FILE: '' }));check('run --serve without an agent token: REFUSED', /^run: REFUSED — COLONY_AGENT_TOKEN_FILE is not set/m.test(o.stdout), o.stdout + o.stderr);writeFileSync(join(WORK, 'sched-1.txt'), sched.text); writeFileSync(join(WORK, 'sched-2.txt'), sched2.text); writeFileSync(join(WORK, 'agent.txt'), agent.text);if (!REAL) check('the REAL claude was never called (the PATH guard logged nothing)', !existsSync(GUARD_LOG), existsSync(GUARD_LOG) ? readFileSync(GUARD_LOG, 'utf8') : '');} catch (e) {check('e2e-agent ran to the end', false, e.stack || String(e));} finally {await cleanup();}console.log(`\n${passes} passed, ${failures} failed (ports ${IDENT_PORT}/${TICKETS_PORT}/${PB + 3}; logs ${WORK}/)`);process.exit(failures ? 1 : 0);
Branches
- mainmain branch