diff --git a/server/.env.example b/server/.env.example index 43a1c73..9f1beb4 100644 --- a/server/.env.example +++ b/server/.env.example @@ -3,5 +3,10 @@ ADMIN_SECRET=replace-with-a-long-random-value TELEGRAM_BOT_TOKEN= FIREBASE_CONFIG= ALLOWED_ORIGINS=https://almaquake-production.up.railway.app +# Optional tuning; safe defaults are used when omitted. +QUAKE_RADIUS_KM=300 +QUAKE_MIN_MAGNITUDE=2.5 +QUAKE_LIMIT=100 +POLL_INTERVAL_MS=60000 # Railway: attach a Volume mounted at /data, then set this variable. DATA_DIR=/data diff --git a/server/src/app.ts b/server/src/app.ts index cd097ff..d1260fe 100644 --- a/server/src/app.ts +++ b/server/src/app.ts @@ -2,44 +2,46 @@ import express from "express"; import cors from "cors"; import { apiRouter } from "./routes"; import { errorHandlerMiddleware } from "./middlewares/errorHandler.middleware"; -import { startMonitorJob } from "./jobs/monitor.job"; +import { startMonitorJob, stopMonitorJob } from "./jobs/monitor.job"; import { startTelegramBot } from "./services/telegram.service"; import { config } from "./config"; import { logStorageConfiguration } from "./services/storage.service"; import { logger } from "./utils/logger.utils"; +import { registerGracefulShutdown } from "./utils/shutdown"; const app = express(); +const isProduction = process.env.NODE_ENV === "production"; const ALLOWED_ORIGINS = process.env.ALLOWED_ORIGINS - ? process.env.ALLOWED_ORIGINS.split(",") - : ["https://almaquake-production.up.railway.app", "exp://", "http://localhost"]; + ? process.env.ALLOWED_ORIGINS.split(",").map((origin) => origin.trim()).filter(Boolean) + : ["https://almaquake-production.up.railway.app"]; + +const isAllowedOrigin = (origin?: string): boolean => { + if (!origin) return true; + if (!isProduction && origin.startsWith("exp://")) return true; + return ALLOWED_ORIGINS.includes(origin); +}; app.use( cors({ origin: (origin, cb) => { - // allow mobile apps (no origin) and listed origins - if (!origin || ALLOWED_ORIGINS.some((o) => origin.startsWith(o))) { - cb(null, true); - } else { - cb(null, false); - } + if (isAllowedOrigin(origin)) cb(null, true); + else cb(null, false); }, }) ); app.use(express.json({ limit: "10kb" })); - -// Mount central API router app.use("/api", apiRouter); - -// Centralized Error Handling Middleware app.use(errorHandlerMiddleware); -app.listen(config.port, () => { +const server = app.listen(config.port, () => { logStorageConfiguration(); logger.info(`🟢 Server started on port ${config.port}`); startMonitorJob(); startTelegramBot(); }); +registerGracefulShutdown(server, stopMonitorJob); + export default app; diff --git a/server/src/config/index.ts b/server/src/config/index.ts index 31d942a..170f553 100644 --- a/server/src/config/index.ts +++ b/server/src/config/index.ts @@ -1,6 +1,22 @@ +const parsePositiveInt = (value: string | undefined, fallback: number): number => { + const parsed = Number(value); + return Number.isInteger(parsed) && parsed > 0 ? parsed : fallback; +}; + +const parseNonNegativeNumber = (value: string | undefined, fallback: number): number => { + const parsed = Number(value); + return Number.isFinite(parsed) && parsed >= 0 ? parsed : fallback; +}; + export const config = { almaty: { lat: 43.2565, lng: 76.9286 }, - quake: { radiusKm: 300, minMagnitude: 2.5, limit: 100 }, - poll: { intervalMs: 60_000 }, - port: Number(process.env.PORT) || 3000, + quake: { + radiusKm: parsePositiveInt(process.env.QUAKE_RADIUS_KM, 300), + minMagnitude: parseNonNegativeNumber(process.env.QUAKE_MIN_MAGNITUDE, 2.5), + limit: parsePositiveInt(process.env.QUAKE_LIMIT, 100), + }, + poll: { + intervalMs: parsePositiveInt(process.env.POLL_INTERVAL_MS, 60_000), + }, + port: parsePositiveInt(process.env.PORT, 3000), } as const; diff --git a/server/src/jobs/monitor.job.ts b/server/src/jobs/monitor.job.ts index aacca6b..366138a 100644 --- a/server/src/jobs/monitor.job.ts +++ b/server/src/jobs/monitor.job.ts @@ -1,44 +1,67 @@ -import { fetchQuakes } from "../services/usgs.service"; -import { sendQuakeAlert } from "../services/telegram.service"; -import { config } from "../config"; +import { fetchQuakes } from "../services/usgs.service"; +import { sendQuakeAlert } from "../services/telegram.service"; +import { config } from "../config"; import { dataFile, readJson, writeJson } from "../services/storage.service"; +import { logger } from "../utils/logger.utils"; const STATE_FILE = dataFile("monitor-state.json"); -let seenIds = new Set(readJson(STATE_FILE, [])); +const MAX_SEEN_IDS = 200; + +let seenIds = new Set(readJson(STATE_FILE, []).slice(0, MAX_SEEN_IDS)); let running = false; +let timer: ReturnType | undefined; +let stopped = false; + +function rememberLatest(quakeIds: string[]): void { + seenIds = new Set(quakeIds.slice(0, MAX_SEEN_IDS)); + writeJson(STATE_FILE, [...seenIds]); +} export async function runMonitor(): Promise { - if (running) return; + if (running || stopped) return; running = true; try { const quakes = await fetchQuakes(); - if (quakes.length === 0) { - return; - } + if (quakes.length === 0) return; if (seenIds.size === 0) { - seenIds = new Set(quakes.map((quake) => quake.id)); - writeJson(STATE_FILE, [...seenIds]); - console.log(`📌 Монитор запущен. Последнее событие: ${quakes[0].id} (M${quakes[0].magnitude})`); + rememberLatest(quakes.map((quake) => quake.id)); + logger.info(`Monitor initialized with latest event ${quakes[0].id} (M${quakes[0].magnitude})`); return; } const newQuakes = quakes.filter((quake) => !seenIds.has(quake.id)).reverse(); for (const quake of newQuakes) { - console.log(`🌍 ОБНАРУЖЕНО НОВОЕ СОБЫТИЕ: M${quake.magnitude} — ${quake.place}`); + logger.info(`New earthquake detected: M${quake.magnitude} — ${quake.place}`); await sendQuakeAlert(quake); } - seenIds = new Set(quakes.map((quake) => quake.id)); - writeJson(STATE_FILE, [...seenIds]); - } catch (error: any) { - console.error("⚠️ Сбой мониторинга (повтор через 60с):", error.message || error); + + rememberLatest(quakes.map((quake) => quake.id)); + } catch (error: unknown) { + const message = error instanceof Error ? error.message : String(error); + logger.error(`Monitoring failure; retrying after ${config.poll.intervalMs}ms: ${message}`); } finally { running = false; } } +function scheduleNextRun(): void { + if (stopped) return; + timer = setTimeout(async () => { + await runMonitor(); + scheduleNextRun(); + }, config.poll.intervalMs); +} + export function startMonitorJob(): void { - // Запуск первой проверки сразу - runMonitor(); - setInterval(runMonitor, config.poll.intervalMs); + stopped = false; + void runMonitor().finally(scheduleNextRun); +} + +export function stopMonitorJob(): void { + stopped = true; + if (timer) { + clearTimeout(timer); + timer = undefined; + } } diff --git a/server/src/utils/shutdown.ts b/server/src/utils/shutdown.ts new file mode 100644 index 0000000..093d602 --- /dev/null +++ b/server/src/utils/shutdown.ts @@ -0,0 +1,31 @@ +import type { Server } from "node:http"; +import { logger } from "./logger.utils"; + +export const registerGracefulShutdown = ( + server: Server, + cleanup: () => void = () => undefined +): void => { + let shuttingDown = false; + + const shutdown = (signal: string): void => { + if (shuttingDown) return; + shuttingDown = true; + cleanup(); + + logger.info(`Received ${signal}; shutting down gracefully`); + + server.close((error) => { + if (error) { + logger.error("HTTP server shutdown failed", error); + process.exitCode = 1; + return; + } + + logger.info("HTTP server closed"); + process.exitCode = 0; + }); + }; + + process.once("SIGTERM", () => shutdown("SIGTERM")); + process.once("SIGINT", () => shutdown("SIGINT")); +};