From cc765eb1e8d56ebc863c1d6be707313de90437e6 Mon Sep 17 00:00:00 2001 From: Costa Fotiadis Date: Fri, 26 Jun 2026 20:44:05 +0300 Subject: [PATCH 01/10] cahe work in progress --- MarketExtension/Data/IQuoteCache.cs | 36 +++++++++ MarketExtension/Data/InMemoryQuoteCache.cs | 76 +++++++++++++++++ MarketExtension/Data/MarketRepository.cs | 94 +++++++++++++++++++++- MarketExtension/Helpers/StateFlow.cs | 5 ++ 4 files changed, 208 insertions(+), 3 deletions(-) create mode 100644 MarketExtension/Data/IQuoteCache.cs create mode 100644 MarketExtension/Data/InMemoryQuoteCache.cs diff --git a/MarketExtension/Data/IQuoteCache.cs b/MarketExtension/Data/IQuoteCache.cs new file mode 100644 index 0000000..a1c5406 --- /dev/null +++ b/MarketExtension/Data/IQuoteCache.cs @@ -0,0 +1,36 @@ +namespace MarketExtension; + +// A process-wide store of the latest DomainQuote 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 an interface so the in-memory implementation now 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 and writes through to it on every fetch; surfaces +// will later OBSERVE it instead of fetching independently (not done yet — see the cache layer plan). +// +// Domain layer: holds DomainQuote (no formatting). Surfaces project to UiQuote as they do today. +internal interface IQuoteCache +{ + // Current cached quote for a symbol, or null if it was never fetched / has been cleared. + // Synchronous snapshot read. + DomainQuote? 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 + // DomainQuote 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(DomainQuote 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/InMemoryQuoteCache.cs b/MarketExtension/Data/InMemoryQuoteCache.cs new file mode 100644 index 0000000..0e23c0e --- /dev/null +++ b/MarketExtension/Data/InMemoryQuoteCache.cs @@ -0,0 +1,76 @@ +using System.Collections.Generic; +using System.Diagnostics.CodeAnalysis; +using System.Linq; +using System.Threading; + +namespace MarketExtension; + +// In-memory IQuoteCache: 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). +// +// Swap this for a database-backed IQuoteCache 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 InMemoryQuoteCache : IQuoteCache +{ + 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 DomainQuote? 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(DomainQuote quote, bool keepLastGood = true) + { + var key = WatchlistStore.Normalize(quote.Symbol); + + MutableStateFlow flow; + DomainQuote? 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..aa83572 100644 --- a/MarketExtension/Data/MarketRepository.cs +++ b/MarketExtension/Data/MarketRepository.cs @@ -1,5 +1,7 @@ +using System; using System.Collections.Generic; using System.Linq; +using System.Reactive.Linq; using System.Threading; using System.Threading.Tasks; @@ -11,20 +13,61 @@ 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 IQuoteCache: 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. Migrating the surfaces onto it is a +// later pass; for now only write-through runs. +internal sealed class MarketRepository { + private readonly IMarketDataProvider[] _providers; + private readonly IQuoteCache _cache; + + // Default: an in-memory cache. The call site in MarketExtensionCommandsProvider uses this. + public MarketRepository(params IMarketDataProvider[] providers) + : this(new InMemoryQuoteCache(), providers) { } + + // Injectable cache — for a future database-backed IQuoteCache (and tests). No other code changes. + public MarketRepository(IQuoteCache cache, IMarketDataProvider[] providers) + { + _cache = cache; + _providers = providers; + + // A demo-mode flip swaps the data SOURCE, so every cached price is now wrong (it came from the + // other source). Clear so it can neither linger nor be preserved by keep-last-good. replay:false — + // construction is a no-op (the cache is empty then). Subscribed here, before any surface, so on a + // flip this Clear runs before a surface's own re-fetch. Never disposed: the repository lives for + // the whole process. + _ = MarketSettingsManager.Instance.DemoModeChanged.Subscribe(_ => _cache.Clear(), replayOnSubscribe: false); + } + // 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; } + // Fetch quotes for a set of instruments AND write them through to the shared cache. Keeps its exact + // public signature/behavior (callers still get the ordered list back); the cache fill is a side effect, + // so every existing caller populates the cache with no change to it. public async Task> GetQuotesAsync( IReadOnlyList instruments, CancellationToken ct = default) + { + var ordered = await FetchMergeAsync(instruments, ct).ConfigureAwait(false); + foreach (var quote in ordered) + _cache.Upsert(quote); // keep-last-good + return ordered; + } + + // Route → fetch (concurrent, one batch per provider) → merge into the caller's instrument order. + // The raw data path with no cache side effects, shared by GetQuotesAsync and 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 +106,51 @@ public async Task> GetQuotesAsync( return ordered; } + // Force a fetch now; results reach observers through the cache (the caller does not read the return). + // The seam the manual Refresh row / demo path call now and the orchestration poll loop will call later. + // 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); + foreach (var quote in ordered) + _cache.Upsert(quote, keepLastGood); + } + + // Observe a FIXED set of instruments as one cache-backed list observable: on subscribe, fetch any + // not-yet-cached symbols once (write-through, fire-and-forget), 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. Defer runs the fetch-missing per subscribe. + public IObservable> ObserveQuotes(IReadOnlyList instruments) + => Observable.Defer(() => + { + var missing = instruments.Where(i => _cache.Get(i.Symbol) is null).ToList(); + if (missing.Count > 0) + _ = RefreshSafe(missing); // fire-and-forget; writes through → observers see it land + + if (instruments.Count == 0) + return Observable.Return>([]); + + return Observable + .CombineLatest(instruments.Select(i => _cache.Observe(i.Symbol).AsObservable())) + .Select(quotes => (IReadOnlyList)quotes.OfType().ToList()); + }); + + // Membership-aware: re-projects (Switch) whenever the set changes — fetching newly-added symbols and + // dropping departed ones for free. The shape the priced pages / dock bands will consume next pass. + public IObservable> ObserveQuotes( + StateFlow> instruments) + => instruments.AsObservable().Select(set => ObserveQuotes(set)).Switch(); + + // RefreshAsync with its exceptions swallowed + logged — for the fire-and-forget fetch-missing path, + // 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); } + } + // 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. 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 From c706a4077af62aabd163a2312425a2834c2c102c Mon Sep 17 00:00:00 2001 From: Costa Fotiadis Date: Fri, 26 Jun 2026 22:02:44 +0300 Subject: [PATCH 02/10] cahe work in progress --- CLAUDE.md | 115 ++++++++++++- MarketExtension/Data/InMemoryQuoteCache.cs | 10 +- MarketExtension/Data/MarketRepository.cs | 161 +++++++++++++++--- .../MarketExtensionCommandsProvider.cs | 5 + MarketExtension/Pages/FavoritesDockPage.cs | 127 +++++--------- 5 files changed, 307 insertions(+), 111 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index b30c2e6..f3bda29 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -82,7 +82,9 @@ 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 `IQuoteCache`** (in-memory by default; an injectable ctor overload `MarketRepository(IQuoteCache, providers[])` for a future DB-backed cache): `GetQuotesAsync`/`RefreshAsync` **write 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. See "Shared quote cache + repository orchestration (in progress)". | +| `Data/IQuoteCache.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). `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/InMemoryQuoteCache.cs` | The **real, live** in-memory `IQuoteCache` (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 `IQuoteCache` later via the repo's injectable ctor — no repository/UI change. | | `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 +93,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`. **Migration in progress:** `MarketRepository` now holds ONE process-lifetime subscription that polls the *observed* set centrally, so a migrated surface (the favorites dock) no longer subscribes here at all — it just observes the cache. Un-migrated surfaces (the `PricedListPage` trio, the portfolio dock, the chart) still subscribe directly for now. | | `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 +233,26 @@ cascading CS0534 ("does not implement … `GetTypeInfo`") onto **every** context ## Current Status / Next Steps +- **Done (this round): shared quote cache + repository orchestration — the de-drift refactor (cache LIVE, favorites dock migrated + verified).** + A new **`IQuoteCache`** (in-memory `InMemoryQuoteCache` 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 + so far: `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. **NOT yet migrated** (still on per-surface fetch+poll, + filling the cache only via write-through): `PortfolioDockPage`, the `PricedListPage` trio + (`WatchlistPage`/`FavoritesPage`/`PortfolioPage`), plus the `SymbolDetailPage` chart (candles, separate path). + **The favorites dock is the template for the rest.** See "Shared quote cache + repository orchestration" for + the full design + the migration checklist. - **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 +576,89 @@ 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 (cache live; favorites dock migrated, rest pending) + +**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:** +- **`IQuoteCache`** (`Data/`) — the shared store abstraction: per-symbol `StateFlow`, keyed by + `WatchlistStore.Normalize`. `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. ✅ **`InMemoryQuoteCache` 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 `IQuoteCache`. `GetQuotesAsync`/`RefreshAsync` write through + (so the cache fills even from un-migrated surfaces). `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. +- **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 so far: `FavoritesDockPage` only — the template.** It is now 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). + +**To migrate the next surface (the checklist):** +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.** Carry over carefully: the **freshness clock** (`_lastFullPriceTicks` / + `RefreshStaleQuotes`) and **`OnPriceCacheUpdated`** (Portfolio's FX `PrimeAsync`) must be driven by + *fetch completion*, NOT by observe emissions — value-equality dedup hides an unchanged re-fetch, so keying + "a fetch happened" off an emission would miss it. `LeadingRows` must recompute on each emission. +7. `PortfolioDockPage` additionally `await`s `CurrencyConverter.PrimeAsync` before rolling up — keep that + inside the observe handler (one paint), not a separate fetch. + +**Transitional truth:** un-migrated surfaces still fetch+poll themselves AND fill the cache via write-through, +so a symbol they share with the dock is briefly fetched by both the repo loop (for the dock) and the surface's +own poll — harmless (the cache dedups by value), and it collapses as each surface migrates. + ### Live price polling (done) +> **Superseded for migrated surfaces** by "Shared quote cache + repository orchestration" above — the favorites +> dock no longer polls itself; the repository's single loop polls the observed set. The model below still +> describes the **un-migrated** surfaces (the `PricedListPage` trio, the portfolio dock, the chart). + 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 +733,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/InMemoryQuoteCache.cs b/MarketExtension/Data/InMemoryQuoteCache.cs index 0e23c0e..98d7df5 100644 --- a/MarketExtension/Data/InMemoryQuoteCache.cs +++ b/MarketExtension/Data/InMemoryQuoteCache.cs @@ -12,8 +12,14 @@ namespace MarketExtension; // System.Threading.Lock is non-reentrant, so updating under the lock could deadlock a re-entrant // handler). // -// Swap this for a database-backed IQuoteCache later via MarketRepository's injectable ctor overload — -// no repository or UI change. +// 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 IQuoteCache 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 " + diff --git a/MarketExtension/Data/MarketRepository.cs b/MarketExtension/Data/MarketRepository.cs index aa83572..8eef4bd 100644 --- a/MarketExtension/Data/MarketRepository.cs +++ b/MarketExtension/Data/MarketRepository.cs @@ -1,6 +1,8 @@ using System; 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; @@ -17,13 +19,24 @@ namespace MarketExtension; // It also owns the shared IQuoteCache: 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. Migrating the surfaces onto it is a -// later pass; for now only write-through runs. +// 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. +// Surfaces still poll/fetch individually today; moving them fully onto observe-only is a later pass. internal sealed class MarketRepository { private readonly IMarketDataProvider[] _providers; private readonly IQuoteCache _cache; + // 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 = []; + // Default: an in-memory cache. The call site in MarketExtensionCommandsProvider uses this. public MarketRepository(params IMarketDataProvider[] providers) : this(new InMemoryQuoteCache(), providers) { } @@ -34,12 +47,19 @@ public MarketRepository(IQuoteCache cache, IMarketDataProvider[] providers) _cache = cache; _providers = providers; - // A demo-mode flip swaps the data SOURCE, so every cached price is now wrong (it came from the - // other source). Clear so it can neither linger nor be preserved by keep-last-good. replay:false — - // construction is a no-op (the cache is empty then). Subscribed here, before any surface, so on a - // flip this Clear runs before a surface's own re-fetch. Never disposed: the repository lives for - // the whole process. - _ = MarketSettingsManager.Instance.DemoModeChanged.Subscribe(_ => _cache.Clear(), replayOnSubscribe: false); + // 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 @@ -118,39 +138,128 @@ public async Task RefreshAsync( _cache.Upsert(quote, keepLastGood); } - // Observe a FIXED set of instruments as one cache-backed list observable: on subscribe, fetch any - // not-yet-cached symbols once (write-through, fire-and-forget), 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. Defer runs the fetch-missing per subscribe. + // 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) - => Observable.Defer(() => + => 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 symbols once + // (write-through, fire-and-forget). 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); var missing = instruments.Where(i => _cache.Get(i.Symbol) is null).ToList(); if (missing.Count > 0) _ = RefreshSafe(missing); // fire-and-forget; writes through → observers see it land - if (instruments.Count == 0) - return Observable.Return>([]); + IObservable> inner = instruments.Count == 0 + ? Observable.Return>([]) + : Observable + .CombineLatest(instruments.Select(i => _cache.Observe(i.Symbol).AsObservable())) + .Select(quotes => (IReadOnlyList)quotes.OfType().ToList()); - return Observable - .CombineLatest(instruments.Select(i => _cache.Observe(i.Symbol).AsObservable())) - .Select(quotes => (IReadOnlyList)quotes.OfType().ToList()); + return new CompositeDisposable( + inner.Subscribe(observer), + Disposable.Create(() => Unregister(instruments))); }); - // Membership-aware: re-projects (Switch) whenever the set changes — fetching newly-added symbols and - // dropping departed ones for free. The shape the priced pages / dock bands will consume next pass. - public IObservable> ObserveQuotes( - StateFlow> instruments) - => instruments.AsObservable().Select(set => ObserveQuotes(set)).Switch(); - - // RefreshAsync with its exceptions swallowed + logged — for the fire-and-forget fetch-missing path, - // where a provider throw would otherwise become an unobserved task exception. + // 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); } } + // --- 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() + { + _cache.Clear(); + var observed = ObservedInstruments(); + 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. diff --git a/MarketExtension/MarketExtensionCommandsProvider.cs b/MarketExtension/MarketExtensionCommandsProvider.cs index 345468b..7eda6d3 100644 --- a/MarketExtension/MarketExtensionCommandsProvider.cs +++ b/MarketExtension/MarketExtensionCommandsProvider.cs @@ -42,6 +42,11 @@ 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). + // + // NOTE: FavoritesDockPage was crashing Command Palette because its ObserveQuotes subscription ran + // the full reactive graph synchronously on the host's band-activation thread. Re-enabled here while + // testing a fix that subscribes OFF that thread (SubscribeOn in FavoritesDockPage). If the crash + // returns, pull this line again. _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..f44a4c1 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,49 @@ 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; + private IDisposable? _quotesSubscription; 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. + _quotesSubscription = _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"); + _quotesSubscription?.Dispose(); + _quotesSubscription = null; + Log.Info("Dock", "stopped observing favorites"); } } @@ -75,18 +73,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 +104,13 @@ 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; RaiseItemsChanged(0); } } From db452fa5b1cbd5d8d11e27fc6104d35f0adc9ed2 Mon Sep 17 00:00:00 2001 From: Costa Fotiadis Date: Sat, 27 Jun 2026 15:28:15 +0300 Subject: [PATCH 03/10] cahe work in progress --- MarketExtension/MarketExtensionCommandsProvider.cs | 10 ++++++---- MarketExtension/Pages/FavoritesDockPage.cs | 14 ++++++++++---- 2 files changed, 16 insertions(+), 8 deletions(-) diff --git a/MarketExtension/MarketExtensionCommandsProvider.cs b/MarketExtension/MarketExtensionCommandsProvider.cs index 7eda6d3..d39dff6 100644 --- a/MarketExtension/MarketExtensionCommandsProvider.cs +++ b/MarketExtension/MarketExtensionCommandsProvider.cs @@ -43,10 +43,12 @@ 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). // - // NOTE: FavoritesDockPage was crashing Command Palette because its ObserveQuotes subscription ran - // the full reactive graph synchronously on the host's band-activation thread. Re-enabled here while - // testing a fix that subscribes OFF that thread (SubscribeOn in FavoritesDockPage). If the crash - // returns, pull this line again. + // 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 f44a4c1..0f038ef 100644 --- a/MarketExtension/Pages/FavoritesDockPage.cs +++ b/MarketExtension/Pages/FavoritesDockPage.cs @@ -33,7 +33,12 @@ internal sealed partial class FavoritesDockPage : ListPage, INotifyItemsChanged private UiQuote[]? _quotes; // latest cache emission, projected for rendering; null before the first private event TypedEventHandler? _itemsChanged; - private IDisposable? _quotesSubscription; + // 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 { @@ -48,14 +53,15 @@ event TypedEventHandler INotifyItemsChanged.Item // 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. - _quotesSubscription = _repository.ObserveQuotes(WatchlistStore.Instance.Favorites).Subscribe(OnQuotesChanged); + _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; - _quotesSubscription?.Dispose(); - _quotesSubscription = null; + foreach (var subscription in _subscriptions) + subscription.Dispose(); + _subscriptions.Clear(); Log.Info("Dock", "stopped observing favorites"); } } From 6deab6173a0aad91d65c234f0c0a09e8f162f0c1 Mon Sep 17 00:00:00 2001 From: Costa Fotiadis Date: Sat, 27 Jun 2026 15:45:14 +0300 Subject: [PATCH 04/10] cahe work in progress --- MarketExtension/Pages/PortfolioDockPage.cs | 181 +++++++++------------ 1 file changed, 79 insertions(+), 102 deletions(-) diff --git a/MarketExtension/Pages/PortfolioDockPage.cs b/MarketExtension/Pages/PortfolioDockPage.cs index 1bf740f..34fcaed 100644 --- a/MarketExtension/Pages/PortfolioDockPage.cs +++ b/MarketExtension/Pages/PortfolioDockPage.cs @@ -16,57 +16,59 @@ 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. + _subscriptions.Add(_repository.ObserveQuotes(PortfolioStore.Instance.Instruments) + .Subscribe(quotes => _ = OnQuotesChangedAsync(quotes))); + 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 +85,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,8 +96,12 @@ 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. + // 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)) @@ -108,82 +113,54 @@ public override IListItem[] GetItems() ]; } - 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. Launched fire-and-forget from the synchronous + // Subscribe lambda (`_ = ...`), with its own try/catch so no exception escapes onto the pool thread. + private async Task OnQuotesChangedAsync(IReadOnlyList quotes) { - _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) + { + _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); + + _portfolio = BuildPortfolio(positions, quotes, preferred); 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 (Exception ex) { - 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); + 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) From 57316593ceef928ad0d0bd9a21d6c45efaad2987 Mon Sep 17 00:00:00 2001 From: Costa Fotiadis Date: Sat, 27 Jun 2026 16:11:38 +0300 Subject: [PATCH 05/10] cahe work in progress --- CLAUDE.md | 64 ++++++++++++++----- MarketExtension/Data/MarketRepository.cs | 73 +++++++++++++++++++--- MarketExtension/Pages/FavoritesDockPage.cs | 1 + MarketExtension/Pages/PortfolioDockPage.cs | 3 + 4 files changed, 119 insertions(+), 22 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index f3bda29..2f753f2 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -233,6 +233,24 @@ cascading CS0534 ("does not implement … `GetTypeInfo`") onto **every** context ## Current Status / Next Steps +- **Done (latest): 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 **`IQuoteCache`** (in-memory `InMemoryQuoteCache` now — **real, no longer stubbed** — swappable for a DB-backed one later) holds the latest `DomainQuote` per symbol as **observable** state. `MarketRepository` is @@ -244,15 +262,15 @@ cascading CS0534 ("does not implement … `GetTypeInfo`") onto **every** context 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 - so far: `FavoritesDockPage`** — now a **pure observer** (no `PollTicker`/demo subscription, no + 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. **NOT yet migrated** (still on per-surface fetch+poll, - filling the cache only via write-through): `PortfolioDockPage`, the `PricedListPage` trio - (`WatchlistPage`/`FavoritesPage`/`PortfolioPage`), plus the `SymbolDetailPage` chart (candles, separate path). - **The favorites dock is the template for the rest.** See "Shared quote cache + repository orchestration" for - the full design + the migration checklist. + activates without hanging CmdPal and shows live prices. (`PortfolioDockPage` was migrated the following round — + see the "Done (latest)" 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. - **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 @@ -617,17 +635,29 @@ orchestrating all fetching/polling, so the same symbol is one cache entry everyo 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 so far: `FavoritesDockPage` only — the template.** It is now a **pure observer**: its whole +**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). +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. **To migrate the next surface (the checklist):** 1. Replace its `MarketRepository.GetQuotesAsync` call + private price cache with a subscription to @@ -642,12 +672,16 @@ live prices (this is the build that confirmed the `ObserveOn` deadlock fix). 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.** Carry over carefully: the **freshness clock** (`_lastFullPriceTicks` / - `RefreshStaleQuotes`) and **`OnPriceCacheUpdated`** (Portfolio's FX `PrimeAsync`) must be driven by - *fetch completion*, NOT by observe emissions — value-equality dedup hides an unchanged re-fetch, so keying - "a fetch happened" off an emission would miss it. `LeadingRows` must recompute on each emission. -7. `PortfolioDockPage` additionally `await`s `CurrencyConverter.PrimeAsync` before rolling up — keep that - inside the observe handler (one paint), not a separate fetch. + 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:** un-migrated surfaces still fetch+poll themselves AND fill the cache via write-through, so a symbol they share with the dock is briefly fetched by both the repo loop (for the dock) and the surface's diff --git a/MarketExtension/Data/MarketRepository.cs b/MarketExtension/Data/MarketRepository.cs index 8eef4bd..cff6ef0 100644 --- a/MarketExtension/Data/MarketRepository.cs +++ b/MarketExtension/Data/MarketRepository.cs @@ -1,4 +1,5 @@ using System; +using System.Collections.Concurrent; using System.Collections.Generic; using System.Linq; using System.Reactive.Concurrency; @@ -37,6 +38,14 @@ internal sealed class MarketRepository 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 InMemoryQuoteCache(), providers) { } @@ -79,8 +88,7 @@ public async Task> GetQuotesAsync( IReadOnlyList instruments, CancellationToken ct = default) { var ordered = await FetchMergeAsync(instruments, ct).ConfigureAwait(false); - foreach (var quote in ordered) - _cache.Upsert(quote); // keep-last-good + WriteThrough(ordered, keepLastGood: true); return ordered; } @@ -134,8 +142,20 @@ 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 for both GetQuotesAsync and RefreshAsync. + private void WriteThrough(IReadOnlyList ordered, bool keepLastGood) + { + var now = Environment.TickCount64; foreach (var quote in ordered) + { _cache.Upsert(quote, keepLastGood); + _lastFetchTicks[WatchlistStore.Normalize(quote.Symbol)] = now; + } } // Observe a FIXED set of instruments as one cache-backed list observable. Both public overloads append @@ -158,8 +178,9 @@ public IObservable> ObserveQuotes( .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 symbols once - // (write-through, fire-and-forget). Then CombineLatest the per-symbol cache flows into one ordered list + // 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). @@ -167,9 +188,25 @@ private IObservable> ObserveQuotesCore(IReadOnlyList< => Observable.Create>(observer => { Register(instruments); - var missing = instruments.Where(i => _cache.Get(i.Symbol) is null).ToList(); - if (missing.Count > 0) - _ = RefreshSafe(missing); // fire-and-forget; writes through → observers see it land + // 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 => _cache.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>([]) @@ -190,6 +227,26 @@ private async Task RefreshSafe(IReadOnlyList instruments) 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 (_cache.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 @@ -211,6 +268,8 @@ private void OnDemoModeFlip() { _cache.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); } diff --git a/MarketExtension/Pages/FavoritesDockPage.cs b/MarketExtension/Pages/FavoritesDockPage.cs index 0f038ef..2c70ba9 100644 --- a/MarketExtension/Pages/FavoritesDockPage.cs +++ b/MarketExtension/Pages/FavoritesDockPage.cs @@ -117,6 +117,7 @@ private void OnQuotesChanged(IReadOnlyList quotes) _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 34fcaed..41eb60e 100644 --- a/MarketExtension/Pages/PortfolioDockPage.cs +++ b/MarketExtension/Pages/PortfolioDockPage.cs @@ -126,6 +126,7 @@ private async Task OnQuotesChangedAsync(IReadOnlyList quotes) // 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) { + Log.Info("Dock", $"portfolio: {positions.Count} holding(s) but no cached prices yet — spinner"); _portfolio = null; IsLoading = true; RaiseItemsChanged(0); @@ -144,6 +145,8 @@ private async Task OnQuotesChangedAsync(IReadOnlyList quotes) await CurrencyConverter.Instance.PrimeAsync(preferred, natives); _portfolio = BuildPortfolio(positions, quotes, preferred); + Log.Info("Dock", + $"portfolio rolled up: {_portfolio.FormatTotalValue()} {_portfolio.FormatTotalChange()} across {positions.Count} holding(s)"); IsLoading = false; RaiseItemsChanged(0); } From 904264b9817b15c19d1cb12eb575150d90851826 Mon Sep 17 00:00:00 2001 From: Costa Fotiadis Date: Sat, 27 Jun 2026 16:36:15 +0300 Subject: [PATCH 06/10] cahe work in progress --- CLAUDE.md | 78 +++-- MarketExtension/Data/MarketRepository.cs | 5 +- MarketExtension/Helpers/CurrencyConverter.cs | 4 +- MarketExtension/Pages/PortfolioPage.cs | 38 ++- MarketExtension/Pages/PricedListPage.cs | 282 ++++++------------- 5 files changed, 165 insertions(+), 242 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index 2f753f2..7d970e6 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -82,7 +82,7 @@ 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**, 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 `IQuoteCache`** (in-memory by default; an injectable ctor overload `MarketRepository(IQuoteCache, providers[])` for a future DB-backed cache): `GetQuotesAsync`/`RefreshAsync` **write 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. See "Shared quote cache + repository orchestration (in progress)". | +| `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 `IQuoteCache`** (in-memory by default; an injectable ctor overload `MarketRepository(IQuoteCache, providers[])` for a future DB-backed cache): `GetQuotesAsync`/`RefreshAsync` **write 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/IQuoteCache.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). `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/InMemoryQuoteCache.cs` | The **real, live** in-memory `IQuoteCache` (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 `IQuoteCache` later via the repo's injectable ctor — no repository/UI change. | | `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). | @@ -94,7 +94,7 @@ the same three layers and the same provider seam — see the "Symbol detail + li | `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. 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`. **Migration in progress:** `MarketRepository` now holds ONE process-lifetime subscription that polls the *observed* set centrally, so a migrated surface (the favorites dock) no longer subscribes here at all — it just observes the cache. Un-migrated surfaces (the `PricedListPage` trio, the portfolio dock, the chart) still subscribe directly for now. | +| `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). | @@ -233,7 +233,28 @@ cascading CS0534 ("does not implement … `GetTypeInfo`") onto **every** context ## Current Status / Next Steps -- **Done (latest): migrated `PortfolioDockPage` to a pure cache observer + added the stale-on-subscribe +- **Done (latest): 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). 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 @@ -267,10 +288,12 @@ cascading CS0534 ("does not implement … `GetTypeInfo`") onto **every** context `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 (latest)" 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. + 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 @@ -594,7 +617,7 @@ 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 (cache live; favorites dock migrated, rest pending) +### 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 @@ -611,7 +634,10 @@ orchestrating all fetching/polling, so the same symbol is one cache entry everyo `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 `IQuoteCache`. `GetQuotesAsync`/`RefreshAsync` write through - (so the cache fills even from un-migrated surfaces). `ObserveQuotes(IReadOnlyList)` and + (the single cache-fill path, now driven by the poll loop + the observe-subscribe fetch; note `GetQuotesAsync` + has **no direct surface callers** anymore — every surface observes, and the observe/poll paths fetch via + `RefreshAsync`/`RefreshSafe`, so the public `GetQuotesAsync` is kept only as coordinator API/for tests). + `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). @@ -659,7 +685,24 @@ handler `await`s `CurrencyConverter.PrimeAsync` before rolling up `UiPortfolio` `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. -**To migrate the next surface (the checklist):** +**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. @@ -683,15 +726,18 @@ Both bands hold subscriptions in a `List` (double-`add`-safe). Buil 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:** un-migrated surfaces still fetch+poll themselves AND fill the cache via write-through, -so a symbol they share with the dock is briefly fetched by both the repo loop (for the dock) and the surface's -own poll — harmless (the cache dedups by value), and it collapses as each surface migrates. +**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 migrated surfaces** by "Shared quote cache + repository orchestration" above — the favorites -> dock no longer polls itself; the repository's single loop polls the observed set. The model below still -> describes the **un-migrated** surfaces (the `PricedListPage` trio, the portfolio dock, the chart). +> **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- diff --git a/MarketExtension/Data/MarketRepository.cs b/MarketExtension/Data/MarketRepository.cs index cff6ef0..ba1297d 100644 --- a/MarketExtension/Data/MarketRepository.cs +++ b/MarketExtension/Data/MarketRepository.cs @@ -25,7 +25,10 @@ namespace MarketExtension; // 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. -// Surfaces still poll/fetch individually today; moving them fully onto observe-only is a later pass. +// 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; 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/Pages/PortfolioPage.cs b/MarketExtension/Pages/PortfolioPage.cs index 51a9e61..3977040 100644 --- a/MarketExtension/Pages/PortfolioPage.cs +++ b/MarketExtension/Pages/PortfolioPage.cs @@ -9,23 +9,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 +92,21 @@ 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. + protected override async Task OnQuotesProjectingAsync(IReadOnlyList quotes) { 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 () => - { + if (natives.Length > 0) await CurrencyConverter.Instance.PrimeAsync(preferred, natives); - RaiseItemsChanged(0); - }); } protected override IListItem[] EmptyState() => diff --git a/MarketExtension/Pages/PricedListPage.cs b/MarketExtension/Pages/PricedListPage.cs index a4cc2ee..ee022e3 100644 --- a/MarketExtension/Pages/PricedListPage.cs +++ b/MarketExtension/Pages/PricedListPage.cs @@ -1,7 +1,6 @@ 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; @@ -10,36 +9,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 +51,27 @@ 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. Fire-and-forget the async handler from the (synchronous) Subscribe lambda so a subclass + // can await an async projection (the Portfolio screen primes FX rates) before the repaint. + _subscriptions.Add(Repository.ObserveQuotes(_instruments) + .Subscribe(quotes => _ = OnQuotesChangedAsync(quotes))); // 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 +79,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 +91,6 @@ protected PricedListPage(MarketRepository repository, StateFlow