o
    8@–jý$  ã                   @   s  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dSd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d/d0„Z%d$efd1d2„Z&d3ed$efd4d5„Z'd$efd6d7„Z(d8efd9d:„Z)d$efd;d<„Z*d=efd>d?„Z+d$efd@dA„Z,dBefdCdD„Z-d$efdEdF„Z.dGefdHdI„Z/d$efdJdK„Z0dLefdMdN„Z1dTdPdQ„Z2G dRdO„ dOƒZ3dS )UaÊ  
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   s8   t ƒ d } t| d  t| d  t| d  t| d  dœS )NÚdata_storageÚpremarket_dirÚintraday_dirÚ
trades_dirÚ	state_dir)Ú	premarketÚintradayÚtradesÚstate)r   ÚBASE_DIR)Úcfg© r   úF/var/www/screener/trade/premarket_backup_2026-09-08_2010/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_str2   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_lock6   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)r,   ÚdataÚtmpÚfr   r   r   Ú_atomic_write_jsonA   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%   r2   ÚloadÚJSONDecodeErrorÚOSErrorÚlogÚwarning)r,   r/   r9   Úer   r   r   Ú
_read_jsonI   s   (ÿ€þrC   Ú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)r/   Ú
)r$   r   r%   Úwriter2   Údumpsr4   )r,   rD   r9   r   r   r   Ú_append_jsonlT   s   "ÿrI   Úsymbolc                 C   ó:   t ƒ d tƒ › d� }t ¡  ¡ | dœ|¥}t||ƒ d S )Nr   ú.jsonl©Ú	timestamprJ   ©r   r!   r   Úutcnowr    rI   ©rJ   rD   r,   r   r   r   Úappend_premarket_snapshot^   ó   rR   úcandidates_20.jsonÚ
candidatesc                 C   s(   t ƒ d tƒ › d|› � }t|| ƒ d S )Nr   Ú_©r   r!   r:   )rU   Úfilenamer,   r   r   r   Úwrite_premarket_candidatesd   s   rY   Ú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   rP   r    r:   )rZ   r,   Úpayloadr   r   r   Úwrite_top_stocksi   s   r_   Úreturnc                  C   s(   t ƒ d d } t| dg iƒ}| dg ¡S )Nr   r[   r]   )r   rC   Úget)r,   r7   r   r   r   Úread_top_stocksp   s   rb   c                 C   rK   )Nr   rL   rM   rO   rQ   r   r   r   Úappend_intraday_snapshotz   rS   rc   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.lockr\   N)r   r+   rC   r   rP   r    r:   )rJ   rD   r,   Úall_datar   r   r   Úwrite_intraday_json€   s   
"ýre   Útradec                 C   ó$   t ƒ d tƒ › d� }t|| ƒ d S )Nr   ú_trades.jsonl)r   r!   rI   )rf   r,   r   r   r   Úappend_trade_record�   ó   ri   r   c                 C   rg   )Nr   z_summary.jsonrW   )r   r,   r   r   r   Úwrite_trades_summary”   rj   rk   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   rh   r;   z(Skipping malformed trade record line in N)r   r!   r<   r%   ÚstripÚappendr2   Úloadsr>   r@   rA   )r,   r   r9   Úliner   r   r   Úload_today_trades™   s(   þú
ÿ
ö
rp   Únamec                 C   s   t ƒ d |  S )Nr   )r   )rq   r   r   r   Ú_state_path½   ó   rr   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Ú )rC   rr   r!   Úitemsra   Ú
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.lockrt   ©r+   rr   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©rC   rr   r   r   r   r   Úload_watchlistÓ   s   r‰   Ú	watchlistc                 C   ó   t tdƒ| ƒ d S )Nr…   ©r:   rr   )rŠ   r   r   r   Úsave_watchlist×   ó   r�   c                   C   ó   t tdƒi ƒS ©Nzcooldowns.jsonrˆ   r   r   r   r   Úload_cooldownsÛ   rs   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_stateè   rs   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ñ   rs   r™   r   c                 C   r�   )Nzsession_state.json.lockr˜   r‚   )r   r   r   r   Úsave_session_stateõ   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›   rr   r   r   r   r   Úacquire_singleton_lockú   s   rœ   c                   @   s*   e Zd Zdefdd„Zdd„ Zdd„ ZdS )	r›   r,   c                 C   s   || _ d | _d S r   )r,   Ú_fh)Úselfr,   r   r   r   Ú__init__  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.)r,   r$   r   r%   r�   r&   r'   r(   ÚLOCK_NBÚBlockingIOErrorÚRuntimeErrorrG   r4   r5   ÚgetpidÚflush)rž   r   r   r   Ú	__enter__  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__  s   þzSingletonLock.__exit__N)Ú__name__Ú
__module__Ú__qualname__r   rŸ   r¥   rª   r   r   r   r   r›      s    )rT   )r`   r›   )4Ú__doc__r2   r5   r&   Útimer   r   Úpathlibr   Ú
contextlibr   Úconfig_loaderr   Úlogger_setupr   r@   Ú__file__Úresolver$   r   r   r   r!   r+   r:   rC   ÚdictrI   r4   rR   ÚlistrY   r_   rb   rc   re   ri   rk   rp   rr   r   rƒ   r‰   r�   r‘   r“   r•   r—   r™   rš   rœ   r›   r   r   r   r   Ú<module>   sT    



$	
