import { TRPCError } from "@trpc/server"; import { ENV } from "./env"; export type HeartbeatJob = { name: string; /** * 6-field cron with seconds (`sec min hour dom mon dow`), UTC, min interval 60s. * Use `0` for the seconds field — e.g. `"0 0 9 * * *"` is daily 09:00 UTC. * See periodic-updates.md. */ cron: string; /** Callback path. MUST start with `/api/scheduled/`. */ path: string; method?: "POST" | "PUT"; payload?: unknown; description?: string; }; /** * Update patch. All fields optional; unset = leave unchanged. * `enable`: true = resume, false = pause; omit = unchanged. * `name` is the (project, owner)-scope key and cannot be changed. */ export type HeartbeatJobUpdate = Partial> & { enable?: boolean; }; export type HeartbeatJobInfo = { taskUid: string; name: string; userId: string; description: string; cronExpression: string; callbackPath: string; callbackMethod: string; callbackPayload: string; isEnable: boolean; createdAt?: string | null; lastExecutedAt?: string | null; nextExecutionAt?: string | null; }; const SERVICE = "webdevtoken.v1.WebDevService"; const buildEndpoint = (rpc: string): string => { if (!ENV.forgeApiUrl) { throw new TRPCError({ code: "INTERNAL_SERVER_ERROR", message: "Heartbeat service URL is not configured (BUILT_IN_FORGE_API_URL).", }); } if (!ENV.forgeApiKey) { throw new TRPCError({ code: "INTERNAL_SERVER_ERROR", message: "Heartbeat service API key is not configured (BUILT_IN_FORGE_API_KEY).", }); } const baseUrl = ENV.forgeApiUrl; const normalizedBase = baseUrl.endsWith("/") ? baseUrl : `${baseUrl}/`; return new URL(`${SERVICE}/${rpc}`, normalizedBase).toString(); }; const callForge = async ( rpc: string, body: Record, userSession: string ): Promise => { const endpoint = buildEndpoint(rpc); const headers: Record = { accept: "application/json", authorization: `Bearer ${ENV.forgeApiKey}`, "content-type": "application/json", "connect-protocol-version": "1", }; // userSession is the decoded `app_session_id` cookie value (NOT the raw // Cookie header). Empty string falls back to the project owner identity. if (userSession) { headers["x-manus-user-session"] = userSession; } let response: Response; try { response = await fetch(endpoint, { method: "POST", headers, body: JSON.stringify(body), }); } catch (error) { throw new TRPCError({ code: "INTERNAL_SERVER_ERROR", message: `Heartbeat ${rpc} network error: ${String(error)}`, }); } if (!response.ok) { const detail = await response.text().catch(() => ""); throw mapForgeError(response, detail, rpc); } return (await response.json()) as T; }; const mapForgeError = ( response: Response, detail: string, rpc: string ): TRPCError => { const status = response.status; let code: TRPCError["code"] = "INTERNAL_SERVER_ERROR"; if (status === 401) code = "UNAUTHORIZED"; else if (status === 403) code = "FORBIDDEN"; else if (status === 404) code = "NOT_FOUND"; else if (status === 400 || status === 422) code = "BAD_REQUEST"; else if (status === 409) code = "CONFLICT"; else if (status === 429) code = "TOO_MANY_REQUESTS"; return new TRPCError({ code, message: `Heartbeat ${rpc} failed (${status})${detail ? `: ${detail}` : ""}`, }); }; const stringifyPayload = (payload: unknown): string => { if (payload === undefined || payload === null) return "{}"; if (typeof payload === "string") return payload; return JSON.stringify(payload); }; const validateCallbackPath = (path: string): void => { if (!path || !path.startsWith("/api/scheduled/")) { throw new TRPCError({ code: "BAD_REQUEST", message: "callback path must start with /api/scheduled/", }); } }; /** * Create a new HTTP cron job. Returns the assigned `taskUid` to persist on * your business row so callbacks can dereference it. */ export async function createHeartbeatJob( job: HeartbeatJob, userSession: string ): Promise<{ taskUid: string; nextExecutionAt?: string | null }> { validateCallbackPath(job.path); return callForge<{ taskUid: string; nextExecutionAt?: string | null }>( "CreateHeartbeatJob", { name: job.name, cronExpression: job.cron, callbackPath: job.path, callbackMethod: job.method ?? "POST", callbackPayload: stringifyPayload(job.payload), description: job.description ?? "", }, userSession ); } /** * Update an existing cron located by `taskUid`. Only fields you pass in * `patch` are mutated. `enable` flips resume/pause; omit to leave alone. */ export async function updateHeartbeatJob( taskUid: string, patch: HeartbeatJobUpdate, userSession: string ): Promise<{ nextExecutionAt?: string | null }> { if (patch.path !== undefined) validateCallbackPath(patch.path); const body: Record = { taskUid }; if (patch.cron !== undefined) body.cronExpression = patch.cron; if (patch.path !== undefined) body.callbackPath = patch.path; if (patch.method !== undefined) body.callbackMethod = patch.method; if (patch.payload !== undefined) { body.callbackPayload = stringifyPayload(patch.payload); } if (patch.description !== undefined) body.description = patch.description; if (patch.enable !== undefined) body.enable = patch.enable; return callForge<{ nextExecutionAt?: string | null }>( "UpdateHeartbeatJob", body, userSession ); } /** Delete a cron located by `taskUid`. Idempotent on caller side. */ export async function deleteHeartbeatJob( taskUid: string, userSession: string ): Promise { await callForge("DeleteHeartbeatJob", { taskUid }, userSession); } /** * List cron jobs owned by the resolved actor (end-user when `userSession` * is set, project owner otherwise) within the current project. * * `actorUserId` in the response echoes whose cron list you got back. End-users * cannot list other users' crons via this SDK; cross-user inspection is * owner-only via the sandbox CLI (`manus-heartbeat list --user-id `). */ export async function listHeartbeatJobs( userSession: string, pagination?: { page?: number; pageSize?: number } ): Promise<{ total: number; actorUserId: string; jobs: HeartbeatJobInfo[] }> { const body: Record = {}; if (pagination?.page !== undefined) body.page = pagination.page; if (pagination?.pageSize !== undefined) body.pageSize = pagination.pageSize; return callForge<{ total: number; actorUserId: string; jobs: HeartbeatJobInfo[]; }>("ListHeartbeatJobs", body, userSession); }