import { Inject, Injectable, Logger } from '@nestjs/common'; import { InjectRepository } from '@nestjs/typeorm'; import { Repository } from 'typeorm'; import Redis from 'ioredis'; import { REDIS } from '../../common/redis/redis.module'; import { Driver } from '../drivers/entities/driver.entity'; import { DriverTrack } from './entities/driver-track.entity'; import { DriverLocationService } from './driver-location.service'; interface TrackPoint { tenant_id: string; driver_id: string; lat: number; lng: number; heading: number; speed: number; status: string; at: number; } /** حدود الدفعة — تحمي القاعدة من دفعة ضخمة بعد انقطاع طويل. */ const MAX_DIRTY_PER_RUN = 500; const MAX_TRACKS_PER_RUN = 2_000; /** * يفرّغ ما تراكم في Redis إلى القاعدة (docs/17 — H3/H4). يعمل في الـworker * على فترات، **لا في مسار الطلب**. * * الفكرة: آلاف النبضات تُكتب مرة واحدة لكل سائق في الدورة، ونقاط المسار * تُدرَج بعبارة واحدة بدل عبارة لكل نقطة. */ @Injectable() export class LocationFlusherService { private readonly logger = new Logger('LocationFlusher'); constructor( @Inject(REDIS) private readonly redis: Redis, @InjectRepository(Driver) private readonly drivers: Repository, @InjectRepository(DriverTrack) private readonly tracks: Repository, ) {} async flush(): Promise<{ snapshots: number; tracks: number }> { const [snapshots, tracksWritten] = await Promise.all([ this.flushSnapshots(), this.flushTracks(), ]); if (snapshots || tracksWritten) { this.logger.debug(`flushed snapshots=${snapshots} tracks=${tracksWritten}`); } return { snapshots, tracks: tracksWritten }; } /** لقطة آخر موقع لكل سائق «متّسخ» — صف واحد لكل سائق مهما كثرت نبضاته. */ private async flushSnapshots(): Promise { // SPOP يسحب ويحذف ذرّياً: نسخة worker أخرى لن تعالج نفس السائق. const members = await this.redis.spop(DriverLocationService.DIRTY_KEY, MAX_DIRTY_PER_RUN); const list = Array.isArray(members) ? members : members ? [members] : []; if (list.length === 0) return 0; let written = 0; for (const member of list) { const [tenantId, driverId] = String(member).split('|'); if (!tenantId || !driverId) continue; try { const h = await this.redis.hgetall(`driver:loc:${tenantId}:${driverId}`); if (!h?.lat) continue; // انتهى عمر المفتاح — لا شيء نكتبه await this.drivers.update( { tenant_id: tenantId, id: driverId }, { last_lat: Number(h.lat), last_lng: Number(h.lng), heading: Number(h.heading ?? 0), speed: Number(h.speed ?? 0), loc_status: h.status ?? null, loc_updated_at: new Date(Number(h.updatedAt ?? Date.now())), }, ); written++; } catch (e: any) { // إعادة السائق للمجموعة كي لا تضيع لقطته بسبب عطل عابر. await this.redis.sadd(DriverLocationService.DIRTY_KEY, member); this.logger.warn(`snapshot ${member} failed: ${e?.message}`); } } return written; } /** نقاط المسار — إدراج مجمّع بعبارة واحدة. */ private async flushTracks(): Promise { const raw = await this.redis.lpop(DriverLocationService.TRACKS_KEY, MAX_TRACKS_PER_RUN); const list = Array.isArray(raw) ? raw : raw ? [raw] : []; if (list.length === 0) return 0; const rows: Partial[] = []; for (const item of list) { try { const p: TrackPoint = JSON.parse(item); rows.push({ tenant_id: p.tenant_id, driver_id: p.driver_id, lat: p.lat, lng: p.lng, heading: p.heading, speed: p.speed, status: p.status, recorded_at: new Date(p.at), }); } catch { this.logger.warn('skipped malformed track point'); } } if (rows.length === 0) return 0; try { await this.tracks.insert(rows); return rows.length; } catch (e: any) { // نُعيدها لرأس الطابور حتى لا يضيع المسار. await this.redis.lpush(DriverLocationService.TRACKS_KEY, ...list); this.logger.warn(`track insert failed, requeued ${list.length}: ${e?.message}`); return 0; } } }