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
44 changes: 44 additions & 0 deletions backend/src/api/middleware/drainProtection.middleware.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,44 @@
import type { FastifyInstance, FastifyRequest, FastifyReply } from "fastify";
import { drainProtocolService } from "../../services/drainProtocol.service.js";

const EXEMPT_PATHS = [
"/health",
"/healthz",
"/readyz",
"/api/v1/health",
"/api/v1/admin/shutdown/drain/start",
"/api/v1/admin/shutdown/drain/status",
"/api/v1/admin/shutdown/drain/cancel",
"/api/v1/admin/shutdown/drain/force",
"/api/v1/admin/shutdown/drain/history",
];

export async function registerDrainProtectionMiddleware(server: FastifyInstance): Promise<void> {
// Track in-flight request lifecycle
server.addHook("onRequest", async (request: FastifyRequest, reply: FastifyReply) => {
drainProtocolService.incrementInFlight();

reply.raw.on("finish", () => {
drainProtocolService.decrementInFlight();
});

const isExempt = EXEMPT_PATHS.some((path) => request.url.startsWith(path));
if (isExempt) {
return;
}

if (drainProtocolService.isDraining()) {
const isMutatingMethod = ["POST", "PUT", "DELETE", "PATCH"].includes(request.method.toUpperCase());

if (isMutatingMethod || drainProtocolService.getMode() === "force") {
reply.header("Retry-After", "30");
return reply.status(503).send({
error: "Service Unavailable",
message: "Server is currently undergoing graceful shutdown drain",
state: drainProtocolService.getState(),
retryAfterSeconds: 30,
});
}
}
});
}
123 changes: 123 additions & 0 deletions backend/src/api/routes/__tests__/drainProtocol.routes.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,123 @@
import { describe, it, expect, beforeEach, vi } from "vitest";
import Fastify from "fastify";
import { drainProtocolRoutes } from "../drainProtocol.routes.js";
import { drainProtocolService } from "../../../services/drainProtocol.service.js";

vi.mock("../../../services/drainProtocol.service.js", () => {
const mockStatus = {
sessionId: "mock-session-id",
nodeId: "node-1",
state: "ACTIVE",
drainMode: "graceful",
inFlightRequests: 0,
activeConnections: 0,
activeStreams: 0,
startedAt: null,
drainedAt: null,
reason: null,
initiatedBy: null,
timeoutSeconds: 30,
};

return {
drainProtocolService: {
getStatus: vi.fn().mockImplementation(() => mockStatus),
startDrain: vi.fn().mockImplementation(async (opts) => ({
...mockStatus,
state: "DRAINED",
reason: opts?.reason || "Graceful drain",
initiatedBy: opts?.initiatedBy || "admin",
})),
cancelDrain: vi.fn().mockImplementation(async (by) => ({
...mockStatus,
state: "ACTIVE",
initiatedBy: by,
})),
forceShutdown: vi.fn().mockImplementation(async (reason) => ({
...mockStatus,
state: "FAILED",
reason,
})),
getDrainHistory: vi.fn().mockResolvedValue([mockStatus]),
},
};
});

vi.mock("../../middleware/auth.js", () => ({
authMiddleware: () => async () => {},
}));

describe("drainProtocolRoutes", () => {
let app: ReturnType<typeof Fastify>;

beforeEach(async () => {
app = Fastify();
await app.register(drainProtocolRoutes);
vi.clearAllMocks();
});

it("GET /status returns current drain status", async () => {
const res = await app.inject({
method: "GET",
url: "/status",
});

expect(res.statusCode).toBe(200);
const body = res.json();
expect(body.nodeId).toBe("node-1");
expect(body.state).toBe("ACTIVE");
});

it("POST /start initiates drain and returns HTTP 202 Accepted", async () => {
const res = await app.inject({
method: "POST",
url: "/start",
payload: {
timeoutSeconds: 15,
reason: "Maintenance shutdown",
},
});

expect(res.statusCode).toBe(202);
expect(drainProtocolService.startDrain).toHaveBeenCalledWith(
expect.objectContaining({
timeoutSeconds: 15,
reason: "Maintenance shutdown",
})
);
});

it("POST /cancel cancels active drain protocol", async () => {
const res = await app.inject({
method: "POST",
url: "/cancel",
payload: { cancelledBy: "operator-admin" },
});

expect(res.statusCode).toBe(200);
expect(drainProtocolService.cancelDrain).toHaveBeenCalledWith("operator-admin");
});

it("POST /force triggers force shutdown", async () => {
const res = await app.inject({
method: "POST",
url: "/force",
payload: { reason: "Urgent abort" },
});

expect(res.statusCode).toBe(200);
expect(drainProtocolService.forceShutdown).toHaveBeenCalledWith("Urgent abort");
});

it("GET /history returns past drain sessions", async () => {
const res = await app.inject({
method: "GET",
url: "/history?limit=5",
});

expect(res.statusCode).toBe(200);
const body = res.json();
expect(body.history).toHaveLength(1);
expect(body.count).toBe(1);
});
});
66 changes: 66 additions & 0 deletions backend/src/api/routes/drainProtocol.routes.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,66 @@
import type { FastifyInstance } from "fastify";
import { z } from "zod";
import { drainProtocolService } from "../../services/drainProtocol.service.js";
import { authMiddleware } from "../middleware/auth.js";
import { sendApiError } from "../utils/response.js";

const startDrainSchema = z.object({
nodeId: z.string().optional(),
timeoutSeconds: z.number().int().min(1).max(300).optional(),
reason: z.string().optional(),
initiatedBy: z.string().optional(),
mode: z.enum(["graceful", "force", "read_only"]).optional(),
metadata: z.record(z.unknown()).optional(),
});

const cancelDrainSchema = z.object({
cancelledBy: z.string().optional(),
});

const forceDrainSchema = z.object({
reason: z.string().optional(),
});

export async function drainProtocolRoutes(server: FastifyInstance): Promise<void> {
const adminAuth = authMiddleware({ requiredScopes: ["admin:write"] });

// Start graceful shutdown drain protocol
server.post("/start", { preHandler: [adminAuth] }, async (request, reply) => {
const parsed = startDrainSchema.safeParse(request.body || {});
if (!parsed.success) {
return sendApiError(reply, 400, "Invalid drain options", { issues: parsed.error.errors });
}

const status = await drainProtocolService.startDrain(parsed.data);
return reply.status(202).send(status);
});

// Get current drain status
server.get("/status", async (_request, reply) => {
const status = drainProtocolService.getStatus();
return reply.status(200).send(status);
});

// Cancel drain and resume normal operation
server.post("/cancel", { preHandler: [adminAuth] }, async (request, reply) => {
const parsed = cancelDrainSchema.safeParse(request.body || {});
const cancelledBy = parsed.success ? parsed.data.cancelledBy || "admin" : "admin";
const status = await drainProtocolService.cancelDrain(cancelledBy);
return reply.status(200).send(status);
});

// Force immediate shutdown
server.post("/force", { preHandler: [adminAuth] }, async (request, reply) => {
const parsed = forceDrainSchema.safeParse(request.body || {});
const reason = parsed.success ? parsed.data.reason || "Force shutdown requested" : "Force shutdown requested";
const status = await drainProtocolService.forceShutdown(reason);
return reply.status(200).send(status);
});

// Get drain history
server.get("/history", { preHandler: [adminAuth] }, async (request, reply) => {
const limit = Number((request.query as { limit?: string })?.limit) || 20;
const history = await drainProtocolService.getDrainHistory(limit);
return reply.status(200).send({ history, count: history.length });
});
}
6 changes: 6 additions & 0 deletions backend/src/api/routes/route-groups/admin-routes.ts
Original file line number Diff line number Diff line change
Expand Up @@ -121,5 +121,11 @@ export async function registerAdminRoutes(server: FastifyInstance): Promise<void
server.register(parseQuarantineQueueRoutes, {
prefix: "/api/v1/admin/quarantine",
});

// #1187 — Graceful Shutdown Drain Protocol
const { drainProtocolRoutes } = await import("../drainProtocol.routes.js");
server.register(drainProtocolRoutes, {
prefix: "/api/v1/admin/shutdown/drain",
});
}

Original file line number Diff line number Diff line change
@@ -0,0 +1,45 @@
import type { Knex } from "knex";

/**
* Migration for Graceful Shutdown Drain Protocol (#1187).
* Creates shutdown_drain_sessions and shutdown_drain_logs tables.
*/
export async function up(knex: Knex): Promise<void> {
await knex.schema.createTable("shutdown_drain_sessions", (t) => {
t.uuid("id").primary().defaultTo(knex.raw("gen_random_uuid()"));
t.string("node_id", 128).notNullable();
t.string("state", 32).notNullable().defaultTo("ACTIVE"); // ACTIVE, DRAINING, DRAINED, CANCELLED, FAILED
t.string("drain_mode", 32).notNullable().defaultTo("graceful"); // graceful, force, read_only
t.string("reason", 500).nullable();
t.string("initiated_by", 128).notNullable().defaultTo("system");
t.integer("timeout_seconds").notNullable().defaultTo(30);
t.integer("pending_jobs_count").notNullable().defaultTo(0);
t.integer("active_connections_count").notNullable().defaultTo(0);
t.integer("active_streams_count").notNullable().defaultTo(0);
t.timestamp("started_at", { useTz: true }).notNullable().defaultTo(knex.fn.now());
t.timestamp("drained_at", { useTz: true }).nullable();
t.timestamp("cancelled_at", { useTz: true }).nullable();
t.jsonb("metadata").notNullable().defaultTo("{}");
t.timestamp("created_at", { useTz: true }).notNullable().defaultTo(knex.fn.now());
t.timestamp("updated_at", { useTz: true }).notNullable().defaultTo(knex.fn.now());

t.index(["node_id", "state"], "idx_shutdown_drain_node_state");
t.index(["created_at"], "idx_shutdown_drain_created");
});

await knex.schema.createTable("shutdown_drain_logs", (t) => {
t.uuid("id").primary().defaultTo(knex.raw("gen_random_uuid()"));
t.uuid("session_id").notNullable().references("id").inTable("shutdown_drain_sessions").onDelete("CASCADE");
t.string("event_type", 64).notNullable(); // DRAIN_INITIATED, JOBS_PAUSED, WS_DRAINED, STREAMS_STOPPED, DRAIN_COMPLETED, DRAIN_CANCELLED, DRAIN_FAILED, FORCE_SHUTDOWN
t.string("message", 500).notNullable();
t.jsonb("details").notNullable().defaultTo("{}");
t.timestamp("timestamp", { useTz: true }).notNullable().defaultTo(knex.fn.now());

t.index(["session_id", "timestamp"], "idx_shutdown_drain_logs_session");
});
}

export async function down(knex: Knex): Promise<void> {
await knex.schema.dropTableIfExists("shutdown_drain_logs");
await knex.schema.dropTableIfExists("shutdown_drain_sessions");
}
8 changes: 7 additions & 1 deletion backend/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,8 @@ import { registerRequestLoggingMiddleware } from "./api/middleware/logging.middl
import { registerTracing } from "./api/middleware/tracing.js";
import { getTelegramBotService } from "./services/telegram.bot.service.js";
import { registerCompatibilityMiddleware } from "./api/compatibility/middleware.js";
import { registerDrainProtectionMiddleware } from "./api/middleware/drainProtection.middleware.js";
import { drainProtocolService } from "./services/drainProtocol.service.js";

export async function buildServer() {
const server = Fastify({
Expand Down Expand Up @@ -75,6 +77,9 @@ export async function buildServer() {
// Register metrics middleware (to capture all requests)
await registerMetrics(server as any);

// Register graceful shutdown drain protection middleware
await registerDrainProtectionMiddleware(server as any);

// Register plugins
await server.register(cors, {
origin: (origin, callback) => {
Expand Down Expand Up @@ -181,7 +186,8 @@ async function start() {

// ─── Graceful shutdown ──────────────────────────────────────────────────────
const shutdown = async (signal: string) => {
logger.info({ signal }, "Shutdown signal received");
logger.info({ signal }, "Shutdown signal received; initiating drain protocol");
await drainProtocolService.startDrain({ reason: `Received signal ${signal}`, initiatedBy: "system" });

// Stop Telegram bot service
const telegramService = getTelegramBotService();
Expand Down
Loading
Loading