Files
tripz-llc/backend/src/worker.ts
T
Hamza-AyedandClaude Opus 4.8 0ed3a3160a feat: المجموعة H — المواقع في Redis، والقاعدة من الـworker فقط
كان كل نبضة موقع (كل 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 <noreply@anthropic.com>
2026-07-17 02:49:45 +03:00

57 lines
2.0 KiB
TypeScript

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 — مهام غير متزامنة خارج مسار الطلب (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 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();