diff --git a/src/jarvis_finance/market_data/prices.py b/src/jarvis_finance/market_data/prices.py index 5506e61..51bb870 100644 --- a/src/jarvis_finance/market_data/prices.py +++ b/src/jarvis_finance/market_data/prices.py @@ -1,88 +1,89 @@ from __future__ import annotations from dataclasses import dataclass, field from datetime import date, datetime, timedelta, timezone 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 from jarvis_finance.audit.log import record_audit_event from jarvis_finance.market_data.catalog import _runtime_secret_value from jarvis_finance.market_data.instruments import resolve_instrument_alerts, update_instrument_metadata from jarvis_finance.quality.alerts import create_alert @dataclass(frozen=True) class EquityPriceQuote: provider_symbol: str currency: str close: Decimal | None provider: str = "mock" provider_market: str | None = None price_timestamp: str | None = None adjusted_close: Decimal | None = None quality_status: str = "fresh" error_message: str | None = None @dataclass(frozen=True) class ProviderCapability: supports_latest: bool supports_historical_as_of: bool supports_exchange_suffix: bool rate_limit_policy: str @dataclass class MarketPriceRefreshResult: asset_class: str total_mappings: int = 0 updated_count: int = 0 skipped_count: int = 0 excluded_count: int = 0 cached_count: int = 0 stale_count: int = 0 warning_count: int = 0 error_count: int = 0 dry_run: bool = False warnings: list[str] = field(default_factory=list) errors: list[str] = field(default_factory=list) class EquityPriceProvider(Protocol): name: str def get_price(self, provider_symbol: str, *, price_date: str | None = None) -> EquityPriceQuote: ... class MockEquityPriceProvider: name = "mock" capability = ProviderCapability(True, True, True, "none") def __init__(self, prices: dict[str, Decimal | str | None], *, currency: str = "USD", timestamps: dict[str, str] | None = None) -> None: self.prices = {k: (Decimal(str(v)) if v is not None else None) for k, v in prices.items()} self.currency = currency.upper() self.timestamps = timestamps or {} def get_price(self, provider_symbol: str, *, price_date: str | None = None) -> EquityPriceQuote: if provider_symbol not in self.prices or self.prices[provider_symbol] is None: return EquityPriceQuote(provider_symbol=provider_symbol, currency=self.currency, close=None, quality_status="missing", error_message="price missing", provider=self.name) return EquityPriceQuote(provider_symbol=provider_symbol, currency=self.currency, close=self.prices[provider_symbol], price_timestamp=self.timestamps.get(provider_symbol), provider=self.name) class FmpEquityPriceProvider: name = "fmp" capability = ProviderCapability(True, True, True, "sequential_retry_after_backoff") def __init__(self, *, api_key: str | None = None, timeout_seconds: float = 8.0) -> None: self.api_key = api_key or _runtime_secret_value(("FMP_API_KEY", "FINANCIAL_MODELING_PREP_API_KEY", "JARVIS_FMP_API_KEY")) self.timeout_seconds = timeout_seconds def _get_json(self, path: str, params: dict[str, str]) -> object: if not self.api_key: raise RuntimeError("fmp_api_key_missing") path = path if path.startswith("/") else "/" + path url = "https://financialmodelingprep.com" + path + "?" + parse.urlencode({**params, "apikey": self.api_key}) req = request.Request(url, headers={"Accept": "application/json"}) @@ -373,186 +374,322 @@ class CompositeEquityPriceProvider: capability = getattr(provider, "capability", ProviderCapability(True, False, False, "unknown")) if price_date and price_date != "latest" and not capability.supports_historical_as_of: continue quote = provider.get_price(provider_symbol, price_date=price_date) if quote.close is not None and quote.quality_status == "fresh": return quote last_quote = quote return last_quote or EquityPriceQuote(provider_symbol=provider_symbol, currency="", close=None, provider=self.name, quality_status="missing", error_message="price_missing") def equity_price_provider_by_name(name: str) -> EquityPriceProvider: key = (name or "").lower().replace("-", "") providers = { "auto": CompositeEquityPriceProvider, "fmp": FmpEquityPriceProvider, "finnhub": FinnhubEquityPriceProvider, "twelvedata": TwelveDataEquityPriceProvider, "massive": MassiveEquityPriceProvider, "eodhd": EodhdEquityPriceProvider, "yfinance": YFinanceEquityPriceProvider, } if key not in providers: raise ValueError("unknown equity price provider") return providers[key]() def provider_capability(name: str) -> ProviderCapability: provider = equity_price_provider_by_name(name) return getattr(provider, "capability", ProviderCapability(True, False, False, "unknown")) def exchange_matches(expected: str | None, actual: str | None) -> bool: expected_key = (expected or "").strip().upper() actual_key = (actual or "").strip().upper() if not expected_key or not actual_key: return False aliases = { "NASDAQ": {"NASDAQ", "NMS", "NGM", "NCM"}, "NYSE": {"NYSE", "NYQ"}, "CBOE": {"CBOE", "BTS"}, "LSE": {"LSE", "LON"}, "FSX": {"FSX", "FRA"}, "SIX": {"SIX", "SWX", "XSWX"}, } return actual_key in aliases.get(expected_key, {expected_key}) def _resolve_market_alerts(conn: Connection, *, instrument_id: str) -> int: return resolve_instrument_alerts(conn, instrument_id=instrument_id, rule_ids=["missing_market_price", "stale_market_price"]) def _previous_price(conn: Connection, *, instrument_id: str, price_date: str, provider: str): return conn.execute( """ SELECT * FROM market_prices WHERE instrument_id=? AND provider=? AND price_date str: if close is None or close <= 0: return "not_checked" prev = _previous_price(conn, instrument_id=instrument_id, price_date=price_date, provider=provider) if not prev or not prev["close"]: return "not_checked" old = Decimal(str(prev["close"])) if old <= 0: return "not_checked" change = abs((close - old) / old) if change > Decimal("0.25"): create_alert(conn, priority="warnung", category="market_data", entity_type="instrument", entity_id=instrument_id, rule_id="corporate_action_suspected", message="Large local price move detected; corporate action review required before high-confidence valuation.", evidence={"previous_price_date": prev["price_date"], "price_date": price_date}, fingerprint="corporate_action_suspected") create_alert(conn, priority="warnung", category="market_data", entity_type="instrument", entity_id=instrument_id, rule_id="split_or_corporate_action_review_required", message="Split/corporate-action review required; no automatic split correction is applied.", evidence={"previous_price_date": prev["price_date"], "price_date": price_date}, fingerprint="split_or_corporate_action_review_required") update_instrument_metadata(conn, instrument_id=instrument_id, corporate_action_status="suspected", split_or_corporate_action_review_required=True, note="Automatic C2 heuristic detected >25% local price move; review required.") return "suspected" 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 _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: + 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=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, + provider_symbol or "", + price_type, + 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"], 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, *, 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, ) -> str: now = utc_now() fetched = fetched_at or now existing = conn.execute( "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( market_price_id, instrument_id, price_date, price_timestamp, close, adjusted_close, currency, provider, provider_symbol, quality_status, created_at, provider_market, error_message, corporate_action_status, fetched_at, price_type, run_id ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(instrument_id, price_date, provider) DO UPDATE SET price_timestamp=excluded.price_timestamp, close=excluded.close, adjusted_close=excluded.adjusted_close, currency=excluded.currency, provider_symbol=excluded.provider_symbol, quality_status=excluded.quality_status, created_at=excluded.created_at, provider_market=excluded.provider_market, error_message=excluded.error_message, corporate_action_status=excluded.corporate_action_status, fetched_at=excluded.fetched_at, price_type=excluded.price_type, run_id=excluded.run_id """, (market_price_id, instrument_id, price_date, price_timestamp, format(close, "f") if close is not None else "", format(adjusted_close, "f") if adjusted_close is not None else None, currency.upper(), provider, provider_symbol, quality_status, now, provider_market, error_message, corp_status, fetched, price_type, run_id), ) if existing and ( str(existing["close"] or "") != (format(close, "f") if close is not None else "") or str(existing["currency"] or "") != currency.upper() ): record_audit_event( conn, source="daily_market_fx_v1", action="market_price_provider_correction", entity_type="market_price", entity_id=str(existing["market_price_id"]), old_values={"close": existing["close"], "currency": existing["currency"], "quality_status": existing["quality_status"]}, new_values={"close": format(close, "f") if close is not None else None, "currency": currency.upper(), "quality_status": quality_status, "run_id": run_id}, confirmed=True, created_by="system", ) if quality_status in {"missing", "stale", "error", "conflict"} or close is None: rule = "stale_market_price" if quality_status == "stale" else "missing_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 def refresh_market_prices( conn: Connection, *, provider: EquityPriceProvider, asset_class: str, price_date: str | None = None, only_missing: bool = False, only_stale: bool = False, only_isin: str | None = None, limit: int | None = None, dry_run: bool = False, ) -> MarketPriceRefreshResult: asset = asset_class.lower() result = MarketPriceRefreshResult(asset_class=asset, dry_run=dry_run) query = """ SELECT m.*, i.asset_class, i.isin, i.instrument_status AS current_instrument_status, i.valuation_policy AS current_valuation_policy FROM instrument_price_mappings m JOIN instruments i ON i.instrument_id=m.instrument_id WHERE LOWER(i.asset_class)=? AND m.mapping_status='mapped' AND m.provider_symbol IS NOT NULL ORDER BY i.isin, m.provider_symbol """ mappings = conn.execute(query, (asset,)).fetchall() if only_isin: mappings = [m for m in mappings if (m["isin"] or "").upper() == only_isin.upper()] if limit is not None: mappings = mappings[:limit] result.total_mappings = len(mappings) effective_date = price_date or utc_now()[:10] for mapping in mappings: if mapping["current_instrument_status"] in {"delisted", "suspended", "merged", "inactive"} or mapping["current_valuation_policy"] == "exclude_from_auto_price_update": result.excluded_count += 1 result.skipped_count += 1 if not dry_run: diff --git a/src/jarvis_finance/services/asset_price_refresh.py b/src/jarvis_finance/services/asset_price_refresh.py index a822fed..1ecf740 100644 --- a/src/jarvis_finance/services/asset_price_refresh.py +++ b/src/jarvis_finance/services/asset_price_refresh.py @@ -1,378 +1,464 @@ 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 from typing import Any, Callable from jarvis_finance.api.schemas.market import QuoteRefreshRequest from jarvis_finance.audit.log import record_audit_event from jarvis_finance.market.providers import CoinGeckoClient, refresh_crypto_prices from jarvis_finance.services.market_service import refresh_equity_quotes_batch from jarvis_finance.services.modelled_wealth import build_modelled_wealth_development from jarvis_finance.storage.database import connect SOURCES = ("equity", "crypto", "fx") PROTECTED_TABLES = ( "accounts", "transactions", "crypto_holdings", "positions_snapshot", "postfinance_snapshot_positions", "truewealth_snapshot_positions", ) +@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, _column: str | None, _database: str | None, _trigger: str | None, ) -> int: if action in {SQLITE_INSERT, SQLITE_UPDATE, SQLITE_DELETE} and table in PROTECTED_TABLES: return SQLITE_DENY return SQLITE_OK def _now() -> str: return datetime.now(UTC).isoformat() def _database_path(conn: Connection) -> str: row = next((row for row in conn.execute("PRAGMA database_list") if str(row[1]) == "main"), None) if not row or not str(row[2] or ""): raise ValueError("asset_refresh_requires_persistent_database") return str(Path(str(row[2])).resolve()) def _protected_fingerprint(conn: Connection) -> str: payload: dict[str, list[dict[str, Any]]] = {} available = { str(row[0]) for row in conn.execute("SELECT name FROM sqlite_master WHERE type='table'").fetchall() } for table in PROTECTED_TABLES: if table not in available: continue rows = conn.execute(f'SELECT * FROM "{table}" ORDER BY rowid').fetchall() payload[table] = [dict(row) for row in rows] return hashlib.sha256( json.dumps(payload, sort_keys=True, separators=(",", ":"), default=str).encode("utf-8") ).hexdigest() def _status_payload(conn: Connection, job_id: str) -> dict[str, Any]: job = conn.execute("SELECT * FROM asset_price_refresh_jobs WHERE job_id=?", (job_id,)).fetchone() if not job: raise ValueError("asset_price_refresh_job_not_found") sources = [ dict(row) for row in conn.execute( "SELECT * FROM asset_price_refresh_sources WHERE job_id=? ORDER BY CASE source WHEN 'equity' THEN 1 WHEN 'crypto' THEN 2 ELSE 3 END", (job_id,), ).fetchall() ] return { "job_id": str(job["job_id"]), "status": str(job["status"]), "requested_at": str(job["requested_at"]), "completed_at": str(job["completed_at"]) if job["completed_at"] else None, "stale_before": str(job["stale_before"]), "progress": {"completed": int(job["progress_completed"]), "total": int(job["progress_total"])}, "sources": [ { "source": str(row["source"]), "status": str(row["status"]), "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, } for row in sources ], "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, } def create_asset_price_refresh_job(conn: Connection, *, stale_hours: int = 24) -> tuple[dict[str, Any], str]: """Persist a queued job only. No provider call occurs before the HTTP response.""" if conn.in_transaction: raise ValueError("asset_refresh_requires_clean_transaction") conn.execute("BEGIN IMMEDIATE") try: if conn.execute( "SELECT 1 FROM asset_price_refresh_jobs WHERE status IN ('queued','running') LIMIT 1" ).fetchone(): raise ValueError("asset_price_refresh_job_already_running") now = datetime.now(UTC) job_id = f"asset-refresh-{uuid.uuid4().hex}" stale_before = (now - timedelta(hours=max(1, min(stale_hours, 720)))).isoformat() conn.execute( """INSERT INTO asset_price_refresh_jobs( job_id,status,requested_at,stale_before,progress_total,progress_completed ) VALUES(?,'queued',?,?,3,0)""", (job_id, now.isoformat(), stale_before), ) conn.executemany( """INSERT INTO asset_price_refresh_sources( job_id,source,status,stale_candidates,updated_count ) VALUES(?,?,'pending',0,0)""", [(job_id, source) for source in SOURCES], ) conn.commit() except Exception: if conn.in_transaction: conn.rollback() raise 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( provider="auto", only_missing=True, stale_before=stale_before, limit=500, max_retries=1, pacing_seconds=0.15, ), ) 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"}) + return SourceRunResult( + stale_candidates=candidates, + updated_count=int(response.updated), + fresh_unchanged_count=int(response.cached), + stale_remaining_count=stale, + failed_count=max(failed, len(response.errors)), + diagnostics=tuple(sorted(set(response.errors)))[: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( """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 ) AND NOT EXISTS( SELECT 1 FROM crypto_prices p WHERE p.asset_id=a.asset_id AND p.fetched_at>=? AND p.quality_status='fresh' AND p.price IS NOT NULL ) ORDER BY a.asset_id""", (stale_before,), ).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) max_age_seconds = max(1, int((datetime.now(UTC) - cutoff.astimezone(UTC)).total_seconds())) result = refresh_crypto_prices( conn, provider=CoinGeckoClient(), currency="CHF", max_age_seconds=max_age_seconds, 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( """SELECT DISTINCT upper(i.currency) currency FROM instruments i WHERE i.is_active=1 AND upper(COALESCE(i.currency,'CHF'))!='CHF' AND NOT EXISTS( SELECT 1 FROM fx_rates f WHERE f.base_currency=upper(i.currency) AND f.quote_currency='CHF' AND f.rate_date>=? AND f.quality_status IN ('fresh','ok') ) ORDER BY currency""", (cutoff_date,), ).fetchall() ] updated = 0 failures = 0 for currency in currencies: try: result = resolve_fx_rate_to_chf( conn, base_currency=currency, rate_date=None, providers=[FrankfurterFxProvider(), TwelveDataFxProvider()], persist=True, resolve_fixed=True, ) updated += int(result.status == "ok") 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)), + 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, } 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) selected = runners or DEFAULT_RUNNERS try: conn.execute("BEGIN IMMEDIATE") claimed = conn.execute( """UPDATE asset_price_refresh_jobs SET status='running' WHERE job_id=? AND status='queued'""", (job_id,), ).rowcount conn.commit() if claimed != 1: return stale_before = str( conn.execute( "SELECT stale_before FROM asset_price_refresh_jobs WHERE job_id=?", (job_id,), ).fetchone()[0] ) protected_before = _protected_fingerprint(conn) 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( "UPDATE asset_price_refresh_sources SET status='running',started_at=? WHERE job_id=? AND source=?", (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" 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=?", (completed, job_id), ) conn.commit() 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() ] conn.execute( """INSERT INTO aggregated_wealth_refresh_snapshots( wealth_snapshot_id,job_id,captured_at,known_wealth_chf,quality_status,source_status_json ) VALUES(?,?,?,?,?,?)""", ( wealth_snapshot_id, job_id, _now(), current.get("value_chf"), "complete" if failures == 0 else "partial", 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", action="asset_prices_refresh_completed", entity_type="asset_price_refresh_job", entity_id=job_id, old_values={}, new_values={ "status": final_status, "source_count": len(SOURCES), "failed_source_count": failures, "wealth_snapshot_created": bool(wealth_snapshot_id), "holdings_mutated": False, "transactions_mutated": False, "trades_created": 0, }, created_by="system", ) conn.execute( """UPDATE asset_price_refresh_jobs SET status=?,completed_at=?,wealth_snapshot_id=?,audit_id=? WHERE job_id=?""", (final_status, _now(), wealth_snapshot_id, audit_id, job_id), ) conn.commit() except Exception as exc: if conn.in_transaction: conn.rollback() audit_id = record_audit_event( conn, source="asset_price_refresh_job_v1", action="asset_prices_refresh_failed", entity_type="asset_price_refresh_job", entity_id=job_id, old_values={}, new_values={"status": "failed", "error_code": str(exc)[:120]}, created_by="system", ) conn.execute( "UPDATE asset_price_refresh_jobs SET status='failed',completed_at=?,audit_id=? WHERE job_id=?", (_now(), audit_id, job_id), ) conn.commit() finally: conn.close() def asset_price_refresh_status(conn: Connection, job_id: str) -> dict[str, Any]: """Stored status only: no provider call, write or lazy refresh.""" return _status_payload(conn, job_id) diff --git a/src/jarvis_finance/services/market_service.py b/src/jarvis_finance/services/market_service.py index e479e34..f6ed215 100644 --- a/src/jarvis_finance/services/market_service.py +++ b/src/jarvis_finance/services/market_service.py @@ -129,161 +129,166 @@ def _business_day_age(earlier: date, later: date) -> int: def _has_fresh_price_for_target( conn: Connection, instrument_id: str, target: date, *, stale_before: str | None = None, ) -> bool: mapping = conn.execute( """SELECT provider,provider_symbol,provider_market,upper(COALESCE(trading_currency,currency,'')) currency FROM instrument_price_mappings WHERE instrument_id=? AND mapping_status='mapped' AND provider_symbol IS NOT NULL ORDER BY CASE provider WHEN 'fmp' THEN 1 ELSE 2 END,updated_at DESC LIMIT 1""", (instrument_id,), ).fetchone() if not mapping: return False rows = conn.execute( """SELECT price_date,currency,provider,provider_symbol,provider_market, COALESCE(fetched_at,created_at,price_timestamp,price_date) freshness_at FROM market_prices WHERE instrument_id=? AND price_date<=? AND close IS NOT NULL AND close!='' AND quality_status='fresh' AND error_message IS NULL ORDER BY price_date DESC,COALESCE(fetched_at,created_at) DESC""", (instrument_id, target.isoformat()), ).fetchall() for row in rows: if stale_before and str(row["freshness_at"] or "") < stale_before: continue actual_date = date.fromisoformat(str(row["price_date"])[:10]) if not 0 <= _business_day_age(actual_date, target) <= 2: continue if str(row["provider_symbol"] or "").upper() != str(mapping["provider_symbol"] or "").upper(): continue if str(row["currency"] or "").upper() != str(mapping["currency"] or "").upper(): continue if not exchange_matches(mapping["provider_market"], row["provider_market"]): continue if str(row["provider"] or "").lower() != str(mapping["provider"] or "").lower(): if str(row["provider"] or "").lower() != "yfinance": continue return True return False def _quote_date(quote: EquityPriceQuote, fallback: date) -> date: if quote.price_timestamp: try: return datetime.fromisoformat(quote.price_timestamp.replace("Z", "+00:00")).date() except ValueError: try: return date.fromisoformat(quote.price_timestamp[:10]) except ValueError: pass return fallback def _validate_historical_quote(quote: EquityPriceQuote, *, mapping: Any, target: date, requested_provider: str) -> EquityPriceQuote: if quote.close is None or quote.close <= 0: return replace(quote, quality_status=_quality_from_error(quote.error_message, quote.quality_status)) quote_date = _quote_date(quote, target) if quote_date > target: return replace(quote, close=None, quality_status="future_price_rejected", error_message="future_price_rejected") if _business_day_age(quote_date, target) > 2: return replace(quote, close=None, quality_status="stale", error_message="historical_price_too_old") expected_currency = str(mapping["trading_currency"] or mapping["currency"] or "").upper() actual_currency = str(quote.currency or "").upper() is_fallback = requested_provider == "auto" and str(quote.provider or "").lower() == "yfinance" if is_fallback: capability = provider_capability(str(quote.provider)) if not capability.supports_historical_as_of: return replace(quote, close=None, quality_status="provider_not_historical", error_message="provider_not_historical") if not actual_currency or (expected_currency and actual_currency != expected_currency): return replace(quote, close=None, quality_status="currency_mismatch", error_message="currency_mismatch") if str(quote.provider_symbol or "").upper() != str(mapping["provider_symbol"] or "").upper(): return replace(quote, close=None, quality_status="symbol_mismatch", error_message="symbol_mismatch") if not exchange_matches(str(mapping["provider_market"] or ""), quote.provider_market): 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: inst, mapping, warnings = _instrument_mapping(conn, instrument_id) if warnings: return MarketQuoteResponse(provider_symbol=None, currency=inst["currency"] if inst else None, quality_status="missing", warnings=warnings) mapping_data = dict(mapping) if mapping else { "provider": inst["data_provider_primary"] or "auto", "provider_symbol": inst["provider_symbol"], "provider_market": inst["exchange"], "trading_currency": inst["currency"], "currency": inst["currency"], } provider_symbol = str(mapping_data["provider_symbol"]) provider_name = (req.provider or "auto").lower() target = _effective_market_date(req.price_date) quote = equity_price_provider_by_name(provider_name).get_price(provider_symbol, price_date=target.isoformat()) 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() if not req.dry_run and quote.close is not None and quality == "fresh": store_market_price( conn, instrument_id=instrument_id, price_date=ts[:10], close=quote.close, currency=quote.currency or inst["currency"] or "CHF", provider=quote.provider, provider_symbol=quote.provider_symbol or provider_symbol, provider_market=quote.provider_market or mapping_data["provider_market"], price_timestamp=ts, adjusted_close=quote.adjusted_close, quality_status=quality, error_message=None, ) upsert_equity_price_point( conn, instrument_id=instrument_id, timestamp=ts, price=quote.close, currency=quote.currency or inst["currency"] or "CHF", provider=quote.provider, provider_symbol=quote.provider_symbol or provider_symbol, interval=req.interval, source_quality=quality, ) conn.commit() return MarketQuoteResponse( latest_price=_decimal_text(quote.close), currency=quote.currency or inst["currency"], close=_decimal_text(quote.close), provider=quote.provider, provider_symbol=quote.provider_symbol or provider_symbol, fetched_at=ts, quality_status=quality, warnings=warnings + ([quote.error_message] if quote.error_message else []), ) def _persistent_database_path(conn: Connection) -> str | None: row = next((row for row in conn.execute("PRAGMA database_list") if str(row[1]) == "main"), None) return str(row[2]) if row and str(row[2] or "") else None def _parallel_equity_worker( db_path: str, instrument_id: str, req: QuoteRefreshRequest ) -> tuple[MarketQuoteResponse, int]: worker = connect(db_path) worker.execute("PRAGMA busy_timeout=10000") try: attempt = 0 quote: MarketQuoteResponse | None = None while attempt <= req.max_retries: attempt += 1 quote = refresh_equity_quote(worker, instrument_id, req) 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: break time.sleep(min(2 ** (attempt - 1), 4)) __HERMES_CWD_8d46a20096ed__/home/agent/.hermes/worktrees/FinanceManager-sprint23.1__HERMES_CWD_8d46a20096ed__