78 lines
3.8 KiB
Python
78 lines
3.8 KiB
Python
|
|
"""Replay the persisted real SPD sample in an isolated database, with no network.
|
||
|
|
|
||
|
|
The one-sector publication is exclusively a local integration fixture; its ranks
|
||
|
|
must never be interpreted as a market-wide ranking.
|
||
|
|
"""
|
||
|
|
import os
|
||
|
|
import json
|
||
|
|
from datetime import date
|
||
|
|
from pathlib import Path
|
||
|
|
from urllib.parse import urlsplit, urlunsplit
|
||
|
|
import psycopg
|
||
|
|
from alembic import command
|
||
|
|
from alembic.config import Config
|
||
|
|
from zhixing_server.bootstrap.config import Settings, get_settings
|
||
|
|
from zhixing_server.modules.sector_radar.application.build import BuildSectorRadar, BuildSectorRadarCommand
|
||
|
|
from zhixing_server.modules.sector_radar.domain.models import SectorType
|
||
|
|
from zhixing_server.modules.sector_radar.domain.source import (
|
||
|
|
SourceSnapshot, SourceResult, TradeCalendarRow, SectorIndexRow, SectorMemberRow,
|
||
|
|
StockBasicRow, SuspendRow, DailyRow, MoneyflowDcRow, MoneyflowRow, build_source_snapshot,
|
||
|
|
)
|
||
|
|
from zhixing_server.modules.sector_radar.infrastructure.postgres import PostgresSectorRadarRepository
|
||
|
|
|
||
|
|
|
||
|
|
def database_url():
|
||
|
|
"""Resolve only the named localhost fixture database without printing secrets."""
|
||
|
|
s=Settings(_env_file='../.env'); u=urlsplit(s.database_url)
|
||
|
|
return urlunsplit((u.scheme,u.netloc.rsplit('@',1)[0]+'@127.0.0.1:5433','/radar_detail_selftest_0906',u.query,''))
|
||
|
|
|
||
|
|
|
||
|
|
class StoredSource:
|
||
|
|
"""Implement provider port using already committed raw snapshots only."""
|
||
|
|
def __init__(self, dsn):
|
||
|
|
self.by_api={}
|
||
|
|
with psycopg.connect(dsn) as c:
|
||
|
|
for r in c.execute('SELECT id,api_name,normalized_params,target_trade_date,partition_key,observed_at,payload,row_count,returned_fields,content_sha256,row_limit,limit_reached FROM sector_radar_source_snapshot').fetchall():
|
||
|
|
snap=SourceSnapshot(r[0],r[1],tuple(sorted(r[2].items())),r[3],r[4],r[5],tuple(r[6]),r[7],tuple(r[8]),r[9],r[10],r[11])
|
||
|
|
self.by_api.setdefault(r[1],[]).append(snap)
|
||
|
|
|
||
|
|
def result(self, api, parser):
|
||
|
|
snaps=tuple(self.by_api[api])
|
||
|
|
return SourceResult(snaps,tuple(parser(row) for s in snaps for row in s.rows))
|
||
|
|
|
||
|
|
def fetch_trade_calendar(self,start,end):
|
||
|
|
return self.result('trade_cal',TradeCalendarRow.from_mapping)
|
||
|
|
def fetch_sector_indices(self,target,kind):
|
||
|
|
if kind is SectorType.INDUSTRY:
|
||
|
|
# Explicitly empty industry scope for this one-concept test fixture.
|
||
|
|
snap=build_source_snapshot(api_name='dc_index',params={'test_scope':'empty_industry'},rows=(),target_trade_date=target)
|
||
|
|
return SourceResult((snap,),())
|
||
|
|
return self.result('dc_index',lambda row:SectorIndexRow.from_mapping(row,kind))
|
||
|
|
def fetch_sector_members(self,target,codes):
|
||
|
|
return self.result('dc_member',SectorMemberRow.from_mapping)
|
||
|
|
def fetch_stock_basics(self):
|
||
|
|
return self.result('stock_basic',StockBasicRow.from_mapping)
|
||
|
|
def fetch_suspensions(self,target):
|
||
|
|
return self.result('suspend_d',SuspendRow.from_mapping)
|
||
|
|
def fetch_daily(self,target):
|
||
|
|
return self.result('daily',DailyRow.from_mapping)
|
||
|
|
def fetch_moneyflow_dc(self,target,codes):
|
||
|
|
return self.result('moneyflow_dc',MoneyflowDcRow.from_mapping)
|
||
|
|
def fetch_moneyflow(self,target):
|
||
|
|
return self.result('moneyflow',MoneyflowRow.from_mapping)
|
||
|
|
|
||
|
|
|
||
|
|
def main():
|
||
|
|
"""Run actual build and HTTP reads, checking independently recomputed facts."""
|
||
|
|
dsn=database_url()
|
||
|
|
os.environ['ZHIXING_DATABASE_URL']=dsn
|
||
|
|
get_settings.cache_clear()
|
||
|
|
command.upgrade(Config('alembic.ini'),'head')
|
||
|
|
repository=PostgresSectorRadarRepository(dsn,max_connections=2)
|
||
|
|
result=BuildSectorRadar(StoredSource(dsn),repository,today=date(2026,9,6)).execute(BuildSectorRadarCommand(trade_date=date(2026,9,4)))
|
||
|
|
print(json.dumps(result.as_dict(),ensure_ascii=False))
|
||
|
|
repository.close()
|
||
|
|
|
||
|
|
if __name__=='__main__':
|
||
|
|
main()
|