Files
iistwin/server/gps/storage.ts

511 lines
17 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 { db } from "../db";
import {
gpsAssets,
gpsPositions,
gpsGeozones,
gpsGeozoneEvents,
gpsGeozoneStates,
gpsGroups,
gpsGroupAssets,
gpsAssetSubscribers,
gpsSettings,
type GpsAsset,
type GpsPosition,
type GpsGeozone,
type GpsGeozoneEvent,
type GpsGroup,
type GpsAssetSubscriber,
type GpsSettings,
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).
* last_seen_at = время ПРИЁМА сервером (new Date() здесь),
* last_recorded_at = время точки по устройству (recordedAt).
*/
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: new Date(),
lastRecordedAt: data.recordedAt,
isOnline: true,
})
.where(eq(gpsAssets.id, assetId));
}
/**
* Отметка «устройство живо» без обновления координат
* (точка отброшена фильтром джиттера, но приём был).
*/
async touchAssetSeen(assetId: number): Promise<void> {
await db
.update(gpsAssets)
.set({ lastSeenAt: new Date(), 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();
}
// =====================
// Настройки организации
// =====================
/** Настройки GPS организации; если строки нет — дефолты */
async getSettings(organizationId: number): Promise<{
jitterMeters: number;
maxJumpSpeedKmh: number;
tripStopMinutes: number;
defaultLat: number | null;
defaultLng: number | null;
defaultZoom: number | null;
}> {
const [row] = await db
.select()
.from(gpsSettings)
.where(eq(gpsSettings.organizationId, organizationId));
return {
jitterMeters: row?.jitterMeters ?? 25,
maxJumpSpeedKmh: row?.maxJumpSpeedKmh ?? 250,
tripStopMinutes: row?.tripStopMinutes ?? 10,
defaultLat: row?.defaultLat ?? null,
defaultLng: row?.defaultLng ?? null,
defaultZoom: row?.defaultZoom ?? null,
};
}
async upsertSettings(
organizationId: number,
settings: {
jitterMeters: number;
maxJumpSpeedKmh: number;
tripStopMinutes: number;
defaultLat: number | null;
defaultLng: number | null;
defaultZoom: number | null;
}
): Promise<GpsSettings> {
const [row] = await db
.insert(gpsSettings)
.values({ organizationId, ...settings, updatedAt: new Date() })
.onConflictDoUpdate({
target: gpsSettings.organizationId,
set: { ...settings, updatedAt: new Date() },
})
.returning();
return row;
}
}
export const gpsStorage = new GpsStorage();