← 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 };