101 lines
3.8 KiB
TypeScript
101 lines
3.8 KiB
TypeScript
import { NextRequest, NextResponse } from "next/server";
|
|
import { getActiveDeployment } from "@/lib/server/activeDeployment";
|
|
|
|
export const runtime = "nodejs";
|
|
|
|
type RouteContext = { params: Promise<{ threadId: string }> };
|
|
|
|
type AsyncTask = {
|
|
task_id: string;
|
|
agent_name: string;
|
|
thread_id: string;
|
|
run_id: string;
|
|
status: string;
|
|
created_at?: string;
|
|
last_updated_at?: string;
|
|
};
|
|
|
|
function tasksFromState(value: unknown): AsyncTask[] {
|
|
if (!value || typeof value !== "object") return [];
|
|
const result: AsyncTask[] = [];
|
|
for (const item of Object.values(value as Record<string, unknown>)) {
|
|
if (!item || typeof item !== "object") continue;
|
|
const task = item as Record<string, unknown>;
|
|
if (typeof task.task_id !== "string" || typeof task.agent_name !== "string") continue;
|
|
result.push({
|
|
task_id: task.task_id,
|
|
agent_name: task.agent_name,
|
|
thread_id: typeof task.thread_id === "string" ? task.thread_id : task.task_id,
|
|
run_id: typeof task.run_id === "string" ? task.run_id : "",
|
|
status: typeof task.status === "string" ? task.status : "unknown",
|
|
created_at: typeof task.created_at === "string" ? task.created_at : undefined,
|
|
last_updated_at: typeof task.last_updated_at === "string" ? task.last_updated_at : undefined,
|
|
});
|
|
}
|
|
return result;
|
|
}
|
|
|
|
async function contextFor(threadId: string) {
|
|
const deployment = await getActiveDeployment();
|
|
if (!deployment.scopeRegistry) throw new Error("Workspace scope service is unavailable.");
|
|
const scope = await deployment.scopeRegistry.getByThread(threadId);
|
|
if (scope.state === "deleting" || scope.state === "deleted") throw new Error("Conversation is not available.");
|
|
const state = await deployment.threadClient.threads.getState(threadId);
|
|
return { deployment, scope, tasks: tasksFromState((state as { values?: { async_tasks?: unknown } }).values?.async_tasks) };
|
|
}
|
|
|
|
async function ownerForTask(
|
|
context: Awaited<ReturnType<typeof contextFor>>,
|
|
task: AsyncTask
|
|
) {
|
|
try {
|
|
return await context.deployment.scopeRegistry!.registerOwner(context.scope.scope_id, {
|
|
ownerType: "async_subagent",
|
|
resourceId: task.thread_id,
|
|
parentOwnerId: context.scope.primary_owner_id,
|
|
state: "active",
|
|
});
|
|
} catch {
|
|
// A retry normally collides with the durable owner. The originating run
|
|
// already carries its owner id; read-only status/detail access can proceed
|
|
// through the parent task record without exposing a browser SDK client.
|
|
return context.deployment.scopeRegistry!.getOwnerByResource(
|
|
context.scope.scope_id,
|
|
task.thread_id
|
|
).catch(() => null);
|
|
}
|
|
}
|
|
|
|
export async function GET(_request: NextRequest, route: RouteContext) {
|
|
try {
|
|
const { threadId } = await route.params;
|
|
const context = await contextFor(threadId);
|
|
const tasks = await Promise.all(
|
|
context.tasks.map(async (task) => {
|
|
await ownerForTask(context, task);
|
|
if (!task.run_id) return { ...task, liveStatus: task.status, startedAt: task.created_at };
|
|
try {
|
|
const run = await context.deployment.threadClient.runs.get(task.thread_id, task.run_id);
|
|
const status = run.status ?? task.status;
|
|
return {
|
|
...task,
|
|
liveStatus: status,
|
|
startedAt: run.created_at ?? task.created_at,
|
|
endedAt: ["success", "error", "timeout", "cancelled", "interrupted"].includes(status)
|
|
? run.updated_at ?? task.last_updated_at
|
|
: undefined,
|
|
};
|
|
} catch {
|
|
return { ...task, liveStatus: "expired", startedAt: task.created_at, endedAt: task.last_updated_at ?? task.created_at };
|
|
}
|
|
})
|
|
);
|
|
return NextResponse.json({ tasks });
|
|
} catch (error) {
|
|
return NextResponse.json(
|
|
{ error: error instanceof Error ? error.message : "Failed to load async tasks." },
|
|
{ status: 400 }
|
|
);
|
|
}
|
|
}
|