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

556 lines
23 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" | "echo";
export interface GpsFilterSettings {
jitterMeters: number;
maxJumpSpeedKmh: number;
}
interface PrevPoint {
lat: number;
lng: number;
recordedAt: Date;
/** Скорость, км/ч (null — неизвестна) */
speed?: number | null;
}
// Сколько точек подряд около новой локации нужно, чтобы подтвердить смену реальности
const CANDIDATE_CONFIRM_COUNT = 3;
// Время жизни кандидата на смену локации (по времени точек)
const CANDIDATE_TTL_MS = 30 * 60 * 1000;
// Окно heartbeat-эхо: статусная точка считается эхом стоянки, если последняя
// принятая точка «в движении» (speed > 10 км/ч) была меньше 5 минут назад
const ECHO_WINDOW_MS = 5 * 60 * 1000;
interface JumpCandidate {
lat: number;
lng: number;
count: number;
firstAt: number; // recordedAt первой точки кандидата, ms
}
/**
* Stateful-фильтр потока точек (stale/jump/echo) с гистерезисом подтверждения.
*
* Проблема latch-up: если последняя принятая точка — стабильный LBS
* (идёт каждую минуту с одной вышки), то при восстановлении GPS реальные
* точки вечно выглядят jump'ом и отбрасываются. Решение — кандидат:
* jump-точки запоминаются; если CANDIDATE_CONFIRM_COUNT точек подряд
* ложатся около одной новой локации (в пределах max(jitterMeters, 100) м) —
* это новая реальность (GPS ожил / объект реально уехал), точка принимается.
* Прыжки между разными вышками (туда-сюда) никогда не подтверждаются.
* Кандидат живёт CANDIDATE_TTL_MS с первой точки, потом сбрасывается.
*
* Heartbeat-эхо: трекер шлёт поминутные «статусные» точки (speed≈0),
* повторяющие координаты стоянки, даже когда объект едет. Walker помнит
* последний кластер стоянки (lastStationary) и время последней принятой
* точки в движении (lastMovingAt). Точка с speed < 1 км/ч рядом со стоянкой
* в пределах ECHO_WINDOW_MS после движения — эхо, отбрасывается. Если движения
* нет дольше окна — правило отпускает и стоянка принимается (это корректно).
*
* Используется и при приёме (processPosition, per-asset in-memory),
* и при отдаче трека (filterTrackPoints).
*/
export class JumpFilterWalker {
private prev: PrevPoint | null;
private candidate: JumpCandidate | null = null;
// Последняя «стоянка»: кластер принятых точек с speed≈0
private lastStationary: { lat: number; lng: number } | null = null;
// recordedAt последней принятой точки «в движении» (speed > 10 км/ч), ms
private lastMovingAt: number | null = null;
constructor(
private settings: GpsFilterSettings,
initialPrev: PrevPoint | null = null
) {
this.prev = initialPrev;
if (initialPrev) this.noteMotion(initialPrev);
}
/** Обновление настроек «на лету» (для долгоживущих per-asset walker'ов) */
setSettings(settings: GpsFilterSettings): void {
this.settings = settings;
}
/** Учёт принятой точки в состоянии стоянки/движения */
private noteMotion(point: PrevPoint): void {
const speed = point.speed;
if (speed === null || speed === undefined) return;
if (speed < 1) {
this.lastStationary = { lat: point.lat, lng: point.lng };
} else if (speed > 10) {
this.lastMovingAt = point.recordedAt.getTime();
}
}
/**
* Проверка очередной точки. При 'accept' walker обновляет prev —
* вызывающий код должен обработать точку как валидную.
*/
check(point: PrevPoint): PointVerdict {
// stale: устаревшая копия/дубликат по времени устройства
if (this.prev && point.recordedAt.getTime() <= this.prev.recordedAt.getTime()) {
return "stale";
}
// Heartbeat-эхо: статусная точка повторяет координаты стоянки на ходу
const speed = point.speed;
if (
speed !== null && speed !== undefined && speed < 1 &&
this.lastStationary &&
this.lastMovingAt !== null &&
point.recordedAt.getTime() - this.lastMovingAt < ECHO_WINDOW_MS &&
haversineMeters(this.lastStationary.lat, this.lastStationary.lng, point.lat, point.lng) <
Math.max(this.settings.jitterMeters, 50)
) {
return "echo";
}
if (!this.isJump(point)) {
// Обычный путь — кандидат не нужен
this.candidate = null;
this.prev = point;
this.noteMotion(point);
return "accept";
}
// Жёсткое противоречие скоростей (LBS-телепорт): отбрасываем сразу,
// такая точка никогда не принимается и не может стать jump-кандидатом
if (this.isContradiction(point)) {
return "jump";
}
// 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;
this.noteMotion(point);
return "accept";
}
return "jump";
}
}
// Точка далеко и от кандидата — заменяем кандидата
this.candidate = { lat: point.lat, lng: point.lng, count: 1, firstAt: now };
return "jump";
}
/** Метрики перемещения относительно последней принятой точки */
private jumpMetrics(point: PrevPoint): { dtSec: number; dist: number; impliedKmh: number } | null {
if (!this.prev) return null;
const dtSec = (point.recordedAt.getTime() - this.prev.recordedAt.getTime()) / 1000;
const dist = haversineMeters(this.prev.lat, this.prev.lng, point.lat, point.lng);
const impliedKmh =
dtSec > 0 ? dist / 1000 / (dtSec / 3600) : dist > this.settings.jitterMeters ? Infinity : 0;
return { dtSec, dist, impliedKmh };
}
/**
* Жёсткое противоречие скоростей (LBS-телепорт к вышке): точка «стоит»
* (speed < 5), но перемещение от предыдущей принятой требует скорости
* выше предельной — физически невозможно. Отбрасывается всегда и не
* может стать jump-кандидатом (иначе стабильный LBS-кластер легализуется).
*/
private isContradiction(point: PrevPoint): boolean {
const m = this.jumpMetrics(point);
if (!m) return false;
const reported = point.speed ?? 0;
const maxJump = this.settings.maxJumpSpeedKmh > 0 ? this.settings.maxJumpSpeedKmh : 250;
return reported < 5 && m.dtSec < 300 && m.impliedKmh > maxJump;
}
/** Нереалистичный скачок скорости относительно prev (LBS-позиции) */
private isJump(point: PrevPoint): boolean {
if (!this.prev || this.settings.maxJumpSpeedKmh <= 0) return false;
const m = this.jumpMetrics(point);
if (!m || m.impliedKmh <= this.settings.maxJumpSpeedKmh) return false;
// Реальная быстрая езда: трекер разгонялся между редкими точками,
// implied согласован с заявленной скоростью (допуск x3.5, минимум 90 км/ч)
const reported = point.speed ?? 0;
if (m.impliedKmh <= Math.max(reported * 3.5, 90)) return false;
return true;
}
}
/**
* Фильтрация сохранённого трека при выдаче (без удаления из БД):
* те же правила stale/jump/echo с подтверждением кандидата, что и при приёме,
* плюс spike-фильтр с lookahead (эхо стоянки между движущимися точками).
*/
export function filterTrackPoints<T extends { lat: number; lng: number; recordedAt: Date; speed?: number | null }>(
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, speed: p.speed }) === "accept") {
accepted.push(p);
} else {
filteredCount++;
}
}
// Spike-фильтр с lookahead: точка-«выброс» к координатам стоянки между
// движущимися точками (dist(prev,cur) большая, а следующая точка снова
// рядом с prev, сама точка почти стоит, соседи едут)
const out: T[] = [];
for (let i = 0; i < accepted.length; i++) {
const cur = accepted[i];
const prev = out[out.length - 1];
const next = accepted[i + 1];
if (prev && next) {
const dPrevCur = haversineMeters(prev.lat, prev.lng, cur.lat, cur.lng);
const dPrevNext = haversineMeters(prev.lat, prev.lng, next.lat, next.lng);
if (
dPrevCur > 250 &&
dPrevNext < dPrevCur / 2 &&
(cur.speed ?? 999) < 2 &&
((prev.speed ?? 0) > 10 || (next.speed ?? 0) > 10)
) {
filteredCount++;
continue;
}
}
out.push(cur);
}
return { points: out, 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, speed: asset.lastSpeed }
: null;
// Фильтры stale/jump/echo с гистерезисом подтверждения (устаревшие копии,
// LBS-скачки, heartbeat-эхо стоянки)
const walker = getAssetWalker(assetId, prev, settings);
const verdict = walker.check({ lat: input.lat, lng: input.lng, recordedAt, speed: input.speed });
if (verdict !== "accept") {
if (verdict === "stale") {
console.log(
`[GPS] skipped stale: asset=${assetId} recordedAt=${recordedAt.toISOString()} <= lastRecordedAt=${prev?.recordedAt.toISOString()}`
);
} else if (verdict === "echo") {
console.log(
`[GPS] skipped echo: asset=${assetId} heartbeat-точка повторяет координаты стоянки (speed=${input.speed ?? "—"} км/ч)`
);
} 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);
}
}
}