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,
"tag": "0035_subscription_dunning",
"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,
userSubscriptions,
} 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 {
type PL24DecodedCategory,
@@ -170,7 +170,8 @@ export class CatalogService {
.where(and(...whereConditions));
if (dbVehicles.length > 0) {
return dbVehicles;
const healed = await this.healStaleBrowseRows(brandName, dbVehicles, whereConditions);
return healed ?? dbVehicles;
}
// Fetch from PL24 for each service
@@ -225,6 +226,118 @@ export class CatalogService {
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.
*/

View File

@@ -17,6 +17,8 @@ function createService(db: any) {
};
const pl24Service = {
getCategories: vi.fn().mockResolvedValue([]),
// P4→P5 onarma yolu bunu çağırır (plv2 Faz 2).
decodeVin: vi.fn().mockResolvedValue(null),
};
const emexService = {
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 { PL24FordLegacyService } from "../integrations/pl24/pl24-ford-legacy.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 { getServiceApiPath } from "../integrations/pl24/pl24.types";
import { classifyNode, foldName, mapToCanonical } from "../jobs/canonical-lexicon";
import { RedisService } from "../redis/redis.service";
import { StorageService } from "../storage/storage.service";
@@ -71,6 +76,20 @@ export class CategoriesService {
.from(categories)
.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 (dbCategories.length === 0 && vehicle.rawData) {
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
* 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(
ids: 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 }),
});
// ─── 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) ──
export const vehicles = pgTable(
"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 { PL24AuthService } from "./pl24-auth.service";
import { PL24BudgetExceededError } from "./pl24-budget.service";
// 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'ı;
@@ -459,3 +460,33 @@ describe("PL24AuthService — 1 saatlik kesinti alarmı (Telegram)", () => {
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 { TelegramService } from "../../common/telegram.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 {
PL24AuthorizeRequest,
@@ -485,6 +485,15 @@ export class PL24AuthService implements OnModuleInit {
return { sessionToken };
} catch (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" };
return { error: err.message };
}
@@ -524,6 +533,9 @@ export class PL24AuthService implements OnModuleInit {
return token;
} catch (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}`);
throw err instanceof UnauthorizedException
? err

View File

@@ -139,3 +139,28 @@ describe("PL24BudgetService — telemetri", () => {
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();
spent = await this.redis.incr(key);
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 {
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 { 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
// 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);
});
});
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.
*/
/** `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. */
const GROUP_WID_PATTERN =
@@ -61,5 +97,55 @@ export function isPl24GroupNode(opts: {
hasSubgroups?: boolean | null;
}): boolean {
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);
}
/**
* 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".
*/
export function normalizeLabel(label: string): string {
return label
.toLocaleLowerCase("tr")
.normalize("NFD")
.replace(/[\u0300-\u036f]/g, "")
.replace(/ı/g, "i")
.replace(/[\s/]+/g, "_")
.trim();
return (
label
.toLocaleLowerCase("tr")
.normalize("NFD")
// \p{M} (all combining marks) rather than the U+0300–U+036F range: the range
// is a character class that can also match a base character followed by a
// 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 ====================
/**
* 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.
*/
@@ -1184,6 +1263,31 @@ export class PL24Service {
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
const prNrRecords = segments.prNr?.records || [];
const prNrByCode: Record<string, string> = {};
@@ -1226,13 +1330,17 @@ export class PL24Service {
// normalizeLabel folds ı→i and strips diacritics, so "Şanzıman kodu" and
// "ŞANZIMAN KODU" both arrive as "sanziman_kodu". PSA uses "AKTARMA
// SİSTEMLERİ" ("5 MEKANİK VİTES KUTUSU"), Subaru "Mission".
const transmissionCode = lookup(
const transmissionCode = lookupDescriptive(
"sanziman_kodu",
"transmission_code",
"vites_kutusu",
"atm,mtm",
"aktarma_sistemleri",
"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",
"mission",
);
@@ -1242,7 +1350,10 @@ export class PL24Service {
const bodyType =
Object.entries(prNrByCode).find(([code]) => code.startsWith("K8"))?.[1] ||
// 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;
// Engine description from prNr D3* (Motor nitelikleri)
@@ -1301,7 +1412,10 @@ export class PL24Service {
damToModelYear(lookup("dam")) ||
extractModelYear(vin) ||
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,
engineCode:
engineCode ||
@@ -2191,11 +2305,22 @@ export class PL24Service {
p5mitsubishi: "/extern/vehicles/vehiclesOverview", // Mitsubishi
p5suzuki: "/extern/vehicle/modelFamilies", // Suzuki
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
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}`;

View File

@@ -1,5 +1,5 @@
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
// 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
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",
WDC: "mercedes_parts",
W1K: "mercedes_parts",
@@ -578,15 +582,21 @@ export const PL24_WMI_SERVICE_MAP: Record<string, string> = {
JTJ: "lexus_parts",
"2T2": "lexus_parts",
// Renault — DISABLED: PL24 has suspended Renault VIN identification ("Bu marka
// için şasi numarası tanımlamasının belirsiz bir süre için mevcut olmayacağını
// üzülerek bildiririz."). renault_parts authorizes fine but every decode throws
// that message → wasted ~1s call AND it trips the PL24 circuit breaker, which
// then skips PL24 for ALL brands. PCAT + EMEX cover Renault. Re-enable when PL24
// restores Renault VIN decode.
// VF1: "renault_parts",
// VF6: "renault_parts",
// VNE: "renault_parts",
// Renault — RE-ENABLED 2026-09-20. It was disabled while PL24 had suspended
// Renault VIN identification ("Bu marka için şasi numarası tanımlamasının
// belirsiz bir süre için mevcut olmayacağını üzülerek bildiririz."), which both
// wasted a call per decode and, back then, tripped the global PL24 breaker.
// BOTH reasons are gone, each verified against prod rather than assumed:
// 1. The suspension is over — live directAccess on /p5renault for a real
// customer VIN (VF14SRCL458170337) returns resultStatus
// VEHICLE_IDENTIFIED, "SYMBOL II/LOGAN II", and the WMI service answers
// {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
UU1: "dacia_parts",
@@ -606,6 +616,10 @@ export const PL24_WMI_SERVICE_MAP: Record<string, string> = {
WMH: "man_parts",
// 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",
JMY: "mmc_parts",
MMB: "mmc_parts",
@@ -615,6 +629,7 @@ export const PL24_WMI_SERVICE_MAP: Record<string, string> = {
JS2: "suzuki_parts",
JS3: "suzuki_parts",
TSM: "suzuki_parts",
JSA: "suzuki_parts", // Suzuki (canlı WMI: suzuki_parts, error:false)
MA3: "suzuki_parts",
MBH: "suzuki_parts",
@@ -638,10 +653,19 @@ export const PL24_WMI_SERVICE_MAP: Record<string, string> = {
// Hyundai
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)
// Kia
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
// Nissan
@@ -675,7 +699,19 @@ export const PL24_WMI_SERVICE_MAP: Record<string, string> = {
// Peugeot (PSA)
VF3: "peugeot_parts", // Peugeot SA (France)
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
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",
PART_PRICE_REFRESH: "part-price-refresh",
VINPIN_DECODE: "vinpin-decode",
RPARTSTORE_DECODE: "rpartstore-decode",
CANONICAL_BACKFILL: "canonical-backfill",
} as const;

View File

@@ -21,6 +21,10 @@ import {
PartPriceRefreshQueueProvider,
} from "./queues/part-price-refresh.queue";
import { QUERY_CLEANUP_QUEUE, QueryCleanupQueueProvider } from "./queues/query-cleanup.queue";
import {
RPARTSTORE_DECODE_QUEUE,
RpartstoreDecodeQueueProvider,
} from "./queues/rpartstore-decode.queue";
import {
SUBSCRIPTION_EXPIRY_QUEUE,
SubscriptionExpiryQueueProvider,
@@ -39,6 +43,7 @@ import { VINPIN_DECODE_QUEUE, VinpinDecodeQueueProvider } from "./queues/vinpin-
ExpertRewardsQueueProvider,
PartPriceRefreshQueueProvider,
VinpinDecodeQueueProvider,
RpartstoreDecodeQueueProvider,
CanonicalBackfillQueueProvider,
PrefetchWorkerService,
],
@@ -52,6 +57,7 @@ import { VINPIN_DECODE_QUEUE, VinpinDecodeQueueProvider } from "./queues/vinpin-
EXPERT_REWARDS_QUEUE,
PART_PRICE_REFRESH_QUEUE,
VINPIN_DECODE_QUEUE,
RPARTSTORE_DECODE_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. */
function currentIstanbulHour(): number {
/** Europe/Istanbul hour (0–23) at an arbitrary instant. */
function istanbulHourAt(ms: number): number {
const hourStr = new Intl.DateTimeFormat("en-US", {
timeZone: "Europe/Istanbul",
hour: "numeric",
hour12: false,
}).format(new Date());
return Number.parseInt(hourStr, 10);
}).format(new Date(ms));
// `% 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.
*/

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", () => {
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 () => {
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 () => {
const { service, redis } = makeDeps({ waiting: 0, limitResults: [] });
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);
});
});
/**
* 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 { DATABASE, type Database } from "../database/database.provider";
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 { RedisService } from "../redis/redis.service";
import { QUEUE_NAMES, getBullConnection } from "./bull.config";
import { backfillContext } from "./prefetch-context";
import {
RateLimitError,
alignToWindow,
checkCooldown,
checkTimeWindow,
initProgress,
@@ -100,6 +101,25 @@ const EST_JOBS_PER_VEHICLE = Number(process.env.PREFETCH_EST_JOBS_PER_VEHICLE) |
*/
const DAILY_FAST_RESERVE = 0.2;
/** 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"];
/** Redis key holding the rolling rescan cursor (last createdAt seen). */
const BACKFILL_CURSOR_KEY = "prefetch:backfill:cursor";
@@ -203,7 +223,7 @@ function jitter(ms: number): number {
}
/** Test-only surface for the pure helpers above. */
export const __testables = { maxDepthFor, jitter };
export const __testables = { maxDepthFor, jitter, isPl24BackfillEnabled };
// ── Phase-1 residue exclusion ──
/**
@@ -354,12 +374,30 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy {
job.name === "prefetch-parts")
) {
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);
// Daily budget AFTER the per-minute gate: a job deferred on the minute
// ceiling above never reaches here, so rate-limited retries don't inflate
// the daily counter — only jobs about to do real work are counted. The
// lane decides which threshold applies (backfill stops at the main limit,
// the user's fast lane may use the full budget).
// 4. Daily budget last: a job deferred on any gate above never reaches
// here, so only jobs about to do real work are counted. The lane
// decides which threshold applies (backfill stops at the main limit,
// the user's fast lane may use the full budget).
await this.checkSourceDailyBudget(data.source, lane);
}
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;
this.logger.log(`[prefetch] Init for vehicle=${vehicleId}, source=${source}`);
await checkCooldown(this.redis, source);
checkTimeWindow(source);
// Cooldown + time-window are enforced in process() before the daily budget
// 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.
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".
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
const [partCheck] = await this.db
.select({ id: parts.id })
@@ -563,8 +601,8 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy {
const { vehicleId, categoryId, source, depth, fast = false } = job.data;
this.logger.log(`[prefetch] Children for category=${categoryId}, depth=${depth}`);
await checkCooldown(this.redis, source);
checkTimeWindow(source);
// Cooldown + time-window are enforced in process() before the daily budget
// is debited — see the comment there; re-checking here would be a no-op.
const depthCeiling = maxDepthFor(source, fast);
if (depth >= depthCeiling) {
@@ -631,8 +669,8 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy {
const { vehicleId, categoryId, source } = job.data;
this.logger.log(`[prefetch] Parts for category=${categoryId}`);
await checkCooldown(this.redis, source);
checkTimeWindow(source);
// Cooldown + time-window are enforced in process() before the daily budget
// is debited — see the comment there; re-checking here would be a no-op.
try {
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.
const eligible: string[] = [];
for (const s of BACKFILL_SOURCES) {
// PL24_TR_DISABLED: tr is dead on PL24's side and ALL pl24 traffic maps to
// the sole surviving de account — keep background backfill off it entirely
// (fast/user lane still flows) so bulk load can't burn the last account.
if (s === "pl24" && process.env.PL24_TR_DISABLED === "true") continue;
// Background backfill is opt-in per source; pl24 defaults to OFF so bulk
// load can never burn the one surviving account by accident.
if (s === "pl24" && !isPl24BackfillEnabled()) continue;
if (await this.redis.exists(`prefetch:activity:${s}`)) continue;
if (cfg.businessHoursOnly !== false && !isWithinTimeWindow(s)) continue;
const mainLimit = this.dailyMainLimit(s);
@@ -875,6 +912,7 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy {
cat: {
id: string;
linkPath: string | null;
linkWid?: string | null;
source: string;
unavailable: boolean;
hasSubgroups?: boolean | null;
@@ -886,7 +924,18 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy {
): Promise<number> {
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
const [partCheck] = await this.db
.select({ id: parts.id })
@@ -951,6 +1000,7 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy {
linkPath: string | null,
source: string,
hasSubgroups?: boolean | null,
linkWid?: string | null,
): boolean {
if (!linkPath) return false;
// 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
// list was case-sensitive, so p5psa/p5volvo's camelCase `/details/vin/
// bomDetails` was never recognised as a leaf and its parts were never
// prefetched (still true today for Mitsubishi/Fiat/Renault).
return isPl24LeafNode({ linkPath, hasSubgroups });
// prefetched.
// `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
// source stay the same as before, so upstream load is unchanged.
const total = SOURCE_RATE_MAX[source] ?? 0;
// PL24_TR_DISABLED: park already-queued pl24 MAIN-lane jobs (long defer, no
// attempt consumed) — the eligibility scan stops producing new ones, this
// stops the existing backlog from draining through the de account.
if (source === "pl24" && lane === "main" && process.env.PL24_TR_DISABLED === "true") {
// Park already-queued pl24 MAIN-lane jobs while background backfill is off
// (long defer, no attempt consumed) — the eligibility scan stops producing
// new ones, this stops an existing backlog from draining through the account.
if (source === "pl24" && lane === "main" && !isPl24BackfillEnabled()) {
throw new RateLimitError(15 * 60_000, "source-rate");
}
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
// 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).
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) {
this.logger.warn(
`[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");

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 { beforeEach, describe, expect, it, vi } from "vitest";
import { catalogVehicles, queryLogs, vinpinDecodes } from "../database/schema/core";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import {
catalogVehicles,
queryLogs,
rpartstoreDecodes,
vinpinDecodes,
} from "../database/schema/core";
import { VehiclesService } from "./vehicles.service";
vi.mock("@sase/shared", () => ({
@@ -93,10 +98,15 @@ function createService(dbOrOverrides: any = {}) {
add: vi.fn().mockResolvedValue(undefined),
};
const rpartstoreQueue = {
add: vi.fn().mockResolvedValue(undefined),
};
const service = new VehiclesService(
db as any,
prefetchQueue as any,
vinpinQueue as any,
rpartstoreQueue as any,
corgiService as any,
pl24Service as any,
vinApiService as any,
@@ -115,6 +125,7 @@ function createService(dbOrOverrides: any = {}) {
partsCatalogsService,
redisService,
vinpinQueue,
rpartstoreQueue,
};
}
@@ -340,9 +351,7 @@ describe("VehiclesService", () => {
selectChain.limit = vi
.fn()
.mockImplementation(() =>
currentTable === vinpinDecodes
? [{ vin: "NM435600006123456", status: "pending" }]
: [],
currentTable === vinpinDecodes ? [{ vin: "NM435600006123456", status: "pending" }] : [],
);
const insertChain = {
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,
plans,
queryLogs,
rpartstoreDecodes,
userBrands,
userSubscriptions,
userVehicles,
@@ -34,10 +35,13 @@ import { PartsCatalogsService } from "../integrations/parts-catalogs/parts-catal
import { PcatCar, PcatVinResult } from "../integrations/parts-catalogs/parts-catalogs.types";
import { PL24Service } from "../integrations/pl24/pl24.service";
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 { isVinpinBrandAllowed } from "../integrations/vinpin/vinpin.constants";
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 { RPARTSTORE_DECODE_QUEUE } from "../jobs/queues/rpartstore-decode.queue";
import { VINPIN_DECODE_QUEUE } from "../jobs/queues/vinpin-decode.queue";
import { RedisService } from "../redis/redis.service";
import { vinCandidateStashKey, vinResolveCacheKeys } from "./vin-cache-keys";
@@ -99,6 +103,7 @@ export class VehiclesService {
@Inject(DATABASE) private db: Database,
@Inject(CATALOG_PREFETCH_FAST_QUEUE) private prefetchQueue: Queue,
@Inject(VINPIN_DECODE_QUEUE) private vinpinQueue: Queue,
@Inject(RPARTSTORE_DECODE_QUEUE) private rpartstoreQueue: Queue,
private corgiService: CorgiService,
private pl24Service: PL24Service,
private vinApiService: VinApiService,
@@ -191,6 +196,19 @@ export class VehiclesService {
// The fallback logs the coverage-gap failure row exactly once, at
// first sighting (job enqueue); not_found/failed/stale fall through
// 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)) {
const vinpinResp = await this.tryVinpinFallback(vin, ident, userId, ctx, startTime);
if (vinpinResp) return vinpinResp;
@@ -791,6 +809,154 @@ export class VehiclesService {
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,
* 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 { processPartPriceRefresh } from "./jobs/processors/part-price-refresh.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 { processTranslation } from "./jobs/processors/translation.processor";
import { processVinpinDecode } from "./jobs/processors/vinpin-decode.processor";
@@ -244,6 +249,54 @@ vinpinDecodeWorker.on("failed", (job, err) => {
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
// + Fiat ePER / Renault Rpartstore / Dialogys windows open) during business hours
// (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();
console.log("[worker] Vinpin warm daemon stopped");
// 2c. Release the RPartStore token/counter Redis client.
rpartstoreRedis.disconnect();
// 3. Close database connection
await sql.end();
console.log("[worker] Database connection closed");

View File

@@ -49,6 +49,7 @@ services:
- PL24_PASSWORD_2=${PL24_PASSWORD_2:-}
- PL24_PROXY_DE=${PL24_PROXY_DE:-}
- PL24_TR_DISABLED=${PL24_TR_DISABLED:-}
- PL24_BACKFILL_ENABLED=${PL24_BACKFILL_ENABLED:-}
- EMEX_USERNAME=${EMEX_USERNAME:-}
- EMEX_PASSWORD=${EMEX_PASSWORD:-}
- EMEX_USE_PROXY=${EMEX_USE_PROXY:-false}
@@ -69,6 +70,14 @@ services:
# Rpartstore outage → Renault decodes go straight to Dialogys (drops the
# per-decode launch-error probe + "Loading application..." stray + budget burn).
- 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_KEY=${POSTAL_API_KEY:-}
- POSTAL_FROM_ADDRESS=${POSTAL_FROM_ADDRESS:-noreply@sase.tr}
@@ -207,6 +216,7 @@ services:
- PL24_PASSWORD_2=${PL24_PASSWORD_2:-}
- PL24_PROXY_DE=${PL24_PROXY_DE:-}
- PL24_TR_DISABLED=${PL24_TR_DISABLED:-}
- PL24_BACKFILL_ENABLED=${PL24_BACKFILL_ENABLED:-}
- EMEX_USERNAME=${EMEX_USERNAME:-}
- EMEX_PASSWORD=${EMEX_PASSWORD:-}
- EMEX_USE_PROXY=${EMEX_USE_PROXY:-false}
@@ -225,6 +235,14 @@ services:
- VINPIN_WARM_DAEMON=${VINPIN_WARM_DAEMON:-false}
# Skip Rpartstore during a known upstream outage (Renault → straight to Dialogys).
- 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_API_URL=${NOVU_API_URL:-https://api.bildirim.semih.ai}
- NOVU_API_KEY=${NOVU_API_KEY:-}

View File

@@ -76,6 +76,26 @@ export const envSchema = z.object({
.transform((v) => v === "true")
.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)
PCAT_USE_PROXY: z.string().default("true"),
PCAT_PROXY_HOST: z.string().default("gw.dataimpulse.com"),