Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
27 changes: 24 additions & 3 deletions backend/src/api/app.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,13 +12,18 @@ import { createServer, Server as HttpServer } from "http";
import swaggerUi from "swagger-ui-express";

import {
createTaskJobHandler,
type DispatchFn,
type PaymentReleaseFn,
} from "../coordinator/coordinator";
import { httpDispatch } from "../coordinator/dispatch";
import type { AgentRegistry } from "../types/agent";
import { getTask } from "../coordinator/taskStore";
import { eventBus } from "../coordinator/eventBus";
import { getTask } from "../coordinator/taskStore";
import { getTaskDb } from "../db/tasks";
import { createPaymentReleaseFn, type StellarReleasePaymentFn } from "../payment";
import { getGlobalJobQueue, JobWorker, type JobQueue } from "../queue";
import { createHeartbeatService, type HeartbeatServiceOptions } from "../services/heartbeat";
import { metricsMiddleware, metricsService } from "../services/metrics";
import type { EventStore } from "../events/eventStore";
import {
attachTaskStream,
Expand All @@ -39,10 +44,20 @@ import { rateLimitMiddleware, registerRateLimitMiddleware, publicLimiter, authed
import { authMiddleware } from "./middleware/auth";
import { createCorsMiddleware } from "./middleware/cors";
import { compressionMiddleware } from "./middleware/compression";
import { createCorsMiddleware } from "./middleware/cors";
import { errorHandler } from "./middleware/errorHandler";
import { readOnlyMiddleware } from "./middleware/readOnly";
import { registerRateLimitMiddleware } from "./middleware/rateLimit";
import { requestId } from "./middleware/requestId";
import { requestLogger } from "./middleware/requestLogger";
import { errorHandler } from "./middleware/errorHandler";
import { versioningMiddleware } from "./middleware/versioning";
import { getOpenapiJson, getOpenapiYaml, openapiSpec, swaggerUiOptions } from "./docs";
import { agentsRouter } from "./routes/agents";
import { createAdminRouter } from "./routes/admin";
import { healthRouter } from "./routes/health";
import { createReconciliationRouter, type ReconciliationRouterOptions } from "./routes/reconciliation";
import { createStatsRouter } from "./routes/stats";
import { attachTaskStream, getStreamConnectionCount, type TaskStreamOptions } from "./routes/stream";
import { createV1TasksRouter } from "./routes/v1/tasks";
import { createV2TasksRouter } from "./routes/v2/tasks";
import { createAuthRouter } from "./routes/auth";
Expand Down Expand Up @@ -114,6 +129,11 @@ export function createApp(opts: AppOptions = {}): {
app.use(requestLogger);
app.use(metricsMiddleware);
app.use(versioningMiddleware);
app.use(
readOnlyMiddleware({
exemptPaths: ["/api/admin", "/api/reconciliation"],
}),
);

if (!opts.disableCompression && config.NODE_ENV !== "test") {
app.use(...compressionMiddleware());
Expand Down Expand Up @@ -188,6 +208,7 @@ export function createApp(opts: AppOptions = {}): {
} else {
return v2TasksRouter(req, res, next);
}
return v2TasksRouter(req, res, next);
});

// ── Admin Queue routes ─────────────────────────────────────────────────────
Expand Down
36 changes: 36 additions & 0 deletions backend/src/api/middleware/readOnly.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
import type { NextFunction, Request, Response } from "express";
import { getReadOnlyState, isReadOnly } from "../../services/adminControl";

const MUTATION_METHODS = new Set(["POST", "PUT", "PATCH", "DELETE"]);

export interface ReadOnlyMiddlewareOptions {
exemptPaths?: string[];
}

export function readOnlyMiddleware(options: ReadOnlyMiddlewareOptions = {}) {
const exemptPaths = options.exemptPaths ?? [];

return (req: Request, res: Response, next: NextFunction): void => {
if (!MUTATION_METHODS.has(req.method)) {
next();
return;
}

if (exemptPaths.some((prefix) => req.path === prefix || req.path.startsWith(`${prefix}/`))) {
next();
return;
}

if (!isReadOnly()) {
next();
return;
}

const state = getReadOnlyState();
res.status(503).json({
error: "READ_ONLY",
message: "Mutations are temporarily disabled by an operator.",
readOnly: state,
});
};
}
201 changes: 200 additions & 1 deletion backend/src/api/routes/admin.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,204 @@
import { Router, Request, Response } from "express";
import { Router, type NextFunction, type Request, type Response } from "express";
import { z } from "zod";
import { getGlobalJobQueue, type JobQueue, type JobStatus } from "../../queue";
import {
actorFromRequest,
auditLogToCsv,
backupDatabases,
getReadOnlyState,
listAdminAuditLog,
listAgentsForAdmin,
recordAdminAudit,
setAgentEnabled,
setReadOnlyState,
vacuumDatabases,
} from "../../services/adminControl";
import {
ReconciliationService,
createDefaultReconciliationService,
} from "../../services/reconciliation";
import type { ReconciliationRouterOptions } from "./reconciliation";
import type { ReconciliationTrigger } from "../../services/reconciliation.types";
import { createLogger } from "../../utils/logger";

const logger = createLogger({ module: "admin" });

const readOnlySchema = z.object({
enabled: z.boolean(),
reason: z.string().max(500).optional(),
});

const agentListSchema = z.object({
status: z.enum(["online", "offline"]).optional(),
});

const auditLogQuerySchema = z.object({
limit: z.coerce.number().int().min(1).max(1000).default(200),
offset: z.coerce.number().int().min(0).default(0),
format: z.enum(["json", "csv"]).default("json"),
});

const reconciliationSchema = z.object({
triggeredBy: z.enum(["manual", "scheduled", "release"]).default("manual"),
});

const backupSchema = z.object({
directory: z.string().min(1).optional(),
});

export interface AdminRouterOptions {
queue?: JobQueue;
reconciliation?: ReconciliationRouterOptions;
}

function asyncHandler(
handler: (req: Request, res: Response, next: NextFunction) => Promise<void>,
) {
return (req: Request, res: Response, next: NextFunction): void => {
void handler(req, res, next).catch(next);
};
}

function auditAdminRequests(req: Request, res: Response, next: NextFunction): void {
res.on("finish", () => {
recordAdminAudit({
at: new Date().toISOString(),
actor: actorFromRequest(req),
action: `${req.method} ${req.baseUrl}${req.path}`,
target: req.params.id,
statusCode: res.statusCode,
requestId:
(res.locals.requestId as string | undefined) ??
(res.locals.correlationId as string | undefined),
details: {
params: req.params,
query: req.query,
body: req.method === "GET" ? undefined : req.body,
},
});
});
next();
}

function getReconciliationService(options?: ReconciliationRouterOptions): ReconciliationService {
return options?.service ?? createDefaultReconciliationService();
}

export function createAdminRouter(options: AdminRouterOptions = {}): Router {
const router = Router();
const jobQueue = options.queue ?? getGlobalJobQueue();
const reconciliationService = getReconciliationService(options.reconciliation);

router.use(auditAdminRequests);

router.get("/read-only", (_req: Request, res: Response) => {
res.json(getReadOnlyState());
});

router.put("/read-only", (req: Request, res: Response) => {
const parsed = readOnlySchema.safeParse(req.body);
if (!parsed.success) {
res.status(400).json({ error: "INVALID_BODY", details: parsed.error.flatten() });
return;
}

const state = setReadOnlyState(
parsed.data.enabled,
actorFromRequest(req),
parsed.data.reason,
);
res.json(state);
});

router.get("/agents", (req: Request, res: Response) => {
const parsed = agentListSchema.safeParse(req.query);
if (!parsed.success) {
res.status(400).json({ error: "INVALID_QUERY", details: parsed.error.flatten() });
return;
}
res.json({ agents: listAgentsForAdmin(parsed.data.status) });
});

router.post("/agents/:id/enable", (req: Request, res: Response) => {
const agent = setAgentEnabled(req.params.id, true);
if (!agent) {
res.status(404).json({ error: "AGENT_NOT_FOUND" });
return;
}
res.json({ enabled: true, agent });
});

router.post("/agents/:id/disable", (req: Request, res: Response) => {
const agent = setAgentEnabled(req.params.id, false);
if (!agent) {
res.status(404).json({ error: "AGENT_NOT_FOUND" });
return;
}
res.json({ enabled: false, agent });
});

router.post(
"/reconciliation/run",
asyncHandler(async (req: Request, res: Response) => {
const parsed = reconciliationSchema.safeParse(req.body ?? {});
if (!parsed.success) {
res.status(400).json({ error: "INVALID_BODY", details: parsed.error.flatten() });
return;
}

const triggeredBy = parsed.data.triggeredBy as ReconciliationTrigger;
const report = await reconciliationService.run(triggeredBy);
res.status(200).json(report);
}),
);

router.post("/maintenance/vacuum", (_req: Request, res: Response) => {
res.json({ results: vacuumDatabases() });
});

router.post(
"/maintenance/backup",
asyncHandler(async (req: Request, res: Response) => {
const parsed = backupSchema.safeParse(req.body ?? {});
if (!parsed.success) {
res.status(400).json({ error: "INVALID_BODY", details: parsed.error.flatten() });
return;
}

const results = await backupDatabases(parsed.data.directory);
res.json({ results });
}),
);

router.get("/audit-log", (req: Request, res: Response) => {
const parsed = auditLogQuerySchema.safeParse(req.query);
if (!parsed.success) {
res.status(400).json({ error: "INVALID_QUERY", details: parsed.error.flatten() });
return;
}

const entries = listAdminAuditLog(parsed.data.limit, parsed.data.offset);
if (parsed.data.format === "csv") {
res.setHeader("Content-Type", "text/csv; charset=utf-8");
res.send(auditLogToCsv(entries));
return;
}
res.json({ entries });
});

router.use("/queue", createAdminQueueRouter(jobQueue));
router.use("/", createAdminQueueRouter(jobQueue));

router.use((err: Error, _req: Request, res: Response, _next: NextFunction) => {
logger.error({ err }, "admin operation failed");
res.status(500).json({
error: "ADMIN_OPERATION_FAILED",
message: err.message,
});
});

return router;
}

export function createAdminQueueRouter(queue?: JobQueue): Router {
const router = Router();
Expand Down
Loading
Loading