diff --git a/src/jarvis_finance/market_data/prices.py b/src/jarvis_finance/market_data/prices.py index 5506e61..4559f09 100644 --- a/src/jarvis_finance/market_data/prices.py +++ b/src/jarvis_finance/market_data/prices.py @@ -6,6 +6,7 @@ from decimal import Decimal, InvalidOperation from sqlite3 import Connection from typing import Protocol from urllib import error, parse, request +import hashlib import json from jarvis_finance.imports.common import stable_id, utc_now @@ -450,6 +451,137 @@ def _corporate_action_status(conn: Connection, *, instrument_id: str, price_date return "none_known" +def _economic_price_payload( + *, + instrument_id: str, + provider: str, + provider_symbol: str | None, + provider_market: str | None, + price_type: str, + close: Decimal | None, + adjusted_close: Decimal | None, + currency: str, + provider_timestamp: str, + quality_status: str, + source_reference: str, +) -> dict[str, str | None]: + return { + "provider": provider.lower(), + "instrument_id": instrument_id, + "provider_symbol": provider_symbol, + "provider_market": provider_market, + "price_type": price_type, + "close": format(close, "f") if close is not None else "", + "adjusted_close": format(adjusted_close, "f") if adjusted_close is not None else None, + "currency": currency.upper(), + "provider_timestamp": provider_timestamp, + "source_reference": source_reference, + "quality_status": quality_status, + } + + +def _normalise_provider_timestamp(value: str) -> str: + """Canonicalise equivalent timestamp spellings without inventing precision.""" + text = value.strip() + if len(text) == 10: + return date.fromisoformat(text).isoformat() + parsed = datetime.fromisoformat(text.replace("Z", "+00:00")) + if parsed.tzinfo is None: + parsed = parsed.replace(tzinfo=timezone.utc) + return parsed.astimezone(timezone.utc).isoformat() + + +def _store_economic_price_observation( + conn: Connection, + *, + instrument_id: str, + provider: str, + provider_symbol: str | None, + provider_market: str | None, + price_type: str, + close: Decimal | None, + adjusted_close: Decimal | None, + currency: str, + provider_timestamp: str, + quality_status: str, + source_reference: str | None, + job_reference: str | None, + created_at: str, +) -> str: + # Serialize identity/version allocation across parallel provider workers. + # The caller commits the complete observation + current-price projection atomically. + if not conn.in_transaction: + conn.execute("BEGIN IMMEDIATE") + canonical_provider_timestamp = _normalise_provider_timestamp(provider_timestamp) + origin_reference = source_reference or ":".join( + part for part in (provider.lower(), provider_symbol or "", provider_market or "") if part + ) + payload = _economic_price_payload( + instrument_id=instrument_id, + provider=provider, + provider_symbol=provider_symbol, + provider_market=provider_market, + price_type=price_type, + close=close, + adjusted_close=adjusted_close, + currency=currency, + provider_timestamp=canonical_provider_timestamp, + quality_status=quality_status, + source_reference=origin_reference, + ) + payload_json = json.dumps(payload, sort_keys=True, separators=(",", ":")) + payload_hash = hashlib.sha256(payload_json.encode()).hexdigest() + source_observation_id = stable_id( + "market-source-observation", + provider.lower(), + instrument_id, + canonical_provider_timestamp, + ) + same = conn.execute( + """SELECT observation_id FROM market_price_observations + WHERE source_observation_id=? AND economic_payload_hash=?""", + (source_observation_id, payload_hash), + ).fetchone() + if same: + return str(same["observation_id"]) + predecessor = conn.execute( + """SELECT observation_id,payload_version,economic_payload_json + FROM market_price_observations WHERE source_observation_id=? + ORDER BY payload_version DESC LIMIT 1""", + (source_observation_id,), + ).fetchone() + version = int(predecessor["payload_version"]) + 1 if predecessor else 1 + observation_id = stable_id("market-observation", source_observation_id, payload_hash) + conn.execute( + """INSERT INTO market_price_observations( + observation_id,source_observation_id,payload_version,supersedes_observation_id, + instrument_id,provider,provider_symbol,provider_market,price_type,close,adjusted_close, + currency,provider_timestamp,source_reference,quality_status,economic_payload_json, + economic_payload_hash,created_at,job_reference + ) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", + ( + observation_id, source_observation_id, version, + str(predecessor["observation_id"]) if predecessor else None, + instrument_id, provider, provider_symbol, provider_market, price_type, + payload["close"], payload["adjusted_close"], payload["currency"], canonical_provider_timestamp, + origin_reference, quality_status, payload_json, payload_hash, created_at, job_reference, + ), + ) + if predecessor: + record_audit_event( + conn, + source="market_price_observation_v2", + action="market_price_observation_corrected", + entity_type="market_price_observation", + entity_id=observation_id, + old_values={"supersedes_observation_id": predecessor["observation_id"]}, + new_values={"source_observation_id": source_observation_id, "payload_version": version}, + confirmed=True, + created_by="system", + ) + return observation_id + + def store_market_price( conn: Connection, *, @@ -467,6 +599,75 @@ def store_market_price( fetched_at: str | None = None, price_type: str = "unadjusted_close", run_id: str | None = None, +) -> str: + owns_transaction = not conn.in_transaction + if owns_transaction: + conn.execute("BEGIN IMMEDIATE") + conn.execute("UPDATE market_price_observation_mutex SET touched=touched WHERE mutex_id=1") + corp_status = ( + _corporate_action_status( + conn, + instrument_id=instrument_id, + price_date=price_date, + provider=provider, + close=close, + ) + if quality_status == "fresh" + else "not_checked" + ) + if not conn.in_transaction: + owns_transaction = True + conn.execute("BEGIN IMMEDIATE") + conn.execute("UPDATE market_price_observation_mutex SET touched=touched WHERE mutex_id=1") + conn.execute("SAVEPOINT market_price_write") + try: + market_price_id = _store_market_price_locked( + conn, + instrument_id=instrument_id, + price_date=price_date, + close=close, + currency=currency, + provider=provider, + provider_symbol=provider_symbol, + provider_market=provider_market, + price_timestamp=price_timestamp, + adjusted_close=adjusted_close, + quality_status=quality_status, + error_message=error_message, + fetched_at=fetched_at, + price_type=price_type, + run_id=run_id, + corp_status=corp_status, + ) + conn.execute("RELEASE SAVEPOINT market_price_write") + conn.commit() + return market_price_id + except Exception: + conn.execute("ROLLBACK TO SAVEPOINT market_price_write") + conn.execute("RELEASE SAVEPOINT market_price_write") + if owns_transaction: + conn.rollback() + raise + + +def _store_market_price_locked( + conn: Connection, + *, + instrument_id: str, + price_date: str, + close: Decimal | None, + currency: str, + provider: str, + provider_symbol: str | None, + provider_market: str | None = None, + price_timestamp: str | None = None, + adjusted_close: Decimal | None = None, + quality_status: str = "fresh", + error_message: str | None = None, + fetched_at: str | None = None, + price_type: str = "unadjusted_close", + run_id: str | None = None, + corp_status: str = "not_checked", ) -> str: now = utc_now() fetched = fetched_at or now @@ -474,8 +675,26 @@ def store_market_price( "SELECT market_price_id, close, currency, quality_status FROM market_prices WHERE instrument_id=? AND price_date=? AND provider=?", (instrument_id, price_date, provider), ).fetchone() - corp_status = _corporate_action_status(conn, instrument_id=instrument_id, price_date=price_date, provider=provider, close=close) if quality_status == "fresh" else "not_checked" - market_price_id = stable_id("marketprice", instrument_id, price_date, provider, provider_symbol or "", now) + provider_observed_at = price_timestamp or price_date + _store_economic_price_observation( + conn, + instrument_id=instrument_id, + provider=provider, + provider_symbol=provider_symbol, + provider_market=provider_market, + price_type=price_type, + close=close, + adjusted_close=adjusted_close, + currency=currency, + provider_timestamp=provider_observed_at, + quality_status=quality_status, + source_reference=None, + job_reference=run_id, + created_at=now, + ) + market_price_id = str(existing["market_price_id"]) if existing else stable_id( + "marketprice", instrument_id, price_date, provider, provider_symbol or "", now + ) conn.execute( """ INSERT INTO market_prices( @@ -520,7 +739,6 @@ def store_market_price( create_alert(conn, priority="warnung", category="market_data", entity_type="instrument", entity_id=instrument_id, rule_id=rule, message="Instrument local market price is not fresh.", evidence={"provider": provider, "provider_symbol": provider_symbol, "price_date": price_date, "quality_status": quality_status}, fingerprint=f"{rule}:{provider}:{provider_symbol}") elif quality_status == "fresh": _resolve_market_alerts(conn, instrument_id=instrument_id) - conn.commit() return market_price_id diff --git a/src/jarvis_finance/services/asset_price_refresh.py b/src/jarvis_finance/services/asset_price_refresh.py index a822fed..fe06cea 100644 --- a/src/jarvis_finance/services/asset_price_refresh.py +++ b/src/jarvis_finance/services/asset_price_refresh.py @@ -3,6 +3,7 @@ from __future__ import annotations import hashlib import json import uuid +from dataclasses import dataclass, field from datetime import UTC, datetime, timedelta from pathlib import Path from sqlite3 import Connection, SQLITE_DELETE, SQLITE_DENY, SQLITE_INSERT, SQLITE_OK, SQLITE_UPDATE @@ -26,6 +27,26 @@ PROTECTED_TABLES = ( ) +@dataclass(frozen=True) +class SourceRunResult: + stale_candidates: int = 0 + updated_count: int = 0 + fresh_unchanged_count: int = 0 + stale_remaining_count: int = 0 + failed_count: int = 0 + diagnostics: tuple[str, ...] = field(default_factory=tuple) + + def __iter__(self): + yield self.stale_candidates + yield self.updated_count + + +def _source_result(value: SourceRunResult | tuple[int, int]) -> SourceRunResult: + if isinstance(value, SourceRunResult): + return value + return SourceRunResult(stale_candidates=int(value[0]), updated_count=int(value[1])) + + def _deny_protected_dml( action: int, table: str | None, @@ -90,6 +111,10 @@ def _status_payload(conn: Connection, job_id: str) -> dict[str, Any]: "stale_candidates": int(row["stale_candidates"]), "updated_count": int(row["updated_count"]), "error_code": str(row["error_code"]) if row["error_code"] else None, + "fresh_unchanged_count": int(row["fresh_unchanged_count"]), + "stale_remaining_count": int(row["stale_remaining_count"]), + "failed_count": int(row["failed_count"]), + "diagnostics": json.loads(str(row["diagnostics_json"] or "[]")), "started_at": str(row["started_at"]) if row["started_at"] else None, "completed_at": str(row["completed_at"]) if row["completed_at"] else None, } @@ -97,6 +122,15 @@ def _status_payload(conn: Connection, job_id: str) -> dict[str, Any]: ], "wealth_snapshot_created": bool(job["wealth_snapshot_id"]), "audit_recorded": bool(job["audit_id"]), + "successful_assets": sum(int(row["updated_count"]) for row in sources), + "fresh_unchanged_assets": sum(int(row["fresh_unchanged_count"]) for row in sources), + "stale_assets": sum(int(row["stale_remaining_count"]) for row in sources), + "failed_assets": sum(int(row["failed_count"]) for row in sources), + "next_action": ( + "Diagnose prüfen und nur betroffene Quelle erneut versuchen." + if any(int(row["failed_count"]) or int(row["stale_remaining_count"]) for row in sources) + else "Keine Aktion nötig; alle verfügbaren Kurse sind aktuell." + ), "provider_calls_on_read": False, } @@ -134,7 +168,7 @@ def create_asset_price_refresh_job(conn: Connection, *, stale_hours: int = 24) - return _status_payload(conn, job_id), _database_path(conn) -def _equity_source(conn: Connection, stale_before: str) -> tuple[int, int]: +def _equity_source(conn: Connection, stale_before: str) -> SourceRunResult: response = refresh_equity_quotes_batch( conn, QuoteRefreshRequest( @@ -147,12 +181,33 @@ def _equity_source(conn: Connection, stale_before: str) -> tuple[int, int]: ), ) candidates = max(0, int(response.total) - int(response.cached)) - if response.errors and response.updated == 0 and candidates > 0: - raise RuntimeError("equity_provider_failed") - return candidates, int(response.updated) + failed = sum(1 for row in response.results if str(row.get("status")) in {"provider_error", "error"}) + stale = sum(1 for row in response.results if str(row.get("status")) in {"stale", "missing"}) + valuation_issues = [ + warning + for warning in response.warnings + if warning in {"portfolio_valuation_partial", "portfolio_valuation_failed"} + ] + return SourceRunResult( + stale_candidates=candidates, + updated_count=int(response.economic_updated), + fresh_unchanged_count=int(response.cached + response.updated - response.economic_updated), + stale_remaining_count=stale, + failed_count=max(failed, len(response.errors)) + len(valuation_issues), + diagnostics=tuple(sorted(set([*response.errors, *valuation_issues])))[:10], + ) -def _crypto_source(conn: Connection, stale_before: str) -> tuple[int, int]: +def _crypto_source(conn: Connection, stale_before: str) -> SourceRunResult: + held_asset_ids = [ + str(row["asset_id"]) + for row in conn.execute( + """SELECT a.asset_id FROM crypto_assets a + WHERE a.is_active=1 AND EXISTS( + SELECT 1 FROM crypto_holdings h WHERE h.asset_id=a.asset_id AND CAST(h.quantity AS REAL)<>0 + ) ORDER BY a.asset_id""" + ).fetchall() + ] stale_asset_ids = [ str(row["asset_id"]) for row in conn.execute( @@ -168,7 +223,7 @@ def _crypto_source(conn: Connection, stale_before: str) -> tuple[int, int]: ).fetchall() ] if not stale_asset_ids: - return 0, 0 + return SourceRunResult(fresh_unchanged_count=len(held_asset_ids)) cutoff = datetime.fromisoformat(stale_before.replace("Z", "+00:00")) if cutoff.tzinfo is None: cutoff = cutoff.replace(tzinfo=UTC) @@ -181,16 +236,30 @@ def _crypto_source(conn: Connection, stale_before: str) -> tuple[int, int]: asset_ids=stale_asset_ids, batch_size=100, ) - if result.error_count and result.updated_count == 0: - raise RuntimeError("crypto_provider_failed") - return len(stale_asset_ids), int(result.updated_count) + return SourceRunResult( + stale_candidates=len(stale_asset_ids), + updated_count=int(result.updated_count), + fresh_unchanged_count=max(0, len(held_asset_ids) - len(stale_asset_ids)) + int(result.cached_count), + stale_remaining_count=int(result.stale_count + result.missing_local_price_count), + failed_count=int(result.error_count), + diagnostics=tuple(result.errors[:10]), + ) -def _fx_source(conn: Connection, stale_before: str) -> tuple[int, int]: +def _fx_source(conn: Connection, stale_before: str) -> SourceRunResult: from jarvis_finance.fx.providers import FrankfurterFxProvider, TwelveDataFxProvider from jarvis_finance.fx.rates import resolve_fx_rate_to_chf cutoff_date = stale_before[:10] + all_currencies = [ + str(row["currency"]).upper() + for row in conn.execute( + """SELECT DISTINCT upper(i.currency) currency + FROM instruments i + WHERE i.is_active=1 AND upper(COALESCE(i.currency,'CHF'))!='CHF' + ORDER BY currency""" + ).fetchall() + ] currencies = [ str(row["currency"]).upper() for row in conn.execute( @@ -207,9 +276,18 @@ def _fx_source(conn: Connection, stale_before: str) -> tuple[int, int]: ).fetchall() ] updated = 0 + provider_unchanged = 0 failures = 0 for currency in currencies: try: + before = conn.execute( + """SELECT rate_date,rate,provider,rate_type,quality_status + FROM fx_rates + WHERE base_currency=? AND quote_currency='CHF' + ORDER BY rate_date DESC,COALESCE(fetched_at,created_at) DESC LIMIT 1""", + (currency,), + ).fetchone() + before_economic = tuple(before) if before is not None else None result = resolve_fx_rate_to_chf( conn, base_currency=currency, @@ -218,16 +296,31 @@ def _fx_source(conn: Connection, stale_before: str) -> tuple[int, int]: persist=True, resolve_fixed=True, ) - updated += int(result.status == "ok") + after = conn.execute( + """SELECT rate_date,rate,provider,rate_type,quality_status + FROM fx_rates + WHERE base_currency=? AND quote_currency='CHF' + ORDER BY rate_date DESC,COALESCE(fetched_at,created_at) DESC LIMIT 1""", + (currency,), + ).fetchone() + after_economic = tuple(after) if after is not None else None + if result.status == "ok" and after_economic != before_economic: + updated += 1 + elif result.status == "ok": + provider_unchanged += 1 except Exception: failures += 1 conn.commit() - if failures and updated == 0: - raise RuntimeError("fx_provider_failed") - return len(currencies), updated + return SourceRunResult( + stale_candidates=len(currencies), + updated_count=updated, + fresh_unchanged_count=max(0, len(all_currencies) - len(currencies)) + provider_unchanged, + stale_remaining_count=failures, + failed_count=failures, + ) -DEFAULT_RUNNERS: dict[str, Callable[[Connection, str], tuple[int, int]]] = { +DEFAULT_RUNNERS: dict[str, Callable[[Connection, str], SourceRunResult | tuple[int, int]]] = { "equity": _equity_source, "crypto": _crypto_source, "fx": _fx_source, @@ -238,7 +331,7 @@ def run_asset_price_refresh( db_path: str, job_id: str, *, - runners: dict[str, Callable[[Connection, str], tuple[int, int]]] | None = None, + runners: dict[str, Callable[[Connection, str], SourceRunResult | tuple[int, int]]] | None = None, ) -> None: """Background worker with source isolation, stored progress and mutation guard.""" conn = connect(db_path) @@ -264,6 +357,9 @@ def run_asset_price_refresh( conn.set_authorizer(_deny_protected_dml) completed = 0 failures = 0 + total_updated = 0 + total_fresh_unchanged = 0 + total_candidates = 0 for source in SOURCES: started = _now() conn.execute( @@ -271,25 +367,42 @@ def run_asset_price_refresh( (started, job_id, source), ) conn.commit() - candidates = updated = 0 + result = SourceRunResult() status = "complete" error_code = None try: - candidates, updated = selected[source](conn, stale_before) - if candidates == 0: + result = _source_result(selected[source](conn, stale_before)) + if result.stale_candidates == 0 and result.fresh_unchanged_count == 0: status = "skipped" + if result.failed_count: + failures += 1 + status = "failed" + error_code = f"{source}_instrument_failures" + elif result.stale_remaining_count: + failures += 1 + status = "failed" + error_code = f"{source}_stale_remaining" except Exception as exc: if conn.in_transaction: conn.rollback() status = "failed" failures += 1 - error_code = str(exc)[:120] or type(exc).__name__ + result = SourceRunResult(failed_count=1, diagnostics=(type(exc).__name__,)) + error_code = f"{source}_refresh_failed" + total_updated += result.updated_count + total_fresh_unchanged += result.fresh_unchanged_count + total_candidates += result.stale_candidates completed += 1 conn.execute( """UPDATE asset_price_refresh_sources - SET status=?,stale_candidates=?,updated_count=?,error_code=?,completed_at=? + SET status=?,stale_candidates=?,updated_count=?,error_code=?,completed_at=?, + fresh_unchanged_count=?,stale_remaining_count=?,failed_count=?,diagnostics_json=? WHERE job_id=? AND source=?""", - (status, candidates, updated, error_code, _now(), job_id, source), + ( + status, result.stale_candidates, result.updated_count, error_code, _now(), + result.fresh_unchanged_count, result.stale_remaining_count, result.failed_count, + json.dumps(list(result.diagnostics)), job_id, source, + ), ) conn.execute( "UPDATE asset_price_refresh_jobs SET progress_completed=? WHERE job_id=?", @@ -299,16 +412,18 @@ def run_asset_price_refresh( if _protected_fingerprint(conn) != protected_before: raise RuntimeError("protected_holdings_or_transactions_mutated") - successful_sources = failures < len(SOURCES) + usable_result = total_updated > 0 or total_fresh_unchanged > 0 or (total_candidates == 0 and failures == 0) wealth_snapshot_id = None - if successful_sources: + if total_updated > 0: model = build_modelled_wealth_development(conn, period="1m") current = model.get("current") or {} wealth_snapshot_id = f"wealth-refresh-{uuid.uuid4().hex}" source_rows = [ dict(row) for row in conn.execute( - "SELECT source,status,stale_candidates,updated_count,error_code FROM asset_price_refresh_sources WHERE job_id=? ORDER BY source", + """SELECT source,status,stale_candidates,updated_count,error_code, + fresh_unchanged_count,stale_remaining_count,failed_count + FROM asset_price_refresh_sources WHERE job_id=? ORDER BY source""", (job_id,), ).fetchall() ] @@ -325,7 +440,7 @@ def run_asset_price_refresh( json.dumps(source_rows, sort_keys=True), ), ) - final_status = "complete" if failures == 0 else "failed" if failures == len(SOURCES) else "partial" + final_status = "complete" if failures == 0 else "partial" if usable_result else "failed" audit_id = record_audit_event( conn, source="asset_price_refresh_job_v1", diff --git a/src/jarvis_finance/services/market_service.py b/src/jarvis_finance/services/market_service.py index e479e34..8594c40 100644 --- a/src/jarvis_finance/services/market_service.py +++ b/src/jarvis_finance/services/market_service.py @@ -206,7 +206,12 @@ def _validate_historical_quote(quote: EquityPriceQuote, *, mapping: Any, target: return replace(quote, close=None, quality_status="exchange_mismatch", error_message="exchange_mismatch") elif actual_currency and expected_currency and actual_currency != expected_currency: return replace(quote, close=None, quality_status="currency_mismatch", error_message="currency_mismatch") - return replace(quote, currency=actual_currency or expected_currency, price_timestamp=quote_date.isoformat(), quality_status="fresh") + return replace( + quote, + currency=actual_currency or expected_currency, + price_timestamp=quote.price_timestamp or quote_date.isoformat(), + quality_status="fresh", + ) def refresh_equity_quote(conn: Connection, instrument_id: str, req: QuoteRefreshRequest) -> MarketQuoteResponse: @@ -227,7 +232,12 @@ def refresh_equity_quote(conn: Connection, instrument_id: str, req: QuoteRefresh quote = _validate_historical_quote(quote, mapping=mapping_data, target=target, requested_provider=provider_name) quality = quote.quality_status if quote.close is not None else _quality_from_error(quote.error_message, quote.quality_status) ts = quote.price_timestamp or target.isoformat() + economic_observation_created = False if not req.dry_run and quote.close is not None and quality == "fresh": + observation_count_before = conn.execute( + "SELECT COUNT(*) FROM market_price_observations WHERE instrument_id=?", + (instrument_id,), + ).fetchone()[0] store_market_price( conn, instrument_id=instrument_id, @@ -242,6 +252,10 @@ def refresh_equity_quote(conn: Connection, instrument_id: str, req: QuoteRefresh quality_status=quality, error_message=None, ) + economic_observation_created = conn.execute( + "SELECT COUNT(*) FROM market_price_observations WHERE instrument_id=?", + (instrument_id,), + ).fetchone()[0] > observation_count_before upsert_equity_price_point( conn, instrument_id=instrument_id, @@ -262,6 +276,7 @@ def refresh_equity_quote(conn: Connection, instrument_id: str, req: QuoteRefresh provider_symbol=quote.provider_symbol or provider_symbol, fetched_at=ts, quality_status=quality, + economic_observation_created=economic_observation_created, warnings=warnings + ([quote.error_message] if quote.error_message else []), ) @@ -326,7 +341,7 @@ def refresh_equity_quotes_batch(conn: Connection, req: QuoteRefreshRequest) -> M row_states.sort(key=lambda item: (item[1], str(item[0]["name"] or ""))) bounded_limit = max(1, min(int(req.limit or 100), 500)) row_states = row_states[:bounded_limit] - updated = skipped = cached = processed = 0 + updated = economic_updated = skipped = cached = processed = 0 would_update = provider_calls = 0 successful_instruments: set[str] = set() result_dates: list[str] = [] @@ -374,7 +389,13 @@ def refresh_equity_quotes_batch(conn: Connection, req: QuoteRefreshRequest) -> M processed += 1 provider_calls += 1 last_call_at = time.monotonic() - quote = refresh_equity_quote(conn, row["instrument_id"], req) + try: + quote = refresh_equity_quote(conn, row["instrument_id"], req) + except Exception as exc: + quote = MarketQuoteResponse( + quality_status="provider_error", + warnings=[type(exc).__name__], + ) if quote.latest_price is not None and quote.quality_status == "fresh": break if quote.quality_status not in {"rate_limited", "network_error"} or attempt > req.max_retries: @@ -389,6 +410,8 @@ def refresh_equity_quotes_batch(conn: Connection, req: QuoteRefreshRequest) -> M would_update += 1 else: updated += 1 + if quote.economic_observation_created: + economic_updated += 1 else: skipped += 1 code = quote.warnings[0] if quote.warnings else quote.quality_status @@ -410,7 +433,18 @@ def refresh_equity_quotes_batch(conn: Connection, req: QuoteRefreshRequest) -> M or (req.dry_run and str(row["instrument_id"]) in successful_instruments) for row in all_rows ) - if valued == coverage_total and coverage_total > 0 and not req.dry_run: + partial_run_exists = bool( + conn.execute( + "SELECT 1 FROM market_data_runs WHERE source_key='daily_market_fx_v4' AND as_of=? AND status='partial' LIMIT 1", + (target,), + ).fetchone() + ) + if ( + valued == coverage_total + and coverage_total > 0 + and not req.dry_run + and (economic_updated > 0 or partial_run_exists) + ): try: from jarvis_finance.services.portfolio_analytics import run_daily_market_valuation @@ -427,6 +461,7 @@ def refresh_equity_quotes_batch(conn: Connection, req: QuoteRefreshRequest) -> M completed_at=utc_now(), total=len(row_states), updated=updated, + economic_updated=economic_updated, skipped=skipped, warnings=warnings, errors=errors, __HERMES_CWD_8d46a20096ed__/home/agent/.hermes/worktrees/FinanceManager-sprint23.1__HERMES_CWD_8d46a20096ed__