Исправлено дублирование сообщений в task-чате: защита от двойного Enter/клика по кнопке отправки

This commit is contained in:
2026-08-24 16:50:05 +03:00
parent 422eafd41f
commit 8b713e66c0
12 changed files with 1674 additions and 1 deletions

View File

@@ -0,0 +1,232 @@
import { gpsStorage } from "./storage";
import { storage } from "../storage";
import { webPushService } from "../services/web-push.service";
import { eventBus, publishNotificationSSE } from "../routes/shared";
import type { GpsAsset, GpsGeozone } from "@shared/schema";
/**
* Сервис обработки GPS-позиций и геозон.
* Детект входа/выхода — сравнение point-in-polygon с сохранённым состоянием
* gps_geozone_states; при смене состояния создаётся событие и рассылаются
* уведомления подписчикам объекта.
*/
/**
* Ray casting: находится ли точка внутри полигона.
* polygon — массив пар [lat, lng]. Без внешних зависимостей.
*/
export function pointInPolygon(lat: number, lng: number, polygon: Array<[number, number]>): boolean {
if (!Array.isArray(polygon) || polygon.length < 3) return false;
let inside = false;
for (let i = 0, j = polygon.length - 1; i < polygon.length; j = i++) {
const [latI, lngI] = polygon[i];
const [latJ, lngJ] = polygon[j];
const intersects =
lngI > lng !== lngJ > lng &&
lat < ((latJ - latI) * (lng - lngI)) / (lngJ - lngI) + latI;
if (intersects) inside = !inside;
}
return inside;
}
export interface ProcessPositionInput {
lat: number;
lng: number;
speed?: number | null;
course?: number | null;
accuracy?: number | null;
recordedAt?: Date;
}
/**
* Основная точка входа: сохраняет позицию, обновляет last_* объекта,
* проверяет геозоны и рассылает SSE/уведомления.
*/
export async function processPosition(
assetId: number,
organizationId: number,
input: ProcessPositionInput
): Promise<void> {
const recordedAt = input.recordedAt ?? new Date();
// 1. Сохраняем точку в историю
await gpsStorage.insertPosition({
assetId,
organizationId,
lat: input.lat,
lng: input.lng,
speed: input.speed ?? null,
course: input.course ?? null,
accuracy: input.accuracy ?? null,
recordedAt,
});
// 2. Обновляем последнюю позицию объекта (и переводим в online)
await gpsStorage.updateAssetLastPosition(assetId, {
lat: input.lat,
lng: input.lng,
speed: input.speed,
course: input.course,
recordedAt,
});
// 3. SSE-событие о перемещении — broadcast по организации,
// чтобы карта у всех открытых клиентов обновилась в реальном времени
eventBus.publishEvent({
type: "gps_asset_moved",
data: {
assetId,
lat: input.lat,
lng: input.lng,
speed: input.speed ?? null,
course: input.course ?? null,
recordedAt: recordedAt.toISOString(),
},
organizationId,
});
// 4. Проверка геозон
await checkGeozones(assetId, organizationId, input, recordedAt);
}
async function checkGeozones(
assetId: number,
organizationId: number,
input: ProcessPositionInput,
recordedAt: Date
): Promise<void> {
const geozones = await gpsStorage.listActiveGeozones(organizationId);
if (geozones.length === 0) return;
const states = await gpsStorage.getGeozoneStatesForAsset(assetId);
const asset = await gpsStorage.getAsset(assetId, organizationId);
if (!asset) return;
for (const geozone of geozones) {
try {
const nowInside = pointInPolygon(input.lat, input.lng, geozone.polygon);
const wasInside = states.get(geozone.id) ?? false;
if (nowInside === wasInside) continue;
// Состояние изменилось — фиксируем и создаём событие
await gpsStorage.upsertGeozoneState(assetId, geozone.id, nowInside);
const eventType = nowInside ? ("enter" as const) : ("exit" as const);
await gpsStorage.insertGeozoneEvent({
organizationId,
geozoneId: geozone.id,
assetId,
event: eventType,
lat: input.lat,
lng: input.lng,
});
const notifyEnabled = nowInside ? geozone.notifyEnter : geozone.notifyExit;
if (notifyEnabled) {
await notifySubscribers(asset, geozone, eventType, input.lat, input.lng, recordedAt);
}
} catch (err) {
console.error(`[GPS] Geozone check failed (asset ${assetId}, geozone ${geozone.id}):`, err);
}
}
}
/**
* Рассылка уведомлений подписчикам объекта о событии геозоны.
* In-app (user_notifications) + Web Push + SSE для мгновенного обновления бейджа.
*/
async function notifySubscribers(
asset: GpsAsset,
geozone: GpsGeozone,
eventType: "enter" | "exit",
lat: number,
lng: number,
recordedAt: Date
): Promise<void> {
const subscribers = await gpsStorage.listSubscribers(asset.id);
const eventText = eventType === "enter" ? "вошёл в зону" : "вышел из зоны";
const title = `GPS: ${asset.name}`;
const message = `${asset.name} ${eventText} «${geozone.name}»`;
for (const sub of subscribers) {
const enabled = eventType === "enter" ? sub.onEnter : sub.onExit;
if (!enabled) continue;
try {
const notification = await storage.createUserNotification({
userId: sub.userId,
organizationId: asset.organizationId,
type: "gps.geozone",
title,
message,
isRead: false,
});
publishNotificationSSE(sub.userId, asset.organizationId, {
type: "gps.geozone",
notificationId: notification.id,
});
await webPushService.sendToUser(sub.userId, asset.organizationId, {
title,
body: message,
icon: "/icon-192.png",
url: "/gps",
tag: `gps-geozone-${asset.id}-${geozone.id}`,
data: {
type: "gps",
event: eventType,
assetId: asset.id.toString(),
geozoneId: geozone.id.toString(),
lat: lat.toString(),
lng: lng.toString(),
recordedAt: recordedAt.toISOString(),
},
});
} catch (err) {
console.error(`[GPS] Notify subscriber ${sub.userId} failed:`, err);
}
}
}
/**
* Уведомление подписчиков о переходе объекта в офлайн (вызывается воркером).
*/
export async function notifyOffline(asset: GpsAsset): Promise<void> {
const subscribers = await gpsStorage.listSubscribers(asset.id);
const title = `GPS: ${asset.name}`;
const message = `${asset.name} не выходит на связь более 15 минут`;
for (const sub of subscribers) {
if (!sub.onOffline) continue;
try {
const notification = await storage.createUserNotification({
userId: sub.userId,
organizationId: asset.organizationId,
type: "gps.offline",
title,
message,
isRead: false,
});
publishNotificationSSE(sub.userId, asset.organizationId, {
type: "gps.offline",
notificationId: notification.id,
});
await webPushService.sendToUser(sub.userId, asset.organizationId, {
title,
body: message,
icon: "/icon-192.png",
url: "/gps",
tag: `gps-offline-${asset.id}`,
data: { type: "gps", event: "offline", assetId: asset.id.toString() },
});
} catch (err) {
console.error(`[GPS] Notify offline subscriber ${sub.userId} failed:`, err);
}
}
}

545
server/gps/routes.ts Normal file
View File

@@ -0,0 +1,545 @@
import { Router } from "express";
import type { Express, Response } from "express";
import { authenticateToken, type AuthenticatedRequest } from "../middleware/auth.middleware";
import { gpsStorage } from "./storage";
import { processPosition } from "./geozone.service";
import { ensureDevice, deleteDevice } from "./traccar.client";
import { storage } from "../storage";
/**
* Роуты GPS-модуля.
* Публичный POST /api/gps/ingest (приём точек от Traccar, секрет в заголовке)
* регистрируется отдельно, ДО роутера с authenticateToken.
*/
const MAX_TRACK_POINTS = 2000;
function parseDateParam(value: unknown): Date | null {
if (typeof value !== "string" || !value) return null;
const d = new Date(value);
return isNaN(d.getTime()) ? null : d;
}
function parseNum(value: unknown): number | null {
if (value === undefined || value === null || value === "") return null;
const n = Number(value);
return Number.isFinite(n) ? n : null;
}
function badRequest(res: Response, error: string) {
return res.status(400).json({ success: false, error });
}
// =====================
// Приём позиций от Traccar (публичный, по секрету)
// =====================
interface NormalizedPoint {
externalId: string;
lat: number;
lng: number;
speed: number | null;
course: number | null;
accuracy: number | null;
recordedAt: Date | null;
}
/**
* Нормализация одной точки. Поддерживает два формата:
* 1. Простой: { externalId, lat, lng, speed?, course?, accuracy?, timestamp? }
* 2. Нативный forward Traccar: { position: {...}, device: { uniqueId } }
*/
function normalizePoint(raw: any): NormalizedPoint | null {
if (!raw || typeof raw !== "object") return null;
// Нативный формат Traccar event forwarding
if (raw.position && raw.device?.uniqueId) {
const p = raw.position;
const lat = parseNum(p.latitude);
const lng = parseNum(p.longitude);
if (lat === null || lng === null) return null;
return {
externalId: String(raw.device.uniqueId),
lat,
lng,
speed: parseNum(p.speed),
course: parseNum(p.course),
accuracy: parseNum(p.accuracy),
recordedAt: parseDateParam(p.fixTime ?? p.deviceTime),
};
}
// Простой формат
const lat = parseNum(raw.lat);
const lng = parseNum(raw.lng);
if (lat === null || lng === null || !raw.externalId) return null;
return {
externalId: String(raw.externalId),
lat,
lng,
speed: parseNum(raw.speed),
course: parseNum(raw.course),
accuracy: parseNum(raw.accuracy),
recordedAt: parseDateParam(raw.timestamp),
};
}
async function handleIngest(req: AuthenticatedRequest, res: Response) {
// Endpoint отключён, если секрет не задан в env
const expectedToken = process.env.GPS_INGEST_TOKEN;
if (!expectedToken) {
return res.status(503).json({ success: false, error: "GPS ingest отключён: не задан GPS_INGEST_TOKEN" });
}
const token = req.headers["x-gps-token"];
if (token !== expectedToken) {
return res.status(401).json({ success: false, error: "Неверный токен" });
}
const items = Array.isArray(req.body) ? req.body : [req.body];
let processed = 0;
let ignored = 0;
for (const item of items) {
const point = normalizePoint(item);
if (!point) {
ignored++;
continue;
}
// Объект ищется по external_id в любой организации;
// невалидный external_id не раскрываем — просто считаем проигнорированным
const asset = await gpsStorage.getAssetByExternalId(point.externalId);
if (!asset || !asset.isActive) {
ignored++;
continue;
}
try {
await processPosition(asset.id, asset.organizationId, {
lat: point.lat,
lng: point.lng,
speed: point.speed,
course: point.course,
accuracy: point.accuracy,
recordedAt: point.recordedAt ?? undefined,
});
processed++;
} catch (err) {
console.error(`[GPS] Ingest error (asset ${asset.id}):`, err);
ignored++;
}
}
res.json({ success: true, processed, ...(ignored > 0 ? { ignored } : {}) });
}
// =====================
// Авторизованные роуты
// =====================
export function registerGpsRoutes(app: Express) {
// Публичный приём точек — без authenticateToken, защита секретом в заголовке
app.post("/api/gps/ingest", (req, res) => {
handleIngest(req as AuthenticatedRequest, res).catch((err) => {
console.error("[GPS] Ingest handler error:", err);
res.status(500).json({ success: false, error: "Внутренняя ошибка" });
});
});
const router = Router();
router.use(authenticateToken);
// =====================
// Конфигурация клиента (ключ Яндекс.Карт из external_services)
// =====================
router.get("/config", async (req: AuthenticatedRequest, res: Response) => {
const orgId = req.user!.organizationId;
let yandexMapsKey: string | null = null;
const service = await storage.getExternalServiceByType("yandex_maps", orgId);
if (service?.isActive && service.apiKey) {
const { decrypt } = await import("../crypto");
yandexMapsKey = decrypt(service.apiKey);
}
res.json({ success: true, yandexMapsKey });
});
// =====================
// Объекты (assets)
// =====================
router.get("/assets", async (req: AuthenticatedRequest, res: Response) => {
const assets = await gpsStorage.listAssets(req.user!.organizationId);
res.json({ success: true, assets });
});
router.post("/assets", async (req: AuthenticatedRequest, res: Response) => {
const orgId = req.user!.organizationId;
const { name, type, externalId, color, userId } = req.body ?? {};
if (!name || typeof name !== "string") {
return badRequest(res, "Название объекта обязательно");
}
if (type !== undefined && !["vehicle", "person"].includes(type)) {
return badRequest(res, "Тип объекта: vehicle или person");
}
const asset = await gpsStorage.createAsset({
organizationId: orgId,
name: name.trim(),
type: type === "person" ? "person" : "vehicle",
externalId: externalId ? String(externalId).trim() : null,
color: color || undefined,
userId: userId ?? null,
});
// Регистрируем устройство в Traccar (no-op, если Traccar не настроен)
if (asset.externalId) {
await ensureDevice(asset.externalId, asset.name);
}
res.status(201).json({ success: true, asset });
});
router.patch("/assets/:id", async (req: AuthenticatedRequest, res: Response) => {
const orgId = req.user!.organizationId;
const id = Number(req.params.id);
const { name, type, externalId, color, isActive, userId } = req.body ?? {};
const updates: Record<string, unknown> = {};
if (name !== undefined) updates.name = String(name).trim();
if (type !== undefined) {
if (!["vehicle", "person"].includes(type)) {
return badRequest(res, "Тип объекта: vehicle или person");
}
updates.type = type;
}
if (externalId !== undefined) updates.externalId = externalId ? String(externalId).trim() : null;
if (color !== undefined) updates.color = color;
if (isActive !== undefined) updates.isActive = Boolean(isActive);
if (userId !== undefined) updates.userId = userId ?? null;
const asset = await gpsStorage.updateAsset(id, orgId, updates);
if (!asset) return res.status(404).json({ success: false, error: "Объект не найден" });
// Если появился/сменился external_id — регистрируем устройство в Traccar
if (updates.externalId && asset.externalId) {
await ensureDevice(asset.externalId, asset.name);
}
res.json({ success: true, asset });
});
router.delete("/assets/:id", async (req: AuthenticatedRequest, res: Response) => {
const orgId = req.user!.organizationId;
const id = Number(req.params.id);
const asset = await gpsStorage.getAsset(id, orgId);
if (!asset) return res.status(404).json({ success: false, error: "Объект не найден" });
if (asset.externalId) {
await deleteDevice(asset.externalId);
}
await gpsStorage.deleteAsset(id, orgId);
res.json({ success: true });
});
// =====================
// Позиции и треки
// =====================
// Текущие last_* всех активных объектов организации
router.get("/assets/positions", async (req: AuthenticatedRequest, res: Response) => {
const assets = await gpsStorage.listCurrentPositions(req.user!.organizationId);
res.json({
success: true,
positions: assets.map((a) => ({
assetId: a.id,
name: a.name,
type: a.type,
color: a.color,
isOnline: a.isOnline,
lat: a.lastLat,
lng: a.lastLng,
speed: a.lastSpeed,
course: a.lastCourse,
recordedAt: a.lastSeenAt,
})),
});
});
// Трек объекта за период (точки ASC; при > 2000 — прореживание)
router.get("/assets/:id/track", async (req: AuthenticatedRequest, res: Response) => {
const orgId = req.user!.organizationId;
const id = Number(req.params.id);
const asset = await gpsStorage.getAsset(id, orgId);
if (!asset) return res.status(404).json({ success: false, error: "Объект не найден" });
const from = parseDateParam(req.query.from);
const to = parseDateParam(req.query.to);
const all = await gpsStorage.getTrack(id, orgId, from ?? undefined, to ?? undefined);
const total = all.length;
let points = all;
let sampled = false;
if (total > MAX_TRACK_POINTS) {
const step = Math.ceil(total / MAX_TRACK_POINTS);
points = all.filter((_, idx) => idx % step === 0);
// Всегда включаем последнюю точку трека
if (points[points.length - 1] !== all[total - 1]) {
points.push(all[total - 1]);
}
sampled = true;
}
res.json({ success: true, points, total, sampled });
});
// =====================
// Геозоны
// =====================
router.get("/geozones", async (req: AuthenticatedRequest, res: Response) => {
const geozones = await gpsStorage.listGeozones(req.user!.organizationId);
res.json({ success: true, geozones });
});
router.post("/geozones", async (req: AuthenticatedRequest, res: Response) => {
const orgId = req.user!.organizationId;
const { name, polygon, color, notifyEnter, notifyExit } = req.body ?? {};
if (!name || typeof name !== "string") {
return badRequest(res, "Название геозоны обязательно");
}
if (!isValidPolygon(polygon)) {
return badRequest(res, "Полигон: массив минимум из 3 пар [lat, lng]");
}
const geozone = await gpsStorage.createGeozone({
organizationId: orgId,
name: name.trim(),
polygon,
color: color || undefined,
notifyEnter: notifyEnter ?? undefined,
notifyExit: notifyExit ?? undefined,
});
res.status(201).json({ success: true, geozone });
});
router.patch("/geozones/:id", async (req: AuthenticatedRequest, res: Response) => {
const orgId = req.user!.organizationId;
const id = Number(req.params.id);
const { name, polygon, color, notifyEnter, notifyExit, isActive } = req.body ?? {};
const updates: Record<string, unknown> = {};
if (name !== undefined) updates.name = String(name).trim();
if (polygon !== undefined) {
if (!isValidPolygon(polygon)) {
return badRequest(res, "Полигон: массив минимум из 3 пар [lat, lng]");
}
updates.polygon = polygon;
}
if (color !== undefined) updates.color = color;
if (notifyEnter !== undefined) updates.notifyEnter = Boolean(notifyEnter);
if (notifyExit !== undefined) updates.notifyExit = Boolean(notifyExit);
if (isActive !== undefined) updates.isActive = Boolean(isActive);
const geozone = await gpsStorage.updateGeozone(id, orgId, updates);
if (!geozone) return res.status(404).json({ success: false, error: "Геозона не найдена" });
res.json({ success: true, geozone });
});
router.delete("/geozones/:id", async (req: AuthenticatedRequest, res: Response) => {
const deleted = await gpsStorage.deleteGeozone(Number(req.params.id), req.user!.organizationId);
if (!deleted) return res.status(404).json({ success: false, error: "Геозона не найдена" });
res.json({ success: true });
});
// События входа/выхода из геозон
router.get("/geozone-events", async (req: AuthenticatedRequest, res: Response) => {
const orgId = req.user!.organizationId;
const events = await gpsStorage.listGeozoneEvents(orgId, {
assetId: req.query.assetId ? Number(req.query.assetId) : undefined,
geozoneId: req.query.geozoneId ? Number(req.query.geozoneId) : undefined,
from: parseDateParam(req.query.from) ?? undefined,
to: parseDateParam(req.query.to) ?? undefined,
limit: req.query.limit ? Math.min(Number(req.query.limit), 1000) : undefined,
});
res.json({ success: true, events });
});
// =====================
// Группы
// =====================
router.get("/groups", async (req: AuthenticatedRequest, res: Response) => {
const groups = await gpsStorage.listGroups(req.user!.organizationId);
res.json({ success: true, groups });
});
router.post("/groups", async (req: AuthenticatedRequest, res: Response) => {
const orgId = req.user!.organizationId;
const { name, assetIds } = req.body ?? {};
if (!name || typeof name !== "string") {
return badRequest(res, "Название группы обязательно");
}
const group = await gpsStorage.createGroup(
orgId,
name.trim(),
Array.isArray(assetIds) ? assetIds.map(Number).filter(Number.isFinite) : []
);
res.status(201).json({ success: true, group });
});
router.patch("/groups/:id", async (req: AuthenticatedRequest, res: Response) => {
const orgId = req.user!.organizationId;
const id = Number(req.params.id);
const { name, assetIds } = req.body ?? {};
const group = await gpsStorage.updateGroup(id, orgId, {
name: name !== undefined ? String(name).trim() : undefined,
assetIds: Array.isArray(assetIds) ? assetIds.map(Number).filter(Number.isFinite) : undefined,
});
if (!group) return res.status(404).json({ success: false, error: "Группа не найдена" });
res.json({ success: true, group });
});
router.delete("/groups/:id", async (req: AuthenticatedRequest, res: Response) => {
const deleted = await gpsStorage.deleteGroup(Number(req.params.id), req.user!.organizationId);
if (!deleted) return res.status(404).json({ success: false, error: "Группа не найдена" });
res.json({ success: true });
});
// =====================
// Подписчики объекта
// =====================
router.get("/assets/:id/subscribers", async (req: AuthenticatedRequest, res: Response) => {
const orgId = req.user!.organizationId;
const id = Number(req.params.id);
const asset = await gpsStorage.getAsset(id, orgId);
if (!asset) return res.status(404).json({ success: false, error: "Объект не найден" });
const subscribers = await gpsStorage.listSubscribers(id);
res.json({ success: true, subscribers });
});
// Полная замена списка подписчиков: [{ userId, onEnter, onExit, onOffline }]
router.put("/assets/:id/subscribers", async (req: AuthenticatedRequest, res: Response) => {
const orgId = req.user!.organizationId;
const id = Number(req.params.id);
const asset = await gpsStorage.getAsset(id, orgId);
if (!asset) return res.status(404).json({ success: false, error: "Объект не найден" });
const list = Array.isArray(req.body) ? req.body : req.body?.subscribers;
if (!Array.isArray(list)) {
return badRequest(res, "Ожидается массив подписчиков [{ userId, onEnter, onExit, onOffline }]");
}
const subscribers = list.map((s: any) => ({
userId: Number(s.userId),
onEnter: Boolean(s.onEnter ?? true),
onExit: Boolean(s.onExit ?? true),
onOffline: Boolean(s.onOffline ?? false),
})).filter((s: { userId: number }) => Number.isFinite(s.userId));
const result = await gpsStorage.replaceSubscribers(id, subscribers);
res.json({ success: true, subscribers: result });
});
// =====================
// PWA-трекинг (личная геолокация пользователя)
// =====================
router.post("/my-position", async (req: AuthenticatedRequest, res: Response) => {
const user = req.user!;
const orgId = user.organizationId;
const { lat, lng, speed, course, accuracy } = req.body ?? {};
const latN = parseNum(lat);
const lngN = parseNum(lng);
if (latN === null || lngN === null) {
return badRequest(res, "Обязательные поля: lat, lng");
}
// Находим или создаём персональный объект пользователя
let asset = await gpsStorage.getPersonAssetByUser(user.id, orgId);
if (!asset) {
asset = await gpsStorage.createAsset({
organizationId: orgId,
name: `${user.firstName} ${user.lastName}`.trim() || user.email,
type: "person",
userId: user.id,
});
}
if (!asset.isActive) {
return res.status(403).json({ success: false, error: "Отслеживание геолокации выключено" });
}
await processPosition(asset.id, orgId, {
lat: latN,
lng: lngN,
speed: parseNum(speed),
course: parseNum(course),
accuracy: parseNum(accuracy),
});
res.json({ success: true, assetId: asset.id });
});
router.get("/my-tracking", async (req: AuthenticatedRequest, res: Response) => {
const user = req.user!;
const asset = await gpsStorage.getPersonAssetByUser(user.id, user.organizationId);
res.json({ success: true, enabled: Boolean(asset?.isActive), assetId: asset?.id ?? null });
});
router.put("/my-tracking", async (req: AuthenticatedRequest, res: Response) => {
const user = req.user!;
const orgId = user.organizationId;
const enabled = Boolean(req.body?.enabled);
let asset = await gpsStorage.getPersonAssetByUser(user.id, orgId);
if (!asset) {
if (!enabled) {
// Выключать нечего — объекта ещё нет
return res.json({ success: true, enabled: false, assetId: null });
}
asset = await gpsStorage.createAsset({
organizationId: orgId,
name: `${user.firstName} ${user.lastName}`.trim() || user.email,
type: "person",
userId: user.id,
isActive: true,
});
} else if (asset.isActive !== enabled) {
asset = (await gpsStorage.updateAsset(asset.id, orgId, { isActive: enabled }))!;
}
res.json({ success: true, enabled, assetId: asset.id });
});
app.use("/api/gps", router);
}
function isValidPolygon(polygon: unknown): polygon is Array<[number, number]> {
return (
Array.isArray(polygon) &&
polygon.length >= 3 &&
polygon.every(
(p) => Array.isArray(p) && p.length === 2 && Number.isFinite(Number(p[0])) && Number.isFinite(Number(p[1]))
)
);
}

443
server/gps/storage.ts Normal file
View File

@@ -0,0 +1,443 @@
import { db } from "../db";
import {
gpsAssets,
gpsPositions,
gpsGeozones,
gpsGeozoneEvents,
gpsGeozoneStates,
gpsGroups,
gpsGroupAssets,
gpsAssetSubscribers,
type GpsAsset,
type GpsPosition,
type GpsGeozone,
type GpsGeozoneEvent,
type GpsGroup,
type GpsAssetSubscriber,
type InsertGpsAsset,
type InsertGpsGeozone,
type InsertGpsPosition,
} from "@shared/schema";
import { eq, and, asc, desc, gte, lte, lt, inArray, sql, isNotNull } from "drizzle-orm";
/**
* Хранилище GPS-модуля. Все публичные методы, работающие с данными
* конкретной организации, обязательно фильтруют по organizationId.
* Исключения (помечены отдельно): поиск по external_id для ingest
* и выборка «протухших» объектов для воркера — они обходят все организации.
*/
export class GpsStorage {
// =====================
// Объекты (assets)
// =====================
async listAssets(organizationId: number) {
const assets = await db
.select()
.from(gpsAssets)
.where(eq(gpsAssets.organizationId, organizationId))
.orderBy(asc(gpsAssets.name));
if (assets.length === 0) return [];
const assetIds = assets.map((a) => a.id);
// Группы каждого объекта
const groupLinks = await db
.select({
assetId: gpsGroupAssets.assetId,
groupId: gpsGroups.id,
groupName: gpsGroups.name,
})
.from(gpsGroupAssets)
.innerJoin(gpsGroups, eq(gpsGroupAssets.groupId, gpsGroups.id))
.where(
and(
inArray(gpsGroupAssets.assetId, assetIds),
eq(gpsGroups.organizationId, organizationId)
)
);
// Число подписчиков каждого объекта
const subscriberCounts = await db
.select({
assetId: gpsAssetSubscribers.assetId,
count: sql<number>`count(*)::int`,
})
.from(gpsAssetSubscribers)
.where(inArray(gpsAssetSubscribers.assetId, assetIds))
.groupBy(gpsAssetSubscribers.assetId);
const groupsByAsset = new Map<number, Array<{ id: number; name: string }>>();
for (const link of groupLinks) {
const list = groupsByAsset.get(link.assetId) ?? [];
list.push({ id: link.groupId, name: link.groupName });
groupsByAsset.set(link.assetId, list);
}
const countByAsset = new Map(subscriberCounts.map((r) => [r.assetId, r.count]));
return assets.map((asset) => ({
...asset,
groups: groupsByAsset.get(asset.id) ?? [],
subscriberCount: countByAsset.get(asset.id) ?? 0,
}));
}
async getAsset(id: number, organizationId: number): Promise<GpsAsset | undefined> {
const [asset] = await db
.select()
.from(gpsAssets)
.where(and(eq(gpsAssets.id, id), eq(gpsAssets.organizationId, organizationId)));
return asset;
}
/**
* Поиск объекта по external_id без фильтра организации.
* Используется только публичным ingest-endpoint'ом (Traccar).
*/
async getAssetByExternalId(externalId: string): Promise<GpsAsset | undefined> {
const [asset] = await db
.select()
.from(gpsAssets)
.where(eq(gpsAssets.externalId, externalId));
return asset;
}
/** Персональный объект (type='person') пользователя в организации */
async getPersonAssetByUser(userId: number, organizationId: number): Promise<GpsAsset | undefined> {
const [asset] = await db
.select()
.from(gpsAssets)
.where(
and(
eq(gpsAssets.userId, userId),
eq(gpsAssets.organizationId, organizationId),
eq(gpsAssets.type, "person")
)
);
return asset;
}
async createAsset(data: InsertGpsAsset): Promise<GpsAsset> {
const [asset] = await db.insert(gpsAssets).values(data).returning();
return asset;
}
async updateAsset(id: number, organizationId: number, updates: Partial<InsertGpsAsset>): Promise<GpsAsset | undefined> {
const [asset] = await db
.update(gpsAssets)
.set(updates)
.where(and(eq(gpsAssets.id, id), eq(gpsAssets.organizationId, organizationId)))
.returning();
return asset;
}
async deleteAsset(id: number, organizationId: number): Promise<boolean> {
const deleted = await db
.delete(gpsAssets)
.where(and(eq(gpsAssets.id, id), eq(gpsAssets.organizationId, organizationId)))
.returning({ id: gpsAssets.id });
return deleted.length > 0;
}
/** Обновление последней позиции объекта (и перевод в online) */
async updateAssetLastPosition(
assetId: number,
data: { lat: number; lng: number; speed?: number | null; course?: number | null; recordedAt: Date }
): Promise<void> {
await db
.update(gpsAssets)
.set({
lastLat: data.lat,
lastLng: data.lng,
lastSpeed: data.speed ?? null,
lastCourse: data.course ?? null,
lastSeenAt: data.recordedAt,
isOnline: true,
})
.where(eq(gpsAssets.id, assetId));
}
async setAssetOnline(assetId: number, isOnline: boolean): Promise<void> {
await db.update(gpsAssets).set({ isOnline }).where(eq(gpsAssets.id, assetId));
}
/** Текущие позиции всех активных объектов организации */
async listCurrentPositions(organizationId: number): Promise<GpsAsset[]> {
return db
.select()
.from(gpsAssets)
.where(
and(
eq(gpsAssets.organizationId, organizationId),
eq(gpsAssets.isActive, true),
isNotNull(gpsAssets.lastLat)
)
);
}
/**
* Онлайн-объекты, от которых давно не было точек (для воркера офлайн-детекта).
* Обходит все организации — вызывается под withSuperAdmin.
*/
async getStaleOnlineAssets(cutoff: Date): Promise<GpsAsset[]> {
return db
.select()
.from(gpsAssets)
.where(
and(
eq(gpsAssets.isActive, true),
eq(gpsAssets.isOnline, true),
isNotNull(gpsAssets.lastSeenAt),
lt(gpsAssets.lastSeenAt, cutoff)
)
);
}
// =====================
// Позиции (треки)
// =====================
async insertPosition(data: InsertGpsPosition): Promise<GpsPosition> {
const [position] = await db.insert(gpsPositions).values(data).returning();
return position;
}
async getTrack(assetId: number, organizationId: number, from?: Date, to?: Date): Promise<GpsPosition[]> {
const conditions = [
eq(gpsPositions.assetId, assetId),
eq(gpsPositions.organizationId, organizationId),
];
if (from) conditions.push(gte(gpsPositions.recordedAt, from));
if (to) conditions.push(lte(gpsPositions.recordedAt, to));
return db
.select()
.from(gpsPositions)
.where(and(...conditions))
.orderBy(asc(gpsPositions.recordedAt));
}
// =====================
// Геозоны
// =====================
async listGeozones(organizationId: number): Promise<GpsGeozone[]> {
return db
.select()
.from(gpsGeozones)
.where(eq(gpsGeozones.organizationId, organizationId))
.orderBy(asc(gpsGeozones.name));
}
async listActiveGeozones(organizationId: number): Promise<GpsGeozone[]> {
return db
.select()
.from(gpsGeozones)
.where(and(eq(gpsGeozones.organizationId, organizationId), eq(gpsGeozones.isActive, true)));
}
async getGeozone(id: number, organizationId: number): Promise<GpsGeozone | undefined> {
const [geozone] = await db
.select()
.from(gpsGeozones)
.where(and(eq(gpsGeozones.id, id), eq(gpsGeozones.organizationId, organizationId)));
return geozone;
}
async createGeozone(data: InsertGpsGeozone): Promise<GpsGeozone> {
const [geozone] = await db.insert(gpsGeozones).values(data).returning();
return geozone;
}
async updateGeozone(id: number, organizationId: number, updates: Partial<InsertGpsGeozone>): Promise<GpsGeozone | undefined> {
const [geozone] = await db
.update(gpsGeozones)
.set(updates)
.where(and(eq(gpsGeozones.id, id), eq(gpsGeozones.organizationId, organizationId)))
.returning();
return geozone;
}
async deleteGeozone(id: number, organizationId: number): Promise<boolean> {
const deleted = await db
.delete(gpsGeozones)
.where(and(eq(gpsGeozones.id, id), eq(gpsGeozones.organizationId, organizationId)))
.returning({ id: gpsGeozones.id });
return deleted.length > 0;
}
// =====================
// События геозон и состояния
// =====================
async insertGeozoneEvent(data: {
organizationId: number;
geozoneId: number;
assetId: number;
event: "enter" | "exit";
lat?: number | null;
lng?: number | null;
}): Promise<GpsGeozoneEvent> {
const [event] = await db.insert(gpsGeozoneEvents).values(data).returning();
return event;
}
async listGeozoneEvents(
organizationId: number,
filters: { assetId?: number; geozoneId?: number; from?: Date; to?: Date; limit?: number }
): Promise<Array<GpsGeozoneEvent & { assetName: string | null; geozoneName: string | null }>> {
const conditions = [eq(gpsGeozoneEvents.organizationId, organizationId)];
if (filters.assetId) conditions.push(eq(gpsGeozoneEvents.assetId, filters.assetId));
if (filters.geozoneId) conditions.push(eq(gpsGeozoneEvents.geozoneId, filters.geozoneId));
if (filters.from) conditions.push(gte(gpsGeozoneEvents.createdAt, filters.from));
if (filters.to) conditions.push(lte(gpsGeozoneEvents.createdAt, filters.to));
const rows = await db
.select({
event: gpsGeozoneEvents,
assetName: gpsAssets.name,
geozoneName: gpsGeozones.name,
})
.from(gpsGeozoneEvents)
.leftJoin(gpsAssets, eq(gpsGeozoneEvents.assetId, gpsAssets.id))
.leftJoin(gpsGeozones, eq(gpsGeozoneEvents.geozoneId, gpsGeozones.id))
.where(and(...conditions))
.orderBy(desc(gpsGeozoneEvents.createdAt))
.limit(filters.limit ?? 200);
return rows.map((r) => ({ ...r.event, assetName: r.assetName, geozoneName: r.geozoneName }));
}
/** Состояния «внутри геозоны» для конкретного объекта */
async getGeozoneStatesForAsset(assetId: number): Promise<Map<number, boolean>> {
const rows = await db
.select()
.from(gpsGeozoneStates)
.where(eq(gpsGeozoneStates.assetId, assetId));
return new Map(rows.map((r) => [r.geozoneId, r.inside]));
}
async upsertGeozoneState(assetId: number, geozoneId: number, inside: boolean): Promise<void> {
await db
.insert(gpsGeozoneStates)
.values({ assetId, geozoneId, inside, updatedAt: new Date() })
.onConflictDoUpdate({
target: [gpsGeozoneStates.assetId, gpsGeozoneStates.geozoneId],
set: { inside, updatedAt: new Date() },
});
}
// =====================
// Группы
// =====================
async listGroups(organizationId: number): Promise<Array<GpsGroup & { assetIds: number[] }>> {
const groups = await db
.select()
.from(gpsGroups)
.where(eq(gpsGroups.organizationId, organizationId))
.orderBy(asc(gpsGroups.name));
if (groups.length === 0) return [];
const links = await db
.select()
.from(gpsGroupAssets)
.where(inArray(gpsGroupAssets.groupId, groups.map((g) => g.id)));
const assetsByGroup = new Map<number, number[]>();
for (const link of links) {
const list = assetsByGroup.get(link.groupId) ?? [];
list.push(link.assetId);
assetsByGroup.set(link.groupId, list);
}
return groups.map((g) => ({ ...g, assetIds: assetsByGroup.get(g.id) ?? [] }));
}
async getGroup(id: number, organizationId: number): Promise<GpsGroup | undefined> {
const [group] = await db
.select()
.from(gpsGroups)
.where(and(eq(gpsGroups.id, id), eq(gpsGroups.organizationId, organizationId)));
return group;
}
async createGroup(organizationId: number, name: string, assetIds: number[]): Promise<GpsGroup> {
const [group] = await db.insert(gpsGroups).values({ organizationId, name }).returning();
await this.setGroupAssets(group.id, assetIds);
return group;
}
async updateGroup(id: number, organizationId: number, updates: { name?: string; assetIds?: number[] }): Promise<GpsGroup | undefined> {
let group: GpsGroup | undefined;
if (updates.name !== undefined) {
[group] = await db
.update(gpsGroups)
.set({ name: updates.name })
.where(and(eq(gpsGroups.id, id), eq(gpsGroups.organizationId, organizationId)))
.returning();
} else {
group = await this.getGroup(id, organizationId);
}
if (!group) return undefined;
if (updates.assetIds !== undefined) {
await this.setGroupAssets(id, updates.assetIds);
}
return group;
}
async deleteGroup(id: number, organizationId: number): Promise<boolean> {
const deleted = await db
.delete(gpsGroups)
.where(and(eq(gpsGroups.id, id), eq(gpsGroups.organizationId, organizationId)))
.returning({ id: gpsGroups.id });
return deleted.length > 0;
}
/** Полная замена состава группы */
private async setGroupAssets(groupId: number, assetIds: number[]): Promise<void> {
await db.delete(gpsGroupAssets).where(eq(gpsGroupAssets.groupId, groupId));
if (assetIds.length > 0) {
await db
.insert(gpsGroupAssets)
.values(assetIds.map((assetId) => ({ groupId, assetId })))
.onConflictDoNothing();
}
}
// =====================
// Подписчики
// =====================
async listSubscribers(assetId: number): Promise<GpsAssetSubscriber[]> {
return db
.select()
.from(gpsAssetSubscribers)
.where(eq(gpsAssetSubscribers.assetId, assetId));
}
/** Полная замена списка подписчиков объекта */
async replaceSubscribers(
assetId: number,
subscribers: Array<{ userId: number; onEnter: boolean; onExit: boolean; onOffline: boolean }>
): Promise<GpsAssetSubscriber[]> {
await db.delete(gpsAssetSubscribers).where(eq(gpsAssetSubscribers.assetId, assetId));
if (subscribers.length === 0) return [];
return db
.insert(gpsAssetSubscribers)
.values(
subscribers.map((s) => ({
assetId,
userId: s.userId,
onEnter: s.onEnter,
onExit: s.onExit,
onOffline: s.onOffline,
}))
)
.returning();
}
}
export const gpsStorage = new GpsStorage();

View File

@@ -0,0 +1,95 @@
/**
* Клиент Traccar REST API.
* Настройка через env: TRACCAR_API_URL, TRACCAR_USER, TRACCAR_PASSWORD.
* Если переменные не заданы — все методы тихо no-op с логом,
* чтобы GPS-модуль работал и без Traccar (например, только PWA-трекинг).
*/
const TRACCAR_API_URL = process.env.TRACCAR_API_URL;
const TRACCAR_USER = process.env.TRACCAR_USER;
const TRACCAR_PASSWORD = process.env.TRACCAR_PASSWORD;
function isConfigured(): boolean {
return Boolean(TRACCAR_API_URL && TRACCAR_USER && TRACCAR_PASSWORD);
}
function authHeaders(): Record<string, string> {
const basic = Buffer.from(`${TRACCAR_USER}:${TRACCAR_PASSWORD}`).toString("base64");
return {
Authorization: `Basic ${basic}`,
"Content-Type": "application/json",
};
}
interface TraccarDevice {
id: number;
uniqueId: string;
name: string;
}
/** Найти устройство по uniqueId; возвращает null, если не найдено */
async function findDevice(uniqueId: string): Promise<TraccarDevice | null> {
const url = `${TRACCAR_API_URL}/api/devices?uniqueId=${encodeURIComponent(uniqueId)}`;
const res = await fetch(url, { headers: authHeaders() });
if (!res.ok) {
throw new Error(`Traccar GET /api/devices: HTTP ${res.status}`);
}
const devices = (await res.json()) as TraccarDevice[];
return devices.find((d) => d.uniqueId === uniqueId) ?? null;
}
/**
* Создаёт устройство в Traccar, если его ещё нет.
* Возвращает id устройства или null при ошибке/отсутствии конфигурации.
*/
export async function ensureDevice(uniqueId: string, name: string): Promise<number | null> {
if (!isConfigured()) {
console.log("[GPS] Traccar не настроен (TRACCAR_API_URL/USER/PASSWORD), ensureDevice пропущен");
return null;
}
try {
const existing = await findDevice(uniqueId);
if (existing) return existing.id;
const res = await fetch(`${TRACCAR_API_URL}/api/devices`, {
method: "POST",
headers: authHeaders(),
body: JSON.stringify({ uniqueId, name }),
});
if (!res.ok) {
const text = await res.text().catch(() => "");
throw new Error(`Traccar POST /api/devices: HTTP ${res.status} ${text}`);
}
const device = (await res.json()) as TraccarDevice;
console.log(`[GPS] Traccar: создано устройство '${name}' (${uniqueId}), id=${device.id}`);
return device.id;
} catch (err) {
console.error("[GPS] Traccar ensureDevice error:", err);
return null;
}
}
/** Удаляет устройство из Traccar по uniqueId (если существует) */
export async function deleteDevice(uniqueId: string): Promise<void> {
if (!isConfigured()) {
console.log("[GPS] Traccar не настроен, deleteDevice пропущен");
return;
}
try {
const existing = await findDevice(uniqueId);
if (!existing) return;
const res = await fetch(`${TRACCAR_API_URL}/api/devices/${existing.id}`, {
method: "DELETE",
headers: authHeaders(),
});
if (!res.ok && res.status !== 404) {
throw new Error(`Traccar DELETE /api/devices/${existing.id}: HTTP ${res.status}`);
}
console.log(`[GPS] Traccar: удалено устройство ${uniqueId}`);
} catch (err) {
console.error("[GPS] Traccar deleteDevice error:", err);
}
}

53
server/gps/worker.ts Normal file
View File

@@ -0,0 +1,53 @@
import { withSuperAdmin } from "../db";
import { gpsStorage } from "./storage";
import { notifyOffline } from "./geozone.service";
import { eventBus } from "../routes/shared";
/**
* Воркер офлайн-детекта GPS-объектов.
* Раз в минуту ищет активные онлайн-объекты, от которых не было точек
* более 15 минут, переводит их в offline и уведомляет подписчиков (on_offline).
* Обратный переход в online происходит в processPosition при новой точке.
* Работает независимо от Traccar (PWA-трекинг тоже покрывается).
*/
const TICK_MS = 60_000;
const OFFLINE_THRESHOLD_MS = 15 * 60_000;
let isRunning = false;
async function tick(): Promise<void> {
const cutoff = new Date(Date.now() - OFFLINE_THRESHOLD_MS);
const staleAssets = await withSuperAdmin(() => gpsStorage.getStaleOnlineAssets(cutoff));
for (const asset of staleAssets) {
try {
await withSuperAdmin(async () => {
await gpsStorage.setAssetOnline(asset.id, false);
await notifyOffline(asset);
});
// SSE — чтобы открытые карты сразу «погасили» маркер
eventBus.publishEvent({
type: "gps_asset_offline",
data: { assetId: asset.id },
organizationId: asset.organizationId,
});
console.log(`[GPS Worker] Объект ${asset.id} ('${asset.name}', org ${asset.organizationId}) переведён в offline`);
} catch (err) {
console.error(`[GPS Worker] Ошибка обработки объекта ${asset.id}:`, err);
}
}
}
export function startGpsWorker(): void {
if (isRunning) return;
isRunning = true;
setInterval(() => {
tick().catch((err) => console.error("[GPS Worker] Tick error:", err));
}, TICK_MS);
console.log("[GPS Worker] started");
}

View File

@@ -8,6 +8,7 @@ import { notificationService } from "./services/notification.service";
import { webPushService } from "./services/web-push.service";
import { startMedScheduleWorker } from "./medschedule/worker";
import { startAutomationScheduler } from "./workers/automation-scheduler";
import { startGpsWorker } from "./gps/worker";
import { storage } from "./storage";
import type { ReminderRecipient } from "@shared/schema";
import { db, withSuperAdmin } from "./db"; // lazy proxy — safe at module load; throws on first use if DATABASE_URL missing
@@ -1099,6 +1100,10 @@ async function runStartupDataPatches() {
startAutomationScheduler();
log('Automation scheduler started');
// Start GPS worker (offline detection for tracked assets)
startGpsWorker();
log('GPS worker started');
// Clean up duplicate subscriptions then backfill defaults for existing users
(async () => {
try {

View File

@@ -41,6 +41,7 @@ import organizationSkinsRouter from "./organization-skins.routes";
import { registerSyncRoutes } from "./sync.routes";
import { registerMedScheduleRoutes } from "../medschedule/routes";
import { registerMedScheduleBotRoutes } from "../medschedule/bot.routes";
import { registerGpsRoutes } from "../gps/routes";
import { registerVpnRoutes } from "../vpn";
import { FIELD_TYPE_DEFINITIONS } from "../../shared/field-types";
export { publishNotificationSSE };
@@ -120,6 +121,8 @@ export async function registerRoutes(app: Express): Promise<Server> {
registerSyncRoutes(app);
registerMedScheduleRoutes(app);
registerMedScheduleBotRoutes(app);
// GPS-модуль: публичный /api/gps/ingest регистрируется внутри до authenticateToken
registerGpsRoutes(app);
registerVpnRoutes(app);
const httpServer = createServer(app);