# Writes are idempotent against that unique constraint: an estimate already stored for the # same (period_end, metric, as_of_date) is updated in place, never duplicated on re-ingestion. estimates_written = _ingest_estimates(db, adapter, security.id, ticker, source) prices_written = _ingest_prices(db, adapter, security.id, ticker, source) db.commit() logger.info("ingest.security.complete", ticker=ticker, provider=adapter.name, is_demo=is_demo, periods=len(income_periods), estimates=estimates_written, prices=prices_written, secondary_provider=(secondary_adapter.name if secondary_adapter else None)) return security.id def _ingest_estimates(db: Session, adapter: ProviderAdapter, security_id: str, ticker: str, source: Source) -> int: """Fetch and persist forward analyst estimates. Returns how many rows were written or updated. Never raises for an adapter that has not implemented estimates (EODHD's is `NotImplementedError`) — that degrades to "no estimates this run", which is exactly what `select_forward_estimate()` is built to report honestly, rather than crashing an ingestion that otherwise succeeded. """ try: rows = adapter.get_estimates(ticker) except NotImplementedError: logger.warning("ingest.estimates.not_implemented", ticker=ticker, provider=adapter.name) return 0 written = 0 for r in rows: if r.consensus_value is None or r.period_end is None: continue existing = db.query(Estimate).filter_by( security_id=security_id, period_end=r.period_end, metric=r.metric, as_of_date=r.as_of_date, ).one_or_none() if existing is None: existing = Estimate( security_id=security_id, period_end=r.period_end, metric=r.metric, as_of_date=r.as_of_date, ) db.add(existing) existing.consensus_value = r.consensus_value existing.num_analysts = r.num_analysts existing.source_id = source.id written += 1 return written # AUDIT FIX (StockLab overhaul, Part 23): this task previously had NO retry policy at all — a # transient rate-limit (`ProviderRateLimitError`, e.g. HTTP 429) or network blip crashed the task # outright with no retry, despite `Settings.PROVIDER_MAX_RETRIES` existing in config and looking # like it was already wired in (it was not referenced anywhere in this file before this fix). # `ProviderAuthError` (bad/expired API key) and `ProviderNotFoundError` (unknown ticker) are # deliberately NOT in `autoretry_for` — retrying a bad API key or a ticker that doesn't exist just # wastes the retry budget on an error that will never resolve itself. @celery_app.task( name="app.workers.ingest.ingest_security_task", bind=True, autoretry_for=(ProviderRateLimitError, ConnectionError, TimeoutError), retry_backoff=True, # exponential backoff between attempts retry_backoff_max=300, # cap backoff at 5 minutes retry_jitter=True, # avoid a thundering herd of simultaneous retries max_retries=get_settings().PROVIDER_MAX_RETRIES, ) def ingest_security_task(self, ticker: str, provider_name: str) -> str: from app.adapters import build_adapter, get_secondary_adapter settings = get_settings() adapter = build_adapter(provider_name, settings) # AUDIT FIX (StockLab overhaul, final engineering pass, Part A6): dual-source conflict