/** * ADDITIVE SRP Hub research orchestration * ================================================================ * Pure logic + HTTP handlers for research job queue, structured * handoffs, auto-routing, decisions, claims, and overnight briefs. * * MUST NOT replace existing Grok/Claude/Astra workflows, Hub routes, * or frozen SRP governance. MUST NOT write construct scores. * * KV prefixes (only): * research:job: * research:idx * research:transitions: * research:decisions: / research:decisions:idx * research:claim: / research:claims:idx * research:brief:latest * research:branch: — idempotent seed lookup * orch:idempotency: * * Never rewrite job: / comments: / SRP docs as side effects of the parser. */ /** Frozen SRP snapshot guard — orchestration MUST NOT write construct scores. */ export const FROZEN_SRP_GUARD = Object.freeze({ L1: { id: 'L1-US-v0.1', score_authorized: true, note: 'Only score-authorized production construct; unchanged by orchestration.' }, M1: { score_authorized: false, status: 'UNKNOWN', note: 'Research-only unauthorized; score_authorized remains 0/false.' }, A2: { status: 'RETIRED', note: 'A2 retired — do not revive via orchestration.' }, 'M1-D': { status: 'RETIRED', note: 'M1-D retired — do not revive via orchestration.' }, 'M1-LONG': { status: 'RESEARCH', note: 'Michigan same-person pilot research only; no score authorization.' }, PRD: { status: 'BACKLOG', note: 'PRD research backlog — feasibility only.' }, Phase: { value: 'UNKNOWN', note: 'Phase/Alert/Load/Sync/PPS remain UNKNOWN.' }, Alert: 'UNKNOWN', Load: 'UNKNOWN', Sync: 'UNKNOWN', PPS: 'UNKNOWN', constraint: 'Orchestration MUST NOT write to construct scores. L1 frozen authorized; M1 UNKNOWN score_authorized=0; A2/M1-D retired; Phase/Alert/etc UNKNOWN.' }); export const STATES = Object.freeze([ 'QUEUED', 'RESEARCH', 'EVIDENCE_READY', 'REVIEW', 'REVISION', 'DECISION', 'IMPLEMENT', 'VALIDATE', 'CLOSED', 'BLOCKED' ]); export const SIGNIFICANCE = Object.freeze([ 'ROUTINE', 'INTERESTING', 'SIGNIFICANT', 'POTENTIALLY_NOVEL', 'REPLICATION_REQUIRED', 'REPLICATED', 'REJECTED' ]); /** * Allowed transitions (documented): * QUEUED→RESEARCH|BLOCKED|CLOSED * RESEARCH→EVIDENCE_READY|BLOCKED|REVISION|CLOSED * EVIDENCE_READY→REVIEW (auto) * REVIEW→REVISION|DECISION|BLOCKED * REVISION→RESEARCH|EVIDENCE_READY * DECISION→IMPLEMENT|CLOSED|QUEUED (spawn children) * IMPLEMENT→VALIDATE|BLOCKED * VALIDATE→CLOSED|REVISION * BLOCKED→QUEUED|RESEARCH (unblock) * Any→BLOCKED if blockers set */ export const ALLOWED_TRANSITIONS = Object.freeze({ QUEUED: Object.freeze(['RESEARCH', 'BLOCKED', 'CLOSED']), RESEARCH: Object.freeze(['EVIDENCE_READY', 'BLOCKED', 'REVISION', 'CLOSED']), EVIDENCE_READY: Object.freeze(['REVIEW', 'BLOCKED']), REVIEW: Object.freeze(['REVISION', 'DECISION', 'BLOCKED']), REVISION: Object.freeze(['RESEARCH', 'EVIDENCE_READY', 'BLOCKED']), DECISION: Object.freeze(['IMPLEMENT', 'CLOSED', 'QUEUED', 'BLOCKED']), IMPLEMENT: Object.freeze(['VALIDATE', 'BLOCKED']), VALIDATE: Object.freeze(['CLOSED', 'REVISION', 'BLOCKED']), BLOCKED: Object.freeze(['QUEUED', 'RESEARCH']), CLOSED: Object.freeze([]) }); const META_KEYS = [ 'STATUS', 'DECISION', 'NEXT_ACTION', 'ASSIGN_TO', 'REVIEW_BY', 'BLOCKERS', 'MODEL_CHANGE', 'EVIDENCE_IDS', 'CONFIDENCE', 'REQUIRES_REPLICATION', 'SIGNIFICANCE' ]; const RESEARCH_IDX = 'research:idx'; const DECISIONS_IDX = 'research:decisions:idx'; const CLAIMS_IDX = 'research:claims:idx'; const BRIEF_LATEST = 'research:brief:latest'; function normAgent(a) { return String(a || '') .trim() .toLowerCase(); } function parseList(v) { if (v == null || v === '') return []; if (Array.isArray(v)) return v.map((x) => String(x).trim()).filter(Boolean); return String(v) .split(/[,;|]/) .map((s) => s.trim()) .filter(Boolean); } function yesish(v) { if (v === true) return true; const s = String(v || '') .trim() .toUpperCase(); return s === 'YES' || s === 'TRUE' || s === '1' || s === 'Y'; } /** Parse leading KEY: VALUE lines; preserve human-readable body/analysis. */ export function parseStructuredHandoff(body, metaIn) { const meta = { ...(metaIn && typeof metaIn === 'object' ? metaIn : {}) }; let text = body == null ? '' : String(body); const lines = text.split(/\r?\n/); let i = 0; while (i < lines.length) { const line = lines[i]; const m = line.match(/^([A-Z_]+)\s*:\s*(.*)$/); if (!m) break; const key = m[1]; if (!META_KEYS.includes(key) && !(metaIn && key in metaIn)) { // Allow known keys only from leading block; stop on unknown to preserve prose if (!META_KEYS.includes(key)) break; } if (META_KEYS.includes(key)) { meta[key] = m[2].trim(); i++; continue; } break; } // Also accept lowercase keys from meta object for (const k of META_KEYS) { const low = k.toLowerCase(); if (meta[low] != null && meta[k] == null) meta[k] = meta[low]; } const remaining = lines.slice(i).join('\n').replace(/^\n+/, ''); return { meta, body: remaining.length ? remaining : text, analysis: remaining.length ? remaining : text }; } /** Material-update filter for James relay. */ export function isMaterialForJames(meta, body) { const m = meta && typeof meta === 'object' ? meta : {}; const text = String(body || ''); const status = String(m.STATUS || m.status || '') .trim() .toUpperCase(); const decision = String(m.DECISION || m.decision || '') .trim() .toUpperCase(); const combined = `${status} ${decision} ${text}`.toLowerCase(); // ACK / heartbeat / acknowledgement-only → not material const ackOnly = /^(ack|acked|acknowledged|acknowledgement|acknowledgment|heartbeat|hb|ping|pong|ok|收到|noted)\b/.test( combined.trim() ) || (/acknowledg?e?ment[_\s-]?only/.test(combined) && !/\b(evidence|objection|decision|blocker|replication|failed|implement)/.test(combined)); if (ackOnly && !m.EVIDENCE_IDS && !m.BLOCKERS && !decision) { // Pure ack with no structured signal if ( !status || status === 'ACK' || status === 'HEARTBEAT' || status === 'OK' || status === 'NOTED' ) { return false; } } // Explicit non-material statuses if (['ACK', 'HEARTBEAT', 'ACKNOWLEDGEMENT', 'ACKNOWLEDGMENT'].includes(status)) { if (!m.EVIDENCE_IDS && !m.BLOCKERS && !decision && !m.MODEL_CHANGE) return false; } // Material signals if (m.EVIDENCE_IDS || m.evidence_ids) return true; if (m.BLOCKERS || m.blockers) return true; if (m.MODEL_CHANGE || m.model_change) return true; if (yesish(m.REQUIRES_REPLICATION || m.requires_replication)) return true; if ( ['POTENTIALLY_NOVEL', 'REPLICATION_REQUIRED', 'REPLICATED', 'REJECTED', 'SIGNIFICANT'].includes( String(m.SIGNIFICANCE || m.significance || '').toUpperCase() ) ) { return true; } if ( decision && ['ACCEPT', 'CONDITIONAL_ACCEPT', 'REJECT', 'CHANGES_REQUESTED'].includes(decision) ) { return true; } if ( [ 'EVIDENCE_READY', 'REVIEW', 'REVISION', 'DECISION', 'IMPLEMENT', 'VALIDATE', 'CLOSED', 'BLOCKED', 'RESEARCH' ].includes(status) ) { return true; } // Body heuristics for material content if ( /\b(new evidence|objection|failed test|blocker|implementation complete|replication|unexpected relationship|model[- ]version change|decision)\b/i.test( text ) ) { return true; } // Default: if structured meta has STATUS/DECISION beyond ack, material; else if body looks like prose ack only if (!status && !decision) { const trimmed = text.trim(); if ( !trimmed || /^(ack|acked|acknowledged|heartbeat|ok|noted)[.!]?\s*$/i.test(trimmed) || /^ack(nowledg(e|ement|ment))?([_\s-]?only)?[.!]?\s*$/i.test(trimmed) ) { return false; } } return !!(status || decision || (text && text.length > 40)); } export function canTransition(from, to, job) { const f = String(from || '').toUpperCase(); const t = String(to || '').toUpperCase(); if (!STATES.includes(f) || !STATES.includes(t)) return false; if (t === 'BLOCKED') return true; // Any→BLOCKED if blockers set (caller may set) const allowed = ALLOWED_TRANSITIONS[f] || []; return allowed.includes(t); } export function pickIndependentReviewer(proposer) { const p = normAgent(proposer); if (p === 'claude') return 'astra'; // prefer ASTRA/CHATGPT if proposer is claude if (p === 'astra' || p === 'chatgpt') return 'claude'; if (p === 'grok') return 'claude'; return 'claude'; // default independent } export function highestValueUnblockedJob(jobs) { const list = Array.isArray(jobs) ? jobs : []; const byId = new Map(list.map((j) => [j.id, j])); const eligible = list.filter((j) => { if (!j) return false; const st = String(j.state || '').toUpperCase(); if (st === 'CLOSED' || st === 'BLOCKED') return false; const blockers = Array.isArray(j.blockers) ? j.blockers.filter(Boolean) : []; if (blockers.length) return false; const deps = Array.isArray(j.dependencies) ? j.dependencies : []; for (const d of deps) { const dep = byId.get(d); // Missing dependency is a gate failure (not treated as unblocked). if (!dep) return false; if (String(dep.state || '').toUpperCase() !== 'CLOSED') return false; } return true; }); eligible.sort((a, b) => { const ev = (Number(b.expected_value) || 0) - (Number(a.expected_value) || 0); if (ev !== 0) return ev; const pr = (Number(b.priority) || 0) - (Number(a.priority) || 0); if (pr !== 0) return pr; return String(a.created || '').localeCompare(String(b.created || '')); }); return eligible[0] || null; } export function assertSrpUntouched(srpState, buildSrpStateFn) { // Helper for tests: ensure frozen snapshot fields unchanged / orchestration didn't mutate if (srpState && typeof srpState === 'object') { if (srpState.constructs && srpState.constructs.L1) { if (srpState.constructs.L1.score_authorized !== true) { throw new Error('L1 score_authorized mutated'); } } if (srpState.constructs && srpState.constructs.M1) { if (srpState.constructs.M1.score_authorized !== false) { throw new Error('M1 score_authorized mutated'); } } } if (typeof buildSrpStateFn === 'function') { // optional injection — call without writing return true; } return true; } async function sha256Hex(str) { const data = new TextEncoder().encode(str); const hash = await crypto.subtle.digest('SHA-256', data); return [...new Uint8Array(hash)].map((b) => b.toString(16).padStart(2, '0')).join(''); } async function getJson(env, key, fallback) { const raw = await env.AI_HUB.get(key); if (!raw) return fallback; try { return JSON.parse(raw); } catch (_) { return fallback; } } async function putJson(env, key, val) { await env.AI_HUB.put(key, JSON.stringify(val)); } async function getResearchIds(env) { const arr = await getJson(env, RESEARCH_IDX, []); return Array.isArray(arr) ? arr : []; } async function putResearchIds(env, ids) { await putJson(env, RESEARCH_IDX, ids); } async function getResearchJob(env, id) { return getJson(env, `research:job:${id}`, null); } /** Dependencies must exist and be CLOSED before work transitions. */ export async function checkDependenciesSatisfied(env, job) { const deps = Array.isArray(job && job.dependencies) ? job.dependencies : []; if (!deps.length) return { ok: true, missing: [], open: [] }; const missing = []; const open = []; for (const d of deps) { const dep = await getResearchJob(env, String(d)); if (!dep) { missing.push(String(d)); continue; } if (String(dep.state || '').toUpperCase() !== 'CLOSED') { open.push({ id: String(d), state: dep.state }); } } if (missing.length || open.length) { return { ok: false, missing, open }; } return { ok: true, missing: [], open: [] }; } async function putResearchJob(env, job) { await putJson(env, `research:job:${job.id}`, job); } async function listResearchJobs(env) { const ids = await getResearchIds(env); const jobs = []; for (const id of ids) { const j = await getResearchJob(env, id); if (j) jobs.push(j); } return jobs; } async function getTransitions(env, jobId) { const arr = await getJson(env, `research:transitions:${jobId}`, []); return Array.isArray(arr) ? arr : []; } async function appendTransition(env, jobId, transition, deps) { const list = await getTransitions(env, jobId); list.push(transition); await putJson(env, `research:transitions:${jobId}`, list); return list; } function makeTransition(deps, { id, from, to, by, reason, evidence_ids, meta }) { return { id: id || (deps.newId ? deps.newId() : crypto.randomUUID()), from, to, at: deps.nowIso ? deps.nowIso() : new Date().toISOString(), by: by || 'system', reason: reason || '', evidence_ids: Array.isArray(evidence_ids) ? evidence_ids : [], meta: meta || {} }; } export async function createResearchJob(env, deps, input) { const ts = deps.nowIso(); const id = (input && input.id) || deps.newId(); const job = { id, title: String((input && input.title) || '').trim() || 'Untitled research', priority: Number(input && input.priority != null ? input.priority : 50), expected_value: Math.max( 0, Math.min(100, Number(input && input.expected_value != null ? input.expected_value : 50)) ), dependencies: Array.isArray(input && input.dependencies) ? input.dependencies.map(String) : [], blockers: Array.isArray(input && input.blockers) ? input.blockers.map(String) : parseList(input && input.blockers), assigned_agent: input && input.assigned_agent != null ? normAgent(input.assigned_agent) : null, reviewer: input && input.reviewer != null ? normAgent(input.reviewer) : null, proposer: input && input.proposer != null ? normAgent(input.proposer) : null, estimated_effort: input && input.estimated_effort != null ? input.estimated_effort : null, research_branch: input && input.research_branch != null ? String(input.research_branch) : null, model_version_affected: input && input.model_version_affected != null ? input.model_version_affected : null, state: String((input && input.state) || 'QUEUED').toUpperCase(), significance: String((input && input.significance) || 'ROUTINE').toUpperCase(), linked_hub_job_id: (input && input.linked_hub_job_id) || null, created: (input && input.created) || ts, updated: ts, result_summary: (input && input.result_summary) || null, decision_record_id: (input && input.decision_record_id) || null, pending_actions: Array.isArray(input && input.pending_actions) ? input.pending_actions : [], parent_job_id: (input && input.parent_job_id) || null, seed_branch_key: (input && input.seed_branch_key) || null }; if (!STATES.includes(job.state)) job.state = 'QUEUED'; if (!SIGNIFICANCE.includes(job.significance)) job.significance = 'ROUTINE'; // Self-reviewer not allowed at create — auto-correct to independent reviewer. const createActor = job.proposer || job.assigned_agent; if (job.reviewer && createActor && normAgent(job.reviewer) === normAgent(createActor)) { job.reviewer = pickIndependentReviewer(createActor); if (normAgent(job.reviewer) === normAgent(createActor)) { job.reviewer = normAgent(createActor) === 'claude' ? 'chatgpt' : 'claude'; } } if (!job.reviewer && createActor) { job.reviewer = pickIndependentReviewer(createActor); } // Dependency existence check at create (missing ids rejected). if (job.dependencies.length) { const depCheck = await checkDependenciesSatisfied(env, job); // At create, OPEN deps are allowed (job stays QUEUED) but MISSING deps are rejected. if (depCheck.missing && depCheck.missing.length) { return { error: 'missing_dependencies', status: 400, missing: depCheck.missing, message: 'All dependency ids must refer to existing research jobs' }; } } // Cannot create directly into RESEARCH/EVIDENCE_READY/REVIEW with unsatisfied deps. if (['RESEARCH', 'EVIDENCE_READY', 'REVIEW', 'DECISION', 'IMPLEMENT', 'VALIDATE'].includes(job.state)) { const depCheck2 = await checkDependenciesSatisfied(env, job); if (!depCheck2.ok) { return { error: 'dependencies_unsatisfied', status: 400, missing: depCheck2.missing, open: depCheck2.open }; } } const ids = await getResearchIds(env); if (!ids.includes(job.id)) { ids.unshift(job.id); await putResearchIds(env, ids); } await putResearchJob(env, job); if (job.research_branch || job.seed_branch_key) { const bk = job.seed_branch_key || job.research_branch; await putJson(env, `research:branch:${bk}`, { id: job.id, at: ts }); } // Initial transition record (append-only history starts here) await appendTransition( env, job.id, makeTransition(deps, { from: null, to: job.state, by: (input && input.by) || 'system', reason: (input && input.reason) || 'created', evidence_ids: [], meta: { created: true } }) ); return job; } export async function transitionResearchJob(env, deps, jobId, { to, by, reason, evidence_ids, meta }) { const job = await getResearchJob(env, jobId); if (!job) return { error: 'not_found', status: 404 }; const target = String(to || '').toUpperCase(); const from = String(job.state || '').toUpperCase(); // Auto-block if blockers set and going anywhere except already BLOCKED handling const blockers = Array.isArray(job.blockers) ? job.blockers.filter(Boolean) : []; let finalTo = target; if (blockers.length && target !== 'BLOCKED' && from !== 'BLOCKED') { // Any→BLOCKED if blockers set (when transitioning) if (meta && meta.force_despite_blockers) { // allow explicit override } else if (target !== 'CLOSED') { // Prefer BLOCKED when blockers present unless unblocking explicitly finalTo = 'BLOCKED'; } } if (!canTransition(from, finalTo, job) && !(finalTo === 'BLOCKED')) { return { error: 'invalid_transition', status: 400, from, to: finalTo, allowed: ALLOWED_TRANSITIONS[from] || [] }; } if (finalTo === 'BLOCKED' && from === 'CLOSED') { return { error: 'invalid_transition', status: 400, from, to: finalTo }; } // Dependency gate: must exist + CLOSED before leaving QUEUED into active work. const workStates = ['RESEARCH', 'EVIDENCE_READY', 'REVIEW', 'REVISION', 'DECISION', 'IMPLEMENT', 'VALIDATE']; if (workStates.includes(finalTo)) { const depCheck = await checkDependenciesSatisfied(env, job); if (!depCheck.ok) { return { error: 'dependencies_unsatisfied', status: 400, from, to: finalTo, missing: depCheck.missing, open: depCheck.open }; } } const tr = makeTransition(deps, { from, to: finalTo, by: by || 'system', reason: reason || '', evidence_ids: evidence_ids || [], meta: meta || {} }); await appendTransition(env, jobId, tr, deps); job.state = finalTo; job.updated = deps.nowIso(); if (meta && meta.assigned_agent) job.assigned_agent = normAgent(meta.assigned_agent); if (meta && meta.reviewer) job.reviewer = normAgent(meta.reviewer); if (meta && meta.significance) job.significance = String(meta.significance).toUpperCase(); if (meta && meta.result_summary) job.result_summary = meta.result_summary; if (meta && meta.blockers != null) { job.blockers = Array.isArray(meta.blockers) ? meta.blockers : parseList(meta.blockers); } await putResearchJob(env, job); return { ok: true, job, transition: tr }; } async function ensureReviewerAssignment(env, deps, job) { const proposer = job.proposer || job.assigned_agent; let reviewer = job.reviewer; if (!reviewer || normAgent(reviewer) === normAgent(proposer)) { reviewer = pickIndependentReviewer(proposer); // never assign proposer as sole reviewer if (normAgent(reviewer) === normAgent(proposer)) { reviewer = normAgent(proposer) === 'claude' ? 'chatgpt' : 'claude'; } } job.reviewer = reviewer; const action = { type: 'review_task', reviewer, job_id: job.id, at: deps.nowIso(), note: `Independent review requested for research job ${job.id}` }; job.pending_actions = Array.isArray(job.pending_actions) ? job.pending_actions : []; const already = job.pending_actions.some( (a) => a && a.type === 'review_task' && a.reviewer === reviewer ); if (!already) job.pending_actions.push(action); await putResearchJob(env, job); if (deps.putInbox) { try { await deps.putInbox(env, reviewer, { id: deps.newId(), from: 'orchestration', to: reviewer, body: `REVIEW_TASK research:${job.id} — ${job.title}\nPlease review evidence and return STATUS/DECISION handoff.`, at: deps.nowIso(), ts: Date.now(), research_job_id: job.id }); } catch (_) { /* optional */ } } return job; } async function appendDecision(env, deps, decision) { const id = decision.id || deps.newId(); const rec = { ...decision, id, at: decision.at || deps.nowIso() }; await putJson(env, `research:decisions:${id}`, rec); const idx = await getJson(env, DECISIONS_IDX, []); const arr = Array.isArray(idx) ? idx : []; if (!arr.includes(id)) { arr.unshift(id); await putJson(env, DECISIONS_IDX, arr); } return rec; } async function listDecisions(env) { const idx = await getJson(env, DECISIONS_IDX, []); const out = []; for (const id of Array.isArray(idx) ? idx : []) { const d = await getJson(env, `research:decisions:${id}`, null); if (d) out.push(d); } return out; } async function createChildJobsFromNext(env, deps, parent, nextJobs, meta) { const created = []; const specs = Array.isArray(nextJobs) ? nextJobs : []; if (!specs.length && meta && meta.NEXT_ACTION) { specs.push({ title: String(meta.NEXT_ACTION).slice(0, 200), assigned_agent: meta.ASSIGN_TO ? normAgent(meta.ASSIGN_TO) : parent.assigned_agent, state: 'QUEUED', priority: parent.priority, expected_value: Math.max(0, (Number(parent.expected_value) || 50) - 5), proposer: parent.proposer || parent.assigned_agent, parent_job_id: parent.id, dependencies: [] }); } for (const spec of specs) { const child = await createResearchJob(env, deps, { ...spec, parent_job_id: parent.id, linked_hub_job_id: spec.linked_hub_job_id || parent.linked_hub_job_id, by: 'orchestration', reason: `spawned_from_decision:${parent.id}` }); created.push(child); } return created; } async function maybeSpawnReplication(env, deps, job, meta, by) { const sig = String(meta.SIGNIFICANCE || job.significance || '').toUpperCase(); const reqRep = yesish(meta.REQUIRES_REPLICATION); if (sig !== 'POTENTIALLY_NOVEL' && !reqRep) return null; // Proposer cannot self-certify REPLICATED const proposer = job.proposer || job.assigned_agent; const independent = pickIndependentReviewer(proposer); const agent = normAgent(independent) === normAgent(proposer) ? 'claude' : independent; if (normAgent(by) === normAgent(proposer) && String(meta.SIGNIFICANCE || '').toUpperCase() === 'REPLICATED') { return { error: 'proposer_cannot_self_certify_replicated', status: 403 }; } const branchKey = `replication:${job.id}`; const existing = await getJson(env, `research:branch:${branchKey}`, null); if (existing && existing.id) { const ej = await getResearchJob(env, existing.id); if (ej) return { job: ej, idempotent: true }; } const rep = await createResearchJob(env, deps, { title: `REPLICATION: ${job.title}`, assigned_agent: agent, reviewer: pickIndependentReviewer(agent), proposer: agent, state: 'QUEUED', significance: 'REPLICATION_REQUIRED', expected_value: Number(job.expected_value) || 70, priority: (Number(job.priority) || 50) + 5, parent_job_id: job.id, linked_hub_job_id: job.linked_hub_job_id, research_branch: branchKey, seed_branch_key: branchKey, by: 'orchestration', reason: `auto_replication_for:${job.id}` }); return { job: rep }; } /** * Apply structured handoff routing rules. */ export async function applyHandoff(env, deps, jobId, { agent, body, meta: metaIn, next_jobs }) { const job = await getResearchJob(env, jobId); if (!job) return { error: 'not_found', status: 404 }; const parsed = parseStructuredHandoff(body, metaIn); const meta = parsed.meta; const by = normAgent(agent); // Idempotency const evidenceSorted = parseList(meta.EVIDENCE_IDS || meta.evidence_ids) .slice() .sort() .join(','); const idemSrc = [ jobId, by, String(meta.STATUS || ''), String(meta.DECISION || ''), String(meta.NEXT_ACTION || ''), evidenceSorted ].join('|'); const hash = await sha256Hex(idemSrc); const idemKey = `orch:idempotency:${hash}`; const prior = await getJson(env, idemKey, null); if (prior) { return { ...prior, idempotent: true }; } // Update job fields from meta if (meta.ASSIGN_TO) job.assigned_agent = normAgent(meta.ASSIGN_TO); if (meta.REVIEW_BY) job.reviewer = normAgent(meta.REVIEW_BY); if (meta.BLOCKERS != null) job.blockers = parseList(meta.BLOCKERS); if (meta.SIGNIFICANCE) job.significance = String(meta.SIGNIFICANCE).toUpperCase(); if (meta.MODEL_CHANGE) job.model_version_affected = meta.MODEL_CHANGE; if (parsed.body) job.result_summary = String(parsed.body).slice(0, 2000); job.updated = deps.nowIso(); // Proposer cannot self-certify REPLICATED if ( String(meta.SIGNIFICANCE || '').toUpperCase() === 'REPLICATED' && by && job.proposer && by === normAgent(job.proposer) ) { return { error: 'proposer_cannot_self_certify_replicated', status: 403 }; } let status = String(meta.STATUS || '').toUpperCase(); const decision = String(meta.DECISION || '').toUpperCase(); // Map DECISION to status hints if (!status && decision === 'CHANGES_REQUESTED') status = 'REVISION'; if (!status && ['ACCEPT', 'CONDITIONAL_ACCEPT', 'REJECT'].includes(decision)) status = 'DECISION'; const sideEffects = { spawned_jobs: [], replication_job: null, decision_record: null, transitions: [] }; // Auto-block if blockers present if (job.blockers && job.blockers.length && job.state !== 'BLOCKED') { const r = await transitionResearchJob(env, deps, jobId, { to: 'BLOCKED', by, reason: 'blockers_set', evidence_ids: parseList(meta.EVIDENCE_IDS), meta: { blockers: job.blockers, handoff: true } }); if (r.transition) sideEffects.transitions.push(r.transition); Object.assign(job, r.job || job); } // Routing: EVIDENCE_READY if (status === 'EVIDENCE_READY' || job.state === 'EVIDENCE_READY') { const evidenceIdsStrict = parseList(meta.EVIDENCE_IDS || meta.evidence_ids); if (status === 'EVIDENCE_READY' && evidenceIdsStrict.length === 0) { return { error: 'evidence_ids_required', status: 400, message: 'EVIDENCE_READY requires non-empty EVIDENCE_IDS before REVIEW routing' }; } if (status === 'EVIDENCE_READY' && job.state !== 'EVIDENCE_READY' && job.state !== 'REVIEW') { const r = await transitionResearchJob(env, deps, jobId, { to: 'EVIDENCE_READY', by, reason: meta.NEXT_ACTION || 'evidence_ready', evidence_ids: evidenceIdsStrict, meta: { handoff: true } }); if (r.error) return r; if (r.transition) sideEffects.transitions.push(r.transition); Object.assign(job, r.job); } // Auto assign reviewer + REVIEW await ensureReviewerAssignment(env, deps, job); if (job.state === 'EVIDENCE_READY') { const r2 = await transitionResearchJob(env, deps, jobId, { to: 'REVIEW', by: 'orchestration', reason: 'auto_route_to_review', evidence_ids: parseList(meta.EVIDENCE_IDS), meta: { reviewer: job.reviewer, handoff: true } }); if (r2.transition) sideEffects.transitions.push(r2.transition); Object.assign(job, r2.job || job); } } // REVIEW → REVISION if ( status === 'REVISION' || decision === 'CHANGES_REQUESTED' || (status === 'REVIEW' && decision === 'CHANGES_REQUESTED') ) { if (job.state === 'REVIEW' || job.state === 'REVISION' || job.state === 'EVIDENCE_READY') { if (job.state !== 'REVISION') { const r = await transitionResearchJob(env, deps, jobId, { to: 'REVISION', by, reason: meta.NEXT_ACTION || 'changes_requested', evidence_ids: parseList(meta.EVIDENCE_IDS), meta: { handoff: true } }); if (r.error && r.error !== 'invalid_transition') return r; if (r.transition) sideEffects.transitions.push(r.transition); if (r.job) Object.assign(job, r.job); } } // REVISION → RESEARCH, ASSIGN_TO=proposer const proposer = job.proposer || job.assigned_agent; job.assigned_agent = normAgent(meta.ASSIGN_TO || proposer); await putResearchJob(env, job); if (job.state === 'REVISION') { const r2 = await transitionResearchJob(env, deps, jobId, { to: 'RESEARCH', by: 'orchestration', reason: 'revision_routes_to_researcher', evidence_ids: parseList(meta.EVIDENCE_IDS), meta: { assigned_agent: job.assigned_agent, handoff: true } }); if (r2.transition) sideEffects.transitions.push(r2.transition); if (r2.job) Object.assign(job, r2.job); } } // Final decision path if (['ACCEPT', 'CONDITIONAL_ACCEPT', 'REJECT'].includes(decision) || status === 'DECISION') { const proposer = job.proposer || job.assigned_agent; const reviewer = job.reviewer; // Proposer cannot be sole final reviewer // Proposer cannot be sole final reviewer if (by && proposer && by === normAgent(proposer)) { const separateAck = !!( meta.REVIEWER_ACK || meta.reviewer_ack || meta.separate_reviewer_ack ); const reviewerIsOther = reviewer && normAgent(reviewer) !== by; if (!separateAck || !reviewerIsOther) { return { error: 'proposer_cannot_self_final_review', status: 403 }; } } if (job.state === 'REVIEW' || job.state === 'DECISION' || job.state === 'EVIDENCE_READY') { if (job.state !== 'DECISION') { const r = await transitionResearchJob(env, deps, jobId, { to: 'DECISION', by, reason: decision || 'decision', evidence_ids: parseList(meta.EVIDENCE_IDS), meta: { handoff: true, decision } }); if (r.error && r.error !== 'invalid_transition') return r; if (r.transition) sideEffects.transitions.push(r.transition); if (r.job) Object.assign(job, r.job); } } const decisionRec = await appendDecision(env, deps, { research_job_id: job.id, decision: decision || 'ACCEPT', evidence_ids: parseList(meta.EVIDENCE_IDS), proposer: proposer, reviewer: reviewer || by, rationale: parsed.body || meta.NEXT_ACTION || reasonFrom(meta), superseded_methodology: meta.superseded_methodology || meta.SUPERSEDED_METHODOLOGY || null, rollback_target: meta.rollback_target || meta.ROLLBACK_TARGET || null, concept_version: meta.concept_version || meta.CONCEPT_VERSION || null, measurement_version: meta.measurement_version || meta.MEASUREMENT_VERSION || null, by, superseded_by: null, parent_decision_id: meta.parent_decision_id || null }); sideEffects.decision_record = decisionRec; job.decision_record_id = decisionRec.id; await putResearchJob(env, job); if (decision === 'ACCEPT' || decision === 'CONDITIONAL_ACCEPT') { const kids = await createChildJobsFromNext(env, deps, job, next_jobs, meta); sideEffects.spawned_jobs = kids; } if (decision === 'REJECT') { const r = await transitionResearchJob(env, deps, jobId, { to: 'CLOSED', by, reason: 'rejected', evidence_ids: parseList(meta.EVIDENCE_IDS), meta: { handoff: true } }); if (r.transition) sideEffects.transitions.push(r.transition); if (r.job) Object.assign(job, r.job); } } // Generic status transition if provided and not already handled if ( status && STATES.includes(status) && status !== job.state && !['EVIDENCE_READY', 'REVISION', 'DECISION'].includes(status) ) { const r = await transitionResearchJob(env, deps, jobId, { to: status, by, reason: meta.NEXT_ACTION || 'handoff_status', evidence_ids: parseList(meta.EVIDENCE_IDS), meta: { handoff: true } }); if (!r.error) { if (r.transition) sideEffects.transitions.push(r.transition); if (r.job) Object.assign(job, r.job); } } // Significance / replication const rep = await maybeSpawnReplication(env, deps, job, meta, by); if (rep && rep.error) return rep; if (rep && rep.job) sideEffects.replication_job = rep.job; const fresh = await getResearchJob(env, jobId); const result = { ok: true, job: fresh, meta, body: parsed.body, material: isMaterialForJames(meta, parsed.body), ...sideEffects, idempotent: false }; await putJson(env, idemKey, result); // James relay on material only if (result.material && deps.relayJames) { try { await deps.relayJames(env, { from: by, body: `research:${jobId} handoff\nSTATUS:${meta.STATUS || ''}\nDECISION:${meta.DECISION || ''}\n${parsed.body || ''}`, research_job_id: jobId }); } catch (_) { /* optional */ } } return result; } function reasonFrom(meta) { return String(meta.NEXT_ACTION || meta.rationale || ''); } /** Seed parallel SRP research branches — idempotent by branch key. */ export async function seedSrpResearchQueue(env, deps) { const seeds = [ { seed_branch_key: 'M1-LONG', research_branch: 'M1-LONG', title: 'M1-LONG Michigan same-person pilot', state: 'RESEARCH', assigned_agent: 'grok', reviewer: 'claude', proposer: 'grok', linked_hub_job_id: 'b73c292a-16e2-4f11-9370-1337d0239791', expected_value: 92, priority: 90, significance: 'SIGNIFICANT', estimated_effort: 'high', model_version_affected: null }, { seed_branch_key: 'S1', research_branch: 'S1', title: 'S1 Affective Tribalization measurement research', state: 'QUEUED', assigned_agent: 'grok', reviewer: 'astra', proposer: 'grok', expected_value: 80, priority: 80, significance: 'INTERESTING', dependencies: [] }, { seed_branch_key: 'E2', research_branch: 'E2', title: 'E2 Elite Fragmentation narrow-sensor research', state: 'QUEUED', assigned_agent: 'grok', reviewer: 'astra', proposer: 'grok', expected_value: 78, priority: 78, significance: 'INTERESTING' }, { seed_branch_key: 'ACLED-B1-B2', research_branch: 'ACLED-B1-B2', title: 'ACLED B1/B2 connector/data preparation', state: 'QUEUED', assigned_agent: 'grok', reviewer: 'claude', proposer: 'grok', expected_value: 55, priority: 40, significance: 'ROUTINE', blockers: [], result_summary: 'May wait on sensors' }, { seed_branch_key: 'PRD', research_branch: 'PRD', title: 'PRD research backlog/source feasibility', state: 'QUEUED', assigned_agent: 'grok', reviewer: 'claude', proposer: 'grok', expected_value: 35, priority: 20, significance: 'ROUTINE', result_summary: 'backlog' }, { seed_branch_key: 'HIST-BACKTEST', research_branch: 'HIST-BACKTEST', title: 'Historical validation/backtest infrastructure', state: 'QUEUED', assigned_agent: 'grok', reviewer: 'claude', proposer: 'grok', expected_value: 60, priority: 50, significance: 'INTERESTING' } ]; const created = []; const existing = []; for (const spec of seeds) { const prev = await getJson(env, `research:branch:${spec.seed_branch_key}`, null); if (prev && prev.id) { const j = await getResearchJob(env, prev.id); if (j) { existing.push(j); continue; } } const job = await createResearchJob(env, deps, { ...spec, by: 'seed', reason: 'seedSrpResearchQueue' }); created.push(job); } return { ok: true, created, existing, count_created: created.length, count_existing: existing.length }; } async function generateBrief(env, deps) { const jobs = await listResearchJobs(env); const highest = highestValueUnblockedJob(jobs); const decisions = await listDecisions(env); const recentDecisions = decisions.slice(0, 10); const active = jobs.filter((j) => !['CLOSED'].includes(String(j.state || '').toUpperCase())); const blocked = jobs.filter((j) => String(j.state || '').toUpperCase() === 'BLOCKED'); const signals = []; for (const j of active) { if (j.significance === 'POTENTIALLY_NOVEL' || j.significance === 'REPLICATION_REQUIRED') { signals.push({ type: 'significance', job_id: j.id, significance: j.significance, title: j.title }); } if (j.blockers && j.blockers.length) { signals.push({ type: 'blocker', job_id: j.id, blockers: j.blockers, title: j.title }); } } for (const d of recentDecisions.slice(0, 5)) { signals.push({ type: 'decision', decision_id: d.id, decision: d.decision, research_job_id: d.research_job_id, rationale: (d.rationale || '').slice(0, 200) }); } const brief = { generated_at: deps.nowIso(), principle: 'SIGNAL not volume', highest_value_unblocked_job: highest, active_count: active.length, blocked_count: blocked.length, signals, frozen_srp_guard: FROZEN_SRP_GUARD.constraint, jobs_summary: active.map((j) => ({ id: j.id, title: j.title, state: j.state, expected_value: j.expected_value, priority: j.priority, assigned_agent: j.assigned_agent, significance: j.significance, blockers: j.blockers })) }; await putJson(env, BRIEF_LATEST, brief); return brief; } async function upsertClaim(env, deps, input) { const id = (input && input.id) || deps.newId(); let claim = await getJson(env, `research:claim:${id}`, null); if (!claim) { claim = { id, claim: (input && input.claim) || '', positions: [], status: 'UNRESOLVED', created: deps.nowIso(), updated: deps.nowIso(), spawned_research_job_id: null }; } if (input.claim) claim.claim = input.claim; if (input.position) { claim.positions = Array.isArray(claim.positions) ? claim.positions : []; claim.positions.push({ agent: normAgent(input.agent || input.position.agent), position: input.position.position || input.position, confidence: input.position.confidence != null ? input.position.confidence : input.confidence, evidence: input.position.evidence || input.evidence || [], objections: input.position.objections || input.objections || [] }); } if (input.status) claim.status = String(input.status).toUpperCase(); // Contested detection const agents = new Set(claim.positions.map((p) => p.agent)); const texts = claim.positions.map((p) => String(p.position || '').toLowerCase()); const hasSupport = texts.some((t) => /support|agree|true|yes/.test(t)); const hasContest = texts.some((t) => /contest|disagree|object|false|no/.test(t)); if (agents.size >= 2 && hasSupport && hasContest) { claim.status = claim.status === 'RESOLVED' ? claim.status : 'CONTESTED'; } // Spawn research job for contested with resolvable empirics if ( claim.status === 'CONTESTED' && input.spawn_research !== false && (input.resolvable_empirics || input.spawn_research) ) { if (!claim.spawned_research_job_id) { const rj = await createResearchJob(env, deps, { title: `Claim resolution: ${String(claim.claim).slice(0, 120)}`, state: 'QUEUED', assigned_agent: 'grok', reviewer: 'claude', proposer: 'orchestration', expected_value: 65, priority: 60, significance: 'INTERESTING', by: 'disagreement_engine', reason: `contested_claim:${id}` }); claim.spawned_research_job_id = rj.id; } } claim.updated = deps.nowIso(); await putJson(env, `research:claim:${id}`, claim); const idx = await getJson(env, CLAIMS_IDX, []); const arr = Array.isArray(idx) ? idx : []; if (!arr.includes(id)) { arr.unshift(id); await putJson(env, CLAIMS_IDX, arr); } return claim; } async function listClaims(env) { const idx = await getJson(env, CLAIMS_IDX, []); const out = []; for (const id of Array.isArray(idx) ? idx : []) { const c = await getJson(env, `research:claim:${id}`, null); if (c) out.push(c); } return out; } export function researchSchema() { return { states: STATES, allowed_transitions: ALLOWED_TRANSITIONS, significance: SIGNIFICANCE, job_fields: [ 'id', 'title', 'priority', 'expected_value', 'dependencies', 'blockers', 'assigned_agent', 'reviewer', 'proposer', 'estimated_effort', 'research_branch', 'model_version_affected', 'state', 'significance', 'linked_hub_job_id', 'created', 'updated', 'result_summary', 'decision_record_id' ], handoff_meta_keys: META_KEYS, kv_prefixes: ['research:', 'orch:'], frozen_srp_guard: FROZEN_SRP_GUARD, routes: [ 'GET /research/queue', 'GET /research/queue/highest', 'GET /research/jobs/:id', 'GET /research/jobs/:id/transitions', 'POST /research/jobs', 'POST /research/jobs/:id/transition', 'POST /research/jobs/:id/handoff', 'GET /research/brief', 'POST /research/brief/generate', 'GET/POST /research/claims', 'GET /research/decisions', 'POST /research/seed', 'GET /research/schema' ] }; } /** * HTTP router — returns Response|null if path doesn't match /research. */ export async function handleResearchRoutes(request, env, deps) { const url = new URL(request.url); let path = url.pathname.replace(/\/+$/, '') || '/'; if (!path.startsWith('/research')) return null; const method = request.method; // GET /research/schema if (path === '/research/schema' && method === 'GET') { return deps.json(researchSchema()); } // GET /research/queue if (path === '/research/queue' && method === 'GET') { const jobs = await listResearchJobs(env); const highest = highestValueUnblockedJob(jobs); return deps.json({ jobs, count: jobs.length, highest_value_unblocked_job: highest }); } // GET /research/queue/highest if (path === '/research/queue/highest' && method === 'GET') { const jobs = await listResearchJobs(env); const highest = highestValueUnblockedJob(jobs); return deps.json({ highest_value_unblocked_job: highest }); } // GET /research/brief if (path === '/research/brief' && method === 'GET') { const brief = await getJson(env, BRIEF_LATEST, null); return deps.json({ brief }); } // POST /research/brief/generate if (path === '/research/brief/generate' && method === 'POST') { const denied = deps.requireAuth(request, env); if (denied) return denied; const brief = await generateBrief(env, deps); return deps.json({ ok: true, brief }); } // GET /research/decisions if (path === '/research/decisions' && method === 'GET') { const decisions = await listDecisions(env); return deps.json({ decisions, count: decisions.length }); } // GET/POST /research/claims if (path === '/research/claims') { if (method === 'GET') { const claims = await listClaims(env); return deps.json({ claims, count: claims.length }); } if (method === 'POST') { const denied = deps.requireAuth(request, env); if (denied) return denied; const { body, error } = await deps.readJson(request); if (error) return deps.json({ error }, 400); const claim = await upsertClaim(env, deps, body || {}); return deps.json({ ok: true, claim }, 201); } } // POST /research/seed if (path === '/research/seed' && method === 'POST') { const denied = deps.requireAuth(request, env); if (denied) return denied; const result = await seedSrpResearchQueue(env, deps); return deps.json(result, 201); } // POST /research/jobs if (path === '/research/jobs' && method === 'POST') { const denied = deps.requireAuth(request, env); if (denied) return denied; const { body, error } = await deps.readJson(request); if (error) return deps.json({ error }, 400); if (!body || !body.title) return deps.json({ error: 'title_required' }, 400); const job = await createResearchJob(env, deps, body); if (job && job.error) return deps.json(job, job.status || 400); return deps.json({ ok: true, job }, 201); } // /research/jobs/:id[/transition|/handoff|/transitions] const jobMatch = path.match(/^\/research\/jobs\/([^/]+)(?:\/(transition|handoff|transitions))?$/); if (jobMatch) { const id = decodeURIComponent(jobMatch[1]); const action = jobMatch[2] || null; if (!action && method === 'GET') { const job = await getResearchJob(env, id); if (!job) return deps.json({ error: 'not_found' }, 404); return deps.json({ job }); } if (action === 'transitions' && method === 'GET') { const job = await getResearchJob(env, id); if (!job) return deps.json({ error: 'not_found' }, 404); const transitions = await getTransitions(env, id); return deps.json({ job_id: id, transitions, count: transitions.length }); } if (action === 'transition' && method === 'POST') { const denied = deps.requireAuth(request, env); if (denied) return denied; const { body, error } = await deps.readJson(request); if (error) return deps.json({ error }, 400); if (!body || !body.to) return deps.json({ error: 'to_required' }, 400); const result = await transitionResearchJob(env, deps, id, { to: body.to, by: body.by, reason: body.reason, evidence_ids: body.evidence_ids, meta: body.meta }); if (result.error) return deps.json(result, result.status || 400); return deps.json(result); } if (action === 'handoff' && method === 'POST') { const denied = deps.requireAuth(request, env); if (denied) return denied; const { body, error } = await deps.readJson(request); if (error) return deps.json({ error }, 400); if (!body || !body.agent) return deps.json({ error: 'agent_required' }, 400); const result = await applyHandoff(env, deps, id, { agent: body.agent, body: body.body != null ? body.body : body.analysis || '', meta: body.meta, next_jobs: body.next_jobs }); if (result.error) return deps.json(result, result.status || 400); return deps.json(result); } } // Path matched /research but no route — return 404 JSON (still "handled") if (path.startsWith('/research')) { return deps.json({ error: 'not_found', path }, 404); } return null; }