Files
iistwin/server/gps/geozone.service.ts

453 lines
18 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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;
}
/** Расстояние между двумя точками (haversine), в метрах. Без зависимостей. */
export function haversineMeters(lat1: number, lng1: number, lat2: number, lng2: number): number {
const R = 6371000; // радиус Земли, м
const toRad = (deg: number) => (deg * Math.PI) / 180;
const dLat = toRad(lat2 - lat1);
const dLng = toRad(lng2 - lng1);
const a =
Math.sin(dLat / 2) ** 2 +
Math.cos(toRad(lat1)) * Math.cos(toRad(lat2)) * Math.sin(dLng / 2) ** 2;
return 2 * R * Math.asin(Math.sqrt(a));
}
export interface ProcessPositionInput {
lat: number;
lng: number;
speed?: number | null;
course?: number | null;
accuracy?: number | null;
recordedAt?: Date;
}
// Вердикт фильтрации точки относительно предыдущей принятой
export type PointVerdict = "accept" | "stale" | "jump";
export interface GpsFilterSettings {
jitterMeters: number;
maxJumpSpeedKmh: number;
}
interface PrevPoint {
lat: number;
lng: number;
recordedAt: Date;
}
// Сколько точек подряд около новой локации нужно, чтобы подтвердить смену реальности
const CANDIDATE_CONFIRM_COUNT = 3;
// Время жизни кандидата на смену локации (по времени точек)
const CANDIDATE_TTL_MS = 30 * 60 * 1000;
interface JumpCandidate {
lat: number;
lng: number;
count: number;
firstAt: number; // recordedAt первой точки кандидата, ms
}
/**
* Stateful-фильтр потока точек (stale/jump) с гистерезисом подтверждения.
*
* Проблема latch-up: если последняя принятая точка — стабильный LBS
* (идёт каждую минуту с одной вышки), то при восстановлении GPS реальные
* точки вечно выглядят jump'ом и отбрасываются. Решение — кандидат:
* jump-точки запоминаются; если CANDIDATE_CONFIRM_COUNT точек подряд
* ложатся около одной новой локации (в пределах max(jitterMeters, 100) м) —
* это новая реальность (GPS ожил / объект реально уехал), точка принимается.
* Прыжки между разными вышками (туда-сюда) никогда не подтверждаются.
* Кандидат живёт CANDIDATE_TTL_MS с первой точки, потом сбрасывается.
*
* Используется и при приёме (processPosition, per-asset in-memory),
* и при отдаче трека (filterTrackPoints).
*/
export class JumpFilterWalker {
private prev: PrevPoint | null;
private candidate: JumpCandidate | null = null;
constructor(
private settings: GpsFilterSettings,
initialPrev: PrevPoint | null = null
) {
this.prev = initialPrev;
}
/** Обновление настроек «на лету» (для долгоживущих per-asset walker'ов) */
setSettings(settings: GpsFilterSettings): void {
this.settings = settings;
}
/**
* Проверка очередной точки. При 'accept' walker обновляет prev —
* вызывающий код должен обработать точку как валидную.
*/
check(point: PrevPoint): PointVerdict {
// stale: устаревшая копия/дубликат по времени устройства
if (this.prev && point.recordedAt.getTime() <= this.prev.recordedAt.getTime()) {
return "stale";
}
if (!this.isJump(point)) {
// Обычный путь — кандидат не нужен
this.candidate = null;
this.prev = point;
return "accept";
}
// jump: логика кандидата на смену локации
const now = point.recordedAt.getTime();
if (this.candidate && now - this.candidate.firstAt > CANDIDATE_TTL_MS) {
this.candidate = null; // старый кандидат протух — не мешает новым
}
if (this.candidate) {
const distToCandidate = haversineMeters(
this.candidate.lat,
this.candidate.lng,
point.lat,
point.lng
);
if (distToCandidate <= Math.max(this.settings.jitterMeters, 100)) {
this.candidate.count++;
if (this.candidate.count >= CANDIDATE_CONFIRM_COUNT) {
// Три подряд точки около новой локации — новая реальность, принимаем
this.candidate = null;
this.prev = point;
return "accept";
}
return "jump";
}
}
// Точка далеко и от кандидата — заменяем кандидата
this.candidate = { lat: point.lat, lng: point.lng, count: 1, firstAt: now };
return "jump";
}
/** Нереалистичный скачок скорости относительно prev (LBS-позиции) */
private isJump(point: PrevPoint): boolean {
if (!this.prev || this.settings.maxJumpSpeedKmh <= 0) return false;
const dtSec = (point.recordedAt.getTime() - this.prev.recordedAt.getTime()) / 1000;
const dist = haversineMeters(this.prev.lat, this.prev.lng, point.lat, point.lng);
if (dtSec <= 0) {
// Защита (после stale-проверки недостижимо, но на всякий случай)
return dist > this.settings.jitterMeters;
}
const impliedKmh = dist / 1000 / (dtSec / 3600);
return impliedKmh > this.settings.maxJumpSpeedKmh;
}
}
/**
* Фильтрация сохранённого трека при выдаче (без удаления из БД):
* те же правила stale/jump с подтверждением кандидата, что и при приёме.
*/
export function filterTrackPoints<T extends { lat: number; lng: number; recordedAt: Date }>(
points: T[],
settings: GpsFilterSettings
): { points: T[]; filteredCount: number } {
const walker = new JumpFilterWalker(settings);
const accepted: T[] = [];
let filteredCount = 0;
for (const p of points) {
if (walker.check({ lat: p.lat, lng: p.lng, recordedAt: p.recordedAt }) === "accept") {
accepted.push(p);
} else {
filteredCount++;
}
}
return { points: accepted, filteredCount };
}
/**
* Основная точка входа: сохраняет позицию, обновляет last_* объекта,
* проверяет геозоны и рассылает SSE/уведомления.
*
* Семантика времени:
* - last_seen_at = ВСЕГДА время ПРИЁМА сервером (онлайн-индикатор на карте);
* - last_recorded_at = время последней ПРИНЯТОЙ точки по устройству (fixTime);
* - gps_positions.recorded_at = время устройства (для треков).
*
* Цепочка фильтров (сравнение всегда с последней принятой точкой):
* 1. stale — устаревшая копия/дубликат по времени устройства;
* 2. jump — нереалистичный скачок скорости (LBS-позиции), с гистерезисом:
* подряд идущие точки около новой локации подтверждают смену реальности;
* 3. jitter — мелкое перемещение с малой скоростью.
* Отброшенная точка не пишется в gps_positions, координаты и геозоны не
* трогаются, но last_seen_at/is_online обновляются (устройство живо).
*/
// In-memory фильтры по объектам (процесс один; при рестарте состояние
// восстанавливается из last_* объекта в БД)
const assetWalkers = new Map<number, JumpFilterWalker>();
function getAssetWalker(assetId: number, prev: PrevPoint | null, settings: GpsFilterSettings): JumpFilterWalker {
let walker = assetWalkers.get(assetId);
if (!walker) {
walker = new JumpFilterWalker(settings, prev);
assetWalkers.set(assetId, walker);
} else {
walker.setSettings(settings);
}
return walker;
}
export async function processPosition(
assetId: number,
organizationId: number,
input: ProcessPositionInput
): Promise<void> {
const recordedAt = input.recordedAt ?? new Date();
const asset = await gpsStorage.getAsset(assetId, organizationId);
if (!asset) return;
const settings = await gpsStorage.getSettings(organizationId);
// Предыдущая принятая точка объекта (если есть) — для инициализации walker'а
const prev: PrevPoint | null =
asset.lastLat !== null && asset.lastLng !== null && asset.lastRecordedAt !== null
? { lat: asset.lastLat, lng: asset.lastLng, recordedAt: asset.lastRecordedAt }
: null;
// Фильтры stale/jump с гистерезисом подтверждения (устаревшие копии и LBS-скачки)
const walker = getAssetWalker(assetId, prev, settings);
const verdict = walker.check({ lat: input.lat, lng: input.lng, recordedAt });
if (verdict !== "accept") {
if (verdict === "stale") {
console.log(
`[GPS] skipped stale: asset=${assetId} recordedAt=${recordedAt.toISOString()} <= lastRecordedAt=${prev?.recordedAt.toISOString()}`
);
} else {
const dist = prev ? haversineMeters(prev.lat, prev.lng, input.lat, input.lng) : 0;
const dtSec = prev ? (recordedAt.getTime() - prev.recordedAt.getTime()) / 1000 : 0;
const impliedKmh = dtSec > 0 ? dist / 1000 / (dtSec / 3600) : Infinity;
console.log(
`[GPS] skipped jump: asset=${assetId} dist=${Math.round(dist)}m dt=${Math.round(dtSec)}s implied=${Math.round(impliedKmh)}km/h > max=${settings.maxJumpSpeedKmh}km/h`
);
}
// Устройство живо — фиксируем факт приёма, точку не сохраняем
await gpsStorage.touchAssetSeen(assetId);
return;
}
// Фильтр джиттера (округление мелких перемещений)
if (settings.jitterMeters > 0 && prev) {
const distance = haversineMeters(prev.lat, prev.lng, input.lat, input.lng);
const speed = input.speed ?? null;
if (distance < settings.jitterMeters && (speed === null || speed < 5)) {
await gpsStorage.touchAssetSeen(assetId);
return;
}
}
// 1. Сохраняем точку в историю (recorded_at — время устройства, для треков)
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. Обновляем последнюю позицию объекта (last_seen_at = время приёма, 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(asset, input, recordedAt);
}
async function checkGeozones(
asset: GpsAsset,
input: ProcessPositionInput,
recordedAt: Date
): Promise<void> {
const { id: assetId, organizationId } = asset;
const geozones = await gpsStorage.listActiveGeozones(organizationId);
if (geozones.length === 0) return;
const states = await gpsStorage.getGeozoneStatesForAsset(assetId);
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);
}
}
}