feat: المجموعة G — Redis خط أول والقاعدة خط احتياط
الأساس: common/cache — كاش-جانبي فوق Redis. مبدأ صارم: فشل Redis لا يُسقط الطلب (يُسجَّل ويُرجَع للقاعدة) — الكاش تحسين أداء لا مصدر حقيقة. - G1: لغة المستخدم في Redis (كانت استعلاماً قبل كل إشعار)؛ تُكتب عند PATCH /users/me فيصير الإصابة دائمة - G2/G3: توكنات FCM في Redis — يُبطَل عند register وعند اكتشاف توكن ميت. البث الجماعي صار بلا استعلامات قاعدة - G4: كاش tenant resolve/config + إبطال عند الإنشاء - G5: كاش التعرفة و ride-types + إبطال صريح عند أي إنشاء - G6: التقييم يتراكم في Redis، والقاعدة تُكتب مرة واحدة يومياً لكل طرف (حجز الكتابة ذرّي بـLua). التأخير مقصود ليبعد الاحتكاك بين الطرفين. أُضيف ratings.target_user_id و users.rating — تقييم السائق للراكب كان يُخزَّن بلا هدف فلا يُجمَّع أبداً - كاش الخرائط (route/reverse/geocode) بإحداثيات مقرَّبة 4 خانات (~11م): نداء انطلق كان يهيمن على p50 (2.9s). لا يُخزَّن الرجوع لخط مستقيم ولا استجابة فاشلة تصحيح: ادّعيت سابقاً أن tenants.resolve يعمل على كل request — خطأ. tenant_id يأتي من الـJWT؛ resolve يُستدعى عند الدخول و3 كنترولرات فقط. الحمل الأثقل فعلياً هو رفع موقع السائق (المجموعة H). هجرة: RatingsBothParties. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.8
parent
857652404b
commit
a83dd7e5e5
@@ -1,5 +1,6 @@
|
||||
import { Injectable, Logger } from '@nestjs/common';
|
||||
import { ConfigService } from '@nestjs/config';
|
||||
import { CacheService, CacheKeys, TTL } from '../../common/cache/cache.service';
|
||||
|
||||
export interface LatLng {
|
||||
lat: number;
|
||||
@@ -21,7 +22,18 @@ export class MapsService {
|
||||
private readonly logger = new Logger('Maps');
|
||||
private readonly avgSpeedKmh = 30;
|
||||
|
||||
constructor(private readonly config: ConfigService) {}
|
||||
constructor(
|
||||
private readonly config: ConfigService,
|
||||
private readonly cache: CacheService,
|
||||
) {}
|
||||
|
||||
/**
|
||||
* يقرّب الإحداثيات لـ4 خانات (~11م) لصنع مفتاح كاش.
|
||||
* بلا تقريب لا يتكرر مفتاح أبداً (كل GPS فريد) فيصير الكاش بلا فائدة.
|
||||
*/
|
||||
private static coordKey(lat: number, lng: number): string {
|
||||
return `${lat.toFixed(4)},${lng.toFixed(4)}`;
|
||||
}
|
||||
|
||||
private get base(): string {
|
||||
return this.config.get<string>('maps.baseUrl') ?? 'https://map-saas.intaleqapp.com';
|
||||
@@ -36,50 +48,87 @@ export class MapsService {
|
||||
return { 'x-api-key': key, 'Content-Type': 'application/json' };
|
||||
}
|
||||
|
||||
/** توجيه حقيقي من انطلق مع رجوع آمن لخط مستقيم. */
|
||||
/**
|
||||
* توجيه حقيقي من انطلق مع رجوع آمن لخط مستقيم.
|
||||
*
|
||||
* كاش Redis (docs/17 — G): نداء انطلق الخارجي هو المهيمن على زمن الاستجابة
|
||||
* (اختبار التحميل: p50 = 2.9s). المسار بين نقطتين لا يتغيّر عملياً، فنخزّنه.
|
||||
* **لا يُخزَّن الرجوع لخط مستقيم** — وإلا ثبّتنا تقديراً رديئاً ليوم كامل
|
||||
* بينما قد يكون انطلق قد عاد للعمل بعد ثوانٍ.
|
||||
*/
|
||||
async route(from: LatLng, to: LatLng, country?: string): Promise<RouteResult> {
|
||||
const key = this.keyFor(country);
|
||||
if (key) {
|
||||
try {
|
||||
const url =
|
||||
`${this.base}/api/maps/route?fromLat=${from.lat}&fromLng=${from.lng}` +
|
||||
`&toLat=${to.lat}&toLng=${to.lng}&steps=false&locale=ar`;
|
||||
const res = await fetch(url, { headers: this.headers(key) });
|
||||
if (res.ok) {
|
||||
const json: any = await res.json();
|
||||
const parsed = MapsService.parseRoute(json);
|
||||
if (parsed) return { ...parsed, provider: 'antlaq' };
|
||||
} else {
|
||||
this.logger.warn(`antlaq route ${res.status} — fallback`);
|
||||
if (!key) return this.straightLine(from, to);
|
||||
|
||||
const cacheKey = CacheKeys.mapsRoute(
|
||||
MapsService.coordKey(from.lat, from.lng),
|
||||
MapsService.coordKey(to.lat, to.lng),
|
||||
);
|
||||
const hit = await this.cache.get<RouteResult>(cacheKey);
|
||||
if (hit) return hit;
|
||||
|
||||
try {
|
||||
const url =
|
||||
`${this.base}/api/maps/route?fromLat=${from.lat}&fromLng=${from.lng}` +
|
||||
`&toLat=${to.lat}&toLng=${to.lng}&steps=false&locale=ar`;
|
||||
const res = await fetch(url, { headers: this.headers(key) });
|
||||
if (res.ok) {
|
||||
const json: any = await res.json();
|
||||
const parsed = MapsService.parseRoute(json);
|
||||
if (parsed) {
|
||||
const result: RouteResult = { ...parsed, provider: 'antlaq' };
|
||||
await this.cache.set(cacheKey, result, TTL.maps);
|
||||
return result;
|
||||
}
|
||||
} catch (e: any) {
|
||||
this.logger.warn(`antlaq route error: ${e?.message} — fallback`);
|
||||
} else {
|
||||
this.logger.warn(`antlaq route ${res.status} — fallback`);
|
||||
}
|
||||
} catch (e: any) {
|
||||
this.logger.warn(`antlaq route error: ${e?.message} — fallback`);
|
||||
}
|
||||
return this.straightLine(from, to);
|
||||
}
|
||||
|
||||
/** بحث العناوين — يُكاش فقط عند غياب تحيّز الموقع (وإلا صار المفتاح فريداً بلا فائدة). */
|
||||
async geocodeSearch(q: string, country: string, opts?: { radius?: number; lat?: number; lng?: number }) {
|
||||
const key = this.keyFor(country);
|
||||
const params = new URLSearchParams({ q, country });
|
||||
if (opts?.radius) params.set('radius', String(opts.radius));
|
||||
if (opts?.lat != null) params.set('lat', String(opts.lat));
|
||||
if (opts?.lng != null) params.set('lng', String(opts.lng));
|
||||
|
||||
const cacheable = opts?.lat == null && opts?.lng == null;
|
||||
const cacheKey = CacheKeys.mapsGeocode(country, q.trim().toLowerCase());
|
||||
if (cacheable) {
|
||||
const hit = await this.cache.get<any>(cacheKey);
|
||||
if (hit) return hit;
|
||||
}
|
||||
|
||||
const res = await fetch(`${this.base}/api/geocoding/search?${params}`, {
|
||||
headers: this.headers(key),
|
||||
});
|
||||
return res.json();
|
||||
const json = await res.json();
|
||||
// لا نُخزّن استجابة فاشلة — وإلا ثبّتنا الخطأ ليوم كامل.
|
||||
if (cacheable && res.ok) await this.cache.set(cacheKey, json, TTL.maps);
|
||||
return json;
|
||||
}
|
||||
|
||||
async reverse(lat: number, lng: number, country?: string) {
|
||||
const key = this.keyFor(country);
|
||||
const cacheKey = CacheKeys.mapsReverse(lat.toFixed(4), lng.toFixed(4));
|
||||
const hit = await this.cache.get<any>(cacheKey);
|
||||
if (hit) return hit;
|
||||
|
||||
const res = await fetch(
|
||||
`${this.base}/api/geocoding/reverse?lat=${lat}&lng=${lng}`,
|
||||
{ headers: this.headers(key) },
|
||||
);
|
||||
return res.json();
|
||||
const json = await res.json();
|
||||
if (res.ok) await this.cache.set(cacheKey, json, TTL.maps);
|
||||
return json;
|
||||
}
|
||||
|
||||
/** إضافة مكان — كتابة، لا تُكاش. تُبطل كاش البحث لأن النتائج تغيّرت. */
|
||||
async addPlace(body: any) {
|
||||
const key = this.config.get<string>('maps.placesApiKey') || this.keyFor(body?.country);
|
||||
const res = await fetch(`${this.base}/api/geocoding/places`, {
|
||||
|
||||
@@ -5,12 +5,19 @@ import { In, Repository } from 'typeorm';
|
||||
import { DeviceToken } from './entities/device-token.entity';
|
||||
import { I18nService } from '../../common/i18n/i18n.service';
|
||||
import { UsersService } from '../users/users.service';
|
||||
import { CacheService, CacheKeys, TTL } from '../../common/cache/cache.service';
|
||||
|
||||
export interface PushOptions {
|
||||
/** رسالة بيانات فقط — لا يعرضها النظام، يتولّاها التطبيق (overlay أندرويد، docs/17 A7). */
|
||||
dataOnly?: boolean;
|
||||
}
|
||||
|
||||
/** ما يُخزَّن في Redis لكل مستخدم — التوكن والمنصّة فقط، لا صف القاعدة كاملاً. */
|
||||
interface CachedToken {
|
||||
token: string;
|
||||
platform: string;
|
||||
}
|
||||
|
||||
/**
|
||||
* إشعارات FCM. تسجيل التوكنات + إرسال أفضل جهد (best-effort).
|
||||
* النصوص تُترجَم حسب لغة المستخدم (docs/17 — A2).
|
||||
@@ -25,19 +32,30 @@ export class NotificationsService {
|
||||
private readonly config: ConfigService,
|
||||
private readonly i18n: I18nService,
|
||||
private readonly users: UsersService,
|
||||
private readonly cache: CacheService,
|
||||
) {}
|
||||
|
||||
async register(tenantId: string, userId: string, token: string, platform = 'android') {
|
||||
const existing = await this.tokens.findOne({ where: { token } });
|
||||
let saved: DeviceToken;
|
||||
if (existing) {
|
||||
// قد يكون التوكن كان لمستخدم آخر على نفس الجهاز — نُبطل كاشه أيضاً.
|
||||
const previousUser = existing.user_id;
|
||||
const previousTenant = existing.tenant_id;
|
||||
existing.tenant_id = tenantId;
|
||||
existing.user_id = userId;
|
||||
existing.platform = platform;
|
||||
return this.tokens.save(existing);
|
||||
saved = await this.tokens.save(existing);
|
||||
if (previousUser !== userId || previousTenant !== tenantId) {
|
||||
await this.cache.del(CacheKeys.deviceTokens(previousTenant, previousUser));
|
||||
}
|
||||
} else {
|
||||
saved = await this.tokens.save(
|
||||
this.tokens.create({ tenant_id: tenantId, user_id: userId, token, platform }),
|
||||
);
|
||||
}
|
||||
return this.tokens.save(
|
||||
this.tokens.create({ tenant_id: tenantId, user_id: userId, token, platform }),
|
||||
);
|
||||
await this.cache.del(CacheKeys.deviceTokens(tenantId, userId));
|
||||
return saved;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -84,14 +102,36 @@ export class NotificationsService {
|
||||
data: Record<string, any> = {},
|
||||
opts: PushOptions = {},
|
||||
) {
|
||||
const rows = await this.tokens.find({ where: { tenant_id: tenantId, user_id: userId } });
|
||||
return this.dispatch(rows, title, body, data, opts);
|
||||
const rows = await this.tokensFor(tenantId, userId);
|
||||
return this.dispatch(tenantId, userId, rows, title, body, data, opts);
|
||||
}
|
||||
|
||||
// ---- داخلي ----
|
||||
|
||||
/**
|
||||
* توكنات المستخدم — Redis خط أول (docs/17 — G2). كانت استعلام قاعدة قبل كل
|
||||
* إشعار. القائمة الفارغة تُخزَّن أيضاً (مستخدم بلا جهاز لا يستحق استعلاماً
|
||||
* متكرراً)؛ `register` و«التوكن الميت» يُبطلان المفتاح.
|
||||
*/
|
||||
private async tokensFor(tenantId: string, userId: string): Promise<CachedToken[]> {
|
||||
const cached = await this.cache.wrap<CachedToken[]>(
|
||||
CacheKeys.deviceTokens(tenantId, userId),
|
||||
TTL.deviceTokens,
|
||||
async () => {
|
||||
const rows = await this.tokens.find({
|
||||
where: { tenant_id: tenantId, user_id: userId },
|
||||
select: { token: true, platform: true },
|
||||
});
|
||||
return rows.map((r) => ({ token: r.token, platform: r.platform }));
|
||||
},
|
||||
);
|
||||
return cached ?? [];
|
||||
}
|
||||
|
||||
private async dispatch(
|
||||
rows: DeviceToken[],
|
||||
tenantId: string,
|
||||
userId: string,
|
||||
rows: CachedToken[],
|
||||
title: string,
|
||||
body: string,
|
||||
data: Record<string, any>,
|
||||
@@ -133,8 +173,10 @@ export class NotificationsService {
|
||||
);
|
||||
|
||||
// توكن ميت = إرسال ضائع لكل رحلة لاحقة؛ نحذفه فور ما يخبرنا FCM.
|
||||
// والكاش يُبطَل معه، وإلا بقي التوكن الميت يُستعمل حتى انتهاء عمر المفتاح.
|
||||
if (stale.length > 0) {
|
||||
await this.tokens.delete({ token: In(stale) });
|
||||
await this.cache.del(CacheKeys.deviceTokens(tenantId, userId));
|
||||
this.logger.debug(`removed ${stale.length} stale token(s)`);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -10,6 +10,7 @@ import {
|
||||
@Entity('ratings')
|
||||
@Index(['tenant_id', 'trip_id'])
|
||||
@Index(['tenant_id', 'target_driver_id'])
|
||||
@Index(['tenant_id', 'target_user_id'])
|
||||
export class Rating {
|
||||
@PrimaryGeneratedColumn('uuid')
|
||||
id: string;
|
||||
@@ -26,9 +27,15 @@ export class Rating {
|
||||
@Column()
|
||||
by_role: string; // rider | driver
|
||||
|
||||
// السائق المُقيَّم (عندما يقيّم الراكب).
|
||||
@Column({ type: 'uuid', nullable: true })
|
||||
target_driver_id: string | null;
|
||||
|
||||
// الراكب المُقيَّم (عندما يقيّم السائق) — بدونه كان تقييم السائق للراكب
|
||||
// يُخزَّن بلا هدف فلا يُجمَّع أبداً (docs/17 — G6).
|
||||
@Column({ type: 'uuid', nullable: true })
|
||||
target_user_id: string | null;
|
||||
|
||||
@Column({ type: 'int' })
|
||||
stars: number;
|
||||
|
||||
|
||||
@@ -0,0 +1,84 @@
|
||||
import RedisMock from 'ioredis-mock';
|
||||
import { RatingAggregateService } from './rating-aggregate.service';
|
||||
|
||||
const TENANT = 't1';
|
||||
const DRIVER = 'd1';
|
||||
|
||||
/** ريبو مزيّف يحاكي SUM/COUNT من جدول ratings. */
|
||||
function fakeRepo(sum = 0, count = 0) {
|
||||
const chain: any = {
|
||||
select: () => chain,
|
||||
addSelect: () => chain,
|
||||
where: () => chain,
|
||||
getRawOne: async () => ({ sum: String(sum), count: String(count) }),
|
||||
};
|
||||
return { createQueryBuilder: () => chain } as any;
|
||||
}
|
||||
|
||||
describe('RatingAggregateService', () => {
|
||||
let redis: any;
|
||||
|
||||
beforeEach(() => {
|
||||
redis = new RedisMock({ keyPrefix: 'tripz:' });
|
||||
});
|
||||
|
||||
it('يبني التجميعة من القاعدة عند أول تقييم ثم يراكم في Redis', async () => {
|
||||
// القاعدة فيها تقييمان بمجموع 8 (متوسط 4)
|
||||
const agg = new RatingAggregateService(redis, fakeRepo(8, 2));
|
||||
|
||||
const r = await agg.add(TENANT, 'driver', DRIVER, 5);
|
||||
expect(r.count).toBe(3);
|
||||
expect(r.avg).toBeCloseTo(13 / 3);
|
||||
});
|
||||
|
||||
it('التقييم الثاني لا يعيد البذر من القاعدة (لا يضيع)', async () => {
|
||||
const repo = fakeRepo(8, 2);
|
||||
const spy = jest.spyOn(repo, 'createQueryBuilder');
|
||||
const agg = new RatingAggregateService(redis, repo);
|
||||
|
||||
await agg.add(TENANT, 'driver', DRIVER, 5);
|
||||
await agg.add(TENANT, 'driver', DRIVER, 3);
|
||||
|
||||
const r = await agg.get(TENANT, 'driver', DRIVER);
|
||||
expect(r!.count).toBe(4);
|
||||
expect(r!.avg).toBeCloseTo(16 / 4);
|
||||
expect(spy).toHaveBeenCalledTimes(1); // بُذرت مرة واحدة فقط
|
||||
});
|
||||
|
||||
it('الكتابة اليومية تُحجز لأول تقييم فقط', async () => {
|
||||
const agg = new RatingAggregateService(redis, fakeRepo());
|
||||
|
||||
const first = await agg.add(TENANT, 'driver', DRIVER, 5);
|
||||
const second = await agg.add(TENANT, 'driver', DRIVER, 4);
|
||||
const third = await agg.add(TENANT, 'driver', DRIVER, 3);
|
||||
|
||||
expect(first.shouldFlush).toBe(true);
|
||||
expect(second.shouldFlush).toBe(false);
|
||||
expect(third.shouldFlush).toBe(false);
|
||||
});
|
||||
|
||||
it('سائق وراكب لهما تجميعتان منفصلتان', async () => {
|
||||
const agg = new RatingAggregateService(redis, fakeRepo());
|
||||
|
||||
await agg.add(TENANT, 'driver', 'x1', 5);
|
||||
await agg.add(TENANT, 'rider', 'x1', 1);
|
||||
|
||||
expect((await agg.get(TENANT, 'driver', 'x1'))!.avg).toBe(5);
|
||||
expect((await agg.get(TENANT, 'rider', 'x1'))!.avg).toBe(1);
|
||||
});
|
||||
|
||||
it('المستأجرون معزولون', async () => {
|
||||
const agg = new RatingAggregateService(redis, fakeRepo());
|
||||
|
||||
await agg.add('tenant-a', 'driver', DRIVER, 5);
|
||||
await agg.add('tenant-b', 'driver', DRIVER, 1);
|
||||
|
||||
expect((await agg.get('tenant-a', 'driver', DRIVER))!.avg).toBe(5);
|
||||
expect((await agg.get('tenant-b', 'driver', DRIVER))!.avg).toBe(1);
|
||||
});
|
||||
|
||||
it('get يرجع null لهدف بلا أي تقييم', async () => {
|
||||
const agg = new RatingAggregateService(redis, fakeRepo(0, 0));
|
||||
expect(await agg.get(TENANT, 'driver', 'never-rated')).toBeNull();
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,147 @@
|
||||
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 { Rating } from './entities/rating.entity';
|
||||
|
||||
export type RatingTarget = 'driver' | 'rider';
|
||||
|
||||
export interface RatingAggregate {
|
||||
avg: number;
|
||||
count: number;
|
||||
}
|
||||
|
||||
/** التجميعة تُجدَّد مع كل تقييم؛ الضياع غير مؤذٍ (تُبنى من جدول ratings). */
|
||||
const AGG_TTL_SEC = 30 * 24 * 3600;
|
||||
|
||||
/**
|
||||
* تجميع التقييمات في Redis (docs/17 — G6).
|
||||
*
|
||||
* كان كل تقييم يشغّل `AVG()` على كامل تقييمات السائق ثم يكتب على القاعدة.
|
||||
* الآن: مجموع/عدد متراكمان في Redis (استجابة فورية)، و**كتابة واحدة على القاعدة
|
||||
* في اليوم** لكل طرف. التأخير مقصود: التقييم الظاهر لا يتحرّك فوراً بعد الرحلة،
|
||||
* فيبعُد الاحتكاك المباشر بين الراكب والسائق (طلب المالك).
|
||||
*
|
||||
* مصدر الحقيقة يبقى جدول `ratings` — التجميعة تُبنى منه عند أول استخدام.
|
||||
*/
|
||||
@Injectable()
|
||||
export class RatingAggregateService {
|
||||
private readonly logger = new Logger('RatingAgg');
|
||||
|
||||
constructor(
|
||||
@Inject(REDIS) private readonly redis: Redis,
|
||||
@InjectRepository(Rating) private readonly repo: Repository<Rating>,
|
||||
) {}
|
||||
|
||||
private key(tenantId: string, target: RatingTarget, targetId: string): string {
|
||||
return `rating:${tenantId}:${target}:${targetId}`;
|
||||
}
|
||||
|
||||
private today(): string {
|
||||
return new Date().toISOString().slice(0, 10);
|
||||
}
|
||||
|
||||
/**
|
||||
* يضيف تقييماً ويُرجع المتوسط الجديد، ويقول هل حان دور الكتابة اليومية.
|
||||
* `shouldFlush` يصير true لأول مُنادٍ في اليوم فقط — البقية يقرؤون من Redis.
|
||||
*/
|
||||
async add(
|
||||
tenantId: string,
|
||||
target: RatingTarget,
|
||||
targetId: string,
|
||||
stars: number,
|
||||
): Promise<RatingAggregate & { shouldFlush: boolean }> {
|
||||
const k = this.key(tenantId, target, targetId);
|
||||
await this.seed(tenantId, target, targetId, k);
|
||||
|
||||
const [sumRes, countRes] = await this.redis
|
||||
.multi()
|
||||
.hincrbyfloat(k, 'sum', stars)
|
||||
.hincrby(k, 'count', 1)
|
||||
.expire(k, AGG_TTL_SEC)
|
||||
.exec();
|
||||
|
||||
const sum = Number(sumRes?.[1] ?? 0);
|
||||
const count = Number(countRes?.[1] ?? 0);
|
||||
const avg = count > 0 ? sum / count : 0;
|
||||
|
||||
return { avg, count, shouldFlush: await this.claimDailyFlush(k) };
|
||||
}
|
||||
|
||||
/** المتوسط الحالي من Redis (بلا لمس القاعدة). */
|
||||
async get(
|
||||
tenantId: string,
|
||||
target: RatingTarget,
|
||||
targetId: string,
|
||||
): Promise<RatingAggregate | null> {
|
||||
const k = this.key(tenantId, target, targetId);
|
||||
await this.seed(tenantId, target, targetId, k);
|
||||
const h = await this.redis.hgetall(k);
|
||||
if (!h?.count) return null;
|
||||
const count = Number(h.count);
|
||||
return { avg: count > 0 ? Number(h.sum) / count : 0, count };
|
||||
}
|
||||
|
||||
// ---- داخلي ----
|
||||
|
||||
/**
|
||||
* يبني التجميعة من القاعدة إن غابت. السكربت يكتب **فقط** إن كان المفتاح
|
||||
* مفقوداً — فمُنادِيان متزامنان لا يمسح أحدهما زيادة الآخر.
|
||||
*/
|
||||
private async seed(
|
||||
tenantId: string,
|
||||
target: RatingTarget,
|
||||
targetId: string,
|
||||
k: string,
|
||||
): Promise<void> {
|
||||
if ((await this.redis.exists(k)) === 1) return;
|
||||
|
||||
const column = target === 'driver' ? 'target_driver_id' : 'target_user_id';
|
||||
const raw = await this.repo
|
||||
.createQueryBuilder('r')
|
||||
.select('COALESCE(SUM(r.stars), 0)', 'sum')
|
||||
.addSelect('COUNT(*)', 'count')
|
||||
.where(`r.tenant_id = :t AND r.${column} = :id`, { t: tenantId, id: targetId })
|
||||
.getRawOne<{ sum: string; count: string }>();
|
||||
|
||||
const script = `
|
||||
if redis.call('EXISTS', KEYS[1]) == 1 then return 0 end
|
||||
redis.call('HSET', KEYS[1], 'sum', ARGV[1], 'count', ARGV[2])
|
||||
redis.call('EXPIRE', KEYS[1], ARGV[3])
|
||||
return 1
|
||||
`;
|
||||
try {
|
||||
await this.redis.eval(
|
||||
script,
|
||||
1,
|
||||
k,
|
||||
String(Number(raw?.sum ?? 0)),
|
||||
String(Number(raw?.count ?? 0)),
|
||||
String(AGG_TTL_SEC),
|
||||
);
|
||||
} catch (e: any) {
|
||||
this.logger.warn(`seed failed for ${k}: ${e?.message}`);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* يحجز الكتابة اليومية: يرجع true لأول مُنادٍ في اليوم فقط.
|
||||
* ذرّي — وإلا كتب كل تقييم على القاعدة وضاع الغرض كله.
|
||||
*/
|
||||
private async claimDailyFlush(k: string): Promise<boolean> {
|
||||
const script = `
|
||||
if redis.call('HGET', KEYS[1], 'flushed_on') == ARGV[1] then return 0 end
|
||||
redis.call('HSET', KEYS[1], 'flushed_on', ARGV[1])
|
||||
return 1
|
||||
`;
|
||||
try {
|
||||
const res = await this.redis.eval(script, 1, k, this.today());
|
||||
return res === 1;
|
||||
} catch (e: any) {
|
||||
// فشل Redis: نكتب على القاعدة (الأسلم أن نكتب مرة زائدة من ألّا نكتب أبداً).
|
||||
this.logger.warn(`flush claim failed for ${k}: ${e?.message}`);
|
||||
return true;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -2,13 +2,15 @@ import { Module } from '@nestjs/common';
|
||||
import { TypeOrmModule } from '@nestjs/typeorm';
|
||||
import { Rating } from './entities/rating.entity';
|
||||
import { RatingsService } from './ratings.service';
|
||||
import { RatingAggregateService } from './rating-aggregate.service';
|
||||
import { RatingsController } from './ratings.controller';
|
||||
import { TripsModule } from '../trips/trips.module';
|
||||
import { DriversModule } from '../drivers/drivers.module';
|
||||
import { UsersModule } from '../users/users.module';
|
||||
|
||||
@Module({
|
||||
imports: [TypeOrmModule.forFeature([Rating]), TripsModule, DriversModule],
|
||||
imports: [TypeOrmModule.forFeature([Rating]), TripsModule, DriversModule, UsersModule],
|
||||
controllers: [RatingsController],
|
||||
providers: [RatingsService],
|
||||
providers: [RatingsService, RatingAggregateService],
|
||||
})
|
||||
export class RatingsModule {}
|
||||
|
||||
@@ -9,6 +9,8 @@ import { Repository } from 'typeorm';
|
||||
import { Rating } from './entities/rating.entity';
|
||||
import { TripsService } from '../trips/trips.service';
|
||||
import { DriversService } from '../drivers/drivers.service';
|
||||
import { UsersService } from '../users/users.service';
|
||||
import { RatingAggregateService } from './rating-aggregate.service';
|
||||
|
||||
@Injectable()
|
||||
export class RatingsService {
|
||||
@@ -16,6 +18,8 @@ export class RatingsService {
|
||||
@InjectRepository(Rating) private readonly repo: Repository<Rating>,
|
||||
private readonly trips: TripsService,
|
||||
private readonly drivers: DriversService,
|
||||
private readonly users: UsersService,
|
||||
private readonly aggregate: RatingAggregateService,
|
||||
) {}
|
||||
|
||||
async rate(
|
||||
@@ -49,7 +53,10 @@ export class RatingsService {
|
||||
});
|
||||
if (dup) throw new BadRequestException('Already rated');
|
||||
|
||||
// الراكب يقيّم السائق، والسائق يقيّم الراكب — الطرفان معاً (docs/17 — G6).
|
||||
const targetDriverId = role === 'rider' ? trip.driver_id : null;
|
||||
const targetUserId = role === 'driver' ? trip.rider_id : null;
|
||||
|
||||
const rating = await this.repo.save(
|
||||
this.repo.create({
|
||||
tenant_id: tenantId,
|
||||
@@ -57,23 +64,22 @@ export class RatingsService {
|
||||
by_user_id: byUserId,
|
||||
by_role: role,
|
||||
target_driver_id: targetDriverId,
|
||||
target_user_id: targetUserId,
|
||||
stars,
|
||||
comment: comment ?? null,
|
||||
}),
|
||||
);
|
||||
|
||||
// تحديث متوسط تقييم السائق عند تقييم الراكب له
|
||||
if (role === 'rider' && targetDriverId) {
|
||||
const raw = await this.repo
|
||||
.createQueryBuilder('r')
|
||||
.select('AVG(r.stars)', 'avg')
|
||||
.where('r.tenant_id = :t AND r.target_driver_id = :d', {
|
||||
t: tenantId,
|
||||
d: targetDriverId,
|
||||
})
|
||||
.getRawOne<{ avg: string }>();
|
||||
if (raw?.avg) {
|
||||
await this.drivers.setRating(tenantId, targetDriverId, Number(raw.avg));
|
||||
// التجميع في Redis؛ والقاعدة تُحدَّث مرة واحدة في اليوم لكل طرف.
|
||||
if (targetDriverId) {
|
||||
const agg = await this.aggregate.add(tenantId, 'driver', targetDriverId, stars);
|
||||
if (agg.shouldFlush) {
|
||||
await this.drivers.setRating(tenantId, targetDriverId, agg.avg);
|
||||
}
|
||||
} else if (targetUserId) {
|
||||
const agg = await this.aggregate.add(tenantId, 'rider', targetUserId, stars);
|
||||
if (agg.shouldFlush) {
|
||||
await this.users.setRating(tenantId, targetUserId, agg.avg);
|
||||
}
|
||||
}
|
||||
return rating;
|
||||
|
||||
@@ -2,26 +2,37 @@ import { Injectable } from '@nestjs/common';
|
||||
import { InjectRepository } from '@nestjs/typeorm';
|
||||
import { Repository } from 'typeorm';
|
||||
import { RideType } from './entities/ride-type.entity';
|
||||
import { CacheService, CacheKeys, TTL } from '../../common/cache/cache.service';
|
||||
|
||||
@Injectable()
|
||||
export class RideTypesService {
|
||||
constructor(
|
||||
@InjectRepository(RideType) private readonly repo: Repository<RideType>,
|
||||
private readonly cache: CacheService,
|
||||
) {}
|
||||
|
||||
listActive(tenantId: string): Promise<RideType[]> {
|
||||
return this.repo.find({
|
||||
where: { tenant_id: tenantId, active: true },
|
||||
order: { sort: 'ASC' },
|
||||
});
|
||||
/** أنواع الرحلات — Redis خط أول (docs/17 — G5)؛ يناديها كل تطبيق عند الإقلاع. */
|
||||
async listActive(tenantId: string): Promise<RideType[]> {
|
||||
const cached = await this.cache.wrap<RideType[]>(
|
||||
CacheKeys.rideTypes(tenantId),
|
||||
TTL.catalog,
|
||||
() =>
|
||||
this.repo.find({
|
||||
where: { tenant_id: tenantId, active: true },
|
||||
order: { sort: 'ASC' },
|
||||
}),
|
||||
);
|
||||
return cached ?? [];
|
||||
}
|
||||
|
||||
findByCode(tenantId: string, code: string): Promise<RideType | null> {
|
||||
return this.repo.findOne({ where: { tenant_id: tenantId, code } });
|
||||
}
|
||||
|
||||
create(data: Partial<RideType>): Promise<RideType> {
|
||||
return this.repo.save(this.repo.create(data));
|
||||
async create(data: Partial<RideType>): Promise<RideType> {
|
||||
const saved = await this.repo.save(this.repo.create(data));
|
||||
await this.cache.del(CacheKeys.rideTypes(saved.tenant_id));
|
||||
return saved;
|
||||
}
|
||||
|
||||
/** يُنشئ النوع إن لم يوجد (للـ seed). */
|
||||
|
||||
@@ -3,32 +3,45 @@ import { InjectRepository } from '@nestjs/typeorm';
|
||||
import { Repository } from 'typeorm';
|
||||
import { Tariff } from './entities/tariff.entity';
|
||||
import { TariffEngine, QuoteInput, QuoteBreakdown } from './tariff.engine';
|
||||
import { CacheService, CacheKeys, TTL } from '../../common/cache/cache.service';
|
||||
|
||||
@Injectable()
|
||||
export class TariffService {
|
||||
constructor(
|
||||
@InjectRepository(Tariff)
|
||||
private readonly repo: Repository<Tariff>,
|
||||
private readonly cache: CacheService,
|
||||
) {}
|
||||
|
||||
/** التعرفة الفعّالة — Redis خط أول (docs/17 — G5)؛ تُقرأ في كل طلب رحلة. */
|
||||
getActive(
|
||||
tenantId: string,
|
||||
city: string,
|
||||
serviceClass: string,
|
||||
): Promise<Tariff | null> {
|
||||
return this.repo.findOne({
|
||||
where: {
|
||||
tenant_id: tenantId,
|
||||
city,
|
||||
service_class: serviceClass,
|
||||
active: true,
|
||||
},
|
||||
order: { version: 'DESC' },
|
||||
});
|
||||
return this.cache.wrap(
|
||||
CacheKeys.tariff(tenantId, city, serviceClass),
|
||||
TTL.catalog,
|
||||
() =>
|
||||
this.repo.findOne({
|
||||
where: {
|
||||
tenant_id: tenantId,
|
||||
city,
|
||||
service_class: serviceClass,
|
||||
active: true,
|
||||
},
|
||||
order: { version: 'DESC' },
|
||||
}),
|
||||
);
|
||||
}
|
||||
|
||||
create(data: Partial<Tariff>): Promise<Tariff> {
|
||||
return this.repo.save(this.repo.create(data));
|
||||
async create(data: Partial<Tariff>): Promise<Tariff> {
|
||||
const saved = await this.repo.save(this.repo.create(data));
|
||||
// تعرفة جديدة = النسخة الفعّالة تغيّرت → الكاش صار كذباً.
|
||||
await this.cache.del(
|
||||
CacheKeys.tariff(saved.tenant_id, saved.city, saved.service_class),
|
||||
);
|
||||
return saved;
|
||||
}
|
||||
|
||||
async quote(
|
||||
|
||||
@@ -2,12 +2,14 @@ import { Injectable } from '@nestjs/common';
|
||||
import { InjectRepository } from '@nestjs/typeorm';
|
||||
import { Repository } from 'typeorm';
|
||||
import { Tenant } from '../../database/entities/tenant.entity';
|
||||
import { CacheService, CacheKeys, TTL } from '../../common/cache/cache.service';
|
||||
|
||||
@Injectable()
|
||||
export class TenantsService {
|
||||
constructor(
|
||||
@InjectRepository(Tenant)
|
||||
private readonly repo: Repository<Tenant>,
|
||||
private readonly cache: CacheService,
|
||||
) {}
|
||||
|
||||
findAll(): Promise<Tenant[]> {
|
||||
@@ -27,16 +29,26 @@ export class TenantsService {
|
||||
*/
|
||||
async resolve(idOrSlug: string): Promise<Tenant | null> {
|
||||
if (!idOrSlug) return null;
|
||||
const bySlug = await this.findBySlug(idOrSlug);
|
||||
if (bySlug) return bySlug;
|
||||
if (TenantsService.UUID_RE.test(idOrSlug)) {
|
||||
return this.repo.findOne({ where: { id: idOrSlug } });
|
||||
}
|
||||
return null;
|
||||
// Redis خط أول (docs/17 — G4): المستأجر يتغيّر نادراً جداً.
|
||||
return this.cache.wrap(CacheKeys.tenant(idOrSlug), TTL.tenant, async () => {
|
||||
const bySlug = await this.findBySlug(idOrSlug);
|
||||
if (bySlug) return bySlug;
|
||||
if (TenantsService.UUID_RE.test(idOrSlug)) {
|
||||
return this.repo.findOne({ where: { id: idOrSlug } });
|
||||
}
|
||||
return null;
|
||||
});
|
||||
}
|
||||
|
||||
create(data: Partial<Tenant>): Promise<Tenant> {
|
||||
return this.repo.save(this.repo.create(data));
|
||||
async create(data: Partial<Tenant>): Promise<Tenant> {
|
||||
const saved = await this.repo.save(this.repo.create(data));
|
||||
await this.invalidate(saved);
|
||||
return saved;
|
||||
}
|
||||
|
||||
/** يُبطل مفتاحَي المستأجر (بالـslug وبالـUUID) — يُنادى بعد أي تعديل عليه. */
|
||||
async invalidate(tenant: Pick<Tenant, 'id' | 'slug'>): Promise<void> {
|
||||
await this.cache.del(CacheKeys.tenant(tenant.slug), CacheKeys.tenant(tenant.id));
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -44,7 +56,8 @@ export class TenantsService {
|
||||
* كل ما يمكن جعله ديناميكياً (نصوص، ميزات، ألوان، دفع) يأتي من هنا — راجع docs/06.
|
||||
*/
|
||||
async config(slug: string) {
|
||||
const t = await this.findBySlug(slug);
|
||||
// عبر resolve ليستفيد من الكاش — كل تطبيق يناديها عند الإقلاع.
|
||||
const t = await this.resolve(slug);
|
||||
if (!t) return null;
|
||||
return {
|
||||
slug: t.slug,
|
||||
|
||||
@@ -34,6 +34,10 @@ export class User {
|
||||
@Column({ default: 'ar' })
|
||||
language: string;
|
||||
|
||||
// تقييم الراكب من السائقين — يُحدَّث مرة واحدة يومياً (docs/17 — G6).
|
||||
@Column({ type: 'numeric', precision: 3, scale: 2, default: 5 })
|
||||
rating: number;
|
||||
|
||||
@CreateDateColumn()
|
||||
created_at: Date;
|
||||
|
||||
|
||||
@@ -2,12 +2,14 @@ import { Injectable } from '@nestjs/common';
|
||||
import { InjectRepository } from '@nestjs/typeorm';
|
||||
import { Repository } from 'typeorm';
|
||||
import { User, UserRole } from './entities/user.entity';
|
||||
import { CacheService, CacheKeys, TTL } from '../../common/cache/cache.service';
|
||||
|
||||
@Injectable()
|
||||
export class UsersService {
|
||||
constructor(
|
||||
@InjectRepository(User)
|
||||
private readonly userRepository: Repository<User>,
|
||||
private readonly cache: CacheService,
|
||||
) {}
|
||||
|
||||
async findByPhone(tenantId: string, phone: string): Promise<User | null> {
|
||||
@@ -43,16 +45,34 @@ export class UsersService {
|
||||
if (Object.keys(patch).length > 0) {
|
||||
await this.userRepository.update({ tenant_id: tenantId, id }, patch);
|
||||
}
|
||||
// اللغة تغيّرت → الكاش صار كذباً. نكتب القيمة الجديدة مباشرة بدل الحذف
|
||||
// (فلاتر يرفع اللغة عند كل إقلاع، فالكتابة هنا تعني إصابة كاش دائمة).
|
||||
if (data.language !== undefined) {
|
||||
await this.cache.set(CacheKeys.userLang(tenantId, id), data.language, TTL.userLang);
|
||||
}
|
||||
return this.findById(tenantId, id);
|
||||
}
|
||||
|
||||
/** لغة الإشعارات للمستخدم — استعلام خفيف (عمود واحد) يُستدعى قبل كل push. */
|
||||
/**
|
||||
* لغة الإشعارات — Redis خط أول (docs/17 — G1). كانت استعلام قاعدة قبل **كل**
|
||||
* إشعار؛ صارت إصابة كاش، والقاعدة احتياط عند أول مرة أو انتهاء العمر.
|
||||
*/
|
||||
async getLanguage(tenantId: string, id: string): Promise<string | null> {
|
||||
const row = await this.userRepository.findOne({
|
||||
where: { tenant_id: tenantId, id },
|
||||
select: { language: true },
|
||||
return this.cache.wrap(CacheKeys.userLang(tenantId, id), TTL.userLang, async () => {
|
||||
const row = await this.userRepository.findOne({
|
||||
where: { tenant_id: tenantId, id },
|
||||
select: { language: true },
|
||||
});
|
||||
return row?.language ?? null;
|
||||
});
|
||||
return row?.language ?? null;
|
||||
}
|
||||
|
||||
/** تقييم الراكب — يُكتب مرة واحدة يومياً من التجميعة (docs/17 — G6). */
|
||||
async setRating(tenantId: string, id: string, rating: number): Promise<void> {
|
||||
await this.userRepository.update(
|
||||
{ tenant_id: tenantId, id },
|
||||
{ rating: Number(rating.toFixed(2)) },
|
||||
);
|
||||
}
|
||||
|
||||
async setRole(tenantId: string, id: string, role: string): Promise<void> {
|
||||
|
||||
Reference in New Issue
Block a user