
    xj                    .   d Z ddlmZ ddlZddlZddlZddlZddlZddlm	Z	 ddl
mZ ddlmZ 	 ddlZn# e$ r dZY nw xY wdZdd	hZ ej        d
          Zd&dZd'd(dZd)dZe	d             Zddd*dZddd+dZd,dZd-d Zddd.d!Zd/d#Zd0d%ZdS )1a  Crash-safe WebUI turn journal helpers.

The journal is deliberately tiny: one JSONL file per session, append-only events,
and read helpers that tolerate malformed lines. Recovery and repair can then
reason about submitted turns without depending on in-memory stream state.
    )annotationsN)contextmanager)Path)Iterable_turn_journal	completedinterruptedz^[A-Za-z0-9_.-]+$returnr   c                 ,    ddl m}  t          |           S )Nr   SESSION_DIR)
api.modelsr   r   r   s    3/Users/bertmahoney/hermes-webui/api/turn_journal.py_default_session_dirr      s"    &&&&&&    
session_idstrsession_dirPath | Nonec                   t          | pd                                          }|r"d|v sd|v st                              |          st	          d          |t          |          nt                      }|t          z  | dz  S )N /\zinvalid session_idz.jsonl)r   strip_SESSION_ID_RE	fullmatch
ValueErrorr   r   TURN_JOURNAL_DIR_NAME)r   r   sidroots       r   _journal_pathr!   "   s    
jB


%
%
'
'C /#**>3K3KC3P3P-... + 74=Q=S=SD''S...88r   c                     t          j        dt          j                               dt          j                    j        d d          S )Nz%Y%m%dT%H%M%SZ-   )timestrftimegmtimeuuiduuid4hex r   r   _make_turn_idr,   *   s>    m,dkmm<<VVtz||?OPSQSPS?TVVVr   c              #  ^  K   t           dV  dS t          j        |                                 t           j                   	 dV  t          j        |                                 t           j                   dS # t          j        |                                 t           j                   w xY w)a  Serialize multi-process journal writes when advisory locks exist.

    ``O_APPEND`` keeps normal same-process appends simple, but a long JSONL event
    can exceed POSIX's small atomic-write boundary.  On Unix, take an advisory
    lock around the single event write+fsync so two WebUI worker processes cannot
    interleave large submitted-message payloads into corrupted JSONL.  Platforms
    without ``fcntl`` keep the previous best-effort append behavior.
    N)_fcntlflockfilenoLOCK_EXLOCK_UN)file_objs    r   _journal_file_lockr4   .   s       ~
L""FN3338X__&&77777X__&&7777s   A9 93B,r   eventdictc                  t          |t                    st          d          t          |                    d          pd                                          }|st          d          t          |          }|                    dd           t          |           |d<   |                    dt                                 |                    d	t          j	                               t          | |
          }|j                            dd           t          j        |dd          dz   }t          j        |t          j        t          j        z  t          j        z  d          }t          j        |dd          5 }t+          |          5  |                    |           |                                 t          j        |                                           ddd           n# 1 swxY w Y   ddd           n# 1 swxY w Y   	 t          j        |j        t          j                  }		 t          j        |	           t          j        |	           n# t          j        |	           w xY wn# t8          $ r Y nw xY w|S )zAppend one turn journal event and fsync it before returning.

    The returned event is the exact payload written, with default ``version``,
    ``session_id``, ``turn_id``, and ``created_at`` fields filled in.
    zevent must be a dictr6   r   zevent is requiredversion   r   turn_id
created_atr5   T)parentsexist_okF),:)ensure_ascii
separators
i  autf-8encodingN)
isinstancer7   	TypeErrorr   getr   r   
setdefaultr,   r%   r!   parentmkdirjsondumpsosopenO_CREATO_APPENDO_WRONLYfdopenr4   writeflushfsyncr0   O_DIRECTORYcloseOSError)
r   r6   r   
event_namepayloadpathlinefdfhdir_fds
             r   append_turn_journal_eventrc   B   s    eT"" 0.///UYYw''-2..4466J .,---5kkGy!$$$
OOGLy-//222|TY[[111===DKdT222:gEjIIIDPD	rzBK/"+=u	E	EB	2sW	-	-	- "## 	" 	"HHTNNNHHJJJHRYY[[!!!	" 	" 	" 	" 	" 	" 	" 	" 	" 	" 	" 	" 	" 	" 	"" " " " " " " " " " " " " " "
bn55	HVHVBHV   Nsa   HAG0$H0G4	4H7G4	8HHH$I; 8I! I; !I77I; ;
JJc               P   t          | |          }g }g }	 |                    d                                          }n## t          $ r t	          |           g g dcY S w xY wt          |d          D ]\  }}|                                s	 t          j        |          }n-# t          j	        $ r |
                    ||d           Y Yw xY wt          |t                    r|
                    |           |
                    ||d           t	          |           ||dS )zDRead a session journal, returning valid events plus malformed lines.r5   rE   rF   )r   events	malformedr:   )start)r_   raw)r!   	read_text
splitlinesFileNotFoundErrorr   	enumerater   rN   loadsJSONDecodeErrorappendrH   r7   )	r   r   r^   re   rf   linesline_norh   r6   s	            r   read_turn_journalrr   l   sd   ===DFIN00;;== N N N!*oo"MMMMMN!%q111 < <yy{{ 		JsOOEE# 	 	 	gc::;;;H	 eT"" 	<MM%    gc::;;;;j//V)TTTs#   (A   A A B$$'CCre   Iterable[dict]"tuple[dict[str, dict], list[dict]]c                2   i }i }| D ]}t          |t                    st          |                    d          pd                                          }|sQt          |          r)|                    |g                               |           |                    |          }|Jt          |                    d          pd          t          |                    d          pd          k    r|||<   d |	                                D             }||fS )a  Return the latest event per ``turn_id`` and any terminal-collision entries.

    The first element is the latest event per turn_id (same overwrite-by-timestamp
    behaviour as before).  The second element is a list of collision records, one
    per turn_id that had more than one terminal event.  Each collision record
    contains ``turn_id`` and the ``events`` list (in ascending created_at order).

    A collision means the same logical turn recorded both ``completed`` and
    ``interrupted`` terminal events -- the derived state still picks the latest
    by timestamp, but callers can now detect and audit the double-terminal
    situation explicitly rather than having it silently collapse.
    r;   r   Nr<   r   c                d    g | ]-\  }}t          |          d k    |t          |d           d.S )r:   c                J    t          |                     d          pd          S )Nr<   r   )floatrJ   )es    r   <lambda>z7derive_turn_journal_states.<locals>.<listcomp>.<lambda>   s     eAEE,DWDWD\[\>]>] r   )key)r;   re   )lensorted).0tidevtss      r   
<listcomp>z.derive_turn_journal_states.<locals>.<listcomp>   sL       Ct99q== 6$4]4]#^#^#^__==r   )
rH   r7   r   rJ   r   is_terminal_turn_eventrK   ro   rx   items)re   statesterminal_eventsr6   r;   previous
collisionss          r   derive_turn_journal_statesr      s2    !F-/O $ $%&& 	eii	**0b117799 	!%(( 	B&&w33::5AAA::g&&uUYY|%<%<%ABBeHLLYeLfLfLkjkFlFlll#F7O (..00  J
 :r   	stream_id
str | Nonec                T   t          |pd                                          }|sd S d }| D ]{}t          |t                    st          |                    d          pd          |k    rAt          |                    d          pd                                          }|r|}||S )Nr   r   r;   )r   r   rH   r7   rJ   )re   r   streamlatestr6   r;   s         r   _latest_turn_id_for_streamr      s    b!!''))F tF  %&& 	uyy%%+,,66eii	**0b117799 	FMr   c                  t          |          }t          |          |d<   |                    d          s=t          | |          }t	          |                    d          pg |          }|r||d<   t          | ||          S )zDAppend a lifecycle event for the turn associated with ``stream_id``.r   r;   r5   re   )r7   r   rJ   rr   r   rc   )r   r   r6   r   r]   journalr;   s          r   $append_turn_journal_event_for_streamr      s     5kkGy>>GK;;y!! )#JKHHH,W[[-B-B-Hb)TT 	)!(GI$ZkRRRRr   	list[str]c                    t          |           t          z  }|                                sg S t          d |                    d          D                       S )Nc              3  L   K   | ]}|                                 |j        V   d S N)is_filestem)r~   r^   s     r   	<genexpr>z0iter_turn_journal_session_ids.<locals>.<genexpr>   s1      VVt||~~V$)VVVVVVr   z*.jsonl)r   r   existsr}   glob)r   journal_dirs     r   iter_turn_journal_session_idsr      sY    {##&;;K 	VV(8(8(C(CVVVVVVr   boolc                \    t          | pi                     d          pd          t          v S )Nr6   r   )r   rJ   _TERMINAL_EVENTS)r6   s    r   r   r      s-      ))/R004DDDr   )r
   r   r   )r   r   r   r   r
   r   )r
   r   )r   r   r6   r7   r   r   r
   r7   )r   r   r   r   r
   r7   )re   rs   r
   rt   )re   rs   r   r   r
   r   )
r   r   r   r   r6   r7   r   r   r
   r7   )r   r   r
   r   )r6   r7   r
   r   ) __doc__
__future__r   rN   rP   rer%   r(   
contextlibr   pathlibr   typingr   fcntlr.   ImportErrorr   r   compiler   r   r!   r,   r4   rc   rr   r   r   r   r   r   r+   r   r   <module>r      s     # " " " " "  				 				   % % % % % %               FFF ( / 011   9 9 9 9 9W W W W 8 8 8.  $	' ' ' ' ' 'T FJ U U U U U U0$ $ $ $L   *  $S S S S S S$W W W WE E E E E Es   5 ??