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
96 changes: 93 additions & 3 deletions src/lib/admin/author.server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, any>` 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,
Expand Down
54 changes: 51 additions & 3 deletions src/lib/uploads/queue.server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down Expand Up @@ -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 };
}
Loading