
    
iO                       d Z ddlmZ ddlZddlmZmZ ddlmZ ddl	m
Z
mZmZmZ ddlZ	 ddlZddlmZ dZe G d d             Ze G d	 d
             Z G d d      Zd*dZd+dZd,dZedk(  r ej>                  d      Z 	 d-	 	 	 	 	 	 	 	 	 d.dZ! ee  e!ddd       e!ddd       e!ddd       e!ddd       e!ddd       g!      Z"e"jG                  d"d d#      Z$ e%d$        e%d%e$jL                  jO                          d&e$jP                   d'e$jR                  d(        e%d)e$jL                  jU                         e jU                         z
  jV                          yy# e$ rZdZdZeZY dZ[<dZ[ww xY w)/u   
Memory-conscious loader and point-in-time query utilities for iVolatility-style
end-of-day SPY (or other) options CSV dumps.

Typical workflow
----------------
**Single CSV**

1. ``IVolatilityLoader.convert_to_parquet("huge.csv", "cache/options.parquet")``

**Several shards** (e.g. ``data_download.csv`` + ``data_download 2.csv`` …)

1. ``paths = IVolatilityLoader.sorted_ivolatility_csv_paths("SPY-Option-Data")``
2. ``IVolatilityLoader.combine_csvs_to_parquet(paths, "cache/options.parquet")``

**Query**

3. ``loader = IVolatilityLoader("cache/options.parquet")``
4. ``chain = loader.get_chain_for_date(pd.Timestamp("2013-02-15"))``
5. ``leg = chain.find_target_leg(target_dte=30, target_delta=-0.10, option_type="P")``

Dependencies: pandas, pyarrow (for Parquet I/O and predicate pushdown).
    )annotationsN)	dataclassfield)Path)AnyClassVarIteratorSequencec                      e Zd ZU dZded<   ded<   ded<   ded<   ded	<   ded
<   ded<   ded<   ded<   ded<   ded<   ded<   y)OptionContractz
    One EOD option quote + greeks on a single ``date`` (as-of / trade date).

    ``option_type`` is normalized to 'C' (call) or 'P' (put).
    pd.Timestampdate
expirationfloatstrikestroption_typebidaskmidivdeltagammathetavegaN)__name__
__module____qualname____doc____annotations__     </opt/rentech/trading_bot/RenTech/core/options_data_loader.pyr   r   2   sH     M	J	J	JILLL
Kr"   r   c                  P    e Zd ZU dZded<    ee      Zded<   	 	 	 	 	 	 	 	 d	dZy)
OptionChainz
    All contracts observed on one trading day (same ``as_of``).

    ``as_of`` is the chain date used for DTE = (expiration - as_of).days in
    ``find_target_leg``.
    r   as_of)default_factoryzlist[OptionContract]	contractsc           	        | j                   st        d      t        |      }| j                   D cg c]  }|j                  |k(  s| }}|s*t        d|d| j                  j                          d      t        j                  | j                        j                         }i }|D ]R  }t        j                  |j                        j                         }	t        |	|z
  j                        }
|
dk  rN|
||	<   T |st        d      t        |j                         fd      \  }}|D cg c]S  }t        j                  |j                        j                         |k(  r!t        j                  |j                         r|U }}|s#t        d	|j                          d
| d d      t        |fd      S c c}w c c}w )a  
        Pick one contract: closest expiration by DTE, then closest |delta - target|.

        Parameters
        ----------
        target_dte
            Desired days-to-expiration (calendar days between chain date and expiry).
        target_delta
            Model delta to match (e.g. -0.10 for a ~10-delta put; calls are usually positive).
        option_type
            'C', 'P', 'Call', 'Put', etc. (normalized to single-letter).
        z-OptionChain is empty; cannot find target leg.zNo contracts with option_type=z on .r   z9No non-expired expirations for the requested option type.c                &    t        | d   z
        S )N   )abs)kv
target_dtes    r#   <lambda>z-OptionChain.find_target_leg.<locals>.<lambda>}   s    #bejFXBY r"   keyz-No contracts with finite delta at expiration u    (DTE≈z, target DTE=z).c                X    t        t        | j                        t              z
        S )N)r-   r   r   )ctarget_deltas    r#   r0   z-OptionChain.find_target_leg.<locals>.<lambda>   s    U177^eL>Q-Q)R r"   )r(   
ValueError_normalize_option_typer   r&   r   pd	Timestamp	normalizer   intdaysminitemsmathisfiniter   )selfr/   r5   r   wantr4   typedr&   expiriesexpdtebest_expbest_dtebuckets    ``           r#   find_target_legzOptionChain.find_target_legT   s   $ ~~LMM%k2 NNDqammt.CDD=dXT$**//J[I\\]^__TZZ(224 -/ 	 A,,q||,668CsU{(()CQwHSM	  XYY !!17YZ( 
||ALL)335Aagg& 
 
 ??P Q"=B@  6RSSE E*
s   GG2AGN)r/   r;   r5   r   r   r   returnr   )	r   r   r   r   r    r   listr(   rJ   r!   r"   r#   r%   r%   H   sO     &+D&AI#A8T8T 8T 	8T
 
8Tr"   r%   c                  L   e Zd ZU dZi ddddddddddddd	d
dd
ddddddddddddddddddZded<   g dZded<   d%dZedddd	 	 	 	 	 	 	 	 	 	 	 d&d       Z	ed'd        Z
edddd	 	 	 	 	 	 	 	 	 	 	 d(d!       Ze	 	 	 	 	 	 	 	 	 	 d)d"       Zd*d#Zd+d$Zy),IVolatilityLoadera  
    Load iVolatility CSV (via Parquet cache) and query single-day chains.

    ``COLUMN_MAP`` maps **raw CSV header** -> **internal column name** used in
    Parquet and in memory. Adjust keys to match your export (e.g. some files use
    ``implied_volatility`` instead of ``iv``).
    r   r   r   zcall/putr   call_putzCall/Putimplied_volatilityr   IVr   r   pricez
mean pricer   r   r   r   rhozClassVar[dict[str, str]]
COLUMN_MAPr   r   r   r   r   r   r   r   r   r   r   r   zClassVar[list[str]]_PARQUET_COLSc                    t        |      j                         | _        | j                  j                  j	                         dvryy)z
        Parameters
        ----------
        file_path
            Path to a **Parquet** file produced by ``convert_to_parquet``. Using Parquet
            keeps random access and row-group filters fast and RAM low.
        >   .pq.parquetN)r   
expanduser	file_pathsuffixlower)rA   r[   s     r#   __init__zIVolatilityLoader.__init__   s=     i335>>  &&(0CC Dr"   i@ )r   r   N)	chunksizedate_columnsfloat32_colsc                  t         t        t        d      t        t	        |       j                         } t	        |      j                         }|j                  j                  dd       |d}t        j                  | |      }d}	 |D ]  }t        j                  |t        j                  ||      }|j                  r7t        j                  j                  |d	      }	|)t        j                   t#        |      |	j$                        }|j'                  |	        	 ||j)                          |S # ||j)                          w w xY w)
a  
        Stream-read a large CSV in chunks, normalize columns, downcast floats, write Parquet.

        * Dates are parsed with ``pd.to_datetime`` (UTC-naive midnight normalized).
        * Greeks, IV, bid/ask/mid/strike use ``float32`` to halve RAM vs float64.
        * ``mid`` is ``(bid+ask)/2`` when not otherwise supplied from ``price``.

        Parameters
        ----------
        chunksize
            Rows per read from CSV; tune based on available RAM and column count.
        NzFconvert_to_parquet requires pyarrow. Install with: pip install pyarrowTparentsexist_ok	r   r   r   r   r   r   r   r   r   r_   
column_mapr`   ra   Fpreserve_index)pqpaImportError_PYARROW_IMPORT_ERRORr   rZ   parentmkdirr8   read_csvrN   _standardize_chunkrT   emptyTablefrom_pandasParquetWriterr   schemawrite_tableclose)
csv_pathparquet_pathr_   r`   ra   readerwriter	raw_chunkstdtables
             r#   convert_to_parquetz$IVolatilityLoader.convert_to_parquet   sA   * :X() >,,.L)446!!$!>
L X;*.	# *	'::0;;!-!-	 ;  99,,S,G>--c,.?NF""5)* ! ! "s   
BD7 7Ec                x    t        |       j                         }t        |j                  d      t              }|S )u  
        Return iVolatility split exports in a stable order: ``data_download.csv`` first,
        then ``data_download 2.csv`` … ``data_download 7.csv`` by numeric suffix.

        Use this when combining multiple shards from the same vendor export.
        zdata_download*.csvr1   )r   rZ   sortedglob_ivolatility_csv_sort_key)	directorydpathss      r#   sorted_ivolatility_csv_pathsz.IVolatilityLoader.sorted_ivolatility_csv_paths  s2     O&&(qvv239RSr"   c                  t         t        t        d      t        t	        |      j                         }|j                  j                  dd       |d}d}d}	 | D ]  }t	        |      j                         }|j                         st        d|       t        j                  ||      }	|	D ]  }
t        j                  |
t        j                  ||	      }|j                  r7t        j                   j#                  |d
      }|)t        j$                  t'        |      |j(                        }|j+                  |       d}  	 ||j-                          	 |st/        d      |S # ||j-                          w w xY w)u'  
        Merge several iVolatility CSV shards into **one** Parquet file (streaming).

        Each file is read in ``chunksize`` rows so memory stays bounded. Schema is
        unified via ``COLUMN_MAP`` (e.g. ``mean price`` → ``price`` for files that
        omit ``is_settlement``).
        NzKcombine_csvs_to_parquet requires pyarrow. Install with: pip install pyarrowTrc   rf   FzCSV not found: rg   rh   rj   z@No rows written from csv_paths (empty files or schema mismatch).)rl   rm   rn   ro   r   rZ   rp   rq   is_fileFileNotFoundErrorr8   rr   rN   rs   rT   rt   ru   rv   rw   r   rx   ry   rz   r6   )	csv_pathsr|   r_   r`   ra   r~   	wrote_anyr{   pathr}   r   r   r   s                r#   combine_csvs_to_parquetz)IVolatilityLoader.combine_csvs_to_parquet  s     :]() L)446!!$!>
L +/		% %H~002||~+odV,DEETY?!' %I+>>!#4#?#?%1%1	 ? C yy HH00U0KE~!#!1!1#l2CU\\!R&&u- $I%%( !_`` ! "s   C/E. .Fc                  | j                  t        |            }|D ]F  }||j                  v st        j                  ||   d      j
                  j                         ||<   H d|j                  v r|d   j                  t              |d<   d|j                  vrd|j                  v rt        j                  |d   d      |d<   nqt        j                  |j                  d      d      t        j                  |j                  d      d      z   d	z  |d<   nt        j                  |d   d      |d<   d
D ].  }||j                  v st        j                  ||   d      ||<   0 g d}|j                  |D cg c]  }||j                  v s| c}      }|j                  r|S |D ](  }||j                  v s||   j                  d      ||<   * t        j                  D ]8  }||j                  vs|dv rt        j                   nt"        j$                  ||<   : |t        j                     j'                         }|S c c}w )zLApply renames, parse dates, compute mid, coerce dtypes, keep Parquet schema.columnscoerceerrorsr   r   rR   r   r   g       @)r   r   r   r   r   r   r   r   )r   r   r   r   r   r   )subsetfloat32)r   r   r   )renamedictr   r8   to_datetimedtr:   mapr7   
to_numericgetdropnart   astyperN   rV   NAr?   nancopy)	dfri   r`   ra   chunkcolr4   needouts	            r#   rs   z$IVolatilityLoader._standardize_chunkb  s0    		$z"2	3 	XCemm#^^E#JxHKKUUWc
	X EMM)#(#7#;#;<R#SE-  %%--'!}}U7^HMe "eii.>x PSUS`S`afajajkpaq  {C  TD  !D  HK   Ke==uhGE%LR 	DAEMM!==q(Ca	D
 M$K1U]]8JQ$KL;;L 	6AEMM! 8??95a	6
 #00 	]A%$%)N$N255TXT\T\a	] %33499;
 %Ls   I
I
c                *   t         t        d      t        t        j                  | j                        }t               }t        |j                        D ]  }|j                  |dg      j                  d      }t        j                  |j                         d      j                  j                         }|j                         j!                         D ]&  }|j#                  t        j$                  |             (  t'        t)        |            S )zBYield unique ``date`` values present in the Parquet file (sorted).zpyarrow is requiredr   r   r   r   r   )rl   rn   ro   ParquetFiler[   setrangenum_row_groupsread_row_groupcolumnr8   r   	to_pandasr   r:   r   uniqueaddr9   iterr   )rA   pfdatesir   sus          r#   iter_chain_datesz"IVolatilityLoader.iter_chain_dates  s    :34:OO^^DNN+#&5r(() 	+A##Ax#8??BCs}}x@CCMMOAXXZ&&( +		",,q/*+		+ F5M""r"   c                0   t         t        d      t        t        j                  |      j                         }dd|fgdd|j                         fgdd|j                         fgg}d}t        j                         }|D ]3  }	 t        j                  | j                  | j                  |d      }d} n ||j                  rt        d      ||j                  rt        |g       S t        j                   |d   d	
      j"                  j                         |d<   ||d   |k(     }|j                  rt        |g       S |j%                         D 	cg c]  \  }}	t'        |	       }
}}	t        ||
      S # t        $ r}|}Y d}~d}~ww xY wc c}	}w )a  
        Load **only** rows for ``target_date`` using Parquet predicate pushdown.

        Uses the PyArrow engine so filters prune row groups where possible. If your
        ``date`` column was written as plain ``date32``/``timestamp`` without timezone,
        equality against a normalized ``pd.Timestamp`` is typically recognized by the
        scanner. We do **not** fall back to a full-file read (that would defeat the
        purpose on multi-GB datasets).
        Nz#get_chain_for_date requires pyarrowr   z==pyarrow)r   filtersenginezCould not apply Parquet row filter on 'date'. Re-run convert_to_parquet so 'date' is timezone-naive datetime64, or pass a filter literal matching the on-disk type.r&   r(   r   r   )rl   rn   ro   r8   r9   r:   to_pydatetimer   	DataFrameread_parquetr[   rV   	Exceptionrt   RuntimeErrorr%   r   r   iterrows_row_to_contract)rA   target_datetsfilter_variantslast_errr   r   e_rr(   s              r#   get_chain_for_datez$IVolatilityLoader.get_chain_for_date  s    :CDJ__\\+&002 dB dB,,./0dBGGI&'=
 &*\\^& 	G__NN ..#$	  	 BHHF 	 88R266^^BvJx@CCMMO6
6
b !88R26657[[]CTQ%a(C	Cy99)  & Ds   /E9F9	FF

F)r[   
str | PathrK   None)r{   r   r|   r   r_   r;   r`   tuple[str, ...]ra   tuple[str, ...] | NonerK   r   )r   r   rK   z
list[Path])r   zSequence[str | Path]r|   r   r_   r;   r`   r   ra   r   rK   r   )
r   pd.DataFrameri   zdict[str, str]r`   r   ra   r   rK   r   )rK   zIterator[pd.Timestamp])r   r   rK   r%   )r   r   r   r   rT   r    rV   r^   staticmethodr   r   r   rs   r   r   r!   r"   r#   rN   rN      s   ,,l, 	(, 	M	, 	M, 	M, 	d, 	d, 	u, 	u, 	, 	g, 	, 	,  	!," 	#,$ 	u%,J( ,*M&  
 !(>/3?? ? 	?
 &? -? 
? ?B 	 	 
 !(>/3A'A A 	A
 &A -A 
A AF .. #. &	.
 &. 
. .`#6:r"   rN   c                    | j                   }|dk(  ryd}|j                  |      r+|t        |      d j                         }	 dt	        |      dfS dd|fS # t
        $ r dd|fcY S w xY w)uP   ``data_download.csv`` first, then ``data_download 2.csv`` … by numeric suffix.data_download)r   r    zdata_download Nr   r   r,   )stem
startswithlenstripr;   r6   )pr   prefixsufs       r#   r   r     s    66DFv3v;=!'')	 s3x$$ q$<  	 q$<	 s   A A'&A'c                   | %t        | t              rt        j                  |       ryt	        |       j                         j                         }|j                  d      ry|j                  d      ry|dv r	|dk(  rdS dS |r|d d S dS )Nr   CP>   PUTCALLr   r,   )
isinstancer   r?   isnanr   r   upperr   )xr   s     r#   r7   r7     s    yZ5)djjmAA||C||CO6ks*s*1Ra52r"   c                   t        t        j                  | d         j                         t        j                  | d         j                         t	        | d         t        | d         t	        | d         t	        | d         t	        | d         t	        | d         t	        | d	         t	        | d
         t	        | d         t	        | d               S )Nr   r   r   r   r   r   r   r   r   r   r   r   rU   )r   r8   r9   r:   r   r   )r   s    r#   r   r     s    \\!F)$..0<<,0::<Qx[!-()!E(O!E(O!E(O4>AgJAgJAgJ1V9 r"   __main__z
2024-06-15r   c                b    t        j                  |       }t        t        |||dddd|ddd      S )	Ng      ?g?g?g      ?g{Gz?g{GzgQ?rU   )r8   r9   r   day)rE   r   r   optexp_tss        r#   mkr     sC     c"
 	
r"   z
2024-07-10g     @@gz
2024-07-17g      @g
ףp=
g     ~@gffffffֿg     `@g)\(z
2024-08-15g     ~@gr      )r/   r5   r   z,Demo: 30 DTE (closest bucket), ~10-delta putz  Selected: exp=z strike=z delta=z.3fz  DTE=)r   r   rK   ztuple[int, int, str])r   r   rK   r   )r   z	pd.SeriesrK   r   )r   )
rE   r   r   r   r   r   r   r   rK   r   ),r   
__future__r   r?   dataclassesr   r   pathlibr   typingr   r   r	   r
   pandasr8   r   rm   pyarrow.parquetparquetrl   ro   rn   r   r   r%   rN   r   r7   r   r   r9   r   r   chainrJ   legprintr   r   r   r   r:   r<   r!   r"   r#   <module>r      s  0 #  (  4 4 !  !   * CT CT CTVC: C:V

* z
",,|
$C 	


 
 	

 

0 |UE*|UE*|UE*|UE*|UE*
	E 

2Es

SC	
89	S^^00238CJJ<wsyyY\o
^_	FCNN,,.@FFG
HIU W  	B	Bs   
E# #E9(E44E9