import { z } from "zod"; import { CreateScheduledTaskRequest, TriggerScheduledTaskRequest, UpdateScheduledTaskRequest, } from "@opengeni/contracts"; import { listScheduledTaskRuns, listScheduledTasks } from "@opengeni/db"; import type { Hono } from "hono"; import { HTTPException } from "hono/http-exception"; import { requireAccessGrant, requireAccessGrantAuthorization, resolveWorkspaceCatalogSettings, } from "@opengeni/core"; import { recordWorkspaceUsage, requireLimit } from "@opengeni/core"; import type { ApiRouteDeps } from "@opengeni/core"; import { captureScheduledTaskRestoreState, createValidatedScheduledTask, manualScheduledTaskTriggerUsageKey, manualScheduledTaskTriggerWorkflowId, scheduledTaskToolsProvided, scheduledTaskForGrant, scheduledTaskRunForGrant, scheduledTaskTriggerToken, requireScheduledTaskForApi, syncCreatedScheduledTask, syncUpdatedScheduledTask, updateScheduledTaskForApi, validateScheduledTaskMachineTarget, validateScheduledTaskTarget, validatedScheduledTaskUpdate, } from "@opengeni/core"; import { boundedLimit } from "../http/common"; import { deleteScheduledTaskWithDurableCleanup } from "../scheduled-task-deletion"; export function registerScheduledTaskRoutes(app: Hono, deps: ApiRouteDeps): void { const { db, workflowClient, objectStorage } = deps; app.post("/v1/workspaces/:workspaceId/scheduled-tasks", async (c) => { const workspaceId = c.req.param("workspaceId"); const authorization = await requireAccessGrantAuthorization( c, deps, workspaceId, "scheduled_tasks:manage", ); const grant = authorization.grant; const rawPayload = await c.req.json(); const parsedPayload = CreateScheduledTaskRequest.safeParse(rawPayload); if (!parsedPayload.success) { throw new HTTPException(400, { message: "invalid scheduled task create request", }); } const payload = parsedPayload.data; const catalogSettings = ( await resolveWorkspaceCatalogSettings(db, deps.settings, { accountId: grant.accountId, workspaceId, }) ).settings; await requireLimit(deps, { accountId: grant.accountId, workspaceId, action: "schedule:create", quantity: 1, }); const task = await createValidatedScheduledTask({ settings: catalogSettings, db, objectStorage, grant, authorization, payload, toolsProvided: scheduledTaskToolsProvided(rawPayload), sessionAuthorization: deps.sessionAuthorization, authorizationSurface: "http", }); await syncCreatedScheduledTask({ db, workflowClient, task }); return c.json(scheduledTaskForGrant(task, grant), 201); }); app.get("/v1/workspaces/:workspaceId/scheduled-tasks", async (c) => { const workspaceId = c.req.param("workspaceId"); const grant = await requireAccessGrant(c, deps, workspaceId, "scheduled_tasks:run"); const sessionId = c.req.query("sessionId"); const offset = Number(c.req.query("offset") ?? 0); if (!Number.isSafeInteger(offset) || offset < 0) { throw new HTTPException(400, { message: "invalid offset" }); } if (sessionId !== undefined) { if (!z.string().uuid().safeParse(sessionId).success) { throw new HTTPException(400, { message: "invalid sessionId" }); } await requireAccessGrant(c, deps, workspaceId, "sessions:control"); } const tasks = await listScheduledTasks( db, workspaceId, boundedLimit(c.req.query("limit")), offset, sessionId, ); return c.json(tasks.map((task) => scheduledTaskForGrant(task, grant))); }); app.get("/v1/workspaces/:workspaceId/scheduled-tasks/:taskId", async (c) => { const workspaceId = c.req.param("workspaceId"); const grant = await requireAccessGrant(c, deps, workspaceId, "scheduled_tasks:run"); const task = await requireScheduledTaskForApi(db, workspaceId, c.req.param("taskId")); return c.json(scheduledTaskForGrant(task, grant)); }); app.patch("/v1/workspaces/:workspaceId/scheduled-tasks/:taskId", async (c) => { const workspaceId = c.req.param("workspaceId"); const authorization = await requireAccessGrantAuthorization( c, deps, workspaceId, "scheduled_tasks:manage", ); const grant = authorization.grant; const taskId = c.req.param("taskId"); const existing = await requireScheduledTaskForApi(db, workspaceId, taskId); const previous = await captureScheduledTaskRestoreState(db, existing); const rawPayload = await c.req.json(); const parsedPayload = UpdateScheduledTaskRequest.safeParse(rawPayload); if (!parsedPayload.success) { throw new HTTPException(400, { message: "invalid scheduled task update request", }); } const payload = parsedPayload.data; const catalogSettings = ( await resolveWorkspaceCatalogSettings(db, deps.settings, { accountId: grant.accountId, workspaceId, }) ).settings; const update = await validatedScheduledTaskUpdate({ settings: catalogSettings, db, objectStorage, grant, existing, authorization, payload, toolsProvided: scheduledTaskToolsProvided(rawPayload), sessionAuthorization: deps.sessionAuthorization, authorizationSurface: "http", }); const task = await updateScheduledTaskForApi( db, workspaceId, taskId, update, payload.agentLearning ? { authorization, request: payload.agentLearning, restoreState: previous } : undefined, ); await syncUpdatedScheduledTask({ db, workflowClient, previous, task }); return c.json(scheduledTaskForGrant(task, grant)); }); app.post("/v1/workspaces/:workspaceId/scheduled-tasks/:taskId/pause", async (c) => { const workspaceId = c.req.param("workspaceId"); const grant = await requireAccessGrant(c, deps, workspaceId, "scheduled_tasks:manage"); const existing = await requireScheduledTaskForApi(db, workspaceId, c.req.param("taskId")); const previous = await captureScheduledTaskRestoreState(db, existing); const task = await updateScheduledTaskForApi(db, workspaceId, existing.id, { status: "paused", }); await syncUpdatedScheduledTask({ db, workflowClient, previous, task }); return c.json(scheduledTaskForGrant(task, grant)); }); app.post("/v1/workspaces/:workspaceId/scheduled-tasks/:taskId/resume", async (c) => { const workspaceId = c.req.param("workspaceId"); const authorization = await requireAccessGrantAuthorization( c, deps, workspaceId, "scheduled_tasks:manage", ); const grant = authorization.grant; const existing = await requireScheduledTaskForApi(db, workspaceId, c.req.param("taskId")); const previous = await captureScheduledTaskRestoreState(db, existing); const catalogSettings = ( await resolveWorkspaceCatalogSettings(db, deps.settings, { accountId: grant.accountId, workspaceId, }) ).settings; const update = await validatedScheduledTaskUpdate({ settings: catalogSettings, db, objectStorage, grant, existing, payload: { status: "active" }, authorization, sessionAuthorization: deps.sessionAuthorization, authorizationSurface: "http", }); const task = await updateScheduledTaskForApi(db, workspaceId, existing.id, update); await syncUpdatedScheduledTask({ db, workflowClient, previous, task }); return c.json(scheduledTaskForGrant(task, grant)); }); app.post("/v1/workspaces/:workspaceId/scheduled-tasks/:taskId/trigger", async (c) => { const workspaceId = c.req.param("workspaceId"); const grant = await requireAccessGrant(c, deps, workspaceId, "scheduled_tasks:run"); // Load the task before the gate so a codex-model scheduled task can be // recognised as codex-billed and skip the credit/cost gates at the edge. const task = await requireScheduledTaskForApi(db, workspaceId, c.req.param("taskId")); if (task.action.kind === "agent_turn") { const catalogSettings = ( await resolveWorkspaceCatalogSettings(db, deps.settings, { accountId: grant.accountId, workspaceId, }) ).settings; await validateScheduledTaskTarget({ db, sessionAuthorization: deps.sessionAuthorization, authorizationSurface: "http", grant, targetSessionId: task.targetSessionId, runMode: task.runMode, variableSetId: task.variableSetId, rigId: task.rigId, agentConfig: task.agentConfig, missingTargetStatus: 404, }); await validateScheduledTaskMachineTarget({ settings: catalogSettings, db, grant, runMode: task.runMode, agentConfig: task.agentConfig, requireOnline: true, }); await requireLimit( { ...deps, settings: catalogSettings }, { accountId: grant.accountId, workspaceId, action: "agent_run:create", quantity: 1, model: task.agentConfig.model ?? catalogSettings.openaiModel, }, ); } // Body is optional (a bare POST is still a valid trigger); only a present, // non-empty body must parse against the contract. const body = await c.req.json().catch(() => ({})); const { triggerId } = TriggerScheduledTaskRequest.parse(body ?? {}); const triggerToken = scheduledTaskTriggerToken(triggerId); const agentRunUsageIdempotencyKey = task.action.kind === "agent_turn" ? manualScheduledTaskTriggerUsageKey(workspaceId, task.id, triggerToken) : `knowledge-source-sync:manual:${workspaceId}:${task.id}:${triggerToken}`; const triggerWorkflowId = manualScheduledTaskTriggerWorkflowId(task.id, triggerToken); await workflowClient.triggerScheduledTask({ task, agentRunUsageIdempotencyKey, triggerWorkflowId, initiator: { kind: "subject", subjectId: grant.subjectId }, }); if (task.action.kind === "agent_turn") { await recordWorkspaceUsage(deps, { accountId: grant.accountId, workspaceId, subjectId: grant.subjectId, eventType: "agent_run.created", quantity: 1, unit: "run", sourceResourceType: "scheduled_task", sourceResourceId: task.id, idempotencyKey: agentRunUsageIdempotencyKey, }); } return c.json(scheduledTaskForGrant(task, grant), 202); }); app.delete("/v1/workspaces/:workspaceId/scheduled-tasks/:taskId", async (c) => { const workspaceId = c.req.param("workspaceId"); const grant = await requireAccessGrant(c, deps, workspaceId, "scheduled_tasks:manage"); await deleteScheduledTaskWithDurableCleanup(deps, { workspaceId, taskId: c.req.param("taskId"), subjectId: grant.subjectId, }); return c.json({ ok: true }); }); app.get("/v1/workspaces/:workspaceId/scheduled-tasks/:taskId/runs", async (c) => { const workspaceId = c.req.param("workspaceId"); const grant = await requireAccessGrant(c, deps, workspaceId, "scheduled_tasks:run"); const task = await requireScheduledTaskForApi(db, workspaceId, c.req.param("taskId")); const taskRuns = await listScheduledTaskRuns( db, workspaceId, task.id, boundedLimit(c.req.query("limit")), ); return c.json(taskRuns.map((run) => scheduledTaskRunForGrant(run, grant))); }); }