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`count(*)::int`, }) .from(gpsAssetSubscribers) .where(inArray(gpsAssetSubscribers.assetId, assetIds)) .groupBy(gpsAssetSubscribers.assetId); const groupsByAsset = new Map>(); 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 { 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 { const [asset] = await db .select() .from(gpsAssets) .where(eq(gpsAssets.externalId, externalId)); return asset; } /** Персональный объект (type='person') пользователя в организации */ async getPersonAssetByUser(userId: number, organizationId: number): Promise { 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 { const [asset] = await db.insert(gpsAssets).values(data).returning(); return asset; } async updateAsset(id: number, organizationId: number, updates: Partial): Promise { 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 { 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 { 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 { await db .update(gpsAssets) .set({ lastSeenAt: new Date(), isOnline: true }) .where(eq(gpsAssets.id, assetId)); } async setAssetOnline(assetId: number, isOnline: boolean): Promise { await db.update(gpsAssets).set({ isOnline }).where(eq(gpsAssets.id, assetId)); } /** Текущие позиции всех активных объектов организации */ async listCurrentPositions(organizationId: number): Promise { return db .select() .from(gpsAssets) .where( and( eq(gpsAssets.organizationId, organizationId), eq(gpsAssets.isActive, true), isNotNull(gpsAssets.lastLat) ) ); } /** * Онлайн-объекты, от которых давно не было точек (для воркера офлайн-детекта). * Обходит все организации — вызывается под withSuperAdmin. */ async getStaleOnlineAssets(cutoff: Date): Promise { 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 { const [position] = await db.insert(gpsPositions).values(data).returning(); return position; } async getTrack(assetId: number, organizationId: number, from?: Date, to?: Date): Promise { 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 { return db .select() .from(gpsGeozones) .where(eq(gpsGeozones.organizationId, organizationId)) .orderBy(asc(gpsGeozones.name)); } async listActiveGeozones(organizationId: number): Promise { return db .select() .from(gpsGeozones) .where(and(eq(gpsGeozones.organizationId, organizationId), eq(gpsGeozones.isActive, true))); } async getGeozone(id: number, organizationId: number): Promise { const [geozone] = await db .select() .from(gpsGeozones) .where(and(eq(gpsGeozones.id, id), eq(gpsGeozones.organizationId, organizationId))); return geozone; } async createGeozone(data: InsertGpsGeozone): Promise { const [geozone] = await db.insert(gpsGeozones).values(data).returning(); return geozone; } async updateGeozone(id: number, organizationId: number, updates: Partial): Promise { 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 { 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 { 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> { 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> { 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 { 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> { 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(); 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 { 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 { 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 { 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 { 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 { 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 { 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 { 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 }> { 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, }; } async upsertSettings( organizationId: number, settings: { jitterMeters: number; maxJumpSpeedKmh: number; tripStopMinutes: number } ): Promise { 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();