214 lines
6.6 KiB
TypeScript
214 lines
6.6 KiB
TypeScript
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<Omit<HeartbeatJob, "name">> & {
|
|
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 <T>(
|
|
rpc: string,
|
|
body: Record<string, unknown>,
|
|
userSession: string
|
|
): Promise<T> => {
|
|
const endpoint = buildEndpoint(rpc);
|
|
const headers: Record<string, string> = {
|
|
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<string, unknown> = { 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<void> {
|
|
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 <uid>`).
|
|
*/
|
|
export async function listHeartbeatJobs(
|
|
userSession: string,
|
|
pagination?: { page?: number; pageSize?: number }
|
|
): Promise<{ total: number; actorUserId: string; jobs: HeartbeatJobInfo[] }> {
|
|
const body: Record<string, unknown> = {};
|
|
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);
|
|
}
|