← back to The Ai Factory
src/pipeline.js
180 lines
// The AI Factory — pipeline runner.
// 11 stages mirroring Site Factory.
// Parallelization (2026-04-30): research ‖ scaffold, document ‖ critic.
const { runIntake } = require('./stages/intake');
const { runResearch } = require('./stages/research');
const { runScaffold } = require('./stages/scaffold');
const { runDocument } = require('./stages/document');
const { runSmokeTest } = require('./stages/smoke_test');
const { runCritic } = require('./stages/critic');
const { runIterate } = require('./stages/iterate');
const STAGES = [
'intake', 'research', 'scaffold', 'wire', 'smoke_test',
'critic', 'iterate', 'document', 'memory_write', 'activate', 'watchdog_hook'
];
async function logEvent(pg, runId, stage, stageName, level, message, payload) {
await pg.query(
'INSERT INTO events (run_id, stage, stage_name, level, message, payload) VALUES ($1,$2,$3,$4,$5,$6)',
[runId, stage, stageName, level, message, payload ? JSON.stringify(payload) : null]
);
}
async function setStage(pg, runId, stage, status) {
await pg.query('UPDATE runs SET current_stage=$1, status=$2, updated_at=now() WHERE id=$3', [stage, status, runId]);
}
async function finishRun(pg, runId, status) {
await pg.query('UPDATE runs SET status=$1, finished_at=now(), updated_at=now() WHERE id=$2', [status, runId]);
}
async function setArtifactMeta(pg, runId, type, name) {
await pg.query('UPDATE runs SET artifact_type=$1, artifact_name=$2, updated_at=now() WHERE id=$3', [type, name, runId]);
}
async function runPipeline(pg, run) {
const runId = run.id;
try {
// === Stage 1: intake ===============================================
await setStage(pg, runId, 1, 'running');
await logEvent(pg, runId, 1, 'intake', 'info', 'starting', null);
const { spec, errors: intakeErrors } = await runIntake({ prompt: run.prompt });
if (intakeErrors.length) {
await logEvent(pg, runId, 1, 'intake', 'error', 'spec validation failed', { spec, errors: intakeErrors });
await finishRun(pg, runId, 'failed_intake');
return;
}
await setArtifactMeta(pg, runId, spec.artifact_type, spec.artifact_name);
await logEvent(pg, runId, 1, 'intake', 'info', 'spec extracted', { spec });
// === Stages 2 + 3 in parallel: research ‖ scaffold =================
await setStage(pg, runId, 3, 'running'); // surface the longer-running one
const [researchResult, scaffold] = await Promise.all([
runResearch({ spec }).catch(e => ({ error: e.message })),
runScaffold({ runId, spec })
]);
// log research outcome
if (researchResult.error) {
await logEvent(pg, runId, 2, 'research', 'warn', 'research failed (non-blocking)', { error: researchResult.error });
} else {
const conflictMsg = researchResult.exact_conflict
? `EXACT NAME CONFLICT — ~/.claude/${spec.artifact_type === 'subagent' ? 'agents' : 'skills'}/${spec.artifact_name} already exists`
: researchResult.matches.length
? `${researchResult.matches.length} similar artifact(s) found`
: 'no conflicts';
await logEvent(pg, runId, 2, 'research',
researchResult.exact_conflict ? 'warn' : 'info',
conflictMsg,
researchResult);
}
// log scaffold outcome
if (scaffold.errors.length) {
await logEvent(pg, runId, 3, 'scaffold', 'warn', 'validation issues', { errors: scaffold.errors, outPath: scaffold.outPath });
} else {
await logEvent(pg, runId, 3, 'scaffold', 'info', 'wrote file', { outPath: scaffold.outPath, bytes: scaffold.content.length });
}
await pg.query(
'INSERT INTO artifacts (run_id, artifact_type, artifact_name, location, status) VALUES ($1,$2,$3,$4,$5)',
[runId, spec.artifact_type, spec.artifact_name, scaffold.outPath, 'draft']
);
// === Stage 4: wire (no-op for v1 subagent/skill artifacts) =========
await setStage(pg, runId, 4, 'running');
await logEvent(pg, runId, 4, 'wire', 'info', 'no-op for v1 artifacts (no external registration needed)', null);
// === Stage 5: smoke_test (cheap, local) ============================
await setStage(pg, runId, 5, 'running');
const smoke = await runSmokeTest({ outPath: scaffold.outPath, spec });
await logEvent(pg, runId, 5, 'smoke_test',
smoke.passing ? 'info' : 'warn',
`${smoke.passed}/${smoke.total} checks passing`,
smoke);
// === Stages 6 + 8 in parallel: critic ‖ document ===================
// Document only runs for skills — for subagents we just skip and Promise.all
// gets back an immediate { skipped: true }.
await setStage(pg, runId, 6, 'running');
await logEvent(pg, runId, 6, 'critic', 'info', 'starting (claude CLI subprocess) ‖ document', null);
const [criticResult, documentResult] = await Promise.all([
runCritic({ artifactType: spec.artifact_type, content: scaffold.content, spec })
.catch(e => ({ approved: false, score: 0, issues: [`critic failed: ${e.message}`], fix_instructions: [] })),
runDocument({ runId, spec, scaffoldOutPath: scaffold.outPath })
.catch(e => ({ skipped: true, error: e.message }))
]);
let review = criticResult;
await logEvent(pg, runId, 6, 'critic',
review.approved ? 'info' : 'warn',
`score=${review.score} approved=${review.approved}`,
{ review });
await logEvent(pg, runId, 8, 'document',
documentResult.skipped ? 'info' : 'info',
documentResult.skipped ? `skipped: ${documentResult.reason || documentResult.error || 'n/a'}` : 'wrote README.md',
documentResult);
// === Stage 7: iterate (only if rejected, with re-critic) ===========
await setStage(pg, runId, 7, 'running');
let finalReview = review;
let didIterate = false;
if (!review.approved && review.fix_instructions.length) {
await logEvent(pg, runId, 7, 'iterate', 'info', `applying ${review.fix_instructions.length} fixes`, { fixes: review.fix_instructions });
try {
const iter = await runIterate({ outPath: scaffold.outPath, content: scaffold.content, spec, review });
scaffold.content = iter.content;
didIterate = true;
await logEvent(pg, runId, 7, 'iterate', 'info', 'rewrote file', { bytes: iter.content.length });
await logEvent(pg, runId, 7, 'iterate', 'info', 're-running critic on v2', null);
try {
const review2 = await runCritic({ artifactType: spec.artifact_type, content: iter.content, spec });
finalReview = review2;
await logEvent(pg, runId, 7, 'iterate', review2.approved ? 'info' : 'warn',
`re-critic: score=${review2.score} approved=${review2.approved} (was ${review.score})`,
{ review: review2 });
} catch (err) {
await logEvent(pg, runId, 7, 'iterate', 'error', 're-critic failed', { error: err.message });
}
} catch (err) {
await logEvent(pg, runId, 7, 'iterate', 'error', 'iterate failed', { error: err.message });
}
} else {
await logEvent(pg, runId, 7, 'iterate', 'info', review.approved ? 'critic approved — no changes' : 'no fix_instructions', null);
}
// Re-run smoke test on revised content if iterate ran
if (didIterate) {
const smoke2 = await runSmokeTest({ outPath: scaffold.outPath, spec });
await logEvent(pg, runId, 5, 'smoke_test',
smoke2.passing ? 'info' : 'warn',
`post-iterate: ${smoke2.passed}/${smoke2.total} checks passing`,
smoke2);
}
// Persist the final review on the artifact row so activate can gate on it.
await pg.query("UPDATE artifacts SET status=$1 WHERE run_id=$2",
[finalReview.approved ? 'awaiting_activation' : 'unapproved', runId]);
// === Stage 9 + 10: deferred to POST /runs/:id/activate =============
await setStage(pg, runId, 9, 'running');
await logEvent(pg, runId, 9, 'memory_write', 'info', 'deferred — fires inside activate handler', null);
await setStage(pg, runId, 10, 'running');
await logEvent(pg, runId, 10, 'activate', 'info', 'awaiting manual activation (POST /runs/:id/activate)', null);
// === Stage 11: watchdog_hook (no-op for v1) ========================
await setStage(pg, runId, 11, 'running');
await logEvent(pg, runId, 11, 'watchdog_hook', 'info', 'no-op for v1 artifacts (only relevant for service-type artifacts in v2)', null);
const finalStatus = finalReview.approved
? (didIterate ? 'awaiting_activation_v2' : 'awaiting_activation')
: 'completed_unapproved';
await finishRun(pg, runId, finalStatus);
} catch (err) {
await logEvent(pg, runId, run.current_stage || 0, STAGES[run.current_stage - 1] || 'unknown', 'error', err.message, { stack: err.stack });
await finishRun(pg, runId, 'failed');
}
}
module.exports = { STAGES, runPipeline };