Исправлено дублирование сообщений в task-чате: защита от двойного Enter/клика по кнопке отправки
This commit is contained in:
@@ -1,443 +0,0 @@
|
||||
import { db } from "../db";
|
||||
import {
|
||||
gpsAssets,
|
||||
gpsPositions,
|
||||
gpsGeozones,
|
||||
gpsGeozoneEvents,
|
||||
gpsGeozoneStates,
|
||||
gpsGroups,
|
||||
gpsGroupAssets,
|
||||
gpsAssetSubscribers,
|
||||
type GpsAsset,
|
||||
type GpsPosition,
|
||||
type GpsGeozone,
|
||||
type GpsGeozoneEvent,
|
||||
type GpsGroup,
|
||||
type GpsAssetSubscriber,
|
||||
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) */
|
||||
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: data.recordedAt,
|
||||
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();
|
||||
}
|
||||
}
|
||||
|
||||
export const gpsStorage = new GpsStorage();
|
||||
Reference in New Issue
Block a user