diff --git a/src/lib/admin/author.server.ts b/src/lib/admin/author.server.ts index cb7505be..4691fbce 100644 --- a/src/lib/admin/author.server.ts +++ b/src/lib/admin/author.server.ts @@ -85,18 +85,108 @@ export async function generateDraft( } return experimental_output; } catch (e: any) { - const msg = e?.message ?? String(e); + const msg = describeAttemptError(e); attempts.push({ model: modelId, error: msg }); - console.error(`[skillforge.author] ${modelId} failed:`, msg); + console.error(`[skillforge.author] ${modelId} structured failed:`, msg); } } - // All fallbacks exhausted. Surface a structured error so the upload + + // Structured-output path failed on every model. Most common cause is + // provider strict-mode rejecting the JSON schema generated from + // PackageDraftSchema (e.g. open `record` slots) — both + // OpenAI and Gemini surface this as opaque "Bad Request" / "no object + // generated". Fall back to plain text: ask the model for raw JSON, + // parse it ourselves, validate with Zod. Loses provider-side schema + // guarantees but lets the user actually get a draft. + try { + const fallbackModel = getGatewayModel("default"); + const { text } = await generateText({ + model: fallbackModel, + system: + META_SYSTEM + + `\n\nReturn ONLY a single JSON object matching PackageDraft. No markdown, no prose, no code fences. ` + + `Required top-level keys: slug, name, type, description, long_description, system_prompt, rules, examples. ` + + `examples must contain >= 2 items. rules has: input_schema (object), output_schema (object), must (array), must_not (array).`, + prompt, + abortSignal: AbortSignal.timeout(PER_ATTEMPT_TIMEOUT_MS), + }); + const json = extractJsonObject(text); + const parsed = PackageDraftSchema.safeParse(json); + if (!parsed.success) { + attempts.push({ + model: "text-fallback", + error: `zod: ${parsed.error.issues.slice(0, 3).map((i) => `${i.path.join(".")}: ${i.message}`).join("; ")}`, + }); + } else { + const out = parsed.data; + if (out.type !== type) out.type = type; + console.warn( + `[skillforge.author] recovered via text fallback after ${attempts.length} structured failure(s):`, + attempts, + ); + return out; + } + } catch (e: any) { + attempts.push({ model: "text-fallback", error: describeAttemptError(e) }); + console.error(`[skillforge.author] text fallback failed:`, attempts[attempts.length - 1]?.error); + } + + // All paths exhausted. Surface a structured error so the upload // pipeline can report it back to the caller instead of silently // recording "failed". const summary = attempts.map((a) => `${a.model}: ${a.error}`).join(" | "); throw new Error(`SkillForge author failed across all fallback models — ${summary}`); } +// AI SDK errors hide the upstream body inside `e.text` (NoObjectGeneratedError), +// `e.responseBody` (HTTP errors), or `e.cause`. Bare `e.message` strips all of +// that and leaves operators chasing "Bad Request" with no context. +function describeAttemptError(e: any): string { + const parts: string[] = []; + if (e?.message) parts.push(String(e.message)); + if (e?.name && !parts[0]?.includes(e.name)) parts.unshift(String(e.name)); + if (e?.text && typeof e.text === "string") parts.push(`text=${e.text.slice(0, 400)}`); + if (e?.responseBody && typeof e.responseBody === "string") + parts.push(`body=${e.responseBody.slice(0, 400)}`); + if (e?.cause?.message && e.cause.message !== e?.message) + parts.push(`cause=${String(e.cause.message).slice(0, 200)}`); + return parts.join(" | ") || String(e); +} + +// Permissive JSON extractor: handles plain JSON, markdown-fenced JSON, and +// preamble/postamble text the model sometimes adds despite instructions. +function extractJsonObject(text: string): unknown { + const stripped = text.trim(); + const tryParse = (s: string) => { + try { return JSON.parse(s); } catch { return undefined; } + }; + const direct = tryParse(stripped); + if (direct !== undefined) return direct; + const fence = stripped.match(/```(?:json)?\s*([\s\S]*?)```/i); + if (fence?.[1]) { + const v = tryParse(fence[1].trim()); + if (v !== undefined) return v; + } + // Last resort: locate the first {...} block by brace balancing. + const start = stripped.indexOf("{"); + if (start >= 0) { + let depth = 0; + for (let i = start; i < stripped.length; i++) { + const c = stripped[i]; + if (c === "{") depth++; + else if (c === "}") { + depth--; + if (depth === 0) { + const v = tryParse(stripped.slice(start, i + 1)); + if (v !== undefined) return v; + break; + } + } + } + } + throw new Error(`text fallback returned non-JSON: ${stripped.slice(0, 200)}`); +} + export async function insertDraftPackage( supabase: any, userId: string, diff --git a/src/lib/uploads/queue.server.ts b/src/lib/uploads/queue.server.ts index bf95184a..818d42a9 100644 --- a/src/lib/uploads/queue.server.ts +++ b/src/lib/uploads/queue.server.ts @@ -35,26 +35,74 @@ export async function enqueueUploadJobs( // Drain up to `limit` queued jobs. Used by the cron endpoint AND // fire-and-forget from the MCP tool so a single-file follow-up upload // also nudges the queue forward. +// Jobs that sit in `processing` longer than this without a finished_at are +// assumed orphaned (worker crashed, Vercel function timed out before the +// status flip). The next drain pass flips them back to `queued` so they +// don't become permanent zombies in the user's /account/packages list. +const STALE_PROCESSING_MS = 2 * 60 * 1000; +const MAX_ATTEMPTS = 3; + export async function drainUploadQueue(limit = 3): Promise<{ processed: number; failed: number; + requeued: number; remaining: number; }> { let processed = 0; let failed = 0; + let requeued = 0; + + // Heal: any job stuck in `processing` past STALE_PROCESSING_MS — the + // worker that claimed it died (Vercel 60s budget, OOM, network flap). + // Flip back to `queued` unless we've already retried MAX_ATTEMPTS times, + // in which case mark `failed` so the user sees a real outcome instead + // of an indefinite spinner. + const staleCutoff = new Date(Date.now() - STALE_PROCESSING_MS).toISOString(); + const { data: stale } = await supabaseAdmin + .from("package_upload_jobs") + .select("id, attempts") + .eq("status", "processing") + .lt("started_at", staleCutoff); + for (const row of (stale ?? []) as Array<{ id: string; attempts: number }>) { + if ((row.attempts ?? 0) >= MAX_ATTEMPTS) { + await supabaseAdmin + .from("package_upload_jobs") + .update({ + status: "failed", + finished_at: new Date().toISOString(), + error: `abandoned after ${row.attempts} attempt(s) — worker crashed or timed out`, + }) + .eq("id", row.id) + .eq("status", "processing"); + failed++; + } else { + await supabaseAdmin + .from("package_upload_jobs") + .update({ status: "queued", started_at: null }) + .eq("id", row.id) + .eq("status", "processing"); + requeued++; + } + } + for (let i = 0; i < limit; i++) { // Claim one job atomically: flip queued → processing if still queued. const { data: claimed } = await supabaseAdmin .from("package_upload_jobs") - .select("id") + .select("id, attempts") .eq("status", "queued") .order("created_at", { ascending: true }) .limit(1) .maybeSingle(); if (!claimed) break; + const nextAttempts = ((claimed as { attempts?: number }).attempts ?? 0) + 1; const { data: locked } = await supabaseAdmin .from("package_upload_jobs") - .update({ status: "processing", started_at: new Date().toISOString(), attempts: 1 }) + .update({ + status: "processing", + started_at: new Date().toISOString(), + attempts: nextAttempts, + }) .eq("id", claimed.id) .eq("status", "queued") .select("id, user_id, filename, content, inferred_type") @@ -105,5 +153,5 @@ export async function drainUploadQueue(limit = 3): Promise<{ .from("package_upload_jobs") .select("id", { count: "exact", head: true }) .eq("status", "queued"); - return { processed, failed, remaining: count ?? 0 }; + return { processed, failed, requeued, remaining: count ?? 0 }; }