diff --git a/skill/scripts/live/codex-worker-supervisor.mjs b/skill/scripts/live/codex-worker-supervisor.mjs index eac8e231e..188bb130d 100644 --- a/skill/scripts/live/codex-worker-supervisor.mjs +++ b/skill/scripts/live/codex-worker-supervisor.mjs @@ -157,6 +157,9 @@ export class CodexLiveWorkerSupervisor { const handled = await this.handleAccept(event, this.base, this.token, { deferReply: event.type === 'accept', }); + if (handled?._acceptResult?.handled !== true) { + this.log(`${event.type} ${event.id} source update failed: ${handled?._acceptResult?.error || 'unhandled'}`); + } if (event.type === 'accept' && handled?._acceptResult?.carbonize === true) { await this.postCleanup(this.base, this.token, { id: event.id, @@ -326,6 +329,7 @@ export class CodexLiveWorkerSupervisor { let publishedFromMessage = false; let publicationPromise = null; let earlyCandidateError = null; + let durableCandidate = null; const publishCandidate = async (answer) => { if (publishedFromMessage || this.isCanceled(event.id)) return; if (!publicationPromise) { @@ -335,34 +339,49 @@ export class CodexLiveWorkerSupervisor { phase: generationPhaseName(phase, 'validating'), durationMs: Date.now() - phaseStartedAt, }); - const applied = applyCodexWorkerOutput({ - output: answer, - prepared, - phase, - expectedVariants: Number(event.count || arrivedVariants), - sessionId: event.id, - scaffold: event.scaffold, - cwd: this.cwd, - maxBytes: this.config.maxArtifactBytes, - }); - if (!prepared.previewMode && (phase === 'second' || phase === 'final')) { + if (!durableCandidate) { const candidatePath = path.resolve(this.cwd, prepared.artifactFile); - const reconciled = reconcilePublishedSourceVariants({ - current: artifact.content, - candidate: fs.readFileSync(candidatePath, 'utf-8'), - priorArrived: Math.max(1, arrivedVariants - 1), + if (!prepared.previewMode && (phase === 'first' || phase === 'second')) { + // A structured agent message and the final turn result can contain + // the same delta. Always apply against the immutable phase input so + // a failed publication/checkpoint retry cannot double-insert it. + fs.writeFileSync(candidatePath, artifact.content, 'utf-8'); + } + const applied = applyCodexWorkerOutput({ + output: answer, + prepared, + phase, + expectedVariants: Number(event.count || arrivedVariants), + sessionId: event.id, + scaffold: event.scaffold, + cwd: this.cwd, + maxBytes: this.config.maxArtifactBytes, }); - if (!reconciled.ok) throw supervisorError(`reconcile_${reconciled.error}`); - fs.writeFileSync(candidatePath, reconciled.content, 'utf-8'); + if (!prepared.previewMode && (phase === 'second' || phase === 'final')) { + const reconciled = reconcilePublishedSourceVariants({ + current: artifact.content, + candidate: fs.readFileSync(candidatePath, 'utf-8'), + priorArrived: Math.max(1, arrivedVariants - 1), + }); + if (!reconciled.ok) throw supervisorError(`reconcile_${reconciled.error}`); + fs.writeFileSync(candidatePath, reconciled.content, 'utf-8'); + } + if (this.isCanceled(event.id)) return; + const published = publishCodexWorkerPhase({ event, prepared, arrivedVariants, cwd: this.cwd }); + durableCandidate = { applied, published, planRecorded: false }; } - if (applied.plan) { - this.sessionStore.appendEvent({ type: 'variant_plan', id: event.id, plan: applied.plan }); + if (durableCandidate.applied.plan && !durableCandidate.planRecorded) { + this.sessionStore.appendEvent({ + type: 'variant_plan', + id: event.id, + plan: durableCandidate.applied.plan, + }); + durableCandidate.planRecorded = true; } if (this.isCanceled(event.id)) return; - const published = publishCodexWorkerPhase({ event, prepared, arrivedVariants, cwd: this.cwd }); await this.publishCheckpoint(this.base, this.token, { event, - published, + published: durableCandidate.published, scaffold: event.scaffold, arrivedVariants, }); @@ -373,7 +392,7 @@ export class CodexLiveWorkerSupervisor { try { await pendingPublication; } catch (error) { - earlyCandidateError = error; + if (!earlyCandidateError) earlyCandidateError = error; } finally { if (publicationPromise === pendingPublication) publicationPromise = null; } @@ -384,7 +403,7 @@ export class CodexLiveWorkerSupervisor { outputSchema: codexWorkerOutputSchemaForPhase( phase, Number(event.count || arrivedVariants), - { sourceDelta: phase === 'second' && !prepared.previewMode }, + { sourceDelta: (phase === 'first' || phase === 'second') && !prepared.previewMode }, ), onAgentMessage: publishCandidate, eventId: event.id, diff --git a/skill/scripts/live/codex-worker.mjs b/skill/scripts/live/codex-worker.mjs index 0693fff96..8d3827b1d 100644 --- a/skill/scripts/live/codex-worker.mjs +++ b/skill/scripts/live/codex-worker.mjs @@ -57,31 +57,13 @@ export const CODEX_WORKER_OUTPUT_SCHEMA = Object.freeze({ required: ['files'], additionalProperties: false, }); -const CODEX_SOURCE_DELTA_OUTPUT_SCHEMA = Object.freeze({ - type: 'object', - properties: { - sourceDelta: { - type: 'object', - properties: { - variantId: { type: 'integer', minimum: 2, maximum: 2 }, - markup: { type: 'string', minLength: 1 }, - css: { type: 'string', minLength: 1 }, - }, - required: ['variantId', 'markup', 'css'], - additionalProperties: false, - }, - }, - required: ['sourceDelta'], - additionalProperties: false, -}); - export function codexWorkerOutputSchemaForPhase( phase, expectedVariants = 3, { sourceDelta = false } = {}, ) { - if (sourceDelta) return CODEX_SOURCE_DELTA_OUTPUT_SCHEMA; const requirePlan = Number(expectedVariants) > 1 && (phase === 'first' || phase === 'atomic'); + if (sourceDelta) return codexSourceDeltaOutputSchema(phase, requirePlan); return { ...CODEX_WORKER_OUTPUT_SCHEMA, properties: requirePlan @@ -91,6 +73,28 @@ export function codexWorkerOutputSchemaForPhase( }; } +function codexSourceDeltaOutputSchema(phase, requirePlan) { + const variantId = phase === 'first' ? 1 : 2; + const sourceDelta = { + type: 'object', + properties: { + variantId: { type: 'integer', minimum: variantId, maximum: variantId }, + markup: { type: 'string', minLength: 1 }, + css: { type: 'string', minLength: 1 }, + }, + required: ['variantId', 'markup', 'css'], + additionalProperties: false, + }; + return { + type: 'object', + properties: requirePlan + ? { sourceDelta, plan: VARIANT_PLAN_SCHEMA } + : { sourceDelta }, + required: requirePlan ? ['sourceDelta', 'plan'] : ['sourceDelta'], + additionalProperties: false, + }; +} + export function resolveCodexWorkerConfig({ env = process.env, liveConfig = {} } = {}) { const configured = liveConfig.experimentalCodexWorker || liveConfig.codexWorker || {}; const envEnabled = parseBoolean(env.IMPECCABLE_LIVE_CODEX_WORKER); @@ -161,7 +165,8 @@ export function buildGenerationTurnInput({ const first = phase === 'first'; const second = phase === 'second'; const component = Boolean(prepared.previewMode); - const sourceDelta = second && !component; + const sourceDelta = !component && (first || second); + const sourceDeltaVariant = first ? 1 : 2; const actionRules = event.action === 'bolder' && count > 1 ? [ 'For /bolder, keep variant 1 low-risk: preserve the selected root’s high-level layout and create impact through controlled hierarchy, proportion, or rhythm. Reserve root recomposition for variant 2 or 3.', @@ -199,12 +204,12 @@ export function buildGenerationTurnInput({ ...phaseRules, ...actionRules, sourceDelta - ? 'Return exactly sourceDelta for variant 2. markup is only the selected root replacement, without an outer data-impeccable wrapper. css is only the complete fenced CSS for variant 2, following event.scaffold.cssAuthoring.' + ? `Return exactly sourceDelta for variant ${sourceDeltaVariant}${first && count > 1 ? ' plus the complete variant plan' : ''}. markup is only the selected root replacement, without an outer data-impeccable wrapper. css is only the complete fenced CSS for variant ${sourceDeltaVariant}, following event.scaffold.cssAuthoring.` : component ? `Return staged component files relative to componentDir. Allowed variant extension: .${artifact.componentExtension}. The supervisor updates manifest.json.` : `Return exactly one file whose path is ${JSON.stringify(prepared.artifactFile)} and whose content is the complete staged source artifact.`, sourceDelta - ? 'Do not repeat the staged artifact, variant 1, style tags, wrapper comments, or any data-impeccable attributes. The supervisor merges and validates this delta transactionally.' + ? `Do not repeat the staged artifact${second ? ', variant 1' : ''}, style tags, wrapper comments, or any data-impeccable attributes. The supervisor merges and validates this delta transactionally.` : component ? 'For the final/atomic phase include params.json keyed by variant number. Never include manifest.json or paths outside componentDir.' : 'Keep the existing session wrapper and markers intact. Add only valid variant blocks and preview CSS inside that wrapper.', @@ -300,18 +305,24 @@ export function applyCodexWorkerOutput({ maxBytes = 2_000_000, }) { const parsed = typeof output === 'string' ? parseWorkerJson(output) : output; - if (!prepared.previewMode && phase === 'second') { + const requirePlan = Number(expectedVariants) > 1 && (phase === 'first' || phase === 'atomic'); + if (requirePlan && !parsed?.plan) throw workerError('worker_output_plan_missing'); + const plan = parsed?.plan ? normalizeVariantPlan(parsed.plan, expectedVariants) : null; + if (!prepared.previewMode && (phase === 'first' || phase === 'second')) { const artifactPath = resolveInside(cwd, prepared.artifactFile); if (!artifactPath) throw workerError('artifact_path_outside_project'); const content = applyCodexSourceDelta({ source: fs.readFileSync(artifactPath, 'utf-8'), delta: parsed?.sourceDelta, sessionId, + expectedVariantId: phase === 'first' ? 1 : 2, styleMode: scaffold?.styleMode || scaffold?.cssAuthoring?.mode || 'scoped', + styleTag: scaffold?.styleTag, + jsx: scaffold?.commentSyntax?.open === '{/*', }); if (Buffer.byteLength(content) > maxBytes) throw workerError('worker_output_too_large'); fs.writeFileSync(artifactPath, content, 'utf-8'); - return { files: [prepared.artifactFile], plan: null, sourceDelta: true }; + return { files: [prepared.artifactFile], plan, sourceDelta: true }; } if (!Array.isArray(parsed?.files) || parsed.files.length === 0) { throw workerError('worker_output_files_missing'); @@ -327,10 +338,6 @@ export function applyCodexWorkerOutput({ totalBytes += Buffer.byteLength(file.content); } if (totalBytes > maxBytes) throw workerError('worker_output_too_large'); - const requirePlan = Number(expectedVariants) > 1 && (phase === 'first' || phase === 'atomic'); - if (requirePlan && !parsed.plan) throw workerError('worker_output_plan_missing'); - const plan = parsed.plan ? normalizeVariantPlan(parsed.plan, expectedVariants) : null; - if (!prepared.previewMode) { if (parsed.files.length !== 1 || parsed.files[0].path !== prepared.artifactFile) { throw workerError('worker_output_source_path_invalid'); @@ -390,12 +397,18 @@ export function applyCodexSourceDelta({ source, delta, sessionId, + expectedVariantId = 2, styleMode = 'scoped', + styleTag = null, + jsx = false, }) { if (!delta || typeof delta !== 'object' || Array.isArray(delta)) { throw workerError('worker_output_source_delta_missing'); } - if (Number(delta.variantId) !== 2) throw workerError('worker_output_source_delta_variant_invalid'); + const variantId = Number(expectedVariantId); + if (![1, 2].includes(variantId) || Number(delta.variantId) !== variantId) { + throw workerError('worker_output_source_delta_variant_invalid'); + } const markup = String(delta.markup || '').trim(); const css = String(delta.css || '').trim(); if (!markup || !css) throw workerError('worker_output_source_delta_empty'); @@ -407,11 +420,12 @@ export function applyCodexSourceDelta({ } const cssVariantRefs = [...css.matchAll(/\[data-impeccable-variant=(?:"([^"]+)"|'([^']+)')\]/g)] .map((match) => match[1] || match[2]); - if (cssVariantRefs.length === 0 || cssVariantRefs.some((variant) => variant !== '2')) { + if (cssVariantRefs.length === 0 || cssVariantRefs.some((variant) => variant !== String(variantId))) { throw workerError('worker_output_source_delta_css_unfenced'); } const astroGlobal = styleMode === 'astro-global-prefixed'; - if (astroGlobal ? /@scope\b/.test(css) : !/@scope\s*\(\s*\[data-impeccable-variant=(?:"2"|'2')\]\s*\)/.test(css)) { + const scopePattern = new RegExp(`@scope\\s*\\(\\s*\\[data-impeccable-variant=(?:"${variantId}"|'${variantId}')\\]\\s*\\)`); + if (astroGlobal ? /@scope\b/.test(css) : !scopePattern.test(css)) { throw workerError('worker_output_source_delta_css_strategy_invalid'); } @@ -419,48 +433,65 @@ export function applyCodexSourceDelta({ if (!id) throw workerError('worker_output_source_delta_session_missing'); const wrapper = findSessionWrapper(source, id); if (!wrapper) throw workerError('worker_output_source_delta_wrapper_missing'); - if (extractSourceVariantBlock(source, 2)) throw workerError('worker_output_source_delta_variant_exists'); + const wrapperSource = source.slice(wrapper.openStart, wrapper.closeEnd); + if (extractSourceVariantBlock(wrapperSource, variantId)) throw workerError('worker_output_source_delta_variant_exists'); const escapedId = escapeRegExp(id); const styleOpen = new RegExp(`', styleContentStart); - if (styleClose < 0 || styleClose > wrapper.closeEnd) { - throw workerError('worker_output_source_delta_style_invalid'); - } - const styleContent = source.slice(styleContentStart, styleClose); - let nextStyleContent; - const firstTick = styleContent.indexOf('`'); - const lastTick = styleContent.lastIndexOf('`'); - if (firstTick >= 0 || lastTick >= 0) { - if (firstTick < 0 || lastTick <= firstTick) { + let merged = source; + let newStyleBlock = null; + if (styleMatch) { + const styleContentStart = styleMatch.index + styleMatch[0].length; + const styleClose = source.indexOf('', styleContentStart); + if (styleClose < 0 || styleClose > wrapper.closeEnd) { throw workerError('worker_output_source_delta_style_invalid'); } - nextStyleContent = styleContent.slice(0, lastTick).trimEnd() - + '\n' + css + '\n' - + styleContent.slice(lastTick); + const styleContent = source.slice(styleContentStart, styleClose); + let nextStyleContent; + const firstTick = styleContent.indexOf('`'); + const lastTick = styleContent.lastIndexOf('`'); + if (firstTick >= 0 || lastTick >= 0) { + if (firstTick < 0 || lastTick <= firstTick) { + throw workerError('worker_output_source_delta_style_invalid'); + } + nextStyleContent = styleContent.slice(0, lastTick).trimEnd() + + '\n' + css + '\n' + + styleContent.slice(lastTick); + } else { + nextStyleContent = styleContent.trimEnd() + '\n' + css + '\n'; + } + merged = source.slice(0, styleContentStart) + nextStyleContent + source.slice(styleClose); } else { - nextStyleContent = styleContent.trimEnd() + '\n' + css + '\n'; + if (variantId !== 1) throw workerError('worker_output_source_delta_style_missing'); + const openingTag = String(styleTag || `'].join('\n') + : [openingTag, css, ''].join('\n'); } - let merged = source.slice(0, styleContentStart) + nextStyleContent + source.slice(styleClose); const nextWrapper = findSessionWrapper(merged, id); if (!nextWrapper) throw workerError('worker_output_source_delta_wrapper_missing'); + const endMarker = findSessionEndMarker(merged, id, nextWrapper); const closeLineStart = merged.lastIndexOf('\n', nextWrapper.closeStart) + 1; const closeLinePrefix = merged.slice(closeLineStart, nextWrapper.closeStart); - const childIndent = nextWrapper.indent + ' '; + const childIndent = endMarker?.indent || nextWrapper.indent + ' '; const contentIndent = childIndent + ' '; const indentedMarkup = markup.split('\n') .map((line) => line.trim() ? contentIndent + line : '') .join('\n'); const variantBlock = [ - `${childIndent}