diff --git a/CLAUDE.md b/CLAUDE.md index b30c2e6..2955e66 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -75,6 +75,7 @@ the same three layers and the same provider seam — see the "Symbol detail + li **Naming convention:** - `Api*` — raw provider DTOs (provider-specific), e.g. `ApiFinnhubQuoteDto`. - `Domain*` — provider-agnostic, **no formatting**; what every provider AND the repository return (`DomainQuote`, `DomainInstrument`, `DomainCandleSeries`). +- `*Entity` — data-layer **storage models** that live BELOW the repository, INSIDE a data source (today just `QuoteEntity` in the cache data source). A structural mirror of the matching `Domain*` that **never escapes its data source**: the repository maps `Domain* ⇄ *Entity` at its boundary (`QuoteEntity.From(domain)` on write, `entity.ToDomainQuote()` on read). Lets the stored shape diverge later (storage-only fields, a DB representation for a DB-backed source) without touching domain or surfaces. (Data → domain is the allowed dependency direction, so an `*Entity` may reference its `Domain*` for mapping — never the reverse.) - `Ui*` — presentation; the ONLY place `FormatPrice()`/`FormatChange()`/SVG live (`UiQuote`, `UiCandleSeries`). - `AssetCategory`, `ChartRange`, `CandleInterval` (enums) stay **unprefixed** — shared vocabulary across all layers. @@ -82,7 +83,10 @@ the same three layers and the same provider seam — see the "Symbol detail + li | File | Role | |---|---| -| `Data/MarketRepository.cs` | **The coordinator the UI depends on.** Routes each `DomainInstrument` to the first `IMarketDataProvider` whose `Supports(AssetCategory)` matches, fans out concurrently, merges into one order-preserving list. Also `SearchAsync(query)` — fans out the free-text lookup to every provider and merges/dedupes by symbol — and `GetCandlesAsync(instrument, range)` — routes the chart history to the first supporting provider (no fan-out; one instrument). No provider for a category → `IsValid:false` quote / invalid candle series. **`ActiveProviders()`**: if any provider is `IsExclusive` (the mock in Demo mode), ALL three operations route to the exclusive set only — so an exclusive source wins everywhere, including the search fan-out. | +| `Data/MarketRepository.cs` | **The coordinator the UI depends on**, and now the **orchestration layer** that owns the shared quote cache + polling. Routes each `DomainInstrument` to the first `IMarketDataProvider` whose `Supports(AssetCategory)` matches, fans out concurrently, merges into one order-preserving list. Also `SearchAsync(query)` and `GetCandlesAsync(instrument, range)`. No provider for a category → `IsValid:false` quote / invalid candle series. **`ActiveProviders()`**: any `IsExclusive` provider (the Demo mock) → ALL operations route to it alone. **Owns one `IQuoteCacheDataSource`** (in-memory by default; an injectable ctor overload `MarketRepository(IQuoteCacheDataSource, providers[])` for a future DB-backed cache): `RefreshAsync` **writes through** every fetch (so the cache fills with NO surface changes), and **`ObserveQuotes(...)`** returns a cache-backed `IObservable>` — for a fixed instrument list OR a membership `StateFlow` (Rx `CombineLatest` + `Switch`), **delivered via `ObserveOn(TaskPoolScheduler)`** so a surface's `RaiseItemsChanged` is never invoked while Rx holds the combiner gate lock (**the deadlock fix** — see the orchestration section) — so surfaces OBSERVE instead of fetching, and two surfaces on the same symbol can't drift. **Owns the single poll loop:** one process-lifetime `PollTicker` subscription whose `RefreshObserved` re-fetches the union of currently-OBSERVED instruments each tick — "observed" tracked by a refcounted registry that `ObserveQuotes` Register/Unregisters on subscribe/dispose (so a hidden surface stops being polled). The demo-flip handler `Clear()`s the cache then refills the observed set from the new source. **All priced quote surfaces (the `PricedListPage` trio + both dock bands) now OBSERVE through here**; only the candle chart stays off the cache. See "Shared quote cache + repository orchestration". | +| `Data/IQuoteCacheDataSource.cs` | The **shared observable quote cache** abstraction (an interface, so the in-memory impl now can be swapped for a DB-backed one later without touching the repository or surfaces). **Stores `QuoteEntity` (the data-layer storage model), NOT `DomainQuote`** — `MarketRepository` maps `QuoteEntity ⇄ DomainQuote` at its boundary, so the storage model never escapes the data source. `Get(symbol)` (sync snapshot, null if uncached), `Observe(symbol)` → a per-symbol `StateFlow` (replays current value, **null until the first fetch**, then pushes changes), `Upsert(quote, keepLastGood=true)` (write-through; **the SINGLE home for keep-last-good** — a transient invalid quote won't overwrite a cached valid one; pass `false` to overwrite unconditionally), `Clear()` (reset every entry to null **keeping observers subscribed** — for a demo source flip). Keyed by `WatchlistStore.Normalize`. | +| `Data/InMemoryQuoteCacheDataSource.cs` | The **real, live** in-memory `IQuoteCacheDataSource` (un-stubbed): a `Lock`-guarded per-symbol `MutableStateFlow` dictionary (lazily created, keyed by `WatchlistStore.Normalize`), value-equality distinct-until-changed (an unchanged re-fetch doesn't re-emit → no churn/flicker). `Get` snapshots under the lock; `Upsert` decides the value under the lock (**the SINGLE home for keep-last-good**) then `Update()`s **outside** it; `Clear()` resets every value to null keeping observers subscribed. **It was briefly stubbed to a no-op after the favorites dock appeared to deadlock CmdPal — but the cache's lock was never the culprit** (`Update()` already fans out outside the lock). The real hang was the host's `RaiseItemsChanged` being delivered **under the Rx combiner gate** in `ObserveQuotes`; fixed there via **`ObserveOn`**, so this cache stays the simple, correct version. Swap for a DB-backed `IQuoteCacheDataSource` later via the repo's injectable ctor — no repository/UI change. | +| `Data/QuoteEntity.cs` | The cache data source's **storage model** (`*Entity` layer): a structural mirror of `DomainQuote` that lives BELOW the repository and **never escapes the data source**. `IQuoteCacheDataSource` stores/emits `QuoteEntity`; `MarketRepository` maps at its boundary — `QuoteEntity.From(domain)` on write-through, `entity.ToDomainQuote()` on read. Identical to `DomainQuote` today, kept separate so the stored shape can diverge later (storage-only fields, a DB representation) without touching domain or surfaces. | | `Data/IMarketDataProvider.cs` | ONE data source: `bool Supports(AssetCategory)` + **`bool IsExclusive`** (default `false`; true → MarketRepository routes every operation to this provider alone — how the Demo-mode mock takes precedence everywhere) + `GetQuotesAsync(instruments, ct)` + `SearchAsync(query, ct)` (free-text symbol lookup → `DomainInstrument`s, **identity only, no prices**; a provider that can't search returns `[]`) + `GetCandlesAsync(instrument, ChartRange, ct)` (price history for the detail chart; **default interface method** returns an invalid series, so non-candle providers opt out for free). | | `Data/TwelveData/TwelveDataMarketDataProvider.cs` | **Primary provider when its key is set** (`Supports` Stock+Crypto+Currency, **gated on `HasTwelveDataApiKey`** → first-match routing falls back to Finnhub/Frankfurter when unset). One API for all three classes; **`/time_series` candles are free-tier**, so charts render on a free key. `GetQuotesAsync` **batches all symbols into one `/quote`** call (tight 8/min limit); `SearchAsync` = `/symbol_search` (results normalized back to neutral symbols). See Twelve Data specifics. | | `Data/Finnhub/FinnhubMarketDataProvider.cs` | Stock+Crypto provider (`Supports` Stock+Crypto), used when no Twelve Data key is set. Maps `ApiFinnhubQuoteDto` → `DomainQuote`; `SearchAsync` calls `/search` (US equities only — see Finnhub specifics). | @@ -91,8 +95,8 @@ the same three layers and the same provider seam — see the "Symbol detail + li | `Data/MockMarketDataProvider.cs` | The **Demo-mode** data source — quotes/candles for every asset class so the whole UI works with no key/network. Registered **first**; `Supports()` is **gated on the `DemoMode` setting** and **`IsExclusive`** also returns it, so in demo mode the repository routes **everything (quotes, candles, search)** to the mock alone — it wins everywhere and no live provider is hit. Off → serves nothing, live providers take over (`SearchAsync` self-gates to `[]` so it stays out of the live search fan-out). **Serves ANY symbol:** a small hand-tuned `Seed` (headline symbols, incl. GBP `HSBA.L`) overrides on top of `SynthesizeQuote` — a stable-hash, category-/currency-inferred quote for any other ticker, so an off-seed real holding still prices. **Search** matches a curated **`Catalog`** of ~120 *real* instruments (never fabricated symbols — a faked match could be persisted and break when demo is off). See "Demo mode (offline testing)". | | `Data/InstrumentCatalog.cs` | Static `DomainInstrument` defaults — the **first-run seed** for `WatchlistStore` (no longer always-shown; removable once seeded). | | `Settings/WatchlistStore.cs` | JSON-persisted tracked instruments, each carrying **two independent flags** `InWatchlist`/`IsFavorite` (favorites = the dock subset). Stores **full `DomainInstrument` identity** so searched non-catalog symbols re-price; an entry with both flags false is dropped. Source-gen `WatchlistItem`/`WatchlistJsonContext` → `market_watchlist.json`; seeds `InstrumentCatalog` on first run, else migrates the legacy `market_favorites.json` (old pins → watchlisted **and** favorited). Exposes its `Watchlist`/`Favorites` subsets as **observable `StateFlow`s** (each mutation calls `PublishState()` → re-publishes both); pages and the dock **subscribe** and re-render themselves. Replaces the old `FavoritesStore` + `Watchlist.cs`. | -| `Helpers/StateFlow.cs` | A tiny **Kotlin-StateFlow analog**, now a **thin wrapper over `System.Reactive`'s `BehaviorSubject`** (the migration off the old hand-rolled version is **done** — AOT/trim is off, so the Rx dependency is taken warning-free; see the Rx done bullet). Used by the **state holders** (`WatchlistStore`, `MarketSettingsManager`) — observable *state with a current value*, which `IObservable` (no `Value`) and a bare `BehaviorSubject` (no read-only face) can't express alone. Public API unchanged: `StateFlow` (read-only `Value` + `Subscribe` with **replay-on-subscribe** + distinct-until-changed; a `Subscribe(onNext, replayOnSubscribe:false)` overload opts out of the replay via `Skip(1)` — used by the priced pages' secondary flows, e.g. `HasAnyApiKey`) and writable `MutableStateFlow` (`Update`). `SetValue` does **source-side** distinct-until-changed (an equal value isn't pushed, so re-publishing an unchanged subset doesn't wake its subscribers); `BehaviorSubject.OnNext` fans handlers out **outside** its lock (handlers re-read the store safely). The old hand-rolled **subscriber-count seam** (`OnActive`/`OnInactive` + the refcount + a custom `Subscription` wrapper) was **removed** once its only user, `PollTicker`, went pure Rx — Rx's `Publish().RefCount()` is the proper home for that, and Rx subscriptions are already idempotent on dispose. Plus `InstrumentListComparer` (symbol-sequence dedup for the store's two flows). | -| `Helpers/PollTicker.cs` | The **live-price poll ticker** — now **pure Rx** (no longer a `StateFlow`). A process-wide singleton `IObservable` = `Observable.Generate(...)` (self-rescheduling timer; per-step delay re-read from `MarketSettingsManager` each iteration so interval/on-off applies without reload, `0` = off idles on a 30 s re-check and the tick is filtered out) wrapped in **`Publish().RefCount()`** — the WhileSubscribed seam that starts the loop on the first subscriber and tears it down on the last (Generate never completes, so disposal is silent and resubscribe restarts cleanly). `Defer`/`Finally` log the start/stop at the refcount edges. Surfaces attach via the **guarded** `PollTicker.Subscribe(onTick)` (the raw stream is private): the handler is wrapped in try/catch (swallow + `Log.Error`) so a throwing tick handler can't trip Rx's `SafeObserver` into tearing down the shared multicast stream (which would kill polling for every surface). Handlers still offload heavy work via `Task.Run`. | +| `Helpers/StateFlow.cs` | A tiny **Kotlin-StateFlow analog**, now a **thin wrapper over `System.Reactive`'s `BehaviorSubject`** (the migration off the old hand-rolled version is **done** — AOT/trim is off, so the Rx dependency is taken warning-free; see the Rx done bullet). Used by the **state holders** (`WatchlistStore`, `MarketSettingsManager`) — observable *state with a current value*, which `IObservable` (no `Value`) and a bare `BehaviorSubject` (no read-only face) can't express alone. Public API unchanged: `StateFlow` (read-only `Value` + `Subscribe` with **replay-on-subscribe** + distinct-until-changed; a `Subscribe(onNext, replayOnSubscribe:false)` overload opts out of the replay via `Skip(1)` — used by the priced pages' secondary flows, e.g. `HasAnyApiKey`) and writable `MutableStateFlow` (`Update`). `SetValue` does **source-side** distinct-until-changed (an equal value isn't pushed, so re-publishing an unchanged subset doesn't wake its subscribers); `BehaviorSubject.OnNext` fans handlers out **outside** its lock (handlers re-read the store safely). The old hand-rolled **subscriber-count seam** (`OnActive`/`OnInactive` + the refcount + a custom `Subscription` wrapper) was **removed** once its only user, `PollTicker`, went pure Rx — Rx's `Publish().RefCount()` is the proper home for that, and Rx subscriptions are already idempotent on dispose. Also exposes **`AsObservable()`** (the raw `BehaviorSubject` stream, replays on subscribe) so the quote cache's per-symbol flows compose with Rx `CombineLatest`/`Switch` in `MarketRepository.ObserveQuotes`. Plus `InstrumentListComparer` (symbol-sequence dedup for the store's two flows). | +| `Helpers/PollTicker.cs` | The **live-price poll ticker** — now **pure Rx** (no longer a `StateFlow`). A process-wide singleton `IObservable` = `Observable.Generate(...)` (self-rescheduling timer; per-step delay re-read from `MarketSettingsManager` each iteration so interval/on-off applies without reload, `0` = off idles on a 30 s re-check and the tick is filtered out) wrapped in **`Publish().RefCount()`** — the WhileSubscribed seam that starts the loop on the first subscriber and tears it down on the last (Generate never completes, so disposal is silent and resubscribe restarts cleanly). `Defer`/`Finally` log the start/stop at the refcount edges. Surfaces attach via the **guarded** `PollTicker.Subscribe(onTick)` (the raw stream is private): the handler is wrapped in try/catch (swallow + `Log.Error`) so a throwing tick handler can't trip Rx's `SafeObserver` into tearing down the shared multicast stream (which would kill polling for every surface). Handlers still offload heavy work via `Task.Run`. **Two subscriber kinds now:** `MarketRepository` holds ONE process-lifetime subscription that polls the *observed* set centrally, so every priced quote surface (the `PricedListPage` trio + both dock bands) observes the cache and no longer subscribes here. The **only** direct subscriber left is the symbol-detail **chart** (a separate candle path with no cache layer — it re-fetches its visible range each tick). | | `Helpers/HttpRetry.cs` | **Shared 429 back-off** at the HTTP seam. `SendAsync(send, tag, ct)` takes a request **thunk** (each retry must re-issue a fresh `HttpResponseMessage`) and, on HTTP `429`, honors a short `Retry-After` else backs off `1s`→`2s` (max 3 attempts, **bails** if the wait would exceed an `8s` cap — per-minute windows don't clear in seconds, so hammering only burns quota). Returns the final response (success / non-429 error / surviving 429) so callers inspect status exactly as before — a **drop-in** at each `GetAsync`. Also the single choke point that feeds `RateLimitSignal` (2xx → off, surviving 429 → on). Used by Finnhub + Twelve Data (not keyless Frankfurter). | | `Helpers/RateLimitSignal.cs` | Process-wide **"are we throttled" flag** — a `MutableStateFlow` behind a read-only `StateFlow` (same state-holder idiom as `WatchlistStore`/`MarketSettingsManager`). `ReportRateLimited()`/`ReportSuccess()` flip it (distinct-until-changed); priced surfaces **subscribe** and re-render. Intentionally **global, not per-symbol** (a free-tier limit is key-wide). | | `Helpers/RateLimitHint.cs` | The **rate-limited banner row** (parallels `ApiKeyHint.StatusRow()`): `Row()` returns an amber "Rate-limited — showing last known prices" `ListItem` (Enter = **`NoOpCommand`**, purely informational — it's the default-selected first row, so it must not navigate) when `RateLimitSignal` is set **and** the `ShowRateLimitErrors` setting is on, else `null`. Pinned to the **top** so it's seen without scrolling: `PricedListPage.GetItems` inserts it at index 0; `SearchPage.SearchItems` inserts it at index 1 (just under the Enter-to-search action, which must stay first). | @@ -231,6 +235,95 @@ cascading CS0534 ("does not implement … `GetTypeInfo`") onto **every** context ## Current Status / Next Steps +- **Done (latest): fixed an out-of-order race in the async-projection observe path (Portfolio screen + Portfolio + dock band).** On a cold cache `ObserveQuotes` emits progressive *partial* snapshots (CombineLatest fills the + set symbol-by-symbol, e.g. `[BABA]` before `[BABA, SPY]` lands). The observe handler was launched + **fire-and-forget** (`.Subscribe(q => _ = OnQuotesChangedAsync(q))`), and for the surfaces that `await` an + async projection (the Portfolio screen + dock band both `await CurrencyConverter.PrimeAsync`) that lambda + returns at the first `await` — so **two handlers run concurrently** and a stale partial could finish LAST and + overwrite the full total. Caught it in logs: the dock rolled up `$7,429.79` (SPY+BABA, correct) and + `$139.49` (BABA-only) at the same millisecond, the partial able to win and stick until the next poll tick. + `ObserveOn`'s serialization did **not** protect this — it serializes the *synchronous* part of `OnNext`, but + the fire-and-forget escapes it past the first `await`. **Fix:** project each emission through the handler with + **`Select(q => Observable.FromAsync(ct => OnQuotesChangedAsync(q, ct))).Switch()`** instead of fire-and-forget + — `Switch` cancels the prior projection's `CancellationToken` the instant a newer emission arrives, and the + handler **checks `ct` before painting** (`if (ct.IsCancellationRequested) return;` after the projection) and + swallows `OperationCanceledException` as the expected supersede path; `ct` is forwarded into `PrimeAsync` so a + superseded FX fetch stops early. Applied in the **base `PricedListPage`** (covers all three priced screens at + once) and in **`PortfolioDockPage`**; `OnQuotesProjectingAsync` gained a `ct` param (only `PortfolioPage` + overrides it). The base `GetItems` already snapshots `_quotes` into a local (no torn read); the dock's + `GetItems` was given the same snapshot. **Watchlist/Favorites are behaviorally unchanged**: their projection + is the default `Task.CompletedTask`, which never yields, so their handler still runs fully synchronously and + in order under `Switch` (the `ct` guard never skips a paint for them). **Decision (recorded as the escape + hatch, NOT built):** the *principled* de-drift fix is a shared **`PortfolioRollup`** observable (a single + resolved-portfolio value both dock + page observe, mirroring the per-symbol quote cache one level up) — but + the `Switch` fix already removes the wrong-value bug, residual dock↔page divergence is transient-only and + self-correcting (they share the quote cache + `CurrencyConverter`), and a shared rollup would pull + `PortfolioPage` off its `PricedListPage` base (re-adding chrome). Build clean (0 warnings). ✅ + **Live-verified on-device** (Portfolio screen + dock band roll up correctly; Watchlist/Favorites unaffected). +- **Done (previous): migrated the `PricedListPage` trio (`WatchlistPage`/`FavoritesPage`/`PortfolioPage`) to + pure cache observers — the QUOTE-cache migration is now COMPLETE (every priced quote surface observes).** + Migrating the shared base moved all three at once. The visible-lifecycle hook is now a single + `Repository.ObserveQuotes(membership StateFlow)` subscription, and the per-surface fetch/poll machinery is + **deleted**: the fields (`_priceCache`/`_cacheLock`/`_snapshot`/`_received`/`_fetchGeneration`/ + `_lastFullPriceTicks`) and the methods (`OnInstrumentsChanged`/`Fetch`/`LoadQuotes`/`PollRefresh`/ + `RefreshStaleQuotes`/`OnDataSourceChanged`), plus the page's own `PollTicker` and `DemoModeChanged` + subscriptions and the keep-last-good guard — the repo owns the poll loop, the demo flip, and + stale-on-subscribe, and the cache owns keep-last-good. Membership changes (incl. a quantity-only Portfolio + edit) re-project for free via the `Switch` overload; empty-vs-loading is read from the membership + `.Value.Count`. The **async-projection** hook generalizes from `PortfolioDockPage`: a new + `protected virtual Task OnQuotesProjectingAsync(quotes)` (default no-op) that `PortfolioPage` overrides to + `await CurrencyConverter.PrimeAsync` before the paint — **replacing** the old `OnPriceCacheUpdated` + + `SnapshotPricedQuotes` hooks (both deleted from the base). The manual **Refresh 🔄** row reroutes onto + `RefreshAsync(keepLastGood:false)` with its own spinner (distinct-until-changed means an unchanged refetch + won't re-emit, so the row owns its IsLoading lifecycle). `WatchlistPage`/`FavoritesPage` needed no changes + (Watchlist keeps its `RelistTriggers => [Favorites]` for the ★). The **only** quote surface left un-migrated + is intentionally none — the lone exception is the symbol-detail **chart**, a separate candle path + (`GetCandlesAsync`, no cache/observe layer) that stays on its own `PollTicker` subscription by design (one + detail page open at a time → no cross-surface drift). **Follow-up:** with no surface calling it, the now-dead + public `MarketRepository.GetQuotesAsync` was **removed** — the observe/poll paths fetch via `RefreshAsync` + (same `FetchMergeAsync` + `WriteThrough`, void return). Build clean (0 warnings); **not yet live-verified + on-device.** See "Shared quote cache + repository orchestration". +- **Done (previous round): migrated `PortfolioDockPage` to a pure cache observer + added the stale-on-subscribe + freshness primitive.** The **second dock band is now a pure observer** like the favorites band, and it proves + the **async-projection** pattern the favorites band didn't need: its observe handler `await`s + `CurrencyConverter.PrimeAsync` before rolling up, so the converted total lands in **one paint** (safe because + `ObserveQuotes` delivers off the Rx gate via `ObserveOn`). It dropped its own `PollTicker` + `DemoModeChanged` + subscriptions, its keep-last-good merge, and its `GetQuotesAsync` call; it observes `PortfolioStore.Instruments` + (which re-emits on **every** mutation, so a quantity-only edit re-rolls via `Switch`) and reads + `PortfolioStore.Positions` in the handler. Both dock bands now hold subscriptions in a `List` (not + a single overwritable field) so a double-`add` can't orphan a subscription — which would otherwise pin its + symbols to the repo poll loop forever. **Stale-on-subscribe primitive (`MarketRepository`):** every + write-through now stamps a per-symbol last-fetch time (`_lastFetchTicks` + a shared `WriteThrough` helper), and + a subscribe refreshes **null OR stale** symbols (cached but aged ≥ one `RefreshInterval` while unobserved) via + `NeedsFetchOnSubscribe`, not just missing ones — the central home for the freshness clock the priced pages keep + per-surface (`_lastFullPriceTicks`/`RefreshStaleQuotes`), so that machinery **deletes** when the trio migrates, + and it closes the "reopen shows a stale price until the next poll tick" gap on the two migrated dock bands. The + dead Portfolio-dock `SubscribeOn` comment in `MarketExtensionCommandsProvider` was corrected to the real + `ObserveOn`-in-the-repo fix. Build clean (0 warnings); **not yet live-verified on-device.** See "Shared quote + cache + repository orchestration". +- **Done (this round): shared quote cache + repository orchestration — the de-drift refactor (cache LIVE, favorites dock migrated + verified).** + A new **`IQuoteCacheDataSource`** (in-memory `InMemoryQuoteCacheDataSource` now — **real, no longer stubbed** — swappable for a + DB-backed one later) holds the latest `DomainQuote` per symbol as **observable** state. `MarketRepository` is + the orchestration layer: every fetch **writes through** to the cache, surfaces can **`ObserveQuotes(...)`** + (cache-backed `IObservable`, fixed set or membership `StateFlow`) instead of fetching, and the repository + runs the **single poll loop** (one lifetime `PollTicker` subscription refreshing the *observed* set each + tick) + the demo-flip Clear/refill. ⚠️ **The deadlock that briefly forced the stub is FIXED:** + `ObserveQuotes` now delivers via **`ObserveOn(TaskPoolScheduler)`**, so a surface's `RaiseItemsChanged` (a + blocking COM call into Command Palette's STA) is never invoked while Rx still holds the `CombineLatest`/`Switch` + gate lock — that gate/STA lock-order cycle was hanging CmdPal whenever the favorites band activated. (Two + dead ends, recorded so nobody re-tries them: pulling the band proved it was the trigger; **`SubscribeOn` did + NOT help** — it moves *subscription*, not *delivery*, and the host call happens during delivery.) **Migrated: + `FavoritesDockPage`** — now a **pure observer** (no `PollTicker`/demo subscription, no + `GetQuotesAsync`/`RefreshAsync`, no keep-last-good; it renders what the cache emits, and membership changes + come free via the `Switch` overload). Build clean (0 warnings). ✅ **Live-verified on-device**: the band + activates without hanging CmdPal and shows live prices. (`PortfolioDockPage` was migrated the following round — + see the "Done (previous round)" bullet above.) **NOT yet migrated** (still on per-surface fetch+poll, filling + the cache only via write-through): the `PricedListPage` trio (`WatchlistPage`/`FavoritesPage`/`PortfolioPage`), + plus the `SymbolDetailPage` chart (candles, separate path). **The dock bands are the template for the rest.** + See "Shared quote cache + repository orchestration" for the full design + the migration checklist. + **(Update: the `PricedListPage` trio is now migrated too — see the top "Done (latest)" bullet; only the + candle chart remains off the cache, deliberately.)** - **Done (this round): removed Sentry entirely — the app ships with NO telemetry.** The optional crash reporter (SDK init in `Program.cs`, the `CaptureException`/`CaptureMessage` calls in `Log.Error`, the `Sentry` package in the csproj + `Directory.Packages.props`) is **gone**. Rationale: it was a one-off in @@ -554,8 +647,129 @@ just **mutates the store** (no manual UI callback: the flows re-render every liv **shows a confirmation toast** (e.g. "Added AAPL to watchlist") while **keeping the palette open** so the user gets explicit feedback and can keep editing. +### Shared quote cache + repository orchestration (DONE — all quote surfaces migrated; only the candle chart stays off the cache by design) + +**Why.** Every priced surface used to fetch the **same** quotes independently and keep its **own** price +cache + its **own** `PollTicker` subscription. Because their fetches landed at different times, two surfaces +showing the same symbol (e.g. the favorites dock vs. another band) displayed **different** prices — they +drifted out of sync. The fix: one shared observable cache that every surface observes, with the repository +orchestrating all fetching/polling, so the same symbol is one cache entry everyone reads. + +**The pieces:** +- **`IQuoteCacheDataSource`** (`Data/`) — the shared store abstraction: per-symbol `StateFlow`, keyed by + `WatchlistStore.Normalize`. **Stores `QuoteEntity` (the `*Entity` storage model in `Data/QuoteEntity.cs`), not + `DomainQuote`** — the repository maps `QuoteEntity.From` on write-through and `entity.ToDomainQuote()` on read, so + the storage model never escapes the data source (the public `ObserveQuotes` still hands surfaces `DomainQuote`). `Observe(symbol)` replays the current value (null until first fetch); + `Upsert(quote, keepLastGood=true)` is the SINGLE home for keep-last-good; `Clear()` resets entries to null + while keeping observers subscribed. Interface, so a DB-backed impl swaps in later via the repo's injectable + ctor. ✅ **`InMemoryQuoteCacheDataSource` is the real, live impl now** (a `Lock`-guarded per-symbol + `MutableStateFlow` dictionary; `Update()` fires outside the lock). Value-equality + distinct-until-changed → an unchanged re-fetch doesn't re-emit (no churn/flicker). +- **`MarketRepository` orchestration** — owns one `IQuoteCacheDataSource`. `RefreshAsync` writes through (the single + cache-fill path, driven by the poll loop + the observe-subscribe fetch via `RefreshSafe`). The old public + `GetQuotesAsync` was **removed** — it had no callers once every surface moved to observe. + `ObserveQuotes(IReadOnlyList)` and + `ObserveQuotes(StateFlow<…>)` return a cache-backed list observable (`CombineLatest`; the membership overload + adds `Switch` so add/unstar re-projects — fetching only new symbols, dropping departed). `RefreshSafe` + wraps fire-and-forget refreshes (logged). +- **The `ObserveOn` delivery seam (the deadlock fix — do NOT remove).** Both `ObserveQuotes` overloads append + **`.ObserveOn(TaskPoolScheduler.Default)`** (the raw graph lives in a private `ObserveQuotesCore`; the + `StateFlow` overload `ObserveOn`s *after* `Switch` so the gate of `Switch` is covered too). Why it's load-bearing: + a cache write-through fans out synchronously, and `CombineLatest`/`Switch` call the subscriber's handler — + ultimately a surface's `RaiseItemsChanged`, a **blocking COM call into Command Palette's STA** — *while Rx + still holds the combiner gate lock*. The host then re-enters the extension and the gate↔STA lock order cycles + → **CmdPal hangs** (this is exactly what made the favorites band crash the host). `ObserveOn` hands each + emission to the scheduler, so surfaces are notified only **after** the gate locks are released — no host call + ever runs under a producer-side lock. The cache itself was a red herring (its `Update()` already fires + outside its lock). ⚠️ **`SubscribeOn` is NOT a substitute** — it moves *where you subscribe*, not *where + notifications fire*; the deadlock is in delivery. Consequence for migrators: **`OnQuotesChanged`-style + handlers run on a pool thread** (the toolkit marshals `RaiseItemsChanged`, same as the priced pages' `Task.Run`). +- **The single poll loop (the ticker moved into the repo)** — the repo holds ONE process-lifetime + `PollTicker.Subscribe(RefreshObserved)`. Each tick refreshes the **distinct union of currently-observed + instruments** in one batched fetch. "Observed" = a refcounted registry (`_observed`, a `Lock`-guarded + `Dictionary`) that `ObserveQuotes` **Register**s on subscribe and **Unregister**s + on dispose (via `Observable.Create` + `CompositeDisposable`/`Disposable.Create`) — so a hidden surface's + symbols stop being polled, and a symbol observed by two surfaces is fetched once. Carrying the full + `DomainInstrument` (not just the symbol) is deliberate: routing needs `Category`, which a bare + `BehaviorSubject.HasObservers` boolean couldn't give for a not-yet-loaded symbol. +- **Stale-on-subscribe (the central freshness clock)** — every write-through stamps a per-symbol last-fetch time + (`_lastFetchTicks`, via the `WriteThrough` helper shared by `GetQuotesAsync`/`RefreshAsync`). On subscribe, + `ObserveQuotesCore` refreshes the symbols `NeedsFetchOnSubscribe` flags: **never-cached** ones always (must + render something), plus — only when `AutoRefreshEnabled` — ones whose cached price **aged ≥ one + `RefreshInterval`** while unobserved. The poll loop keeps *observed* symbols stamped within the interval, so + only symbols that went quiet while their surfaces were hidden trip the stale check — that's what makes a hidden + surface refresh on reopen instead of showing a stale price until the next tick. Replaces the priced pages' + per-surface `_lastFullPriceTicks`/`RefreshStaleQuotes` clock. +- **Demo-mode flip** is centralized: the repo's `DemoModeChanged` handler `Clear()`s the cache then + `RefreshSafe`es the observed set from the new source (visible surfaces repaint at once; a hidden one + refetches on next open via `ObserveQuotes`' fetch-missing). + +**Migrated: both dock bands.** `FavoritesDockPage` (the simple template) is a **pure observer**: its whole +lifecycle is `_repository.ObserveQuotes(WatchlistStore.Instance.Favorites).Subscribe(OnQuotesChanged)` and +nothing else — no `PollTicker`/demo subscriptions, no `GetQuotesAsync`/`RefreshAsync`, no keep-last-good code. +`GetItems` disambiguates "no favorites" from "still loading" via `Favorites.Value.Count == 0` (an empty cache +emission with favorites present = spinner, not the empty-state row). Membership changes are free (the `Switch` +overload). ✅ Build clean; ✅ **live-verified on-device** — the band activates without hanging CmdPal and shows +live prices (this is the build that confirmed the `ObserveOn` deadlock fix). **`PortfolioDockPage`** is the +**async-projection template**: same pure-observer shape, but it observes `PortfolioStore.Instruments` and its +handler `await`s `CurrencyConverter.PrimeAsync` before rolling up `UiPortfolio` (one paint), reading +`PortfolioStore.Positions` for the quantities; it disambiguates empty-vs-loading from `Positions.Value.Count`. +Both bands hold subscriptions in a `List` (double-`add`-safe). Build clean; not yet live-verified. + +**Migrated: the `PricedListPage` trio (`WatchlistPage`/`FavoritesPage`/`PortfolioPage`) — the migration is now +COMPLETE for quotes.** Migrating the shared base moved all three: the visible-lifecycle hook is a single +`Repository.ObserveQuotes(_instruments).Subscribe(q => _ = OnQuotesChangedAsync(q))`, and the old per-surface +machinery is **deleted** (the `_priceCache`/`_cacheLock`/`_snapshot`/`_received`/`_fetchGeneration`/ +`_lastFullPriceTicks` fields and the `OnInstrumentsChanged`/`Fetch`/`LoadQuotes`/`PollRefresh`/ +`RefreshStaleQuotes`/`OnDataSourceChanged` methods, plus the page's `PollTicker` + `DemoModeChanged` +subscriptions and the keep-last-good guard). The base now stores the latest emission as `_quotes` +(`UiQuote[]?`, null pre-first-emission, dock-style) and disambiguates empty-vs-loading from +`_instruments.Value.Count`. The async-projection seam generalizes `PortfolioDockPage`'s shape into a base hook +`protected virtual Task OnQuotesProjectingAsync(quotes)` (default no-op) that `PortfolioPage` overrides to +`await CurrencyConverter.PrimeAsync` before the paint — this **replaced** the old `OnPriceCacheUpdated` + +`SnapshotPricedQuotes` hooks (deleted). The manual **Refresh 🔄** row reroutes onto +`RefreshAsync(keepLastGood:false)` and owns its own `IsLoading` spinner (since an unchanged refetch won't +re-emit). `WatchlistPage`/`FavoritesPage` were untouched (Watchlist keeps `RelistTriggers => [Favorites]` for +the ★). Build clean (0 warnings); **not yet live-verified on-device.** The symbol-detail **chart** is the only +surface still on its own `PollTicker` — a separate candle path with no cache (see "Transitional truth" below). + +**The migration checklist (now fully applied — kept as the record of how it was done):** +1. Replace its `MarketRepository.GetQuotesAsync` call + private price cache with a subscription to + `_repository.ObserveQuotes()`; render from the emitted `IReadOnlyList`. +2. **Delete** its `PollTicker.Subscribe(...)` and its `DemoModeChanged` subscription — the repo does both now. +3. **Delete** its keep-last-good guard (the cache owns it) and its manual reconcile (the `Switch` overload + handles add/remove). +4. Disambiguate empty-vs-loading from the membership count (as the dock does). +5. **Threading is already handled — don't add your own.** `ObserveQuotes` delivers via `ObserveOn` (see the + delivery-seam bullet above), so your emission handler runs on a **pool thread** with no Rx lock held — do + NOT wrap the subscribe in `Task.Run` and do NOT add `SubscribeOn`/`ObserveOn` yourself. Just `Subscribe` + and call `RaiseItemsChanged` from the handler (the toolkit marshals it). Re-introducing a synchronous + delivery path (or doing heavy work inside the handler under an Rx gate) is what re-creates the CmdPal hang. +6. ⚠️ **`PricedListPage` is shared by `WatchlistPage`/`FavoritesPage`/`PortfolioPage` — migrating the base + moves all three at once.** Two things to handle: the **freshness clock** (`_lastFullPriceTicks` / + `RefreshStaleQuotes`) is now **handled centrally** by the repo's stale-on-subscribe primitive + (`_lastFetchTicks` + `NeedsFetchOnSubscribe` refreshing null-OR-stale on subscribe), so the base can just + **DROP** its per-surface clock rather than re-wire it. **`OnPriceCacheUpdated`** (Portfolio's FX `PrimeAsync`) + should fold **into the observe handler** — `await` the prime before rendering, exactly as `PortfolioDockPage` + now does — not a separate fetch-completion hook (an unchanged, value-equality-deduped re-fetch won't re-emit, + but the prime keys off the emission's currencies, which only change when the set does anyway). `LeadingRows` + must recompute on each emission. +7. `PortfolioDockPage` is the **worked example** of the async-projection handler (`await PrimeAsync` then roll + up, one paint) — copy its shape for the Portfolio screen rather than reinventing it. + +**Transitional truth (now essentially resolved):** every priced QUOTE surface observes the cache, so there is +no longer a quote surface that double-fetches a shared symbol. The only thing still on its own `PollTicker` +subscription is the symbol-detail **chart**, and that's a separate **candle** path (`GetCandlesAsync`) with no +cache/observe layer — left that way deliberately (only one detail page is ever open, so no cross-surface drift +to fix). Migrating candles onto an `ObserveCandles`-style layer is possible but unplanned. + ### Live price polling (done) +> **Superseded for all quote surfaces** by "Shared quote cache + repository orchestration" above — the priced +> list pages and both dock bands no longer poll themselves; the repository's single loop polls the observed set. +> The per-surface model below now describes **only the symbol-detail chart** (a separate candle path that still +> subscribes `PollTicker` directly and keeps its own per-range cache + keep-last-good). + Priced surfaces used to fetch **once** when they became visible (the StateFlow replay-on-subscribe that drives the first load). They now also auto-refresh on a timer while visible: **default 10 min, settings- configurable** (0 = off). Built on the ticker's subscriber-count lifecycle (now Rx `Publish().RefCount()`). @@ -630,6 +844,12 @@ configurable** (0 = off). Built on the ticker's subscriber-count lifecycle (now ### Dock refresh on favorites change (done) +> **Implementation superseded by the cache migration** (see "Shared quote cache + repository orchestration"): +> `FavoritesDockPage` no longer calls `RefreshQuotes()` on a `Favorites` subscription — it observes +> `MarketRepository.ObserveQuotes(WatchlistStore.Instance.Favorites)`, and the `Switch` overload re-projects +> on membership change. The user-facing behavior below (a pinned band updates instantly on a favorite change) +> is unchanged; the mechanism is now the cache observe stream. + A pinned `FavoritesDockPage` updates the instant a favorite changes anywhere, instead of going stale until reopened. Implemented via the observable layer: diff --git a/MarketExtension/Data/IQuoteCacheDataSource.cs b/MarketExtension/Data/IQuoteCacheDataSource.cs new file mode 100644 index 0000000..ff82fb7 --- /dev/null +++ b/MarketExtension/Data/IQuoteCacheDataSource.cs @@ -0,0 +1,40 @@ +namespace MarketExtension; + +// A process-wide store of the latest QuoteEntity per symbol, exposed as OBSERVABLE state every surface +// can subscribe to — so two surfaces showing the same symbol read the SAME entry and can never drift +// apart (the out-of-sync prices the favorites dock and the screens show today). Keyed by +// WatchlistStore.Normalize, the one cache key used everywhere. +// +// This is the quote CACHE DATA SOURCE — one of the sources MarketRepository orchestrates over (alongside +// the IMarketDataProviders). It's an interface so the storage mechanism stays an implementation detail: +// the in-memory implementation can be swapped for a database-backed one later (the "local cache that will +// probably be a database layer") without touching MarketRepository or any UI surface. MarketRepository +// owns one instance, writes through to it on every fetch, and every priced surface OBSERVES it instead of +// fetching independently. +// +// Data layer: holds QuoteEntity (the storage model, no formatting), NOT DomainQuote. MarketRepository maps +// QuoteEntity <-> DomainQuote at its boundary, so the entity never escapes the data source; surfaces still +// observe DomainQuote off the repository and project it to UiQuote as they do today. +internal interface IQuoteCacheDataSource +{ + // Current cached quote for a symbol, or null if it was never fetched / has been cleared. + // Synchronous snapshot read. + QuoteEntity? Get(string symbol); + + // Per-symbol observable entry: replays the current value (null until the first fetch lands) on + // subscribe, then pushes each change. Lazily created on first access. Distinct-until-changed via + // QuoteEntity value equality (an identical re-fetch does NOT re-emit). This is the seam a surface + // waits on for the first fetch: subscribe → null (render a spinner) → the quote when Upsert lands. + StateFlow Observe(string symbol); + + // Write a freshly fetched quote through to the cache. keepLastGood:true (the default) drops a + // transient invalid quote (e.g. a 429 mapped to IsValid:false) when a valid quote is already + // cached — the SINGLE home for the "keep last good" guard currently copy-pasted across the priced + // surfaces. keepLastGood:false overwrites unconditionally (a hard refresh / data-source flip). + void Upsert(QuoteEntity quote, bool keepLastGood = true); + + // Reset every entry to null, KEEPING observer subscriptions live (so observers re-emit a "loading" + // state and a re-fetch refills). Used when the data SOURCE flips (demo mode) so prices from the old + // source can neither linger nor be preserved by keep-last-good. + void Clear(); +} diff --git a/MarketExtension/Data/InMemoryQuoteCacheDataSource.cs b/MarketExtension/Data/InMemoryQuoteCacheDataSource.cs new file mode 100644 index 0000000..7334a4a --- /dev/null +++ b/MarketExtension/Data/InMemoryQuoteCacheDataSource.cs @@ -0,0 +1,82 @@ +using System.Collections.Generic; +using System.Diagnostics.CodeAnalysis; +using System.Linq; +using System.Threading; + +namespace MarketExtension; + +// In-memory IQuoteCacheDataSource: one MutableStateFlow per normalized symbol, lazily created. +// The single shared store of live quotes; MarketRepository owns one instance for the whole process. +// Mirrors the WatchlistStore state-holder idiom — a Lock guarding a dictionary, snapshot under the +// lock then Update() OUTSIDE it (Update fans handlers out that may re-enter Get/Observe, and +// System.Threading.Lock is non-reentrant, so updating under the lock could deadlock a re-entrant +// handler). +// +// Deadlock note (the reason this was briefly stubbed): the hang was NOT this cache's lock — Update() +// already fires outside it. It was that the per-symbol fan-out reached a surface's RaiseItemsChanged +// while Rx still held the CombineLatest/Switch gate lock in MarketRepository.ObserveQuotes, and that +// host call re-entered Command Palette's STA → a lock-ordering cycle. That is fixed at the seam: +// ObserveQuotes now delivers via ObserveOn, so surfaces are notified only AFTER the Rx gate locks are +// released (no host call under a producer-side lock). This cache can therefore stay the simple, +// correct in-memory version. Swap it for a database-backed IQuoteCacheDataSource later via MarketRepository's +// injectable ctor overload — no repository or UI change. +[SuppressMessage("Reliability", "CA1001:Types that own disposable fields should be disposable", + Justification = "Owns per-symbol MutableStateFlow (BehaviorSubject-backed) entries for " + + "the life of the process via the single MarketRepository; they are intentionally never " + + "completed or disposed — mirrors the StateFlow singleton convention (see StateFlow.cs).")] +internal sealed class InMemoryQuoteCacheDataSource : IQuoteCacheDataSource +{ + private readonly Lock _lock = new(); + + // key = WatchlistStore.Normalize(symbol). Grows only with distinct observed symbols (watchlist + + // favorites + portfolio membership — tens), so no eviction is needed. + private readonly Dictionary> _entries = []; + + public QuoteEntity? Get(string symbol) + { + lock (_lock) + return _entries.TryGetValue(WatchlistStore.Normalize(symbol), out var flow) ? flow.Value : null; + } + + public StateFlow Observe(string symbol) + { + lock (_lock) + return GetOrCreate(WatchlistStore.Normalize(symbol)); + } + + public void Upsert(QuoteEntity quote, bool keepLastGood = true) + { + var key = WatchlistStore.Normalize(quote.Symbol); + + MutableStateFlow flow; + QuoteEntity? next = quote; + lock (_lock) + { + flow = GetOrCreate(key); + // Keep-last-good: a transient invalid quote must not overwrite a price that was fine. + // Decide the value to write under the lock (reads flow.Value); Update fires outside. + if (keepLastGood && !quote.IsValid && flow.Value is { IsValid: true }) + next = flow.Value; // Update below no-ops (distinct-until-changed) + } + + flow.Update(next); // fan handlers out OUTSIDE the lock (handlers re-read the cache) + } + + public void Clear() + { + List> flows; + lock (_lock) + flows = [.. _entries.Values]; // keep entries so observers stay subscribed; reset the values + + foreach (var flow in flows) + flow.Update(null); + } + + // Caller must hold _lock. + private MutableStateFlow GetOrCreate(string key) + { + if (!_entries.TryGetValue(key, out var flow)) + _entries[key] = flow = new MutableStateFlow(null); // default comparer = record value equality + return flow; + } +} diff --git a/MarketExtension/Data/MarketRepository.cs b/MarketExtension/Data/MarketRepository.cs index 1ae253a..1a0d2f1 100644 --- a/MarketExtension/Data/MarketRepository.cs +++ b/MarketExtension/Data/MarketRepository.cs @@ -1,5 +1,10 @@ +using System; +using System.Collections.Concurrent; using System.Collections.Generic; using System.Linq; +using System.Reactive.Concurrency; +using System.Reactive.Disposables; +using System.Reactive.Linq; using System.Threading; using System.Threading.Tasks; @@ -11,20 +16,78 @@ namespace MarketExtension; // results are merged into one list in the original instrument order. Instruments no provider can // serve become invalid placeholders. Adding a data source = pass another provider to the // constructor; nothing else changes. -internal sealed class MarketRepository(params IMarketDataProvider[] providers) +// +// It also owns the shared IQuoteCacheDataSource: every fetch writes through to it (so the cache fills with no +// surface changes), and surfaces can OBSERVE a set of instruments as one cache-backed list observable +// instead of each fetching independently (the orchestration API — ObserveQuotes/RefreshAsync). This is +// how two surfaces showing the same symbol stop drifting apart. +// +// Polling lives HERE, not per surface: the repository runs ONE poll loop (a single PollTicker +// subscription for its lifetime) that, each tick, refreshes the distinct union of all instruments +// currently being observed through the cache — so the observed prices stay fresh from one shared fetch. +// Every priced QUOTE surface (the Watchlist/Favorites/Portfolio pages + both dock bands) is now a pure +// observer through here — none fetch or poll themselves. The lone exception is the symbol-detail CHART, +// a separate candle path (GetCandlesAsync) with no cache/observe layer: it still subscribes PollTicker +// directly, which is fine since only one detail page is ever open (no cross-surface drift to fix). +internal sealed class MarketRepository { + private readonly IMarketDataProvider[] _providers; + private readonly IQuoteCacheDataSource _cacheSource; + + // Refcounted registry of the instruments currently being observed via ObserveQuotes — keyed by + // Normalize(symbol), each carrying the latest DomainInstrument (needed for provider routing) and a + // subscriber count. The poll loop refreshes the distinct union each tick; a symbol observed by two + // surfaces has count 2 and is fetched once. + private readonly Lock _observedLock = new(); + private readonly Dictionary _observed = []; + + // Per-symbol last fetch-ATTEMPT time (Environment.TickCount64), keyed by Normalize(symbol). Stamped on + // every write-through — success OR a kept-last-good failure — so it reflects "last time we tried", which + // is what staleness means. A subscribe reads it (NeedsFetchOnSubscribe) to tell a stale cached price from + // a fresh one and refresh only the stale ones: the central home for the freshness clock the priced pages + // used to keep per-surface (RefreshStaleQuotes). Bounded by tracked symbols (tens) like the cache itself, + // so no eviction is needed. + private readonly ConcurrentDictionary _lastFetchTicks = new(); + + // Default: an in-memory cache. The call site in MarketExtensionCommandsProvider uses this. + public MarketRepository(params IMarketDataProvider[] providers) + : this(new InMemoryQuoteCacheDataSource(), providers) { } + + // Injectable cache — for a future database-backed IQuoteCacheDataSource (and tests). No other code changes. + public MarketRepository(IQuoteCacheDataSource cacheSource, IMarketDataProvider[] providers) + { + _cacheSource = cacheSource; + _providers = providers; + + // A demo-mode flip swaps the data SOURCE, so every cached price is now wrong (it came from the other + // source). Clear it, then refresh whatever is currently observed so visible surfaces repaint from the + // new source at once (a hidden surface refetches on its next open via ObserveQuotes' fetch-missing). + // replay:false — construction is a no-op (the cache is empty then). Subscribed before any surface; + // never disposed (process-lifetime). + _ = MarketSettingsManager.Instance.DemoModeChanged.Subscribe(_ => OnDemoModeFlip(), replayOnSubscribe: false); + + // The single poll loop. While the repository is alive (the process lifetime), each PollTicker tick + // refreshes the union of currently-observed instruments (RefreshObserved) in one batched fetch and + // writes the results back through the cache, so every observing surface repaints in sync. PollTicker + // honors the refresh-interval setting and idles when it's off (and an empty observed set is a no-op), + // so a permanent subscription is safe. Never disposed — process-lifetime, like the cache itself. + _ = PollTicker.Subscribe(RefreshObserved); + } + // The providers that should serve right now. When any provider declares itself exclusive (e.g. the // mock in Demo mode), ONLY those serve — every operation routes to them alone, so the exclusive source // takes precedence everywhere, including the search fan-out that Supports() can't gate. Otherwise the // full ordered set, routed normally by first-match Supports. private IReadOnlyList ActiveProviders() { - var exclusive = providers.Where(p => p.IsExclusive).ToList(); - return exclusive.Count > 0 ? exclusive : providers; + var exclusive = _providers.Where(p => p.IsExclusive).ToList(); + return exclusive.Count > 0 ? exclusive : _providers; } - public async Task> GetQuotesAsync( - IReadOnlyList instruments, CancellationToken ct = default) + // Route → fetch (concurrent, one batch per provider) → merge into the caller's instrument order. + // The raw data path with no cache side effects, used by RefreshAsync. + private async Task> FetchMergeAsync( + IReadOnlyList instruments, CancellationToken ct) { // Route each instrument to the first ACTIVE provider that can serve its asset class. var active = ActiveProviders(); @@ -63,6 +126,194 @@ public async Task> GetQuotesAsync( return ordered; } + // Force a fetch now; results reach observers through the cache (the caller does not read the return). + // The single fetch seam — the manual Refresh row, the demo flip, the observe-subscribe fetch, and the + // poll loop all call it (via RefreshSafe for the fire-and-forget paths). + // keepLastGood:false = a HARD refresh that overwrites even with an invalid quote (e.g. after a source flip, + // where a stale "good" value would be wrong); the default keeps the last good price through a bad fetch. + public async Task RefreshAsync( + IReadOnlyList instruments, bool keepLastGood = true, CancellationToken ct = default) + { + var ordered = await FetchMergeAsync(instruments, ct).ConfigureAwait(false); + WriteThrough(ordered, keepLastGood); + } + + // Write a freshly fetched batch through to the cache AND stamp each symbol's last fetch-attempt time, so a + // later subscribe can distinguish a stale cached price from a fresh one (see NeedsFetchOnSubscribe). The + // single write-through path; called by RefreshAsync (the one fetch path). + private void WriteThrough(IReadOnlyList ordered, bool keepLastGood) + { + var now = Environment.TickCount64; + foreach (var quote in ordered) + { + _cacheSource.Upsert(QuoteEntity.From(quote), keepLastGood); // domain → storage at the boundary + _lastFetchTicks[WatchlistStore.Normalize(quote.Symbol)] = now; + } + } + + // Observe a FIXED set of instruments as one cache-backed list observable. Both public overloads append + // ObserveOn(TaskPoolScheduler) — this is the DEADLOCK FIX. Without it, a cache write-through fans out + // synchronously and the CombineLatest/Switch combiner calls the subscriber's handler (ultimately a + // surface's RaiseItemsChanged → a blocking COM call into Command Palette's STA) WHILE Rx still holds the + // combiner gate lock; the host then re-enters the extension and the gate/STA lock order cycles → hang. + // ObserveOn hands each emission to the scheduler, so subscribers are notified only AFTER the Rx gate + // locks are released — no host call ever runs under a producer-side lock. (Surfaces already expect + // off-host-thread notifications: the toolkit marshals RaiseItemsChanged, same as the priced pages' Task.Run.) + public IObservable> ObserveQuotes(IReadOnlyList instruments) + => ObserveQuotesCore(instruments).ObserveOn(TaskPoolScheduler.Default); + + // Membership-aware: re-projects (Switch) whenever the set changes — fetching newly-added symbols and + // dropping departed ones for free. ObserveOn AFTER Switch so the band is notified off the Switch gate too + // (Core is used here, not the public overload, so there's a single ObserveOn hop, after Switch). + public IObservable> ObserveQuotes( + StateFlow> instruments) + => instruments.AsObservable().Select(set => ObserveQuotesCore(set)).Switch() + .ObserveOn(TaskPoolScheduler.Default); + + // The cache-backed list observable WITHOUT the ObserveOn hop — the raw graph. On subscribe: register the + // instruments as observed (so the poll loop keeps them fresh) and fetch any not-yet-cached OR gone-stale + // symbols (write-through, fire-and-forget; see NeedsFetchOnSubscribe). Then CombineLatest the per-symbol + // cache flows into one ordered list + // (symbols not yet loaded — null entries — are dropped) that re-emits whenever any member quote changes. + // On dispose: unregister. Observable.Create (not Defer) so the dispose hook can unregister. The public + // overloads above add ObserveOn so the host is never called under the combiner gate (see that note). + private IObservable> ObserveQuotesCore(IReadOnlyList instruments) + => Observable.Create>(observer => + { + Register(instruments); + // Fetch the symbols this subscribe should refresh: never-cached ones (we must render something) + // plus — when auto-refresh is on — any whose cached price has aged past one interval while + // unobserved. The poll loop keeps OBSERVED symbols fresh; this catches a hidden surface reopening + // on a stale cache, the central replacement for the priced pages' old per-surface RefreshStaleQuotes. + var toFetch = instruments.Where(NeedsFetchOnSubscribe).ToList(); + if (toFetch.Count > 0) + { + // Re-derive the missing/stale split for the log only (NeedsFetchOnSubscribe stays the single + // decision): a to-fetch symbol is "missing" if it's still uncached, else it's a stale refresh. + var missing = toFetch.Count(i => _cacheSource.Get(i.Symbol) is null); + Log.Info("Repository", + $"observe subscribe: refreshing {toFetch.Count}/{instruments.Count} " + + $"({missing} missing, {toFetch.Count - missing} stale) [{string.Join(", ", toFetch.Select(i => i.Symbol))}]"); + _ = RefreshSafe(toFetch); // fire-and-forget; writes through → observers see it land + } + else if (instruments.Count > 0) + { + Log.Info("Repository", $"observe subscribe: all {instruments.Count} symbol(s) fresh in cache — no fetch"); + } + + IObservable> inner = instruments.Count == 0 + ? Observable.Return>([]) + : Observable + .CombineLatest(instruments.Select(i => _cacheSource.Observe(i.Symbol).AsObservable())) + // storage → domain at the boundary; OfType drops the null (not-yet-loaded) entries. + .Select(entities => (IReadOnlyList)entities + .OfType().Select(e => e.ToDomainQuote()).ToList()); + + return new CompositeDisposable( + inner.Subscribe(observer), + Disposable.Create(() => Unregister(instruments))); + }); + + // RefreshAsync with its exceptions swallowed + logged — for the fire-and-forget fetch-missing path and + // the poll loop, where a provider throw would otherwise become an unobserved task exception. + private async Task RefreshSafe(IReadOnlyList instruments) + { + try { await RefreshAsync(instruments).ConfigureAwait(false); } + catch (Exception ex) { Log.Error("Repository", "background refresh failed", ex); } + } + + // Should this instrument be (re)fetched at subscribe time? Never-cached → always (we must render + // something). Otherwise only when auto-refresh is on AND the cached price has aged past one refresh + // interval — i.e. it went stale while no surface was observing it. With auto-refresh off (interval 0) an + // already-cached price is left as-is, so "off" spends no calls. Reads the stamp WriteThrough sets. + private bool NeedsFetchOnSubscribe(DomainInstrument instrument) + { + if (_cacheSource.Get(instrument.Symbol) is null) + return true; + + var settings = MarketSettingsManager.Instance; + if (!settings.AutoRefreshEnabled) + return false; + + // Cached but somehow unstamped (e.g. a value that survived a cache Clear) → refresh to be safe. + if (!_lastFetchTicks.TryGetValue(WatchlistStore.Normalize(instrument.Symbol), out var ticks)) + return true; + + return TimeSpan.FromMilliseconds(Environment.TickCount64 - ticks) >= settings.RefreshInterval; + } + + // --- The poll loop + the observed-instrument registry ---------------------------------------------- + + // One PollTicker tick: refresh the distinct union of all currently-observed instruments in one batched + // fetch (keep-last-good via RefreshAsync's default, so a transient bad poll doesn't blank a good price). + // Nothing observed → nothing to do. + private void RefreshObserved() + { + var instruments = ObservedInstruments(); + if (instruments.Count == 0) + return; + Log.Info("Repository", $"poll tick: refreshing {instruments.Count} observed instrument(s)"); + _ = RefreshSafe(instruments); + } + + // Demo mode flipped the data source: drop the now-wrong cached prices, then refill the observed set from + // the new source so visible surfaces repaint immediately (post-Clear all entries are null, so this writes + // the new-source quotes through unconditionally). + private void OnDemoModeFlip() + { + _cacheSource.Clear(); + var observed = ObservedInstruments(); + Log.Info("Repository", + $"demo mode flipped — cache cleared; refreshing {observed.Count} observed instrument(s) from the new source"); + if (observed.Count > 0) + _ = RefreshSafe(observed); + } + + // Called when a surface subscribes to ObserveQuotes: ref up each instrument so the poll loop fetches it. + private void Register(IReadOnlyList instruments) + { + lock (_observedLock) + foreach (var instrument in instruments) + { + var key = WatchlistStore.Normalize(instrument.Symbol); + if (_observed.TryGetValue(key, out var entry)) + { + entry.Count++; + entry.Instrument = instrument; // keep the latest identity (name/category) + } + else + { + _observed[key] = new ObservedInstrument(instrument); + } + } + } + + // Called when an ObserveQuotes subscription is disposed: ref down, dropping symbols no surface observes. + private void Unregister(IReadOnlyList instruments) + { + lock (_observedLock) + foreach (var instrument in instruments) + { + var key = WatchlistStore.Normalize(instrument.Symbol); + if (_observed.TryGetValue(key, out var entry) && --entry.Count <= 0) + _observed.Remove(key); + } + } + + // Snapshot of the distinct instruments with at least one observer — what the poll loop refreshes. + private IReadOnlyList ObservedInstruments() + { + lock (_observedLock) + return [.. _observed.Values.Select(e => e.Instrument)]; + } + + // A registry entry: the latest observed identity for a symbol + how many subscriptions reference it. + private sealed class ObservedInstrument(DomainInstrument instrument) + { + public DomainInstrument Instrument { get; set; } = instrument; + public int Count { get; set; } = 1; + } + // Free-text instrument lookup. Fans out to every provider and merges their matches into one // list, deduped by symbol and preserving first-seen (provider) order. Identity only — callers // fetch quotes separately, so a search stays one call per provider regardless of match count. @@ -83,7 +334,7 @@ public async Task> SearchAsync( } // Historical candles for one instrument over a ChartRange — routed to the first provider that - // supports its asset class (mirrors GetQuotesAsync). No provider → an invalid series so the chart + // supports its asset class (mirrors the quote routing). No provider → an invalid series so the chart // renders an "unavailable" state rather than throwing. public async Task GetCandlesAsync( DomainInstrument instrument, ChartRange range, CancellationToken ct = default) diff --git a/MarketExtension/Data/QuoteEntity.cs b/MarketExtension/Data/QuoteEntity.cs new file mode 100644 index 0000000..8607c38 --- /dev/null +++ b/MarketExtension/Data/QuoteEntity.cs @@ -0,0 +1,32 @@ +namespace MarketExtension; + +// Data layer (the cache data source's storage model): a structural mirror of DomainQuote that lives BELOW +// the repository. IQuoteCacheDataSource stores and emits QuoteEntity, never DomainQuote, so the storage +// model never escapes the data source — MarketRepository maps QuoteEntity <-> DomainQuote at its boundary +// (From on write-through, ToDomainQuote on read). It's deliberately identical to DomainQuote today; keeping +// it a separate type means the cache's stored shape can diverge from the domain model later (storage-only +// fields, a different on-disk representation for a DB-backed source) without touching the domain layer or +// any surface. See DomainQuote for the field semantics (Currency, IsValid, the GBX/pence rule). +// +// Data → domain is the allowed dependency direction, so this type may reference DomainQuote for mapping; +// DomainQuote must never reference QuoteEntity. The mappers live here (like UiQuote.From) but are called +// only by the repository — "do the mapping at the boundary". +internal sealed record QuoteEntity( + string Symbol, + string Name, + AssetCategory Category, + decimal Price, + decimal Change, + decimal ChangePercent, + bool IsValid = true, + string Currency = "USD") +{ + // Domain → storage. The single map-IN seam; the repository calls this on write-through. + public static QuoteEntity From(DomainQuote q) => + new(q.Symbol, q.Name, q.Category, q.Price, q.Change, q.ChangePercent, q.IsValid, q.Currency); + + // Storage → domain. The single map-OUT seam; the repository calls this on read so the entity never + // escapes the data source. + public DomainQuote ToDomainQuote() => + new(Symbol, Name, Category, Price, Change, ChangePercent, IsValid, Currency); +} diff --git a/MarketExtension/Helpers/CurrencyConverter.cs b/MarketExtension/Helpers/CurrencyConverter.cs index 5c1a5e6..c59e92f 100644 --- a/MarketExtension/Helpers/CurrencyConverter.cs +++ b/MarketExtension/Helpers/CurrencyConverter.cs @@ -15,8 +15,8 @@ namespace MarketExtension; // // Two-phase, because the screen renders synchronously but FX is a network fetch: // * PrimeAsync(preferred, natives) — fetch (and cache) every native→preferred rate that isn't already -// fresh, in ONE batched Frankfurter call. Called off the price-load path (see PricedListPage's -// OnPriceCacheUpdated hook), then the page re-renders. +// fresh, in ONE batched Frankfurter call. Awaited before rendering a quote emission (see PortfolioPage's +// OnQuotesProjectingAsync override), so the converted values land in the same paint. // * TryGetRate(from, to) — a synchronous cache read used while building the rows. Returns the rate, or // null when it isn't known yet OR the currency pair isn't ECB-supported (the row then shows native // value only and is excluded from the converted total). from == to short-circuits to 1. diff --git a/MarketExtension/Helpers/StateFlow.cs b/MarketExtension/Helpers/StateFlow.cs index 8c7ecd4..7f657fc 100644 --- a/MarketExtension/Helpers/StateFlow.cs +++ b/MarketExtension/Helpers/StateFlow.cs @@ -72,6 +72,11 @@ void Guarded(T value) return source.Subscribe(Guarded); } + // The underlying stream, for Rx composition (CombineLatest/Switch in MarketRepository.ObserveQuotes). + // Replays the current value on subscribe, like Subscribe(). Bypasses the Guarded wrapper deliberately — + // composition operators manage their own errors, and these flows never OnError. + public IObservable AsObservable() => _subject.AsObservable(); + // Writable entry point for subclasses. Returns true if the value actually changed (i.e. listeners // were notified). Distinct-until-changed at the source: an equal value is not pushed, so re-publishing // an unchanged subset doesn't wake its subscribers. OnNext invokes handlers OUTSIDE the subject's diff --git a/MarketExtension/MarketExtensionCommandsProvider.cs b/MarketExtension/MarketExtensionCommandsProvider.cs index 345468b..d39dff6 100644 --- a/MarketExtension/MarketExtensionCommandsProvider.cs +++ b/MarketExtension/MarketExtensionCommandsProvider.cs @@ -42,6 +42,13 @@ public MarketExtensionCommandsProvider() // Dock bands, each pinnable from the Dock: a ticker strip of favorited instruments, and a // one-line portfolio total (value + daily P&L in the preferred currency). + // + // History: FavoritesDockPage once crashed Command Palette because its ObserveQuotes subscription + // delivered the reactive graph synchronously while Rx held the CombineLatest/Switch gate lock, so + // RaiseItemsChanged's COM call re-entered the host's STA and the gate/STA lock order cycled. Fixed in + // MarketRepository.ObserveQuotes via ObserveOn(TaskPoolScheduler) — surfaces are notified only after + // the gate locks release. (SubscribeOn was tried first and did NOT help: it moves where you subscribe, + // not where notifications fire.) Both bands are live-verified. _dockBands = [ new CommandItem(new FavoritesDockPage(_repository)) { Title = Resources.Command_Markets }, new CommandItem(new PortfolioDockPage(_repository)) { Title = Resources.Command_MarketsPortfolio }, diff --git a/MarketExtension/Pages/FavoritesDockPage.cs b/MarketExtension/Pages/FavoritesDockPage.cs index 52d7365..2c70ba9 100644 --- a/MarketExtension/Pages/FavoritesDockPage.cs +++ b/MarketExtension/Pages/FavoritesDockPage.cs @@ -1,7 +1,6 @@ using System; using System.Collections.Generic; using System.Linq; -using System.Threading.Tasks; using Microsoft.CommandPalette.Extensions; using Microsoft.CommandPalette.Extensions.Toolkit; using Windows.Foundation; @@ -15,50 +14,55 @@ namespace MarketExtension; // // Because the band's command is an IListPage, the host renders each item from GetItems() as // its own button within the one band (see reference/dock-support.md). We use the project's -// INotifyItemsChanged on-load refresh so the band re-reads favorites + quotes every time it -// becomes visible. Live polling while visible (a timer) + a real data source come with the -// API phase — see reference/dock-support.md for the OnLoad lifecycle to add then. +// INotifyItemsChanged on-load refresh so the band re-renders every time it becomes visible. +// +// A PURE OBSERVER of the shared quote cache: while visible it subscribes to the repository's cache-backed +// quote stream for the favorites set (MarketRepository.ObserveQuotes) and renders whatever it emits — and +// nothing else. It does NOT fetch, poll, or handle demo-mode flips itself: the repository owns all of that +// (its single poll loop refreshes the observed set on a timer, and it refills the cache on a source flip), +// so this band can never drift out of sync with any other surface observing the same symbols. Subscribing +// also registers favorites as "observed" so the repository keeps them fresh; disposing on hide unregisters. +// +// Threading: ObserveQuotes delivers via ObserveOn (see MarketRepository) so OnQuotesChanged — and therefore +// RaiseItemsChanged — runs on a pool thread with NO Rx gate lock held. That is what makes this safe: an +// earlier revision delivered synchronously under the CombineLatest/Switch gate, so RaiseItemsChanged's COM +// call into the host re-entered while the gate was held → an STA/gate lock-order deadlock that hung CmdPal. internal sealed partial class FavoritesDockPage : ListPage, INotifyItemsChanged { private readonly MarketRepository _repository; - private UiQuote[]? _quotes; + private UiQuote[]? _quotes; // latest cache emission, projected for rendering; null before the first private event TypedEventHandler? _itemsChanged; - private IDisposable? _subscription; - private IDisposable? _pollSubscription; - private IDisposable? _demoModeSubscription; + // Subscriptions held in a list (not a single field) so a double-`add` without an intervening `remove` + // can't orphan a subscription: a single field would be OVERWRITTEN by the second add, losing the first's + // reference so it's never disposed — and because ObserveQuotes registers its symbols as "observed" on + // subscribe and only unregisters on dispose, that orphan would pin those symbols to the repository's poll + // loop forever. Dispose-all-and-clear in `remove` keeps every subscribe balanced. Matches PricedListPage. + private readonly List _subscriptions = []; event TypedEventHandler INotifyItemsChanged.ItemsChanged { add { _itemsChanged += value; - // Observe the favorites flow while the band is visible: its replay paints the band the moment - // it opens, and any later star/unstar from a palette page refreshes it at once (no waiting - // for a reopen). Disposed in `remove` so a hidden band does no work and doesn't leak. - _subscription = WatchlistStore.Instance.Favorites.Subscribe(_ => RefreshQuotes()); - // Live polling: each tick silently re-prices favorites in place (no spinner). The ticker is a - // pure event stream (no replay), so becoming visible doesn't double-fetch — the favorites - // subscription above already paints. - _pollSubscription = PollTicker.Subscribe(PollRefresh); - // Demo mode flips the data source — re-price from the new one at once if the band is open. A - // hidden band already re-fetches on its next open (the favorites-flow replay calls RefreshQuotes), - // so this only needs to cover a currently-visible/pinned band. replay:false — opening already - // paints via the favorites subscription above. - _demoModeSubscription = MarketSettingsManager.Instance.DemoModeChanged - .Subscribe(_ => RefreshQuotes(), replayOnSubscribe: false); - Log.Info("Poll", $"Dock: started polling [{string.Join(", ", WatchlistStore.Instance.Favorites.Value.Select(i => i.Symbol))}]"); + // Observe the cache-backed quote stream for the favorites set while the band is visible. The + // membership-aware overload replays the current favorites on subscribe (painting from cache the + // moment it opens — fetching only symbols not already cached), re-projects on any star/unstar, + // and re-emits whenever a member quote changes in the cache (so the repository's poll/demo + // refresh repaints the band). Delivery is off-thread (ObserveOn in the repository), so the first + // emission lands after this accessor returns — no synchronous RaiseItemsChanged inside the host's + // subscription. Disposed in `remove` — a hidden band does no work, doesn't leak, and unregisters + // its symbols so the repository stops polling them. + _subscriptions.Add(_repository.ObserveQuotes(WatchlistStore.Instance.Favorites).Subscribe(OnQuotesChanged)); + Log.Info("Dock", $"observing favorites [{string.Join(", ", WatchlistStore.Instance.Favorites.Value.Select(i => i.Symbol))}]"); } remove { _itemsChanged -= value; - _subscription?.Dispose(); - _subscription = null; - _pollSubscription?.Dispose(); - _pollSubscription = null; - _demoModeSubscription?.Dispose(); - _demoModeSubscription = null; - Log.Info("Poll", "Dock: stopped polling"); + foreach (var subscription in _subscriptions) + subscription.Dispose(); + _subscriptions.Clear(); + Log.Info("Dock", "stopped observing favorites"); } } @@ -75,18 +79,24 @@ public FavoritesDockPage(MarketRepository repository) public override IListItem[] GetItems() { - if (_quotes is null) - return []; - var favorites = _quotes; + if (favorites is null) + return []; // before the first cache emission — nothing to show yet if (favorites.Length == 0) { - return [new ListItem(new NoOpCommand { Id = "com.costafotiadis.market.dock.empty" }) + // An empty list means "no favorites" only when the membership is actually empty; otherwise the + // prices just haven't landed in the cache yet, so show nothing (the spinner) rather than the + // empty-state row. + if (WatchlistStore.Instance.Favorites.Value.Count == 0) { - Title = Resources.Favorites_Empty_Title, - Subtitle = Resources.Favorites_Empty_Subtitle, - }]; + return [new ListItem(new NoOpCommand { Id = "com.costafotiadis.market.dock.empty" }) + { + Title = Resources.Favorites_Empty_Title, + Subtitle = Resources.Favorites_Empty_Subtitle, + }]; + } + return []; } return favorites @@ -100,50 +110,14 @@ public override IListItem[] GetItems() .ToArray(); } - internal void RefreshQuotes() - { - _quotes = null; - IsLoading = true; - RaiseItemsChanged(0); - Task.Run(() => LoadQuotes(silent: false)); - } - - // Live-poll refresh: re-price favorites WITHOUT clearing _quotes or showing a spinner, so the current - // prices stay on the band until LoadQuotes swaps the new ones in (no flicker). Skips work before the - // first paint — the favorites subscription already has a load in flight then. `silent: true` also keeps - // a transient bad poll from blanking a good price (see LoadQuotes). - internal void PollRefresh() - { - if (_quotes is null) - return; - Log.Info("Poll", $"Dock: re-pricing favorites [{string.Join(", ", WatchlistStore.Instance.Favorites.Value.Select(i => i.Symbol))}]"); - Task.Run(() => LoadQuotes(silent: true)); - } - - // `silent` = a background poll (no spinner): in that mode, don't let a transient bad result (e.g. a - // Finnhub 429 mapped to an invalid quote) overwrite a price that was fine a moment ago. This mirrors - // PricedListPage.LoadQuotes' keep-last-good guard — without it a single failed poll would replace every - // UiQuote with an invalid one and the band would blank (symbol-only buttons) until the extension is - // reloaded, because a pinned dock never re-fetches except on a poll tick. - private async Task LoadQuotes(bool silent) + // A new cache emission for the favorites set: project to UiQuote for rendering and repaint. Runs on a + // pool thread (ObserveOn) — no Rx gate lock is held here, so RaiseItemsChanged's host call is safe. + private void OnQuotesChanged(IReadOnlyList quotes) { - // The dock shows only favorites — price exactly that subset (a snapshot of the flow's value). - var prior = _quotes; // on-screen prices, for the keep-last-good merge - IEnumerable fetched = - (await _repository.GetQuotesAsync(WatchlistStore.Instance.Favorites.Value)).Select(UiQuote.From); - - if (silent && prior is not null) - { - var lastGood = prior - .Where(q => q.IsValid) - .GroupBy(q => q.Symbol, StringComparer.OrdinalIgnoreCase) - .ToDictionary(g => g.Key, g => g.First(), StringComparer.OrdinalIgnoreCase); - fetched = fetched.Select(q => - !q.IsValid && lastGood.TryGetValue(q.Symbol, out var good) ? good : q); - } - - _quotes = [.. fetched]; - IsLoading = false; + _quotes = [.. quotes.Select(UiQuote.From)]; + // Spinner only while favorites exist but their prices haven't filled the cache yet. + IsLoading = _quotes.Length == 0 && WatchlistStore.Instance.Favorites.Value.Count > 0; + Log.Info("Dock", $"favorites painted: {_quotes.Length} quote(s)"); RaiseItemsChanged(0); } } diff --git a/MarketExtension/Pages/PortfolioDockPage.cs b/MarketExtension/Pages/PortfolioDockPage.cs index 1bf740f..5fc9109 100644 --- a/MarketExtension/Pages/PortfolioDockPage.cs +++ b/MarketExtension/Pages/PortfolioDockPage.cs @@ -1,6 +1,8 @@ using System; using System.Collections.Generic; using System.Linq; +using System.Reactive.Linq; +using System.Threading; using System.Threading.Tasks; using Microsoft.CommandPalette.Extensions; using Microsoft.CommandPalette.Extensions.Toolkit; @@ -16,57 +18,69 @@ namespace MarketExtension; // Unlike the favorites band (one button per instrument), this band renders ONE summary button — the same // totals row the Portfolio screen pins on top — so the dock shows the bottom line at a glance. // -// Live lifecycle mirrors FavoritesDockPage: the project's INotifyItemsChanged on-load refresh re-reads -// holdings + quotes every time the band becomes visible, it subscribes to PortfolioStore.Positions (so a -// quantity/membership change repaints at once) and to PollTicker (so prices refresh on the timer), and it -// disposes both when hidden so a hidden band does no work. Multi-currency: the band prices every holding, -// primes the native->preferred FX rates (CurrencyConverter / Frankfurter), and rolls them up via -// UiPortfolio — exactly like PortfolioPage, but it AWAITS the FX prime inside the async load so the -// converted total is ready in one paint (no progressive native-only first frame needed on a one-row band). +// A PURE OBSERVER of the shared quote cache, like FavoritesDockPage: while visible it subscribes to the +// repository's cache-backed quote stream for the holdings (MarketRepository.ObserveQuotes over +// PortfolioStore.Instruments) and renders whatever it emits. It does NOT fetch, poll, or handle demo-mode +// flips itself — the repository owns all of that (its single poll loop refreshes the observed set on a timer, +// and it refills the cache on a source flip), and the cache owns keep-last-good — so this band can never +// drift out of sync with the Portfolio screen. Because PortfolioStore.Instruments re-emits on EVERY mutation +// (reference equality), the membership-aware ObserveQuotes Switch also re-fires on a quantity-only edit, so +// the handler re-reads PortfolioStore.Positions and re-rolls the total. +// +// The one wrinkle over the favorites band is async PROJECTION: rolling up needs native->preferred FX rates, +// an async fetch, so the observe handler AWAITS CurrencyConverter.PrimeAsync before building the total — one +// paint, no progressive native-only frame (a one-row band has nothing to gain from that). This is safe +// because ObserveQuotes delivers via ObserveOn (see MarketRepository): the handler runs on a pool thread with +// NO Rx gate lock held, so RaiseItemsChanged's COM call into the host can't re-enter under a producer-side +// lock (the deadlock the favorites band hit before the ObserveOn fix). internal sealed partial class PortfolioDockPage : ListPage, INotifyItemsChanged { - private const string PortfolioGlyph = ""; // Segoe MDL2 Bank glyph U+E825 (matches PortfolioPage's summary row) + private const string PortfolioGlyph = "\uE825"; // Segoe MDL2 Bank glyph U+E825 (matches PortfolioPage's summary row) private readonly MarketRepository _repository; - private UiQuote[]? _quotes; // last priced holdings, for the keep-last-good poll merge - private UiPortfolio? _portfolio; // the rolled-up total for rendering; null until first load lands + private UiPortfolio? _portfolio; // the rolled-up total for rendering; null until the first roll-up lands private event TypedEventHandler? _itemsChanged; - private IDisposable? _subscription; - private IDisposable? _pollSubscription; - private IDisposable? _demoModeSubscription; + // Subscriptions in a list (not a single field) so a double-`add` without an intervening `remove` can't + // orphan a subscription — which would also leave its symbols pinned to the repository's poll loop. Same + // reasoning and pattern as FavoritesDockPage / PricedListPage. + private readonly List _subscriptions = []; event TypedEventHandler INotifyItemsChanged.ItemsChanged { add { _itemsChanged += value; - // Observe the holdings flow while the band is visible: its replay paints the band the moment it - // opens, and any later add/edit/remove from the detail page refreshes it at once. Positions uses - // reference equality, so even a quantity-only edit re-emits and re-rolls the total. - _subscription = PortfolioStore.Instance.Positions.Subscribe(_ => RefreshPortfolio()); - // Live polling: each tick silently re-prices holdings in place (no spinner). The ticker is a pure - // event stream (no replay), so becoming visible doesn't double-fetch — the positions subscription - // above already paints. - _pollSubscription = PollTicker.Subscribe(PollRefresh); - // Demo mode flips the data source — re-roll the total from the new one at once if the band is - // open. A hidden band already re-fetches on its next open (the positions-flow replay calls - // RefreshPortfolio), so this only needs to cover a currently-visible/pinned band. replay:false — - // opening already paints via the positions subscription above. - _demoModeSubscription = MarketSettingsManager.Instance.DemoModeChanged - .Subscribe(_ => RefreshPortfolio(), replayOnSubscribe: false); - Log.Info("Poll", $"PortfolioDock: started polling [{string.Join(", ", PortfolioStore.Instance.Positions.Value.Select(p => p.Instrument.Symbol))}]"); + // Observe the cache-backed quote stream for the holdings while the band is visible. The + // membership-aware overload replays the holdings' cached quotes on subscribe (painting at once, + // fetching only symbols not already cached), re-projects on any add/remove/quantity edit + // (Instruments re-emits on every mutation), and re-emits whenever a member quote changes in the + // cache (so the repository's poll/demo refresh repaints the band). Delivery is off-thread + // (ObserveOn in the repository). Disposed in `remove` — a hidden band does no work and unregisters + // its symbols so the repository stops polling them. No PollTicker / DemoModeChanged subscription + // and no keep-last-good merge here anymore: the repository polls the observed set and refills on a + // source flip, and the cache holds the last good price through a transient bad fetch. + // Project each emission through the async roll-up with Switch, NOT fire-and-forget. The cache fills + // symbol-by-symbol on a cold start, so ObserveQuotes emits progressive partial snapshots (e.g. just + // BABA before SPY lands); each roll-up AWAITS the FX prime, so two handlers could otherwise race and + // a stale partial could finish last and overwrite the full total. Switch cancels the prior + // projection's CancellationToken the instant a newer emission arrives, and OnQuotesChangedAsync + // honors that token before it paints — so only the latest emission ever reaches the band. (Plain + // `_ = OnQuotesChangedAsync(...)` escaped ObserveOn's serialization because it returns at the first + // await, which is exactly the bug this fixes.) + _subscriptions.Add(_repository.ObserveQuotes(PortfolioStore.Instance.Instruments) + .Select(quotes => Observable.FromAsync(ct => OnQuotesChangedAsync(quotes, ct))) + .Switch() + .Subscribe()); + Log.Info("Dock", $"observing portfolio [{string.Join(", ", PortfolioStore.Instance.Positions.Value.Select(p => p.Instrument.Symbol))}]"); } remove { _itemsChanged -= value; - _subscription?.Dispose(); - _subscription = null; - _pollSubscription?.Dispose(); - _pollSubscription = null; - _demoModeSubscription?.Dispose(); - _demoModeSubscription = null; - Log.Info("Poll", "PortfolioDock: stopped polling"); + foreach (var subscription in _subscriptions) + subscription.Dispose(); + _subscriptions.Clear(); + Log.Info("Dock", "stopped observing portfolio"); } } @@ -83,10 +97,9 @@ public PortfolioDockPage(MarketRepository repository) public override IListItem[] GetItems() { - if (_quotes is null) - return []; // not loaded yet - - if (_portfolio is null || !_portfolio.HasHoldings) + // No holdings at all → the empty-state row. Read membership (not _portfolio) so this shows only when + // the portfolio is genuinely empty, never during the pre-first-emission window. + if (PortfolioStore.Instance.Positions.Value.Count == 0) { return [new ListItem(new NoOpCommand { Id = "com.costafotiadis.market.dock.portfolio.empty" }) { @@ -95,95 +108,92 @@ public override IListItem[] GetItems() }]; } - // The one summary button: total value as the title, today's P&L (and any "not converted" note) as the - // subtitle. Clicking opens the full Portfolio screen for the breakdown. + // Snapshot the field once: a roll-up handler runs on a pool thread and can reassign (or null, via the + // spinner branch) _portfolio between the null-check and the reads below, which would otherwise tear the + // rendered row or NRE. + var portfolio = _portfolio; + + // Holdings exist but the roll-up isn't ready yet (prices still filling the cache) → let the spinner show. + if (portfolio is null) + return []; + + // The one summary button: total value as the title, today's P&L (and any total-return / "not + // converted" notes) as the subtitle. Clicking opens the full Portfolio screen for the breakdown. return [ new ListItem(new PortfolioPage(_repository)) { - Title = Strings.Format(Resources.Portfolio_TotalsRow_Title, _portfolio.FormatTotalValue()), - Subtitle = _portfolio.FormatTotalChange() + _portfolio.FormatTotalReturnNote() + _portfolio.FormatUnconvertedNote(), + Title = Strings.Format(Resources.Portfolio_TotalsRow_Title, portfolio.FormatTotalValue()), + Subtitle = portfolio.FormatTotalChange() + portfolio.FormatTotalReturnNote() + portfolio.FormatUnconvertedNote(), Icon = new IconInfo(PortfolioGlyph), }, ]; } - internal void RefreshPortfolio() + // A new cache emission for the holdings: roll up the total and repaint. Runs on a pool thread (ObserveOn) + // with no Rx gate lock held, so RaiseItemsChanged's host call is safe, and awaiting the FX prime here is + // what makes the converted total land in a single paint. Driven by Switch (see the subscribe): `ct` is + // cancelled the moment a newer emission supersedes this one, so it's checked before every paint — a stale + // or partial snapshot bails instead of overwriting the latest total. Own try/catch so no exception escapes + // onto the pool thread; OperationCanceledException is swallowed quietly (it's the expected supersede path). + private async Task OnQuotesChangedAsync(IReadOnlyList quotes, CancellationToken ct) { - _quotes = null; - _portfolio = null; - IsLoading = true; - RaiseItemsChanged(0); - Task.Run(() => LoadPortfolio(silent: false)); - } + try + { + var positions = PortfolioStore.Instance.Positions.Value; - // Live-poll refresh: re-price holdings WITHOUT clearing _quotes/_portfolio or showing a spinner, so the - // current total stays on the band until LoadPortfolio swaps the new one in (no flicker). Skips work - // before the first paint — the positions subscription already has a load in flight then. - internal void PollRefresh() - { - if (_quotes is null) - return; - Log.Info("Poll", $"PortfolioDock: re-pricing holdings [{string.Join(", ", PortfolioStore.Instance.Positions.Value.Select(p => p.Instrument.Symbol))}]"); - Task.Run(() => LoadPortfolio(silent: true)); - } + // Holdings exist but their prices haven't landed in the cache yet → keep the spinner, don't roll up. + if (positions.Count > 0 && quotes.Count == 0) + { + if (ct.IsCancellationRequested) return; // a newer emission already supersedes this — don't blank + Log.Info("Dock", $"portfolio: {positions.Count} holding(s) but no cached prices yet — spinner"); + _portfolio = null; + IsLoading = true; + RaiseItemsChanged(0); + return; + } - // Price the current holdings, convert into the preferred currency, and roll up the total. `silent` = a - // background poll (no spinner): in that mode, don't let a transient bad result (e.g. a 429 mapped to an - // invalid quote) overwrite a price that was fine a moment ago, which would otherwise drop that holding - // from the total and make it jump. Same keep-last-good guard as FavoritesDockPage / PricedListPage. - private async Task LoadPortfolio(bool silent) - { - var positions = PortfolioStore.Instance.Positions.Value; // snapshot the holdings to price - if (positions.Count == 0) - { - _quotes = []; // loaded-but-empty (distinct from null = not-loaded), so GetItems shows the empty row - _portfolio = null; + // Prime native->preferred FX rates BEFORE rolling up, so the converted total is ready this paint. + // PrimeAsync skips currencies it already has fresh, so a steady portfolio does no extra network. + var preferred = MarketSettingsManager.Instance.PortfolioCurrency; + var natives = quotes + .Where(q => q.IsValid) + .Select(q => q.Currency) + .Distinct(StringComparer.OrdinalIgnoreCase) + .ToArray(); + if (natives.Length > 0) + await CurrencyConverter.Instance.PrimeAsync(preferred, natives, ct); + + // A newer emission arrived while we were priming → it's the one that should paint. Drop this + // (possibly partial / now-stale) roll-up so it can't overwrite the fresher total — the whole point + // of the Switch above; without this check the stale handler would still write _portfolio. + if (ct.IsCancellationRequested) return; + + var portfolio = BuildPortfolio(positions, quotes, preferred); + _portfolio = portfolio; + Log.Info("Dock", + $"portfolio rolled up: {portfolio.FormatTotalValue()} {portfolio.FormatTotalChange()} across {positions.Count} holding(s)"); IsLoading = false; RaiseItemsChanged(0); - return; } - - var instruments = positions.Select(p => p.Instrument).ToList(); - var prior = _quotes; // on-screen prices, for the keep-last-good merge - IEnumerable fetched = - (await _repository.GetQuotesAsync(instruments)).Select(UiQuote.From); - - if (silent && prior is not null) + catch (OperationCanceledException) { - var lastGood = prior - .Where(q => q.IsValid) - .GroupBy(q => q.Symbol, StringComparer.OrdinalIgnoreCase) - .ToDictionary(g => g.Key, g => g.First(), StringComparer.OrdinalIgnoreCase); - fetched = fetched.Select(q => - !q.IsValid && lastGood.TryGetValue(q.Symbol, out var good) ? good : q); + // Superseded by a newer emission while priming FX — expected (Switch cancelled us), ignore. + } + catch (Exception ex) + { + Log.Error("Dock", "portfolio roll-up failed", ex); } - - var quotes = fetched.ToArray(); - _quotes = quotes; - - // Prime native->preferred FX rates before rolling up, so the converted total is ready this paint. - var preferred = MarketSettingsManager.Instance.PortfolioCurrency; - var natives = quotes - .Where(q => q.IsValid) - .Select(q => q.Source.Currency) - .Distinct(StringComparer.OrdinalIgnoreCase) - .ToArray(); - if (natives.Length > 0) - await CurrencyConverter.Instance.PrimeAsync(preferred, natives); - - _portfolio = BuildPortfolio(positions, quotes, preferred); - IsLoading = false; - RaiseItemsChanged(0); } // Zip each holding's quantity with its priced quote and the (cached) native->preferred rate, then roll // the set up into the totals — the same composition PortfolioPage.LeadingRows does, just eager here. - private static UiPortfolio BuildPortfolio(IReadOnlyList positions, UiQuote[] quotes, string preferred) + private static UiPortfolio BuildPortfolio( + IReadOnlyList positions, IReadOnlyList quotes, string preferred) { var quoteBySymbol = quotes .GroupBy(q => q.Symbol, StringComparer.OrdinalIgnoreCase) - .ToDictionary(g => g.Key, g => g.First().Source, StringComparer.OrdinalIgnoreCase); + .ToDictionary(g => g.Key, g => g.First(), StringComparer.OrdinalIgnoreCase); var uiPositions = new List(); foreach (var p in positions) diff --git a/MarketExtension/Pages/PortfolioPage.cs b/MarketExtension/Pages/PortfolioPage.cs index 51a9e61..be69bec 100644 --- a/MarketExtension/Pages/PortfolioPage.cs +++ b/MarketExtension/Pages/PortfolioPage.cs @@ -1,6 +1,7 @@ using System; using System.Collections.Generic; using System.Linq; +using System.Threading; using System.Threading.Tasks; using Microsoft.CommandPalette.Extensions; using Microsoft.CommandPalette.Extensions.Toolkit; @@ -9,23 +10,24 @@ namespace MarketExtension; // Top-level screen: the user's portfolio holdings, priced live, with a totals summary pinned on top. -// Reached from the Markets hub. Subclasses PricedListPage (like Watchlist/Favorites) to inherit all of -// its caching / polling / reconcile / keep-last-good plumbing — it just observes PortfolioStore.Instruments -// (the symbols to price) instead of a watchlist subset. +// Reached from the Markets hub. Subclasses PricedListPage (like Watchlist/Favorites) to inherit its shared +// quote-cache observation — it just observes PortfolioStore.Instruments (the symbols to price) instead of a +// watchlist subset. // // Each row shows the holding (symbol · quantity) with its market value and today's P&L; Enter opens the // shared SymbolDetailPage, which is where holdings are added/edited/removed (Add to Portfolio / Edit // holding / Remove), consistent with every other list. The quantity for a row is read from PortfolioStore // at render time — the same way WatchlistPage reads IsFavorite — and the totals summary is built from the // full priced set via the LeadingRows hook. A quantity-only edit re-emits PortfolioStore.Instruments -// (default-equality flow), so the base re-renders with no fetch and the new quantity + total show at once. +// (default-equality flow), so the membership-aware ObserveQuotes Switch re-fires with no fetch and the new +// quantity + total show at once. // // Multi-currency: holdings priced in various native currencies are converted into the user's // PortfolioCurrency setting. Conversion rates are an async FX fetch (CurrencyConverter / Frankfurter), but -// GetItems is synchronous — so the rates are PRIMED off the price-load path (OnPriceCacheUpdated, the base's -// post-update hook) and read from the converter's cache while rendering. Until a rate lands a row shows its -// native value only and sits out of the total; a holding in a currency the FX provider can't convert stays -// that way (surfaced as "N not converted" on the summary). +// GetItems is synchronous — so the rates are PRIMED in OnQuotesProjectingAsync (the base awaits it before +// rendering an emission) and read from the converter's cache while rendering, so the converted values land +// in the same paint. Until a rate lands a row shows its native value only and sits out of the total; a +// holding in a currency the FX provider can't convert stays that way (surfaced as "N not converted"). internal sealed partial class PortfolioPage : PricedListPage { private const string PortfolioGlyph = ""; // Segoe MDL2 Bank @@ -91,26 +93,23 @@ private static UiPosition MakePosition(DomainQuote quote, decimal qty, decimal? return UiPosition.From(quote, qty, preferred, rate, costBasis); } - // After each price update, ensure the FX rates for every native currency now present are fetched into - // the converter's cache, then re-render so the converted values + total appear. Runs off the UI/render - // path; the converter skips currencies it already has fresh, so a steady portfolio does no extra network. - protected override void OnPriceCacheUpdated() + // Before each emission is rendered, ensure the FX rates for every native currency now present are fetched + // into the converter's cache, so the converted values + total are ready in the same paint. The base awaits + // this then projects + repaints once (no RaiseItemsChanged / Task.Run here — already on a pool thread). The + // converter skips currencies it already has fresh, so a steady portfolio does no extra network. `ct` is + // cancelled when a newer emission supersedes this one (Switch in the base), so it's forwarded to the prime + // to stop a superseded FX fetch early. + protected override async Task OnQuotesProjectingAsync(IReadOnlyList quotes, CancellationToken ct) { var preferred = MarketSettingsManager.Instance.PortfolioCurrency; - var natives = SnapshotPricedQuotes() + var natives = quotes .Where(q => q.IsValid) - .Select(q => q.Source.Currency) + .Select(q => q.Currency) .Distinct(StringComparer.OrdinalIgnoreCase) .ToArray(); - if (natives.Length == 0) - return; - - Task.Run(async () => - { - await CurrencyConverter.Instance.PrimeAsync(preferred, natives); - RaiseItemsChanged(0); - }); + if (natives.Length > 0) + await CurrencyConverter.Instance.PrimeAsync(preferred, natives, ct); } protected override IListItem[] EmptyState() => diff --git a/MarketExtension/Pages/PricedListPage.cs b/MarketExtension/Pages/PricedListPage.cs index a4cc2ee..1c502a9 100644 --- a/MarketExtension/Pages/PricedListPage.cs +++ b/MarketExtension/Pages/PricedListPage.cs @@ -1,6 +1,7 @@ using System; using System.Collections.Generic; using System.Linq; +using System.Reactive.Linq; using System.Threading; using System.Threading.Tasks; using Microsoft.CommandPalette.Extensions; @@ -10,36 +11,40 @@ namespace MarketExtension; -// Shared base for the two priced list screens (Watchlist, Favorites). Each OBSERVES a different -// WatchlistStore membership flow (passed to the ctor) and prices its instruments. Wiring: +// Shared base for the priced list screens (Watchlist, Favorites, Portfolio). Each OBSERVES a different +// membership flow (passed to the ctor — WatchlistStore.Watchlist/Favorites or PortfolioStore.Instruments) +// and renders its instruments priced live. // -// * The page subscribes to its instrument flow in the INotifyItemsChanged `add` accessor (CmdPal's -// de-facto "page became visible" hook) and disposes in `remove`. StateFlow replays the current set -// on subscribe, which drives the initial price load — so navigating here always re-prices. -// * When membership changes (a row added/removed from anywhere), the flow pushes the new set and the -// page reconciles LOCALLY against a per-symbol price cache: rows that left are dropped with no -// network call, and only newly-added symbols are fetched. This keeps edits cheap against the -// Finnhub rate limit and makes a removed row disappear at once (rather than lingering until reload). -// * The Refresh row (and, later, the live-poll timer) forces a full re-price of the whole set. +// A PURE OBSERVER of the shared quote cache, like the two dock bands (FavoritesDockPage / PortfolioDockPage): // -// Typing only re-filters the already-loaded quotes locally — these screens never hit the network on +// * In the INotifyItemsChanged `add` accessor (CmdPal's de-facto "page became visible" hook) the page +// subscribes to MarketRepository.ObserveQuotes(membership flow) and disposes in `remove`. The +// membership-aware overload replays the current set's cached quotes on subscribe (the initial paint, +// fetching only symbols not already cached/stale), re-projects on any membership change (Switch — so a +// removed row drops and a new one is fetched with no page-side reconcile), and re-emits whenever a member +// quote changes in the cache (so the repository's single poll loop / demo-flip refill repaint the page). +// * The page does NOT fetch, poll, or handle demo-mode flips itself — the repository owns all of that, and +// the cache owns keep-last-good — so these screens can never drift out of sync with each other or the dock. +// * Delivery is off-thread (ObserveOn in the repository): the emission handler runs on a pool thread with no +// Rx gate lock held, so RaiseItemsChanged's host call is safe (the deadlock the favorites band hit before +// the ObserveOn fix). Do NOT add Task.Run / SubscribeOn / ObserveOn on this path. +// +// Typing only re-filters the already-emitted quotes locally — these screens never hit the network on a // keystroke (only the SearchPage talks to /search, and only on Enter). internal abstract partial class PricedListPage : DynamicListPage, INotifyItemsChanged { protected readonly MarketRepository Repository; private readonly StateFlow> _instruments; + // Subscriptions in a list (not a single field) so a double-`add` without an intervening `remove` can't + // orphan a subscription — which would also leave its symbols pinned to the repository's poll loop. Same + // reasoning and pattern as the dock bands. private readonly List _subscriptions = []; - // Per-symbol price cache (key = WatchlistStore.Normalize(symbol)) + the current instrument set, so a - // membership change can reconcile without re-pricing instruments we already hold. Also guards the - // in-flight fetch token swap. - private readonly Lock _cacheLock = new(); - private readonly Dictionary _priceCache = []; - private IReadOnlyList _snapshot = []; - private bool _received; // have we processed the first flow emission yet? - private int _fetchGeneration; // bumped per fetch so a superseded one's late result is discarded - private long _lastFullPriceTicks; // Environment.TickCount64 at the last WHOLE-set re-price; 0 = never priced + // The latest emission projected for rendering, in set order (null entries already dropped by the cache). + // null until the first emission lands. Written on a pool thread, read on the host thread; reference + // assignment is atomic and we never mutate the array in place, so no lock (same as the dock bands). + private UiQuote[]? _quotes; private event TypedEventHandler? _itemsChanged; @@ -48,25 +53,34 @@ event TypedEventHandler INotifyItemsChanged.Item add { _itemsChanged += value; - // Primary flow: drives which rows exist + their pricing (replays the current set at once). - _subscriptions.Add(_instruments.Subscribe(OnInstrumentsChanged)); + // Observe the cache-backed quote stream for this membership set while the page is visible. Replays + // current quotes on subscribe (initial paint), re-projects on membership change, re-emits on price + // change. Project each emission through the async handler with Switch (NOT fire-and-forget): a + // subclass can await an async projection (the Portfolio screen primes FX rates) before the repaint, + // and on a cold cache ObserveQuotes emits progressive partial snapshots (CombineLatest fills symbol + // by symbol), so two awaited handlers could otherwise race and a stale partial finish last. Switch + // cancels the prior projection's CancellationToken the instant a newer emission arrives, and + // OnQuotesChangedAsync honors it before painting — so only the latest emission ever reaches the page. + // (Watchlist/Favorites never await — their projection is a no-op — so this is a no-op for them too.) + _subscriptions.Add(Repository.ObserveQuotes(_instruments) + .Select(quotes => Observable.FromAsync(ct => OnQuotesChangedAsync(quotes, ct))) + .Switch() + .Subscribe()); // Secondary flows: a change just re-renders (e.g. the Watchlist's ★ when favorites change). foreach (var trigger in RelistTriggers) _subscriptions.Add(trigger.Subscribe(_ => RaiseItemsChanged(0))); - // Live polling: while visible, each tick silently re-prices the set in place. The ticker is a - // pure event stream (no replay), so opening the page doesn't double-fetch (the membership replay - // above already did the first load); the ticker runs its timer only while a surface is subscribed. - _subscriptions.Add(PollTicker.Subscribe(PollRefresh)); // Re-render when a key is added/cleared so the missing-key hint appears/disappears at once. - // replay:false — the membership flow above already drives the initial paint. + // replay:false — the observe subscription above already drives the initial paint. _subscriptions.Add(MarketSettingsManager.Instance.HasAnyApiKey .Subscribe(_ => RaiseItemsChanged(0), replayOnSubscribe: false)); // Re-render when throttling starts/stops so the rate-limited banner appears/disappears at once. // replay:false — same reason as above. _subscriptions.Add(RateLimitSignal.Instance.IsRateLimited .Subscribe(_ => RaiseItemsChanged(0), replayOnSubscribe: false)); - Log.Info("Poll", $"{Title}: started polling [{string.Join(", ", _snapshot.Select(i => i.Symbol))}] " + - $"(every {MarketSettingsManager.Instance.RefreshMinutes} min, 0=off)"); + // Spinner only on the very first open of a non-empty set; a reopen keeps the last data showing + // until the cache replay lands. + IsLoading = _quotes is null && _instruments.Value.Count > 0; + Log.Info("Poll", $"{Title}: observing [{string.Join(", ", _instruments.Value.Select(i => i.Symbol))}]"); } remove { @@ -74,7 +88,7 @@ event TypedEventHandler INotifyItemsChanged.Item foreach (var subscription in _subscriptions) subscription.Dispose(); _subscriptions.Clear(); - Log.Info("Poll", $"{Title}: stopped polling [{string.Join(", ", _snapshot.Select(i => i.Symbol))}]"); + Log.Info("Poll", $"{Title}: stopped observing"); } } @@ -86,29 +100,6 @@ protected PricedListPage(MarketRepository repository, StateFlow