diff --git a/FIXPLAN-2026-07-09.md b/FIXPLAN-2026-07-09.md new file mode 100644 index 0000000..8bdfcc6 --- /dev/null +++ b/FIXPLAN-2026-07-09.md @@ -0,0 +1,320 @@ +# Fix- und Datenreparatur-Plan (Stand 2026-07-09, Übergabe an Gemini) + +> **Abnahmekriterium für alle Code-Änderungen:** `dotnet test src/Predictalytics.Application.Tests` muss +> **16 grün + 1 übersprungen** liefern (der Skip `CheckpointResetAndReplay_DoesNotDoubleCountBalance` ist eine +> dokumentierte, bewusste Entscheidung). Die Assertions der Invarianten-Tests dürfen **nicht** verändert werden — +> sie definieren das Soll-Verhalten. Wenn ein Test rot wird, ist der Code falsch, nicht der Test. + +## Hintergrund + +Die Engine-Fixes vom 09.07. sind korrekt (Tests grün). Die im WebUI sichtbaren Probleme haben drei andere Ursachen: + +1. Die Buttons der **Trader-Detailseite** nutzen alte, Job-lose Endpoints (die Listen-Buttons nutzen bereits das Job-System). +2. Die **abgeleiteten Daten in der DB stammen aus der Bug-Ära** (Snapshots/Positionen wurden von den alten, fehlerhaften + Engine-Versionen berechnet). Beispiel aus dem Live-System: `PnL30d = 244,0K` bei `TotalPnL = 158,5K`, weil der + Basis-Snapshot `-85,5K` enthält (korrupter Altwert). Kein Code-Fix ändert das — die Daten müssen einmalig repariert werden. +3. **Deadlocks + Shutdown-Fehlerkaskaden** in den Workern (unbatchtes Reconciliation-UPDATE, fehlende Cancellation-Behandlung). + +**Ein DB-Reset ist NICHT nötig.** Die Rohdaten (`Trades`) sind größtenteils intakt; Positionen, Analytics, Snapshots und +Scores sind abgeleitet und lokal neu berechenbar. Nur Trader, deren Alt-Trades die Retention bereits gelöscht/kompaktiert +hat, brauchen einen gezielten API-Re-Import (kleine Teilmenge, siehe Teil B). + +--- + +## Teil A — Code-Fixes + +### A1. Trader-Detailseite: Buttons auf das Job-System umstellen +**Problem:** +- `manualUpdateTrader` in `src/Predictalytics.Api/wwwroot/js/app.js` (~Zeile 110) baut die URL mit **Backslashes**: + `` fetch(`\api\traders\${id}\refresh`) `` — in JS-Template-Literals ist `\t` ein Tab und `\${id}` unterdrückt die + Interpolation. Der Request geht als Müll-URL raus. +- Der „Analyze"-Button der Detailseite (~Zeile 441) ruft `POST /api/traders/{id}/force-analyze` — läuft **synchron** im + API-Request, legt **keinen** `BackgroundJob` an (der Alert behauptet es aber) und führt nur die PnL-Engine aus, + weder `CopytradingEstimator` noch KI. + +**Fix:** +- `btn-sync-trader` → `POST /api/jobs/sync/{id}`, `btn-analyze-trader` → `POST /api/jobs/analyze/{id}` + (bestehende Funktionen `queueHistorySync(id)` / `queueTraderAnalysis(id)` wiederverwenden). +- Alert-Texte ehrlich machen (Job-Id anzeigen oder auf die Jobs-Seite verweisen). +- Die Endpoints `/{id}/force-analyze` und `/{id}/refresh` entweder entfernen oder intern auf Job-Enqueue umbauen — + es darf nur noch **einen** Auslöse-Pfad geben. + +### A2. TradeReconciliationWorker: Bulk-UPDATE batchen, Fehler pro Markt behandeln +**Problem:** Das eine große `UPDATE Trades ... INNER JOIN ... WHERE MarketOutcomeId IS NULL` läuft über die gesamte +Tabelle, hält minutenlang Locks und produziert Deadlocks mit den Insert-Workern. Außerdem verwirft der eine +try/catch um den ganzen Batch bei jedem Einzelfehler (z. B. ein fehlgeschlagener `GetMarketAsync`) die komplette Restarbeit. + +**Fix:** +- UPDATE in Batches. Achtung: MySQL erlaubt kein `LIMIT` bei Multi-Table-UPDATE — Pattern mit Subquery verwenden: + ```sql + UPDATE Trades t + JOIN ( + SELECT t2.Id, o.Id AS OutcomeId, o.Label, m.Id AS MarketDbId + FROM Trades t2 + JOIN MarketOutcomes o ON t2.AssetId = o.TokenId + JOIN Markets m ON o.MarketId = m.Id + WHERE t2.MarketOutcomeId IS NULL AND t2.AssetId != '' + LIMIT 5000 + ) x ON t.Id = x.Id + SET t.MarketOutcomeId = x.OutcomeId, t.Outcome = x.Label, t.DbMarketId = x.MarketDbId; + ``` + In einer Schleife ausführen, bis 0 Zeilen betroffen sind (mit kurzem Delay zwischen den Batches). +- `GetMarketAsync`-Fehler pro Markt fangen und loggen — die restlichen Märkte des Batches weiterverarbeiten. +- Das Checkpoint-Reset (`LastAppliedTradeId = 0`) weiterhin **nur** für Positionen mit `IsHistoryPruned = 0` + (ist bereits so umgesetzt — nicht regressieren, Test `PrunedPositionWithResetCheckpoint_DoesNotDoubleCount` wacht darüber). + +### A3. Deadlock-Retry in `TradeRepository.AddRangeAsync` +`MySqlException` mit `Number == 1213` (Deadlock) oder `1205` (Lock wait timeout) → bis zu 3 Versuche mit Backoff +(250 ms / 500 ms / 1 s). Chunk-Größe von 1000 auf 500 Zeilen reduzieren. Bei endgültigem Fehlschlag: Fehler loggen +inkl. Anzahl verlorener Zeilen. + +### A4. Saubere Cancellation in allen Worker-Loops +**Problem:** Beim Stoppen des Servers wirft jede laufende Operation `OperationCanceledException`; der +`TraderAnalyticsWorker` fängt das **pro Trader** als ERROR und nudelt durch den restlichen 500er-Batch +(→ hunderte Fehlerlog-Einträge pro Shutdown, verzögerter Stopp). + +**Fix (in TraderAnalyticsWorker, PollingWorker, TradeHistoryWorker, TradeReconciliationWorker, TradeContextEnrichmentWorker):** +- Vor jeder Batch-Iteration: `if (ct.IsCancellationRequested) break;` +- `catch (OperationCanceledException) when (ct.IsCancellationRequested)` separat behandeln: + als Information loggen („shutting down"), Schleife beenden — **nicht** als Error. + +### A5. Hängengebliebene Jobs wiederbeleben +**Problem:** Jobs, die beim Shutdown `InProgress` waren, bleiben für immer stecken (`GetNextPendingJobAsync` holt nur `Pending`). + +**Fix:** Beim Start der Job-verarbeitenden Worker (oder einmal pro Zyklus): Jobs mit `Status = InProgress` und +`StartedAt < UtcNow - 15min` zurück auf `Pending` setzen (Log-Hinweis). + +### A6. Deep-Resync-Fähigkeit (Voraussetzung für die Datenreparatur in Teil B) +**Problem:** `PolymarketApiClient.GetTradesAsync` macht genau **einen** Request (`/activity?user=X&limit=1000`, +keine Pagination; die API cappt vermutlich ohnehin bei 500). Der „INITIAL FULL sync" holt also nur die jüngsten +~500–1000 Aktivitäten. Für die Reparatur der Retention-/Kompaktierungs-Opfer brauchen wir die **komplette** Historie. + +**Fix:** +- Neue Methode `GetTradesPagedAsync(wallet, ...)` mit **Timestamp-basierter Pagination**: erste Seite normal laden, + Folgeseiten mit `&end=<ältester Timestamp der Vorseite - 1>` bis eine leere Seite kommt. (Timestamp-Pagination ist + robuster als `offset`, da Offset-Limits der API umgangen werden.) `limit=500` verwenden. Jede Seite über den + vorhandenen `IRateLimiter` drosseln. +- Neuer `JobType.DeepResync` (Migration für Enum nicht nötig, Enum ist int): Der `TradeHistoryWorker` behandelt ihn wie + `HistorySync`, lädt aber ALLE Seiten. +- **Vor** dem Import im DeepResync-Pfad für den Trader aufräumen (sonst Doppelzählung!): + 1. `DELETE FROM Trades WHERE TraderId = @id AND PlatformTradeId LIKE 'COMPACT_%'` + (Re-Import bringt die Original-Trades zurück; die Aggregate dürfen nicht zusätzlich existieren), + 2. alle `TraderPositions` des Traders löschen (**inklusive** `IsHistoryPruned = 1` — die Konserve wird durch den + vollständigen Re-Import ersetzt), + 3. nach erfolgreichem Import: `IsInitialImportComplete = true`, `LastTradesUpdatedAt = now`, `LastAnalyzedAt = NULL`. +- Endpoints: + - `POST /api/jobs/deep-resync/{traderId}` (einzeln), + - `POST /api/jobs/deep-resync-pruned?take=25` — enqueued DeepResync-Jobs für Trader mit `IsHistoryPruned`-Positionen + oder `COMPACT_`-Trades, Watchlist zuerst, dann nach `TotalTrades` absteigend. +- WebUI: Button „Deep Resync" auf der Jobs-Seite neben „Analyze Backlog". + +### A7. Retention pausierbar machen +Neues Config-Flag `RetentionSettings:Enabled` (Default `true`), das der `TradeRetentionWorker` pro Zyklus prüft. +Während der Datenreparatur steht es auf `false` — sonst prunt/kompaktiert die tägliche Runde die frisch +re-importierten Alt-Trades wieder weg, bevor die Engine sie eingerechnet hat. + +### A8. ⚠️ NEU (2026-07-10, höchste Priorität): ResolutionOutcome existiert in der Gamma-API nicht — alle Gewinner werden als Totalverlust gebucht + +**Empirisch gegen die Live-API verifiziert:** Die Antwort von `gamma-api.polymarket.com/markets` enthält +**weder** ein Feld `resolution_outcome` (so mappt es `GammaMarketResponse` aktuell) **noch** `resolutionOutcome` +**noch** `resolved`. Folgen im Bestand und in jeder Neuberechnung: + +- `Market.ResolutionOutcome` ist für **jeden** Markt `NULL` → `MarketOutcomeHelper.IsWinningOutcome` liefert immer + `false` → jeder Redeem und jeder virtuelle Payout bucht Auszahlung **0** → **jeder aufgelöste Markt ist ein + Totalverlust**. Das erzeugt exakt das Live-Bild: WinRate 0 %, Quality Edge 0.0, negative Total-PnL. +- `IsResolved = raw.Resolved || raw.Closed` degeneriert zu `IsResolved = closed`. Märkte, die für den Handel + geschlossen, aber noch nicht UMA-aufgelöst sind, werden **vorzeitig** zu Payout 0 ausgebucht. + +**Wie man den Gewinner wirklich erkennt** (Live-API-Beispiele): Nach der Auflösung rasten die `outcomePrices` +auf `["1","0"]` / `["0","1"]` ein (liegen bei uns bereits in `MarketOutcome.CurrentPrice`), und es gibt das Feld +`umaResolutionStatus` (String, `"resolved"` bei aufgelösten Märkten; bei sehr alten Märkten fehlt es). + +**Fix (drei Teile):** +1. **Model:** In `GammaMarketResponse` das tote `resolution_outcome`-Mapping entfernen, + `[JsonPropertyName("umaResolutionStatus")] public string? UmaResolutionStatus` ergänzen. +2. **Mapper (`MapGammaMarket`):** + - Preise parsen, dann: `pricesSnapped = alle Outcome-Preise ≤ 0.02 oder ≥ 0.98` (und mindestens ein Preis ≥ 0.98). + - `IsResolved = raw.UmaResolutionStatus == "resolved" || (raw.Closed && pricesSnapped)`. + - `ResolutionOutcome = Label des Outcomes mit Preis ≥ 0.98` (nur wenn `IsResolved`; sonst `NULL`). +3. **Engine-Absicherung (Defense in depth, weil der Bestand NULL-Werte enthält):** Redeem-Buchung und virtueller + Payout dürfen nur settlen, wenn das Ergebnis entscheidbar ist: `ResolutionOutcome` gesetzt **oder** ein + Outcome-Preis des Marktes ≥ 0.98 (dann gilt das Outcome mit Preis ≥ 0.98 als Gewinner, z. B. via erweitertem + `MarketOutcomeHelper`). Ist der Markt „resolved", aber nichts entscheidbar (Preise nicht eingerastet) → + **Position offen lassen** (kein Payout zu 0!). + +**Abnahme:** Zwei neue rote Invarianten-Tests in `PositionPnLEngineTests.cs` müssen grün werden, ohne die +Assertions zu ändern: +- `RecalculateTraderPositionsAsync_ResolvedMarketWithoutResolutionOutcome_PaysWinnerViaSnappedPrice` + (aktuell: RealizedPnl −40 statt +60) +- `RecalculateTraderPositionsAsync_ClosedButUnresolvedMarket_DoesNotBookPrematurePayout` + (aktuell: Position wird zu 0 ausgebucht statt offen zu bleiben) + +### A9. Trader-Namen aus der Activity-API übernehmen (Suche nach Benutzername) + +**Problem:** Manuell hinzugefügte (und über Markt-Trades entdeckte) Trader behalten für immer den +Platzhalter-Namen `0x2005d16a...` — die Suche findet sie nur über die Adresse, nicht über den Polymarket-Namen +(Beispiel: `0x2005d16a84ceefa912d4e380cd32e7ff827875ea` heißt auf Polymarket „RN1"). + +**Empirisch verifiziert:** Jede Zeile der `/activity`-Antwort enthält bereits `name` („RN1") und `pseudonym` +(„Scary-Edible") — die Felder werden nur nicht gemappt und damit bei jedem Sync weggeworfen. + +**Fix:** +1. `PolymarketTradeResponse`: `[JsonPropertyName("name")] public string? Name` und + `[JsonPropertyName("pseudonym")] public string? Pseudonym` ergänzen. +2. `Trade`: transientes Feld `[NotMapped] public string? TransientDisplayName` (analog `TransientWallet`); + im `PolymarketProvider`-Mapping mit `name`, Fallback `pseudonym`, befüllen. +3. `PollingWorker` und `TradeHistoryWorker`: nach dem Fetch, wenn ein nicht-leerer `TransientDisplayName` + vorliegt und vom aktuellen `DisplayName` abweicht → `trader.DisplayName` aktualisieren + (die Plattform ist die Quelle der Wahrheit; Platzhalter wie `0x…` heilen sich damit von selbst). +4. `DiscoveryService.ImportTraderAsync` (manuelles Hinzufügen): direkt beim Import die erste Activity-Seite + abrufen und den Namen setzen, statt des Wallet-Präfixes. +5. Die Suche (`TraderRepository.SearchAsync`) durchsucht `DisplayName` bereits — funktioniert danach automatisch + für Name **und** Adresse. + +**Empfohlener Beifang im selben Handgriff:** Die Antwort enthält auch `usdcSize` (echter Cash-Betrag — wichtig für +korrekte Split/Merge/Redeem-Buchungen) und `outcomeIndex` (robustes Outcome-Matching ohne Label-Vergleich). +Mindestens im Response-Model mit erfassen; Persistierung von `usdcSize` auf `Trade` (Migration) als eigener +kleiner Folge-Task. + +### A10. Watchlist end-to-end reparieren + `api()`-Helper-Bug (betrifft auch die KI-Analyse!) + +**Problem 1 — der zentrale JS-Helper verwirft alle Fetch-Optionen:** +```js +// app.js Zeile 133 — options-Parameter fehlt komplett: +async function api(endpoint) { + const res = await fetch(`${API_BASE}${endpoint}`); // ← { method: 'POST' } wird ignoriert! +``` +Jeder Aufruf der Form `api(url, { method: 'POST'|'DELETE' })` degradiert still zu einem **GET** → 404 → +der Fehler wird im catch geschluckt (`return null`). Betroffen: **Watchlist-Toggle** (Zeile ~449) und +**KI-Analyse-Button** (Zeile ~605). Deshalb „passiert nichts" beim Watchlist-Button — und deshalb steht überall +„Not analyzed yet". + +**Problem 2 — Backslash im Route-Template (gleiche Tippfehler-Familie wie in app.js):** +`TraderEndpoints.cs` Zeile ~52: `group.MapPost("\{id:int}/ai-analysis", ...)` — die Route ist mit dem +Backslash unerreichbar. Der KI-Analyse-Endpoint ist damit **serverseitig ebenfalls tot** (doppelt kaputt). + +**Problem 3 — es gibt keine Watchlist-Ansicht:** Der Toggle-Button existiert, aber nirgendwo im WebUI kann man +die beobachteten Trader sehen. `WatchlistService.GetAllAsync` existiert im Backend, hat aber weder Endpoint noch UI. + +**Fix:** +1. `api()`-Helper reparieren: + ```js + async function api(endpoint, options = {}) { + try { + const res = await fetch(`${API_BASE}${endpoint}`, options); + if (!res.ok) throw new Error(`HTTP ${res.status}`); + const text = await res.text(); + return text ? JSON.parse(text) : true; // leere 200er (Results.Ok()) nicht crashen lassen + } catch (err) { console.error(`API Error [${endpoint}]:`, err); return null; } + } + ``` +2. Route-Template fixen: `"\{id:int}/ai-analysis"` → `"/{id:int}/ai-analysis"`. + **Danach das gesamte Projekt nach weiteren Backslash-Pfaden absuchen** (`grep -rn '"\\{' src/` und + `grep -n '\\\\api' wwwroot/js/app.js`) — das ist jetzt der dritte Fall dieser Fehlerklasse. +3. Neuer Endpoint `GET /api/watchlist`: liefert Watchlist-Einträge mit Trader-Kerndaten + (TraderId, DisplayName, Label, Notes, CreatedAt, TotalPnl, WinRate, CopytradingScore). +4. WebUI: Nav-Punkt „Watchlist" + Seite mit Tabelle (Spalten wie Traders-Liste, plus Label/Notes und + Remove-Button; Zeilenklick öffnet die Detailseite). Der Toggle auf der Detailseite muss nach dem Klick + sichtbar den Zustand wechseln („Watchlist (Add)" ↔ „Watchlist (Remove)"). + +--- + +## Teil B — Datenreparatur ohne DB-Reset (Reihenfolge strikt einhalten) + +**Warum kein Reset nötig ist:** `Trades` = Rohdaten, größtenteils intakt. `TraderPositions`, `TraderAnalytics`, +`TraderDailySnapshots`, `TraderScores`, `TraderCategoryPerformances` = abgeleitet, lokal neu berechenbar. +`Markets/Events` = unvollständig (Erbe des 10%-Sampling-Bugs), aber per Marktsync günstig nachladbar. +Nur Pruned-/Compacted-Trader brauchen API-Re-Import — das ist eine kleine Teilmenge, nicht die ganze Import-Woche. + +### B0. Diagnose (Umfang bestimmen — SQL führt Richard selbst aus) +```sql +-- Wie viele Trader brauchen Deep-Resync? +SELECT COUNT(*) AS PrunedPositions, COUNT(DISTINCT TraderId) AS BetroffeneTrader +FROM TraderPositions WHERE IsHistoryPruned = 1; +SELECT COUNT(DISTINCT TraderId) FROM Trades WHERE PlatformTradeId LIKE 'COMPACT_%'; +-- Reconciliation-Backlog und Snapshot-Bestand +SELECT COUNT(*) FROM Trades WHERE MarketOutcomeId IS NULL; +SELECT COUNT(*) FROM TraderDailySnapshots; +``` + +### B1. Vorbereitung +Teil A deployen → Worker stoppen → **DB-Dump als Sicherung** → `RetentionSettings:Enabled = false`. + +### B2. Voll-Marktsync +Manuellen Market-Sync (inkl. geschlossener Märkte) einmal komplett durchlaufen lassen — schließt die Markt-Lücken, +an denen die Trade-Verlinkung bisher scheiterte, **und befüllt nach A8 erstmals `ResolutionOutcome`/korrektes +`IsResolved` für den gesamten Marktbestand**. Kostet nur Events-Endpoint-Requests (einige hundert), keine Import-Woche. + +Optional als Sofort-Backfill vor dem Sync (nutzt die bereits gespeicherten, eingerasteten Preise): +```sql +UPDATE Markets m +JOIN MarketOutcomes o ON o.MarketId = m.Id AND o.CurrentPrice >= 0.98 +SET m.ResolutionOutcome = o.Label +WHERE m.IsResolved = 1 AND (m.ResolutionOutcome IS NULL OR m.ResolutionOutcome = ''); +``` + +### B3. SQL-Reparatur der abgeleiteten Daten (Richard führt aus, Worker sind aus) +```sql +-- Vergiftete Fenster-Basis komplett verwerfen (heilt über den Fallback + neue Snapshots) +TRUNCATE TABLE TraderDailySnapshots; + +-- Abgeleitete Kategorien-Statistik neu aufbauen lassen +DELETE FROM TraderCategoryPerformances; + +-- Positionen mit vollständiger lokaler Historie löschen → Engine baut sie mit gefixtem Code neu +DELETE FROM TraderPositions WHERE IsHistoryPruned = 0; +-- (IsHistoryPruned = 1 absichtlich behalten: Konserve bis zum Deep-Resync in B5) + +-- Analytics nullen +UPDATE TraderAnalytics SET OverallPnL=0, PnL30d=0, PnL7d=0, PnL24h=0, + OverallWinRate=0, WinRate30d=0, WinRate7d=0, WinRate24h=0, + CurrentBalance=0, EstimatedBankroll=0, Trades30d=0, + CopytradingScore=0, CopytradingQualityScore=0, CopytradingCopyabilityScore=0; + +-- Re-Analyse für alle triggern +UPDATE Traders SET LastAnalyzedAt = NULL; + +-- Optional (einmalig, teuer — außerhalb der Stoßzeiten): Zähler geradeziehen +UPDATE Traders t SET TotalTrades = (SELECT COUNT(*) FROM Trades tr WHERE tr.TraderId = t.Id); +``` + +### B4. Worker starten, Backlog abarbeiten lassen +Reconciliation (jetzt gebatcht) verlinkt die Orphans; der `TraderAnalyticsWorker` rechnet alle Trader neu +(Fortschritt über Jobs-Seite/Analyze-Backlog-Button sichtbar). Die Fenster-PnL läuft anfangs über den Fallback und +gewinnt mit jedem Tag Snapshot-Präzision — nach 30 Tagen voll da. Das ist korrekt und erwartbar. + +### B5. Deep-Resync der betroffenen Trader +Für alle Trader aus B0 (Pruned/Compacted): `POST /api/jobs/deep-resync-pruned` in Häppchen (z. B. 25er-Batches), +über Tage verteilt — der RateLimiter drosselt automatisch. Watchlist-Trader zuerst. +Bis ein Trader dran war, zeigt er die (möglicherweise leicht verzerrte) Pruned-Konserve — akzeptierter Zwischenzustand. + +### B6. Retention wieder aktivieren +`RetentionSettings:Enabled = true`. Ab jetzt entsteht die Pruned-Konserve auf Basis der **korrekten** Engine — +zukünftiges Pruning ist damit verlustfrei im Sinne der PnL-Summen. + +### B7. Verifikation +1. Invarianten-SQL: + ```sql + -- Fenster-PnL darf für Trader ohne Trades im Fenster nicht = Lifetime sein + SELECT COUNT(*) FROM TraderAnalytics a + WHERE ABS(a.PnL30d) > 0 AND a.Trades30d = 0; + ``` +2. Plausibilitäts-Stichprobe gegen Polymarkets eigene Zahlen: kleiner Dev-Endpoint + `GET /api/dev/verify-positions/{traderId}`, der `GetTraderPositionsAsync` (Polymarkets `/positions` liefert + deren berechnete `size`/`avgPrice`/`percentPnl`) mit unseren `TraderPositions` vergleicht und Abweichungen + > 5 % listet. 10–20 aktive Trader stichproben. +3. 24 h Logs beobachten: keine 1213-Deadlocks, keine ERR-Kaskaden bei Shutdown, Jobs-Seite zeigt Durchsatz. + +--- + +## Teil C — Abnahmekriterien (gesamt) + +1. `dotnet test`: **18 grün + 1 skip** (inkl. der beiden A8-Tests), Assertions unverändert. +2. „Sync"/„Analyze" auf der Detailseite erzeugen sichtbare Einträge auf der Jobs-Seite, die auch abgearbeitet werden. +3. Nach B3/B4: kein Trader mehr mit `|PnL30d| > 0` bei `Trades30d = 0`; PnL30d/Total-Verhältnisse plausibel. +4. Quality Edge / Copyability auf der Detailseite ≠ 0 für analysierte Trader mit verlinkten Trades. +5. 24 h Betrieb ohne Deadlock-Errors und ohne Shutdown-Fehlerkaskaden. +6. Die Suche nach „RN1" findet den Trader `0x2005d16a84ceefa912d4e380cd32e7ff827875ea` (nach dessen nächstem Sync). *(A9)* +7. Watchlist: Toggle auf der Detailseite wechselt sichtbar den Zustand; die neue Watchlist-Seite listet die + beobachteten Trader; Remove funktioniert. *(A10)* +8. „Run Deep Analysis" (KI) füllt die AI Strategy Analysis auf der Detailseite tatsächlich. *(A10)* diff --git a/src/Predictalytics.Api/ApiConfiguration.cs b/src/Predictalytics.Api/ApiConfiguration.cs index a0e6b75..c9da848 100644 --- a/src/Predictalytics.Api/ApiConfiguration.cs +++ b/src/Predictalytics.Api/ApiConfiguration.cs @@ -31,6 +31,8 @@ public static class ApiConfiguration app.MapAlertEndpoints(); app.MapMarketEndpoints(); app.MapJobEndpoints(); + app.MapDevEndpoints(); + app.MapWatchlistEndpoints(); // Health check app.MapGet("/api/health", () => Results.Ok(new { Status = "OK", Timestamp = DateTime.UtcNow })); diff --git a/src/Predictalytics.Api/Endpoints/DevEndpoints.cs b/src/Predictalytics.Api/Endpoints/DevEndpoints.cs new file mode 100644 index 0000000..47ffe35 --- /dev/null +++ b/src/Predictalytics.Api/Endpoints/DevEndpoints.cs @@ -0,0 +1,97 @@ +using Microsoft.AspNetCore.Builder; +using Microsoft.AspNetCore.Http; +using Microsoft.AspNetCore.Routing; +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.DependencyInjection; +using Predictalytics.Domain.Entities; +using Predictalytics.Domain.Interfaces; +using Predictalytics.Infrastructure.Data; + +namespace Predictalytics.Api.Endpoints; + +public static class DevEndpoints +{ + public static void MapDevEndpoints(this IEndpointRouteBuilder routes) + { + var group = routes.MapGroup("/api/dev"); + + group.MapGet("/verify-positions/{traderId:int}", async (int traderId, AppDbContext db, IEnumerable providers, CancellationToken ct) => + { + var trader = await db.Traders.FindAsync(new object[] { traderId }, ct); + if (trader == null) return Results.NotFound("Trader not found"); + + var provider = providers.FirstOrDefault(p => p.Platform == trader.Platform && p.IsImplemented); + if (provider == null) return Results.BadRequest("Platform provider not found"); + + var apiPositions = await provider.GetTraderPositionsAsync(trader.PlatformUserId, ct); + var dbPositions = await db.TraderPositions + .Include(p => p.MarketOutcome) + .Where(p => p.TraderId == traderId && p.SharesHeld > 0) + .ToListAsync(ct); + + var apiPosDict = apiPositions.ToDictionary(p => p.AssetId ?? ""); + var dbPosDict = dbPositions.ToDictionary(p => p.MarketOutcome?.TokenId ?? ""); + + var mismatches = new List(); + + // Check DB -> API + foreach (var kvp in dbPosDict) + { + if (string.IsNullOrEmpty(kvp.Key)) continue; + if (!apiPosDict.TryGetValue(kvp.Key, out var apiPos)) + { + mismatches.Add(new { Token = kvp.Key, DB = kvp.Value.SharesHeld, API = 0, Reason = "Missing in API" }); + } + else if (Math.Abs(kvp.Value.SharesHeld - apiPos.Size) > 0.01m) + { + mismatches.Add(new { Token = kvp.Key, DB = kvp.Value.SharesHeld, API = apiPos.Size, Reason = "Size mismatch" }); + } + } + + // Check API -> DB + foreach (var kvp in apiPosDict) + { + if (string.IsNullOrEmpty(kvp.Key)) continue; + if (!dbPosDict.ContainsKey(kvp.Key)) + { + mismatches.Add(new { Token = kvp.Key, DB = 0, API = kvp.Value.Size, Reason = "Missing in DB" }); + } + } + + return Results.Ok(new + { + TraderId = traderId, + Mismatches = mismatches, + Match = mismatches.Count == 0, + DbPositionsCount = dbPosDict.Count, + ApiPositionsCount = apiPosDict.Count + }); + }); + + group.MapPost("/repair-db", async (AppDbContext db, CancellationToken ct) => + { + await db.Database.ExecuteSqlRawAsync(@" + -- 1. Fake PnL Trades löschen + DELETE FROM Trades WHERE Type = 4 AND Payout = 0 AND Size = 0; + + -- 2. Positionen löschen, da sie durch Sync neu aufgebaut werden + DELETE FROM TraderPositions; + + -- 3. Sync State von Tradern zurücksetzen (DeepResync forcieren) + UPDATE Traders SET + IsInitialImportComplete = 0, + LastTradesUpdatedAt = NULL, + LastPositionsUpdatedAt = NULL, + TotalPnl = 0, + WinRate = 0, + TotalTrades = 0, + EstimatedBankroll = 0; + + -- 4. Jobs abbrechen + UPDATE Jobs SET Status = 5 WHERE Status IN (1, 2); + ", ct); + + return Results.Ok("DB repaired. DeepResync needed."); + }); + } +} diff --git a/src/Predictalytics.Api/Endpoints/JobEndpoints.cs b/src/Predictalytics.Api/Endpoints/JobEndpoints.cs index 97c0469..7a063cf 100644 --- a/src/Predictalytics.Api/Endpoints/JobEndpoints.cs +++ b/src/Predictalytics.Api/Endpoints/JobEndpoints.cs @@ -54,6 +54,38 @@ public static class JobEndpoints return Results.Ok(job.Id); }); + group.MapPost("/deep-resync/{traderId:int}", async (int traderId, IJobRepository repo, CancellationToken ct) => + { + var job = new BackgroundJob + { + JobType = JobType.DeepResync, + Status = JobStatus.Pending, + TraderId = traderId + }; + await repo.AddAsync(job, ct); + return Results.Ok(job.Id); + }); + + group.MapPost("/deep-resync-pruned", async (Predictalytics.Infrastructure.Data.AppDbContext db, IJobRepository repo, CancellationToken ct) => + { + var traderIds = await Microsoft.EntityFrameworkCore.EntityFrameworkQueryableExtensions.ToListAsync( + db.TraderPositions.Where(tp => tp.IsHistoryPruned).Select(tp => tp.TraderId).Distinct(), ct); + + int count = 0; + foreach (var tid in traderIds) + { + var job = new BackgroundJob + { + JobType = JobType.DeepResync, + Status = JobStatus.Pending, + TraderId = tid + }; + await repo.AddAsync(job, ct); + count++; + } + return Results.Ok(new { Count = count }); + }); + group.MapPost("/analyze-backlog", async (int? take, IJobRepository repo, Predictalytics.Infrastructure.Data.AppDbContext db, CancellationToken ct) => { int batchSize = take ?? 50; diff --git a/src/Predictalytics.Api/Endpoints/TraderEndpoints.cs b/src/Predictalytics.Api/Endpoints/TraderEndpoints.cs index 0351840..a64cbf9 100644 --- a/src/Predictalytics.Api/Endpoints/TraderEndpoints.cs +++ b/src/Predictalytics.Api/Endpoints/TraderEndpoints.cs @@ -35,18 +35,7 @@ public static class TraderEndpoints await svc.SetManualOverrideAsync(id, score, ct); return Results.Ok(); }); - - group.MapPost("/{id:int}/refresh", async (int id, IAnalyticsService svc, CancellationToken ct) => - { - await svc.TriggerTradeSyncAsync(id, ct); - return Results.Ok(); - }); - group.MapPost("/{id:int}/force-analyze", async (int id, IAnalyticsService svc, CancellationToken ct) => - { - await svc.ForceAnalyzeTraderAsync(id, ct); - return Results.Ok(); - }); group.MapPost("/{id:int}/watchlist", async (int id, WatchlistService svc, CancellationToken ct) => { diff --git a/src/Predictalytics.Api/Endpoints/WatchlistEndpoints.cs b/src/Predictalytics.Api/Endpoints/WatchlistEndpoints.cs new file mode 100644 index 0000000..45306eb --- /dev/null +++ b/src/Predictalytics.Api/Endpoints/WatchlistEndpoints.cs @@ -0,0 +1,29 @@ +using Microsoft.AspNetCore.Mvc; +using Predictalytics.Application.Services; + +namespace Predictalytics.Api.Endpoints; + +public static class WatchlistEndpoints +{ + public static void MapWatchlistEndpoints(this IEndpointRouteBuilder routes) + { + var group = routes.MapGroup("/api/watchlist").WithTags("Watchlist"); + + group.MapGet("/", async (WatchlistService svc, CancellationToken ct) => + { + var list = await svc.GetAllAsync(ct); + return Results.Ok(list.Select(w => new + { + w.TraderId, + w.Trader.DisplayName, + w.Trader.PlatformUserId, + w.Trader.TotalPnl, + w.Trader.WinRate, + CopytradingScore = w.Trader.CurrentScore?.CombinedScore ?? 0, + w.Label, + w.Notes, + CreatedAt = w.AddedAt + })); + }); + } +} diff --git a/src/Predictalytics.Api/wwwroot/index.html b/src/Predictalytics.Api/wwwroot/index.html index 877e2b5..76c962b 100644 --- a/src/Predictalytics.Api/wwwroot/index.html +++ b/src/Predictalytics.Api/wwwroot/index.html @@ -26,6 +26,10 @@ Traders + + + Watchlist + Markets @@ -187,6 +191,32 @@ + +
+
+

Watchlist

+
+
+
+ + + + + + + + + + + + + + +
TraderPlatformScoreWin RatePnLLabelAddedActions
+
+
+
+
@@ -286,6 +316,9 @@ + diff --git a/src/Predictalytics.Api/wwwroot/js/app.js b/src/Predictalytics.Api/wwwroot/js/app.js index f800681..3e0a4e9 100644 --- a/src/Predictalytics.Api/wwwroot/js/app.js +++ b/src/Predictalytics.Api/wwwroot/js/app.js @@ -27,6 +27,7 @@ document.querySelectorAll('.nav-item[data-page]').forEach(item => { if (page === 'alerts') loadAlerts(); if (page === 'markets') loadMarkets(); if (page === 'jobs') loadJobs(); + if (page === 'watchlist') loadWatchlist(); }); }); @@ -106,14 +107,7 @@ async function manualAddTrader() { } } -async function manualUpdateTrader(id) { - const res = await fetch(`/api/traders/${id}/refresh`, { method: 'POST' }); - if (res.ok) { - alert('Sync triggered manually. Data will update in a few minutes.'); - } else { - alert('Failed to trigger sync.'); - } -} + let currentPlatform = 'All'; let currentSort = 'default'; @@ -134,14 +128,16 @@ function refreshActivePage() { if (activePage === 'page-dashboard') loadDashboard(); else if (activePage === 'page-traders') loadTraders(); else if (activePage === 'page-markets') loadMarkets(); + else if (activePage === 'page-watchlist') loadWatchlist(); } // ─── API Helpers ─── -async function api(endpoint) { +async function api(endpoint, options = {}) { try { - const res = await fetch(`${API_BASE}${endpoint}`); + const res = await fetch(`${API_BASE}${endpoint}`, options); if (!res.ok) throw new Error(`HTTP ${res.status}`); - return await res.json(); + const text = await res.text(); + return text ? JSON.parse(text) : true; } catch (err) { console.error(`API Error [${endpoint}]:`, err); return null; @@ -349,6 +345,41 @@ async function loadTraders() { `).join(''); } +async function loadWatchlist() { + const data = await api(`/api/watchlist`); + if (!data) return; + + const tbody = document.getElementById('watchlistBody'); + tbody.innerHTML = data.map(w => ` + + +
${w.displayName.substring(0, 2).toUpperCase()}
+
+ ${w.displayName}
+ ${w.platformUserId.substring(0, 8)}... +
+ + ${w.platform ?? 'Unknown'} + ${Number(w.copytradingScore || 0).toFixed(1)} + ${fmt.pct(w.winRate)} + ${fmt.pnl(w.totalPnl)} + ${w.label || ''} + ${new Date(w.createdAt).toLocaleDateString()} + + + + + + `).join('') || 'Your watchlist is empty.'; +} + +async function removeFromWatchlist(id) { + if (confirm('Remove this trader from watchlist?')) { + await api(`/api/traders/${id}/watchlist`, { method: 'DELETE' }); + loadWatchlist(); + } +} + // ─── Alerts Page ─── async function loadAlerts() { const data = await api('/api/alerts?count=50'); @@ -429,22 +460,21 @@ async function viewTrader(id) { const syncBtn = document.getElementById('btn-sync-trader'); if (syncBtn) { - syncBtn.onclick = () => manualUpdateTrader(id); + syncBtn.onclick = () => queueHistorySync(id); + } + + const deepSyncBtn = document.getElementById('btn-deep-resync-trader'); + if (deepSyncBtn) { + deepSyncBtn.onclick = () => { + if (confirm("Are you sure? This will delete all compacted trades and reset positions, then fetch all historical trades via pagination.")) { + queueDeepResync(id); + } + }; } const analyzeBtn = document.getElementById('btn-analyze-trader'); if (analyzeBtn) { - analyzeBtn.onclick = async () => { - analyzeBtn.disabled = true; - analyzeBtn.textContent = '...'; - try { - await api(`/api/traders/${id}/force-analyze`, { method: 'POST' }); - alert('Deep Analysis queued! Please wait a moment and then refresh.'); - } finally { - analyzeBtn.disabled = false; - analyzeBtn.textContent = '⚙ Analyze'; - } - }; + analyzeBtn.onclick = () => queueTraderAnalysis(id); } const aiBtn = document.getElementById('btn-ai-analysis'); @@ -658,6 +688,15 @@ async function queueHistorySync(id) { } } +async function queueDeepResync(id) { + const res = await fetch(`/api/jobs/deep-resync/${id}`, { method: 'POST' }); + if (res.ok) { + alert('Deep Resync job queued successfully.'); + } else { + alert('Failed to queue deep resync.'); + } +} + async function queueTraderAnalysis(id) { const res = await fetch(`/api/jobs/analyze/${id}`, { method: 'POST' }); if (res.ok) { diff --git a/src/Predictalytics.Application.Tests/Services/PositionPnLEngineTests.cs b/src/Predictalytics.Application.Tests/Services/PositionPnLEngineTests.cs index ea902d2..0ab51cc 100644 --- a/src/Predictalytics.Application.Tests/Services/PositionPnLEngineTests.cs +++ b/src/Predictalytics.Application.Tests/Services/PositionPnLEngineTests.cs @@ -654,4 +654,128 @@ public class PositionPnLEngineTests Assert.Equal(0m, analytics.PnL24h); } } + + // ═════════════════════════════════════════════════════════════════════════ + // Invariant tests added 2026-07-10 (review round 5). + // Root cause found via live Gamma API: the response contains NO + // "resolution_outcome" and NO "resolved" field. Market.ResolutionOutcome is + // therefore ALWAYS NULL in our DB, and IsResolved is effectively just + // "closed". Winner detection must use the snapped outcomePrices (winner→1, + // loser→0, persisted in MarketOutcome.CurrentPrice) and/or + // "umaResolutionStatus". EXPECTED TO BE RED until fixed. + // ═════════════════════════════════════════════════════════════════════════ + + /// + /// Defect 11a: A resolved market whose ResolutionOutcome string is NULL + /// (which is ALL markets today) must still pay out winners. The winner is + /// identifiable by its snapped price (CurrentPrice ≈ 1). Booking payout 0 + /// for every winner is what currently makes every trader show + /// WinRate 0% and negative PnL. + /// + [Fact] + public async Task RecalculateTraderPositionsAsync_ResolvedMarketWithoutResolutionOutcome_PaysWinnerViaSnappedPrice() + { + // Arrange + var dbName = Guid.NewGuid().ToString(); + using (var db = CreateDbContext(dbName)) + { + var trader = new Trader { Id = 1, PlatformUserId = "0x1", DisplayName = "Trader 1" }; + var market = new Market + { + Id = 10, PlatformMarketId = 1L, Question = "Q?", + IsResolved = true, + ResolutionOutcome = null // ← reality: Gamma never delivers this field + }; + // Snapped prices after resolution: this outcome won. + market.Outcomes.Add(new MarketOutcome { Id = 100, MarketId = 10, Label = "Yes", TokenId = "t100", CurrentPrice = 1.00m }); + market.Outcomes.Add(new MarketOutcome { Id = 101, MarketId = 10, Label = "No", TokenId = "t101", CurrentPrice = 0.00m }); + db.Traders.Add(trader); + db.Markets.Add(market); + + db.Trades.Add(new Trade + { + Id = 10, TraderId = 1, DbMarketId = 10, MarketOutcomeId = 100, + Side = TradeSide.Buy, Price = 0.40m, Size = 100m, Amount = 40m, + ExecutedAt = DateTime.UtcNow.AddDays(-3) + }); + await db.SaveChangesAsync(); + } + + // Act + using (var db = CreateDbContext(dbName)) + { + var pnlEngine = new PositionPnLEngine(db, NullLogger.Instance); + await pnlEngine.RecalculateTraderPositionsAsync(1); + } + + // Assert: virtual payout 100 × (1.00 − 0.40) = +60 — NOT −40. + using (var db = CreateDbContext(dbName)) + { + var pos = await db.TraderPositions.SingleAsync(p => p.TraderId == 1 && p.MarketOutcomeId == 100); + Assert.Equal(60m, pos.RealizedPnl); + Assert.Equal(0m, pos.SharesHeld); + + var analytics = await db.TraderAnalytics.SingleAsync(a => a.TraderId == 1); + Assert.Equal(60m, analytics.OverallPnL); + } + } + + /// + /// Defect 11b: "closed" is NOT "resolved". The Gamma mapper sets + /// IsResolved = closed, so markets that closed for trading but are still + /// awaiting UMA resolution (prices NOT snapped, e.g. 0.70) currently get a + /// virtual payout of 0 → every open position is booked as a total loss + /// days before the real outcome is known. The engine must only settle a + /// position when the outcome is actually decidable (ResolutionOutcome set, + /// or prices snapped to 0/1); otherwise the position stays open with + /// unrealized PnL. + /// + [Fact] + public async Task RecalculateTraderPositionsAsync_ClosedButUnresolvedMarket_DoesNotBookPrematurePayout() + { + // Arrange + var dbName = Guid.NewGuid().ToString(); + using (var db = CreateDbContext(dbName)) + { + var trader = new Trader { Id = 1, PlatformUserId = "0x1", DisplayName = "Trader 1" }; + var market = new Market + { + Id = 10, PlatformMarketId = 1L, Question = "Q?", + IsResolved = true, // ← buggy mapper sets this for merely CLOSED markets + ResolutionOutcome = null + }; + // Prices NOT snapped → UMA has not resolved yet, outcome undecided. + market.Outcomes.Add(new MarketOutcome { Id = 100, MarketId = 10, Label = "Yes", TokenId = "t100", CurrentPrice = 0.70m }); + market.Outcomes.Add(new MarketOutcome { Id = 101, MarketId = 10, Label = "No", TokenId = "t101", CurrentPrice = 0.30m }); + db.Traders.Add(trader); + db.Markets.Add(market); + + db.Trades.Add(new Trade + { + Id = 10, TraderId = 1, DbMarketId = 10, MarketOutcomeId = 100, + Side = TradeSide.Buy, Price = 0.40m, Size = 100m, Amount = 40m, + ExecutedAt = DateTime.UtcNow.AddDays(-3) + }); + await db.SaveChangesAsync(); + } + + // Act + using (var db = CreateDbContext(dbName)) + { + var pnlEngine = new PositionPnLEngine(db, NullLogger.Instance); + await pnlEngine.RecalculateTraderPositionsAsync(1); + } + + // Assert: position must remain OPEN (no premature settlement at payout 0). + using (var db = CreateDbContext(dbName)) + { + var pos = await db.TraderPositions.SingleAsync(p => p.TraderId == 1 && p.MarketOutcomeId == 100); + Assert.Equal(100m, pos.SharesHeld); + Assert.Equal(0m, pos.RealizedPnl); + + // Unrealized: 100 × (0.70 − 0.40) = +30 + var analytics = await db.TraderAnalytics.SingleAsync(a => a.TraderId == 1); + Assert.Equal(30m, analytics.OverallPnL); + } + } } diff --git a/src/Predictalytics.Application/Services/DiscoveryService.cs b/src/Predictalytics.Application/Services/DiscoveryService.cs index b21f1b3..6dbbfde 100644 --- a/src/Predictalytics.Application/Services/DiscoveryService.cs +++ b/src/Predictalytics.Application/Services/DiscoveryService.cs @@ -64,11 +64,33 @@ public class DiscoveryService : IDiscoveryService var existing = await _traderRepo.GetByPlatformIdAsync(platform, platformUserId, ct); if (existing != null) return existing.Id; + if (string.IsNullOrEmpty(displayName) || displayName.StartsWith("0x", StringComparison.OrdinalIgnoreCase)) + { + var provider = _providers.FirstOrDefault(p => p.Platform == platform); + if (provider != null && provider.IsImplemented) + { + try + { + await _rateLimiter.WaitAsync(platform, ct); + var trades = await provider.GetTraderTradesAsync(platformUserId, 50, ct); + var realName = trades.FirstOrDefault(t => !string.IsNullOrEmpty(t.TransientDisplayName))?.TransientDisplayName; + if (!string.IsNullOrEmpty(realName)) + { + displayName = realName; + } + } + catch (Exception ex) + { + _logger.LogWarning(ex, "Failed to fetch real display name for {Wallet} during import", platformUserId); + } + } + } + var trader = new Trader { Platform = platform, PlatformUserId = platformUserId, - DisplayName = string.IsNullOrEmpty(displayName) ? platformUserId[..8] + "..." : displayName, + DisplayName = string.IsNullOrEmpty(displayName) ? (platformUserId.Length > 8 ? platformUserId[..8] + "..." : platformUserId) : displayName, IsAutoDiscovered = isAutoDiscovered, CreatedAt = DateTime.UtcNow }; diff --git a/src/Predictalytics.Domain/Entities/Trade.cs b/src/Predictalytics.Domain/Entities/Trade.cs index 3f7bece..60d9051 100644 --- a/src/Predictalytics.Domain/Entities/Trade.cs +++ b/src/Predictalytics.Domain/Entities/Trade.cs @@ -87,6 +87,7 @@ public class Trade /// [NotMapped] public string? TransientWallet { get; set; } + [NotMapped] public string? TransientDisplayName { get; set; } // ── Navigation ──────────────────────────────────────────────────────── public Trader Trader { get; set; } = null!; diff --git a/src/Predictalytics.Domain/Enums/JobType.cs b/src/Predictalytics.Domain/Enums/JobType.cs index d74a9a9..6ec9b61 100644 --- a/src/Predictalytics.Domain/Enums/JobType.cs +++ b/src/Predictalytics.Domain/Enums/JobType.cs @@ -4,5 +4,6 @@ public enum JobType { HistorySync, TraderAnalysis, - ContextEnrichment + ContextEnrichment, + DeepResync } diff --git a/src/Predictalytics.Domain/Interfaces/IJobRepository.cs b/src/Predictalytics.Domain/Interfaces/IJobRepository.cs index e96caa6..2bc7e70 100644 --- a/src/Predictalytics.Domain/Interfaces/IJobRepository.cs +++ b/src/Predictalytics.Domain/Interfaces/IJobRepository.cs @@ -14,4 +14,5 @@ public interface IJobRepository Task GetNextPendingJobAsync(JobType type, CancellationToken ct = default); Task AddAsync(BackgroundJob job, CancellationToken ct = default); Task UpdateAsync(BackgroundJob job, CancellationToken ct = default); + Task ResetHungJobsAsync(CancellationToken ct = default); } diff --git a/src/Predictalytics.Domain/Interfaces/IPlatformProvider.cs b/src/Predictalytics.Domain/Interfaces/IPlatformProvider.cs index 6092657..981a61a 100644 --- a/src/Predictalytics.Domain/Interfaces/IPlatformProvider.cs +++ b/src/Predictalytics.Domain/Interfaces/IPlatformProvider.cs @@ -21,6 +21,9 @@ public interface IPlatformProvider /// Fetch recent trades for a specific trader. Task> GetTraderTradesAsync(string platformUserId, int limit = 50, CancellationToken ct = default); + /// Fetch all historical trades for a specific trader via pagination. + Task> GetTradesPagedAsync(string platformUserId, int limitPerRequest = 500, CancellationToken ct = default); + /// Fetch current positions/holdings for a trader. Task> GetTraderPositionsAsync(string platformUserId, CancellationToken ct = default); @@ -54,7 +57,8 @@ public record TraderPositionInfo( decimal Size, decimal AveragePrice, decimal CurrentValue, - decimal PnlPercent + decimal PnlPercent, + string? AssetId = null ); /// diff --git a/src/Predictalytics.Infrastructure/Data/Repositories/JobRepository.cs b/src/Predictalytics.Infrastructure/Data/Repositories/JobRepository.cs index 3df463c..1a8ba23 100644 --- a/src/Predictalytics.Infrastructure/Data/Repositories/JobRepository.cs +++ b/src/Predictalytics.Infrastructure/Data/Repositories/JobRepository.cs @@ -52,4 +52,14 @@ public class JobRepository : IJobRepository _db.BackgroundJobs.Update(job); await _db.SaveChangesAsync(ct); } + + public async Task ResetHungJobsAsync(CancellationToken ct = default) + { + var cutoff = DateTime.UtcNow.AddMinutes(-15); + await _db.BackgroundJobs + .Where(j => j.Status == JobStatus.InProgress && j.StartedAt < cutoff) + .ExecuteUpdateAsync(s => s + .SetProperty(j => j.Status, JobStatus.Pending) + .SetProperty(j => j.StartedAt, (DateTime?)null), ct); + } } diff --git a/src/Predictalytics.Infrastructure/Data/Repositories/TradeRepository.cs b/src/Predictalytics.Infrastructure/Data/Repositories/TradeRepository.cs index 8b70567..ddcf5fa 100644 --- a/src/Predictalytics.Infrastructure/Data/Repositories/TradeRepository.cs +++ b/src/Predictalytics.Infrastructure/Data/Repositories/TradeRepository.cs @@ -68,7 +68,7 @@ public class TradeRepository : ITradeRepository t.TransactionHash = StringHelper.Truncate(t.TransactionHash, 66); } - foreach (var chunk in tradeList.Chunk(1000)) + foreach (var chunk in tradeList.Chunk(500)) { var sb = new System.Text.StringBuilder("INSERT IGNORE INTO Trades (PlatformTradeId, MarketId, AssetId, Outcome, Side, Price, Size, Amount, ExecutedAt, TransactionHash, TraderId, MarketOutcomeId, DbMarketId, Platform, IsContextEnriched) VALUES "); var parameters = new List(); @@ -99,8 +99,37 @@ public class TradeRepository : ITradeRepository parameters.Add(t.IsContextEnriched); } - var rowsInserted = await _db.Database.ExecuteSqlRawAsync(sb.ToString(), parameters.ToArray(), ct); - _logger.LogInformation("Inserted {RowsInserted} trades into the database.", rowsInserted); + int maxRetries = 3; + var backoffs = new[] { 250, 500, 1000 }; + for (int retry = 0; retry <= maxRetries; retry++) + { + try + { + var rowsInserted = await _db.Database.ExecuteSqlRawAsync(sb.ToString(), parameters.ToArray(), ct); + _logger.LogInformation("Inserted {RowsInserted} trades into the database.", rowsInserted); + break; + } + catch (Exception ex) + { + var mysqlEx = ex as MySqlConnector.MySqlException ?? ex.InnerException as MySqlConnector.MySqlException; + if (mysqlEx != null && (mysqlEx.Number == 1213 || mysqlEx.Number == 1205)) + { + if (retry == maxRetries) + { + _logger.LogError(ex, "Failed to insert {Count} trades after {Retries} retries due to deadlocks.", chunk.Length, maxRetries); + } + else + { + _logger.LogWarning("Deadlock detected during trade insertion. Retrying in {Delay}ms... (Attempt {Attempt}/{Max})", backoffs[retry], retry + 1, maxRetries); + await Task.Delay(backoffs[retry], ct); + } + } + else + { + throw; + } + } + } } } diff --git a/src/Predictalytics.Infrastructure/Providers/Azuro/AzuroProvider.cs b/src/Predictalytics.Infrastructure/Providers/Azuro/AzuroProvider.cs index d05a46f..2278538 100644 --- a/src/Predictalytics.Infrastructure/Providers/Azuro/AzuroProvider.cs +++ b/src/Predictalytics.Infrastructure/Providers/Azuro/AzuroProvider.cs @@ -22,6 +22,9 @@ public class AzuroProvider : IPlatformProvider public Task> GetTraderTradesAsync(string platformUserId, int limit = 50, CancellationToken ct = default) { using var _ = PlatformLogContext.Push(PlatformName); _logger.LogWarning("Provider not yet implemented"); return Task.FromResult>(Array.Empty()); } + public Task> GetTradesPagedAsync(string platformUserId, int limitPerRequest = 500, CancellationToken ct = default) + { using var _ = PlatformLogContext.Push(PlatformName); _logger.LogWarning("Provider not yet implemented"); return Task.FromResult>(Array.Empty()); } + public Task> GetTraderPositionsAsync(string platformUserId, CancellationToken ct = default) { using var _ = PlatformLogContext.Push(PlatformName); _logger.LogWarning("Provider not yet implemented"); return Task.FromResult>(Array.Empty()); } diff --git a/src/Predictalytics.Infrastructure/Providers/Limitless/LimitlessProvider.cs b/src/Predictalytics.Infrastructure/Providers/Limitless/LimitlessProvider.cs index 24e8eb0..52eb981 100644 --- a/src/Predictalytics.Infrastructure/Providers/Limitless/LimitlessProvider.cs +++ b/src/Predictalytics.Infrastructure/Providers/Limitless/LimitlessProvider.cs @@ -25,6 +25,11 @@ public class LimitlessProvider : IPlatformProvider _logger = logger; } + public Task> GetTradesPagedAsync(string platformUserId, int limitPerRequest = 500, CancellationToken ct = default) + { + throw new NotImplementedException(); + } + public async Task> GetTraderTradesAsync(string platformUserId, int limit = 50, CancellationToken ct = default) { using var _ = PlatformLogContext.Push(PlatformName); diff --git a/src/Predictalytics.Infrastructure/Providers/Polymarket/PolymarketApiClient.cs b/src/Predictalytics.Infrastructure/Providers/Polymarket/PolymarketApiClient.cs index 8ecdd91..5c9a32e 100644 --- a/src/Predictalytics.Infrastructure/Providers/Polymarket/PolymarketApiClient.cs +++ b/src/Predictalytics.Infrastructure/Providers/Polymarket/PolymarketApiClient.cs @@ -46,6 +46,39 @@ public class PolymarketApiClient return await ExecuteWithRetryAsync>(_client, url, "Data", ct) ?? []; } + public async Task> GetTradesPagedAsync(string walletAddress, int limit = 500, CancellationToken ct = default) + { + var allTrades = new List(); + long? endTimestamp = null; + + while (true) + { + var url = $"/activity?user={walletAddress}&limit={limit}"; + if (endTimestamp.HasValue) + { + url += $"&end={endTimestamp.Value}"; + } + + var batch = await ExecuteWithRetryAsync>(_client, url, "Data", ct); + if (batch == null || batch.Count == 0) + { + break; + } + + allTrades.AddRange(batch); + + if (batch.Count < limit) + { + break; + } + + var oldestTimestamp = batch.Min(t => t.Timestamp); + endTimestamp = oldestTimestamp - 1; + } + + return allTrades; + } + public async Task> GetMarketTradesAsync(string conditionId, int limit = 1000, int offset = 0, CancellationToken ct = default) { var url = $"/trades?condition_id={conditionId}&limit={limit}&offset={offset}"; diff --git a/src/Predictalytics.Infrastructure/Providers/Polymarket/PolymarketModels.cs b/src/Predictalytics.Infrastructure/Providers/Polymarket/PolymarketModels.cs index 963f9cf..1cfe6cc 100644 --- a/src/Predictalytics.Infrastructure/Providers/Polymarket/PolymarketModels.cs +++ b/src/Predictalytics.Infrastructure/Providers/Polymarket/PolymarketModels.cs @@ -59,10 +59,20 @@ public class PolymarketTradeResponse [JsonConverter(typeof(FlexibleDoubleConverter))] public double Size { get; set; } + [JsonPropertyName("usdcSize")] + [JsonConverter(typeof(FlexibleDoubleConverter))] + public double UsdcSize { get; set; } + [JsonPropertyName("price")] [JsonConverter(typeof(FlexibleDoubleConverter))] public double Price { get; set; } + [JsonPropertyName("outcomeIndex")] + public int? OutcomeIndex { get; set; } + + [JsonPropertyName("name")] public string? Name { get; set; } + [JsonPropertyName("pseudonym")] public string? Pseudonym { get; set; } + [JsonPropertyName("outcome")] public string Outcome { get; set; } = ""; [JsonPropertyName("timestamp")] @@ -137,7 +147,7 @@ public class GammaMarketResponse [JsonPropertyName("closed")] public bool Closed { get; set; } [JsonPropertyName("active")] public bool Active { get; set; } [JsonPropertyName("resolved")] public bool Resolved { get; set; } - [JsonPropertyName("resolution_outcome")] public string? ResolutionOutcome { get; set; } + [JsonPropertyName("umaResolutionStatus")] public string? UmaResolutionStatus { get; set; } [JsonPropertyName("negRisk")] public bool NegRisk { get; set; } [JsonPropertyName("closedTime")] public string? ClosedTime { get; set; } [JsonPropertyName("takerFee")] [JsonConverter(typeof(FlexibleDoubleConverter))] public double TakerFee { get; set; } diff --git a/src/Predictalytics.Infrastructure/Providers/Polymarket/PolymarketProvider.cs b/src/Predictalytics.Infrastructure/Providers/Polymarket/PolymarketProvider.cs index d12447c..32d45e8 100644 --- a/src/Predictalytics.Infrastructure/Providers/Polymarket/PolymarketProvider.cs +++ b/src/Predictalytics.Infrastructure/Providers/Polymarket/PolymarketProvider.cs @@ -56,13 +56,50 @@ public class PolymarketProvider : IPlatformProvider TransactionHash = r.TransactionHash?.ToLowerInvariant(), TraderId = 0, TransientWallet = wallet, + TransientDisplayName = !string.IsNullOrEmpty(r.Name) ? r.Name : r.Pseudonym }; }).ToList(); return mappedTrades.GroupBy(t => t.PlatformTradeId, StringComparer.OrdinalIgnoreCase).Select(g => g.First()).ToList(); } + public async Task> GetTradesPagedAsync(string platformUserId, int limitPerRequest = 500, CancellationToken ct = default) + { + using var _ = PlatformLogContext.Push(PlatformName); + _logger.LogDebug("Fetching all paged trades for {Wallet}", platformUserId); + var raw = await _api.GetTradesPagedAsync(platformUserId, limitPerRequest, ct); + _logger.LogInformation("Fetched {Count} trades total for {Wallet}", raw.Count, platformUserId); + var mappedTrades = raw.Select(r => + { + var wallet = !string.IsNullOrEmpty(r.User) ? r.User : + !string.IsNullOrEmpty(r.ProxyWallet) ? r.ProxyWallet : + platformUserId; + var side = MapTradeSide(r); + var sideStr = side.ToString().ToUpperInvariant(); + return new Trade + { + Platform = PlatformType.Polymarket, + PlatformTradeId = string.IsNullOrEmpty(r.TransactionHash) + ? $"{r.Timestamp}_{wallet}_{r.Asset}_{sideStr}" + : $"{r.TransactionHash.ToLowerInvariant()}_{wallet}_{r.Asset}_{sideStr}", + MarketId = r.ConditionId ?? "", + AssetId = r.Asset ?? "", + Outcome = r.Outcome ?? "", + Side = side, + Price = (decimal)r.Price, + Size = (decimal)r.Size, + Amount = (decimal)(r.Price * r.Size), + ExecutedAt = DateTimeOffset.FromUnixTimeSeconds(r.Timestamp).UtcDateTime, + TransactionHash = r.TransactionHash?.ToLowerInvariant(), + TraderId = 0, + TransientWallet = wallet, + TransientDisplayName = !string.IsNullOrEmpty(r.Name) ? r.Name : r.Pseudonym + }; + }).ToList(); + + return mappedTrades.GroupBy(t => t.PlatformTradeId, StringComparer.OrdinalIgnoreCase).Select(g => g.First()).ToList(); + } public async Task> GetMarketTradesAsync(string platformMarketId, int limit = 1000, CancellationToken ct = default) { var raw = await _api.GetMarketTradesAsync(platformMarketId, limit, 0, ct); @@ -90,6 +127,7 @@ public class PolymarketProvider : IPlatformProvider TransactionHash = r.TransactionHash, TraderId = 0, TransientWallet = wallet, + TransientDisplayName = !string.IsNullOrEmpty(r.Name) ? r.Name : r.Pseudonym }; }).ToList(); @@ -106,7 +144,8 @@ public class PolymarketProvider : IPlatformProvider return raw.Select(r => new TraderPositionInfo( platformUserId, r.Market, r.Question, r.Outcome, (decimal)r.Size, (decimal)r.AvgPrice, - (decimal)r.CurrentValue, (decimal)r.PercentPnl + (decimal)r.CurrentValue, (decimal)r.PercentPnl, + r.AssetId )).ToList(); } @@ -300,8 +339,6 @@ public class PolymarketProvider : IPlatformProvider EndDate = DateTime.TryParse(raw.EndDate ?? raw.EndDateIso, out var med) ? med : null, CreatedAt = DateTime.TryParse(raw.CreatedAt, out var mcd) ? mcd : DateTime.UtcNow, DbCreatedAt = DateTime.UtcNow, - IsResolved = raw.Resolved || raw.Closed, - ResolutionOutcome = raw.ResolutionOutcome, LastUpdatedAt = DateTime.UtcNow, FeeRateBps = (decimal)(raw.TakerFee * 10000), IsNegRisk = raw.NegRisk, @@ -313,6 +350,11 @@ public class PolymarketProvider : IPlatformProvider var outcomePrices = ParseJsonStringArray(raw.OutcomePrices); var tokenIds = ParseJsonStringArray(raw.ClobTokenIds); + bool pricesSnapped = false; + bool hasHighPrice = false; + bool allSnapped = outcomePrices.Count > 0; + string? snappedWinnerLabel = null; + for (int i = 0; i < outcomeLabels.Count; i++) { decimal price = 0; @@ -320,8 +362,6 @@ public class PolymarketProvider : IPlatformProvider decimal.TryParse(outcomePrices[i], System.Globalization.NumberStyles.Any, System.Globalization.CultureInfo.InvariantCulture, out price); - string tokenId = i < tokenIds.Count ? tokenIds[i] : ""; - var label = outcomeLabels[i]; if ((label.Equals("Yes", StringComparison.OrdinalIgnoreCase) || label.Equals("No", StringComparison.OrdinalIgnoreCase)) && !string.IsNullOrEmpty(raw.GroupItemTitle)) @@ -329,6 +369,21 @@ public class PolymarketProvider : IPlatformProvider label = $"{raw.GroupItemTitle} - {label}"; } + if (price <= 0.02m || price >= 0.98m) + { + if (price >= 0.98m) + { + hasHighPrice = true; + snappedWinnerLabel = label; + } + } + else + { + allSnapped = false; + } + + string tokenId = i < tokenIds.Count ? tokenIds[i] : ""; + market.Outcomes.Add(new MarketOutcome { Label = label, @@ -338,6 +393,11 @@ public class PolymarketProvider : IPlatformProvider }); } + pricesSnapped = allSnapped && hasHighPrice; + + market.IsResolved = raw.UmaResolutionStatus == "resolved" || (raw.Closed && pricesSnapped); + market.ResolutionOutcome = market.IsResolved ? snappedWinnerLabel : null; + return market; } diff --git a/src/Predictalytics.Infrastructure/Services/PositionPnLEngine.cs b/src/Predictalytics.Infrastructure/Services/PositionPnLEngine.cs index d76cd3e..9560ba3 100644 --- a/src/Predictalytics.Infrastructure/Services/PositionPnLEngine.cs +++ b/src/Predictalytics.Infrastructure/Services/PositionPnLEngine.cs @@ -177,16 +177,13 @@ public class PositionPnLEngine : IPositionPnLEngine break; case TradeSide.Redeem: - var market = trade.MarketOutcome.Market; - var isResolved = market?.IsResolved ?? false; - var resolutionOutcome = market?.ResolutionOutcome; - var isWinner = isResolved && Predictalytics.Domain.Helpers.MarketOutcomeHelper.IsWinningOutcome(trade.MarketOutcome, resolutionOutcome); - - var payout = isWinner ? 1.00m : 0.00m; - currentBalance += (pos.SharesHeld * payout); - pos.RealizedPnl += pos.SharesHeld * (payout - pos.AvgCost); - pos.SharesHeld = 0; - pos.AvgCost = 0; + if (TryDeterminePayout(trade.MarketOutcome, out var redeemPayout)) + { + currentBalance += (pos.SharesHeld * redeemPayout); + pos.RealizedPnl += pos.SharesHeld * (redeemPayout - pos.AvgCost); + pos.SharesHeld = 0; + pos.AvgCost = 0; + } break; case TradeSide.Split: @@ -241,13 +238,9 @@ public class PositionPnLEngine : IPositionPnLEngine { if (pos.SharesHeld > 0 && pos.MarketOutcome?.Market != null) { - var market = pos.MarketOutcome.Market; - if (market.IsResolved) + if (TryDeterminePayout(pos.MarketOutcome, out var virtualPayout)) { - var isWinner = Predictalytics.Domain.Helpers.MarketOutcomeHelper.IsWinningOutcome(pos.MarketOutcome, market.ResolutionOutcome); - var payout = isWinner ? 1.00m : 0.00m; - - var virtualPnlDelta = pos.SharesHeld * (payout - pos.AvgCost); + var virtualPnlDelta = pos.SharesHeld * (virtualPayout - pos.AvgCost); pos.RealizedPnl += virtualPnlDelta; pos.SharesHeld = 0; pos.AvgCost = 0; @@ -553,4 +546,25 @@ public class PositionPnLEngine : IPositionPnLEngine return result; } + + private bool TryDeterminePayout(MarketOutcome outcome, out decimal payout) + { + payout = 0m; + var market = outcome.Market; + if (market == null || !market.IsResolved) return false; + + if (!string.IsNullOrWhiteSpace(market.ResolutionOutcome)) + { + payout = Predictalytics.Domain.Helpers.MarketOutcomeHelper.IsWinningOutcome(outcome, market.ResolutionOutcome) ? 1.00m : 0.00m; + return true; + } + + if (market.Outcomes != null && market.Outcomes.Any(o => o.CurrentPrice >= 0.98m)) + { + payout = outcome.CurrentPrice >= 0.98m ? 1.00m : 0.00m; + return true; + } + + return false; + } } diff --git a/src/Predictalytics.Worker/Services/PollingWorker.cs b/src/Predictalytics.Worker/Services/PollingWorker.cs index 8cad0d0..8e1c223 100644 --- a/src/Predictalytics.Worker/Services/PollingWorker.cs +++ b/src/Predictalytics.Worker/Services/PollingWorker.cs @@ -64,10 +64,16 @@ public class PollingWorker : BackgroundService await rateLimiter.WaitAsync(trader.Platform, stoppingToken); var trades = await provider.GetTraderTradesAsync(trader.PlatformUserId, 100, stoppingToken); - // ── Filter: skip trades with empty PlatformTradeId ── var validTrades = trades .Where(t => !string.IsNullOrWhiteSpace(t.PlatformTradeId)) .ToList(); + + var updatedName = validTrades.FirstOrDefault(t => !string.IsNullOrEmpty(t.TransientDisplayName))?.TransientDisplayName; + if (!string.IsNullOrEmpty(updatedName) && !string.Equals(trader.DisplayName, updatedName, StringComparison.OrdinalIgnoreCase)) + { + trader.DisplayName = updatedName; + await traderRepo.UpdateAsync(trader, stoppingToken); + } if (validTrades.Count < trades.Count) { @@ -147,6 +153,11 @@ public class PollingWorker : BackgroundService trader.LastPolledAt = DateTime.UtcNow; await traderRepo.UpdateAsync(trader, stoppingToken); } + catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) + { + _logger.LogInformation("Cancellation requested during polling, stopping batch."); + break; + } catch (Exception ex) { _logger.LogError(ex, "Error polling trader {Trader} on {Platform}", diff --git a/src/Predictalytics.Worker/Services/TradeContextEnrichmentWorker.cs b/src/Predictalytics.Worker/Services/TradeContextEnrichmentWorker.cs index b07b3d4..40de277 100644 --- a/src/Predictalytics.Worker/Services/TradeContextEnrichmentWorker.cs +++ b/src/Predictalytics.Worker/Services/TradeContextEnrichmentWorker.cs @@ -107,6 +107,11 @@ public class TradeContextEnrichmentWorker : BackgroundService updatedCount++; } } + catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) + { + _logger.LogInformation("Cancellation requested during enrichment, stopping batch."); + break; + } catch (Exception ex) { _logger.LogError(ex, "Failed to enrich asset {AssetId}", assetId); diff --git a/src/Predictalytics.Worker/Services/TradeHistoryWorker.cs b/src/Predictalytics.Worker/Services/TradeHistoryWorker.cs index cb39987..53f2696 100644 --- a/src/Predictalytics.Worker/Services/TradeHistoryWorker.cs +++ b/src/Predictalytics.Worker/Services/TradeHistoryWorker.cs @@ -31,6 +31,13 @@ public class TradeHistoryWorker : BackgroundService protected override async Task ExecuteAsync(CancellationToken stoppingToken) { _logger.LogInformation("📜 TradeHistoryWorker started (update cooldown: {Hours}h)", CooldownHours); + + using (var scope = _services.CreateScope()) + { + var jobRepo = scope.ServiceProvider.GetRequiredService(); + await jobRepo.ResetHungJobsAsync(stoppingToken); + } + await Task.Delay(15000, stoppingToken); while (!stoppingToken.IsCancellationRequested) @@ -43,7 +50,11 @@ public class TradeHistoryWorker : BackgroundService using (var scope = _services.CreateScope()) { var jobRepo = scope.ServiceProvider.GetRequiredService(); - activeJob = await jobRepo.GetNextPendingJobAsync(Predictalytics.Domain.Enums.JobType.HistorySync, stoppingToken); + activeJob = await jobRepo.GetNextPendingJobAsync(Predictalytics.Domain.Enums.JobType.DeepResync, stoppingToken); + if (activeJob == null) + { + activeJob = await jobRepo.GetNextPendingJobAsync(Predictalytics.Domain.Enums.JobType.HistorySync, stoppingToken); + } var repo = scope.ServiceProvider.GetRequiredService(); @@ -109,12 +120,35 @@ public class TradeHistoryWorker : BackgroundService using var platformCtx = PlatformLogContext.Push(provider.PlatformName); await rateLimiter.WaitAsync(trader.Platform, ct); - bool isInitial = !trader.IsInitialImportComplete; - _logger.LogInformation("{Trader}: Starting {Type} sync", trader.DisplayName, isInitial ? "INITIAL FULL" : "INCREMENTAL"); + bool isDeepResync = activeJob != null && activeJob.JobType == Predictalytics.Domain.Enums.JobType.DeepResync; + bool isInitial = !trader.IsInitialImportComplete || isDeepResync; + _logger.LogInformation("{Trader}: Starting {Type} sync", trader.DisplayName, isDeepResync ? "DEEP RESYNC" : (isInitial ? "INITIAL FULL" : "INCREMENTAL")); - var fetchedTrades = await provider.GetTraderTradesAsync(trader.PlatformUserId, TradesPerFetch, ct); + if (isDeepResync) + { + var db = scope.ServiceProvider.GetRequiredService(); + await Microsoft.EntityFrameworkCore.RelationalDatabaseFacadeExtensions.ExecuteSqlRawAsync(db.Database, "DELETE FROM Trades WHERE TraderId = {0} AND PlatformTradeId LIKE 'COMPACT_%'", t.Id); + await Microsoft.EntityFrameworkCore.RelationalDatabaseFacadeExtensions.ExecuteSqlRawAsync(db.Database, "DELETE FROM TraderPositions WHERE TraderId = {0}", t.Id); + } + + IReadOnlyList fetchedTrades; + if (isDeepResync) + { + fetchedTrades = await provider.GetTradesPagedAsync(trader.PlatformUserId, 500, ct); + } + else + { + fetchedTrades = await provider.GetTraderTradesAsync(trader.PlatformUserId, TradesPerFetch, ct); + } var validTrades = fetchedTrades.Where(tr => !string.IsNullOrWhiteSpace(tr.PlatformTradeId)).ToList(); + var updatedName = validTrades.FirstOrDefault(t => !string.IsNullOrEmpty(t.TransientDisplayName))?.TransientDisplayName; + if (!string.IsNullOrEmpty(updatedName) && !string.Equals(trader.DisplayName, updatedName, StringComparison.OrdinalIgnoreCase)) + { + trader.DisplayName = updatedName; + await traderRepo.UpdateAsync(trader, ct); + } + var fetchedTradeIds = validTrades.Select(tr => tr.PlatformTradeId).ToList(); var knownTradeIds = await tradeRepo.GetKnownPlatformTradeIdsAsync(trader.Platform, trader.Id, fetchedTradeIds, ct); @@ -230,6 +264,7 @@ public class TradeHistoryWorker : BackgroundService trader.IsInitialImportComplete = true; trader.LastTradesUpdatedAt = DateTime.UtcNow; trader.LastPolledAt = DateTime.UtcNow; + if (isDeepResync) trader.LastAnalyzedAt = null; await traderRepo.UpdateAsync(trader, ct); if (activeJob != null && activeJob.TraderId == trader.Id) diff --git a/src/Predictalytics.Worker/Services/TradeReconciliationWorker.cs b/src/Predictalytics.Worker/Services/TradeReconciliationWorker.cs index b7fc394..cf4b72f 100644 --- a/src/Predictalytics.Worker/Services/TradeReconciliationWorker.cs +++ b/src/Predictalytics.Worker/Services/TradeReconciliationWorker.cs @@ -88,10 +88,17 @@ public class TradeReconciliationWorker : BackgroundService var newEvents = new List(); foreach (var marketId in missingMarketIds) { - var market = await polyProvider.GetMarketAsync(marketId, ct); - if (market != null && market.Event != null) + try { - newEvents.Add(market.Event); + var market = await polyProvider.GetMarketAsync(marketId, ct); + if (market != null && market.Event != null) + { + newEvents.Add(market.Event); + } + } + catch (Exception ex) + { + _logger.LogError(ex, "Failed to fetch market {MarketId} during reconciliation", marketId); } } if (newEvents.Count > 0) @@ -108,15 +115,27 @@ public class TradeReconciliationWorker : BackgroundService .Distinct() .ToListAsync(ct); - reconciledCount = await Microsoft.EntityFrameworkCore.RelationalDatabaseFacadeExtensions.ExecuteSqlRawAsync(db.Database, @" - UPDATE Trades t - INNER JOIN MarketOutcomes o ON t.AssetId = o.TokenId - INNER JOIN Markets m ON o.MarketId = m.Id - SET t.MarketOutcomeId = o.Id, - t.Outcome = o.Label, - t.DbMarketId = m.Id - WHERE t.MarketOutcomeId IS NULL AND t.AssetId != ''; - ", ct); + int totalReconciled = 0; + int batchAffected; + do + { + batchAffected = await Microsoft.EntityFrameworkCore.RelationalDatabaseFacadeExtensions.ExecuteSqlRawAsync(db.Database, @" + UPDATE Trades t + JOIN ( + SELECT t2.Id, o.Id AS OutcomeId, o.Label, m.Id AS MarketDbId + FROM Trades t2 + JOIN MarketOutcomes o ON t2.AssetId = o.TokenId + JOIN Markets m ON o.MarketId = m.Id + WHERE t2.MarketOutcomeId IS NULL AND t2.AssetId != '' + LIMIT 5000 + ) x ON t.Id = x.Id + SET t.MarketOutcomeId = x.OutcomeId, t.Outcome = x.Label, t.DbMarketId = x.MarketDbId; + ", ct); + + totalReconciled += batchAffected; + if (batchAffected > 0) await Task.Delay(100, ct); + } while (batchAffected > 0); + reconciledCount = totalReconciled; if (reconciledCount > 0 && pairsToReset.Any()) { @@ -164,6 +183,10 @@ public class TradeReconciliationWorker : BackgroundService .ExecuteUpdateAsync(s => s.SetProperty(p => p.LastAnalyzedAt, (DateTime?)null), ct); } } + catch (OperationCanceledException) when (ct.IsCancellationRequested) + { + _logger.LogInformation("Cancellation requested during reconciliation, stopping batch."); + } catch (Exception ex) { _logger.LogError(ex, "Error executing bulk reconciliation SQL"); diff --git a/src/Predictalytics.Worker/Services/TradeRetentionWorker.cs b/src/Predictalytics.Worker/Services/TradeRetentionWorker.cs index d2656c7..1a430fb 100644 --- a/src/Predictalytics.Worker/Services/TradeRetentionWorker.cs +++ b/src/Predictalytics.Worker/Services/TradeRetentionWorker.cs @@ -44,6 +44,7 @@ public class TradeRetentionWorker : BackgroundService { await RunOptimizationAsync(stoppingToken); } + catch (OperationCanceledException) { break; } catch (Exception ex) { _logger.LogError(ex, "Error occurred executing TradeRetentionWorker cycle."); @@ -61,6 +62,13 @@ public class TradeRetentionWorker : BackgroundService var db = scope.ServiceProvider.GetRequiredService(); // Load configuration values + bool isEnabled = _config.GetValue("RetentionSettings:Enabled", true); + if (!isEnabled) + { + _logger.LogInformation("🧹 TradeRetentionWorker: Retention is disabled in configuration. Skipping optimization."); + return; + } + var retentionDays = _config.GetValue("RetentionSettings:RetentionDays", 90); var compactionDays = _config.GetValue("RetentionSettings:CompactionDays", 14); diff --git a/src/Predictalytics.Worker/Services/TraderAnalyticsWorker.cs b/src/Predictalytics.Worker/Services/TraderAnalyticsWorker.cs index 73d06f5..ef7ace4 100644 --- a/src/Predictalytics.Worker/Services/TraderAnalyticsWorker.cs +++ b/src/Predictalytics.Worker/Services/TraderAnalyticsWorker.cs @@ -88,6 +88,7 @@ public class TraderAnalyticsWorker : BackgroundService foreach (var id in traderIds) { + if (ct.IsCancellationRequested) break; try { using var traderScope = _services.CreateScope(); @@ -149,6 +150,11 @@ public class TraderAnalyticsWorker : BackgroundService await updateJobRepo.UpdateAsync(activeJob, ct); } } + catch (OperationCanceledException) when (ct.IsCancellationRequested) + { + _logger.LogInformation("Cancellation requested during analysis, stopping batch."); + break; + } catch (Exception ex) { _logger.LogError(ex, "Error recalculating positions/PnL for trader {TraderId}", id);