From 0ed3a3160ae909824a3ae7678bae11d9a2398f08 Mon Sep 17 00:00:00 2001 From: Hamza-Ayed Date: Fri, 17 Jul 2026 02:49:45 +0300 Subject: [PATCH] =?UTF-8?q?feat:=20=D8=A7=D9=84=D9=85=D8=AC=D9=85=D9=88?= =?UTF-8?q?=D8=B9=D8=A9=20H=20=E2=80=94=20=D8=A7=D9=84=D9=85=D9=88=D8=A7?= =?UTF-8?q?=D9=82=D8=B9=20=D9=81=D9=8A=20Redis=D8=8C=20=D9=88=D8=A7=D9=84?= =?UTF-8?q?=D9=82=D8=A7=D8=B9=D8=AF=D8=A9=20=D9=85=D9=86=20=D8=A7=D9=84?= =?UTF-8?q?=D9=80worker=20=D9=81=D9=82=D8=B7?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit كان كل نبضة موقع (كل 1-3 ثوانٍ لكل سائق متصل) تكلّف استعلاماً + كتابة صف السائق في Postgres — أثقل حمل كان على القاعدة، وأثقل حتى من سيرو. - H1/H2: نبضة الموقع = Redis فقط، بعتبات سيرو (10م / 1.0 سرعة / 5° اتجاه). سائق واقف = EXPIRE واحد: لا كتابة، لا نقطة مسار، ولا بثّ للراكب - H3: لقطة آخر موقع على drivers (heading/speed/loc_status/loc_updated_at) يكتبها الـworker دورياً. لم نُنشئ جدولاً منفصلاً — drivers أصلاً صف واحد لكل سائق - H4: driver_tracks — تتراكم في قائمة Redis ويُدرجها الـworker بعبارة واحدة - H5: فهرسا available/busy — البحث يمسح المتاحين فقط. السائق يصير busy عند القبول ويعود available عند الإنهاء/الإلغاء - H8: المطابقة والاحتيال و«الرحلات المتاحة» تقرأ الموقع الحيّ من Redis. كشف تزوير الوصول خاصةً يجب ألّا يحكم بلقطة دورية الـworker صار حقيقياً (كان setInterval فارغاً): سياق Nest مستقل، تفريغ كل 5s، قفل يمنع تراكب الدورات، وتفريغة أخيرة عند SIGTERM. مؤجَّل من H: H6 (سلوك السائق) و H7 (ساعات العمل) — تحليلات تُبنى من driver_tracks لاحقاً. هجرة: DriverLocations. Co-Authored-By: Claude Opus 4.8 --- backend/src/app.module.ts | 2 + .../1721830000000-DriverLocations.ts | 41 ++++ .../src/modules/drivers/drivers.controller.ts | 14 +- backend/src/modules/drivers/drivers.module.ts | 4 +- .../src/modules/drivers/drivers.service.ts | 53 +++-- .../modules/drivers/entities/driver.entity.ts | 15 ++ backend/src/modules/fraud/fraud.service.ts | 15 +- .../locations/driver-location.service.spec.ts | 112 ++++++++++ .../locations/driver-location.service.ts | 203 ++++++++++++++++++ .../locations/entities/driver-track.entity.ts | 37 ++++ .../locations/location-flusher.service.ts | 124 +++++++++++ .../src/modules/locations/locations.module.ts | 18 ++ .../src/modules/matching/matching.service.ts | 56 +++-- backend/src/modules/trips/trips.module.ts | 2 + backend/src/modules/trips/trips.service.ts | 52 +++-- backend/src/realtime/realtime.gateway.ts | 22 +- backend/src/worker.module.ts | 35 +++ backend/src/worker.ts | 55 ++++- docs/17-backend-backlog.md | 32 +-- 19 files changed, 794 insertions(+), 98 deletions(-) create mode 100644 backend/src/database/migrations/1721830000000-DriverLocations.ts create mode 100644 backend/src/modules/locations/driver-location.service.spec.ts create mode 100644 backend/src/modules/locations/driver-location.service.ts create mode 100644 backend/src/modules/locations/entities/driver-track.entity.ts create mode 100644 backend/src/modules/locations/location-flusher.service.ts create mode 100644 backend/src/modules/locations/locations.module.ts create mode 100644 backend/src/worker.module.ts diff --git a/backend/src/app.module.ts b/backend/src/app.module.ts index d1f54e4..489dda2 100644 --- a/backend/src/app.module.ts +++ b/backend/src/app.module.ts @@ -14,6 +14,7 @@ import { UsersModule } from './modules/users/users.module'; import { AuthModule } from './modules/auth/auth.module'; import { DriversModule } from './modules/drivers/drivers.module'; import { MatchingModule } from './modules/matching/matching.module'; +import { LocationsModule } from './modules/locations/locations.module'; import { TariffModule } from './modules/tariff/tariff.module'; import { MapsModule } from './modules/maps/maps.module'; import { TripsModule } from './modules/trips/trips.module'; @@ -69,6 +70,7 @@ import { GeminiModule } from './integrations/gemini/gemini.module'; AuthModule, MapsModule, MatchingModule, + LocationsModule, TariffModule, RideTypesModule, DriversModule, diff --git a/backend/src/database/migrations/1721830000000-DriverLocations.ts b/backend/src/database/migrations/1721830000000-DriverLocations.ts new file mode 100644 index 0000000..7dab299 --- /dev/null +++ b/backend/src/database/migrations/1721830000000-DriverLocations.ts @@ -0,0 +1,41 @@ +import { MigrationInterface, QueryRunner } from 'typeorm'; + +/** + * المواقع (docs/17 — H3/H4): + * - أعمدة لقطة الموقع على `drivers` — يكتبها الـworker دورياً، لا كل نبضة. + * - `driver_tracks` — نقاط المسار التاريخية، إدراج مجمّع من الـworker. + */ +export class DriverLocations1721830000000 implements MigrationInterface { + public async up(q: QueryRunner): Promise { + await q.query(`ALTER TABLE tripz_drivers ADD COLUMN IF NOT EXISTS heading real`); + await q.query(`ALTER TABLE tripz_drivers ADD COLUMN IF NOT EXISTS speed real`); + await q.query(`ALTER TABLE tripz_drivers ADD COLUMN IF NOT EXISTS loc_status varchar`); + await q.query(`ALTER TABLE tripz_drivers ADD COLUMN IF NOT EXISTS loc_updated_at timestamptz`); + + await q.query(` + CREATE TABLE IF NOT EXISTS tripz_driver_tracks ( + id bigserial PRIMARY KEY, + tenant_id uuid NOT NULL, + driver_id uuid NOT NULL, + lat double precision NOT NULL, + lng double precision NOT NULL, + heading real, + speed real, + status varchar, + recorded_at timestamptz NOT NULL + ) + `); + await q.query(` + CREATE INDEX IF NOT EXISTS "IDX_tripz_driver_tracks_driver_time" + ON tripz_driver_tracks (tenant_id, driver_id, recorded_at) + `); + } + + public async down(q: QueryRunner): Promise { + await q.query(`DROP TABLE IF EXISTS tripz_driver_tracks`); + await q.query(`ALTER TABLE tripz_drivers DROP COLUMN IF EXISTS loc_updated_at`); + await q.query(`ALTER TABLE tripz_drivers DROP COLUMN IF EXISTS loc_status`); + await q.query(`ALTER TABLE tripz_drivers DROP COLUMN IF EXISTS speed`); + await q.query(`ALTER TABLE tripz_drivers DROP COLUMN IF EXISTS heading`); + } +} diff --git a/backend/src/modules/drivers/drivers.controller.ts b/backend/src/modules/drivers/drivers.controller.ts index c859250..f624466 100644 --- a/backend/src/modules/drivers/drivers.controller.ts +++ b/backend/src/modules/drivers/drivers.controller.ts @@ -29,10 +29,18 @@ export class DriversController { @Post('location') location( @CurrentUser() user: AuthUser, - @Body('lat') lat: number, - @Body('lng') lng: number, + @Body() body: { lat: number; lng: number; heading?: number; speed?: number }, ) { - return this.drivers.updateLocation(user.tenantId, user.userId, Number(lat), Number(lng)); + return this.drivers.updateLocation( + user.tenantId, + user.userId, + Number(body.lat), + Number(body.lng), + { + heading: body.heading == null ? undefined : Number(body.heading), + speed: body.speed == null ? undefined : Number(body.speed), + }, + ); } // للأدمن/المشغّل — تبسيط P1: أي مستخدم مصادَق؛ يُقيَّد بدور لاحقاً. diff --git a/backend/src/modules/drivers/drivers.module.ts b/backend/src/modules/drivers/drivers.module.ts index b511e9f..7b872b5 100644 --- a/backend/src/modules/drivers/drivers.module.ts +++ b/backend/src/modules/drivers/drivers.module.ts @@ -4,10 +4,10 @@ import { Driver } from './entities/driver.entity'; import { DriversService } from './drivers.service'; import { DriversController } from './drivers.controller'; import { UsersModule } from '../users/users.module'; -import { MatchingModule } from '../matching/matching.module'; +import { LocationsModule } from '../locations/locations.module'; @Module({ - imports: [TypeOrmModule.forFeature([Driver]), UsersModule, MatchingModule], + imports: [TypeOrmModule.forFeature([Driver]), UsersModule, LocationsModule], controllers: [DriversController], providers: [DriversService], exports: [DriversService], diff --git a/backend/src/modules/drivers/drivers.service.ts b/backend/src/modules/drivers/drivers.service.ts index e1ff14a..7abe992 100644 --- a/backend/src/modules/drivers/drivers.service.ts +++ b/backend/src/modules/drivers/drivers.service.ts @@ -3,7 +3,7 @@ import { InjectRepository } from '@nestjs/typeorm'; import { Repository } from 'typeorm'; import { Driver } from './entities/driver.entity'; import { UsersService } from '../users/users.service'; -import { MatchingService } from '../matching/matching.service'; +import { DriverLocationService } from '../locations/driver-location.service'; @Injectable() export class DriversService { @@ -11,7 +11,7 @@ export class DriversService { @InjectRepository(Driver) private readonly repo: Repository, private readonly users: UsersService, - private readonly matching: MatchingService, + private readonly locations: DriverLocationService, ) {} findByUser(tenantId: string, userId: string): Promise { @@ -62,7 +62,7 @@ export class DriversService { return this.repo.save(driver); } - /** يبدّل حالة الاتصال؛ عند الاتصال يُضاف لفهرس GEO، وعند الفصل يُزال. */ + /** يبدّل حالة الاتصال؛ عند الاتصال يدخل فهرس المتاحين، وعند الفصل يُزال. */ async setOnline( tenantId: string, userId: string, @@ -72,35 +72,42 @@ export class DriversService { if (!driver) throw new NotFoundException('Driver profile not found'); driver.is_online = online; await this.repo.save(driver); - if (!online) { - await this.matching.removeDriver(tenantId, driver.service_class, driver.id); - } else if (driver.last_lat != null && driver.last_lng != null) { - await this.matching.addDriver( - tenantId, - driver.service_class, - driver.id, - driver.last_lat, - driver.last_lng, - ); - } + await this.locations.setAvailability( + tenantId, + driver.id, + driver.service_class, + online ? 'available' : 'off', + ); return driver; } - /** تحديث الموقع الجاري: Redis GEO + لقطة في Postgres. */ + /** + * نبضة موقع — **Redis فقط** (docs/17 — H1). كانت تكتب صف السائق في Postgres + * على كل نبضة (كل ثانية/ثلاث لكل سائق متصل) — أثقل حمل كان على القاعدة. + * اللقطة الدائمة يكتبها الـworker دورياً. + */ async updateLocation( tenantId: string, userId: string, lat: number, lng: number, - ): Promise<{ ok: boolean }> { + extra: { heading?: number; speed?: number } = {}, + ): Promise<{ ok: boolean; significant: boolean }> { const driver = await this.findByUser(tenantId, userId); if (!driver) throw new NotFoundException('Driver profile not found'); - driver.last_lat = lat; - driver.last_lng = lng; - await this.repo.save(driver); - if (driver.is_online) { - await this.matching.addDriver(tenantId, driver.service_class, driver.id, lat, lng); - } - return { ok: true }; + if (!driver.is_online) return { ok: true, significant: false }; + + const res = await this.locations.update(tenantId, driver.id, driver.service_class, { + lat, + lng, + heading: extra.heading, + speed: extra.speed, + }); + return { ok: true, significant: res.significant }; + } + + /** الموقع الحيّ للسائق — من Redis، واللقطة احتياط. */ + liveLocation(tenantId: string, driverId: string) { + return this.locations.get(tenantId, driverId); } } diff --git a/backend/src/modules/drivers/entities/driver.entity.ts b/backend/src/modules/drivers/entities/driver.entity.ts index 67c3a67..fd355b3 100644 --- a/backend/src/modules/drivers/entities/driver.entity.ts +++ b/backend/src/modules/drivers/entities/driver.entity.ts @@ -54,12 +54,27 @@ export class Driver { @Column({ type: 'numeric', precision: 3, scale: 2, default: 5 }) rating: number; + // ---- لقطة الموقع (مكافئ `car_locations` عند سيرو) ---- + // الموقع الحيّ يعيش في Redis؛ هذه الأعمدة **يكتبها الـworker دورياً فقط**، + // لا الطلب. لا تقرأها للحيّ — استخدم DriverLocationService (docs/17 — H1/H3). @Column({ type: 'double precision', nullable: true }) last_lat: number; @Column({ type: 'double precision', nullable: true }) last_lng: number; + @Column({ type: 'real', nullable: true }) + heading: number | null; + + @Column({ type: 'real', nullable: true }) + speed: number | null; + + @Column({ type: 'varchar', nullable: true }) + loc_status: string | null; // available | busy | off + + @Column({ type: 'timestamptz', nullable: true }) + loc_updated_at: Date | null; + @CreateDateColumn() created_at: Date; diff --git a/backend/src/modules/fraud/fraud.service.ts b/backend/src/modules/fraud/fraud.service.ts index a11ad6a..b86b258 100644 --- a/backend/src/modules/fraud/fraud.service.ts +++ b/backend/src/modules/fraud/fraud.service.ts @@ -70,19 +70,20 @@ export class FraudService { } } - /** يتحقق من قرب السائق عند "وصل"؛ يبلّغ إن كان بعيداً. */ + /** + * يتحقق من قرب السائق عند "وصل"؛ يبلّغ إن كان بعيداً. + * `driverAt` هو الموقع **الحيّ** من Redis — لا لقطة القاعدة الدورية، وإلا + * حكمنا على تزوير الوصول بموقع عمره ثوانٍ (docs/17 — H1). + */ async checkArrivedProximity( tenantId: string, driverUserId: string, - driver: { last_lat?: number | null; last_lng?: number | null }, + driverAt: { lat: number; lng: number } | null, pickup: { lat: number; lng: number }, tripId: string, ) { - if (driver.last_lat == null || driver.last_lng == null) return; - const km = MapsService.haversineKm( - { lat: driver.last_lat, lng: driver.last_lng }, - pickup, - ); + if (!driverAt) return; + const km = MapsService.haversineKm(driverAt, pickup); if (km * 1000 > this.ARRIVED_MAX_M) { await this.flag( tenantId, diff --git a/backend/src/modules/locations/driver-location.service.spec.ts b/backend/src/modules/locations/driver-location.service.spec.ts new file mode 100644 index 0000000..9484f83 --- /dev/null +++ b/backend/src/modules/locations/driver-location.service.spec.ts @@ -0,0 +1,112 @@ +import RedisMock from 'ioredis-mock'; +import { DriverLocationService } from './driver-location.service'; +import { MatchingService } from '../matching/matching.service'; + +const TENANT = 't1'; +const DRIVER = 'd1'; +const CLASS = 'economy'; + +/** نقطة في وسط عمّان + إزاحة بالأمتار (تقريب كافٍ للاختبار). */ +const BASE = { lat: 31.9539, lng: 35.9106 }; +const metersToLat = (m: number) => m / 111_320; + +describe('DriverLocationService', () => { + let redis: any; + let matching: MatchingService; + let svc: DriverLocationService; + let driversRepo: any; + + beforeEach(() => { + redis = new RedisMock({ keyPrefix: 'tripz:' }); + matching = new MatchingService(redis); + driversRepo = { findOne: jest.fn().mockResolvedValue(null) }; + svc = new DriverLocationService(redis, driversRepo, matching); + }); + + it('أول نبضة ذات دلالة وتُخزَّن', async () => { + const r = await svc.update(TENANT, DRIVER, CLASS, BASE); + expect(r.significant).toBe(true); + + const pos = await svc.get(TENANT, DRIVER); + expect(pos!.lat).toBeCloseTo(BASE.lat); + expect(pos!.status).toBe('available'); + }); + + it('سائق واقف: النبضة التالية غير ذات دلالة ولا تضيف مساراً', async () => { + await svc.update(TENANT, DRIVER, CLASS, BASE); + const before = await redis.llen('loc:tracks'); + + // إزاحة مترين فقط — تحت عتبة 10م + const r = await svc.update(TENANT, DRIVER, CLASS, { + lat: BASE.lat + metersToLat(2), + lng: BASE.lng, + }); + + expect(r.significant).toBe(false); + expect(await redis.llen('loc:tracks')).toBe(before); + }); + + it('تحرّك أكثر من 10م = نبضة ذات دلالة + نقطة مسار', async () => { + await svc.update(TENANT, DRIVER, CLASS, BASE); + const before = await redis.llen('loc:tracks'); + + const r = await svc.update(TENANT, DRIVER, CLASS, { + lat: BASE.lat + metersToLat(25), + lng: BASE.lng, + }); + + expect(r.significant).toBe(true); + expect(await redis.llen('loc:tracks')).toBe(before + 1); + }); + + it('تغيّر السرعة وحده يكفي رغم الوقوف', async () => { + await svc.update(TENANT, DRIVER, CLASS, { ...BASE, speed: 0 }); + const r = await svc.update(TENANT, DRIVER, CLASS, { ...BASE, speed: 12 }); + expect(r.significant).toBe(true); + }); + + it('السائق «متّسخ» بعد نبضة ذات دلالة — ليكتبه الـworker لاحقاً', async () => { + await svc.update(TENANT, DRIVER, CLASS, BASE); + expect(await redis.smembers(DriverLocationService.DIRTY_KEY)).toContain(`${TENANT}|${DRIVER}`); + }); + + it('المتاح يظهر في المطابقة، والمشغول يختفي منها', async () => { + await svc.update(TENANT, DRIVER, CLASS, BASE); + expect(await matching.findNearby(TENANT, CLASS, BASE.lat, BASE.lng)).toHaveLength(1); + + await svc.setAvailability(TENANT, DRIVER, CLASS, 'busy'); + expect(await matching.findNearby(TENANT, CLASS, BASE.lat, BASE.lng)).toHaveLength(0); + + await svc.setAvailability(TENANT, DRIVER, CLASS, 'available'); + expect(await matching.findNearby(TENANT, CLASS, BASE.lat, BASE.lng)).toHaveLength(1); + }); + + it('الفصل (off) يزيله من الفهرس ومن Redis', async () => { + await svc.update(TENANT, DRIVER, CLASS, BASE); + await svc.setAvailability(TENANT, DRIVER, CLASS, 'off'); + + expect(await matching.findNearby(TENANT, CLASS, BASE.lat, BASE.lng)).toHaveLength(0); + expect(await svc.get(TENANT, DRIVER)).toBeNull(); + }); + + it('عند غياب المفتاح يرجع للقطة القاعدة', async () => { + driversRepo.findOne.mockResolvedValue({ + last_lat: 31.9, + last_lng: 35.9, + heading: 90, + speed: 0, + loc_status: 'available', + service_class: CLASS, + loc_updated_at: new Date(), + }); + + const pos = await svc.get(TENANT, 'cold-driver'); + expect(pos!.lat).toBe(31.9); + expect(driversRepo.findOne).toHaveBeenCalled(); + }); + + it('المستأجرون معزولون في المطابقة', async () => { + await svc.update('tenant-a', DRIVER, CLASS, BASE); + expect(await matching.findNearby('tenant-b', CLASS, BASE.lat, BASE.lng)).toHaveLength(0); + }); +}); diff --git a/backend/src/modules/locations/driver-location.service.ts b/backend/src/modules/locations/driver-location.service.ts new file mode 100644 index 0000000..82ceac6 --- /dev/null +++ b/backend/src/modules/locations/driver-location.service.ts @@ -0,0 +1,203 @@ +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 { MatchingService, DriverAvailability } from '../matching/matching.service'; +import { Driver } from '../drivers/entities/driver.entity'; +import { MapsService } from '../maps/maps.service'; + +/** الموقع الحيّ كما يعيش في Redis. */ +export interface LivePosition { + lat: number; + lng: number; + heading: number; + speed: number; + status: DriverAvailability; + serviceClass: string; + updatedAt: number; +} + +export interface LocationPing { + lat: number; + lng: number; + heading?: number; + speed?: number; + status?: DriverAvailability; +} + +/** + * عتبات مأخوذة من `driver_socket.php` عند سيرو — سائق واقف يرسل كل ثانية + * يجب ألّا يكلّف شيئاً تقريباً. + */ +const MIN_MOVE_M = 10; // أقل من هذا = لا GEOADD ولا نقطة مسار +const SPEED_DELTA = 1.0; +const HEADING_DELTA = 5; +const LOC_TTL_SEC = 300; // سائق صامت 5 دقائق يسقط من الفهرس +const MAX_TRACK_BUFFER = 50_000; // سقف أمان لطابور المسارات + +@Injectable() +export class DriverLocationService { + private readonly logger = new Logger('DriverLocation'); + + /** طابور المسارات ومجموعة «المتّسخين» — يستهلكهما الـworker (docs/17 H3/H4). */ + static readonly DIRTY_KEY = 'loc:dirty'; + static readonly TRACKS_KEY = 'loc:tracks'; + + constructor( + @Inject(REDIS) private readonly redis: Redis, + @InjectRepository(Driver) private readonly drivers: Repository, + private readonly matching: MatchingService, + ) {} + + private key(tenantId: string, driverId: string): string { + return `driver:loc:${tenantId}:${driverId}`; + } + + /** + * نبضة موقع من السائق — **Redis فقط، لا قاعدة** (docs/17 — H1/H2). + * يرجع `true` إن كانت النبضة ذات دلالة (تحرّك/تغيّر فعلي) — عندها فقط + * تُحدَّث الفهارس ويُسجَّل أثر، ويعرف المُنادي أن البثّ يستحق. + */ + async update( + tenantId: string, + driverId: string, + serviceClass: string, + ping: LocationPing, + ): Promise<{ significant: boolean; position: LivePosition }> { + const k = this.key(tenantId, driverId); + const prev = await this.readHash(k); + + const heading = ping.heading ?? prev?.heading ?? 0; + const speed = ping.speed ?? prev?.speed ?? 0; + const status = ping.status ?? prev?.status ?? 'available'; + + const movedM = prev + ? MapsService.haversineKm({ lat: prev.lat, lng: prev.lng }, { lat: ping.lat, lng: ping.lng }) * 1000 + : Infinity; + + const significant = + !prev || + movedM >= MIN_MOVE_M || + Math.abs(speed - prev.speed) >= SPEED_DELTA || + Math.abs(heading - prev.heading) >= HEADING_DELTA || + status !== prev.status; + + const position: LivePosition = { + lat: ping.lat, + lng: ping.lng, + heading, + speed, + status, + serviceClass, + updatedAt: Date.now(), + }; + + if (!significant) { + // واقف مكانه: نجدّد العمر فقط حتى لا يسقط من الفهرس. لا كتابة، لا مسار. + await this.redis.expire(k, LOC_TTL_SEC); + return { significant: false, position: prev! }; + } + + const pipe = this.redis.multi(); + pipe.hset(k, { + lat: String(position.lat), + lng: String(position.lng), + heading: String(heading), + speed: String(speed), + status, + serviceClass, + updatedAt: String(position.updatedAt), + }); + pipe.expire(k, LOC_TTL_SEC); + // «متّسخ» = يحتاج كتابة لقطة على القاعدة لاحقاً (مجموعة → لا تكرار). + pipe.sadd(DriverLocationService.DIRTY_KEY, `${tenantId}|${driverId}`); + // نقطة مسار — تُدرَج دفعةً واحدة من الـworker. + pipe.rpush( + DriverLocationService.TRACKS_KEY, + JSON.stringify({ + tenant_id: tenantId, + driver_id: driverId, + lat: position.lat, + lng: position.lng, + heading, + speed, + status, + at: position.updatedAt, + }), + ); + pipe.ltrim(DriverLocationService.TRACKS_KEY, -MAX_TRACK_BUFFER, -1); + await pipe.exec(); + + await this.matching.setPosition(tenantId, serviceClass, driverId, ping.lat, ping.lng, status); + return { significant: true, position }; + } + + /** + * الموقع الحيّ — Redis خط أول، واللقطة في القاعدة احتياط عند غياب المفتاح + * (سائق لم يرسل منذ 5 دقائق، أو Redis أُعيد تشغيله). + */ + async get(tenantId: string, driverId: string): Promise { + const live = await this.readHash(this.key(tenantId, driverId)); + if (live) return live; + + const row = await this.drivers.findOne({ + where: { tenant_id: tenantId, id: driverId }, + select: { + last_lat: true, + last_lng: true, + heading: true, + speed: true, + loc_status: true, + service_class: true, + loc_updated_at: true, + }, + }); + if (!row || row.last_lat == null || row.last_lng == null) return null; + return { + lat: Number(row.last_lat), + lng: Number(row.last_lng), + heading: Number(row.heading ?? 0), + speed: Number(row.speed ?? 0), + status: (row.loc_status ?? 'off') as DriverAvailability, + serviceClass: row.service_class, + updatedAt: row.loc_updated_at ? row.loc_updated_at.getTime() : 0, + }; + } + + /** يبدّل توفّر السائق (متاح/مشغول) وينقله بين الفهرسين (docs/17 — H5). */ + async setAvailability( + tenantId: string, + driverId: string, + serviceClass: string, + status: DriverAvailability, + ): Promise { + const k = this.key(tenantId, driverId); + const pos = await this.readHash(k); + if (status === 'off') { + await this.redis.del(k); + await this.matching.removeDriver(tenantId, serviceClass, driverId); + return; + } + await this.redis.multi().hset(k, 'status', status).expire(k, LOC_TTL_SEC).exec(); + if (pos) { + await this.matching.setPosition(tenantId, serviceClass, driverId, pos.lat, pos.lng, status); + } + } + + // ---- داخلي ---- + + private async readHash(k: string): Promise { + const h = await this.redis.hgetall(k); + if (!h?.lat) return null; + return { + lat: Number(h.lat), + lng: Number(h.lng), + heading: Number(h.heading ?? 0), + speed: Number(h.speed ?? 0), + status: (h.status ?? 'available') as DriverAvailability, + serviceClass: h.serviceClass ?? 'economy', + updatedAt: Number(h.updatedAt ?? 0), + }; + } +} diff --git a/backend/src/modules/locations/entities/driver-track.entity.ts b/backend/src/modules/locations/entities/driver-track.entity.ts new file mode 100644 index 0000000..d1abbf3 --- /dev/null +++ b/backend/src/modules/locations/entities/driver-track.entity.ts @@ -0,0 +1,37 @@ +import { Column, Entity, Index, PrimaryGeneratedColumn } from 'typeorm'; + +/** + * نقاط مسار السائق (مكافئ `car_tracks` عند سيرو) — للتتبع وفضّ النزاعات. + * الجدول: tripz_driver_tracks. **يُكتب دفعةً واحدة من الـworker، لا من الطلب** + * (docs/17 — H4). + */ +@Entity('driver_tracks') +@Index(['tenant_id', 'driver_id', 'recorded_at']) +export class DriverTrack { + @PrimaryGeneratedColumn('increment') + id: string; + + @Column({ type: 'uuid' }) + tenant_id: string; + + @Column({ type: 'uuid' }) + driver_id: string; + + @Column({ type: 'double precision' }) + lat: number; + + @Column({ type: 'double precision' }) + lng: number; + + @Column({ type: 'real', nullable: true }) + heading: number | null; + + @Column({ type: 'real', nullable: true }) + speed: number | null; + + @Column({ type: 'varchar', nullable: true }) + status: string | null; + + @Column({ type: 'timestamptz' }) + recorded_at: Date; +} diff --git a/backend/src/modules/locations/location-flusher.service.ts b/backend/src/modules/locations/location-flusher.service.ts new file mode 100644 index 0000000..9789306 --- /dev/null +++ b/backend/src/modules/locations/location-flusher.service.ts @@ -0,0 +1,124 @@ +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; + } + } +} diff --git a/backend/src/modules/locations/locations.module.ts b/backend/src/modules/locations/locations.module.ts new file mode 100644 index 0000000..f697125 --- /dev/null +++ b/backend/src/modules/locations/locations.module.ts @@ -0,0 +1,18 @@ +import { Module } from '@nestjs/common'; +import { TypeOrmModule } from '@nestjs/typeorm'; +import { Driver } from '../drivers/entities/driver.entity'; +import { DriverTrack } from './entities/driver-track.entity'; +import { DriverLocationService } from './driver-location.service'; +import { LocationFlusherService } from './location-flusher.service'; +import { MatchingModule } from '../matching/matching.module'; + +/** + * وحدة مستقلة عمداً: تستورد MatchingModule، وتستوردها DriversModule و + * TripsModule — لو عاشت داخل drivers لصارت دورة استيراد مع matching. + */ +@Module({ + imports: [TypeOrmModule.forFeature([Driver, DriverTrack]), MatchingModule], + providers: [DriverLocationService, LocationFlusherService], + exports: [DriverLocationService, LocationFlusherService], +}) +export class LocationsModule {} diff --git a/backend/src/modules/matching/matching.service.ts b/backend/src/modules/matching/matching.service.ts index 79bec18..9b29a62 100644 --- a/backend/src/modules/matching/matching.service.ts +++ b/backend/src/modules/matching/matching.service.ts @@ -7,41 +7,59 @@ export interface NearbyDriver { distanceKm: number; } +/** توفّر السائق — يحدّد الفهرس الذي يعيش فيه (docs/17 — H5). */ +export type DriverAvailability = 'available' | 'busy' | 'off'; + /** - * المطابقة الجغرافية عبر Redis GEO (docs/09). مفتاح لكل مستأجر: - * geo:drivers:{tenantId} (ببادئة tripz: تلقائياً من ioredis). - * نستخدم redis.call لتفادي تعقيد أنماط ioredis المتغيّرة. + * المطابقة الجغرافية عبر Redis GEO (docs/09). مفتاح لكل + * (مستأجر × فئة خدمة × توفّر): geo:drivers:{tenantId}:{class}:{available|busy} + * (ببادئة tripz: تلقائياً من ioredis). + * + * فصل available عن busy (نمط سيرو): البحث يمسح المتاحين فقط بدل أن يجلب + * الجميع ثم يصفّي المشغولين. */ @Injectable() export class MatchingService { constructor(@Inject(REDIS) private readonly redis: Redis) {} - // فهرس GEO لكل (مستأجر × فئة خدمة) — المطابقة تحترم نوع الرحلة. - private key(tenantId: string, serviceClass: string): string { - return `geo:drivers:${tenantId}:${serviceClass}`; + private key(tenantId: string, serviceClass: string, status: 'available' | 'busy'): string { + return `geo:drivers:${tenantId}:${serviceClass}:${status}`; } - async addDriver( + /** + * يضع السائق في الفهرس الموافق لحالته ويزيله من الآخر — عملية واحدة + * تمنع بقاءه في الفهرسين معاً. + */ + async setPosition( tenantId: string, serviceClass: string, driverId: string, lat: number, lng: number, - ) { - await this.redis.call( - 'GEOADD', - this.key(tenantId, serviceClass), - String(lng), - String(lat), - driverId, - ); + status: DriverAvailability, + ): Promise { + if (status === 'off') { + await this.removeDriver(tenantId, serviceClass, driverId); + return; + } + const target = this.key(tenantId, serviceClass, status); + const other = this.key(tenantId, serviceClass, status === 'available' ? 'busy' : 'available'); + await this.redis + .multi() + .geoadd(target, lng, lat, driverId) + .zrem(other, driverId) + .exec(); } - async removeDriver(tenantId: string, serviceClass: string, driverId: string) { - await this.redis.call('ZREM', this.key(tenantId, serviceClass), driverId); + async removeDriver(tenantId: string, serviceClass: string, driverId: string): Promise { + await this.redis + .multi() + .zrem(this.key(tenantId, serviceClass, 'available'), driverId) + .zrem(this.key(tenantId, serviceClass, 'busy'), driverId) + .exec(); } - /** أقرب السائقين من نفس فئة الخدمة ضمن نصف قطر (كم). */ + /** أقرب السائقين **المتاحين** من نفس فئة الخدمة ضمن نصف قطر (كم). */ async findNearby( tenantId: string, serviceClass: string, @@ -52,7 +70,7 @@ export class MatchingService { ): Promise { const res = (await this.redis.call( 'GEOSEARCH', - this.key(tenantId, serviceClass), + this.key(tenantId, serviceClass, 'available'), 'FROMLONLAT', String(lng), String(lat), diff --git a/backend/src/modules/trips/trips.module.ts b/backend/src/modules/trips/trips.module.ts index 383348e..009784e 100644 --- a/backend/src/modules/trips/trips.module.ts +++ b/backend/src/modules/trips/trips.module.ts @@ -13,6 +13,7 @@ import { RealtimeModule } from '../../realtime/realtime.module'; import { FraudModule } from '../fraud/fraud.module'; import { WalletModule } from '../wallet/wallet.module'; import { UsersModule } from '../users/users.module'; +import { LocationsModule } from '../locations/locations.module'; @Module({ imports: [ @@ -25,6 +26,7 @@ import { UsersModule } from '../users/users.module'; FraudModule, WalletModule, UsersModule, + LocationsModule, ], controllers: [TripsController], providers: [TripsService, TripStateService], diff --git a/backend/src/modules/trips/trips.service.ts b/backend/src/modules/trips/trips.service.ts index 6a30a6a..b971841 100644 --- a/backend/src/modules/trips/trips.service.ts +++ b/backend/src/modules/trips/trips.service.ts @@ -19,6 +19,7 @@ import { WalletService } from '../wallet/wallet.service'; import { NotificationsService } from '../notifications/notifications.service'; import { TripState, TripStateService } from './trip-state.service'; import { UsersService } from '../users/users.service'; +import { DriverLocationService } from '../locations/driver-location.service'; /** رسوم الإلغاء حسب مرحلة الرحلة (docs/04). */ const CANCEL_FEE_BY_STAGE: Partial> = { @@ -71,6 +72,7 @@ export class TripsService { private readonly notifications: NotificationsService, private readonly state: TripStateService, private readonly users: UsersService, + private readonly locations: DriverLocationService, ) {} get(tenantId: string, id: string): Promise { @@ -182,6 +184,9 @@ export class TripsService { } const trip = await this.getOr404(tenantId, tripId); + // مشغول الآن: يخرج من فهرس المتاحين فلا تُعرض عليه رحلة أخرى (H5). + await this.locations.setAvailability(tenantId, driver.id, driver.service_class, 'busy'); + await this.recordEvent(trip, 'searching', 'assigned', 'driver'); const driverUser = await this.users.findById(tenantId, driverUserId); await this.notifyParties(tenantId, trip, 'assigned', driverUserId, { @@ -212,18 +217,16 @@ export class TripsService { const patch: Partial = { status: toStatus }; - // كشف احتيال عند نقاط حسّاسة - if (toStatus === 'driver_arrived' && snapshot.driver_id) { - const drv = await this.drivers.findById(tenantId, snapshot.driver_id); - if (drv) { - await this.fraud.checkArrivedProximity( - tenantId, - drv.user_id, - { last_lat: drv.last_lat, last_lng: drv.last_lng }, - { lat: snapshot.origin_lat, lng: snapshot.origin_lng }, - tripId, - ); - } + // كشف احتيال عند نقاط حسّاسة — الموقع الحيّ من Redis، لا صف السائق. + if (toStatus === 'driver_arrived' && snapshot.driver_id && snapshot.driver_user_id) { + const at = await this.locations.get(tenantId, snapshot.driver_id); + await this.fraud.checkArrivedProximity( + tenantId, + snapshot.driver_user_id, + at ? { lat: at.lat, lng: at.lng } : null, + { lat: snapshot.origin_lat, lng: snapshot.origin_lng }, + tripId, + ); } if (toStatus === 'completed') { patch.completed_at = new Date(); @@ -249,7 +252,7 @@ export class TripsService { } const trip = await this.getOr404(tenantId, tripId); - await this.syncState(trip, snapshot.driver_user_id, toStatus); + await this.syncState(trip, snapshot, toStatus); await this.recordEvent(trip, snapshot.status, toStatus, actor); // تسوية المحفظة عند الدفع (payment_method === wallet) @@ -296,7 +299,7 @@ export class TripsService { if (!res.affected) throw new BadRequestException('Trip is no longer cancellable'); const trip = await this.getOr404(tenantId, tripId); - await this.syncState(trip, snapshot.driver_user_id, 'cancelled'); + await this.syncState(trip, snapshot, 'cancelled'); await this.recordEvent(trip, snapshot.status, 'cancelled', actor); // الإلغاء وهي searching: العرض معلّق عند سائقين — أزِله. @@ -322,7 +325,9 @@ export class TripsService { if (driver.verification_status !== 'approved') { throw new ForbiddenException('Driver not approved'); } - if (driver.last_lat == null || driver.last_lng == null) return []; + // موقعه الحيّ من Redis — لقطة القاعدة قد تتأخّر ثوانٍ عن الواقع. + const at = await this.locations.get(tenantId, driver.id); + if (!at) return []; const open = await this.trips.find({ where: { @@ -337,7 +342,7 @@ export class TripsService { return open .map((t) => ({ trip: t, - distanceKm: haversineKm(driver.last_lat!, driver.last_lng!, t.origin_lat, t.origin_lng), + distanceKm: haversineKm(at.lat, at.lng, t.origin_lat, t.origin_lng), })) .filter((r) => r.distanceKm <= radiusKm) .sort((a, b) => a.distanceKm - b.distanceKm) @@ -390,11 +395,22 @@ export class TripsService { } } - private async syncState(trip: Trip, driverUserId: string | null, to: TripStatus) { + private async syncState(trip: Trip, snapshot: TripState, to: TripStatus) { if (TERMINAL.includes(to)) { await this.state.clear(trip.tenant_id, trip.id); } else { - await this.state.save(trip, driverUserId); + await this.state.save(trip, snapshot.driver_user_id); + } + + // انتهى دور السائق (أنهى أو أُلغيت) → يعود لفهرس المتاحين (docs/17 — H5). + // 'completed' يحرّره وإن لم يُدفع بعد: القيادة انتهت فعلاً. + if (snapshot.driver_id && (to === 'completed' || TERMINAL.includes(to))) { + await this.locations.setAvailability( + trip.tenant_id, + snapshot.driver_id, + snapshot.service_class, + 'available', + ); } } diff --git a/backend/src/realtime/realtime.gateway.ts b/backend/src/realtime/realtime.gateway.ts index 2d150d8..e9436f4 100644 --- a/backend/src/realtime/realtime.gateway.ts +++ b/backend/src/realtime/realtime.gateway.ts @@ -57,20 +57,34 @@ export class RealtimeGateway implements OnGatewayConnection { return { ok: true }; } - /** موقع السائق الحي: يحدّث Redis GEO ويبثّ لغرفة الرحلة. */ + /** + * موقع السائق الحي: Redis فقط (لا قاعدة — docs/17 H1)، والبثّ لغرفة الرحلة + * **فقط عند حركة ذات دلالة**؛ سائق واقف يرسل كل ثانية لا يُغرق الراكب + * بتحديثات متطابقة. + */ @SubscribeMessage('driver:location') async onDriverLocation( @ConnectedSocket() c: Socket, - @MessageBody() body: { lat: number; lng: number; tripId?: string }, + @MessageBody() + body: { lat: number; lng: number; heading?: number; speed?: number; tripId?: string }, ) { const u = c.data.user; if (!u || u.role !== 'driver' || body?.lat == null || body?.lng == null) return; - await this.drivers.updateLocation(u.tenantId, u.userId, body.lat, body.lng); - if (body.tripId) { + + const { significant } = await this.drivers.updateLocation( + u.tenantId, + u.userId, + body.lat, + body.lng, + { heading: body.heading, speed: body.speed }, + ); + if (body.tripId && significant) { c.to(this.tripRoom(u.tenantId, body.tripId)).emit('driver:location', { tripId: body.tripId, lat: body.lat, lng: body.lng, + heading: body.heading ?? 0, + speed: body.speed ?? 0, }); } return { ok: true }; diff --git a/backend/src/worker.module.ts b/backend/src/worker.module.ts new file mode 100644 index 0000000..33756dd --- /dev/null +++ b/backend/src/worker.module.ts @@ -0,0 +1,35 @@ +import { Module } from '@nestjs/common'; +import { ConfigModule, ConfigService } from '@nestjs/config'; +import { TypeOrmModule } from '@nestjs/typeorm'; +import configuration from './config/configuration'; +import { RedisModule } from './common/redis/redis.module'; +import { LocationsModule } from './modules/locations/locations.module'; + +/** + * وحدة الـworker — عملية منفصلة (`node dist/worker.js`) لا تخدم HTTP. + * تحمل ما يحتاجه العمل الدوري فقط: قاعدة + Redis + المواقع. + */ +@Module({ + imports: [ + ConfigModule.forRoot({ isGlobal: true, load: [configuration] }), + TypeOrmModule.forRootAsync({ + inject: [ConfigService], + useFactory: (cfg: ConfigService) => ({ + type: 'postgres' as const, + host: cfg.get('db.host'), + port: cfg.get('db.port'), + database: cfg.get('db.name'), + username: cfg.get('db.user'), + password: cfg.get('db.password'), + // بركة أصغر من الـAPI: الـworker يكتب دفعات قليلة لا طلبات متزامنة. + extra: { max: parseInt(process.env.WORKER_DB_POOL_MAX ?? '5', 10) }, + entityPrefix: cfg.get('db.tablePrefix'), + synchronize: false, + autoLoadEntities: true, + }), + }), + RedisModule, + LocationsModule, + ], +}) +export class WorkerModule {} diff --git a/backend/src/worker.ts b/backend/src/worker.ts index 3c68259..a2e5eec 100644 --- a/backend/src/worker.ts +++ b/backend/src/worker.ts @@ -1,17 +1,56 @@ import 'reflect-metadata'; import { Logger } from '@nestjs/common'; +import { NestFactory } from '@nestjs/core'; +import { WorkerModule } from './worker.module'; +import { LocationFlusherService } from './modules/locations/location-flusher.service'; /** - * نقطة دخول الـ worker (BullMQ) — مهام غير متزامنة: انتهاء صلاحية العروض، - * الإشعارات، التسويات، تجميع usage للفوترة (راجع docs/09). - * حالياً هيكل فقط؛ المعالِجات تُضاف في P1. + * نقطة دخول الـ worker — مهام غير متزامنة خارج مسار الطلب (docs/09، docs/17 H3/H4). + * حالياً: تفريغ لقطات المواقع ونقاط المسار من Redis إلى القاعدة دورياً. */ +const FLUSH_INTERVAL_MS = parseInt(process.env.LOC_FLUSH_INTERVAL_MS ?? '5000', 10); + async function bootstrap() { const log = new Logger('Worker'); - const prefix = process.env.QUEUE_PREFIX ?? 'tripz_'; - const redisDb = process.env.REDIS_DB ?? '3'; - log.log(`Tripz worker up. queuePrefix=${prefix} redisDb=${redisDb}`); - // TODO(P1): سجّل معالِجات BullMQ هنا. - setInterval(() => void 0, 1 << 30); + const app = await NestFactory.createApplicationContext(WorkerModule, { + logger: ['error', 'warn', 'log'], + }); + const flusher = app.get(LocationFlusherService); + + log.log( + `Tripz worker up. locationFlush=${FLUSH_INTERVAL_MS}ms ` + + `queuePrefix=${process.env.QUEUE_PREFIX ?? 'tripz_'} redisDb=${process.env.REDIS_DB ?? '3'}`, + ); + + // قفل بسيط: دورة لا تبدأ قبل أن تنتهي سابقتها — دفعة بطيئة يجب ألّا تتراكب + // مع التالية فتُدرَج النقاط مرتين أو تُغرق القاعدة. + let running = false; + const timer = setInterval(async () => { + if (running) return; + running = true; + try { + await flusher.flush(); + } catch (e: any) { + log.error(`flush cycle failed: ${e?.message}`); + } finally { + running = false; + } + }, FLUSH_INTERVAL_MS); + + const shutdown = async (signal: string) => { + log.log(`${signal} — draining…`); + clearInterval(timer); + // تفريغة أخيرة حتى لا تضيع اللقطات المعلّقة عند إعادة النشر. + try { + await flusher.flush(); + } catch (e: any) { + log.warn(`final flush failed: ${e?.message}`); + } + await app.close(); + process.exit(0); + }; + process.on('SIGTERM', () => void shutdown('SIGTERM')); + process.on('SIGINT', () => void shutdown('SIGINT')); } + bootstrap(); diff --git a/docs/17-backend-backlog.md b/docs/17-backend-backlog.md index b592eda..d96c54e 100644 --- a/docs/17-backend-backlog.md +++ b/docs/17-backend-backlog.md @@ -97,9 +97,9 @@ --- -## المجموعة H — المواقع والتتبع (من `loction_server` في سيرو) +## المجموعة H — المواقع والتتبع (من `loction_server` في سيرو) — ✅ منفَّذة > ملاحظة المالك: «لحد الآن مش شايف السائق أو الراكب يرفع موقعه، ولا جداول location». -> عندنا الآن: `driver:location` عبر WebSocket → Redis GEO + **حفظ صف السائق في Postgres على كل نبضة** (`drivers.updateLocation` يعمل `repo.save`) — هذا أثقل حتى من سيرو. +> كان عندنا: `driver:location` عبر WebSocket → Redis GEO + **حفظ صف السائق في Postgres على كل نبضة** (`drivers.updateLocation` يعمل `repo.save`) — أثقل حتى من سيرو. **ما يفعله سيرو (مقروء من الكود):** - `driver_socket.php` — سوكيت مخصص للمواقع، **لا يلمس القاعدة إطلاقاً**؛ Redis فقط عبر **pipeline كل 500ms** (`REDIS_BATCH_INTERVAL`). @@ -108,17 +108,21 @@ - `driver:profile:{id}` hash في Redis (heading/speed/status) — المطابقة تقرأ منه بلا قاعدة. - عروض الرحلة: `setex` للعرض + `sadd` لمجموعة المعروض عليهم + `expire` — **نفس نمطنا في A5** ✅. -| # | البند | التفصيل | -|---|-------|---------| -| H1 | **إيقاف كتابة الموقع على القاعدة لكل نبضة** | `drivers.updateLocation` يكتب Postgres كل ثانية/ثلاث — يُنقل إلى Redis فقط. | -| H2 | **عتبات + batching** | تبنّي عتبات سيرو (10م/1.0 سرعة/5 اتجاه) + pipeline كل 500ms. | -| H3 | **جدول `car_locations` مكافئ** | صف واحد لكل سائق (آخر موقع) — يُكتب دورياً من worker لا من الطلب. عند سيرو: `ON DUPLICATE KEY UPDATE` + عمود `point` SRID 4326 عبر trigger (عندنا PostGIS متاح). | -| H4 | **جدول `car_tracks` (المسار)** | سجل نقاط تاريخي للتتبع/النزاعات — إدراج مجمّع (batch insert) من worker. | -| H5 | **فهرس available/busy** | تمييز السائق المشغول عن المتاح في فهرس GEO. | -| H6 | **`driver_behavior`** | max_speed · avg_speed · hard_brakes · total_distance · behavior_score لكل رحلة. | -| H7 | **`driver_daily_work` / `driver_daily_summary`** | ساعات عمل السائق (total_seconds باليوم + last_point_at + last_status). | -| H8 | **ربط المواقع بالمطابقة والـdispatch** | المطابقة تقرأ `driver:profile` من Redis؛ لوحة dispatch ترى الأسطول حياً. | -| — | **Geofence** | مؤجَّل بقرار المالك («خليها لوقتها») — موجود في سيرو (`get_location_area_links`, `LocationIntelligenceEngine`). | +| # | البند | التفصيل | الحالة | +|---|-------|---------|--------| +| H1 | **إيقاف كتابة الموقع على القاعدة لكل نبضة** | `drivers.updateLocation` كان يكتب Postgres كل ثانية/ثلاث لكل سائق متصل. | ✅ Redis فقط. الكتابة الدائمة صارت من الـworker | +| H2 | **عتبات + batching** | عتبات سيرو: 10م حركة · 1.0 سرعة · 5° اتجاه. سائق واقف = `EXPIRE` فقط (بلا كتابة ولا نقطة مسار ولا بثّ). | ✅ + البثّ للراكب صار عند الحركة ذات الدلالة فقط | +| H3 | **لقطة «آخر موقع» لكل سائق** | صف واحد لكل سائق يكتبه الـworker دورياً. **قرار: لم نُنشئ جدولاً منفصلاً** — `drivers` أصلاً صف واحد لكل سائق؛ أُضيفت `heading`/`speed`/`loc_status`/`loc_updated_at` بجانب `last_lat/last_lng`. | ✅ | +| H4 | **`driver_tracks` (المسار)** | نقاط تاريخية للتتبع/النزاعات — تتراكم في قائمة Redis ويُدرجها الـworker **بعبارة واحدة**. | ✅ جدول `tripz_driver_tracks` | +| H5 | **فهرس available/busy** | `geo:drivers:{tenant}:{class}:{available\|busy}` — البحث يمسح المتاحين فقط. السائق يصير busy عند القبول ويعود available عند الإنهاء/الإلغاء. | ✅ | +| H8 | **ربط المواقع بالمطابقة** | المطابقة والاحتيال و«الرحلات المتاحة» صاروا يقرؤون الموقع الحيّ من Redis. | ✅ | +| H6 | **`driver_behavior`** | max_speed · avg_speed · hard_brakes · total_distance · behavior_score لكل رحلة. | ⏳ **مؤجَّل** — تحليلات فوق المسار، تُبنى من `driver_tracks` لاحقاً | +| H7 | **`driver_daily_work` / `driver_daily_summary`** | ساعات عمل السائق (total_seconds باليوم + last_point_at + last_status). | ⏳ **مؤجَّل** | +| — | **Geofence** | مؤجَّل بقرار المالك («خليها لوقتها») — موجود في سيرو (`get_location_area_links`, `LocationIntelligenceEngine`). | ⏳ | + +**الـworker صار حقيقياً**: كان هيكلاً فارغاً (`setInterval` بلا عمل). الآن `NestFactory.createApplicationContext(WorkerModule)` — قاعدة + Redis + المواقع، يفرّغ كل 5s (`LOC_FLUSH_INTERVAL_MS`)، بقفل يمنع تراكب الدورات، وتفريغة أخيرة عند SIGTERM حتى لا تضيع اللقطات عند النشر. + +**الأثر**: سائق متصل كان يكلّف استعلاماً + كتابة قاعدة لكل نبضة (كل 1–3 ثوانٍ). الآن: نبضة السائق الواقف = أمر `EXPIRE` واحد، والمتحرّك = كتابة Redis + كتابة قاعدة واحدة كل 5 ثوانٍ **مهما بلغ عدد نبضاته**. --- @@ -182,7 +186,7 @@ 1. ~~**A** (الزمن الحقيقي + FCM + Redis + race)~~ ✅ **منفَّذة ومُثبَتة على السيرفر**. 2. ~~**I1** (سباق المحفظة)~~ ✅ **منفَّذة ومُثبَتة بتزامن حقيقي**. 3. ~~**G** (Redis خط أول: لغة/توكنات/tenant/تعرفة/تقييم/كاش الخرائط)~~ ✅ **منفَّذة**. -4. **H** (المواقع: إيقاف الكتابة لكل نبضة + batching + tracks) — **التالي، وهو الحمل الأثقل فعلياً**. +4. ~~**H** (المواقع: إيقاف الكتابة لكل نبضة + batching + tracks)~~ ✅ **منفَّذة** (عدا H6/H7 — تحليلات مؤجَّلة). 5. **B** (نموذج الرحلة: started/انتظار/فصل السعر/الطوابع/stops). 6. **I** الباقي (المدفوعات: جداول لكل طريقة + OTP/بصمة/HMAC للسحب). 7. **C** (بيانات المركبة والسائق + ai_data).