o
    ¸ªjh6  ã                   @   sB  d Z ddlZddlZddlZddlZddlmZmZ ddlmZ ddl	m
Z
 ddlmZ ddlmZ edƒZeeƒ ¡ jZd	d
„ Zdd„ Zdd„ Ze
defdd„ƒZdefdd„Zdefdd„Zdedefdd„Zdedefdd„ZdYdefdd „Zdefd!d"„Z d#efd$d%„Z!d&efd'd(„Z"dedefd)d*„Z#dedefd+d,„Z$d-efd.d/„Z%defd0d1„Z&defd2d3„Z'd4efd5d6„Z(d&efd7d8„Z)d9ed&efd:d;„Z*d&efd<d=„Z+d>efd?d@„Z,d&efdAdB„Z-dCefdDdE„Z.d&efdFdG„Z/dHefdIdJ„Z0d&efdKdL„Z1dMefdNdO„Z2d&efdPdQ„Z3dRefdSdT„Z4dZdVdW„Z5G dXdU„ dUƒZ6dS )[aÊ  
data_store.py

All persistence lives here: premarket snapshots, intraday snapshots,
trade records, and bot state (positions, watchlists, cooldowns).

Design notes:
- Uses simple JSON-lines files per day for premarket/intraday data so
  they can be analyzed later with pandas without a database dependency.
- State files (positions.json, watchlist.json, cooldowns.json) are
  written atomically (write to tmp file, then os.replace) to avoid
  corruption if the process is killed mid-write.
- A single filelock-style guard (using a lock file + fcntl) protects
  state writes in case of accidental concurrent processes, matching
  a lesson learned in prior systems where duplicate processes raced
  on shared state.
é    N)ÚdatetimeÚdate)ÚPath)Úcontextmanager)Ú
get_config)Ú
get_loggerÚ
data_storec                  C   sL   t ƒ d } t| d  t| d  t| d  t| d  t| d  t| d  dœS )	NÚdata_storageÚpremarket_dirÚintraday_dirÚ
trades_dirÚfast_engine_dirÚ	state_dirÚbreakout_scanner_dir)Ú	premarketÚintradayÚtradesÚfast_engineÚstateÚbreakout_scanner)r   ÚBASE_DIR)Úcfg© r   údata_store.pyÚ_dirs#   s   






úr   c                  C   s"   t ƒ  ¡ D ]	} | jddd� qd S )NT©ÚparentsÚexist_ok)r   ÚvaluesÚmkdir)Údr   r   r   Úensure_dirs/   s   ÿr!   c                   C   s   t  ¡  ¡ S ©N)r   ÚtodayÚ	isoformatr   r   r   r   Ú
_today_str4   s   r%   Ú	lock_pathc                 c   sz   � | j jddd� t| dƒ�%}t |tj¡ zd V  W t |tj¡ nt |tj¡ w W d   ƒ d S 1 s6w   Y  d S )NTr   Úw)Úparentr   ÚopenÚfcntlÚflockÚLOCK_EXÚLOCK_UN)r&   Úlfr   r   r   Ú
_file_lock8   s   €""ûr/   Úpathc                 C   sl   | j jddd� |  | jd ¡}t|dƒ�}tj||dtd� W d   ƒ n1 s)w   Y  t 	|| ¡ d S )NTr   z.tmpr'   é   )ÚindentÚdefault)
r(   r   Úwith_suffixÚsuffixr)   ÚjsonÚdumpÚstrÚosÚreplace)r0   ÚdataÚtmpÚfr   r   r   Ú_atomic_write_jsonC   s   ÿr>   c              
   C   s�   |   ¡ s|S zt| dƒ�}t |¡W  d   ƒ W S 1 sw   Y  W d S  tjtfyG } zt d| › d|› d�¡ |W  Y d }~S d }~ww )NÚrzFailed to read z: z. Returning default.)Úexistsr)   r6   ÚloadÚJSONDecodeErrorÚOSErrorÚlogÚwarning)r0   r3   r=   Úer   r   r   Ú
_read_jsonK   s   (ÿ€þrG   Úrecordc                 C   sX   | j jddd� t| dƒ�}| tj|td�d ¡ W d   ƒ d S 1 s%w   Y  d S )NTr   Úa)r3   Ú
)r(   r   r)   Úwriter6   Údumpsr8   )r0   rH   r=   r   r   r   Ú_append_jsonlV   s   "ÿrM   Úsymbolc                 C   ó:   t ƒ d tƒ › d� }t ¡  ¡ | dœ|¥}t||ƒ d S )Nr   ú.jsonl©Ú	timestamprN   ©r   r%   r   Úutcnowr$   rM   ©rN   rH   r0   r   r   r   Úappend_premarket_snapshot`   ó   rV   úcandidates_20.jsonÚ
candidatesc                 C   s(   t ƒ d tƒ › d|› � }t|| ƒ d S )Nr   Ú_©r   r%   r>   )rY   Úfilenamer0   r   r   r   Úwrite_premarket_candidatesf   s   r]   c                 C   s$   t ƒ d tƒ › d� }t|| ƒ dS )a¢  [FEATURE 2026-09-16] breakout_scanner.py's own ranked candidate
    list -- deliberately a SEPARATE file from write_premarket_candidates
    above (different directory, different name) so it can never be
    confused with or overwrite sip_bot's own live candidate output.
    Date-stamped filename, same convention as write_premarket_candidates,
    so a day's breakout_scanner run and sip_bot run can be diffed later.r   z_scanner.jsonNr[   )rY   r0   r   r   r   Ú!write_breakout_scanner_candidatesk   s   r^   Útop_listc                 C   s.   t ƒ d d }t ¡  ¡ | dœ}t||ƒ dS )zDWrites the canonical top_stocks.json referenced by the whole system.r   útop_stocks.json)Ú
updated_atÚstocksN)r   r   rT   r$   r>   )r_   r0   Úpayloadr   r   r   Úwrite_top_stocksv   s   rd   Úreturnc                  C   s(   t ƒ d d } t| dg iƒ}| dg ¡S )Nr   r`   rb   )r   rG   Úget)r0   r;   r   r   r   Úread_top_stocks}   s   rg   c                 C   rO   )Nr   rP   rQ   rS   rU   r   r   r   Úappend_intraday_snapshot‡   rW   rh   c                 C   sr   t ƒ d d }tt ƒ d d ƒ� t|i ƒ}dt ¡  ¡ i|¥|| < t||ƒ W d  ƒ dS 1 s2w   Y  dS )z×Mirrors the latest intraday state for a symbol into intraday.json
    (a single mutable snapshot file, distinct from the historical .jsonl log),
    matching the project's requested intraday.json state file pattern.r   zintraday.jsonzintraday.json.lockra   N)r   r/   rG   r   rT   r$   r>   )rN   rH   r0   Úall_datar   r   r   Úwrite_intraday_json�   s   
"ýrj   Útradec                 C   ó$   t ƒ d tƒ › d� }t|| ƒ d S )Nr   ú_trades.jsonl)r   r%   rM   )rk   r0   r   r   r   Úappend_trade_recordœ   ó   rn   c                 C   s8   t ƒ d tƒ › d� }dt ¡  ¡ i| ¥} t|| ƒ dS )a^  
    [FEATURE 2026-09-11, extended 2026-09-11 v3] One JSON line per
    ranked candidate, every poll cycle, for EVERY health-eligible
    premarket_20 symbol -- not just the top-5 shortlist that gets a real
    fast_entry_gate decision. `shortlisted`/`confidence_rank` distinguish
    "ranked but not evaluated this cycle" (`shortlisted: False`,
    `state`/`should_enter`: None) from "made the top-5, got a real
    BUY/WAIT/REJECT" (`shortlisted: True`). Per explicit instruction
    ("log all bot decisions so we can review and adjust what is not
    working", later extended to "make sure we have all the indicator
    values logged during the life cycle of the symbol"): the whole point
    is being able to load this file into pandas afterward and see
    exactly what the engine saw for every symbol at every poll, whether
    or not it ever became a real trade (which append_trade_record above
    already covers) or even made the shortlist. Each record also carries
    a full raw-indicator snapshot (fast_pipeline.snapshot_dict() --
    every stream_features.StreamFeatures/FastPredictionReading field,
    not just the classified labels) alongside the classified summary
    fields kept here for quick filtering without unpacking the nested
    snapshot. Same per-day .jsonl pattern as premarket/intraday
    snapshots -- pandas-readable, no database dependency. Volume note:
    at ~20-30 eligible symbols x poll_interval_seconds_intraday=5s, this
    is meaningfully larger than the old shortlisted-only version (was
    ~5 records/cycle, now ~20-30/cycle) -- expect a materially bigger
    file per trading day.
    r   rP   rR   NrS   ©rH   r0   r   r   r   Úappend_fast_engine_decision¥   s   rq   c                 C   s:   t ƒ d dtƒ › d� }dt ¡  ¡ i| ¥} t|| ƒ dS )a®  
    [FEATURE 2026-09-11] One JSON line per open-position poll where
    position_manager._check_momentum_fade_exit() successfully computed
    features -- added the same day that check was introduced, after
    discovering there was no way to validate its thresholds
    (velocity_stall_threshold_pct_per_sec / pressure_threshold /
    min_ticks_sub, see config.json's stop.momentum_fade_exit) against
    real data: a full replay of 2026-09-11's 16 trades against the new
    check turned out to be impossible, because append_fast_engine_
    decision() above only dumps CANDIDATE-scan snapshots, which stop the
    moment a symbol becomes an open position (fast_pipeline.py has no
    reason to keep ranking something already held) -- so no tick-level
    velocity_sub/pressure.score history exists for any open-position hold
    window from that day. This closes that gap going forward: every poll
    where the check actually evaluates a reading (not the early-return
    guard cases -- disabled, missing bars_sub/quote, insufficient_data --
    which aren't informative for tuning) gets one row here, whether or
    not it fired, so the thresholds can be tuned against real fills
    instead of Yahoo-1-minute-bar approximations next time. Separate file
    from fast_engine's per-day dump -- different schema/purpose (a single
    open position's fade check vs. ranking all eligible candidates) even
    though it lives in the same directory.
    r   Úmomentum_fade_rP   rR   NrS   rp   r   r   r   Úappend_momentum_fade_snapshotÅ   s   rs   r   c                 C   rl   )Nr   z_summary.jsonr[   )r   r0   r   r   r   Úwrite_trades_summaryâ   ro   rt   c               
   C   s¦   t ƒ d tƒ › d� } |  ¡ sg S g }t| dƒ�1}|D ]%}| ¡ }|s$qz
| t |¡¡ W q tjy@   t	 
d| › �¡ Y qw W d  ƒ |S 1 sLw   Y  |S )uÜ  
    Reads every trade record appended today via append_trade_record() â€”
    the complete, per-trade history for the day.

    [BUGFIX 2026-08-17] This is deliberately separate from
    PositionManager.positions, which is keyed by symbol and therefore
    only retains the MOST RECENT trade for any symbol traded more than
    once in a day. monitor.py's end-of-day summary previously built
    itself from positions.values() and silently undercounted days with
    repeated same-symbol round trips (verified: a session with 19 actual
    closed trades logged a summary of only 10 â€” one per unique symbol).
    Use this function instead for anything that needs a true count of
    the day's trades or an accurate total P/L.
    r   rm   r?   z(Skipping malformed trade record line in N)r   r%   r@   r)   ÚstripÚappendr6   ÚloadsrB   rD   rE   )r0   r   r=   Úliner   r   r   Úload_today_tradesç   s(   þú
ÿ
ö
ry   Únamec                 C   s   t ƒ d |  S )Nr   )r   )rz   r   r   r   Ú_state_path  ó   r{   c                  C   s`   t tdƒi ƒ} tƒ }i }|  ¡ D ]\}}| d¡dv r |||< q| dd¡ |¡r-|||< q|S )Núpositions.jsonÚstatus)r)   ÚclosingÚ
entry_timeÚ )rG   r{   r%   Úitemsrf   Ú
startswith)Úrawr#   ÚfilteredÚsymÚpr   r   r   Úload_positions  s   
€rˆ   Ú	positionsc                 C   ó@   t tdƒƒ� ttdƒ| ƒ W d   ƒ d S 1 sw   Y  d S )Nzpositions.json.lockr}   ©r/   r{   r>   )r‰   r   r   r   Úsave_positions  ó   "ÿrŒ   c                   C   s   t tdƒg g dœƒS )Núwatchlist.json)Úpremarket_20Úfinal_10©rG   r{   r   r   r   r   Úload_watchlist!  s   r’   Ú	watchlistc                 C   ó   t tdƒ| ƒ d S )NrŽ   ©r>   r{   )r“   r   r   r   Úsave_watchlist%  ó   r–   c                   C   ó   t tdƒi ƒS ©Nzcooldowns.jsonr‘   r   r   r   r   Úload_cooldowns)  r|   rš   Ú	cooldownsc                 C   r”   r™   r•   )r›   r   r   r   Úsave_cooldowns-  r—   rœ   c                   C   r˜   )Núintraday_health.jsonr‘   r   r   r   r   Úload_health_state6  r|   rž   Úhealth_statec                 C   rŠ   )Nzintraday_health.json.lockr�   r‹   )rŸ   r   r   r   Úsave_health_state:  r�   r    c                   C   r˜   )Núsession_state.jsonr‘   r   r   r   r   Úload_session_state?  r|   r¢   r   c                 C   rŠ   )Nzsession_state.json.lockr¡   r‹   )r   r   r   r   Úsave_session_stateC  r�   r£   ÚSingletonLockc                   C   s   t tdƒƒS )u�   Ensures only one monitor.py instance runs at a time â€” prevents the
    duplicate-process-racing-the-same-account failure class.zmonitor.pid.lock)r¤   r{   r   r   r   r   Úacquire_singleton_lockH  s   r¥   c                   @   s*   e Zd Zdefdd„Zdd„ Zdd„ ZdS )	r¤   r0   c                 C   s   || _ d | _d S r"   )r0   Ú_fh)Úselfr0   r   r   r   Ú__init__O  s   
zSingletonLock.__init__c                 C   st   | j jjddd� t| j dƒ| _zt | jtjtjB ¡ W n t	y(   t
dƒ‚w | j tt ¡ ƒ¡ | j ¡  | S )NTr   r'   z{Another monitor.py instance already holds the singleton lock. Refusing to start a second instance against the same account.)r0   r(   r   r)   r¦   r*   r+   r,   ÚLOCK_NBÚBlockingIOErrorÚRuntimeErrorrK   r8   r9   ÚgetpidÚflush)r§   r   r   r   Ú	__enter__S  s   ÿÿ
zSingletonLock.__enter__c                 C   s(   | j rt | j tj¡ | j  ¡  d S d S r"   )r¦   r*   r+   r-   Úclose)r§   Úexc_typeÚexc_valÚexc_tbr   r   r   Ú__exit__a  s   þzSingletonLock.__exit__N)Ú__name__Ú
__module__Ú__qualname__r   r¨   r®   r³   r   r   r   r   r¤   N  s    )rX   )re   r¤   )7Ú__doc__r6   r9   r*   Útimer   r   Úpathlibr   Ú
contextlibr   Úconfig_loaderr   Úlogger_setupr   rD   Ú__file__Úresolver(   r   r   r!   r%   r/   r>   rG   ÚdictrM   r8   rV   Úlistr]   r^   rd   rg   rh   rj   rn   rq   rs   rt   ry   r{   rˆ   rŒ   r’   r–   rš   rœ   rž   r    r¢   r£   r¥   r¤   r   r   r   r   Ú<module>   sZ    


	 $	
