Skip to content
5 changes: 5 additions & 0 deletions server/.env.example
Original file line number Diff line number Diff line change
Expand Up @@ -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
30 changes: 16 additions & 14 deletions server/src/app.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
22 changes: 19 additions & 3 deletions server/src/config/index.ts
Original file line number Diff line number Diff line change
@@ -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;
61 changes: 42 additions & 19 deletions server/src/jobs/monitor.job.ts
Original file line number Diff line number Diff line change
@@ -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<string[]>(STATE_FILE, []));
const MAX_SEEN_IDS = 200;

let seenIds = new Set(readJson<string[]>(STATE_FILE, []).slice(0, MAX_SEEN_IDS));
let running = false;
let timer: ReturnType<typeof setTimeout> | 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<void> {
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;
}
}
31 changes: 31 additions & 0 deletions server/src/utils/shutdown.ts
Original file line number Diff line number Diff line change
@@ -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"));
};
Loading