19 Commits

Author SHA1 Message Date
71568e2274 feat(rpartstore): Renault/Dacia için RPartStore VIN-decode fallback kaynağı
Some checks are pending
QA Gate (P0/P1) / Test affected app (pull_request) Waiting to run
rpartstore.renault.com (Renault Group bayi portalı) yeni decode kaynağı olarak
eklendi. Portal REST değil STOMP 1.2 over WebSocket konuşuyor; login Okta OIE
(IDX JSON API + PKCE), tarayıcı gerekmiyor. Sıralama Renault/Dacia için:
pcat + PL24 + emex yarışı → rpartstore → (bayrağı açıksa) Vinpin son çare.

- integrations/rpartstore: STOMP codec, Okta login, oturum (trace-id ile
  çok-mesajlı cevap eşleme), istemci (Redis'te 1 saatlik token), Renault/Dacia
  yönlendirme, RPartStore → PL24 catalog_vehicles model eşleyici (roman rakam,
  karoseri niteleyici, yanlış-pazar cezası).
- rpartstore_decodes tablosu (migration 0037) + rpartstore-decode BullMQ
  kuyruğu; worker concurrency 1 + 1 iş / 6 s limiter (portal 2 arama / 10 s).
- Günlük sert kota RPARTSTORE_DAILY_CAP (varsayılan 10, İstanbul günü):
  işlemci göndermeden önce Redis INCR ile rezervasyon yapar, aşan VIN'ler
  `capped` olur ve 24 saat sonra tekrar denenir; API kota doluysa kuyruğa
  hiç almaz. Kısa vadeli rate-limit cevabında retryAfter kadar bekleyip bir
  kez tekrar dener.
- VehiclesService.tryRpartstoreFallback: Vinpin fallback ile aynı sözleşme
  (decoding/catalogVehicle cevapları, tek başarısızlık log satırı).
- Env: RPARTSTORE_ENABLED/USER/PASS/DAILY_CAP/BROKER_URL/APP_VERSION
  (compose api+worker). Vinpin koda dokunulmadan pasif kalır (VINPIN_ENABLED).

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-09-25 13:58:09 +03:00
41147f3f05 fix(pl24): kendi bütçe frenimiz "hesap banlandı" alarmı üretmesin
Some checks failed
QA Gate (P0/P1) / Test affected app (pull_request) Has been cancelled
Prod 2026-09-21 19:10 (İstanbul) itibarıyla Telegram "🔴 PL24 giriş
yapılamıyor — hesap banlanmış olabilir" dedi. Hesap sağlamdı:
kullanıcı tarayıcıdan giriş yapabildi, aynı gün 143 başarılı auth
çağrısı ve **0 × 401** vardı. Gerçek sebep alarmın kendi metninde:
"PL24 daily HTTP budget exhausted (user lane: 1205/1200)" — yani
upstream değil, BİZİM günlük tavanımız.

Zincir: `attemptLogin` içindeki `budget.consume()` fırlatıyor → dış
`catch` bunu `{ error }` olarak düzleştiriyor → `loginWithLock`
`!result.sessionToken` dalına giriyor → `breakerFail()` + kesinti
sayacı → bir saat sonra ban alarmı. Yani kendi frenimiz hem devre
kesiciyi tetikliyor hem de operatörü yanlış yere çağırıyor.

- `PL24BudgetExceededError` artık `attemptLogin` ve
  `authorizeServiceForAccount` içinde olduğu gibi yukarı çıkıyor;
  giriş/yetkilendirme hatasına çevrilmiyor. Devre kesici tetiklenmiyor,
  kesinti sayacı başlamıyor, Telegram susuyor. Çağıran taraf PL24'ü
  atlayıp pcat/emex'e düşüyor — bugün de öyle oluyordu, ama artık
  sessizce ve doğru gerekçeyle.
- Şerit bazlı sayaç eklendi: `pl24:http:<gün>:{user,worker}`. Paylaşılan
  toplam tavanın NE ZAMAN dolduğunu söylüyordu ama KİMİN harcadığını
  söylemiyordu; tavanı veriyle boyutlandırmak için bu şart. Bugün
  toplam 1200'e vurdu ve 16:10 UTC'den sonra her kullanıcı decode'u
  reddedildi, ama ne kadarının ısıtma prefetch'i ne kadarının ekran
  başında bekleyen biri olduğu `proxy_logs`'tan çıkarılamıyordu.

5 yeni test. api 650 test geçiyor.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-09-21 20:52:52 +03:00
e03d721c37 Merge pull request 'feat(pl24): Faz 3 — backfill anahtarı + Mitsubishi parça-listesi düzeltmesi' (#270) from dev into main 2026-09-21 20:47:52 +03:00
142ec3a150 feat(pl24): Faz 3 — backfill anahtarı + Mitsubishi parça-listesi düzeltmesi
Some checks failed
QA Gate (P0/P1) / Test affected app (pull_request) Has been cancelled
Faz 3 backfill'i açmadan önce hacmi ölçtüm ve raporun işaret ettiğinden
çok daha büyük bir israf çıktı.

**Mitsubishi'nin parça listesi grup sanılıyordu.**
`/p5mitsubishi/extern/details/vinDetails` yanıtı `partno`/`qty` taşıyan
bir PARÇA listesi (canlı doğrulama: 16 kayıt), ama her kaydın kendi
linki `partInfoTable` ve ne wid ne yol sınıflandırıcıda karşılık
buluyordu. Sonuç (prod ölçümü): 2.193 parça listesi grup düğümüne
döndü, içlerindeki 19.576 tekil parça ("SCREW,LOCK CYLINDER",
"BOLT,STEERING COLUMN WASHER") kategori olarak kaydedildi. Bu 19.576
sahte düğümün TOPLAM 2 tanesinde parça var ve hepsi her prefetch
turunda yeniden çekiliyor. İkisi de %100 Mitsubishi.

- `detailsTable` artık yaprak, `isPl24PartDetailNode()` ile
  `partInfoTable` hiç kuyruklanmıyor.
- `isLeafLinkPath` artık `linkWid`'i de geçiriyor. Okuma yolu bu
  güvenilir sinyali hep kullanıyordu ama kuyruklama yolu düşürüyordu —
  sınıflandırıcı `detailsTable`'ı öğrense bile burada yine grup
  sayılacaktı.
- Migration 0036: 19.576 sahte kategori siliniyor. Okuma yolu bir
  düğümün ÖNCE çocuklarına baktığı için bu silme düzeltmenin parçası,
  ayrı temizlik değil. Kuru çalıştırma: 19.576 kategori, 9 araç, 2
  parça, Mitsubishi dışı 0.

**Backfill anahtarı.** `PL24_BACKFILL_ENABLED` eklendi, varsayılan
KAPALI. Eski `PL24_TR_DISABLED` adı "tr hesabı öldü" diyordu ama işi
"toplu yükü tek sağ kalan hesaptan uzak tut"tu. Eski değişken hâlâ
kapatabiliyor — yarım deploy musluğu sessizce açamasın. Kullanıcı
tetikli fast-lane bu anahtardan etkilenmiyor.

9 yeni test. api 647 test geçiyor.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-09-20 10:51:42 +03:00
e8ffa6a5db Merge pull request 'fix(pl24): model listesi giriş noktasını catmeta'dan çöz (Volvo browse)' (#269) from dev into main 2026-09-20 10:18:50 +03:00
bd9f4bd3f8 style(pl24): normalizeLabel'daki yanıltıcı karakter sınıfını düzelt
Some checks failed
QA Gate (P0/P1) / Test affected app (pull_request) Has been cancelled
`[̀-ͯ]` bir aralık sınıfı; biome bunu "bir taban karakter +
birleştirici karakter çiftini de eşleyebilir" diye işaretliyordu.
NFD zaten her aksanı kendi işaretine ayırdığı için doğrudan `\p{M}`
(tüm birleştirici işaretler) hem doğru hem daha geniş.

Gerçek Volvo/PSA/Subaru yakalamalarıyla çıktı birebir aynı kaldı.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-09-20 09:32:24 +03:00
d75c468341 fix(pl24): model listesi giriş noktasını catmeta'dan çöz
Prod'da browse onarması Peugeot'yu taşıdı (64 bayat satır → 39 güncel
model) ama Volvo'da "upstream returned no models" ile korumalı dala
düştü. Sebep: `BACKEND_MODEL_PATH`'te p5volvo yok ve yedek keşif yedi
bilinen yolu deniyor — Volvo'nunki `/extern/vehicles/models`, yani
ÇOĞUL "vehicles", listede yok. Volvo/Polestar browse bu yüzden sıfır
model tohumluyor ve emekli P4 listesini sunmaya devam ediyordu.

Otoritatif kaynak her P5 backend'inin kendi `/extern/catmeta`
yanıtındaki `data.catalogEntryPoint.path`. Canlı doğrulandı:
  p5volvo → /p5volvo/extern/vehicles/models  (60 model: XC90, V40…)
  p5psa   → /p5psa/extern/vehicle/catalogs

İki değişiklik:
- p5psa ve p5volvo `BACKEND_MODEL_PATH`'e sabitlendi (sıfır ek istek).
  p5psa zaten çalışıyordu ama her çağrıda beş boşa probe isteğiyle
  yeniden keşfediliyordu.
- `resolveModelPathFromCatmeta()`: haritada olmayan backend artık yedi
  kör probe yerine tek otoritatif catmeta çağrısı yapıyor. Sonuç
  backend başına 30 gün cache'leniyor (olumsuz sonuç dahil), yani
  backend ömrü boyunca tek istek; PL24 bir ucu taşırsa kendiliğinden
  düzeliyor. Hata durumunda null → eski yedeğe düşer.

8 yeni test. api 638 test geçiyor.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-09-20 09:30:53 +03:00
a4e22e1bb8 Merge pull request 'fix(pl24): bütçe kilitlenmesi, sabitlenmiş browse satırları, Volvo alanları + WMI kapsamı' (#268) from dev into main 2026-09-20 09:15:26 +03:00
c2f4d8f421 feat(pl24): WMI haritasını canlı servis taramasıyla genişlet, Renault'yu aç
Some checks failed
QA Gate (P0/P1) / Test affected app (pull_request) Has been cancelled
Prod'daki 4.799 aracın 1.406'sının WMI'si haritada yoktu. Hangilerinin
PL24'te gerçekten karşılığı olduğu tahmin edilmedi — canlı
`/pl24-wmi/ext/api/2.0/decode` servisine gerçek prod VIN'leriyle tek tek
soruldu (2026-09-20).

**Renault yeniden açıldı (450 araç).** İki gerekçe de artık geçersiz,
ikisi de prod'a karşı doğrulandı:
- Askı kalkmış: `/p5renault` directAccess, gerçek müşteri aracı
  VF14SRCL458170337 için `VEHICLE_IDENTIFIED` + "SYMBOL II/LOGAN II"
  döndü; WMI servisi de `{service: renault_parts, error: false}` diyor.
- Devre kesici artık kesin olumsuz yanıtları saymıyor, yalnız geçici
  taşıma hatalarını (`vehicles.service` `isTransient`).

**Yeni eşlemeler (toplam 730 araç):** NLH/TMA/NLJ/KMF → hyundai_parts,
KNE/KNC → kia_parts, MMC/XMC → mmc_parts, JSA → suzuki_parts,
NMB → **mercestrucks**_parts (servis binek Mercedes'e değil kamyon
kataloğuna çözüyor; bu 30 araç "Mercedes-Benz" etiketliydi ama binek
kataloğunda yok).

**Kasıtlı olarak EKLENMEYENLER** — servis HTTP 410 "no brands found"
dedi, tıpkı NM4 ve VR7 gibi: JHM/SHH/SHS/MAK/NLA (Honda), KL1
(Chevrolet). Eklemek yalnız boşuna istek üretir ve pcat/emex/vinpin
fallback'ini geciktirir. JMZ (Mazda) 410 değil ama fordp/fordt
döndürüyor (Ford-Mazda platform ortaklığı); marka tutarsız olduğu için
eklenmedi.

6 yeni test; "olmayanlar" da kilitlendi. Ayrıca her eşlemenin katalog
tanımı olduğunu doğrulayan bütünlük testi. api 623 test geçiyor.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-09-20 08:53:07 +03:00
72f8c2a3e4 fix(pl24): Volvo alanlarında iç kodu değil okunur değeri kullan
Some checks failed
QA Gate (P0/P1) / Test affected app (pull_request) Has been cancelled
Volvo vinfoBasic aynı özniteliği iki kez veriyor: "Şanzıman kodu" = "B"
ile "Şanzıman" = "6-PSHIFT 2WD / MPS6", "Satis tipi" = "42" ile "Türü"
= "S80". VAG için doğru olan kod-önce sırası (VAG'ın "Şanzıman kodu"
zaten anlamlı: "MQ200") bu yüzden Volvo'da şanzımanı tek harf, seriyi
çıplak sayı yapıyordu.

`lookupDescriptive()` eklendi: listedeki ilk *açıklayıcı* değeri döner
(2 karakterden uzun ve salt rakam değil), hiçbiri açıklayıcı değilse
ilk mevcut değere düşer — yani yalnız kod üreten backend'ler aynen
eskisi gibi davranır. Şanzıman ve seri bu yardımcıya geçti.

Ayrıca eksik İngilizce etiketler eklendi: `transmission` (İngilizce
yanıtta şanzıman yine tek harfe düşüyordu) ve `body_style`.

Ölçüm (gerçek keşif yakalamaları, YV1AS84ABD1168166 S80):
- şanzıman "B" → "6-PSHIFT 2WD / MPS6"
- seri "42" → "S80"
- kaporta İngilizce yanıtta null → "Sedan"
- PSA İngilizce yanıtta şanzıman null → "BVM5"
PSA Türkçe, Subaru ve diğer markalarda çıktı değişmedi.

6 yeni test + iki dilde gerçek Volvo fixture'ı. api 617 test geçiyor.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-09-20 08:45:16 +03:00
0928559ea4 fix(catalog): browse onarmasında yalnız bayat satırları sil
Some checks failed
QA Gate (P0/P1) / Test affected app (pull_request) Has been cancelled
Bir servis karışık olabilir: bazı satırları yeni mimariyle yeniden
listelenmiş, bazıları hâlâ eski. Silme servis adına göre yapıldığında
zaten taşınmış satırlar da gidiyordu — üstteki insert onları
`onConflictDoNothing` ile atladığı için geri gelmiyorlar, yani temelli
kayıp. Silme artık yalnız gerçekten bayat olan satırlara uygulanıyor.

Hangi satırların silindiğini doğrulayan test eklendi (drizzle
`inArray` parametrelerini okuyarak).

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-09-20 08:40:26 +03:00
eb18c9a122 fix(catalog): sabitlenmiş browse satırlarını görüntülendikçe yenile
Some checks failed
QA Gate (P0/P1) / Test affected app (pull_request) Has been cancelled
`catalog_vehicles` kalıcı bir cache: `getModels` markanın tek bir satırı
varsa erken dönüyor, dolayısıyla `fetchVehicleList` o marka için bir
daha hiç çağrılmıyor. PSA ve Volvo hâlâ P4 iken listelenmiş 188 satır
(prod 2026-09-20: 127 LEGACY_PSA + 61 LEGACY_VOLVO) bu yüzden kalıcı
olarak çakılı kalmıştı:

- Peugeot/Citroën browse donmuş 2024-02-13 anlık görüntüsünü sunuyor,
- Volvo/Polestar browse HTTP 503 dönen bir uca gidiyor.

Kod düzeltmesi bu satırlara hiçbir zaman ulaşmıyordu.

`isStaleBrowseArchitecture()` + `CatalogService.healStaleBrowseRows()`:
satırın kayıtlı mimarisi servis tablosundakiyle uyuşmuyorsa marka bir
kez yeniden listeleniyor, yeni satırlar yazılıp eskiler aynı
transaction'da siliniyor. VIN tarafındaki onarma ile aynı temkinli
kurallar: marka başına günde bir deneme (Redis kilidi), başarısız veya
boş listede eski satırlar korunuyor, silme ancak yenisi elde edilince
ve servis bazında yapılıyor.

8 yeni test; api paketi 610 test geçiyor.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-09-20 08:35:09 +03:00
ae8a8046bf refactor(prefetch): kapı sırasını en ucuzdan pahalıya diz
Some checks failed
QA Gate (P0/P1) / Test affected app (pull_request) Has been cancelled
Pencere kontrolü saf saat aritmetiği, hiç I/O yapmıyor ve pencere
dışındaki iş zaten hiçbir faydalı iş yapamıyor — dolayısıyla dakikalık
Redis sayacından da önce gelmeli. Yeni sıra: pencere → cooldown →
dakikalık tavan → günlük bütçe. Böylece pencere dışında uyanan iş
hiçbir sayacı kirletmiyor.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-09-20 08:31:55 +03:00
7644d91407 fix(prefetch): PL24 bütçe/iş-saati kilitlenmesini kır
Some checks failed
QA Gate (P0/P1) / Test affected app (pull_request) Has been cancelled
Prod'da PL24 prefetch tamamen durmuştu: günlük bütçe 600/600 dolu
görünürken gün boyunca SIFIR katalog isteği ve SIFIR yeni kategori
üretiliyordu (ölçüm 2026-09-20).

İki hata birlikte kapalı bir döngü kuruyordu:

1. Günlük bütçe `process()` içinde düşülüyor, iş-saati (ve cooldown)
   kapısı ise her handler'ın başında duruyordu. Pencere dışında uyanan
   bir iş önce bütçeden bir birim yiyor, sonra `time-window` fırlatıp
   hiçbir iş yapmadan erteleniyordu.

2. Bütçesi biten kaynak "bir sonraki UTC gece yarısı"na erteleniyordu.
   PREFETCH_PL24_START=9 ile bu an 03:00 Europe/Istanbul'a denk gelir —
   pencere açılmadan altı saat önce. Uyanan iş yine pencereye takılıyor,
   yine erteleniyor; taze günlük bütçe daha pencere açılmadan bu boş
   uyanmalarla tükeniyordu.

Düzeltme:
- cooldown + iş-saati kapıları `process()` içinde, bütçe düşülmeden
  ÖNCE çalışıyor; handler'lardaki kopyaları kaldırıldı.
- `alignToWindow()` eklendi: bütçe ertelemesi pencerenin içine
  hizalanıyor. Pencere tanımlı değilse (varsayılan 0–24) no-op.
- `currentIstanbulHour` artık `istanbulHourAt`'e deleg ediyor ve h24
  döngüsünün gece yarısı için ürettiği "24" değeri `% 24` ile
  normalleniyor (aksi halde saat hiçbir pencereye düşmez).

8 yeni regresyon testi; api paketi 602 test geçiyor.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-09-20 08:29:55 +03:00
40c4d33068 Merge pull request 'fix(pl24): VR7 WMI'sini haritadan çıkar — PL24'te böyle bir marka yok' (#267) from dev into main 2026-09-20 08:16:37 +03:00
semih
fe91b504de fix(pl24): VR7 WMI'sini haritadan çıkar — PL24'te böyle bir marka yok
Some checks failed
QA Gate (P0/P1) / Test affected app (pull_request) Has been cancelled
Faz 2 / açık soru kapatıldı (analiz: /home/s/ss/plv2.md, bulgu types_wmi-04).

Rapor "VR7 muhtemelen Citroën, peugeot_parts yanlış" diyordu ve doğrulama
istiyordu. Canlı portal WMI servisi (2026-09-19, /pl24-wmi/ext/api/2.0/decode,
pl24-wmidata kapsamlı token) kesin yanıtı verdi:

  HTTP 410 — "error resolving VIN: no brands found for WMI = VR7"

Yani VR7 ne Peugeot ne Citroën: PL24'ün WMI veritabanında HİÇ YOK, tıpkı Tofaş
NM4 gibi. Raporun önerisi (VR7 → citroen_parts) uygulansaydı hata aynen devam
edecekti.

Prod'daki 17 VR7 aracının tamamı model çözülmeden ("Peugeot Peugeot") ve 0
kategoriyle kaydedilmişti. Haritada tutmak yalnız her sorguda boşuna bir PL24
isteği üretiyor ve gerçek kaynağa (pcat/emex/vinpin) geçişi geciktiriyor —
canlı gözlem: P4→P5 onarma denemesi de bu VIN'lerde "kayıt bulunamadı" ile
düşüyor.

Test: pl24.types.spec'e kapsam-dışı WMI kilidi (VR7 ve NM4 haritada olmamalı) +
canlı doğrulanmış yönlendirmeler (JF1→subaru, ZAR→alfa, VXK→psa_opel).

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-09-20 00:35:03 +03:00
3b63ae02c6 Merge pull request 'feat(pl24): bayat P4 mimarisindeki araçları görüntülendikçe P5'e taşı' (#265) from dev into main 2026-09-20 00:26:01 +03:00
58ba612707 Merge pull request 'feat(internal-admin): expired/cancelled kullanıcı için yeni abonelik başlatma ucu' (#266) from feat/admin-start-subscription into main 2026-09-19 01:37:19 +03:00
semih
ddc223f745 feat(pl24): bayat P4 mimarisindeki araçları görüntülendikçe P5'e taşı
Some checks failed
QA Gate (P0/P1) / Test affected app (pull_request) Has been cancelled
Faz 2 / veri göçü (analiz: /home/s/ss/plv2.md, bulgu psa-07 / consumers_jobs-05).

SORUN: PSA ve Volvo upstream'de P5'e taşındı, ama daha önce decode edilmiş
araçlar hâlâ eski yollara işaret ediyor: 324 PSA aracı (302'si bir kullanıcıya
bağlı, 33.535 kategori) donmuş 2024-02-13 anlık görüntüsünden, 36 Volvo aracı
(34'ü kullanıcılı, 10.195 kategori) ölü P4 ucundan besleniyor. `decodeVin`'in
db_hit kısa devresi yüzünden kod düzeltmesi bu satırlara ASLA ulaşmıyor; ayrıca
`categories` tablosundaki (vehicle_id, name, source) tekil indeksi yüzünden yeni
ağaç eskisinin üzerine yazılamıyor.

YAKLAŞIM: toplu yeniden decode fırtınası yerine **görüntülendikçe kendi kendini
onarma**. Tek hayatta kalan hesaba 360 aracı arka arkaya sormak yerine, her araç
ilk açılışında bir kez taşınır.

- `isStalePl24Architecture()` (pl24-tree.ts): saklı catalogPath `/psa`|`/volvo`
  ama servis bugün P5 ise bayat. Hâlâ P4 olan markalar (Ford/Opel/Hyundai/Kia/
  Nissan) ve zaten P5 olan kayıtlar dokunulmaz.
- `CategoriesService.healStalePl24Architecture()`: cache'i atlayarak yeniden
  decode eder, eski ağacı SİLİP yenisini tek transaction'da yazar, raw_data'yı
  yeni catalogInfo ile günceller, `fully_fetched`'i düşürür ve
  `cat:tree` / `prefetch:complete` / `prefetch:noresult` anahtarlarını temizler.

Bilinçli olarak temkinli:
- Yeniden decode başarısız olur veya kategori dönmezse ESKİ ağaç korunur
  (bayat veri, veri yokluğundan iyidir).
- Deneme `pl24:heal:<vehicleId>` ile günde bir kez (setNx) — çözülemeyen bir VIN
  her sayfa görüntülemesinde PL24'e gidemez.
- Eski satırlar ancak yeni ağaç elde edildikten SONRA siliniyor (tekil indeks).
- PL24 bütçesi/hız sınırı zaten üstte: her onarma ~2 upstream isteği.

Test: `isStalePl24Architecture` 4 test (bayat tespiti, zaten-P5, hâlâ-P4 markalar,
eksik bilgi) + onarma akışı 4 test (ağaç değişir, P5 kayda dokunulmaz, boş decode
eski ağacı korur, günlük kilit). 221 test geçti; tsc + biome temiz.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-09-17 16:53:15 +03:00
45 changed files with 4413 additions and 67 deletions

View File

@@ -0,0 +1,20 @@
-- Mitsubishi'nin parça listesi grup sanılıyordu.
--
-- `/p5mitsubishi/extern/details/vinDetails` yanıtı `partno`/`qty` taşıyan bir
-- PARÇA listesi, ama her kaydın kendi linki `partInfoTable` (parça-detay) ve ne
-- wid ne de yol sınıflandırıcıda karşılık buluyordu. Sonuç: 2.193 parça listesi
-- grup düğümüne dönüştü ve içlerindeki 19.576 tekil parça ("SCREW,LOCK CYLINDER",
-- "BOLT,STEERING COLUMN WASHER") kategori olarak kaydedildi. Ölçüm (prod,
-- 2026-09-20): bu 19.576 sahte düğümün TOPLAM 2 tanesinde parça var, ve her
-- prefetch turunda yeniden çekiliyorlar.
--
-- Sınıflandırıcı düzeltildi (detailsTable artık yaprak, partInfoTable hiç
-- kuyruklanmıyor), ama okuma yolu bir düğümün ÖNCE çocuklarına bakıyor: sahte
-- çocuklar dururken parça listesi asla çekilmez. Bu yüzden satırların silinmesi
-- düzeltmenin parçası, ayrı bir temizlik değil.
--
-- Güvenli: yalnız pl24 kaynaklı ve yalnız bu iki imzayı taşıyan satırlar; ikisi
-- de prod'da %100 Mitsubishi. Silinen ~2 parça satırı üst listeden yeniden gelir.
DELETE FROM "categories"
WHERE "source" = 'pl24'
AND ("link_wid" = 'partInfoTable' OR "link_path" ILIKE '%/details/vinpartinfo%');

View File

@@ -0,0 +1,34 @@
-- RPartStore (rpartstore.renault.com, Renault Group dealer portal) VIN-decode
-- fallback cache (feature-flagged: RPARTSTORE_ENABLED). One row per VIN. The
-- rpartstore-decode BullMQ job asks the RPartStore BFF (STOMP over WebSocket)
-- for Renault/Dacia VINs the normal chain (pcat/PL24/emex) can't identify, then
-- matches the decoded model to an EXISTING PL24 catalog_vehicle. RPartStore is
-- ONLY a decode oracle here — parts are served from PL24's existing catalog.
-- status: 'pending' | 'decoded' | 'not_found' | 'capped' | 'failed'.
CREATE TABLE IF NOT EXISTS "rpartstore_decodes" (
"vin" text PRIMARY KEY NOT NULL,
"status" text DEFAULT 'pending' NOT NULL,
"brand_name" text,
"model" text,
"model_code" text,
"family_code" text,
"model_year" text,
"engine" text,
"gearbox" text,
"energy_type" text,
"manufacturing_date" text,
"catalog_vehicle_id" uuid,
"raw" jsonb,
"attempts" integer DEFAULT 0 NOT NULL,
"created_at" timestamp with time zone DEFAULT now() NOT NULL,
"updated_at" timestamp with time zone,
"decoded_at" timestamp with time zone
);
--> statement-breakpoint
DO $$ BEGIN
ALTER TABLE "rpartstore_decodes" ADD CONSTRAINT "rpartstore_decodes_catalog_vehicle_id_catalog_vehicles_id_fk" FOREIGN KEY ("catalog_vehicle_id") REFERENCES "public"."catalog_vehicles"("id") ON DELETE set null ON UPDATE no action;
EXCEPTION
WHEN duplicate_object THEN null;
END $$;
--> statement-breakpoint
CREATE INDEX IF NOT EXISTS "rpartstore_decodes_status_idx" ON "rpartstore_decodes" USING btree ("status");

View File

@@ -253,6 +253,20 @@
"when": 1786435493000, "when": 1786435493000,
"tag": "0035_subscription_dunning", "tag": "0035_subscription_dunning",
"breakpoints": true "breakpoints": true
},
{
"idx": 36,
"version": "7",
"when": 1786521893000,
"tag": "0036_pl24_mitsubishi_partinfo_cleanup",
"breakpoints": true
},
{
"idx": 37,
"version": "7",
"when": 1790332800000,
"tag": "0037_rpartstore_decodes",
"breakpoints": true
} }
] ]
} }

View File

@@ -0,0 +1,192 @@
import { describe, expect, it, vi } from "vitest";
import { isStaleBrowseArchitecture } from "../integrations/pl24/pl24-tree";
import { CatalogService } from "./catalog.service";
/**
* Regression lock for the pinned-browse-rows bug (plv2.md, finding consumers-09).
*
* `catalog_vehicles` is a permanent cache: `getModels` returns early as soon as
* a brand has any row, so `fetchVehicleList` never runs again for that brand.
* The 188 rows listed while PSA and Volvo were still P4 (127 LEGACY_PSA + 61
* LEGACY_VOLVO on prod, 2026-09-20) were therefore stuck serving the frozen
* 2024-02-13 PSA snapshot and a Volvo endpoint that answers HTTP 503.
*/
type Row = {
id: string;
serviceName: string;
architecture: string | null;
brandId: string | null;
};
const psaRow = (id: string): Row => ({
id,
serviceName: "peugeot_parts",
architecture: "LEGACY_PSA",
brandId: "b1",
});
const freshRow = (id: string): Row => ({
id,
serviceName: "peugeot_parts",
architecture: "P5_MODERN",
brandId: "b1",
});
function makeService(opts: {
lockFree?: boolean;
fetched?: Array<{ model: string; vehicleId: string }>;
fetchThrows?: boolean;
refetched?: Row[];
}) {
const deleted: string[][] = [];
const inserted: unknown[][] = [];
const tx = {
insert: () => ({
values: (v: unknown[]) => {
inserted.push(v);
return { onConflictDoNothing: async () => undefined };
},
}),
delete: () => ({
where: async (cond: unknown) => {
// drizzle's inArray() puts its bound values in a `queryChunks` entry
// that is itself an array of Param objects; pull the plain ids back out
// so a test can assert WHICH rows were removed, not just how many.
const chunks = (cond as { queryChunks?: unknown[] })?.queryChunks ?? [];
const params = chunks.find((c): c is unknown[] => Array.isArray(c)) ?? [];
const ids = params
.map((param) => (param as { value?: unknown })?.value)
.filter((v): v is string => typeof v === "string");
deleted.push(ids);
return undefined;
},
}),
};
const db = {
select: () => ({ from: () => ({ where: async () => opts.refetched ?? [] }) }),
transaction: async (cb: (t: typeof tx) => Promise<void>) => cb(tx),
};
const redis = { setNx: vi.fn(async () => opts.lockFree !== false) };
const pl24Service = {
fetchVehicleList: vi.fn(async () => {
if (opts.fetchThrows) throw new Error("upstream 503");
return opts.fetched ?? [];
}),
};
const svc = new CatalogService(
db as never,
redis as never,
pl24Service as never,
{} as never,
) as never as {
healStaleBrowseRows: (brand: string, rows: Row[], where: unknown[]) => Promise<Row[] | null>;
};
return { svc, redis, pl24Service, deleted, inserted };
}
describe("isStaleBrowseArchitecture", () => {
it("flags a row whose stored architecture is not the current one", () => {
expect(
isStaleBrowseArchitecture({
storedArchitecture: "LEGACY_PSA",
currentArchitecture: "P5_MODERN",
}),
).toBe(true);
expect(
isStaleBrowseArchitecture({
storedArchitecture: "LEGACY_VOLVO",
currentArchitecture: "P5_MODERN",
}),
).toBe(true);
});
it("leaves a matching row alone", () => {
expect(
isStaleBrowseArchitecture({
storedArchitecture: "P5_MODERN",
currentArchitecture: "P5_MODERN",
}),
).toBe(false);
// Brands that are still P4 upstream must not be touched.
expect(
isStaleBrowseArchitecture({
storedArchitecture: "LEGACY_FORD",
currentArchitecture: "LEGACY_FORD",
}),
).toBe(false);
});
it("does nothing without both sides (unknown service / unlabelled row)", () => {
expect(
isStaleBrowseArchitecture({ storedArchitecture: null, currentArchitecture: "P5_MODERN" }),
).toBe(false);
expect(
isStaleBrowseArchitecture({ storedArchitecture: "LEGACY_PSA", currentArchitecture: null }),
).toBe(false);
});
});
describe("healStaleBrowseRows", () => {
it("does not touch PL24 when every row is current", async () => {
const { svc, pl24Service, redis } = makeService({});
const out = await svc.healStaleBrowseRows("Peugeot", [freshRow("a"), freshRow("b")], []);
expect(out).toBeNull();
expect(pl24Service.fetchVehicleList).not.toHaveBeenCalled();
expect(redis.setNx).not.toHaveBeenCalled();
});
it("re-lists a stale brand and replaces its rows", async () => {
const refetched = [freshRow("new1"), freshRow("new2")];
const { svc, pl24Service, deleted, inserted } = makeService({
fetched: [
{ model: "208", vehicleId: "p5-208" },
{ model: "308", vehicleId: "p5-308" },
],
refetched,
});
const out = await svc.healStaleBrowseRows("Peugeot", [psaRow("old1"), psaRow("old2")], []);
expect(pl24Service.fetchVehicleList).toHaveBeenCalledWith("peugeot_parts");
expect(inserted).toHaveLength(1);
expect(inserted[0]).toHaveLength(2);
expect(deleted).toHaveLength(1);
expect(out).toEqual(refetched);
});
it("deletes only the stale rows when a service holds a mix", async () => {
// A row already re-listed under the new architecture conflicts on insert and
// is skipped, so deleting it would drop it for good.
const { svc, deleted } = makeService({
fetched: [{ model: "208", vehicleId: "p5-208" }],
refetched: [freshRow("keep")],
});
await svc.healStaleBrowseRows("Peugeot", [psaRow("old1"), freshRow("keep")], []);
expect(deleted).toHaveLength(1);
expect(deleted).toEqual([["old1"]]);
});
it("keeps the stale rows when the re-list throws", async () => {
const { svc, deleted, inserted } = makeService({ fetchThrows: true });
const out = await svc.healStaleBrowseRows("Volvo", [psaRow("old1")], []);
expect(out).toBeNull();
expect(inserted).toHaveLength(0);
expect(deleted).toHaveLength(0);
});
it("keeps the stale rows when upstream returns an empty model list", async () => {
const { svc, deleted, inserted } = makeService({ fetched: [] });
const out = await svc.healStaleBrowseRows("Volvo", [psaRow("old1")], []);
expect(out).toBeNull();
expect(inserted).toHaveLength(0);
expect(deleted).toHaveLength(0);
});
it("attempts at most once a day per brand", async () => {
const { svc, pl24Service } = makeService({ lockFree: false });
const out = await svc.healStaleBrowseRows("Peugeot", [psaRow("old1")], []);
expect(out).toBeNull();
expect(pl24Service.fetchVehicleList).not.toHaveBeenCalled();
});
});

View File

@@ -19,7 +19,7 @@ import {
userBrands, userBrands,
userSubscriptions, userSubscriptions,
} from "../database/schema/core"; } from "../database/schema/core";
import { isPl24LeafNode } from "../integrations/pl24/pl24-tree"; import { isPl24LeafNode, isStaleBrowseArchitecture } from "../integrations/pl24/pl24-tree";
import { PL24Service } from "../integrations/pl24/pl24.service"; import { PL24Service } from "../integrations/pl24/pl24.service";
import { import {
type PL24DecodedCategory, type PL24DecodedCategory,
@@ -170,7 +170,8 @@ export class CatalogService {
.where(and(...whereConditions)); .where(and(...whereConditions));
if (dbVehicles.length > 0) { if (dbVehicles.length > 0) {
return dbVehicles; const healed = await this.healStaleBrowseRows(brandName, dbVehicles, whereConditions);
return healed ?? dbVehicles;
} }
// Fetch from PL24 for each service // Fetch from PL24 for each service
@@ -225,6 +226,118 @@ export class CatalogService {
return allVehicles; return allVehicles;
} }
/**
* Re-list a brand whose stored browse rows were created under a retired PL24
* architecture, and replace them with the current ones.
*
* WHY (plv2.md, finding consumers-09): `getModels` treats `catalog_vehicles`
* as a permanent cache — one row for the brand and `fetchVehicleList` is never
* called again. The rows listed while PSA and Volvo were still P4 are pinned
* forever, so Peugeot/Citroën browse keeps serving the frozen 2024-02-13
* snapshot and Volvo/Polestar browse keeps calling an endpoint that answers
* HTTP 503. A code fix alone never reaches these rows.
*
* Deliberately conservative, mirroring `healStalePl24Architecture` on the
* VIN side:
* - only rows whose stored architecture disagrees with the service table;
* - one attempt per brand per day (Redis lock) so a permanently failing
* re-list cannot hammer the one surviving account on every page view;
* - a failed or empty re-list leaves the old rows in place — stale data
* beats an empty catalog;
* - the stale rows are deleted only once the replacements are committed,
* and per service, so one broken sub-catalog cannot wipe a working one.
*
* Deleting a browse row cascades to its categories and parts. That is
* intended: those rows describe the retired tree and would otherwise survive
* as unreachable orphans under the new listing.
*
* Returns the refreshed rows, or null when nothing was migrated (caller keeps
* what it already had).
*/
private async healStaleBrowseRows(
brandName: string,
dbVehicles: (typeof catalogVehicles.$inferSelect)[],
whereConditions: ReturnType<typeof eq>[],
): Promise<(typeof catalogVehicles.$inferSelect)[] | null> {
const isStale = (v: typeof catalogVehicles.$inferSelect) =>
isStaleBrowseArchitecture({
storedArchitecture: v.architecture,
currentArchitecture: PL24_SERVICE_CATALOGS[v.serviceName]?.architecture,
});
const staleServices = [...new Set(dbVehicles.filter(isStale).map((v) => v.serviceName))];
if (staleServices.length === 0) return null;
if (!(await this.redis.setNx(`pl24:browse-heal:${brandName}`, "1", 86_400))) return null;
this.logger.log(
`[pl24-browse-heal] ${brandName}: ${staleServices.join(", ")} listed under a retired architecture — re-listing`,
);
let migrated = 0;
for (const svc of staleServices) {
// ONLY the stale rows. A service can hold a mix — some rows already
// re-listed under the new architecture — and those must survive: the
// insert below skips them on conflict, so deleting them here would drop
// them for good.
const staleIds = dbVehicles
.filter((v) => v.serviceName === svc && isStale(v))
.map((v) => v.id);
try {
const fetched = await this.pl24Service.fetchVehicleList(svc);
if (fetched.length === 0) {
this.logger.warn(
`[pl24-browse-heal] ${svc}: upstream returned no models — keeping ${staleIds.length} stale row(s)`,
);
continue;
}
const config = PL24_SERVICE_CATALOGS[svc];
const brandId = dbVehicles.find((v) => v.serviceName === svc)?.brandId ?? null;
await this.db.transaction(async (tx) => {
await tx
.insert(catalogVehicles)
.values(
fetched.map((v) => ({
source: "pl24" as const,
serviceName: svc,
brandName,
brandId,
model: v.model,
year: v.year || null,
engine: v.engine || null,
bodyType: v.bodyType || null,
transmission: v.transmission || null,
market: v.market || null,
serviceVehicleId: v.vehicleId,
catalogPath: v.catalogPath || null,
architecture: config?.architecture || "P5_MODERN",
metadata: v.metadata || null,
categoriesFetched: false,
updatedAt: new Date(),
})),
)
.onConflictDoNothing();
if (staleIds.length > 0) {
await tx.delete(catalogVehicles).where(inArray(catalogVehicles.id, staleIds));
}
});
migrated += staleIds.length;
this.logger.log(
`[pl24-browse-heal] ${svc}: ${staleIds.length} stale row(s) → ${fetched.length} fresh model(s)`,
);
} catch (err) {
this.logger.warn(
`[pl24-browse-heal] ${svc} re-list failed, keeping stale rows: ${(err as Error).message}`,
);
}
}
if (migrated === 0) return null;
return await this.db
.select()
.from(catalogVehicles)
.where(and(...whereConditions));
}
/** /**
* Get a single catalog vehicle by ID. * Get a single catalog vehicle by ID.
*/ */

View File

@@ -17,6 +17,8 @@ function createService(db: any) {
}; };
const pl24Service = { const pl24Service = {
getCategories: vi.fn().mockResolvedValue([]), getCategories: vi.fn().mockResolvedValue([]),
// P4→P5 onarma yolu bunu çağırır (plv2 Faz 2).
decodeVin: vi.fn().mockResolvedValue(null),
}; };
const emexService = { const emexService = {
isSupported: vi.fn().mockReturnValue(false), isSupported: vi.fn().mockReturnValue(false),
@@ -462,3 +464,81 @@ describe("CategoriesService", () => {
}); });
}); });
}); });
// ── P4 → P5 kendi kendini onarma (plv2 Faz 2, bulgu psa-07) ──
// PSA/Volvo upstream'de P5'e taşındı; daha önce decode edilmiş araçlar donmuş
// 2024-02-13 PSA anlık görüntüsüne / ölü P4 Volvo ucuna işaret etmeye devam
// ediyor ve decodeVin'in db_hit kısa devresi kod düzeltmesini onlara ulaştırmıyor.
describe("CategoriesService — bayat PL24 mimarisini onarma", () => {
const makeVehicle = (catalogPath: string) => ({
id: "veh-1",
vin: "VF3MCBHZWHS150390",
rawData: { catalogInfo: { serviceName: "peugeot_parts", catalogPath } },
});
const setup = (catalogPath: string, decodeResult: unknown) => {
const { service, redis, pl24Service } = createService({} as never);
const p = service as unknown as {
healStalePl24Architecture(
id: string,
v: { vin: string | null; rawData: unknown },
n: number,
): Promise<boolean>;
db: unknown;
};
// setNx: ilk denemeye izin ver (günlük tek deneme kilidi)
(redis as unknown as { setNx: unknown }).setNx = vi.fn().mockResolvedValue(true);
(redis as unknown as { del: unknown }).del = vi.fn().mockResolvedValue(undefined);
pl24Service.decodeVin = vi.fn().mockResolvedValue(decodeResult);
return { service, p, redis, pl24Service, vehicle: makeVehicle(catalogPath) };
};
it("P4 PSA yolundaki araç yeniden decode edilir ve ağaç değişir", async () => {
const fresh = {
catalogInfo: { serviceName: "peugeot_parts", catalogPath: "/p5psa" },
categories: [
{
code: "_FCT0001",
nameTr: "Mekanik",
nameEn: "Mechanical",
linkPath: "/p5psa/x",
linkWid: "mainGroupTable",
},
],
};
const { p, pl24Service, vehicle } = setup("/psa/peugeot_parts", fresh);
const inserted: unknown[] = [];
(p as unknown as { db: unknown }).db = {
transaction: async (fn: (tx: unknown) => Promise<void>) => {
await fn({
delete: () => ({ where: async () => undefined }),
insert: () => ({ values: async (rows: unknown[]) => inserted.push(...rows) }),
update: () => ({ set: () => ({ where: async () => undefined }) }),
});
},
};
await expect(p.healStalePl24Architecture("veh-1", vehicle, 120)).resolves.toBe(true);
expect(pl24Service.decodeVin).toHaveBeenCalledWith("VF3MCBHZWHS150390");
expect(inserted).toHaveLength(1);
expect((inserted[0] as { name: string }).name).toBe("Mekanik");
});
it("zaten P5 olan araca dokunulmaz (decode çağrılmaz)", async () => {
const { p, pl24Service, vehicle } = setup("/p5psa", null);
await expect(p.healStalePl24Architecture("veh-1", vehicle, 6)).resolves.toBe(false);
expect(pl24Service.decodeVin).not.toHaveBeenCalled();
});
it("yeniden decode boş dönerse eski ağaç korunur", async () => {
const { p, vehicle } = setup("/psa/peugeot_parts", { categories: [] });
await expect(p.healStalePl24Architecture("veh-1", vehicle, 120)).resolves.toBe(false);
});
it("günde bir denenir (setNx kilidi)", async () => {
const { p, redis, pl24Service, vehicle } = setup("/psa/peugeot_parts", { categories: [] });
(redis as unknown as { setNx: unknown }).setNx = vi.fn().mockResolvedValue(false);
await expect(p.healStalePl24Architecture("veh-1", vehicle, 120)).resolves.toBe(false);
expect(pl24Service.decodeVin).not.toHaveBeenCalled();
});
});

View File

@@ -21,8 +21,13 @@ import { PartsCatalogsService } from "../integrations/parts-catalogs/parts-catal
import { PcatGroup } from "../integrations/parts-catalogs/parts-catalogs.types"; import { PcatGroup } from "../integrations/parts-catalogs/parts-catalogs.types";
import { PL24FordLegacyService } from "../integrations/pl24/pl24-ford-legacy.service"; import { PL24FordLegacyService } from "../integrations/pl24/pl24-ford-legacy.service";
import { PL24PsaService } from "../integrations/pl24/pl24-psa.service"; import { PL24PsaService } from "../integrations/pl24/pl24-psa.service";
import { isPl24GroupNode, isPl24LeafNode } from "../integrations/pl24/pl24-tree"; import {
isPl24GroupNode,
isPl24LeafNode,
isStalePl24Architecture,
} from "../integrations/pl24/pl24-tree";
import { PL24Service } from "../integrations/pl24/pl24.service"; import { PL24Service } from "../integrations/pl24/pl24.service";
import { getServiceApiPath } from "../integrations/pl24/pl24.types";
import { classifyNode, foldName, mapToCanonical } from "../jobs/canonical-lexicon"; import { classifyNode, foldName, mapToCanonical } from "../jobs/canonical-lexicon";
import { RedisService } from "../redis/redis.service"; import { RedisService } from "../redis/redis.service";
import { StorageService } from "../storage/storage.service"; import { StorageService } from "../storage/storage.service";
@@ -71,6 +76,20 @@ export class CategoriesService {
.from(categories) .from(categories)
.where(eq(categories.vehicleId, vehicleId)); .where(eq(categories.vehicleId, vehicleId));
// ── P4 → P5 self-healing (plv2.md Faz 2, bulgu psa-07) ──
// PSA and Volvo moved to P5 upstream; rows decoded before that still point at
// the frozen 2024-02-13 PSA snapshot or the dead P4 Volvo endpoint, and
// decodeVin's db_hit short-circuit means no code fix ever reaches them. Heal
// one vehicle per view instead of a bulk re-decode storm against the single
// surviving account.
const healed = await this.healStalePl24Architecture(vehicleId, vehicle, dbCategories.length);
if (healed) {
dbCategories = await this.db
.select()
.from(categories)
.where(eq(categories.vehicleId, vehicleId));
}
// If no categories in DB, fetch from PL24 // If no categories in DB, fetch from PL24
if (dbCategories.length === 0 && vehicle.rawData) { if (dbCategories.length === 0 && vehicle.rawData) {
const rawData = vehicle.rawData as any; const rawData = vehicle.rawData as any;
@@ -988,6 +1007,104 @@ export class CategoriesService {
* The per-node trail excludes the node itself, ordered root-first. Used by the * The per-node trail excludes the node itself, ordered root-first. Used by the
* catalog search so each hit can show where it sits in the tree. * catalog search so each hit can show where it sits in the tree.
*/ */
/**
* Re-decode a vehicle whose stored catalogInfo still points at a retired PL24
* architecture (P4 PSA/Volvo) and replace its category tree with the P5 one.
*
* Returns true when the tree was rebuilt. Deliberately conservative:
* - only PSA/Volvo rows whose service is P5 today are touched;
* - a failed or empty re-decode leaves the old tree in place (stale data beats
* no data), and the attempt is remembered for a day so a permanently
* unresolvable VIN cannot re-hit PL24 on every page view;
* - the old rows are deleted only once the new tree is in hand, because the
* unique (vehicle_id, name, source) index would otherwise reject the insert.
*/
private async healStalePl24Architecture(
vehicleId: string,
vehicle: { vin: string | null; rawData: unknown },
existingCategoryCount: number,
): Promise<boolean> {
const rawData = vehicle.rawData as {
catalogInfo?: { serviceName?: string; catalogPath?: string };
} | null;
const catalogInfo = rawData?.catalogInfo;
const serviceName = catalogInfo?.serviceName;
if (!vehicle.vin || !serviceName) return false;
if (
!isStalePl24Architecture({
catalogPath: catalogInfo?.catalogPath,
currentApiPath: getServiceApiPath(serviceName),
})
) {
return false;
}
const attemptKey = `pl24:heal:${vehicleId}`;
if (!(await this.redis.setNx(attemptKey, "1", 86_400))) return false;
this.logger.log(
`[pl24-heal] ${vehicle.vin} (${serviceName}) was decoded on ${catalogInfo?.catalogPath} — re-decoding on P5`,
);
let decoded: Awaited<ReturnType<PL24Service["decodeVin"]>> | null = null;
try {
await this.redis.del(`pl24:vehicle:${vehicle.vin}`);
decoded = await this.pl24Service.decodeVin(vehicle.vin);
} catch (err) {
this.logger.warn(`[pl24-heal] ${vehicle.vin} re-decode failed: ${(err as Error).message}`);
return false;
}
if (!decoded?.categories?.length) {
this.logger.warn(
`[pl24-heal] ${vehicle.vin} re-decode returned no categories — keeping old tree`,
);
return false;
}
await this.db.transaction(async (tx) => {
await tx.delete(categories).where(eq(categories.vehicleId, vehicleId));
const seen = new Set<string>();
const rows = decoded.categories
.filter((c) => {
const name = c.nameTr || c.nameEn;
if (!name || seen.has(name)) return false;
seen.add(name);
return true;
})
.map((c) => ({
vehicleId,
catalogVehicleId: null as string | null,
name: c.nameTr || c.nameEn,
nameOriginal: c.nameEn,
parentId: null as string | null,
externalId: c.code,
linkPath: c.linkPath || null,
linkWid: c.linkWid || null,
source: "pl24" as const,
}));
if (rows.length) await tx.insert(categories).values(rows);
await tx
.update(vehicles)
.set({
rawData: { ...(rawData ?? {}), catalogInfo: decoded.catalogInfo },
fullyFetched: false,
fullyFetchedAt: null,
})
.where(eq(vehicles.id, vehicleId));
});
await Promise.all([
this.redis.del(`cat:tree:${vehicleId}`),
this.redis.del(`prefetch:complete:${vehicleId}`),
this.redis.del(`prefetch:noresult:${vehicleId}`),
]);
this.logger.log(
`[pl24-heal] ${vehicle.vin} migrated to ${decoded.catalogInfo?.catalogPath}: ${existingCategoryCount} stale → ${decoded.categories.length} fresh categories`,
);
return true;
}
private async buildBreadcrumbs( private async buildBreadcrumbs(
ids: string[], ids: string[],
): Promise<Map<string, Array<{ id: string; name: string }>>> { ): Promise<Map<string, Array<{ id: string; name: string }>>> {

View File

@@ -409,6 +409,36 @@ export const vinpinDecodes = pgTable("vinpin_decodes", {
decodedAt: timestamp("decoded_at", { withTimezone: true }), decodedAt: timestamp("decoded_at", { withTimezone: true }),
}); });
// ─── RPartStore decodes (Renault/Dacia decode-oracle cache — one row per VIN) ──
// Feature-flagged (RPARTSTORE_ENABLED). Populated by the rpartstore-decode BullMQ
// job which asks the RPartStore BFF (rpartstore.renault.com, STOMP/WebSocket) for
// Renault/Dacia VINs the normal chain (pcat/PL24/emex) can't identify. On success
// the decoded model is matched to an EXISTING PL24 catalog_vehicle and parts are
// served from there — RPartStore is ONLY a decode oracle. Hard daily cap on
// searches (RPARTSTORE_DAILY_CAP); over-cap VINs are 'capped' and retried after 24 h.
export const rpartstoreDecodes = pgTable("rpartstore_decodes", {
vin: text("vin").primaryKey(),
// 'pending' | 'decoded' | 'not_found' | 'capped' | 'failed'
status: text("status").notNull().default("pending"),
brandName: text("brand_name"),
model: text("model"),
modelCode: text("model_code"),
familyCode: text("family_code"),
modelYear: text("model_year"),
engine: text("engine"),
gearbox: text("gearbox"),
energyType: text("energy_type"),
manufacturingDate: text("manufacturing_date"),
catalogVehicleId: uuid("catalog_vehicle_id").references(() => catalogVehicles.id, {
onDelete: "set null",
}),
raw: jsonb("raw"),
attempts: integer("attempts").notNull().default(0),
createdAt: timestamp("created_at", { withTimezone: true }).defaultNow().notNull(),
updatedAt: timestamp("updated_at", { withTimezone: true }),
decodedAt: timestamp("decoded_at", { withTimezone: true }),
});
// ─── Vehicles (shared config — one record per VIN) ── // ─── Vehicles (shared config — one record per VIN) ──
export const vehicles = pgTable( export const vehicles = pgTable(
"vehicles", "vehicles",

View File

@@ -0,0 +1,198 @@
{
"link": {
"wid": "mainGroupTable",
"path": "/p5volvo/extern/groups/vin/mainGroup?lang=en&model=308&modelYear=1620&partnerGroup=46&serviceName=volvo_parts&upds=2026-08-24--11-02&vin=YV1AS84ABD1168166"
},
"segments": {
"vinfoBasic": {
"records": [
{
"values": {
"description": "Vehicle Identification No.",
"value": "YV1AS84ABD1168166"
}
},
{
"values": {
"description": "Year",
"value": "2013"
}
},
{
"values": {
"description": "Model",
"value": "S80 (07\\-)"
}
},
{
"values": {
"description": "Km/h or mph",
"value": "K"
}
},
{
"values": {
"description": "Factory code",
"value": "21"
}
},
{
"values": {
"description": "Steering gear prod no",
"value": "31360538"
}
},
{
"values": {
"description": "Structure week",
"value": "201236"
}
},
{
"values": {
"description": "Type",
"value": "S80"
}
},
{
"values": {
"description": "Chassis",
"value": "168166"
}
},
{
"values": {
"description": "Partner group",
"value": "Europe"
}
},
{
"values": {
"description": "Upholstery/interior code",
"value": "210100"
}
},
{
"values": {
"description": "Upholstery/interior",
"value": "LEATHER/ANTHR/QRTZCEIL/NV"
}
},
{
"values": {
"description": "Exterior color",
"value": "61400"
}
},
{
"values": {
"description": "Exterior color",
"value": "WHITE SOLID ICE WHITE"
}
},
{
"values": {
"description": "Body style code",
"value": "0"
}
},
{
"values": {
"description": "Body style",
"value": "Sedan"
}
},
{
"values": {
"description": "Special vehicle code",
"value": " "
}
},
{
"values": {
"description": "Special vehicles",
"value": " "
}
},
{
"values": {
"description": "Sales type",
"value": "42"
}
},
{
"values": {
"description": "Sales type",
"value": "SALES VERSION 42"
}
},
{
"values": {
"description": "Market code",
"value": "49"
}
},
{
"values": {
"description": "Market",
"value": "TR"
}
},
{
"values": {
"description": "Engine Code",
"value": "84"
}
},
{
"values": {
"description": "Engine",
"value": "D4162T"
}
},
{
"values": {
"description": "Engine part no",
"value": "6906309"
}
},
{
"values": {
"description": "Engine serial no",
"value": "00000000000004138845 / 0ELD61 2208122233587"
}
},
{
"values": {
"description": "Transmission Code",
"value": "B"
}
},
{
"values": {
"description": "Transmission",
"value": "6\\-PSHIFT 2WD / MPS6"
}
},
{
"values": {
"description": "Transmission part no",
"value": "1285041"
}
},
{
"values": {
"description": "Transmission serial no",
"value": "00AWBB1 170812170654"
}
},
{
"values": {
"description": "Chassis code",
"value": "35659B7A6276"
}
}
]
}
}
}

View File

@@ -0,0 +1,198 @@
{
"link": {
"wid": "mainGroupTable",
"path": "/p5volvo/extern/groups/vin/mainGroup?lang=tr&model=308&modelYear=1620&partnerGroup=46&serviceName=volvo_parts&upds=2026-08-24--11-02&vin=YV1AS84ABD1168166"
},
"segments": {
"vinfoBasic": {
"records": [
{
"values": {
"description": "Sasi numarasi",
"value": "YV1AS84ABD1168166"
}
},
{
"values": {
"description": "Model yili",
"value": "2013"
}
},
{
"values": {
"description": "Model",
"value": "S80 (07\\-)"
}
},
{
"values": {
"description": "Hız Birimi",
"value": "K"
}
},
{
"values": {
"description": "Fabrika Kodu",
"value": "21"
}
},
{
"values": {
"description": "Direksiyon Kutusu Üretim No",
"value": "31360538"
}
},
{
"values": {
"description": "Üretim Haftası",
"value": "201236"
}
},
{
"values": {
"description": "Türü",
"value": "S80"
}
},
{
"values": {
"description": "Şasi",
"value": "168166"
}
},
{
"values": {
"description": "Ortak grubu",
"value": "Europe"
}
},
{
"values": {
"description": "Döşeme/iç mekan kodu",
"value": "210100"
}
},
{
"values": {
"description": "Döşeme",
"value": "LEATHER/ANTHR/QRTZCEIL/NV"
}
},
{
"values": {
"description": "Dis rengi",
"value": "61400"
}
},
{
"values": {
"description": "Dis rengi",
"value": "WHITE SOLID ICE WHITE"
}
},
{
"values": {
"description": "Karoseri Tipi Kodu",
"value": "0"
}
},
{
"values": {
"description": "Kaporta Stili",
"value": "Sedan"
}
},
{
"values": {
"description": "Özel Araç Kodu",
"value": " "
}
},
{
"values": {
"description": "Özel araçlar",
"value": " "
}
},
{
"values": {
"description": "Satis tipi",
"value": "42"
}
},
{
"values": {
"description": "Satis tipi",
"value": "SALES VERSION 42"
}
},
{
"values": {
"description": "Piyasa Kodu",
"value": "49"
}
},
{
"values": {
"description": "Market",
"value": "TR"
}
},
{
"values": {
"description": "Motor kodu",
"value": "84"
}
},
{
"values": {
"description": "Motor",
"value": "D4162T"
}
},
{
"values": {
"description": "Motor Parça No",
"value": "6906309"
}
},
{
"values": {
"description": "Motor Seri Numarası",
"value": "00000000000004138845 / 0ELD61 2208122233587"
}
},
{
"values": {
"description": "Şanzıman kodu",
"value": "B"
}
},
{
"values": {
"description": "Şanzıman",
"value": "6\\-PSHIFT 2WD / MPS6"
}
},
{
"values": {
"description": "Şanzıman Parça No",
"value": "1285041"
}
},
{
"values": {
"description": "Şanzıman Seri No",
"value": "00AWBB1 170812170654"
}
},
{
"values": {
"description": "Şasi Kodu",
"value": "35659B7A6276"
}
}
]
}
}
}

View File

@@ -1,5 +1,6 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { PL24AuthService } from "./pl24-auth.service"; import { PL24AuthService } from "./pl24-auth.service";
import { PL24BudgetExceededError } from "./pl24-budget.service";
// Oturum modeli (1 login → PL24TOKEN çerezi → yalnız authorize ile yenileme), // Oturum modeli (1 login → PL24TOKEN çerezi → yalnız authorize ile yenileme),
// PL24_TR_DISABLED köprüsü ve login devre kesici. Ağ yok: global fetch stub'ı; // PL24_TR_DISABLED köprüsü ve login devre kesici. Ağ yok: global fetch stub'ı;
@@ -459,3 +460,33 @@ describe("PL24AuthService — 1 saatlik kesinti alarmı (Telegram)", () => {
expect(redis.store.has("pl24:auth:fail-since:de")).toBe(false); expect(redis.store.has("pl24:auth:fail-since:de")).toBe(false);
}); });
}); });
/**
* Kendi günlük tavanımız dolduğunda bu bir GİRİŞ HATASI değildir.
*
* Prod 2026-09-21: tavan dolunca `attemptLogin` içindeki `budget.consume()`
* fırlattı, dış `catch` bunu `{ error }` düzleştirdi, `loginWithLock` giriş
* hatası sanıp devre kesiciyi tetikledi ve kesinti sayacını başlattı — bir saat
* sonra Telegram "hesap banlanmış olabilir" dedi. Hesap sağlamdı (tarayıcıdan
* giriş çalışıyordu, aynı gün 143 başarılı auth ve 0 × 401).
*/
describe("bütçe reddi ban alarmı üretmemeli", () => {
it("bütçe hatası olduğu gibi yukarı çıkar, giriş hatasına çevrilmez", async () => {
const { svc, redis, telegram, budget } = makeService();
budget.consume = async () => {
throw new PL24BudgetExceededError("user", 1205, 1200);
};
const fetchSpy = vi.fn();
vi.stubGlobal("fetch", fetchSpy);
await expect((svc as never as Exposed).login("de")).rejects.toBeInstanceOf(
PL24BudgetExceededError,
);
// Kesinti sayacı başlamamalı → Telegram susmalı.
expect(telegram.sent).toHaveLength(0);
const failKeys = Object.keys(redis.store ?? {}).filter((k) => k.includes("fail-since"));
expect(failKeys).toHaveLength(0);
vi.unstubAllGlobals();
});
});

View File

@@ -26,7 +26,7 @@ import { Injectable, Logger, type OnModuleInit, UnauthorizedException } from "@n
import { ConfigService } from "@nestjs/config"; import { ConfigService } from "@nestjs/config";
import { TelegramService } from "../../common/telegram.service"; import { TelegramService } from "../../common/telegram.service";
import { RedisService } from "../../redis/redis.service"; import { RedisService } from "../../redis/redis.service";
import { PL24BudgetService } from "./pl24-budget.service"; import { PL24BudgetExceededError, PL24BudgetService } from "./pl24-budget.service";
import { PL24_DEFAULTS, PL24_ENDPOINTS, PL24_USER_AGENT } from "./pl24.constants"; import { PL24_DEFAULTS, PL24_ENDPOINTS, PL24_USER_AGENT } from "./pl24.constants";
import { import {
PL24AuthorizeRequest, PL24AuthorizeRequest,
@@ -485,6 +485,15 @@ export class PL24AuthService implements OnModuleInit {
return { sessionToken }; return { sessionToken };
} catch (error) { } catch (error) {
const err = error as Error; const err = error as Error;
// OUR OWN daily cap refusing to spend is not a login failure. Flattening it
// into `{ error }` made `loginWithLock` treat it as one: it tripped the
// login circuit breaker and started the outage timer, so an hour later the
// Telegram alert claimed the account might be banned — while the account
// was fine (browser login worked, telemetry showed 0 × 401 and 143
// successful auth calls that same day). Observed on prod 2026-09-21:
// "PL24 daily HTTP budget exhausted (user lane: 1205/1200)". Let it through
// as itself so the caller can skip PL24 without anyone being paged.
if (error instanceof PL24BudgetExceededError) throw error;
if (err.name === "TimeoutError") return { error: "giris zaman asimina ugradi" }; if (err.name === "TimeoutError") return { error: "giris zaman asimina ugradi" };
return { error: err.message }; return { error: err.message };
} }
@@ -524,6 +533,9 @@ export class PL24AuthService implements OnModuleInit {
return token; return token;
} catch (error) { } catch (error) {
const err = error as Error; const err = error as Error;
// Same reasoning as attemptLogin: a self-imposed budget stop is not an
// authorization problem and must not be reported as one.
if (error instanceof PL24BudgetExceededError) throw error;
this.logger.error(`Service authorization error (${account}): ${err.message}`); this.logger.error(`Service authorization error (${account}): ${err.message}`);
throw err instanceof UnauthorizedException throw err instanceof UnauthorizedException
? err ? err

View File

@@ -139,3 +139,28 @@ describe("PL24BudgetService — telemetri", () => {
expect(telemetry.events[0].statusCode).toBeNull(); expect(telemetry.events[0].statusCode).toBeNull();
}); });
}); });
/**
* Bütçe reddi bir giriş hatası DEĞİLDİR (plv2.md, bulgu auth-16).
*
* Prod 2026-09-21: günlük tavan dolunca `attemptLogin` içindeki
* `budget.consume()` fırlattı, dış `catch` bunu `{ error }` düzleştirdi,
* `loginWithLock` bunu giriş hatası sanıp devre kesiciyi tetikledi ve kesinti
* sayacını başlattı → bir saat sonra Telegram "hesap banlanmış olabilir" dedi.
* Oysa hesap sağlamdı: tarayıcıdan giriş çalışıyordu, aynı gün 143 başarılı
* auth çağrısı ve 0 × 401 vardı.
*/
describe("PL24BudgetExceededError — kendi frenimiz, upstream hatası değil", () => {
it("kendi hata sınıfını taşır, düz Error değil", () => {
const err = new PL24BudgetExceededError("user", 1205, 1200);
expect(err).toBeInstanceOf(PL24BudgetExceededError);
expect(err).toBeInstanceOf(Error);
});
it("mesajı hangi şeridin ve hangi sayının durduğunu söyler", () => {
const err = new PL24BudgetExceededError("user", 1205, 1200);
expect(err.message).toContain("1205");
expect(err.message).toContain("1200");
expect(err.message.toLowerCase()).toContain("user");
});
});

View File

@@ -64,6 +64,14 @@ export class PL24BudgetService {
const key = this.dayKey(); const key = this.dayKey();
spent = await this.redis.incr(key); spent = await this.redis.incr(key);
if (spent === 1) await this.redis.expire(key, 8 * 86_400); if (spent === 1) await this.redis.expire(key, 8 * 86_400);
// Per-lane counters. The shared total tells us WHEN the cap was hit but not
// WHO spent it, and sizing the cap needs that split: on 2026-09-21 the
// total hit 1200 and every user decode after 16:10 UTC was refused, with
// no way to tell from `proxy_logs` how much of it was warm-up prefetch
// versus somebody waiting on a screen.
const laneKey = `${key}:${backfill ? "worker" : "user"}`;
const laneSpent = await this.redis.incr(laneKey);
if (laneSpent === 1) await this.redis.expire(laneKey, 8 * 86_400);
} catch { } catch {
return; // Redis down → never block PL24 on telemetry return; // Redis down → never block PL24 on telemetry
} }

View File

@@ -0,0 +1,127 @@
import { beforeEach, describe, expect, it, vi } from "vitest";
import { PL24Service } from "./pl24.service";
/**
* Model-list entry point resolution (plv2.md, bulgu p5core-08).
*
* `BACKEND_MODEL_PATH` eksik kaldığında eski davranış yedi bilinen yolu sırayla
* denemekti: paylaşılan günlük bütçeden çağrı başına yedi isteğe kadar, ve
* listede olmayan bir yol için sessiz sıfır sonuç. Volvo tam olarak buydu —
* `/extern/vehicles/models` (çoğul "vehicles") listede yok, dolayısıyla
* Volvo/Polestar browse sıfır model tohumlayıp emekli P4 listesini sunmaya
* devam ediyordu.
*
* Otoritatif kaynak her P5 backend'inin kendi `/extern/catmeta` yanıtındaki
* `data.catalogEntryPoint.path`. Canlı doğrulama (2026-09-20):
* p5volvo → /p5volvo/extern/vehicles/models (60 model)
* p5psa → /p5psa/extern/vehicle/catalogs
*/
type Resolver = {
resolveModelPathFromCatmeta: (
base: string,
svc: string,
headers: Record<string, string>,
) => Promise<string | null>;
baseUrl: string;
language: string;
redis: {
get: (k: string) => Promise<string | null>;
set: (k: string, v: string, ttl?: number) => Promise<unknown>;
};
logger: { log: (m: string) => void; warn: (m: string) => void };
};
function makeService(opts: {
cached?: string | null;
catmeta?: unknown;
status?: number;
throws?: boolean;
}) {
const sets: Array<[string, string]> = [];
const svc = Object.create(PL24Service.prototype) as unknown as Resolver;
svc.baseUrl = "https://pl24.test";
svc.language = "tr";
svc.redis = {
get: vi.fn(async () => opts.cached ?? null),
set: vi.fn(async (k: string, v: string) => {
sets.push([k, v]);
return undefined;
}),
};
svc.logger = { log: vi.fn(), warn: vi.fn() };
const fetchMock = vi.fn(async () => {
if (opts.throws) throw new Error("network down");
return {
ok: (opts.status ?? 200) < 400,
status: opts.status ?? 200,
json: async () => opts.catmeta ?? {},
};
});
vi.stubGlobal("fetch", fetchMock);
return { svc, sets, fetchMock };
}
const volvoMeta = {
data: {
catalogEntryPoint: {
wid: "modelsTable",
path: "/p5volvo/extern/vehicles/models?lang=en&serviceName=volvo_parts",
},
},
};
describe("resolveModelPathFromCatmeta", () => {
beforeEach(() => vi.unstubAllGlobals());
it("catmeta'daki giriş noktasını backend'e göreli yola indirger", async () => {
const { svc } = makeService({ catmeta: volvoMeta });
const out = await svc.resolveModelPathFromCatmeta("/p5volvo", "volvo_parts", {});
expect(out).toBe("/extern/vehicles/models");
});
it("çözülen yolu 30 gün cache'ler", async () => {
const { svc, sets } = makeService({ catmeta: volvoMeta });
await svc.resolveModelPathFromCatmeta("/p5volvo", "volvo_parts", {});
expect(sets[0][0]).toBe("pl24:modelpath:p5volvo");
expect(sets[0][1]).toBe("/extern/vehicles/models");
});
it("cache'lenmiş yol için upstream'e hiç gitmez", async () => {
const { svc, fetchMock } = makeService({ cached: "/extern/vehicles/models" });
const out = await svc.resolveModelPathFromCatmeta("/p5volvo", "volvo_parts", {});
expect(out).toBe("/extern/vehicles/models");
expect(fetchMock).not.toHaveBeenCalled();
});
it("olumsuz sonuç da cache'lenir — her browse'da yeniden denenmez", async () => {
const { svc, fetchMock } = makeService({ cached: "none" });
expect(await svc.resolveModelPathFromCatmeta("/p5x", "x_parts", {})).toBeNull();
expect(fetchMock).not.toHaveBeenCalled();
});
it("catmeta yoksa veya şekil beklenmedikse null döner", async () => {
const { svc, sets } = makeService({ catmeta: { data: {} } });
expect(await svc.resolveModelPathFromCatmeta("/p5x", "x_parts", {})).toBeNull();
expect(sets[0][1]).toBe("none");
});
it("HTTP hatasında null döner", async () => {
const { svc } = makeService({ status: 403, catmeta: volvoMeta });
expect(await svc.resolveModelPathFromCatmeta("/p5volvo", "volvo_parts", {})).toBeNull();
});
it("ağ hatası fırlatmaz", async () => {
const { svc } = makeService({ throws: true });
await expect(
svc.resolveModelPathFromCatmeta("/p5volvo", "volvo_parts", {}),
).resolves.toBeNull();
});
it("/extern/ ile başlamayan yolu kabul etmez", async () => {
const { svc } = makeService({
catmeta: { data: { catalogEntryPoint: { path: "https://baska.site/kotu" } } },
});
expect(await svc.resolveModelPathFromCatmeta("/p5volvo", "volvo_parts", {})).toBeNull();
});
});

View File

@@ -1,5 +1,10 @@
import { describe, expect, it } from "vitest"; import { describe, expect, it } from "vitest";
import { isPl24GroupNode, isPl24LeafNode } from "./pl24-tree"; import {
isPl24GroupNode,
isPl24LeafNode,
isPl24PartDetailNode,
isStalePl24Architecture,
} from "./pl24-tree";
// Canlı P5 yanıtlarından (2026-09-16 keşfi, plv2-artefakt/) alınan gerçek // Canlı P5 yanıtlarından (2026-09-16 keşfi, plv2-artefakt/) alınan gerçek
// wid + path çiftleri. Bu dosya "0 parça" sınıfı hatanın regresyon kilidi. // wid + path çiftleri. Bu dosya "0 parça" sınıfı hatanın regresyon kilidi.
@@ -110,3 +115,72 @@ describe("isPl24LeafNode — canlı P5 şekilleri", () => {
).toBe(true); ).toBe(true);
}); });
}); });
describe("isStalePl24Architecture — P4→P5 göç tespiti", () => {
it("eski PSA/Volvo yolu + artık P5 olan servis → bayat", () => {
expect(
isStalePl24Architecture({ catalogPath: "/psa/peugeot_parts", currentApiPath: "/p5psa" }),
).toBe(true);
expect(isStalePl24Architecture({ catalogPath: "/volvo", currentApiPath: "/p5volvo" })).toBe(
true,
);
});
it("zaten P5 kaydı bayat değil", () => {
expect(isStalePl24Architecture({ catalogPath: "/p5psa", currentApiPath: "/p5psa" })).toBe(
false,
);
});
it("hâlâ P4 olan markalar (Ford/Opel/Hyundai) dokunulmaz", () => {
expect(isStalePl24Architecture({ catalogPath: "/ford", currentApiPath: "/ford" })).toBe(false);
expect(isStalePl24Architecture({ catalogPath: "/opel", currentApiPath: "/opel" })).toBe(false);
});
it("eksik bilgi → bayat sayma (güvenli taraf)", () => {
expect(isStalePl24Architecture({ catalogPath: null, currentApiPath: "/p5psa" })).toBe(false);
expect(isStalePl24Architecture({ catalogPath: "/psa/x", currentApiPath: null })).toBe(false);
});
});
/**
* Mitsubishi: parça listesi grup sanılıyordu (plv2.md, bulgu consumers-13).
*
* `/p5mitsubishi/extern/details/vinDetails` yanıtı `partno`/`qty` taşıyan bir
* PARÇA listesi (canlı doğrulama 2026-09-20: 16 kayıt), ama her kaydın linki
* `partInfoTable`. Ne wid ne yol sınıflandırıcıda karşılık bulmuyordu → 2.193
* parça listesi gruba, içlerindeki 19.576 tekil parça kategoriye dönüşmüştü;
* hepsinde toplam 2 parça vardı ve her prefetch turunda yeniden çekiliyorlardı.
*/
describe("Mitsubishi detailsTable / partInfoTable", () => {
const detailsPath =
"/p5mitsubishi/extern/details/vinDetails?bomDetails=133_110D00125Y&mainGroup=33";
const partInfoPath = "/p5mitsubishi/extern/details/vinpartinfo?bomDetails=142_7103K22Y5T";
it("detailsTable bir parça listesi — yaprak", () => {
expect(isPl24LeafNode({ linkWid: "detailsTable", linkPath: detailsPath })).toBe(true);
expect(isPl24GroupNode({ linkWid: "detailsTable", linkPath: detailsPath })).toBe(false);
});
it("partInfoTable tekil parça detayı — ne grup ne liste", () => {
expect(isPl24PartDetailNode({ linkWid: "partInfoTable", linkPath: partInfoPath })).toBe(true);
expect(isPl24GroupNode({ linkWid: "partInfoTable", linkPath: partInfoPath })).toBe(false);
});
it("wid yoksa yol kalıbı parça-detayını yine yakalar", () => {
expect(isPl24PartDetailNode({ linkPath: partInfoPath })).toBe(true);
});
it("gerçek parça uçlarını parça-detayı sanmaz", () => {
// Volvo'nun partinfo yaprağı ve VW'nin bom listesi etkilenmemeli.
expect(
isPl24PartDetailNode({
linkWid: "partinfo",
linkPath: "/p5volvo/extern/partinfo/vin?partno=1",
}),
).toBe(false);
expect(
isPl24PartDetailNode({ linkWid: "bomlist", linkPath: "/p5vwag/extern/bom/vin?x=1" }),
).toBe(false);
});
});

View File

@@ -22,8 +22,44 @@
* drill. Path matching stays as a fallback for stored rows without a wid. * drill. Path matching stays as a fallback for stored rows without a wid.
*/ */
/** `link.wid` values that identify a parts (BOM) node across every P5 backend. */ /**
const LEAF_WIDS = new Set(["bomlist", "bomoverviewlist", "servicepartsitemstable", "partinfo"]); * `link.wid` values that identify a parts (BOM) node across every P5 backend.
*
* `detailstable` is Mitsubishi's: `/p5mitsubishi/extern/details/vinDetails`
* answers with 16-ish records carrying `partno`/`qty`, i.e. it IS the parts
* list — but each record's own link is a `partInfoTable` per-part detail, and
* neither the wid nor the path matched anything here, so the whole list was
* drilled as a group. Prod on 2026-09-20: 2,193 Mitsubishi parts lists turned
* into group nodes and their 19,576 individual parts ("SCREW,LOCK CYLINDER",
* "BOLT,STEERING COLUMN WASHER") became categories — 19,576 fake tree nodes
* with 2 parts between them, each re-fetched on every prefetch pass.
*/
const LEAF_WIDS = new Set([
"bomlist",
"bomoverviewlist",
"servicepartsitemstable",
"partinfo",
"detailstable",
]);
/**
* Nodes that describe ONE part rather than a list of them. They are neither a
* group to drill nor a list to fetch: the parent's own response already carried
* the part. Queueing them buys nothing and costs one upstream request each —
* 19,576 of them on prod before this was recognised.
*/
const PART_DETAIL_WIDS = new Set(["partinfotable"]);
const PART_DETAIL_PATH = /\/details\/vinpartinfo\b/i;
/** True when this node is a single part's detail view, not a listing. */
export function isPl24PartDetailNode(opts: {
linkPath?: string | null;
linkWid?: string | null;
}): boolean {
const wid = opts.linkWid?.toLowerCase().trim();
if (wid && PART_DETAIL_WIDS.has(wid)) return true;
return PART_DETAIL_PATH.test(opts.linkPath ?? "");
}
/** `link.wid` values that identify a drillable group node. */ /** `link.wid` values that identify a drillable group node. */
const GROUP_WID_PATTERN = const GROUP_WID_PATTERN =
@@ -61,5 +97,55 @@ export function isPl24GroupNode(opts: {
hasSubgroups?: boolean | null; hasSubgroups?: boolean | null;
}): boolean { }): boolean {
if (!opts.linkPath && !opts.linkWid) return false; if (!opts.linkPath && !opts.linkWid) return false;
// A per-part detail node is not a group; drilling it returns nothing.
if (isPl24PartDetailNode(opts)) return false;
return !isPl24LeafNode(opts); return !isPl24LeafNode(opts);
} }
/**
* True when a stored vehicle's catalogInfo points at an architecture the service
* no longer uses — i.e. the row was decoded before PSA/Volvo moved to P5.
*
* These vehicles keep serving a tree built from the frozen P4 PSA snapshot
* (upds 2024-02-13) or the dead P4 Volvo endpoint (HTTP 503), and `decodeVin`'s
* db_hit short-circuit means a code fix never reaches them. Detecting the
* mismatch at read time lets each vehicle heal itself on first view instead of
* needing a bulk re-decode storm against the one surviving account.
*/
export function isStalePl24Architecture(opts: {
catalogPath?: string | null;
currentApiPath?: string | null;
}): boolean {
const stored = opts.catalogPath?.toLowerCase() ?? "";
const current = opts.currentApiPath?.toLowerCase() ?? "";
if (!stored || !current) return false;
// Only the two migrated legacy backends; unknown/other paths are left alone.
const storedIsLegacyPsaOrVolvo = stored.startsWith("/psa") || stored.startsWith("/volvo");
if (!storedIsLegacyPsaOrVolvo) return false;
return current.startsWith("/p5");
}
/**
* True when a stored `catalog_vehicles` browse row was listed under an
* architecture the service table no longer uses.
*
* Browse rows are a PERMANENT cache: `CatalogService.getModels` returns early
* whenever the brand already has rows, so `fetchVehicleList` is never called
* again for that brand. The 188 rows listed while PSA and Volvo were still P4
* (127 LEGACY_PSA + 61 LEGACY_VOLVO on prod, 2026-09-20) are therefore pinned
* forever: Peugeot/Citroën browse serves the frozen 2024-02-13 snapshot and
* Volvo/Polestar browse serves an endpoint that answers HTTP 503. Detecting the
* mismatch at list time lets the brand re-list itself once, the same way
* `isStalePl24Architecture` heals a VIN-decoded vehicle.
*/
export function isStaleBrowseArchitecture(opts: {
storedArchitecture?: string | null;
currentArchitecture?: string | null;
}): boolean {
const stored = opts.storedArchitecture?.trim();
const current = opts.currentArchitecture?.trim();
// An unknown service (no config) or an unlabelled row is left alone: without a
// current architecture to compare against there is nothing to migrate TO.
if (!stored || !current) return false;
return stored !== current;
}

View File

@@ -0,0 +1,68 @@
import { readFileSync } from "node:fs";
import { join } from "node:path";
import { describe, expect, it } from "vitest";
import { PL24Service } from "./pl24.service";
/**
* Volvo P5 vinfoBasic regression lock (plv2.md, finding p5core-09).
*
* Volvo emits BOTH a bare internal code and a readable description for the same
* attribute — "Şanzıman kodu" = "B" next to "Şanzıman" = "6-PSHIFT 2WD / MPS6",
* and "Satis tipi" = "42" next to "Türü" = "S80". The key order the parser needs
* for VAG (whose "Şanzıman kodu" IS the useful value, e.g. "MQ200") therefore
* rendered a single letter as the gearbox and a bare number as the series.
*
* Fixtures are the real 2026-09-16 discovery captures for VIN
* YV1AS84ABD1168166 (S80), trimmed to the segments the parser reads, in both
* the Turkish and the English label locale.
*/
const fixture = (name: string) =>
JSON.parse(readFileSync(join(__dirname, "__fixtures__", `${name}.json`), "utf-8"));
// parseVehicleResponse only touches the response and the service name.
const parse = (data: unknown, serviceName = "volvo_parts") => {
const svc = Object.create(PL24Service.prototype) as unknown as {
parseVehicleResponse: (v: string, d: unknown, s: string) => Record<string, unknown>;
isDaimlerService: (s: string) => boolean;
};
svc.isDaimlerService = () => false;
return svc.parseVehicleResponse("YV1AS84ABD1168166", data, serviceName);
};
describe("Volvo P5 vinfoBasic — kod yerine açıklama", () => {
it("Türkçe etiketlerde şanzımanı okunur değerinden alır", () => {
const out = parse(fixture("p5_volvo_tr"));
expect(out.transmission).toBe("6-PSHIFT 2WD / MPS6");
// Tek harflik "Şanzıman kodu" artık kazanmıyor.
expect(out.transmission).not.toBe("B");
});
it("İngilizce etiketlerde de aynı sonucu verir", () => {
const out = parse(fixture("p5_volvo_en"));
expect(out.transmission).toBe("6-PSHIFT 2WD / MPS6");
});
it("seri, çıplak satış kodu yerine gerçek tipi döner", () => {
expect(parse(fixture("p5_volvo_tr")).series).toBe("S80");
expect(parse(fixture("p5_volvo_en")).series).toBe("S80");
expect(parse(fixture("p5_volvo_tr")).series).not.toBe("42");
});
it("kaporta stili iki dilde de çözülür ve '0' kodu sızmaz", () => {
expect(parse(fixture("p5_volvo_tr")).bodyType).toBe("Sedan");
expect(parse(fixture("p5_volvo_en")).bodyType).toBe("Sedan");
});
it("motor alanları bozulmadan kalır", () => {
const out = parse(fixture("p5_volvo_tr"));
expect(out.engineCode).toBe("84");
expect(out.engineType).toBe("D4162T");
});
it("model ve yıl korunur", () => {
const out = parse(fixture("p5_volvo_tr"));
expect(out.model).toBe("S80 (07-)");
expect(out.year).toBe(2013);
});
});

View File

@@ -66,13 +66,20 @@ const LEGACY_ARCH_SOURCE_TAG: Record<string, string> = {
* single spelling matches both "Model yılı" and "MODEL YILI". * single spelling matches both "Model yılı" and "MODEL YILI".
*/ */
export function normalizeLabel(label: string): string { export function normalizeLabel(label: string): string {
return label return (
.toLocaleLowerCase("tr") label
.normalize("NFD") .toLocaleLowerCase("tr")
.replace(/[\u0300-\u036f]/g, "") .normalize("NFD")
.replace(/ı/g, "i") // \p{M} (all combining marks) rather than the U+0300–U+036F range: the range
.replace(/[\s/]+/g, "_") // is a character class that can also match a base character followed by a
.trim(); // combining one, which biome flags as misleading. NFD has already split every
// accent into its own mark, so matching marks directly is both correct and
// broader (Turkish, Latin-Extended, anything PL24 sends).
.replace(/\p{M}/gu, "")
.replace(/ı/g, "i")
.replace(/[\s/]+/g, "_")
.trim()
);
} }
/** /**
@@ -1137,6 +1144,78 @@ export class PL24Service {
// ==================== PRIVATE: Response parsers ==================== // ==================== PRIVATE: Response parsers ====================
/**
* Ask a P5 backend where its own model list lives, instead of guessing.
*
* Every P5 backend serves `/extern/catmeta`, whose `data.catalogEntryPoint.path`
* is the authoritative browse entry point — e.g. p5psa answers
* `/p5psa/extern/vehicle/catalogs`, p5volvo `/p5volvo/extern/vehicles/models`.
* The old behaviour, when `BACKEND_MODEL_PATH` had no entry, was to fire the
* seven known paths in turn and keep whichever returned rows. That costs up to
* seven upstream requests per call on a shared daily budget, and it silently
* fails for any backend whose path is not already on the list — which is how
* Volvo/Polestar browse ended up seeding zero models while still serving the
* retired P4 listing.
*
* The resolved path is cached per backend (30 days) so this costs one request
* for a backend's whole lifetime, and it self-heals if PL24 moves an endpoint.
* Returns null on any failure; the caller then falls back as before.
*/
private async resolveModelPathFromCatmeta(
catalogBase: string,
serviceName: string,
headers: Record<string, string>,
): Promise<string | null> {
const cacheKey = `pl24:modelpath:${catalogBase.replace(/^\//, "")}`;
try {
const cached = await this.redis.get(cacheKey);
if (cached) return cached === "none" ? null : cached;
} catch {
// Redis down — resolve live rather than failing the listing.
}
let resolved: string | null = null;
try {
const url = `${this.baseUrl}${catalogBase}/extern/catmeta?serviceName=${serviceName}&country=DE&lang=${this.language}`;
const res = await fetch(url, { method: "GET", headers, signal: AbortSignal.timeout(15000) });
if (res.ok) {
const meta = (await res.json()) as {
data?: { catalogEntryPoint?: { path?: string } };
};
const full = meta.data?.catalogEntryPoint?.path;
if (full) {
// catmeta returns the absolute path with query string
// ("/p5volvo/extern/vehicles/models?lang=en&serviceName=…"); the caller
// appends its own lang/serviceName, so keep only the backend-relative
// path segment.
const withoutQuery = full.split("?")[0];
const relative = withoutQuery.startsWith(catalogBase)
? withoutQuery.slice(catalogBase.length)
: withoutQuery;
if (relative.startsWith("/extern/")) resolved = relative;
}
}
} catch (err) {
this.logger.warn(
`fetchVehicleList: catmeta lookup failed for ${serviceName}: ${(err as Error).message}`,
);
}
if (resolved) {
this.logger.log(
`fetchVehicleList: ${serviceName} model path resolved from catmeta → ${resolved}`,
);
}
try {
// Cache the negative too, so a backend without a usable catmeta does not
// re-probe on every browse.
await this.redis.set(cacheKey, resolved ?? "none", 30 * 86_400);
} catch {
// best effort
}
return resolved;
}
/** /**
* Parse vehicle data from directAccess response. * Parse vehicle data from directAccess response.
*/ */
@@ -1184,6 +1263,31 @@ export class PL24Service {
return null; return null;
}; };
/**
* Like `lookup`, but prefers a human-readable value over a bare internal
* code when the record carries both.
*
* Volvo's vinfoBasic has BOTH "Şanzıman kodu" ("B") and "Şanzıman"
* ("6-PSHIFT 2WD / MPS6"), and BOTH "Satis tipi" ("42") and a second
* "Satis tipi" ("SALES VERSION 42"). The code-first key order that VAG needs
* (its "Şanzıman kodu" IS the useful value, e.g. "MQ200") therefore rendered
* a single letter as the Volvo gearbox and a bare number as its series.
*
* Falls back to the first present value, so a backend that only ever emits
* codes behaves exactly as before.
*/
const lookupDescriptive = (...keys: string[]): string | null => {
let firstPresent: string | null = null;
for (const k of keys) {
const v = vehicleData[k]?.trim();
if (!v) continue;
if (firstPresent === null) firstPresent = v;
// Two characters or fewer, or digits only → an index, not a description.
if (v.length > 2 && !/^\d+$/.test(v)) return v;
}
return firstPresent;
};
// Extract prNr records for richer vehicle attributes // Extract prNr records for richer vehicle attributes
const prNrRecords = segments.prNr?.records || []; const prNrRecords = segments.prNr?.records || [];
const prNrByCode: Record<string, string> = {}; const prNrByCode: Record<string, string> = {};
@@ -1226,13 +1330,17 @@ export class PL24Service {
// normalizeLabel folds ı→i and strips diacritics, so "Şanzıman kodu" and // normalizeLabel folds ı→i and strips diacritics, so "Şanzıman kodu" and
// "ŞANZIMAN KODU" both arrive as "sanziman_kodu". PSA uses "AKTARMA // "ŞANZIMAN KODU" both arrive as "sanziman_kodu". PSA uses "AKTARMA
// SİSTEMLERİ" ("5 MEKANİK VİTES KUTUSU"), Subaru "Mission". // SİSTEMLERİ" ("5 MEKANİK VİTES KUTUSU"), Subaru "Mission".
const transmissionCode = lookup( const transmissionCode = lookupDescriptive(
"sanziman_kodu", "sanziman_kodu",
"transmission_code", "transmission_code",
"vites_kutusu", "vites_kutusu",
"atm,mtm", "atm,mtm",
"aktarma_sistemleri", "aktarma_sistemleri",
"sanziman", "sanziman",
// Volvo in the English locale: "Transmission code" ("B") plus the readable
// "Transmission" ("6-PSHIFT 2WD / MPS6"). Without the plain key the
// English response falls back to the single-letter code.
"transmission",
"sanziman_numarasi", "sanziman_numarasi",
"mission", "mission",
); );
@@ -1242,7 +1350,10 @@ export class PL24Service {
const bodyType = const bodyType =
Object.entries(prNrByCode).find(([code]) => code.startsWith("K8"))?.[1] || Object.entries(prNrByCode).find(([code]) => code.startsWith("K8"))?.[1] ||
// PSA "GÖVDE TİPİ" ("4 KAPILI SEDAN"), Volvo "Kaporta Stili" ("Sedan"). // PSA "GÖVDE TİPİ" ("4 KAPILI SEDAN"), Volvo "Kaporta Stili" ("Sedan").
lookup("karoseri", "body", "body_type", "govde_tipi", "kaporta_stili") || // Volvo: TR "Kaporta Stili" / EN "Body style" ("Sedan"). Its
// "Karoseri Tipi Kodu" / "Body style code" is a bare "0" — never matched
// here because the lookup is exact-key, and that is deliberate.
lookup("karoseri", "body", "body_type", "govde_tipi", "kaporta_stili", "body_style") ||
null; null;
// Engine description from prNr D3* (Motor nitelikleri) // Engine description from prNr D3* (Motor nitelikleri)
@@ -1301,7 +1412,10 @@ export class PL24Service {
damToModelYear(lookup("dam")) || damToModelYear(lookup("dam")) ||
extractModelYear(vin) || extractModelYear(vin) ||
0, 0,
series: lookup("seri", "satis_tipi", "sales_type", "turu"), // Volvo's "Satis tipi" is the bare sales code ("42") with the readable
// form ("SALES VERSION 42") in a duplicate row, and its "Türü" is the real
// series ("S80") — so prefer whichever of these is actually descriptive.
series: lookupDescriptive("seri", "satis_tipi", "sales_type", "turu", "type"),
bodyType, bodyType,
engineCode: engineCode:
engineCode || engineCode ||
@@ -2191,11 +2305,22 @@ export class PL24Service {
p5mitsubishi: "/extern/vehicles/vehiclesOverview", // Mitsubishi p5mitsubishi: "/extern/vehicles/vehiclesOverview", // Mitsubishi
p5suzuki: "/extern/vehicle/modelFamilies", // Suzuki p5suzuki: "/extern/vehicle/modelFamilies", // Suzuki
p5man: "/extern/model/categories", // MAN trucks p5man: "/extern/model/categories", // MAN trucks
// Pinned 2026-09-20 from each backend's own catmeta `catalogEntryPoint`
// (see resolveModelPathFromCatmeta). Before this, p5psa needed five wasted
// probe requests to rediscover its path on every call, and p5volvo found
// nothing at all — none of the seven guessed paths matches its plural
// `/extern/vehicles/models`, so Volvo/Polestar browse silently seeded zero
// models and kept serving the dead P4 listing.
p5psa: "/extern/vehicle/catalogs", // Peugeot, Citroën, DS, psa_opel, psa_vauxhall
p5volvo: "/extern/vehicles/models", // Volvo, Polestar — NOTE: plural "vehicles"
}; };
// catalogBase is like "/p5vwag" — strip leading slash for map lookup // catalogBase is like "/p5vwag" — strip leading slash for map lookup
const backendKey = catalogBase.replace(/^\//, ""); const backendKey = catalogBase.replace(/^\//, "");
const modelPath = BACKEND_MODEL_PATH[backendKey] ?? "/extern/vehicle/modelfamilies"; const modelPath =
BACKEND_MODEL_PATH[backendKey] ??
(await this.resolveModelPathFromCatmeta(catalogBase, serviceName, headers)) ??
"/extern/vehicle/modelfamilies";
const url = `${this.baseUrl}${catalogBase}${modelPath}?lang=${this.language}&serviceName=${serviceName}`; const url = `${this.baseUrl}${catalogBase}${modelPath}?lang=${this.language}&serviceName=${serviceName}`;

View File

@@ -1,5 +1,5 @@
import { describe, expect, it } from "vitest"; import { describe, expect, it } from "vitest";
import { PL24_WMI_SERVICE_MAP, SERVICE_TO_BRAND } from "./pl24.types"; import { PL24_SERVICE_CATALOGS, PL24_WMI_SERVICE_MAP, SERVICE_TO_BRAND } from "./pl24.types";
// Q6 routing additions (undecoded-vin-rca.md). Both target services are already // Q6 routing additions (undecoded-vin-rca.md). Both target services are already
// live in prod (nissan_parts: JN1 success; mercedesvans_parts: WDF decodes today), // live in prod (nissan_parts: JN1 success; mercedesvans_parts: WDF decodes today),
@@ -27,3 +27,71 @@ describe("PL24_WMI_SERVICE_MAP — Q6 routing additions", () => {
} }
}); });
}); });
// PL24'te bulunmayan WMI'ler: haritada yer almamalı, yoksa her sorguda boşuna
// upstream isteği üretir ve gerçek kaynağa (pcat/emex/vinpin) geçişi geciktirir.
describe("PL24_WMI_SERVICE_MAP — kapsam dışı WMI'ler", () => {
it("VR7 haritada değil (canlı: HTTP 410 'no brands found for WMI = VR7')", () => {
expect(PL24_WMI_SERVICE_MAP.VR7).toBeUndefined();
});
it("NM4 (Tofaş) haritada değil (canlı: HTTP 410)", () => {
expect(PL24_WMI_SERVICE_MAP.NM4).toBeUndefined();
});
it("canlı doğrulanmış WMI'ler doğru katalogda", () => {
expect(PL24_WMI_SERVICE_MAP.JF1).toBe("subaru_parts");
expect(PL24_WMI_SERVICE_MAP.ZAR).toBe("alfa_parts");
expect(PL24_WMI_SERVICE_MAP.VXK).toBe("psa_opel_parts");
});
});
/**
* 2026-09-20 canlı WMI servisi taraması (plv2.md, bulgu types_wmi-06/09).
* Prod'daki 1.406 haritasız araçtan hangilerinin PL24'te gerçekten karşılığı
* olduğu `/pl24-wmi/ext/api/2.0/decode` ile tek tek soruldu; bu test o cevapları
* kilitliyor — özellikle "yok" cevaplarını, çünkü onları haritaya eklemek
* boşuna istek üretir.
*/
describe("WMI haritası — canlı WMI servisi taraması (2026-09-20)", () => {
it("Renault yeniden açıldı (VIN tanımlama askısı kalktı)", () => {
expect(PL24_WMI_SERVICE_MAP.VF1).toBe("renault_parts");
expect(PL24_WMI_SERVICE_MAP.VF6).toBe("renault_parts");
expect(PL24_WMI_SERVICE_MAP.VNE).toBe("renault_parts");
});
it("canlı servisin çözdüğü yeni WMI'ler haritada", () => {
expect(PL24_WMI_SERVICE_MAP.NLH).toBe("hyundai_parts");
expect(PL24_WMI_SERVICE_MAP.TMA).toBe("hyundai_parts");
expect(PL24_WMI_SERVICE_MAP.NLJ).toBe("hyundai_parts");
expect(PL24_WMI_SERVICE_MAP.KMF).toBe("hyundai_parts");
expect(PL24_WMI_SERVICE_MAP.KNE).toBe("kia_parts");
expect(PL24_WMI_SERVICE_MAP.KNC).toBe("kia_parts");
expect(PL24_WMI_SERVICE_MAP.MMC).toBe("mmc_parts");
expect(PL24_WMI_SERVICE_MAP.XMC).toBe("mmc_parts");
expect(PL24_WMI_SERVICE_MAP.JSA).toBe("suzuki_parts");
});
it("NMB binek Mercedes değil, kamyon kataloğuna gider", () => {
expect(PL24_WMI_SERVICE_MAP.NMB).toBe("mercedestrucks_parts");
// Binek WMI'leri bozulmadı.
expect(PL24_WMI_SERVICE_MAP.WDD).toBe("mercedes_parts");
});
it("PL24'te olmayan WMI'ler haritaya EKLENMEDİ (HTTP 410)", () => {
// Honda ve Chevrolet'nin PL24'te kataloğu yok; eklemek boşuna istek olurdu.
for (const wmi of ["JHM", "SHH", "SHS", "MAK", "NLA", "KL1", "NM4", "VR7"]) {
expect(PL24_WMI_SERVICE_MAP[wmi]).toBeUndefined();
}
});
it("JMZ (Mazda) eklenmedi — servis Ford kataloğu öneriyor, marka tutarsız", () => {
expect(PL24_WMI_SERVICE_MAP.JMZ).toBeUndefined();
});
it("eşlenen her servisin katalog tanımı var", () => {
for (const [wmi, svc] of Object.entries(PL24_WMI_SERVICE_MAP)) {
expect(PL24_SERVICE_CATALOGS[svc], `${wmi} → ${svc}`).toBeDefined();
}
});
});

View File

@@ -540,6 +540,10 @@ export const PL24_WMI_SERVICE_MAP: Record<string, string> = {
// Mercedes-Benz // Mercedes-Benz
WDB: "mercedes_parts", WDB: "mercedes_parts",
// NMB: canlı WMI servisi bunu mercedes_parts'a DEĞİL mercedestrucks_parts'a
// çözüyor (error:false) — 30 prod aracı "Mercedes-Benz" etiketliydi ama binek
// kataloğunda yok.
NMB: "mercedestrucks_parts",
WDD: "mercedes_parts", WDD: "mercedes_parts",
WDC: "mercedes_parts", WDC: "mercedes_parts",
W1K: "mercedes_parts", W1K: "mercedes_parts",
@@ -578,15 +582,21 @@ export const PL24_WMI_SERVICE_MAP: Record<string, string> = {
JTJ: "lexus_parts", JTJ: "lexus_parts",
"2T2": "lexus_parts", "2T2": "lexus_parts",
// Renault — DISABLED: PL24 has suspended Renault VIN identification ("Bu marka // Renault — RE-ENABLED 2026-09-20. It was disabled while PL24 had suspended
// için şasi numarası tanımlamasının belirsiz bir süre için mevcut olmayacağını // Renault VIN identification ("Bu marka için şasi numarası tanımlamasının
// üzülerek bildiririz."). renault_parts authorizes fine but every decode throws // belirsiz bir süre için mevcut olmayacağını üzülerek bildiririz."), which both
// that message → wasted ~1s call AND it trips the PL24 circuit breaker, which // wasted a call per decode and, back then, tripped the global PL24 breaker.
// then skips PL24 for ALL brands. PCAT + EMEX cover Renault. Re-enable when PL24 // BOTH reasons are gone, each verified against prod rather than assumed:
// restores Renault VIN decode. // 1. The suspension is over — live directAccess on /p5renault for a real
// VF1: "renault_parts", // customer VIN (VF14SRCL458170337) returns resultStatus
// VF6: "renault_parts", // VEHICLE_IDENTIFIED, "SYMBOL II/LOGAN II", and the WMI service answers
// VNE: "renault_parts", // {service: renault_parts, valid: true, error: false}.
// 2. The breaker no longer counts definitive upstream negatives, only
// transient transport faults (see vehicles.service `isTransient`).
// 450 prod vehicles carry VF1 and had no PL24 catalog at all.
VF1: "renault_parts",
VF6: "renault_parts",
VNE: "renault_parts",
// Dacia // Dacia
UU1: "dacia_parts", UU1: "dacia_parts",
@@ -606,6 +616,10 @@ export const PL24_WMI_SERVICE_MAP: Record<string, string> = {
WMH: "man_parts", WMH: "man_parts",
// Mitsubishi // Mitsubishi
// XMC / MMC: canlı WMI servisi ikisini de mmc_parts'a çözüyor ve P5 doğruluyor
// (error:false) — 26 prod aracı haritasızdı.
MMC: "mmc_parts", // Mitsubishi Japan
XMC: "mmc_parts", // Mitsubishi (diğer pazarlar)
JMB: "mmc_parts", JMB: "mmc_parts",
JMY: "mmc_parts", JMY: "mmc_parts",
MMB: "mmc_parts", MMB: "mmc_parts",
@@ -615,6 +629,7 @@ export const PL24_WMI_SERVICE_MAP: Record<string, string> = {
JS2: "suzuki_parts", JS2: "suzuki_parts",
JS3: "suzuki_parts", JS3: "suzuki_parts",
TSM: "suzuki_parts", TSM: "suzuki_parts",
JSA: "suzuki_parts", // Suzuki (canlı WMI: suzuki_parts, error:false)
MA3: "suzuki_parts", MA3: "suzuki_parts",
MBH: "suzuki_parts", MBH: "suzuki_parts",
@@ -638,10 +653,19 @@ export const PL24_WMI_SERVICE_MAP: Record<string, string> = {
// Hyundai // Hyundai
KMH: "hyundai_parts", // Hyundai Korea Motor House KMH: "hyundai_parts", // Hyundai Korea Motor House
// Aşağıdakiler 2026-09-20'de canlı WMI servisinden alındı; hepsi hyundai_parts.
// TMA iki marka döndürüyor (hyundai_parts + kia_parts) — Hyundai Assan (Türkiye)
// üretimi olduğu için hyundai seçildi.
NLH: "hyundai_parts", // Hyundai Assan / Türkiye (106 prod aracı)
TMA: "hyundai_parts", // Hyundai Türkiye (49)
NLJ: "hyundai_parts", // Hyundai (8)
KMF: "hyundai_parts", // Hyundai (7)
TMK: "hyundai_parts", // Hyundai (Turkey/other markets) TMK: "hyundai_parts", // Hyundai (Turkey/other markets)
// Kia // Kia
KNA: "kia_parts", // Kia (worldwide production) KNA: "kia_parts", // Kia (worldwide production)
KNE: "kia_parts", // Kia (canlı WMI, 31 prod aracı)
KNC: "kia_parts", // Kia (canlı WMI, 8)
U5Y: "kia_parts", // Kia Slovakia U5Y: "kia_parts", // Kia Slovakia
// Nissan // Nissan
@@ -675,7 +699,19 @@ export const PL24_WMI_SERVICE_MAP: Record<string, string> = {
// Peugeot (PSA) // Peugeot (PSA)
VF3: "peugeot_parts", // Peugeot SA (France) VF3: "peugeot_parts", // Peugeot SA (France)
VR3: "peugeot_parts", // Peugeot (newer WMI) VR3: "peugeot_parts", // Peugeot (newer WMI)
VR7: "peugeot_parts", // Peugeot (newer WMI) // VR7: PL24'ün WMI veritabanında HİÇ YOK — canlı doğrulama (2026-09-19,
// /pl24-wmi/ext/api/2.0/decode): HTTP 410 "no brands found for WMI = VR7",
// tıpkı Tofaş NM4 gibi. Önceden peugeot_parts'a yönlendiriliyordu; prod'daki
// 17 VR7 aracının hepsi model çözülmeden ("Peugeot Peugeot") ve 0 kategoriyle
// kaydedilmişti. Haritada tutmak yalnız boşuna PL24 isteği üretir ve
// pcat/emex/vinpin fallback'ini geciktirir.
// PL24'TE HİÇ OLMAYANLAR — canlı WMI servisi HTTP 410 "no brands found for WMI"
// dedi (2026-09-20), tıpkı NM4 ve VR7 gibi. Haritaya eklenmemeleri kasıtlı:
// eklemek yalnız boşuna istek üretir ve pcat/emex/vinpin fallback'ini geciktirir.
// JHM, SHH, SHS, MAK, NLA (Honda) · KL1 (Chevrolet)
// JMZ (Mazda) 410 DEĞİL ama fordp/fordt döndürüyor (Ford-Mazda platform ortaklığı
// dönemi); marka tutarsız olduğu için eklenmedi — bir Mazda 3 Ford kataloğunda yok.
// Volvo // Volvo
YV1: "volvo_parts", // Volvo Cars (Sweden) YV1: "volvo_parts", // Volvo Cars (Sweden)

View File

@@ -0,0 +1,148 @@
import { afterEach, describe, expect, it, vi } from "vitest";
import { RpartstoreAuthError, decodeJwtExpiry, loginRpartstore } from "./rpartstore.auth";
const b64url = (s: string): string => Buffer.from(s).toString("base64url");
const makeJwt = (claims: Record<string, unknown>): string =>
`${b64url('{"alg":"RS256"}')}.${b64url(JSON.stringify(claims))}.sig`;
function jsonResponse(
body: unknown,
init: { status?: number; headers?: Record<string, string> } = {},
) {
return new Response(JSON.stringify(body), {
status: init.status ?? 200,
headers: { "content-type": "application/json", ...(init.headers ?? {}) },
});
}
describe("loginRpartstore (Okta IDX, browser-free)", () => {
afterEach(() => vi.unstubAllGlobals());
it("walks authorize → introspect → identify → answer → redirect → token and returns the access token", async () => {
const exp = Math.floor(Date.now() / 1000) + 3600;
const jwt = makeJwt({ sub: "G123326", exp });
const calls: { url: string; body?: string; cookie?: string }[] = [];
let redirectState = "";
const fetchImpl = vi.fn(async (url: string | URL, init?: RequestInit) => {
const u = String(url);
const headers = (init?.headers ?? {}) as Record<string, string>;
calls.push({
url: u,
body: init?.body ? String(init.body) : undefined,
cookie: headers.Cookie,
});
if (u.includes("/v1/authorize")) {
redirectState = new URL(u).searchParams.get("state") ?? "";
return new Response("<script>var stateToken = '02.id.abc\\x2Ddef';</script>", {
status: 200,
headers: { "set-cookie": "JSESSIONID=js1; Path=/; HttpOnly" },
});
}
if (u.endsWith("/idp/idx/introspect")) return jsonResponse({ stateHandle: "sh-1" });
if (u.endsWith("/idp/idx/identify")) return jsonResponse({ stateHandle: "sh-2" });
if (u.endsWith("/idp/idx/challenge/answer")) {
return jsonResponse({
success: { href: "https://sso.renault.com/login/token/redirect?stateToken=st" },
});
}
if (u.includes("/login/token/redirect")) {
return new Response(null, {
status: 302,
headers: {
location: `https://rpartstore.renault.com/idp-redirect?code=CODE1&state=${redirectState}`,
},
});
}
if (u.endsWith("/v1/token")) {
return jsonResponse({
token_type: "Bearer",
expires_in: 3600,
access_token: jwt,
scope: "openid",
});
}
throw new Error(`unexpected url ${u}`);
});
const token = await loginRpartstore({
username: "G123326",
password: "pw",
fetchImpl: fetchImpl as unknown as typeof fetch,
});
expect(token.accessToken).toBe(jwt);
expect(token.subject).toBe("G123326");
expect(token.expiresAt).toBe(exp * 1000);
// stateToken unescaped (\x2D → "-") and carried into introspect
expect(calls[1].body).toBe(JSON.stringify({ stateToken: "02.id.abc-def" }));
// identify uses the introspect handle, answer the identify handle
expect(JSON.parse(calls[2].body!)).toEqual({ identifier: "G123326", stateHandle: "sh-1" });
expect(JSON.parse(calls[3].body!)).toEqual({
credentials: { passcode: "pw" },
stateHandle: "sh-2",
});
// cookies from the authorize page are replayed on later hops
expect(calls[3].cookie).toContain("JSESSIONID=js1");
// PKCE: the token exchange sends the code and a verifier, never the password
const tokenBody = new URLSearchParams(calls.at(-1)!.body);
expect(tokenBody.get("grant_type")).toBe("authorization_code");
expect(tokenBody.get("code")).toBe("CODE1");
expect(tokenBody.get("code_verifier")).toBeTruthy();
expect(calls.at(-1)!.body).not.toContain("pw");
});
it("fails with a challenge error when Okta returns no success redirect (bad password / MFA)", async () => {
const fetchImpl = vi.fn(async (url: string | URL) => {
const u = String(url);
if (u.includes("/v1/authorize")) return new Response("stateToken = 'x'", { status: 200 });
if (u.endsWith("/introspect")) return jsonResponse({ stateHandle: "sh" });
if (u.endsWith("/identify")) return jsonResponse({ stateHandle: "sh" });
if (u.endsWith("/challenge/answer")) {
return jsonResponse({ messages: { value: [{ message: "Authentication failed" }] } });
}
throw new Error(`unexpected url ${u}`);
});
await expect(
loginRpartstore({
username: "u",
password: "bad",
fetchImpl: fetchImpl as unknown as typeof fetch,
}),
).rejects.toMatchObject({ name: "RpartstoreAuthError", step: "challenge" });
});
it("rejects an OAuth state mismatch on the callback", async () => {
const fetchImpl = vi.fn(async (url: string | URL) => {
const u = String(url);
if (u.includes("/v1/authorize")) return new Response("stateToken = 'x'", { status: 200 });
if (u.endsWith("/introspect")) return jsonResponse({ stateHandle: "sh" });
if (u.endsWith("/identify")) return jsonResponse({ stateHandle: "sh" });
if (u.endsWith("/challenge/answer"))
return jsonResponse({ success: { href: "https://sso.renault.com/r" } });
if (u === "https://sso.renault.com/r") {
return new Response(null, {
status: 302,
headers: { location: "https://rpartstore.renault.com/idp-redirect?code=C&state=forged" },
});
}
throw new Error(`unexpected url ${u}`);
});
await expect(
loginRpartstore({
username: "u",
password: "p",
fetchImpl: fetchImpl as unknown as typeof fetch,
}),
).rejects.toBeInstanceOf(RpartstoreAuthError);
});
it("decodes exp/sub from the access token", () => {
expect(decodeJwtExpiry(makeJwt({ sub: "G1", exp: 1790334740 }))).toEqual({
exp: 1790334740,
sub: "G1",
});
expect(() => decodeJwtExpiry(makeJwt({ sub: "G1" }))).toThrow(RpartstoreAuthError);
});
});

View File

@@ -0,0 +1,238 @@
/**
* Browser-free Okta (OIE / IDX) login for rpartstore.renault.com.
*
* Flow (verified live 2026-09-25, ~1.8 s, no MFA / captcha):
* 1. GET {issuer}/v1/authorize?...PKCE... → hosted page HTML containing `stateToken`
* 2. POST /idp/idx/introspect {stateToken} → stateHandle
* 3. POST /idp/idx/identify {identifier, stateHandle}
* 4. POST /idp/idx/challenge/answer {credentials:{passcode}, stateHandle} → success.href
* 5. GET success.href (follow redirects manually) → …/idp-redirect?code=…&state=…
* 6. POST {issuer}/v1/token (authorization_code + code_verifier) → access_token (1 h, no refresh_token)
*/
import { createHash, randomBytes } from "node:crypto";
export interface RpartstoreAuthConfig {
username: string;
password: string;
issuer?: string;
clientId?: string;
redirectUri?: string;
scope?: string;
/** Injectable for tests. */
fetchImpl?: typeof fetch;
}
export interface RpartstoreToken {
accessToken: string;
/** Epoch ms. */
expiresAt: number;
subject: string;
}
export class RpartstoreAuthError extends Error {
constructor(
message: string,
readonly step: string,
readonly detail?: unknown,
) {
super(message);
this.name = "RpartstoreAuthError";
}
}
const DEFAULTS = {
issuer: "https://sso.renault.com/oauth2/aus133y6mks4ptDss417",
clientId: "irn-72795_ope_pkce_4hcafvxlbcil",
redirectUri: "https://rpartstore.renault.com/idp-redirect",
scope: "openid alliance_profile apis.default",
};
const USER_AGENT =
"Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/148.0.0.0 Safari/537.36";
const ION = "application/ion+json; okta-version=1.0.0";
const b64url = (buf: Buffer): string => buf.toString("base64url");
/** Tiny cookie jar: Okta needs `idx`/`JSESSIONID` cookies across the IDX steps. */
class CookieJar {
private readonly jar = new Map<string, string>();
absorb(res: Response): void {
const setCookies: string[] =
typeof (res.headers as { getSetCookie?: () => string[] }).getSetCookie === "function"
? (res.headers as unknown as { getSetCookie: () => string[] }).getSetCookie()
: [];
for (const sc of setCookies) {
const kv = sc.split(";")[0];
const i = kv.indexOf("=");
if (i > 0) this.jar.set(kv.slice(0, i).trim(), kv.slice(i + 1).trim());
}
}
header(): string {
return Array.from(this.jar.entries())
.map(([k, v]) => `${k}=${v}`)
.join("; ");
}
}
export function decodeJwtExpiry(accessToken: string): { exp: number; sub: string } {
const payload = JSON.parse(
Buffer.from(accessToken.split(".")[1], "base64url").toString("utf8"),
) as {
exp?: number;
sub?: string;
};
if (!payload.exp) throw new RpartstoreAuthError("access token has no exp claim", "token");
return { exp: payload.exp, sub: payload.sub ?? "" };
}
export async function loginRpartstore(cfg: RpartstoreAuthConfig): Promise<RpartstoreToken> {
const issuer = cfg.issuer ?? DEFAULTS.issuer;
const clientId = cfg.clientId ?? DEFAULTS.clientId;
const redirectUri = cfg.redirectUri ?? DEFAULTS.redirectUri;
const scope = cfg.scope ?? DEFAULTS.scope;
const doFetch = cfg.fetchImpl ?? fetch;
const jar = new CookieJar();
const request = async (url: string, init: RequestInit = {}): Promise<Response> => {
const res = await doFetch(url, {
redirect: "manual",
...init,
headers: {
"User-Agent": USER_AGENT,
Cookie: jar.header(),
...(init.headers as Record<string, string> | undefined),
},
});
jar.absorb(res);
return res;
};
const verifier = b64url(randomBytes(48));
const challenge = b64url(createHash("sha256").update(verifier).digest());
const state = randomBytes(16).toString("hex");
const nonce = randomBytes(16).toString("hex");
// 1. hosted authorize page → stateToken
const authorize = new URL(`${issuer}/v1/authorize`);
for (const [k, v] of Object.entries({
client_id: clientId,
redirect_uri: redirectUri,
response_type: "code",
scope,
state,
nonce,
code_challenge: challenge,
code_challenge_method: "S256",
response_mode: "query",
})) {
authorize.searchParams.set(k, v);
}
const authorizeRes = await request(authorize.toString());
const html = await authorizeRes.text();
const stateTokenMatch =
html.match(/stateToken\s*=\s*'([^']+)'/) ?? html.match(/"stateToken":"([^"]+)"/);
if (authorizeRes.status !== 200 || !stateTokenMatch) {
throw new RpartstoreAuthError(
`authorize page did not expose a stateToken (HTTP ${authorizeRes.status})`,
"authorize",
);
}
const stateToken = stateTokenMatch[1].replace(/\\x([0-9A-Fa-f]{2})/g, (_, hex: string) =>
String.fromCharCode(Number.parseInt(hex, 16)),
);
const idx = async (
path: string,
body: Record<string, unknown>,
step: string,
): Promise<Record<string, any>> => {
const res = await request(`https://sso.renault.com/idp/idx/${path}`, {
method: "POST",
headers: { "Content-Type": ION, Accept: ION, Origin: "https://sso.renault.com" },
body: JSON.stringify(body),
});
const json = (await res.json().catch(() => ({}))) as Record<string, any>;
if (!res.ok) {
throw new RpartstoreAuthError(
`Okta ${step} failed (HTTP ${res.status})`,
step,
json.messages ?? json,
);
}
return json;
};
// 2-4. IDX remediation
const introspected = await idx("introspect", { stateToken }, "introspect");
const identified = await idx(
"identify",
{ identifier: cfg.username, stateHandle: introspected.stateHandle },
"identify",
);
const answered = await idx(
"challenge/answer",
{
credentials: { passcode: cfg.password },
stateHandle: identified.stateHandle ?? introspected.stateHandle,
},
"challenge",
);
const successHref: string | undefined = answered.success?.href;
if (!successHref) {
throw new RpartstoreAuthError(
"Okta did not return a success redirect (wrong password, locked account or MFA now required)",
"challenge",
answered.messages,
);
}
// 5. follow the redirect chain until the app callback carries ?code=
let next = successHref;
let code: string | undefined;
for (let hop = 0; hop < 6 && !code; hop += 1) {
const res = await request(next);
const location = res.headers.get("location");
if (!location) break;
const target = new URL(location, next);
const gotCode = target.searchParams.get("code");
if (gotCode) {
if (target.searchParams.get("state") !== state)
throw new RpartstoreAuthError("OAuth state mismatch", "redirect");
code = gotCode;
}
next = target.toString();
}
if (!code)
throw new RpartstoreAuthError("redirect chain ended without an authorization code", "redirect");
// 6. PKCE token exchange
const tokenRes = await request(`${issuer}/v1/token`, {
method: "POST",
headers: {
"Content-Type": "application/x-www-form-urlencoded",
Origin: "https://rpartstore.renault.com",
},
body: new URLSearchParams({
grant_type: "authorization_code",
redirect_uri: redirectUri,
code,
code_verifier: verifier,
client_id: clientId,
}).toString(),
});
const token = (await tokenRes.json().catch(() => ({}))) as {
access_token?: string;
error?: string;
error_description?: string;
};
if (!tokenRes.ok || !token.access_token) {
throw new RpartstoreAuthError(
`token exchange failed (HTTP ${tokenRes.status}): ${token.error_description ?? token.error ?? ""}`,
"token",
);
}
const { exp, sub } = decodeJwtExpiry(token.access_token);
return { accessToken: token.access_token, expiresAt: exp * 1000, subject: sub };
}

View File

@@ -0,0 +1,157 @@
import { type RpartstoreToken, loginRpartstore } from "./rpartstore.auth";
import { canonicalRpartstoreBrand } from "./rpartstore.routing";
import {
RpartstoreBffError,
RpartstoreRateLimitError,
RpartstoreSession,
} from "./rpartstore.session";
import type {
RpartstoreDecoded,
RpartstoreVehicle,
SearchVehicleResponsePayload,
} from "./rpartstore.types";
export const RPARTSTORE_DEFAULTS = {
brokerUrl: "wss://1po-bff.renault-edh.com/ws",
appVersion: "1.34.0.6",
webLanguage: "tr",
country: "TR",
} as const;
/** Where the (1 h, non-refreshable) Okta access token is cached between decodes. */
export interface TokenStore {
get(): Promise<RpartstoreToken | null>;
set(token: RpartstoreToken, ttlSeconds: number): Promise<void>;
clear(): Promise<void>;
}
export interface RpartstoreClientOptions {
username: string;
password: string;
tokenStore: TokenStore;
brokerUrl?: string;
appVersion?: string;
webLanguage?: string;
country?: string;
/** Injectable for tests. */
login?: typeof loginRpartstore;
sessionFactory?: (opts: ConstructorParameters<typeof RpartstoreSession>[0]) => RpartstoreSession;
logger?: { log: (msg: string) => void; warn: (msg: string) => void };
}
/** Refresh the token this many ms before its `exp`. */
const TOKEN_SKEW_MS = 5 * 60 * 1000;
export class RpartstoreClient {
private readonly login: typeof loginRpartstore;
private readonly sessionFactory: NonNullable<RpartstoreClientOptions["sessionFactory"]>;
constructor(private readonly opts: RpartstoreClientOptions) {
this.login = opts.login ?? loginRpartstore;
this.sessionFactory = opts.sessionFactory ?? ((o) => new RpartstoreSession(o));
}
/** Cached token when still valid (with skew), otherwise a fresh Okta login. */
async getToken(force = false): Promise<RpartstoreToken> {
if (!force) {
const cached = await this.opts.tokenStore.get();
if (cached && cached.expiresAt - Date.now() > TOKEN_SKEW_MS) return cached;
}
const token = await this.login({ username: this.opts.username, password: this.opts.password });
const ttl = Math.max(60, Math.floor((token.expiresAt - Date.now() - TOKEN_SKEW_MS) / 1000));
await this.opts.tokenStore.set(token, ttl);
this.opts.logger?.log(`[rpartstore] logged in as ${token.subject}, token ttl ${ttl}s`);
return token;
}
/**
* One VIN search on a fresh STOMP session (volume is capped at a handful per
* day, so a persistent socket buys nothing). Resolves to the decoded vehicle,
* `null` when RPartStore answers NOT_FOUND, and throws `RpartstoreRateLimitError`
* / other errors for the caller to handle. A CONNECT failure is retried once
* with a forced re-login (expired or revoked token).
*/
async searchVin(vin: string): Promise<RpartstoreDecoded | null> {
let token = await this.getToken();
const session = await this.openSession(token.accessToken).catch(async (err: Error) => {
this.opts.logger?.warn(`[rpartstore] CONNECT failed (${err.message}); re-login and retry`);
await this.opts.tokenStore.clear();
token = await this.getToken(true);
return this.openSession(token.accessToken);
});
try {
const country = this.opts.country ?? RPARTSTORE_DEFAULTS.country;
const lang = this.opts.webLanguage ?? RPARTSTORE_DEFAULTS.webLanguage;
const msg = await session.request<SearchVehicleResponsePayload>(
"/vehicles/search/vin-or-vrn",
{
requestId: crypto.randomUUID(),
value: vin,
queryType: "VIN",
userContext: {
webLanguage: lang,
documentLanguage: lang,
documentCountryLanguage: country,
documentFallbackLanguage: lang,
documentFallbackCountryLanguage: country,
userCountry: country,
r1Country: country,
},
countryCode: country,
includeEstimate: false,
},
{ expect: ["1PO/CATALOG/SEARCH_VEHICLE_RESPONSE"], timeoutMs: 15_000 },
);
const vehicle = msg.payload?.vehicles?.[0];
if (!vehicle) return null;
return normalizeVehicle(vehicle);
} catch (err) {
if (err instanceof RpartstoreBffError) {
const errorType = (err.payload as { errorType?: string } | undefined)?.errorType;
if (err.type === "1PO/CATALOG/SEARCH_VEHICLE_ERROR" && errorType === "NOT_FOUND")
return null;
if (err.type === "1PO/CATALOG/SEARCH_VEHICLE_NOT_COVERED_IN_COUNTRY") return null;
}
if (err instanceof RpartstoreRateLimitError) throw err;
throw err;
} finally {
session.close();
}
}
private async openSession(accessToken: string): Promise<RpartstoreSession> {
const session = this.sessionFactory({
brokerUrl: this.opts.brokerUrl ?? RPARTSTORE_DEFAULTS.brokerUrl,
accessToken,
profile: this.opts.username,
appVersion: this.opts.appVersion ?? RPARTSTORE_DEFAULTS.appVersion,
webLanguage: this.opts.webLanguage ?? RPARTSTORE_DEFAULTS.webLanguage,
logger: this.opts.logger
? { debug: () => undefined, warn: (m) => this.opts.logger?.warn(`[rpartstore] ${m}`) }
: undefined,
});
await session.connect();
return session;
}
}
export function normalizeVehicle(v: RpartstoreVehicle): RpartstoreDecoded {
const dh = v.dataHubVehicle ?? {};
const brandName = canonicalRpartstoreBrand(v.vehicleBrand) ?? v.vehicleBrand;
const clean = (s: string | undefined): string | null => {
const t = (s ?? "").trim();
return t.length > 0 ? t : null;
};
return {
brandName,
model: clean(v.model),
modelCode: clean(dh.modelTypeCode),
familyCode: clean(dh.familyCode),
modelYear: clean(dh.modelYear),
engine: clean(dh.engine),
gearbox: clean(dh.gearbox),
energyType: clean(dh.energyType),
manufacturingDate: clean(v.manufacturingDate),
raw: v,
};
}

View File

@@ -0,0 +1,140 @@
import { describe, expect, it } from "vitest";
import {
type RpartstoreCatalogCandidate,
pickRpartstoreCatalogMatch,
rpartstoreModelTokens,
} from "./rpartstore.matcher";
// Real PL24 `catalog_vehicles.model` values for Renault (prod, 2026-09-25).
const RENAULT = [
"ALASKAN",
"ARKANA EUROPE / XM3",
"ARKANA RUSYA",
"AUSTRAL/ESPACE VI/RAFALE",
"CAPTUR / QM3",
"CAPTUR II EUROPE/SYMBIOZ",
"CAPTUR II ÇİN",
"CAPTUR/KAPTUR",
"CLIO 4 / LUTECIA 4",
"CLIO 5/LUTECIA 5",
"DUSTER 1",
"DUSTER 2",
"DUSTER III / BIGSTER",
"EXPRESS",
"FLUENCE / FLUENCE Z.E.",
"KADJAR",
"KADJAR ÇİN",
"KANGOO 1",
"KANGOO 2 / KANGOO Z.E.",
"KANGOO 3",
"LATITUDE / SAFRANE 2",
"LOGAN\\-SANDERO 1/TONDAR 1",
"LOGAN\\-SANDERO 3 / TALIANT",
"MASTER 3",
"MASTER 4 VAN",
"MEGANE 1",
"MEGANE 2 / SCENIC 2",
"MEGANE 3 / SCENIC 3",
"MEGANE 4",
"MEGANE 4 SEDAN",
"RENAULT 5 EXPRESS / RAPID",
"RENAULT 9 / 11",
"TRAFIC 2",
"TRAFIC 3",
"X62 CHINE",
];
const DACIA = [
"DOKKER",
"DUSTER 1",
"DUSTER 2",
"DUSTER 3/BIGSTER",
"LOGAN\\-SANDERO 1/TONDAR 1",
"LOGAN\\-SANDERO 2/SYMBOL 2",
"LOGAN\\-SANDERO 3 / TALIANT",
];
const cands = (models: string[]): RpartstoreCatalogCandidate[] =>
models.map((model) => ({ id: model, model }));
describe("rpartstoreModelTokens", () => {
it("drops the model code in parentheses and converts roman generations", () => {
expect(rpartstoreModelTokens("Clio IV / Lutecia IV (B98)")).toEqual([
"CLIO",
"4",
"LUTECIA",
"4",
]);
expect(rpartstoreModelTokens("Duster III (SUV)")).toEqual(["DUSTER", "3"]);
expect(rpartstoreModelTokens("Megane I Classic (L64)")).toEqual(["MEGANE", "1", "CLASSIC"]);
});
it("keeps single-digit generation tokens and folds diacritics", () => {
expect(rpartstoreModelTokens("MEGANE 4 SEDAN")).toEqual(["MEGANE", "4", "SEDAN"]);
expect(rpartstoreModelTokens("KADJAR ÇİN")).toEqual(["KADJAR", "CIN"]);
expect(rpartstoreModelTokens("LOGAN\\-SANDERO 3 / TALIANT")).toEqual([
"LOGAN",
"SANDERO",
"3",
"TALIANT",
]);
});
});
describe("pickRpartstoreCatalogMatch (prod benchmark models → PL24 Renault catalog)", () => {
const expectPick = (model: string, expected: string | null, pool = RENAULT) =>
expect(pickRpartstoreCatalogMatch(model, cands(pool))).toBe(expected);
it("matches roman-numeral generations to digit generations", () => {
expectPick("Clio IV / Lutecia IV (B98)", "CLIO 4 / LUTECIA 4");
expectPick("Trafic III (J82)", "TRAFIC 3");
expectPick("Master III (F62)", "MASTER 3");
expectPick("Kangoo II (K61)", "KANGOO 2 / KANGOO Z.E.");
expectPick("Logan III (LJF)", "LOGAN\\-SANDERO 3 / TALIANT");
expectPick("Duster III (SUV)", "DUSTER III / BIGSTER");
});
it("prefers the body-qualified catalog when the decode carries the qualifier, and the plain one otherwise", () => {
expectPick("Megane IV Sedan (LFF)", "MEGANE 4 SEDAN");
expectPick("Megane IV (BFB)", "MEGANE 4");
expectPick("Megane I Classic (L64)", "MEGANE 1");
});
it("never picks a wrong-market catalog when a mainstream one exists", () => {
expectPick("Kadjar (HFE)", "KADJAR");
expectPick("Captur II (HJB)", "CAPTUR II EUROPE/SYMBIOZ");
});
it("does not let a newer generation steal a decode without a generation", () => {
expectPick("Captur (J87)", "CAPTUR / QM3");
expectPick("Express (KJK)", "EXPRESS");
});
it("matches multi-name catalogs and pre-2000 models", () => {
expectPick("Latitude / Safrane II (L43)", "LATITUDE / SAFRANE 2");
expectPick("Fluence (L38)", "FLUENCE / FLUENCE Z.E.");
expectPick("Renault 9 / 11 (L42)", "RENAULT 9 / 11");
});
it("matches Dacia models within the Dacia catalog", () => {
expectPick("Duster I (H79)", "DUSTER 1", DACIA);
expectPick("Duster II (HJD)", "DUSTER 2", DACIA);
expectPick("Dokker (K67)", "DOKKER", DACIA);
expectPick("Logan I (L90)", "LOGAN\\-SANDERO 1/TONDAR 1", DACIA);
});
it("returns null when nothing qualifies or the model is missing", () => {
expectPick("Arkana (LJL)", "ARKANA EUROPE / XM3");
expectPick("Twingo Z.E.", null);
expect(pickRpartstoreCatalogMatch(undefined, cands(RENAULT))).toBeNull();
expect(pickRpartstoreCatalogMatch("", cands(RENAULT))).toBeNull();
});
it("breaks exact ties by the fuller catalog", () => {
const pool: RpartstoreCatalogCandidate[] = [
{ id: "a", model: "KADJAR", categoryCount: 12 },
{ id: "b", model: "KADJAR", categoryCount: 44 },
];
expect(pickRpartstoreCatalogMatch("Kadjar (HFE)", pool)).toBe("b");
});
});

View File

@@ -0,0 +1,133 @@
import { MARKET_QUALIFIERS } from "../vinpin/vinpin.matcher";
/**
* Map an RPartStore model label ("Megane IV Sedan (LFF)", "Clio IV / Lutecia IV
* (B98)", "Duster III (SUV)") onto an EXISTING PL24 `catalog_vehicles` row of the
* same brand ("MEGANE 4 SEDAN", "CLIO 4 / LUTECIA 4", "DUSTER III / BIGSTER").
*
* Differences from the Vinpin matcher that made a dedicated picker necessary:
* - RPartStore writes generations as roman numerals, PL24 mostly as digits
* ("IV" vs "4") — both sides are normalised to digits.
* - Generation digits are single characters; the Vinpin tokenizer drops tokens
* shorter than 2 chars, which would make "MEGANE 4" ≡ "MEGANE".
* - Body qualifiers (SEDAN, CLASSIC, …) are optional on the catalog side:
* "Megane I Classic" must still match "MEGANE 1".
* - PL24 Renault/Dacia rows carry no year, so there is no year scoring.
*/
export interface RpartstoreCatalogCandidate {
id: string;
model: string | null;
/** Tiebreak only: fuller catalog wins. */
categoryCount?: number | null;
}
const ROMAN: Record<string, string> = {
I: "1",
II: "2",
III: "3",
IV: "4",
V: "5",
VI: "6",
VII: "7",
VIII: "8",
};
/** Body / trim words that PL24 may omit; they only add a small bonus when both sides have them. */
const BODY_QUALIFIERS = new Set<string>([
"SEDAN",
"CLASSIC",
"HATCHBACK",
"HB",
"ESTATE",
"GRANDTOUR",
"SPORTTOURER",
"SW",
"BREAK",
"VAN",
"COMBI",
"KOMBI",
"CABRIO",
"COUPE",
"SASI",
"CHASSIS",
"PICKUP",
"PHASE",
"PH",
]);
/** Tokenize a model label: fold diacritics, drop parenthetical codes/years, split on
* non-alphanumerics, convert roman generation numerals to digits. Keeps 1-char
* numeric tokens (generations) and drops 1-char alpha noise ("Z.E." → "ZE" kept
* as-is is not needed for matching). */
export function rpartstoreModelTokens(s: string | null | undefined): string[] {
if (!s) return [];
return s
.toUpperCase()
.normalize("NFD")
.replace(/\p{M}/gu, "")
.replace(/\([^)]*\)/g, " ")
.replace(/[^A-Z0-9]+/g, " ")
.split(" ")
.map((t) => t.trim())
.filter((t) => t.length > 0)
.map((t) => ROMAN[t] ?? t)
.filter((t) => t.length >= 2 || /^\d$/.test(t));
}
const isNumeric = (t: string): boolean => /^\d+$/.test(t);
/**
* Returns the id of the best candidate or null. A candidate qualifies when every
* "core" decoded token (everything except body qualifiers) appears among its
* tokens. Score, strongest first: market-qualifier penalty (wrong-market catalogs
* only win when alone), exact-match bonus, body-qualifier overlap bonus, a hard
* penalty for candidates that add a generation number the decode lacks
* ("Captur" must not pick "CAPTUR II"), a mild penalty per extra token, then the
* fuller catalog as tiebreak.
*/
export function pickRpartstoreCatalogMatch(
model: string | null | undefined,
candidates: RpartstoreCatalogCandidate[],
): string | null {
const decoded = rpartstoreModelTokens(model);
if (decoded.length === 0) return null;
const core = decoded.filter((t) => !BODY_QUALIFIERS.has(t));
const required = core.length > 0 ? core : decoded;
const decodedSet = new Set(decoded);
const decodedHasGeneration = decoded.some(isNumeric);
let best: { id: string; score: number; categoryCount: number } | null = null;
for (const c of candidates) {
const tokens = rpartstoreModelTokens(c.model);
if (tokens.length === 0) continue;
const tokenSet = new Set(tokens);
if (!required.every((t) => tokenSet.has(t))) continue;
const extras = tokens.filter((t) => !decodedSet.has(t));
const marketExtras = extras.filter((t) => MARKET_QUALIFIERS.has(t)).length;
const generationExtras = decodedHasGeneration ? 0 : extras.filter(isNumeric).length;
const qualifierOverlap = decoded.filter(
(t) => BODY_QUALIFIERS.has(t) && tokenSet.has(t),
).length;
const exact = extras.length === 0 && tokens.length === decoded.length;
const score =
1000 -
marketExtras * 5000 -
generationExtras * 150 -
extras.length * 20 +
qualifierOverlap * 100 +
(exact ? 300 : 0);
const categoryCount = c.categoryCount ?? 0;
if (
!best ||
score > best.score ||
(score === best.score && categoryCount > best.categoryCount)
) {
best = { id: c.id, score, categoryCount };
}
}
return best?.id ?? null;
}

View File

@@ -0,0 +1,32 @@
import { getBrandFromWmi } from "@sase/shared";
/** WMIs routed to RPartStore even when the shared WMI table is silent. VF1/VF2 =
* Renault (France), VF6 = Renault (Trucks/Sofasa), UU1 = Dacia (Romania). VF7 is
* Citroën and is deliberately NOT here. */
const RPARTSTORE_WMIS = new Set<string>(["VF1", "VF2", "VF6", "UU1"]);
const RPARTSTORE_BRANDS = new Set<string>(["renault", "dacia"]);
/** Canonical brand names as they appear in `catalog_vehicles.brand_name`. */
export function canonicalRpartstoreBrand(
brand: string | null | undefined,
): "Renault" | "Dacia" | null {
const b = (brand ?? "").trim().toLowerCase();
if (b === "renault") return "Renault";
if (b === "dacia") return "Dacia";
return null;
}
/**
* Should this VIN be offered to the RPartStore decode fallback? Renault/Dacia
* only: decided by the WMI first (shared table, then the explicit allowlist),
* with the identified browse brand as a last resort (Dacia-badged cars built
* under a Renault WMI still say "Dacia" in the identification).
*/
export function isRpartstoreVin(vin: string, browseBrand?: string | null): boolean {
const wmi = (vin || "").toUpperCase().slice(0, 3);
const wmiBrand = getBrandFromWmi(wmi);
if (wmiBrand && RPARTSTORE_BRANDS.has(wmiBrand.toLowerCase())) return true;
if (RPARTSTORE_WMIS.has(wmi)) return true;
return !!browseBrand && RPARTSTORE_BRANDS.has(browseBrand.trim().toLowerCase());
}

View File

@@ -0,0 +1,256 @@
/**
* One STOMP-over-WebSocket session against the RPartStore BFF.
*
* Request/response is correlated by the `trace-id` header we send and the `traceId`
* header the server echoes. One request can yield several MESSAGE frames
* (e.g. a VIN search → SEARCH_VEHICLE_RESPONSE, COMMAND_PROCESSING_RESPONSE,
* EXPLODED_TREE_RESPONSE), so a request resolves on the first frame whose `type`
* matches one of the expected terminal types, and errors on a frame from
* `/user/queue/error` or a `*_RATE_LIMIT_EXCEEDED` type.
*
* Uses Node 22's built-in WebSocket (no Origin / subprotocol required by the server).
*/
import { randomUUID } from "node:crypto";
import { decodeFrame, encodeFrame } from "./rpartstore.stomp";
export interface BffMessage<T = unknown> {
type: string;
payload: T;
}
export interface SessionOptions {
brokerUrl: string;
accessToken: string;
profile: string;
appVersion: string;
webLanguage: string;
connectTimeoutMs?: number;
/** STOMP heart-beat interval (ms) negotiated with the server. */
heartbeatMs?: number;
onClose?: (reason: string) => void;
logger?: { debug: (msg: string) => void; warn: (msg: string) => void };
}
export interface RequestOptions {
/** Message `type`s that complete the request. */
expect: string[];
timeoutMs?: number;
}
export class RpartstoreRateLimitError extends Error {
constructor(
readonly retryAfterSeconds: number,
readonly limitType: string,
) {
super(`RPartStore rate limit (${limitType}), retry after ${retryAfterSeconds}s`);
this.name = "RpartstoreRateLimitError";
}
}
export class RpartstoreBffError extends Error {
constructor(
readonly type: string,
readonly payload: unknown,
) {
super(`RPartStore BFF error ${type}`);
this.name = "RpartstoreBffError";
}
}
interface Pending {
expect: Set<string>;
resolve: (msg: BffMessage) => void;
reject: (err: Error) => void;
timer: NodeJS.Timeout;
}
export class RpartstoreSession {
private ws: WebSocket | null = null;
private readonly pending = new Map<string, Pending>();
private heartbeatTimer: NodeJS.Timeout | null = null;
private closed = false;
constructor(private readonly opts: SessionOptions) {}
get isOpen(): boolean {
return !this.closed && this.ws !== null && this.ws.readyState === WebSocket.OPEN;
}
/** Opens the socket, sends CONNECT and subscribes to the user queues. Resolves on CONNECTED. */
connect(): Promise<void> {
const {
brokerUrl,
accessToken,
profile,
appVersion,
webLanguage,
connectTimeoutMs = 10_000,
heartbeatMs = 60_000,
} = this.opts;
return new Promise<void>((resolve, reject) => {
let settled = false;
const fail = (err: Error): void => {
if (!settled) {
settled = true;
reject(err);
}
this.teardown(err.message);
};
const timer = setTimeout(() => fail(new Error("STOMP CONNECT timeout")), connectTimeoutMs);
let ws: WebSocket;
try {
ws = new WebSocket(brokerUrl, ["v12.stomp"]);
} catch (err) {
clearTimeout(timer);
reject(err instanceof Error ? err : new Error(String(err)));
return;
}
this.ws = ws;
ws.onopen = () => {
ws.send(
encodeFrame("CONNECT", {
"trace-id": randomUUID(),
"x-auth-token": accessToken,
"selected-profile": profile,
"app-version": appVersion,
"web-language": webLanguage,
"accept-version": "1.2,1.1,1.0",
"heart-beat": `${heartbeatMs},${heartbeatMs}`,
}),
);
};
ws.onmessage = (event: MessageEvent) => {
const frame = decodeFrame(typeof event.data === "string" ? event.data : String(event.data));
if (!frame) return; // heartbeat
if (frame.command === "CONNECTED") {
clearTimeout(timer);
ws.send(encodeFrame("SUBSCRIBE", { id: "sub-0", destination: "/user/queue/main" }));
ws.send(encodeFrame("SUBSCRIBE", { id: "sub-1", destination: "/user/queue/error" }));
this.startHeartbeat(heartbeatMs);
settled = true;
resolve();
return;
}
if (frame.command === "ERROR") {
fail(
new Error(`STOMP ERROR: ${frame.headers.message ?? ""} ${frame.body.slice(0, 200)}`),
);
return;
}
if (frame.command === "MESSAGE") this.onMessage(frame.headers, frame.body);
};
ws.onerror = () => fail(new Error("WebSocket error"));
ws.onclose = (event: { code: number; reason: string }) => {
clearTimeout(timer);
const reason = `WebSocket closed (${event.code}${event.reason ? ` ${event.reason}` : ""})`;
if (!settled) fail(new Error(reason));
else this.teardown(reason);
};
});
}
/** Publishes `{payload}` to `/app/<path>` and resolves on the first expected message type. */
request<T = unknown>(
path: string,
payload: unknown,
options: RequestOptions,
): Promise<BffMessage<T>> {
const ws = this.ws;
if (!this.isOpen || !ws) return Promise.reject(new Error("STOMP session is not connected"));
const traceId = randomUUID();
const body = JSON.stringify({ payload });
const { expect, timeoutMs = 15_000 } = options;
return new Promise<BffMessage<T>>((resolve, reject) => {
const timer = setTimeout(() => {
this.pending.delete(traceId);
reject(new Error(`RPartStore request ${path} timed out after ${timeoutMs}ms`));
}, timeoutMs);
this.pending.set(traceId, {
expect: new Set(expect),
resolve: (msg) => resolve(msg as BffMessage<T>),
reject,
timer,
});
ws.send(encodeFrame("SEND", { destination: `/app${path}`, "trace-id": traceId }, body));
this.opts.logger?.debug(`→ ${path} trace=${traceId}`);
});
}
close(): void {
this.teardown("closed by client");
}
private onMessage(headers: Record<string, string>, body: string): void {
const traceId = headers.traceId;
const pending = traceId ? this.pending.get(traceId) : undefined;
let msg: BffMessage;
try {
msg = JSON.parse(body) as BffMessage;
} catch {
this.opts.logger?.warn(
`unparseable BFF message on ${headers.destination ?? "?"} trace=${traceId ?? "?"}`,
);
return;
}
if (!pending) return; // unsolicited push (search-history refresh etc.)
const isError = headers.destination === "/user/queue/error";
if (/RATE_LIMIT_EXCEEDED$/.test(msg.type)) {
const p = msg.payload as { retryAfterSeconds?: number; limitType?: string };
this.settle(traceId, pending, (pend) =>
pend.reject(
new RpartstoreRateLimitError(p.retryAfterSeconds ?? 10, p.limitType ?? "SHORT_TERM"),
),
);
return;
}
if (isError) {
this.settle(traceId, pending, (pend) =>
pend.reject(new RpartstoreBffError(msg.type, msg.payload)),
);
return;
}
if (pending.expect.has(msg.type)) {
this.settle(traceId, pending, (pend) => pend.resolve(msg));
}
// other frames on the same trace (COMMAND_PROCESSING_RESPONSE, EXPLODED_TREE_RESPONSE…) are ignored here
}
private settle(traceId: string, pending: Pending, fn: (p: Pending) => void): void {
clearTimeout(pending.timer);
this.pending.delete(traceId);
fn(pending);
}
private startHeartbeat(intervalMs: number): void {
this.stopHeartbeat();
this.heartbeatTimer = setInterval(() => {
if (this.isOpen) this.ws?.send("\n");
}, intervalMs);
}
private stopHeartbeat(): void {
if (this.heartbeatTimer) clearInterval(this.heartbeatTimer);
this.heartbeatTimer = null;
}
private teardown(reason: string): void {
if (this.closed) return;
this.closed = true;
this.stopHeartbeat();
for (const [traceId, p] of this.pending) {
clearTimeout(p.timer);
p.reject(new Error(`STOMP session ended: ${reason}`));
this.pending.delete(traceId);
}
const ws = this.ws;
this.ws = null;
if (ws && ws.readyState !== WebSocket.CLOSED) {
try {
ws.close();
} catch {
/* ignore */
}
}
this.opts.onClose?.(reason);
}
}

View File

@@ -0,0 +1,50 @@
import { describe, expect, it } from "vitest";
import { decodeFrame, encodeFrame } from "./rpartstore.stomp";
describe("rpartstore STOMP codec", () => {
it("encodes a SEND frame exactly like the RPartStore web app", () => {
const body = '{"payload":{"requestId":"r","value":"VF1RFE00653633190","queryType":"VIN"}}';
const frame = encodeFrame(
"SEND",
{ destination: "/app/vehicles/search/vin-or-vrn", "trace-id": "t-1" },
body,
);
expect(frame).toBe(
`SEND\ndestination:/app/vehicles/search/vin-or-vrn\ntrace-id:t-1\ncontent-length:${Buffer.byteLength(body)}\n\n${body}\u0000`,
);
});
it("does not escape CONNECT header values (STOMP 1.2 rule) but escapes others", () => {
expect(encodeFrame("CONNECT", { "x-auth-token": "a:b" })).toBe(
"CONNECT\nx-auth-token:a:b\n\n\u0000",
);
expect(encodeFrame("SEND", { destination: "/app/x", note: "a:b\nc" })).toContain(
"note:a\\cb\\nc",
);
});
it("decodes a live MESSAGE frame with headers and JSON body", () => {
const raw =
"MESSAGE\ntraceId:f68cc229\ncontent-type:text/plain;charset=UTF-8\ndestination:/user/queue/main\nsubscription:sub-0\nmessage-id:975afec7-1\ncontent-length:64\n\n" +
'{"type":"1PO/COMMON/COMMAND_PROCESSING_RESPONSE","payload":true}\u0000';
const frame = decodeFrame(raw);
expect(frame?.command).toBe("MESSAGE");
expect(frame?.headers.traceId).toBe("f68cc229");
expect(frame?.headers["content-type"]).toBe("text/plain;charset=UTF-8");
expect(JSON.parse(frame!.body)).toEqual({
type: "1PO/COMMON/COMMAND_PROCESSING_RESPONSE",
payload: true,
});
});
it("treats heartbeat newlines as no frame", () => {
expect(decodeFrame("\n")).toBeNull();
expect(decodeFrame("")).toBeNull();
});
it("keeps the first value of a repeated header", () => {
const frame = decodeFrame("MESSAGE\nk:first\nk:second\n\nbody\u0000");
expect(frame?.headers.k).toBe("first");
expect(frame?.body).toBe("body");
});
});

View File

@@ -0,0 +1,69 @@
/**
* Minimal STOMP 1.2 codec used by the RPartStore BFF client.
* The BFF speaks STOMP over a plain WebSocket (`wss://1po-bff.renault-edh.com/ws`):
* CONNECT → CONNECTED, SUBSCRIBE /user/queue/main + /user/queue/error,
* SEND /app/<path> with a `trace-id` header, and 1..N MESSAGE frames back
* carrying the same `traceId` header and a `{type, payload}` JSON body.
*/
export interface StompFrame {
command: string;
headers: Record<string, string>;
body: string;
}
const NULL = "\u0000";
/** STOMP 1.2 header value escaping (RFC: `\\`, `\n` → `\\n`, `:` → `\\c`, `\r` → `\\r`). */
function escapeHeader(value: string): string {
return value
.replace(/\\/g, "\\\\")
.replace(/\r/g, "\\r")
.replace(/\n/g, "\\n")
.replace(/:/g, "\\c");
}
function unescapeHeader(value: string): string {
return value
.replace(/\\n/g, "\n")
.replace(/\\r/g, "\r")
.replace(/\\c/g, ":")
.replace(/\\\\/g, "\\");
}
export function encodeFrame(
command: string,
headers: Record<string, string | number>,
body = "",
): string {
const lines = [command];
for (const [k, v] of Object.entries(headers)) {
lines.push(`${k}:${command === "CONNECT" ? String(v) : escapeHeader(String(v))}`);
}
if (body && headers["content-length"] === undefined) {
lines.push(`content-length:${Buffer.byteLength(body, "utf8")}`);
}
return `${lines.join("\n")}\n\n${body}${NULL}`;
}
/** Heartbeat frames are a bare newline; they decode to `null`. */
export function decodeFrame(raw: string): StompFrame | null {
if (raw === "\n" || raw === "\r\n" || raw === "") return null;
const headerEnd = raw.indexOf("\n\n");
if (headerEnd < 0) return null;
const headerBlock = raw.slice(0, headerEnd).split("\n");
const command = headerBlock[0].trim();
const headers: Record<string, string> = {};
for (const line of headerBlock.slice(1)) {
const i = line.indexOf(":");
if (i < 0) continue;
const key = line.slice(0, i);
// STOMP: the first occurrence of a repeated header wins.
if (headers[key] === undefined)
headers[key] =
command === "CONNECTED" ? line.slice(i + 1) : unescapeHeader(line.slice(i + 1));
}
let body = raw.slice(headerEnd + 2);
if (body.endsWith(NULL)) body = body.slice(0, -1);
return { command, headers, body };
}

View File

@@ -0,0 +1,54 @@
/**
* Shapes of the RPartStore BFF `SEARCH_VEHICLE_RESPONSE` payload (DATAHUB
* catalog source, TR market). Captured live 2026-09-25; see
* /home/s/ss/rpartstore-dogrudan-vin-decode-2026-09-25.md for the protocol notes.
*/
export interface RpartstoreDataHubVehicle {
name?: string; // "RENAULT Kadjar (HFE)"
modelType?: string; // "SUV"
modelTypeCode?: string; // "HFE"
familyCode?: string; // "XFE"
bodyType?: string; // "HFE"
engine?: string; // "1.5 DCI DİZEL MOTOR [K9K]"
gearbox?: string; // "6 VİTESLİ KAVRAMA VİTES KUTUSU:DC4 [DC4]"
energyType?: string; // "MOTORIN"
powerKw?: string; // "066 KW POWER" | ""
modelYear?: string; // "2015"
vehicleAge?: number;
wheelbaseLength?: string;
roofHeight?: string;
}
export interface RpartstoreVehicle {
catalogSource: string; // "DATAHUB"
vin: string;
vehicleKey: string;
model: string | undefined; // "Kadjar (HFE)" — undefined on some pre-2000 cars
vehicleBrand: string; // "RENAULT" | "DACIA"
country: string; // "TR"
manufacturingDate?: string; // "2015-07-28"
imageUrl?: string;
dataHubVehicle?: RpartstoreDataHubVehicle;
vehicleIdentifiedBy?: string;
}
export interface SearchVehicleResponsePayload {
vehicles: RpartstoreVehicle[];
requestId: string;
searchedCountry: string;
}
/** Normalised decode result stored in `rpartstore_decodes`. */
export interface RpartstoreDecoded {
brandName: string; // "Renault" | "Dacia" (canonical casing, matches catalog_vehicles.brand_name)
model: string | null; // "Kadjar (HFE)"
modelCode: string | null; // "HFE"
familyCode: string | null; // "XFE"
modelYear: string | null; // "2015"
engine: string | null;
gearbox: string | null;
energyType: string | null;
manufacturingDate: string | null;
raw: RpartstoreVehicle;
}

View File

@@ -36,5 +36,6 @@ export const QUEUE_NAMES = {
EXPERT_REWARDS: "expert-rewards", EXPERT_REWARDS: "expert-rewards",
PART_PRICE_REFRESH: "part-price-refresh", PART_PRICE_REFRESH: "part-price-refresh",
VINPIN_DECODE: "vinpin-decode", VINPIN_DECODE: "vinpin-decode",
RPARTSTORE_DECODE: "rpartstore-decode",
CANONICAL_BACKFILL: "canonical-backfill", CANONICAL_BACKFILL: "canonical-backfill",
} as const; } as const;

View File

@@ -21,6 +21,10 @@ import {
PartPriceRefreshQueueProvider, PartPriceRefreshQueueProvider,
} from "./queues/part-price-refresh.queue"; } from "./queues/part-price-refresh.queue";
import { QUERY_CLEANUP_QUEUE, QueryCleanupQueueProvider } from "./queues/query-cleanup.queue"; import { QUERY_CLEANUP_QUEUE, QueryCleanupQueueProvider } from "./queues/query-cleanup.queue";
import {
RPARTSTORE_DECODE_QUEUE,
RpartstoreDecodeQueueProvider,
} from "./queues/rpartstore-decode.queue";
import { import {
SUBSCRIPTION_EXPIRY_QUEUE, SUBSCRIPTION_EXPIRY_QUEUE,
SubscriptionExpiryQueueProvider, SubscriptionExpiryQueueProvider,
@@ -39,6 +43,7 @@ import { VINPIN_DECODE_QUEUE, VinpinDecodeQueueProvider } from "./queues/vinpin-
ExpertRewardsQueueProvider, ExpertRewardsQueueProvider,
PartPriceRefreshQueueProvider, PartPriceRefreshQueueProvider,
VinpinDecodeQueueProvider, VinpinDecodeQueueProvider,
RpartstoreDecodeQueueProvider,
CanonicalBackfillQueueProvider, CanonicalBackfillQueueProvider,
PrefetchWorkerService, PrefetchWorkerService,
], ],
@@ -52,6 +57,7 @@ import { VINPIN_DECODE_QUEUE, VinpinDecodeQueueProvider } from "./queues/vinpin-
EXPERT_REWARDS_QUEUE, EXPERT_REWARDS_QUEUE,
PART_PRICE_REFRESH_QUEUE, PART_PRICE_REFRESH_QUEUE,
VINPIN_DECODE_QUEUE, VINPIN_DECODE_QUEUE,
RPARTSTORE_DECODE_QUEUE,
CANONICAL_BACKFILL_QUEUE, CANONICAL_BACKFILL_QUEUE,
], ],
}) })

View File

@@ -74,14 +74,21 @@ export async function checkCooldown(redis: RedisService, source: string): Promis
} }
} }
/** Current hour (0–23) in Europe/Istanbul. */ /** Europe/Istanbul hour (0–23) at an arbitrary instant. */
function currentIstanbulHour(): number { function istanbulHourAt(ms: number): number {
const hourStr = new Intl.DateTimeFormat("en-US", { const hourStr = new Intl.DateTimeFormat("en-US", {
timeZone: "Europe/Istanbul", timeZone: "Europe/Istanbul",
hour: "numeric", hour: "numeric",
hour12: false, hour12: false,
}).format(new Date()); }).format(new Date(ms));
return Number.parseInt(hourStr, 10); // `% 24` because the h24 hour cycle renders midnight as "24", which would put
// the hour outside every window and silently park the source forever.
return Number.parseInt(hourStr, 10) % 24;
}
/** Current hour (0–23) in Europe/Istanbul. */
function currentIstanbulHour(): number {
return istanbulHourAt(Date.now());
} }
/** /**
@@ -115,6 +122,35 @@ export function checkTimeWindow(source: string): void {
} }
} }
/**
* The first instant at or after `fromMs` that falls inside `source`'s scrape
* window. Returns `fromMs` unchanged when the window is disabled (the default
* 0–24) or when `fromMs` is already inside it.
*
* WHY (plv2.md, finding consumers_jobs-03 — the budget/window deadlock):
* `checkSourceDailyBudget` defers a spent source to the next UTC midnight. With
* PREFETCH_PL24_START=9 that midnight lands at 03:00 Europe/Istanbul — six hours
* BEFORE the window opens. The woken job therefore did no work, threw
* `time-window`, and was deferred again to 09:00 — by which point the fresh
* daily budget had already been spent by the same stampede of no-op wake-ups.
* Measured on prod 2026-09-20: 600/600 pl24 budget consumed, 0 catalog requests
* and 0 new categories for the whole day. Landing the deferral inside the window
* breaks the cycle.
*/
export function alignToWindow(source: string, fromMs: number): number {
if (source !== "pl24") return fromMs;
if (PL24_WINDOW_START <= 0 && PL24_WINDOW_END >= 24) return fromMs;
let t = fromMs;
// Step by the hour rather than constructing a local-midnight date: DST-safe and
// free of month/year rollover edge cases. 48 steps covers any window shape.
for (let i = 0; i < 48; i++) {
const h = istanbulHourAt(t);
if (h >= PL24_WINDOW_START && h < PL24_WINDOW_END) return t;
t += 3_600_000;
}
return t;
}
/** /**
* Milliseconds until the next 09:00 Europe/Istanbul. * Milliseconds until the next 09:00 Europe/Istanbul.
*/ */

View File

@@ -0,0 +1,286 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
/**
* Regression lock for the PL24 budget/window DEADLOCK (plv2.md, consumers_jobs-03).
*
* Two bugs combined to kill PL24 prefetch outright on prod:
* 1. `checkSourceDailyBudget` was debited in `process()` while the
* business-hours gate still sat at the top of each handler — so a job that
* woke outside the window paid a budget unit to do nothing.
* 2. A spent source was deferred to the next UTC midnight, which is 03:00
* Europe/Istanbul — six hours BEFORE a 09:00 window opens. Every deferred
* job therefore woke early, burned a unit of the fresh daily budget, hit
* the window gate and was deferred again.
*
* Measured on prod 2026-09-20 before the fix: `prefetch:daily:pl24` at 600/600
* with ZERO catalog requests and ZERO new pl24 categories for the whole day.
*
* The window constants are read at module load, so every case imports the
* modules fresh with the env already set.
*/
const REAL_ENV = { ...process.env };
/** 2026-09-20 00:00 UTC = 03:00 Europe/Istanbul — outside a 09:00–18:00 window. */
const OUTSIDE_MS = Date.UTC(2026, 8, 20, 0, 0, 0);
/** 2026-09-20 09:00 UTC = 12:00 Europe/Istanbul — inside it. */
const INSIDE_MS = Date.UTC(2026, 8, 20, 9, 0, 0);
function istanbulHour(ms: number): number {
return (
Number.parseInt(
new Intl.DateTimeFormat("en-US", {
timeZone: "Europe/Istanbul",
hour: "numeric",
hour12: false,
}).format(new Date(ms)),
10,
) % 24
);
}
async function loadModules(env: Record<string, string | undefined>) {
vi.resetModules();
for (const [k, v] of Object.entries(env)) {
if (v === undefined) delete process.env[k];
else process.env[k] = v;
}
const utils = await import("./prefetch-utils");
const worker = await import("./prefetch-worker.service");
return { utils, worker };
}
/** Minimal deps: only what `process()` touches before dispatching a job. */
function makeDeps(redisOverrides: Record<string, unknown> = {}) {
const incrCalls: string[] = [];
const redis = {
exists: vi.fn(async (..._a: unknown[]) => false),
get: vi.fn(async (..._a: unknown[]): Promise<string | null> => null),
set: vi.fn(async (..._a: unknown[]) => undefined),
del: vi.fn(async (..._a: unknown[]) => undefined),
incr: vi.fn(async (k: string) => {
incrCalls.push(k);
return 1;
}),
expire: vi.fn(async (..._a: unknown[]) => undefined),
ttl: vi.fn(async (..._a: unknown[]) => -2), // no cooldown key
setNx: vi.fn(async (..._a: unknown[]) => true),
getJson: vi.fn(async (..._a: unknown[]): Promise<unknown> => null),
setJson: vi.fn(async (..._a: unknown[]) => undefined),
...redisOverrides,
};
const queue = {
name: "catalog-prefetch",
add: vi.fn(async (..._a: unknown[]) => undefined),
getJob: vi.fn(async (..._a: unknown[]): Promise<unknown> => null),
};
const categoriesService = {
getCategoryWithParts: vi.fn(async (..._a: unknown[]) => ({ parts: [] })),
getChildren: vi.fn(async (..._a: unknown[]) => []),
};
return { redis, queue, categoriesService, incrCalls };
}
function makeJob(name: string, data: Record<string, unknown>) {
return {
name,
data,
moveToDelayed: vi.fn(async (..._a: unknown[]) => undefined),
attemptsMade: 0,
};
}
const dailyIncrs = (keys: string[]) => keys.filter((k) => k.startsWith("prefetch:daily:pl24"));
describe("alignToWindow", () => {
afterEach(() => {
process.env = { ...REAL_ENV };
});
it("pushes a pre-window instant into the configured window", async () => {
const { utils } = await loadModules({ PREFETCH_PL24_START: "9", PREFETCH_PL24_END: "18" });
const aligned = utils.alignToWindow("pl24", OUTSIDE_MS);
expect(aligned).toBeGreaterThan(OUTSIDE_MS);
const h = istanbulHour(aligned);
expect(h).toBeGreaterThanOrEqual(9);
expect(h).toBeLessThan(18);
});
it("leaves an in-window instant untouched", async () => {
const { utils } = await loadModules({ PREFETCH_PL24_START: "9", PREFETCH_PL24_END: "18" });
expect(utils.alignToWindow("pl24", INSIDE_MS)).toBe(INSIDE_MS);
});
it("is a no-op for sources that have no window", async () => {
const { utils } = await loadModules({ PREFETCH_PL24_START: "9", PREFETCH_PL24_END: "18" });
expect(utils.alignToWindow("emex", OUTSIDE_MS)).toBe(OUTSIDE_MS);
expect(utils.alignToWindow("parts-catalogs", OUTSIDE_MS)).toBe(OUTSIDE_MS);
});
it("is a no-op when the window is disabled (the default 0–24)", async () => {
const { utils } = await loadModules({
PREFETCH_PL24_START: undefined,
PREFETCH_PL24_END: undefined,
});
expect(utils.alignToWindow("pl24", OUTSIDE_MS)).toBe(OUTSIDE_MS);
});
});
describe("process() gate order — window before daily budget", () => {
beforeEach(() => {
vi.useFakeTimers();
});
afterEach(() => {
vi.useRealTimers();
process.env = { ...REAL_ENV };
});
it("does NOT debit the daily budget for a job deferred by the window", async () => {
const { worker } = await loadModules({
PREFETCH_PL24_START: "9",
PREFETCH_PL24_END: "18",
PL24_TR_DISABLED: undefined,
});
vi.setSystemTime(OUTSIDE_MS);
const { redis, queue, categoriesService, incrCalls } = makeDeps();
const svc = new worker.PrefetchWorkerService(
queue as never,
queue as never,
categoriesService as never,
redis as never,
{ payload: vi.fn(async () => ({})) } as never,
{} as never,
);
const job = makeJob("prefetch-parts", {
vehicleId: "v1",
categoryId: "c1",
source: "pl24",
fast: true,
});
await expect(
(svc as never as { process: (j: unknown, t?: string) => Promise<void> }).process(job, "tok"),
).rejects.toThrow();
expect(job.moveToDelayed).toHaveBeenCalled();
expect(dailyIncrs(incrCalls)).toHaveLength(0);
// The per-minute counter is not charged either: the window gate is pure
// clock arithmetic and runs before any Redis write.
expect(incrCalls).toHaveLength(0);
// …and the handler never ran, so nothing was fetched upstream either.
expect(categoriesService.getCategoryWithParts).not.toHaveBeenCalled();
});
it("defers the window-blocked job to an instant inside the window", async () => {
const { worker } = await loadModules({
PREFETCH_PL24_START: "9",
PREFETCH_PL24_END: "18",
PL24_TR_DISABLED: undefined,
});
vi.setSystemTime(OUTSIDE_MS);
const { redis, queue, categoriesService } = makeDeps();
const svc = new worker.PrefetchWorkerService(
queue as never,
queue as never,
categoriesService as never,
redis as never,
{ payload: vi.fn(async () => ({})) } as never,
{} as never,
);
const job = makeJob("prefetch-parts", {
vehicleId: "v1",
categoryId: "c1",
source: "pl24",
fast: true,
});
await expect(
(svc as never as { process: (j: unknown, t?: string) => Promise<void> }).process(job, "tok"),
).rejects.toThrow();
const [when] = job.moveToDelayed.mock.calls[0] as [number];
const h = istanbulHour(when);
expect(h).toBeGreaterThanOrEqual(9);
expect(h).toBeLessThan(18);
});
it("debits the daily budget once the window is open", async () => {
const { worker } = await loadModules({
PREFETCH_PL24_START: "9",
PREFETCH_PL24_END: "18",
PL24_TR_DISABLED: undefined,
});
vi.setSystemTime(INSIDE_MS);
const { redis, queue, categoriesService, incrCalls } = makeDeps();
const svc = new worker.PrefetchWorkerService(
queue as never,
queue as never,
categoriesService as never,
redis as never,
{ payload: vi.fn(async () => ({})) } as never,
{} as never,
);
const job = makeJob("prefetch-parts", {
vehicleId: "v1",
categoryId: "c1",
source: "pl24",
fast: true,
});
await (svc as never as { process: (j: unknown, t?: string) => Promise<void> }).process(
job,
"tok",
);
expect(dailyIncrs(incrCalls)).toHaveLength(1);
expect(categoriesService.getCategoryWithParts).toHaveBeenCalledWith("c1");
});
});
describe("checkSourceDailyBudget — spent source retries inside the window", () => {
beforeEach(() => {
vi.useFakeTimers();
});
afterEach(() => {
vi.useRealTimers();
process.env = { ...REAL_ENV };
});
it("never parks a spent source at 03:00 Istanbul again", async () => {
const { worker } = await loadModules({
PREFETCH_PL24_START: "9",
PREFETCH_PL24_END: "18",
PREFETCH_DAILY_PL24: "600",
});
vi.setSystemTime(INSIDE_MS);
// Counter already at the fast-lane ceiling → the next job must be deferred.
const { redis, queue, categoriesService } = makeDeps({
get: vi.fn(async (k: string) => (String(k).startsWith("prefetch:daily:pl24") ? "600" : null)),
});
const svc = new worker.PrefetchWorkerService(
queue as never,
queue as never,
categoriesService as never,
redis as never,
{ payload: vi.fn(async () => ({})) } as never,
{} as never,
);
const job = makeJob("prefetch-parts", {
vehicleId: "v1",
categoryId: "c1",
source: "pl24",
fast: true,
});
await expect(
(svc as never as { process: (j: unknown, t?: string) => Promise<void> }).process(job, "tok"),
).rejects.toThrow();
const [when] = job.moveToDelayed.mock.calls[0] as [number];
// Past the UTC rollover…
expect(when).toBeGreaterThan(Date.UTC(2026, 8, 21, 0, 0, 0));
// …and inside the window, not at 03:00 Istanbul like the old deferral.
const h = istanbulHour(when);
expect(h).toBeGreaterThanOrEqual(9);
expect(h).toBeLessThan(18);
});
});

View File

@@ -225,7 +225,18 @@ describe("PrefetchWorkerService — fast lane (lifo) + backlog gating", () => {
}); });
describe("checkSourceRate — per-source rate limit", () => { describe("checkSourceRate — per-source rate limit", () => {
type CSR = { checkSourceRate: (s: string) => Promise<void> }; type CSR = { checkSourceRate: (s: string, lane?: "main" | "fast") => Promise<void> };
// pl24's MAIN lane is parked whenever background backfill is off, which is
// the default — so these ceiling tests turn it on explicitly to exercise the
// rate limiter itself rather than the backfill gate (covered separately).
beforeEach(() => {
process.env.PL24_BACKFILL_ENABLED = "true";
process.env.PL24_TR_DISABLED = undefined;
});
afterEach(() => {
process.env.PL24_BACKFILL_ENABLED = undefined;
});
it("passes when under the source ceiling", async () => { it("passes when under the source ceiling", async () => {
const { service, redis } = makeDeps({ waiting: 0, limitResults: [] }); const { service, redis } = makeDeps({ waiting: 0, limitResults: [] });
@@ -241,6 +252,24 @@ describe("PrefetchWorkerService — fast lane (lifo) + backlog gating", () => {
}); });
}); });
it("pl24 MAIN lane parkta iken hiç sayaç harcamaz (varsayılan)", async () => {
process.env.PL24_BACKFILL_ENABLED = undefined;
const { service, redis } = makeDeps({ waiting: 0, limitResults: [] });
await expect((service as never as CSR).checkSourceRate("pl24", "main")).rejects.toMatchObject(
{ cause: "source-rate" },
);
expect(redis.incr).not.toHaveBeenCalled();
});
it("kullanıcı (fast) şeridi backfill anahtarından etkilenmez", async () => {
process.env.PL24_BACKFILL_ENABLED = undefined;
const { service, redis } = makeDeps({ waiting: 0, limitResults: [] });
redis.incr.mockResolvedValueOnce(1);
await expect(
(service as never as CSR).checkSourceRate("pl24", "fast"),
).resolves.toBeUndefined();
});
it("is unlimited (no counter) for a source without a configured ceiling", async () => { it("is unlimited (no counter) for a source without a configured ceiling", async () => {
const { service, redis } = makeDeps({ waiting: 0, limitResults: [] }); const { service, redis } = makeDeps({ waiting: 0, limitResults: [] });
await (service as never as CSR).checkSourceRate("unknown-source"); await (service as never as CSR).checkSourceRate("unknown-source");
@@ -339,3 +368,36 @@ describe("PrefetchWorkerService — PL24 derinlik tavanı", () => {
expect(fn("parts-catalogs", true)).toBeGreaterThan(1); expect(fn("parts-catalogs", true)).toBeGreaterThan(1);
}); });
}); });
/**
* PL24 arka plan backfill anahtarı (Faz 3) + Mitsubishi parça-detay kırpması.
*/
describe("PL24 backfill anahtarı", () => {
const ENV = { ...process.env };
afterEach(() => {
process.env = { ...ENV };
});
it("varsayılan KAPALI — değişken hiç yoksa arka plan akmaz", async () => {
const { __testables } = await import("./prefetch-worker.service");
process.env.PL24_BACKFILL_ENABLED = undefined;
process.env.PL24_TR_DISABLED = undefined;
expect(__testables.isPl24BackfillEnabled()).toBe(false);
});
it("yalnız açık 'true' ile açılır", async () => {
const { __testables } = await import("./prefetch-worker.service");
process.env.PL24_TR_DISABLED = undefined;
process.env.PL24_BACKFILL_ENABLED = "true";
expect(__testables.isPl24BackfillEnabled()).toBe(true);
process.env.PL24_BACKFILL_ENABLED = "1";
expect(__testables.isPl24BackfillEnabled()).toBe(false);
});
it("eski PL24_TR_DISABLED hâlâ kapatabilir (yarım deploy musluğu açamaz)", async () => {
const { __testables } = await import("./prefetch-worker.service");
process.env.PL24_BACKFILL_ENABLED = "true";
process.env.PL24_TR_DISABLED = "true";
expect(__testables.isPl24BackfillEnabled()).toBe(false);
});
});

View File

@@ -10,13 +10,14 @@ import { and, asc, eq, gt, inArray, isNull, notExists, sql } from "drizzle-orm";
import { CategoriesService } from "../categories/categories.service"; import { CategoriesService } from "../categories/categories.service";
import { DATABASE, type Database } from "../database/database.provider"; import { DATABASE, type Database } from "../database/database.provider";
import { categories, parts, vehicles } from "../database/schema/core"; import { categories, parts, vehicles } from "../database/schema/core";
import { isPl24LeafNode } from "../integrations/pl24/pl24-tree"; import { isPl24LeafNode, isPl24PartDetailNode } from "../integrations/pl24/pl24-tree";
import { PostHogService } from "../posthog/posthog.service"; import { PostHogService } from "../posthog/posthog.service";
import { RedisService } from "../redis/redis.service"; import { RedisService } from "../redis/redis.service";
import { QUEUE_NAMES, getBullConnection } from "./bull.config"; import { QUEUE_NAMES, getBullConnection } from "./bull.config";
import { backfillContext } from "./prefetch-context"; import { backfillContext } from "./prefetch-context";
import { import {
RateLimitError, RateLimitError,
alignToWindow,
checkCooldown, checkCooldown,
checkTimeWindow, checkTimeWindow,
initProgress, initProgress,
@@ -100,6 +101,25 @@ const EST_JOBS_PER_VEHICLE = Number(process.env.PREFETCH_EST_JOBS_PER_VEHICLE) |
*/ */
const DAILY_FAST_RESERVE = 0.2; const DAILY_FAST_RESERVE = 0.2;
/** Only these decode sources have catalogs worth prefetching. */ /** Only these decode sources have catalogs worth prefetching. */
/**
* Whether the PL24 *background* backfill lane may run. The user-triggered fast
* lane is never gated by this.
*
* Default is OFF. Two PL24 accounts were banned while bulk background load ran
* against them, so the background lane has to be switched on deliberately and
* watched, never inherited from an unset variable.
*
* `PL24_BACKFILL_ENABLED` replaces the old `PL24_TR_DISABLED`, whose name said
* "the tr account is dead" while its actual job was "keep bulk load off the one
* surviving account". The old variable is still honoured so a half-applied
* deploy cannot silently open the tap: it can only keep the lane closed.
*/
function isPl24BackfillEnabled(): boolean {
if (process.env.PL24_BACKFILL_ENABLED !== "true") return false;
// Legacy kill switch still wins while it is explicitly set.
return process.env.PL24_TR_DISABLED !== "true";
}
const BACKFILL_SOURCES = ["pl24", "emex", "parts-catalogs"]; const BACKFILL_SOURCES = ["pl24", "emex", "parts-catalogs"];
/** Redis key holding the rolling rescan cursor (last createdAt seen). */ /** Redis key holding the rolling rescan cursor (last createdAt seen). */
const BACKFILL_CURSOR_KEY = "prefetch:backfill:cursor"; const BACKFILL_CURSOR_KEY = "prefetch:backfill:cursor";
@@ -203,7 +223,7 @@ function jitter(ms: number): number {
} }
/** Test-only surface for the pure helpers above. */ /** Test-only surface for the pure helpers above. */
export const __testables = { maxDepthFor, jitter }; export const __testables = { maxDepthFor, jitter, isPl24BackfillEnabled };
// ── Phase-1 residue exclusion ── // ── Phase-1 residue exclusion ──
/** /**
@@ -354,12 +374,30 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy {
job.name === "prefetch-parts") job.name === "prefetch-parts")
) { ) {
const lane = (job.data as { fast?: boolean }).fast ? "fast" : "main"; const lane = (job.data as { fast?: boolean }).fast ? "fast" : "main";
// GATE ORDER IS LOAD-BEARING — cheapest first, and every gate that can
// reject the job must run BEFORE any counter is debited.
//
// The business-hours window and the cooldown used to sit at the top of
// each handler, i.e. AFTER the per-minute counter and the daily budget
// had already been charged. So a job that woke outside the window paid a
// budget unit to do nothing. Combined with a deferral target of "next UTC
// midnight" (= 03:00 Europe/Istanbul, six hours before a 09:00 window
// opens) that closed a loop: the whole daily allowance was burned by
// no-op wake-ups before the window ever opened, so the source never ran
// again. Measured on prod 2026-09-20 — pl24 at 600/600 with 0 catalog
// requests and 0 new categories for the day.
//
// 1. Window: pure clock arithmetic, no I/O, and an out-of-window job can
// never do useful work — so nothing else is worth spending on it.
checkTimeWindow(data.source);
// 2. Cooldown: one Redis TTL read. Pauses the whole worker (see catch).
await checkCooldown(this.redis, data.source);
// 3. Per-minute ceiling.
await this.checkSourceRate(data.source, lane); await this.checkSourceRate(data.source, lane);
// Daily budget AFTER the per-minute gate: a job deferred on the minute // 4. Daily budget last: a job deferred on any gate above never reaches
// ceiling above never reaches here, so rate-limited retries don't inflate // here, so only jobs about to do real work are counted. The lane
// the daily counter — only jobs about to do real work are counted. The // decides which threshold applies (backfill stops at the main limit,
// lane decides which threshold applies (backfill stops at the main limit, // the user's fast lane may use the full budget).
// the user's fast lane may use the full budget).
await this.checkSourceDailyBudget(data.source, lane); await this.checkSourceDailyBudget(data.source, lane);
} }
if (data.source === "parts-catalogs" && PCAT_PACE_MS > 0) { if (data.source === "parts-catalogs" && PCAT_PACE_MS > 0) {
@@ -421,8 +459,8 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy {
const { vehicleId, source, fast = false } = job.data; const { vehicleId, source, fast = false } = job.data;
this.logger.log(`[prefetch] Init for vehicle=${vehicleId}, source=${source}`); this.logger.log(`[prefetch] Init for vehicle=${vehicleId}, source=${source}`);
await checkCooldown(this.redis, source); // Cooldown + time-window are enforced in process() before the daily budget
checkTimeWindow(source); // is debited — see the comment there; re-checking here would be a no-op.
// Already flagged as poison (tree exceeded CATEGORY_CAP on a prior run) — skip. // Already flagged as poison (tree exceeded CATEGORY_CAP on a prior run) — skip.
if (await this.redis.exists(this.poisonKey(vehicleId))) { if (await this.redis.exists(this.poisonKey(vehicleId))) {
@@ -505,7 +543,7 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy {
// inflated progress.total makes the chain never reach "finished". // inflated progress.total makes the chain never reach "finished".
queued += await this.queueCategoryJob(child, vehicleId, source, 1, fast); queued += await this.queueCategoryJob(child, vehicleId, source, 1, fast);
} }
} else if (this.isLeafLinkPath(cat.linkPath, cat.source, cat.hasSubgroups)) { } else if (this.isLeafLinkPath(cat.linkPath, cat.source, cat.hasSubgroups, cat.linkWid)) {
// Leaf — check if parts already fetched // Leaf — check if parts already fetched
const [partCheck] = await this.db const [partCheck] = await this.db
.select({ id: parts.id }) .select({ id: parts.id })
@@ -563,8 +601,8 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy {
const { vehicleId, categoryId, source, depth, fast = false } = job.data; const { vehicleId, categoryId, source, depth, fast = false } = job.data;
this.logger.log(`[prefetch] Children for category=${categoryId}, depth=${depth}`); this.logger.log(`[prefetch] Children for category=${categoryId}, depth=${depth}`);
await checkCooldown(this.redis, source); // Cooldown + time-window are enforced in process() before the daily budget
checkTimeWindow(source); // is debited — see the comment there; re-checking here would be a no-op.
const depthCeiling = maxDepthFor(source, fast); const depthCeiling = maxDepthFor(source, fast);
if (depth >= depthCeiling) { if (depth >= depthCeiling) {
@@ -631,8 +669,8 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy {
const { vehicleId, categoryId, source } = job.data; const { vehicleId, categoryId, source } = job.data;
this.logger.log(`[prefetch] Parts for category=${categoryId}`); this.logger.log(`[prefetch] Parts for category=${categoryId}`);
await checkCooldown(this.redis, source); // Cooldown + time-window are enforced in process() before the daily budget
checkTimeWindow(source); // is debited — see the comment there; re-checking here would be a no-op.
try { try {
await this.categoriesService.getCategoryWithParts(categoryId); await this.categoriesService.getCategoryWithParts(categoryId);
@@ -721,10 +759,9 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy {
// runaway was built. This also preserves the fast-lane reserve for real users. // runaway was built. This also preserves the fast-lane reserve for real users.
const eligible: string[] = []; const eligible: string[] = [];
for (const s of BACKFILL_SOURCES) { for (const s of BACKFILL_SOURCES) {
// PL24_TR_DISABLED: tr is dead on PL24's side and ALL pl24 traffic maps to // Background backfill is opt-in per source; pl24 defaults to OFF so bulk
// the sole surviving de account — keep background backfill off it entirely // load can never burn the one surviving account by accident.
// (fast/user lane still flows) so bulk load can't burn the last account. if (s === "pl24" && !isPl24BackfillEnabled()) continue;
if (s === "pl24" && process.env.PL24_TR_DISABLED === "true") continue;
if (await this.redis.exists(`prefetch:activity:${s}`)) continue; if (await this.redis.exists(`prefetch:activity:${s}`)) continue;
if (cfg.businessHoursOnly !== false && !isWithinTimeWindow(s)) continue; if (cfg.businessHoursOnly !== false && !isWithinTimeWindow(s)) continue;
const mainLimit = this.dailyMainLimit(s); const mainLimit = this.dailyMainLimit(s);
@@ -875,6 +912,7 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy {
cat: { cat: {
id: string; id: string;
linkPath: string | null; linkPath: string | null;
linkWid?: string | null;
source: string; source: string;
unavailable: boolean; unavailable: boolean;
hasSubgroups?: boolean | null; hasSubgroups?: boolean | null;
@@ -886,7 +924,18 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy {
): Promise<number> { ): Promise<number> {
if (cat.unavailable) return 0; if (cat.unavailable) return 0;
if (this.isLeafLinkPath(cat.linkPath, cat.source, cat.hasSubgroups)) { // A per-part detail node (Mitsubishi `partInfoTable` /details/vinpartinfo) is
// neither a group nor a listing: its parent's response already carried the
// part. Queueing it costs one upstream request and returns nothing. Prod had
// 19,576 of these, with 2 parts between them.
if (
cat.source === "pl24" &&
isPl24PartDetailNode({ linkPath: cat.linkPath, linkWid: cat.linkWid })
) {
return 0;
}
if (this.isLeafLinkPath(cat.linkPath, cat.source, cat.hasSubgroups, cat.linkWid)) {
// Leaf — check if already has parts // Leaf — check if already has parts
const [partCheck] = await this.db const [partCheck] = await this.db
.select({ id: parts.id }) .select({ id: parts.id })
@@ -951,6 +1000,7 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy {
linkPath: string | null, linkPath: string | null,
source: string, source: string,
hasSubgroups?: boolean | null, hasSubgroups?: boolean | null,
linkWid?: string | null,
): boolean { ): boolean {
if (!linkPath) return false; if (!linkPath) return false;
// EMEX: Vehicle.aspx group nodes are parents to drill; Unit.aspx (hierarchical // EMEX: Vehicle.aspx group nodes are parents to drill; Unit.aspx (hierarchical
@@ -969,8 +1019,13 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy {
// PL24: one shared classifier (integrations/pl24/pl24-tree). The old inline // PL24: one shared classifier (integrations/pl24/pl24-tree). The old inline
// list was case-sensitive, so p5psa/p5volvo's camelCase `/details/vin/ // list was case-sensitive, so p5psa/p5volvo's camelCase `/details/vin/
// bomDetails` was never recognised as a leaf and its parts were never // bomDetails` was never recognised as a leaf and its parts were never
// prefetched (still true today for Mitsubishi/Fiat/Renault). // prefetched.
return isPl24LeafNode({ linkPath, hasSubgroups }); // `linkWid` is passed through on purpose: it is the reliable cross-brand
// marker and the read path has always used it, but this queueing path used
// to drop it and fall back to path matching alone. That is why Mitsubishi's
// `detailsTable` parts list was queued as a group here even after the shared
// classifier learned about it.
return isPl24LeafNode({ linkPath, hasSubgroups, linkWid });
} }
/** /**
@@ -1167,10 +1222,10 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy {
// (observed: fresh BMW init couldn't get a single pcat slot). Totals per // (observed: fresh BMW init couldn't get a single pcat slot). Totals per
// source stay the same as before, so upstream load is unchanged. // source stay the same as before, so upstream load is unchanged.
const total = SOURCE_RATE_MAX[source] ?? 0; const total = SOURCE_RATE_MAX[source] ?? 0;
// PL24_TR_DISABLED: park already-queued pl24 MAIN-lane jobs (long defer, no // Park already-queued pl24 MAIN-lane jobs while background backfill is off
// attempt consumed) — the eligibility scan stops producing new ones, this // (long defer, no attempt consumed) — the eligibility scan stops producing
// stops the existing backlog from draining through the de account. // new ones, this stops an existing backlog from draining through the account.
if (source === "pl24" && lane === "main" && process.env.PL24_TR_DISABLED === "true") { if (source === "pl24" && lane === "main" && !isPl24BackfillEnabled()) {
throw new RateLimitError(15 * 60_000, "source-rate"); throw new RateLimitError(15 * 60_000, "source-rate");
} }
if (total <= 0) return; if (total <= 0) return;
@@ -1214,11 +1269,17 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy {
// deferred job wakes in the SAME millisecond (observed: 11495 jobs all at // deferred job wakes in the SAME millisecond (observed: 11495 jobs all at
// 00:00:01 UTC) — the promotion lands as one burst and the pressure signal // 00:00:01 UTC) — the promotion lands as one burst and the pressure signal
// flaps. Other sources keep flowing (per-job defer, not a worker pause). // flaps. Other sources keep flowing (per-job defer, not a worker pause).
const msLeft = dayMs - (now % dayMs) + 1000 + Math.floor(Math.random() * 45 * 60_000); const rollover = now + dayMs - (now % dayMs) + 1000 + Math.floor(Math.random() * 45 * 60_000);
// Land the retry INSIDE the source's scrape window. The UTC rollover alone
// is 03:00 Europe/Istanbul, so with a 09:00 window every deferred job woke
// six hours early, failed the window check and was deferred again — the
// other half of the deadlock fixed in process(). alignToWindow is a no-op
// when no window is configured (the default).
const msLeft = Math.max(1000, alignToWindow(source, rollover) - now);
if (n === limit) { if (n === limit) {
this.logger.warn( this.logger.warn(
`[prefetch] ${source} daily budget hit (lane=${lane}, ${n}/${limit} of ${max}) — ` + `[prefetch] ${source} daily budget hit (lane=${lane}, ${n}/${limit} of ${max}) — ` +
`deferring ~${Math.round(msLeft / 3_600_000)}h until the window rolls`, `deferring ~${Math.round(msLeft / 3_600_000)}h to the next in-window slot`,
); );
} }
throw new RateLimitError(msLeft, "source-rate"); throw new RateLimitError(msLeft, "source-rate");

View File

@@ -0,0 +1,243 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { RpartstoreRateLimitError } from "../../integrations/rpartstore/rpartstore.session";
import type { RpartstoreDecoded } from "../../integrations/rpartstore/rpartstore.types";
import {
type RedisLike,
istanbulDayKey,
processRpartstoreDecode,
redisTokenStore,
rpartstoreDailyKey,
} from "./rpartstore-decode.processor";
/** In-memory ioredis stand-in covering the surface the processor uses. */
function fakeRedis(): RedisLike & { store: Map<string, string>; ttls: Map<string, number> } {
const store = new Map<string, string>();
const ttls = new Map<string, number>();
return {
store,
ttls,
async get(k) {
return store.get(k) ?? null;
},
async set(k, v, _mode, ttl) {
store.set(k, v);
ttls.set(k, ttl);
},
async del(k) {
store.delete(k);
},
async incr(k) {
const n = Number(store.get(k) ?? 0) + 1;
store.set(k, String(n));
return n;
},
async expire(k, s) {
ttls.set(k, s);
},
};
}
/** Chainable drizzle mock: records every `.set()` payload and answers selects with `rows`. */
function fakeDb(rows: unknown[] = []) {
const sets: Record<string, unknown>[] = [];
const update = vi.fn(() => ({
set: vi.fn((payload: Record<string, unknown>) => {
sets.push(payload);
return { where: vi.fn().mockResolvedValue(undefined) };
}),
}));
const select = vi.fn(() => ({
from: vi.fn(() => ({ where: vi.fn().mockResolvedValue(rows) })),
}));
return { db: { update, select } as any, sets };
}
const decodedKadjar: RpartstoreDecoded = {
brandName: "Renault",
model: "Kadjar (HFE)",
modelCode: "HFE",
familyCode: "XFE",
modelYear: "2015",
engine: "1.5 DCI DİZEL MOTOR [K9K]",
gearbox: "DC4",
energyType: "MOTORIN",
manufacturingDate: "2015-07-28",
raw: {
catalogSource: "DATAHUB",
vin: "VF1RFE00653633190",
vehicleKey: "VF1RFE00653633190",
model: "Kadjar (HFE)",
vehicleBrand: "RENAULT",
country: "TR",
},
};
const job = (vin: string) =>
({ id: `j-${vin}`, data: { vin }, attemptsMade: 0, opts: { attempts: 1 } }) as any;
const NOW = new Date("2026-09-25T20:30:00+03:00");
describe("processRpartstoreDecode", () => {
beforeEach(() => {
process.env.RPARTSTORE_ENABLED = "true";
vi.spyOn(console, "log").mockImplementation(() => undefined);
vi.spyOn(console, "warn").mockImplementation(() => undefined);
vi.spyOn(console, "error").mockImplementation(() => undefined);
});
afterEach(() => {
process.env.RPARTSTORE_ENABLED = "false";
vi.restoreAllMocks();
});
it("is a no-op when RPARTSTORE_ENABLED is not true", async () => {
process.env.RPARTSTORE_ENABLED = "false";
const { db, sets } = fakeDb();
const searchVin = vi.fn();
const r = await processRpartstoreDecode(job("VF1X"), {
db,
redis: fakeRedis(),
decoder: () => ({ searchVin }),
dailyCap: 10,
});
expect(r.skipped).toBe(true);
expect(searchVin).not.toHaveBeenCalled();
expect(sets).toEqual([]);
});
it("decodes, matches the PL24 catalog vehicle and records the decoded row", async () => {
const { db, sets } = fakeDb([
{ id: "cv-kadjar", model: "KADJAR", categoryCount: 44 },
{ id: "cv-kadjar-cn", model: "KADJAR ÇİN", categoryCount: 22 },
]);
const redis = fakeRedis();
const searchVin = vi.fn().mockResolvedValue(decodedKadjar);
const r = await processRpartstoreDecode(job("VF1RFE00653633190"), {
db,
redis,
decoder: () => ({ searchVin }),
dailyCap: 10,
now: () => NOW,
});
expect(r).toEqual({ status: "decoded", catalogVehicleId: "cv-kadjar" });
expect(searchVin).toHaveBeenCalledWith("VF1RFE00653633190");
// one cap unit reserved on the Istanbul day, with an expiry
expect(redis.store.get(rpartstoreDailyKey(NOW))).toBe("1");
expect(redis.ttls.get(rpartstoreDailyKey(NOW))).toBeGreaterThan(0);
const final = sets.at(-1)!;
expect(final).toMatchObject({
status: "decoded",
brandName: "Renault",
model: "Kadjar (HFE)",
modelCode: "HFE",
modelYear: "2015",
catalogVehicleId: "cv-kadjar",
});
});
it("records not_found when RPartStore has no vehicle for the VIN", async () => {
const { db, sets } = fakeDb();
const r = await processRpartstoreDecode(job("VF1NOPE"), {
db,
redis: fakeRedis(),
decoder: () => ({ searchVin: vi.fn().mockResolvedValue(null) }),
dailyCap: 10,
});
expect(r.status).toBe("not_found");
expect(sets.at(-1)).toMatchObject({ status: "not_found" });
});
it("stops sending once the daily cap is spent and marks the row capped", async () => {
const redis = fakeRedis();
const key = rpartstoreDailyKey(NOW);
redis.store.set(key, "10"); // ten searches already sent today
const { db, sets } = fakeDb();
const searchVin = vi.fn().mockResolvedValue(decodedKadjar);
const r = await processRpartstoreDecode(job("VF1CAP"), {
db,
redis,
decoder: () => ({ searchVin }),
dailyCap: 10,
now: () => NOW,
});
expect(r.status).toBe("capped");
expect(searchVin).not.toHaveBeenCalled();
expect(sets.at(-1)).toMatchObject({ status: "capped" });
});
it("allows exactly `dailyCap` searches per day", async () => {
const redis = fakeRedis();
const searchVin = vi.fn().mockResolvedValue(null);
const statuses: string[] = [];
for (let i = 0; i < 12; i += 1) {
const { db } = fakeDb();
const r = await processRpartstoreDecode(job(`VF1${i}`), {
db,
redis,
decoder: () => ({ searchVin }),
dailyCap: 10,
now: () => NOW,
});
statuses.push(r.status);
}
expect(searchVin).toHaveBeenCalledTimes(10);
expect(statuses.filter((s) => s === "capped")).toHaveLength(2);
});
it("waits retryAfterSeconds and retries once on a short-term rate limit without a second reservation", async () => {
const redis = fakeRedis();
const { db } = fakeDb();
const searchVin = vi
.fn()
.mockRejectedValueOnce(new RpartstoreRateLimitError(10, "SHORT_TERM"))
.mockResolvedValueOnce(null);
const sleep = vi.fn().mockResolvedValue(undefined);
const r = await processRpartstoreDecode(job("VF1RL"), {
db,
redis,
decoder: () => ({ searchVin }),
dailyCap: 10,
now: () => NOW,
sleep,
});
expect(r.status).toBe("not_found");
expect(searchVin).toHaveBeenCalledTimes(2);
expect(sleep).toHaveBeenCalledWith(10_500);
expect(redis.store.get(rpartstoreDailyKey(NOW))).toBe("1");
});
it("marks the row failed and rethrows on an infrastructure error (single attempt)", async () => {
const { db, sets } = fakeDb();
await expect(
processRpartstoreDecode(job("VF1ERR"), {
db,
redis: fakeRedis(),
decoder: () => {
throw new Error("RPARTSTORE_USER / RPARTSTORE_PASS are not configured");
},
dailyCap: 10,
}),
).rejects.toThrow(/not configured/);
expect(sets.at(-1)).toMatchObject({ status: "failed" });
});
});
describe("istanbulDayKey", () => {
it("counts the cap against the Istanbul calendar day, not UTC", () => {
// 23:30 UTC on the 25th is already the 26th in Istanbul (UTC+3).
expect(istanbulDayKey(new Date("2026-09-25T23:30:00Z"))).toBe("2026-09-26");
expect(istanbulDayKey(new Date("2026-09-25T20:59:00Z"))).toBe("2026-09-25");
});
});
describe("redisTokenStore", () => {
it("round-trips a token with its ttl and ignores garbage", async () => {
const redis = fakeRedis();
const store = redisTokenStore(redis);
await store.set({ accessToken: "t", expiresAt: 1, subject: "s" }, 120);
expect(await store.get()).toEqual({ accessToken: "t", expiresAt: 1, subject: "s" });
expect(redis.ttls.get("rpartstore:token")).toBe(120);
redis.store.set("rpartstore:token", "{not json");
expect(await store.get()).toBeNull();
await store.clear();
expect(redis.store.has("rpartstore:token")).toBe(false);
});
});

View File

@@ -0,0 +1,252 @@
import { Job } from "bullmq";
import { and, eq, sql } from "drizzle-orm";
import { PostgresJsDatabase } from "drizzle-orm/postgres-js";
import { catalogVehicles, rpartstoreDecodes } from "../../database/schema/core";
import type { RpartstoreToken } from "../../integrations/rpartstore/rpartstore.auth";
import { RpartstoreClient, type TokenStore } from "../../integrations/rpartstore/rpartstore.client";
import {
type RpartstoreCatalogCandidate,
pickRpartstoreCatalogMatch,
} from "../../integrations/rpartstore/rpartstore.matcher";
import { RpartstoreRateLimitError } from "../../integrations/rpartstore/rpartstore.session";
import type { RpartstoreDecoded } from "../../integrations/rpartstore/rpartstore.types";
type Database = PostgresJsDatabase<Record<string, unknown>>;
export interface RpartstoreDecodeJobData {
vin: string;
}
/** Minimal ioredis surface the processor needs (also satisfied by RedisService.getClient()). */
export interface RedisLike {
get(key: string): Promise<string | null>;
set(key: string, value: string, mode: "EX", ttlSeconds: number): Promise<unknown>;
del(key: string): Promise<unknown>;
incr(key: string): Promise<number>;
expire(key: string, seconds: number): Promise<unknown>;
}
export interface RpartstoreDecoder {
searchVin(vin: string): Promise<RpartstoreDecoded | null>;
}
export interface RpartstoreProcessorDeps {
db: Database;
redis: RedisLike;
/** Built lazily so a missing credential only fails the job, not worker boot. */
decoder: () => RpartstoreDecoder;
/** Hard cap on VIN searches sent to RPartStore per Istanbul calendar day. */
dailyCap: number;
now?: () => Date;
sleep?: (ms: number) => Promise<void>;
}
export const RPARTSTORE_TOKEN_KEY = "rpartstore:token";
export const RPARTSTORE_DAILY_KEY_PREFIX = "rpartstore:daily:";
/** Counter keys live two days so a late-night job never sees a vanished key. */
const DAILY_KEY_TTL_SECONDS = 2 * 24 * 60 * 60;
/** "YYYY-MM-DD" in Europe/Istanbul — the day the cap is counted against. */
export function istanbulDayKey(date: Date): string {
const parts = new Intl.DateTimeFormat("en-CA", {
timeZone: "Europe/Istanbul",
year: "numeric",
month: "2-digit",
day: "2-digit",
}).formatToParts(date);
const get = (t: string): string => parts.find((p) => p.type === t)?.value ?? "";
return `${get("year")}-${get("month")}-${get("day")}`;
}
export const rpartstoreDailyKey = (date: Date): string =>
`${RPARTSTORE_DAILY_KEY_PREFIX}${istanbulDayKey(date)}`;
/** Redis-backed cache for the 1 h Okta access token (shared by worker restarts). */
export function redisTokenStore(redis: RedisLike): TokenStore {
return {
async get() {
const raw = await redis.get(RPARTSTORE_TOKEN_KEY);
if (!raw) return null;
try {
const t = JSON.parse(raw) as RpartstoreToken;
return t.accessToken && t.expiresAt ? t : null;
} catch {
return null;
}
},
async set(token, ttlSeconds) {
await redis.set(RPARTSTORE_TOKEN_KEY, JSON.stringify(token), "EX", ttlSeconds);
},
async clear() {
await redis.del(RPARTSTORE_TOKEN_KEY);
},
};
}
export function buildRpartstoreClient(redis: RedisLike): RpartstoreClient {
const username = process.env.RPARTSTORE_USER;
const password = process.env.RPARTSTORE_PASS;
if (!username || !password) {
throw new Error("RPARTSTORE_USER / RPARTSTORE_PASS are not configured");
}
return new RpartstoreClient({
username,
password,
tokenStore: redisTokenStore(redis),
brokerUrl: process.env.RPARTSTORE_BROKER_URL || undefined,
appVersion: process.env.RPARTSTORE_APP_VERSION || undefined,
logger: { log: (m) => console.log(m), warn: (m) => console.warn(m) },
});
}
export function rpartstoreDailyCapFromEnv(): number {
const n = Number(process.env.RPARTSTORE_DAILY_CAP);
return Number.isFinite(n) && n >= 0 ? Math.floor(n) : 10;
}
/**
* RPartStore decode processor (worker, concurrency 1, ≥6 s between jobs via the
* queue limiter — RPartStore allows 2 VIN searches per 10 s).
*
* Reserves one unit of the daily cap BEFORE sending anything: `INCR` on the
* Istanbul-day key; when the reservation lands above the cap the row is marked
* `capped` and nothing is sent, so at most `dailyCap` searches reach RPartStore
* per day even under concurrent enqueues. A rate-limited search waits
* `retryAfterSeconds` and is retried once without a second reservation.
*
* Outcomes recorded in `rpartstore_decodes`: decoded (with an optional PL24
* catalog_vehicle match), not_found (definitive), capped (retry tomorrow),
* failed (infra/auth error — user-retriable after 24 h).
*/
export async function processRpartstoreDecode(
job: Job<RpartstoreDecodeJobData>,
deps: RpartstoreProcessorDeps,
): Promise<{ status: string; catalogVehicleId: string | null; skipped?: boolean }> {
const { vin } = job.data;
const { db, redis } = deps;
const now = deps.now ?? (() => new Date());
const sleep = deps.sleep ?? ((ms: number) => new Promise<void>((r) => setTimeout(r, ms)));
if (process.env.RPARTSTORE_ENABLED !== "true") {
console.log(`[rpartstore-decode] disabled (RPARTSTORE_ENABLED!=true) — job ${job.id} no-op`);
return { status: "pending", catalogVehicleId: null, skipped: true };
}
await db
.update(rpartstoreDecodes)
.set({ attempts: sql`${rpartstoreDecodes.attempts} + 1`, updatedAt: now() })
.where(eq(rpartstoreDecodes.vin, vin));
// Daily cap reservation.
const dayKey = rpartstoreDailyKey(now());
const reserved = await redis.incr(dayKey);
if (reserved === 1) await redis.expire(dayKey, DAILY_KEY_TTL_SECONDS);
if (reserved > deps.dailyCap) {
await db
.update(rpartstoreDecodes)
.set({ status: "capped", updatedAt: now() })
.where(eq(rpartstoreDecodes.vin, vin));
console.warn(
`[rpartstore-decode] ${vin} → capped (${reserved - 1}/${deps.dailyCap} searches used on ${dayKey})`,
);
return { status: "capped", catalogVehicleId: null };
}
console.log(
`[rpartstore-decode] job ${job.id} decoding ${vin} (${reserved}/${deps.dailyCap} today)`,
);
try {
const decoder = deps.decoder();
let decoded: RpartstoreDecoded | null;
try {
decoded = await decoder.searchVin(vin);
} catch (err) {
if (!(err instanceof RpartstoreRateLimitError)) throw err;
const waitMs = Math.min(Math.max(err.retryAfterSeconds, 1), 30) * 1000 + 500;
console.warn(
`[rpartstore-decode] ${vin} rate-limited (${err.limitType}); retrying in ${waitMs}ms`,
);
await sleep(waitMs);
decoded = await decoder.searchVin(vin);
}
if (!decoded) {
await db
.update(rpartstoreDecodes)
.set({ status: "not_found", updatedAt: now() })
.where(eq(rpartstoreDecodes.vin, vin));
console.log(`[rpartstore-decode] ${vin} → not_found`);
return { status: "not_found", catalogVehicleId: null };
}
const catalogVehicleId = await matchCatalogVehicle(db, decoded);
await db
.update(rpartstoreDecodes)
.set({
status: "decoded",
brandName: decoded.brandName,
model: decoded.model,
modelCode: decoded.modelCode,
familyCode: decoded.familyCode,
modelYear: decoded.modelYear,
engine: decoded.engine,
gearbox: decoded.gearbox,
energyType: decoded.energyType,
manufacturingDate: decoded.manufacturingDate,
catalogVehicleId,
raw: decoded.raw,
decodedAt: now(),
updatedAt: now(),
})
.where(eq(rpartstoreDecodes.vin, vin));
console.log(
`[rpartstore-decode] ${vin} → decoded ${decoded.brandName} "${decoded.model ?? "?"}" ${decoded.modelYear ?? ""} catalog_vehicle=${catalogVehicleId ?? "null"}`,
);
return { status: "decoded", catalogVehicleId };
} catch (error) {
const message = error instanceof Error ? error.message : String(error);
console.error(`[rpartstore-decode] job ${job.id} failed for ${vin}: ${message}`);
const isLastAttempt = job.attemptsMade + 1 >= (job.opts.attempts ?? 1);
if (isLastAttempt) {
await db
.update(rpartstoreDecodes)
.set({ status: "failed", updatedAt: now() })
.where(eq(rpartstoreDecodes.vin, vin));
}
throw error;
}
}
/** Match within the decoded brand first, then the sister brand (Renault ⇄ Dacia
* share platforms and TR badges differ from the WMI). */
async function matchCatalogVehicle(
db: Database,
decoded: RpartstoreDecoded,
): Promise<string | null> {
if (!decoded.model) return null;
const sister = decoded.brandName.toLowerCase() === "dacia" ? "renault" : "dacia";
for (const brand of [decoded.brandName.toLowerCase(), sister]) {
const rows = await db
.select({
id: catalogVehicles.id,
model: catalogVehicles.model,
categoryCount: sql<number>`(
SELECT count(*)::int FROM categories
WHERE categories.catalog_vehicle_id = ${catalogVehicles.id}
)`,
})
.from(catalogVehicles)
.where(
and(
sql`lower(${catalogVehicles.brandName}) = ${brand}`,
eq(catalogVehicles.source, "pl24"),
),
);
const id = pickRpartstoreCatalogMatch(decoded.model, rows as RpartstoreCatalogCandidate[]);
if (id) return id;
}
return null;
}

View File

@@ -0,0 +1,32 @@
import { Provider } from "@nestjs/common";
import { type JobsOptions, Queue } from "bullmq";
import { QUEUE_NAMES, getBullConnection, getBullTelemetry } from "../bull.config";
export const RPARTSTORE_DECODE_QUEUE = "RPARTSTORE_DECODE_QUEUE";
/**
* `attempts: 1` — the processor records a definitive outcome per job
* (decoded / not_found / capped / failed) and the daily cap must never be
* burned by automatic re-runs. A `failed` or `capped` row is re-enqueued by the
* decode path itself after 24 h (see VehiclesService.tryRpartstoreFallback).
*/
export const RPARTSTORE_DECODE_JOB_OPTIONS: JobsOptions = {
attempts: 1,
removeOnComplete: { count: 50 },
removeOnFail: { count: 100 },
};
/** RPartStore (Renault/Dacia) VIN-decode fallback queue. One shared dealer
* account → the worker consumes it with concurrency 1 and a 1-job-per-6 s
* limiter (RPartStore allows 2 searches per 10 s). Jobs carry `{ vin }`. */
export const RpartstoreDecodeQueueProvider: Provider = {
provide: RPARTSTORE_DECODE_QUEUE,
useFactory: () => {
const telemetry = getBullTelemetry();
return new Queue(QUEUE_NAMES.RPARTSTORE_DECODE, {
connection: getBullConnection(),
...(telemetry ? { telemetry } : {}),
defaultJobOptions: RPARTSTORE_DECODE_JOB_OPTIONS,
});
},
};

View File

@@ -1,6 +1,11 @@
import { BadRequestException, ForbiddenException, NotFoundException } from "@nestjs/common"; import { BadRequestException, ForbiddenException, NotFoundException } from "@nestjs/common";
import { beforeEach, describe, expect, it, vi } from "vitest"; import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { catalogVehicles, queryLogs, vinpinDecodes } from "../database/schema/core"; import {
catalogVehicles,
queryLogs,
rpartstoreDecodes,
vinpinDecodes,
} from "../database/schema/core";
import { VehiclesService } from "./vehicles.service"; import { VehiclesService } from "./vehicles.service";
vi.mock("@sase/shared", () => ({ vi.mock("@sase/shared", () => ({
@@ -93,10 +98,15 @@ function createService(dbOrOverrides: any = {}) {
add: vi.fn().mockResolvedValue(undefined), add: vi.fn().mockResolvedValue(undefined),
}; };
const rpartstoreQueue = {
add: vi.fn().mockResolvedValue(undefined),
};
const service = new VehiclesService( const service = new VehiclesService(
db as any, db as any,
prefetchQueue as any, prefetchQueue as any,
vinpinQueue as any, vinpinQueue as any,
rpartstoreQueue as any,
corgiService as any, corgiService as any,
pl24Service as any, pl24Service as any,
vinApiService as any, vinApiService as any,
@@ -115,6 +125,7 @@ function createService(dbOrOverrides: any = {}) {
partsCatalogsService, partsCatalogsService,
redisService, redisService,
vinpinQueue, vinpinQueue,
rpartstoreQueue,
}; };
} }
@@ -340,9 +351,7 @@ describe("VehiclesService", () => {
selectChain.limit = vi selectChain.limit = vi
.fn() .fn()
.mockImplementation(() => .mockImplementation(() =>
currentTable === vinpinDecodes currentTable === vinpinDecodes ? [{ vin: "NM435600006123456", status: "pending" }] : [],
? [{ vin: "NM435600006123456", status: "pending" }]
: [],
); );
const insertChain = { const insertChain = {
values: vi.fn().mockReturnThis(), values: vi.fn().mockReturnThis(),
@@ -992,3 +1001,168 @@ describe("VehiclesService", () => {
}); });
}); });
}); });
describe("VehiclesService — RPartStore fallback (Renault/Dacia, after the pcat/PL24/emex race)", () => {
const VIN = "VF1RFE00653633190";
function noCatalogDb() {
const selectChain = {
from: vi.fn().mockReturnThis(),
where: vi.fn().mockReturnThis(),
limit: vi.fn().mockReturnValue([]),
};
const insertChain = {
values: vi.fn().mockReturnThis(),
returning: vi.fn().mockReturnValue([]),
onConflictDoNothing: vi.fn().mockReturnThis(),
};
return {
select: vi.fn().mockReturnValue(selectChain),
insert: vi.fn().mockReturnValue(insertChain),
};
}
function renaultIdent(svc: ReturnType<typeof createService>) {
svc.corgiService.decodeVin.mockReturnValue({
isKnown: true,
brandName: "Renault",
modelYear: 2015,
});
svc.vinApiService.decodeVin.mockResolvedValue({
make: "RENAULT",
model: "Kadjar",
modelYear: "2015",
});
}
let prevRparts: string | undefined;
let prevVinpin: string | undefined;
beforeEach(() => {
prevRparts = process.env.RPARTSTORE_ENABLED;
prevVinpin = process.env.VINPIN_ENABLED;
process.env.VINPIN_ENABLED = "false";
vi.mocked(isValidVin).mockReturnValue(true);
});
afterEach(() => {
process.env.RPARTSTORE_ENABLED = prevRparts ?? "false";
process.env.VINPIN_ENABLED = prevVinpin ?? "false";
});
it("[off] leaves the no-catalog path untouched and never touches the queue", async () => {
process.env.RPARTSTORE_ENABLED = "false";
const db = noCatalogDb();
const svc = createService(db);
renaultIdent(svc);
const result: any = await svc.service.decodeVin(VIN, "u1");
expect(result).toMatchObject({ noCatalog: { brandName: "Renault" }, vin: VIN });
expect(svc.rpartstoreQueue.add).not.toHaveBeenCalled();
expect(db.insert.mock.calls.some((c: unknown[]) => c[0] === rpartstoreDecodes)).toBe(false);
});
it("[on] records a pending row, enqueues a day-scoped job and answers `decoding` for an unseen Renault VIN", async () => {
process.env.RPARTSTORE_ENABLED = "true";
const db = noCatalogDb();
const svc = createService(db);
renaultIdent(svc);
const result: any = await svc.service.decodeVin(VIN, "u1");
expect(result).toMatchObject({ decoding: { vin: VIN }, vin: VIN });
expect(result.decoding.display).toContain("Renault");
expect(svc.rpartstoreQueue.add).toHaveBeenCalledTimes(1);
const [name, data, opts] = svc.rpartstoreQueue.add.mock.calls[0];
expect(name).toBe("rpartstore-decode");
expect(data).toEqual({ vin: VIN });
expect(opts.jobId).toMatch(new RegExp(`^rpartstore-${VIN}-\\d{4}-\\d{2}-\\d{2}$`));
expect(db.insert.mock.calls.some((c: unknown[]) => c[0] === rpartstoreDecodes)).toBe(true);
// Exactly ONE coverage-gap failure row at first sighting.
expect(db.insert.mock.calls.filter((c: unknown[]) => c[0] === queryLogs)).toHaveLength(1);
// Vinpin (disabled) is never consulted.
expect(svc.vinpinQueue.add).not.toHaveBeenCalled();
});
it("[on] does not enqueue a non-Renault VIN", async () => {
process.env.RPARTSTORE_ENABLED = "true";
const db = noCatalogDb();
const svc = createService(db);
svc.corgiService.decodeVin.mockReturnValue({
isKnown: true,
brandName: "Fiat",
modelYear: 2018,
});
svc.vinApiService.decodeVin.mockResolvedValue({
make: "FIAT",
model: "Tipo",
modelYear: "2018",
});
const result: any = await svc.service.decodeVin("NM435600006123456", "u1");
expect(result).toMatchObject({ noCatalog: { brandName: "Fiat" } });
expect(svc.rpartstoreQueue.add).not.toHaveBeenCalled();
});
it("[on] falls through to noCatalog without queueing once today's cap is spent", async () => {
process.env.RPARTSTORE_ENABLED = "true";
process.env.RPARTSTORE_DAILY_CAP = "10";
const db = noCatalogDb();
const svc = createService(db);
renaultIdent(svc);
svc.redisService.get.mockImplementation(async (key: string) =>
key.startsWith("rpartstore:daily:") ? "10" : null,
);
const result: any = await svc.service.decodeVin(VIN, "u1");
expect(result).toMatchObject({ noCatalog: { brandName: "Renault" } });
expect(svc.rpartstoreQueue.add).not.toHaveBeenCalled();
expect(db.insert.mock.calls.some((c: unknown[]) => c[0] === rpartstoreDecodes)).toBe(false);
process.env.RPARTSTORE_DAILY_CAP = "";
});
it("[on] a decoded row with a catalog match answers `catalogVehicle` and logs one success row", async () => {
process.env.RPARTSTORE_ENABLED = "true";
const rpartsRow = {
vin: VIN,
status: "decoded",
catalogVehicleId: "cv-1",
createdAt: new Date(),
updatedAt: new Date(),
};
const cv = { id: "cv-1", brandId: null, brandName: "Renault", model: "KADJAR", year: null };
let selectCalls = 0;
const selectChain = {
from: vi.fn().mockReturnThis(),
where: vi.fn().mockReturnThis(),
limit: vi.fn().mockImplementation(() => {
// Call order inside decodeVin: vehicles lookup → rpartstore row → catalog vehicle.
selectCalls += 1;
if (selectCalls === 2) return [rpartsRow];
if (selectCalls === 3) return [cv];
return [];
}),
};
const insertChain = {
values: vi.fn().mockReturnThis(),
returning: vi.fn().mockReturnValue([]),
onConflictDoNothing: vi.fn().mockReturnThis(),
};
const db = {
select: vi.fn().mockReturnValue(selectChain),
insert: vi.fn().mockReturnValue(insertChain),
};
const svc = createService(db);
renaultIdent(svc);
const result: any = await svc.service.decodeVin(VIN, "u1");
expect(result).toEqual({
catalogVehicle: { id: "cv-1", brandName: "Renault", model: "KADJAR", year: null },
vin: VIN,
});
expect(svc.rpartstoreQueue.add).not.toHaveBeenCalled();
const logRows = db.insert.mock.calls.filter((c: unknown[]) => c[0] === queryLogs);
expect(logRows).toHaveLength(1);
expect(insertChain.values.mock.calls.at(-1)?.[0]).toMatchObject({
source: "rpartstore",
success: true,
});
});
});

View File

@@ -16,6 +16,7 @@ import {
parts, parts,
plans, plans,
queryLogs, queryLogs,
rpartstoreDecodes,
userBrands, userBrands,
userSubscriptions, userSubscriptions,
userVehicles, userVehicles,
@@ -34,10 +35,13 @@ import { PartsCatalogsService } from "../integrations/parts-catalogs/parts-catal
import { PcatCar, PcatVinResult } from "../integrations/parts-catalogs/parts-catalogs.types"; import { PcatCar, PcatVinResult } from "../integrations/parts-catalogs/parts-catalogs.types";
import { PL24Service } from "../integrations/pl24/pl24.service"; import { PL24Service } from "../integrations/pl24/pl24.service";
import { SERVICE_TO_BRAND } from "../integrations/pl24/pl24.types"; import { SERVICE_TO_BRAND } from "../integrations/pl24/pl24.types";
import { isRpartstoreVin } from "../integrations/rpartstore/rpartstore.routing";
import { VinApiService } from "../integrations/vin-api/vin-api.service"; import { VinApiService } from "../integrations/vin-api/vin-api.service";
import { isVinpinBrandAllowed } from "../integrations/vinpin/vinpin.constants"; import { isVinpinBrandAllowed } from "../integrations/vinpin/vinpin.constants";
import { PrefetchSource } from "../jobs/prefetch.types"; import { PrefetchSource } from "../jobs/prefetch.types";
import { istanbulDayKey, rpartstoreDailyKey } from "../jobs/processors/rpartstore-decode.processor";
import { CATALOG_PREFETCH_FAST_QUEUE } from "../jobs/queues/catalog-prefetch.queue"; import { CATALOG_PREFETCH_FAST_QUEUE } from "../jobs/queues/catalog-prefetch.queue";
import { RPARTSTORE_DECODE_QUEUE } from "../jobs/queues/rpartstore-decode.queue";
import { VINPIN_DECODE_QUEUE } from "../jobs/queues/vinpin-decode.queue"; import { VINPIN_DECODE_QUEUE } from "../jobs/queues/vinpin-decode.queue";
import { RedisService } from "../redis/redis.service"; import { RedisService } from "../redis/redis.service";
import { vinCandidateStashKey, vinResolveCacheKeys } from "./vin-cache-keys"; import { vinCandidateStashKey, vinResolveCacheKeys } from "./vin-cache-keys";
@@ -99,6 +103,7 @@ export class VehiclesService {
@Inject(DATABASE) private db: Database, @Inject(DATABASE) private db: Database,
@Inject(CATALOG_PREFETCH_FAST_QUEUE) private prefetchQueue: Queue, @Inject(CATALOG_PREFETCH_FAST_QUEUE) private prefetchQueue: Queue,
@Inject(VINPIN_DECODE_QUEUE) private vinpinQueue: Queue, @Inject(VINPIN_DECODE_QUEUE) private vinpinQueue: Queue,
@Inject(RPARTSTORE_DECODE_QUEUE) private rpartstoreQueue: Queue,
private corgiService: CorgiService, private corgiService: CorgiService,
private pl24Service: PL24Service, private pl24Service: PL24Service,
private vinApiService: VinApiService, private vinApiService: VinApiService,
@@ -191,6 +196,19 @@ export class VehiclesService {
// The fallback logs the coverage-gap failure row exactly once, at // The fallback logs the coverage-gap failure row exactly once, at
// first sighting (job enqueue); not_found/failed/stale fall through // first sighting (job enqueue); not_found/failed/stale fall through
// to the failure log below, unchanged. // to the failure log below, unchanged.
// RPartStore decode-oracle fallback for Renault/Dacia (flag-gated,
// hard daily cap). Ordered AFTER the pcat/PL24/emex race (we only get
// here when it returned nothing) and BEFORE Vinpin, which stays wired
// as a last resort behind its own flag. Same accounting contract as
// the Vinpin fallback below: one failure row at first sighting, polls
// log nothing, a resolved decode logs its own success row.
if (
process.env.RPARTSTORE_ENABLED === "true" &&
isRpartstoreVin(vin, ident.browseBrand)
) {
const rpartsResp = await this.tryRpartstoreFallback(vin, ident, userId, ctx, startTime);
if (rpartsResp) return rpartsResp;
}
if (process.env.VINPIN_ENABLED === "true" && isVinpinBrandAllowed(ident.browseBrand)) { if (process.env.VINPIN_ENABLED === "true" && isVinpinBrandAllowed(ident.browseBrand)) {
const vinpinResp = await this.tryVinpinFallback(vin, ident, userId, ctx, startTime); const vinpinResp = await this.tryVinpinFallback(vin, ident, userId, ctx, startTime);
if (vinpinResp) return vinpinResp; if (vinpinResp) return vinpinResp;
@@ -791,6 +809,154 @@ export class VehiclesService {
return null; return null;
} }
/** `failed` / `capped` RPartStore rows become eligible for one more attempt after this long. */
private static readonly RPARTSTORE_RETRY_AFTER_MS = 24 * 60 * 60 * 1000;
/**
* RPartStore decode-oracle fallback for a Renault/Dacia VIN the race couldn't
* identify (feature-flagged, guarded by the caller on RPARTSTORE_ENABLED +
* `isRpartstoreVin`). Mirrors the Vinpin fallback contract:
*
* - decoded + catalog_vehicle_id → resolve that EXISTING PL24 catalog vehicle
* and return `catalogVehicle` (parts come from PL24, not RPartStore).
* - pending → `{ decoding }` (frontend polls).
* - no row → insert 'pending', enqueue a decode job, return `{ decoding }`.
* - capped / failed older than 24 h → re-enqueue (a fresh cap reservation),
* return `{ decoding }`; younger → null.
* - not_found / decoded-without-match / cap exhausted today → null → caller
* falls through to Vinpin (if enabled) and then noCatalog, unchanged.
*
* The API side never talks to RPartStore; it only checks today's counter so a
* VIN is not queued (and left `pending`) when the daily cap is already spent.
*/
private async tryRpartstoreFallback(
vin: string,
ident: { browseBrand: string | null; display: string },
userId: string,
ctx: ResolveContext,
startTime: number,
// biome-ignore lint/suspicious/noExplicitAny: heterogeneous short-circuit response shapes
): Promise<any | null> {
try {
const [row] = await this.db
.select()
.from(rpartstoreDecodes)
.where(eq(rpartstoreDecodes.vin, vin))
.limit(1);
if (row) {
if (row.status === "decoded" && row.catalogVehicleId) {
const [cv] = await this.db
.select()
.from(catalogVehicles)
.where(eq(catalogVehicles.id, row.catalogVehicleId))
.limit(1);
if (cv) {
if (cv.brandId) await this.checkBrandAccess(userId, cv.brandId);
ctx.timings.rpartstore_resolved = 1;
await this.logQuery(
userId,
vin,
cv.brandId,
"rpartstore",
true,
Date.now() - startTime,
undefined,
ctx.timings,
);
return {
catalogVehicle: {
id: cv.id,
brandName: cv.brandName,
model: cv.model,
year: cv.year,
},
vin,
};
}
return null;
}
if (row.status === "pending") {
ctx.timings.rpartstore_pending = 1;
return { decoding: { vin, display: ident.display }, vin };
}
const retryable = row.status === "capped" || row.status === "failed";
const ageMs = Date.now() - (row.updatedAt ?? row.createdAt).getTime();
if (!retryable || ageMs < VehiclesService.RPARTSTORE_RETRY_AFTER_MS) {
// not_found / decoded-without-match / recent capped|failed → existing path.
return null;
}
if (await this.isRpartstoreCapSpent()) {
ctx.timings.rpartstore_capped = 1;
return null;
}
await this.db
.update(rpartstoreDecodes)
.set({ status: "pending", updatedAt: new Date() })
.where(eq(rpartstoreDecodes.vin, vin));
await this.enqueueRpartstoreDecode(vin);
ctx.timings.rpartstore_requeued = 1;
return { decoding: { vin, display: ident.display }, vin };
}
if (await this.isRpartstoreCapSpent()) {
ctx.timings.rpartstore_capped = 1;
return null;
}
await this.db
.insert(rpartstoreDecodes)
.values({ vin, status: "pending" })
.onConflictDoNothing();
await this.enqueueRpartstoreDecode(vin);
ctx.timings.rpartstore_enqueued = 1;
// The ONE coverage-gap failure row for this VIN (see the Vinpin fallback).
await this.logQuery(
userId,
vin,
null,
"none",
false,
Date.now() - startTime,
`No catalog — identified as ${ident.display}`,
ctx.timings,
);
return { decoding: { vin, display: ident.display }, vin };
} catch (err) {
this.logger.warn(`RPartStore fallback failed for ${vin}: ${(err as Error).message}`);
return null;
}
}
/** Today's RPartStore search counter (Istanbul day, maintained by the worker)
* is already at the cap → don't queue. Fails open on Redis trouble. */
private async isRpartstoreCapSpent(): Promise<boolean> {
const cap = Number(process.env.RPARTSTORE_DAILY_CAP);
const limit = Number.isFinite(cap) && cap >= 0 ? Math.floor(cap) : 10;
try {
const used = Number((await this.redis.get(rpartstoreDailyKey(new Date()))) ?? 0);
return used >= limit;
} catch {
return false;
}
}
private async enqueueRpartstoreDecode(vin: string): Promise<void> {
// Day-scoped jobId: dedupes concurrent requests for the same VIN today while
// letting a capped/failed VIN be re-queued tomorrow (BullMQ ignores a re-add
// whose jobId still exists among kept completed/failed jobs).
await this.rpartstoreQueue.add(
"rpartstore-decode",
{ vin },
{
jobId: `rpartstore-${vin}-${istanbulDayKey(new Date())}`,
removeOnComplete: true,
removeOnFail: false,
},
);
}
/** /**
* Vinpin ePER decode-oracle fallback for a no-catalog VIN (feature-flagged, * Vinpin ePER decode-oracle fallback for a no-catalog VIN (feature-flagged,
* guarded by the caller on VINPIN_ENABLED + brand allowlist). * guarded by the caller on VINPIN_ENABLED + brand allowlist).

View File

@@ -15,6 +15,11 @@ import { processExpertRewards } from "./jobs/processors/expert-rewards.processor
import { processLifecycleEmails } from "./jobs/processors/lifecycle-email.processor"; import { processLifecycleEmails } from "./jobs/processors/lifecycle-email.processor";
import { processPartPriceRefresh } from "./jobs/processors/part-price-refresh.processor"; import { processPartPriceRefresh } from "./jobs/processors/part-price-refresh.processor";
import { processQueryCleanup } from "./jobs/processors/query-cleanup.processor"; import { processQueryCleanup } from "./jobs/processors/query-cleanup.processor";
import {
buildRpartstoreClient,
processRpartstoreDecode,
rpartstoreDailyCapFromEnv,
} from "./jobs/processors/rpartstore-decode.processor";
import { processSubscriptionExpiry } from "./jobs/processors/subscription-expiry.processor"; import { processSubscriptionExpiry } from "./jobs/processors/subscription-expiry.processor";
import { processTranslation } from "./jobs/processors/translation.processor"; import { processTranslation } from "./jobs/processors/translation.processor";
import { processVinpinDecode } from "./jobs/processors/vinpin-decode.processor"; import { processVinpinDecode } from "./jobs/processors/vinpin-decode.processor";
@@ -244,6 +249,54 @@ vinpinDecodeWorker.on("failed", (job, err) => {
workers.push(vinpinDecodeWorker); workers.push(vinpinDecodeWorker);
// RPartStore Decode Worker (Renault/Dacia VINs the pcat/PL24/emex race couldn't
// identify — runs BEFORE Vinpin). One shared dealer account: concurrency 1 plus
// a 1-job-per-6 s limiter (the portal allows 2 VIN searches per 10 s), and the
// processor enforces the hard RPARTSTORE_DAILY_CAP. Strict no-op when
// RPARTSTORE_ENABLED!=true. Its own ioredis client caches the 1 h Okta token
// and the per-day counter (BullMQ's connection is not for app data).
const rpartstoreRedis = new Redis({
host: process.env.REDIS_HOST || "localhost",
port: Number(process.env.REDIS_PORT) || 6379,
password: process.env.REDIS_PASSWORD || undefined,
maxRetriesPerRequest: null,
lazyConnect: true,
});
let rpartstoreClient: ReturnType<typeof buildRpartstoreClient> | null = null;
const rpartstoreDecodeWorker = new Worker(
QUEUE_NAMES.RPARTSTORE_DECODE,
async (job) => {
return processRpartstoreDecode(job, {
db,
redis: rpartstoreRedis,
decoder: () => {
rpartstoreClient ??= buildRpartstoreClient(rpartstoreRedis);
return rpartstoreClient;
},
dailyCap: rpartstoreDailyCapFromEnv(),
});
},
{
connection,
concurrency: 1,
limiter: { max: 1, duration: 6_000 },
...(telemetry ? { telemetry } : {}),
},
);
rpartstoreDecodeWorker.on("completed", (job, result) => {
console.log(`[worker] rpartstore-decode job ${job.id} completed → ${result?.status}`);
});
rpartstoreDecodeWorker.on("failed", (job, err) => {
console.error(`[worker] rpartstore-decode job ${job?.id} failed: ${err.message}`);
Sentry.captureException(err, {
tags: { queue: QUEUE_NAMES.RPARTSTORE_DECODE, jobId: job?.id },
});
});
workers.push(rpartstoreDecodeWorker);
// Vinpin warm-session daemon: holds the single Vinpin seat warm (browser + login // Vinpin warm-session daemon: holds the single Vinpin seat warm (browser + login
// + Fiat ePER / Renault Rpartstore / Dialogys windows open) during business hours // + Fiat ePER / Renault Rpartstore / Dialogys windows open) during business hours
// (08:00–21:00 Europe/Istanbul), keepalive-nudged every ~75s, so decodes run on // (08:00–21:00 Europe/Istanbul), keepalive-nudged every ~75s, so decodes run on
@@ -355,6 +408,9 @@ async function shutdown(signal: string) {
await vinpinDaemon.stop(); await vinpinDaemon.stop();
console.log("[worker] Vinpin warm daemon stopped"); console.log("[worker] Vinpin warm daemon stopped");
// 2c. Release the RPartStore token/counter Redis client.
rpartstoreRedis.disconnect();
// 3. Close database connection // 3. Close database connection
await sql.end(); await sql.end();
console.log("[worker] Database connection closed"); console.log("[worker] Database connection closed");

View File

@@ -49,6 +49,7 @@ services:
- PL24_PASSWORD_2=${PL24_PASSWORD_2:-} - PL24_PASSWORD_2=${PL24_PASSWORD_2:-}
- PL24_PROXY_DE=${PL24_PROXY_DE:-} - PL24_PROXY_DE=${PL24_PROXY_DE:-}
- PL24_TR_DISABLED=${PL24_TR_DISABLED:-} - PL24_TR_DISABLED=${PL24_TR_DISABLED:-}
- PL24_BACKFILL_ENABLED=${PL24_BACKFILL_ENABLED:-}
- EMEX_USERNAME=${EMEX_USERNAME:-} - EMEX_USERNAME=${EMEX_USERNAME:-}
- EMEX_PASSWORD=${EMEX_PASSWORD:-} - EMEX_PASSWORD=${EMEX_PASSWORD:-}
- EMEX_USE_PROXY=${EMEX_USE_PROXY:-false} - EMEX_USE_PROXY=${EMEX_USE_PROXY:-false}
@@ -69,6 +70,14 @@ services:
# Rpartstore outage → Renault decodes go straight to Dialogys (drops the # Rpartstore outage → Renault decodes go straight to Dialogys (drops the
# per-decode launch-error probe + "Loading application..." stray + budget burn). # per-decode launch-error probe + "Loading application..." stray + budget burn).
- VINPIN_RPARTSTORE_ENABLED=${VINPIN_RPARTSTORE_ENABLED:-true} - VINPIN_RPARTSTORE_ENABLED=${VINPIN_RPARTSTORE_ENABLED:-true}
# RPartStore (Renault/Dacia) VIN-decode fallback — after pcat/PL24/emex, before Vinpin.
# Hard daily cap on searches sent to rpartstore.renault.com (default 10).
- RPARTSTORE_ENABLED=${RPARTSTORE_ENABLED:-false}
- RPARTSTORE_USER=${RPARTSTORE_USER:-}
- RPARTSTORE_PASS=${RPARTSTORE_PASS:-}
- RPARTSTORE_DAILY_CAP=${RPARTSTORE_DAILY_CAP:-10}
- RPARTSTORE_BROKER_URL=${RPARTSTORE_BROKER_URL:-wss://1po-bff.renault-edh.com/ws}
- RPARTSTORE_APP_VERSION=${RPARTSTORE_APP_VERSION:-1.34.0.6}
- POSTAL_API_URL=${POSTAL_API_URL:-} - POSTAL_API_URL=${POSTAL_API_URL:-}
- POSTAL_API_KEY=${POSTAL_API_KEY:-} - POSTAL_API_KEY=${POSTAL_API_KEY:-}
- POSTAL_FROM_ADDRESS=${POSTAL_FROM_ADDRESS:-noreply@sase.tr} - POSTAL_FROM_ADDRESS=${POSTAL_FROM_ADDRESS:-noreply@sase.tr}
@@ -207,6 +216,7 @@ services:
- PL24_PASSWORD_2=${PL24_PASSWORD_2:-} - PL24_PASSWORD_2=${PL24_PASSWORD_2:-}
- PL24_PROXY_DE=${PL24_PROXY_DE:-} - PL24_PROXY_DE=${PL24_PROXY_DE:-}
- PL24_TR_DISABLED=${PL24_TR_DISABLED:-} - PL24_TR_DISABLED=${PL24_TR_DISABLED:-}
- PL24_BACKFILL_ENABLED=${PL24_BACKFILL_ENABLED:-}
- EMEX_USERNAME=${EMEX_USERNAME:-} - EMEX_USERNAME=${EMEX_USERNAME:-}
- EMEX_PASSWORD=${EMEX_PASSWORD:-} - EMEX_PASSWORD=${EMEX_PASSWORD:-}
- EMEX_USE_PROXY=${EMEX_USE_PROXY:-false} - EMEX_USE_PROXY=${EMEX_USE_PROXY:-false}
@@ -225,6 +235,14 @@ services:
- VINPIN_WARM_DAEMON=${VINPIN_WARM_DAEMON:-false} - VINPIN_WARM_DAEMON=${VINPIN_WARM_DAEMON:-false}
# Skip Rpartstore during a known upstream outage (Renault → straight to Dialogys). # Skip Rpartstore during a known upstream outage (Renault → straight to Dialogys).
- VINPIN_RPARTSTORE_ENABLED=${VINPIN_RPARTSTORE_ENABLED:-true} - VINPIN_RPARTSTORE_ENABLED=${VINPIN_RPARTSTORE_ENABLED:-true}
# RPartStore (Renault/Dacia) VIN-decode fallback — after pcat/PL24/emex, before Vinpin.
# Hard daily cap on searches sent to rpartstore.renault.com (default 10).
- RPARTSTORE_ENABLED=${RPARTSTORE_ENABLED:-false}
- RPARTSTORE_USER=${RPARTSTORE_USER:-}
- RPARTSTORE_PASS=${RPARTSTORE_PASS:-}
- RPARTSTORE_DAILY_CAP=${RPARTSTORE_DAILY_CAP:-10}
- RPARTSTORE_BROKER_URL=${RPARTSTORE_BROKER_URL:-wss://1po-bff.renault-edh.com/ws}
- RPARTSTORE_APP_VERSION=${RPARTSTORE_APP_VERSION:-1.34.0.6}
# Novu lifecycle e-mail automation — the worker fires trial-ending + win-back # Novu lifecycle e-mail automation — the worker fires trial-ending + win-back
- NOVU_API_URL=${NOVU_API_URL:-https://api.bildirim.semih.ai} - NOVU_API_URL=${NOVU_API_URL:-https://api.bildirim.semih.ai}
- NOVU_API_KEY=${NOVU_API_KEY:-} - NOVU_API_KEY=${NOVU_API_KEY:-}

View File

@@ -76,6 +76,26 @@ export const envSchema = z.object({
.transform((v) => v === "true") .transform((v) => v === "true")
.default("false"), .default("false"),
// RPartStore (rpartstore.renault.com) Renault/Dacia VIN-decode fallback — runs
// AFTER the pcat/PL24/emex race, before Vinpin. Off by default; the worker
// needs the dealer credentials. RPARTSTORE_DAILY_CAP is a hard ceiling on VIN
// searches sent per Istanbul day (the portal itself also rate-limits 2/10 s).
RPARTSTORE_ENABLED: z
.string()
.transform((v) => v === "true")
.default("false"),
RPARTSTORE_USER: z.string().optional(),
RPARTSTORE_PASS: z.string().optional(),
RPARTSTORE_DAILY_CAP: z.coerce.number().int().min(0).default(10),
RPARTSTORE_BROKER_URL: z.preprocess(
(v) => (typeof v === "string" && v.trim() === "" ? undefined : v),
z.string().url().default("wss://1po-bff.renault-edh.com/ws"),
),
RPARTSTORE_APP_VERSION: z.preprocess(
(v) => (typeof v === "string" && v.trim() === "" ? undefined : v),
z.string().default("1.34.0.6"),
),
// Parts-Catalogs (Playwright JWT capture + DataImpulse proxy) // Parts-Catalogs (Playwright JWT capture + DataImpulse proxy)
PCAT_USE_PROXY: z.string().default("true"), PCAT_USE_PROXY: z.string().default("true"),
PCAT_PROXY_HOST: z.string().default("gw.dataimpulse.com"), PCAT_PROXY_HOST: z.string().default("gw.dataimpulse.com"),