Checkpoint: Migration complète vers authentification locale (JWT + bcrypt) : schéma DB (10 tables), routes tRPC (auth, users, etablissements, inventaire, capex, opex, parametres), page de login avec branding Itinova/Santinova, protection des routes, onglets Établissements et Utilisateurs branchés sur tRPC, import inventaire/établissements/utilisateurs via tRPC, 14 tests Vitest passants.
This commit is contained in:
213
server/_core/heartbeat.ts
Normal file
213
server/_core/heartbeat.ts
Normal file
@@ -0,0 +1,213 @@
|
||||
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);
|
||||
}
|
||||
Reference in New Issue
Block a user