Каждая точка ingest (~1/сек) broadcast'илась SSE всем клиентам орг., клиент на каждое событие делал invalidateQueries без throttle → 2+ GET/сек → исчерпание per-user бакета apiLimiter (1500/15мин) → все API устройства 429 (формы не грузились). Модель теперь как у нормальных трекеров: точки пишутся в БД, карта читает по интервалу. +4 unit-теста троттла (server/gps/sse-throttle.ts)
803 lines
36 KiB
TypeScript
803 lines
36 KiB
TypeScript
import { gpsStorage } from "./storage";
|
||
import { OSRM_URL } from "./osrm";
|
||
import { shouldPublishGpsAssetMoved } from "./sse-throttle";
|
||
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;
|
||
/** Протокольный признак «GPS невалиден» (Traccar valid=false) — LBS-точка */
|
||
forceLbs?: boolean;
|
||
}
|
||
|
||
// Вердикт фильтрации точки относительно предыдущей принятой
|
||
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;
|
||
|
||
// Траекторный фильтр при выдаче (проход после walker'а, до spike-фильтра):
|
||
// ловит «мяч» вышек — LBS-точки, скачущие по области десятки км, где каждая
|
||
// по отдельности в пределах off_road_meters от какой-нибудь дороги.
|
||
const TRAJ_WINDOW_SIZE = 10; // скользящее окно последних принятых точек
|
||
const TRAJ_ANCHOR_DIST_M = 1500; // дальше этого от якоря (медиана окна) — подозрительная
|
||
const TRAJ_CONFIRM_DIST_M = 800; // ближе этого к подозрительной — подтверждение переезда
|
||
const TRAJ_LOOKAHEAD = 5; // сколько следующих точек потока смотрим для подтверждения
|
||
const TRAJ_MIN_HISTORY = 3; // меньше точек в окне — фильтр ещё не работает
|
||
|
||
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++;
|
||
}
|
||
}
|
||
|
||
// Траекторный фильтр по окну: якорь = медианы lat/lng последних 10 принятых
|
||
// точек. Точка дальше 1500 м от якоря — подозрительная: складывается в
|
||
// pending и ждёт lookahead. Если среди следующих 5 точек потока ≥ 2 лежат
|
||
// в пределах 800 м от неё — реальный переезд в новый район: pending
|
||
// принимаются, якорь съезжает. Если поток вернулся к якорю — pending
|
||
// отбрасываются как мусор (считаются в filteredCount).
|
||
// Применяется ко ВСЕМ точкам (и GPS, и LBS): GPS-выбросы тоже случаются.
|
||
// O(n·window) — для ~4000 точек нормально.
|
||
const trajectory: T[] = [];
|
||
const window: T[] = [];
|
||
let pending: T[] = [];
|
||
const pushToWindow = (p: T) => {
|
||
window.push(p);
|
||
if (window.length > TRAJ_WINDOW_SIZE) window.shift();
|
||
};
|
||
|
||
for (let i = 0; i < accepted.length; i++) {
|
||
const p = accepted[i];
|
||
|
||
// Истории мало — окно ещё не работает, принимаем всё
|
||
if (window.length < TRAJ_MIN_HISTORY) {
|
||
trajectory.push(p);
|
||
pushToWindow(p);
|
||
continue;
|
||
}
|
||
|
||
const anchor = medianAnchor(window);
|
||
if (haversineMeters(anchor.lat, anchor.lng, p.lat, p.lng) <= TRAJ_ANCHOR_DIST_M) {
|
||
// Точка логически вписывается в траекторию
|
||
if (pending.length > 0) {
|
||
filteredCount += pending.length;
|
||
pending = [];
|
||
}
|
||
trajectory.push(p);
|
||
pushToWindow(p);
|
||
continue;
|
||
}
|
||
|
||
// Подозрительная (улетела от якоря) — смотрим вперёд по потоку
|
||
let confirm = 0;
|
||
const lookaheadEnd = Math.min(i + 1 + TRAJ_LOOKAHEAD, accepted.length);
|
||
for (let j = i + 1; j < lookaheadEnd; j++) {
|
||
if (haversineMeters(p.lat, p.lng, accepted[j].lat, accepted[j].lng) <= TRAJ_CONFIRM_DIST_M) {
|
||
confirm++;
|
||
}
|
||
}
|
||
|
||
if (confirm >= 2) {
|
||
// Реальный переезд: принимаем накопленный pending и текущую точку,
|
||
// окно скользит на новый район
|
||
for (const q of pending) {
|
||
trajectory.push(q);
|
||
pushToWindow(q);
|
||
}
|
||
pending = [];
|
||
trajectory.push(p);
|
||
pushToWindow(p);
|
||
} else {
|
||
// Ждём: либо переезд подтвердится следующими точками, либо поток
|
||
// вернётся к якорю и pending будет отброшен
|
||
pending.push(p);
|
||
}
|
||
}
|
||
// Хвост pending в конце потока — мусор (к якорю так и не вернулись)
|
||
filteredCount += pending.length;
|
||
|
||
// Spike-фильтр с lookahead: точка-«выброс» к координатам стоянки между
|
||
// движущимися точками (dist(prev,cur) большая, а следующая точка снова
|
||
// рядом с prev, сама точка почти стоит, соседи едут)
|
||
const out: T[] = [];
|
||
for (let i = 0; i < trajectory.length; i++) {
|
||
const cur = trajectory[i];
|
||
const prev = out[out.length - 1];
|
||
const next = trajectory[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 };
|
||
}
|
||
|
||
/**
|
||
* Якорь траектории: медиана lat и медиана lng окна (по отдельности).
|
||
* Медиана устойчива к одиночным выбросам внутри окна, в отличие от среднего.
|
||
*/
|
||
function medianAnchor(window: Array<{ lat: number; lng: number }>): { lat: number; lng: number } {
|
||
const mid = Math.floor(window.length / 2);
|
||
const median = (values: number[]) => {
|
||
const sorted = [...values].sort((a, b) => a - b);
|
||
return sorted.length % 2 === 1 ? sorted[mid] : (sorted[mid - 1] + sorted[mid]) / 2;
|
||
};
|
||
return {
|
||
lat: median(window.map((p) => p.lat)),
|
||
lng: median(window.map((p) => p.lng)),
|
||
};
|
||
}
|
||
|
||
/**
|
||
* Основная точка входа: сохраняет позицию, обновляет 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 обновляются (устройство живо).
|
||
*
|
||
* Режим «потерян GPS»: идущие stale re-send'ы — сигнал, что трекер не видит
|
||
* спутники. Пока последний stale моложе 10 минут, свежие точки помечаются
|
||
* is_lbs (координаты сотовой вышки): сохраняются и отдаются клиенту
|
||
* (пунктирное отображение), но геозоны по ним не проверяются. Дополнительно
|
||
* LBS-точки дальше off_road_meters от ближайшей дороги (OSRM nearest)
|
||
* отбрасываются совсем. Первый настоящий GPS-фикс после окна stale
|
||
* выключает режим (gps_lost = false).
|
||
*/
|
||
|
||
// In-memory фильтры по объектам (процесс один; при рестарте состояние
|
||
// восстанавливается из last_* объекта в БД)
|
||
const assetWalkers = new Map<number, JumpFilterWalker>();
|
||
|
||
// РЕШЕНО (2026-08-31): слой «stale re-send => LBS-режим» убран — он помечал
|
||
// ~100% точек (любая ретрансляция пакета открывала 10-мин окно LBS).
|
||
// Источник истины: протокольный флаг valid=false (Traccar) + физика фильтров.
|
||
|
||
// Таймаут запроса к OSRM Nearest API (фильтр «далеко от дороги»)
|
||
const OSRM_NEAREST_TIMEOUT_MS = 2000;
|
||
// Warn про недоступность OSRM для nearest пишем один раз до следующего успешного ответа
|
||
let osrmNearestWarned = false;
|
||
|
||
/**
|
||
* Проверка «далеко от дороги» через OSRM Nearest API (только для LBS-режима).
|
||
* Возвращает расстояние до ближайшей дороги (м), если оно БОЛЬШЕ порога
|
||
* (точку надо отбросить), иначе null.
|
||
* При недоступности OSRM/таймауте — null: приём точки не блокируем.
|
||
* code=NoSegment — ближайшая дорога вне радиуса поиска OSRM (порядка км),
|
||
* надёжный сигнал LBS: возвращаем Infinity (точно offroad).
|
||
*/
|
||
export async function offRoadDistance(lat: number, lng: number, thresholdMeters: number): Promise<number | null> {
|
||
let data: { code?: string; waypoints?: Array<{ distance?: number }> };
|
||
try {
|
||
// ВНИМАНИЕ: OSRM ждёт координаты в порядке lng,lat
|
||
const res = await fetch(`${OSRM_URL}/nearest/v1/driving/${lng},${lat}`, {
|
||
signal: AbortSignal.timeout(OSRM_NEAREST_TIMEOUT_MS),
|
||
});
|
||
data = (await res.json()) as typeof data;
|
||
if (osrmNearestWarned) {
|
||
console.log("[GPS] OSRM снова доступен (nearest)");
|
||
osrmNearestWarned = false;
|
||
}
|
||
} catch (err) {
|
||
// OSRM недоступен — пропускаем проверку, точка идёт обычным путём
|
||
if (!osrmNearestWarned) {
|
||
console.warn(`[GPS] OSRM недоступен для nearest (${OSRM_URL}):`, err instanceof Error ? err.message : err);
|
||
osrmNearestWarned = true;
|
||
}
|
||
return null;
|
||
}
|
||
|
||
// Дорога вне радиуса поиска — точно далеко от дороги
|
||
if (data.code === "NoSegment") return Infinity;
|
||
if (data.code !== "Ok") return null;
|
||
|
||
const dist = data.waypoints?.[0]?.distance;
|
||
if (typeof dist !== "number") return null;
|
||
return dist > thresholdMeters ? dist : null;
|
||
}
|
||
|
||
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;
|
||
|
||
// Гидратация in-memory состояния после рестарта: если в БД объект помечен
|
||
// «потерян GPS», а lastStaleAt пуст — считаем, что stale был только что
|
||
// (осторожный режим до свежих событий)
|
||
|
||
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;
|
||
|
||
// Протокольный флаг «GPS невалиден» (Traccar valid=false): это точка по вышкам.
|
||
// Сохраняем как примерную (is_lbs), маркер двигаем, но walker/геозоны/пробег
|
||
// не трогаем. Заодно это сигнал потери GPS — как stale re-send.
|
||
if (input.forceLbs) {
|
||
if (!asset.gpsLost) {
|
||
await gpsStorage.setGpsLost(assetId, true);
|
||
console.log(`[GPS] gps_lost ON: asset=${assetId} (valid=false от протокола)`);
|
||
}
|
||
await gpsStorage.insertPosition({
|
||
assetId,
|
||
organizationId,
|
||
lat: input.lat,
|
||
lng: input.lng,
|
||
speed: input.speed ?? null,
|
||
course: input.course ?? null,
|
||
accuracy: input.accuracy ?? null,
|
||
isLbs: true,
|
||
recordedAt,
|
||
});
|
||
await gpsStorage.updateAssetLastPosition(assetId, {
|
||
lat: input.lat,
|
||
lng: input.lng,
|
||
speed: input.speed,
|
||
course: input.course,
|
||
recordedAt,
|
||
isLbs: true,
|
||
});
|
||
// SSE-broadcast троттлим (не чаще раза в 3 сек на объект) — запись в БД выше не трогаем
|
||
if (shouldPublishGpsAssetMoved(assetId)) {
|
||
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(),
|
||
isLbs: true,
|
||
},
|
||
organizationId,
|
||
});
|
||
}
|
||
return;
|
||
}
|
||
|
||
// Фильтры 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") {
|
||
// Stale-копии просто пропускаем. lastGpsFixAt здесь НЕ трогаем:
|
||
// бейдж «Нет GPS» = «нет точки с валидным GPS дольше N часов» (настройка)
|
||
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;
|
||
}
|
||
|
||
// Точка прошла фильтры и не помечена протоколом как LBS — это валидный GPS:
|
||
// двигаем «время последнего GPS-фикса» вперёд (им питается бейдж «Нет GPS»)
|
||
await gpsStorage.updateAssetGpsFix(assetId, recordedAt, true);
|
||
|
||
// Режим «потерян GPS»: если последний stale re-send был свежее 10 минут,
|
||
// свежая точка — это координаты сотовой вышки (LBS), а не GPS-фикс.
|
||
// Такие точки сохраняем и показываем (пунктир на карте), но в логику
|
||
// геозон не пускаем. jump/echo/джиттер-отбросы режим не переключают.
|
||
// Сюда доходят только точки без протокольного признака LBS
|
||
// (forceLbs-точки обработаны и вышли выше) — значит, обычная GPS-точка
|
||
const isLbs = false;
|
||
|
||
// Фильтр «далеко от дороги» — ТОЛЬКО в режиме потери GPS (LBS-точки):
|
||
// координаты сотовых вышек обычно >100 м от ближайшей дороги, настоящий
|
||
// GPS — почти никогда (p90 ≈ 7 м). Вне LBS-режима проверку не выполняем:
|
||
// при живом GPS заезды в поля/дворы легальны.
|
||
if (isLbs && settings.offRoadMeters > 0) {
|
||
const distToRoad = await offRoadDistance(input.lat, input.lng, settings.offRoadMeters);
|
||
if (distToRoad !== null) {
|
||
console.log(
|
||
`[GPS] skipped offroad: asset=${assetId} dist_to_road=${distToRoad === Infinity ? "> радиуса поиска OSRM" : `${Math.round(distToRoad)}m`} > порога ${settings.offRoadMeters}m`
|
||
);
|
||
// Точку не сохраняем, координаты и геозоны не трогаем — только «устройство живо»
|
||
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,
|
||
isLbs,
|
||
recordedAt,
|
||
});
|
||
|
||
// 2. Обновляем последнюю позицию объекта (last_seen_at = время приёма, online)
|
||
await gpsStorage.updateAssetLastPosition(assetId, {
|
||
lat: input.lat,
|
||
lng: input.lng,
|
||
speed: input.speed,
|
||
course: input.course,
|
||
recordedAt,
|
||
isLbs,
|
||
});
|
||
|
||
// Настоящий GPS-фикс после режима «потерян GPS» — выключаем режим
|
||
if (!isLbs && asset.gpsLost) {
|
||
await gpsStorage.setGpsLost(assetId, false);
|
||
console.log(`[GPS] gps_lost OFF: asset=${assetId} (свежий GPS-фикс)`);
|
||
}
|
||
|
||
// 3. SSE-событие о перемещении — broadcast по организации,
|
||
// чтобы карта у всех открытых клиентов обновилась в реальном времени.
|
||
// Троттлим не чаще раза в 3 сек на объект, иначе при потоке точек
|
||
// ~1/сек клиентская инвалидация исчерпывает apiLimiter (429).
|
||
if (shouldPublishGpsAssetMoved(assetId)) {
|
||
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(),
|
||
isLbs,
|
||
},
|
||
organizationId,
|
||
});
|
||
}
|
||
|
||
// 4. Проверка геозон — только для настоящих GPS-точек;
|
||
// LBS-координаты вышек в логику геозон не пускаем
|
||
if (!isLbs) {
|
||
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);
|
||
}
|
||
}
|
||
}
|