Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
217 changes: 123 additions & 94 deletions src/commands/start.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,10 +8,11 @@ import { findProjectRoot, loadConfig } from '../lib/config.js';
import * as git from '../lib/git.js';
import * as logger from '../lib/logger.js';
import { requireAnthropicApiKey } from '../lib/requirements.js';
import { createState, hasActiveTask, loadState, saveIterationLog, saveState } from '../lib/state.js';
import { createState, hasActiveTask, loadState, saveIterationLog, saveState, updateStep } from '../lib/state.js';
import { runStepsConcurrently, type StepResult } from '../lib/runner.js';
import { submitActiveTask } from '../lib/submit-task.js';
import { buildInitialStepState, ensureTaskDirectory, loadAndValidateTask, moveTaskFileAtomically } from '../lib/tasks.js';
import type { ArbiterResult, ReviewComment } from '../types/index.js';
import type { ArbiterResult, ReviewComment, TaskStep } from '../types/index.js';

const POLL_INTERVAL_MS = 5_000;
const POLL_TIMEOUT_MS = 10 * 60_000;
Expand Down Expand Up @@ -67,107 +68,133 @@ export async function runStart(taskFile: string, options: StartCommandOptions):
}

const claude = new ClaudeClient(process.env.ANTHROPIC_API_KEY ?? '');
const total = task.steps.length;

for (let i = 0; i < task.steps.length; i += 1) {
const step = task.steps[i];
const stepState = state.steps[i];
const stepResults = await runStepsConcurrently(
task.steps,
{ maxConcurrent: config.maxConcurrent },
async (step: TaskStep): Promise<StepResult> => {
const scopedLogger = logger.withPrefix(`[${step.service}]`);
const latestState = options.dryRun ? state : loadState(projectRoot);
const stepState = latestState?.steps.find((item) => item.service === step.service);

if (stepState.status === 'done') {
continue;
}

if (step.depends_on && step.depends_on.length > 0) {
for (const depService of step.depends_on) {
const depState = state.steps.find((item) => item.service === depService);
if (depState?.status !== 'done') {
fatalAndExit(`Step dependency '${depService}' for service '${step.service}' is not done.`);
}
if (!stepState) {
return { service: step.service, status: 'failed', error: 'step_state_missing' };
}
}

logger.step(i + 1, total, `${step.service}: ${task.title}`);
if (stepState.status === 'done') {
return { service: step.service, status: 'done', sessionId: stepState.session_id };
}

const serviceCfg = config.services.find((service) => service.name === step.service);
if (!serviceCfg) {
fatalAndExit(`Unknown service in step: ${step.service}`);
}
const serviceCfg = config.services.find((service) => service.name === step.service);
if (!serviceCfg) {
return { service: step.service, status: 'failed', error: `Unknown service in step: ${step.service}` };
}

const serviceRoot = path.resolve(projectRoot, serviceCfg.path);
const branch = git.getBranchName(task.id, step.service);
const serviceRoot = path.resolve(projectRoot, serviceCfg.path);
const branch = stepState.branch ?? git.getBranchName(task.id, step.service);

try {
if (!options.dryRun) {
if (options.resume) {
if (stepState.branch) {
await git.checkoutBranch(stepState.branch, serviceRoot);
} else {
await git.createBranch(branch, serviceRoot);
}
} else {
await git.createBranch(branch, serviceRoot);
}

await updateStep(projectRoot, task.id, step.service, {
status: 'in_progress',
branch,
});
}

if (!options.dryRun) {
if (options.resume) {
if (stepState.branch) {
await git.checkoutBranch(stepState.branch, serviceRoot);
} else {
await git.createBranch(branch, serviceRoot);
if (options.dryRun) {
scopedLogger.info(`[dry-run] Would run codex cloud implementation for service ${step.service}`);
return { service: step.service, status: 'done' };
}
} else {
await git.createBranch(branch, serviceRoot);
}
}

stepState.status = 'in_progress';
stepState.branch = branch;
if (!options.dryRun) {
saveState(projectRoot, state);
}
scopedLogger.info('Submitting to Codex Cloud...');
const submissionSession = stepState.session_id ?? (await codex.submitTask(step.spec, { cwd: serviceRoot }));
await updateStep(projectRoot, task.id, step.service, { session_id: submissionSession });

const execution = await runCloudReviewLoop({
taskId: task.id,
service: step.service,
spec: step.spec,
sessionId: submissionSession,
stepState: {
iteration: stepState.iteration,
session_id: submissionSession,
},
projectRoot,
config,
claude,
verbose: options.verbose,
log: scopedLogger,
});

await updateStep(projectRoot, task.id, step.service, {
lastReviewComments: execution.lastReviewComments,
lastArbiterResult: execution.lastArbiterResult,
iteration: execution.finalIteration,
session_id: execution.sessionId,
});

if (execution.lastArbiterResult.decision === 'escalate') {
logger.escalation({
taskId: task.id,
service: step.service,
iteration: execution.finalIteration,
spec: step.spec,
diff: '',
reviewComments: execution.lastReviewComments,
arbiterReasoning: execution.lastArbiterResult.reasoning,
summary: execution.lastArbiterResult.summary,
});

await updateStep(projectRoot, task.id, step.service, { status: 'escalated' });
return { service: step.service, status: 'escalated', sessionId: execution.sessionId };
}

if (options.dryRun) {
logger.info(`[dry-run] Would run codex cloud implementation for service ${step.service}`);
} else {
const submissionSession = stepState.session_id ?? (await codex.submitTask(step.spec, {cwd: serviceRoot}));
stepState.session_id = submissionSession;
saveState(projectRoot, state);
await updateStep(projectRoot, task.id, step.service, { status: 'done' });
scopedLogger.success('Review passed — ready for PR');
return { service: step.service, status: 'done', sessionId: execution.sessionId };
} catch (error: unknown) {
const message = error instanceof Error ? error.message : String(error);
if (!options.dryRun) {
await updateStep(projectRoot, task.id, step.service, { status: 'failed' });
}
scopedLogger.error(message);
return { service: step.service, status: 'failed', error: message };
}
},
);

const execution = await runCloudReviewLoop({
taskId: task.id,
service: step.service,
spec: step.spec,
sessionId: submissionSession,
stepState,
projectRoot,
config,
claude,
verbose: options.verbose,
});

stepState.lastReviewComments = execution.lastReviewComments;
stepState.lastArbiterResult = execution.lastArbiterResult;
stepState.iteration = execution.finalIteration;
stepState.session_id = execution.sessionId;
if (!options.dryRun) {
for (const result of stepResults) {
if (result.status === 'failed' && result.error === 'dependency_failed') {
await updateStep(projectRoot, task.id, result.service, { status: 'failed' });
}
}
}

if (stepState.lastArbiterResult?.decision === 'escalate') {
logger.escalation({
taskId: task.id,
service: step.service,
iteration: stepState.iteration,
spec: step.spec,
diff: '',
reviewComments: stepState.lastReviewComments ?? [],
arbiterReasoning: stepState.lastArbiterResult.reasoning,
summary: stepState.lastArbiterResult.summary,
});

stepState.status = 'escalated';
state.status = 'escalated';

if (!options.dryRun) {
saveState(projectRoot, state);
const blockedDir = ensureTaskDirectory(projectRoot, 'blocked');
state.taskPath = moveTaskFileAtomically(state.taskPath, blockedDir);
saveState(projectRoot, state);
}
const hasEscalation = stepResults.some((result) => result.status === 'escalated');
const hasFailure = stepResults.some((result) => result.status === 'failed');

process.exit(1);
}
state = loadState(projectRoot) ?? state;

stepState.status = 'done';
if (hasEscalation || hasFailure) {
state.status = hasEscalation ? 'escalated' : 'blocked';
if (!options.dryRun) {
saveState(projectRoot, state);
const blockedDir = ensureTaskDirectory(projectRoot, 'blocked');
state.taskPath = moveTaskFileAtomically(state.taskPath, blockedDir);
saveState(projectRoot, state);
}
process.exit(1);
}

state.status = 'review';
Expand All @@ -194,19 +221,20 @@ async function runCloudReviewLoop(opts: {
service: string;
spec: string;
sessionId: string;
stepState: {iteration: number; session_id?: string};
stepState: { iteration: number; session_id?: string };
projectRoot: string;
config: ReturnType<typeof loadConfig>;
claude: ClaudeClient;
verbose?: boolean;
}): Promise<{sessionId: string; finalIteration: number; lastReviewComments: ReviewComment[]; lastArbiterResult: ArbiterResult}> {
log: logger.Logger;
}): Promise<{ sessionId: string; finalIteration: number; lastReviewComments: ReviewComment[]; lastArbiterResult: ArbiterResult }> {
let sessionId = opts.sessionId;
let iteration = opts.stepState.iteration;

for (;;) {
logger.iteration(iteration + 1, opts.config.review.max_iterations);
opts.log.iteration(iteration + 1, opts.config.review.max_iterations);

logger.info(`Polling codex cloud session ${sessionId}`);
opts.log.info(`Polling codex cloud session ${sessionId}`);
const status = await codex.pollStatus(sessionId, {
intervalMs: POLL_INTERVAL_MS,
timeoutMs: POLL_TIMEOUT_MS,
Expand All @@ -225,18 +253,19 @@ async function runCloudReviewLoop(opts: {
};
}

logger.info(`Retrieving diff for session ${sessionId}`);
opts.log.success('Codex task completed');
opts.log.info(`Retrieving diff for session ${sessionId}`);
const diff = await codex.getDiff(sessionId);

logger.info(`Requesting reviewer analysis (model: ${opts.config.review.model})`);
opts.log.info(`Running review (iteration ${String(iteration + 1)})...`);
const review = await opts.claude.runReviewer({
spec: opts.spec,
diff,
model: opts.config.review.model,
});
logger.reviewSummary(review.comments);
opts.log.reviewSummary(review.comments);

logger.info(`Requesting arbiter decision (model: ${opts.config.review.model})`);
opts.log.info(`Requesting arbiter decision (model: ${opts.config.review.model})`);
const arbiter = await opts.claude.runArbiter({
spec: opts.spec,
diff,
Expand Down Expand Up @@ -288,7 +317,7 @@ async function runCloudReviewLoop(opts: {
};
}

logger.info('Arbiter requested fixes, resuming codex cloud session');
opts.log.warn(`Review requested fixes (iteration ${String(iteration + 1)}/${String(opts.config.review.max_iterations)})`);
sessionId = await codex.resumeTask(sessionId, arbiter.feedback_for_codex);
opts.stepState.session_id = sessionId;
iteration += 1;
Expand Down
15 changes: 15 additions & 0 deletions src/lib/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,19 @@ function parseReview(value: unknown): ReviewConfig {
};
}


function parseMaxConcurrent(value: unknown): number | undefined {
if (value === undefined) {
return undefined;
}

if (typeof value !== 'number' || !Number.isInteger(value) || value <= 0) {
throw new Error('maxConcurrent must be a positive integer');
}

return value;
}

function parseCodex(value: unknown): CodexConfig {
if (value === undefined) {
return { model: DEFAULT_CODEX_MODEL };
Expand Down Expand Up @@ -147,11 +160,13 @@ export function loadConfig(projectRoot: string): VexdoConfig {
const services = parseServices(readObjectField(parsed, 'services'));
const review = parseReview(readObjectField(parsed, 'review'));
const codex = parseCodex(readObjectField(parsed, 'codex'));
const maxConcurrent = parseMaxConcurrent(readObjectField(parsed, 'maxConcurrent'));

return {
version: 1,
services,
review,
codex,
maxConcurrent,
};
}
66 changes: 66 additions & 0 deletions src/lib/logger.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,16 @@ import type { ReviewComment } from '../types/index.js';

let verboseEnabled = false;

export interface Logger {
info: (message: string) => void;
success: (message: string) => void;
warn: (message: string) => void;
error: (message: string) => void;
debug: (message: string) => void;
iteration: (n: number, max: number) => void;
reviewSummary: (comments: ReviewComment[]) => void;
}

function safeLog(method: 'log' | 'error', message: string): void {
try {
if (method === 'error') {
Expand Down Expand Up @@ -145,3 +155,59 @@ export function reviewSummary(comments: ReviewComment[]): void {
safeLog('log', `- ${comment.severity}${location}: ${comment.comment}`);
}
}

function prefixed(prefix: string, message: string): string {
return `${prefix} ${message}`;
}

export function withPrefix(prefix: string): Logger {
return {
info: (message: string) => {
info(prefixed(prefix, message));
},
success: (message: string) => {
success(prefixed(prefix, message));
},
warn: (message: string) => {
warn(prefixed(prefix, message));
},
error: (message: string) => {
error(prefixed(prefix, message));
},
debug: (message: string) => {
debug(prefixed(prefix, message));
},
iteration: (n: number, max: number) => {
safeLog('log', prefixed(prefix, pc.gray(`Iteration ${String(n)}/${String(max)}`)));
},
reviewSummary: (comments: ReviewComment[]) => {
const counts = {
critical: 0,
important: 0,
minor: 0,
noise: 0,
};

for (const comment of comments) {
counts[comment.severity] += 1;
}

safeLog(
'log',
prefixed(
prefix,
`${pc.bold('Review:')} ${String(counts.critical)} critical ${String(counts.important)} important ${String(counts.minor)} minor`,
),
);

for (const comment of comments) {
if (comment.severity === 'noise') {
continue;
}

const location = comment.file ? ` (${comment.file}${comment.line ? `:${String(comment.line)}` : ''})` : '';
safeLog('log', prefixed(prefix, `- ${comment.severity}${location}: ${comment.comment}`));
}
},
};
}
Loading
Loading