
    <jE                       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mZmZ ddl	m
Z
 dZdZdZdZd	Zd	Zd
ZdZdZdZdZdZdZdZdZdZdZdZdZdZdZdZ dZ!ejD                  jG                  e!ddd      Z$ejD                  jG                  e!ddd      Z%ejD                  jG                  e!ddd      Z&ejD                  jG                  e!ddd       Z'ejD                  jG                  e!ddd!      Z(ejD                  jG                  e!ddd"      Z)g d#Z*d0d1d$Z+d2d%Z,d3d&Z-d4d'Z.d5d(Z/d6d)Z0d7d*Z1d8d+Z2d9d,Z3dddddddd-	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 d:d.Z4d;d/Z5y)<u   spawn_safety_governor.py — task-2774: 무한 spawn 방지 게이트(governor).

stdlib only, 외부 의존 0. 순수 함수 + ledger I/O.
governor 자체는 절대 새 프로세스/재귀를 생성하지 않는다.
    )annotationsN)datetimetimezone)Optional         
   g      ?)	completedfailedblockedcrashALLOWBLOCKQUEUETRIPNON_TERMINAL	DEDUP_HITSINGLEFLIGHT_LIVERATE_EXCEEDEDLOOP_EXCEEDEDCIRCUIT_TRIPCOUNTERS_UNREADABLECOUNTERS_SAVE_FAILEDLEDGER_UNREADABLEz/home/jay/workspacememoryeventszcallback_4tuple_index.jsonlstatezspawn_decisions.jsonlzspawn_counters.jsonzanu_session_alive.lockzspawn_circuit_tripped.markerp0b_driver_enabled)evaluate_spawn_task_familyRATE_PER_MINRATE_PER_HOURLOOP_PER_TASKLOOP_PER_FAMILYCIRCUIT_WINDOW_MINCIRCUIT_MAX_SPAWNCIRCUIT_FAIL_RATETERMINAL_STATESr   r   r   r   REASON_ALLOWREASON_NON_TERMINALREASON_DEDUP_HITREASON_SINGLEFLIGHT_LIVEREASON_RATE_EXCEEDEDREASON_LOOP_EXCEEDEDREASON_CIRCUIT_TRIPREASON_COUNTERS_UNREADABLEREASON_COUNTERS_SAVE_FAILEDREASON_LEDGER_UNREADABLECANONICAL_ROOTDEFAULT_DEDUP_LEDGERDEFAULT_DECISIONS_LOGDEFAULT_COUNTERS_PATHDEFAULT_LIVE_LOCKDEFAULT_CIRCUIT_MARKERDEFAULT_P0B_FLAGc                    t        | t              rt        | j                  ||      xs |      S t        t	        | ||      xs |      S )uM   candidate dict 또는 attribute 접근 가능 객체에서 필드 값 반환.)
isinstancedictstrgetgetattr)	candidatenamedefaults      2/home/jay/workspace/utils/spawn_safety_governor.py	_cand_getrE   Q   s@    )T"9==w/:7;;wy$0;G<<    c                `    t        j                  d| xs d      }|r|j                  d      S | S )u   task_id에서 base 'task-N' 추출.

    예: task-2774+1 → task-2774, task-2774-r2 → task-2774, task-2774 → task-2774.
    z^(task-\d+) r   )rematchgroup)task_idms     rD   r!   r!   X   s-    
 	B/A1771:''rF   c                 H    t        j                  t        j                        S )u+   현재 UTC datetime(timezone-aware) 반환.)r   nowr   utc rF   rD   _now_utcrR   a   s    <<%%rF   c                $    | j                  d      S )u#   datetime → ISO8601 UTC 문자열.z%Y-%m-%dT%H:%M:%SZ)strftime)dts    rD   _isorV   f   s    ;;+,,rF   c                    	 t        | dd      5 }t        j                  |      }ddd       t	        t
              st        d      |S # 1 sw Y   &xY w# t        $ r i cY S w xY w)uN  spawn_counters.json 로드.

    (a) 파일 부재(첫 실행, 정상) → 빈 dict 반환(ALLOW 경로 보존).
    (b) 파일 존재하나 손상/파싱오류/권한·IO 오류(위험) → 예외 전파 → 호출부에서
        fail-closed(COUNTERS_UNREADABLE → BLOCK) 처리. 빈 dict 로 무시 금지(fail-open 차단).
    rutf-8encodingNz-spawn_counters.json is not a dict (corrupted))openjsonloadFileNotFoundErrorr<   r=   
ValueError)counters_pathfdatas      rD   _load_countersrd   k   sg    -w7 	 199Q<D	 
 dD!HIIK	  	  	s'   A AA AA A%$A%c                   t         j                  j                  |       xs d}t        j                  |d       t	        j
                  |d      \  }}	 t        j                  |dd      5 }t        j                  ||d	
       ddd       t        j                  ||        y# 1 sw Y    xY w# t        $ r' 	 t        j                  |        # t        $ r Y  w xY ww xY w)uL  spawn_counters.json 원자적 저장 (tempfile + os.replace).

    저장 실패(디스크 풀/권한/IO 오류) 시 예외를 상위로 전파한다(silent ignore 금지).
    호출부는 저장 실패를 fail-closed(COUNTERS_SAVE_FAILED → BLOCK)로 처리하여
    count 미증가로 인한 무한 spawn 을 차단한다.
    .Texist_okz.tmp)dirsuffixwrY   rZ   Fensure_asciiN)ospathdirnamemakedirstempfilemkstempfdopenr]   dumpreplace	Exceptionunlink)ra   countersdir_fdtmp_pathrb   s         rD   _save_countersr}   ~   s     77??=)0SDKKt$##V<LB	YYr31 	7QIIh6	7


8]+	7 	7  	IIh 	  		sH   B2 .B&B2 &B/+B2 2	C"<CC"	CC"CC"c                0   	 t         j                  j                  |       xs d}t        j                  |d       t	        | dd      5 }|j                  t        j                  |d      d	z          d
d
d
       y
# 1 sw Y   y
xY w# t        $ r Y y
w xY w)u<   spawn_decisions.jsonl에 결정 1줄 append (silent drop 0).rf   Trg   arY   rZ   Frl   
N)	rn   ro   rp   rq   r\   writer]   dumpsrw   )decisions_pathrecordrz   rb   s       rD   _append_decisionr      s    ww~.5#
D4(.#8 	CAGGDJJvE:TAB	C 	C 	C s0   AB	 
*A=4B	 =BB	 B	 		BBc              #  >  K   	 t        | dd      }	 |D ]@  }|j                         }|s	 t        j                  |      }t        |t              s=| B 	 |j                          y# t        $ r Y yw xY w# t
        $ r Y ow xY w# |j                          w xY ww)u  JSONL 파일을 라인 단위로 스트리밍하며 파싱된 dict 를 yield (bounded read).

    전량 메모리 적재 금지 — 호출부는 일치 시 즉시 break 하여 조기탈출한다.
    (a) 파일 부재(FileNotFoundError) → 빈 제너레이터(정상 초기 상태).
    (b) 권한/IO 오류 등 기타 open 예외 → 상위로 전파 → 호출부 fail-closed
        (LEDGER_UNREADABLE → BLOCK). 손상된 개별 라인은 skip(append-only 부분쓰기 허용).
    rX   rY   rZ   N)	r\   r_   stripr]   loadsrw   r<   r=   close)ro   rb   lineobjs       rD   _iter_jsonl_linesr      s     sW- 		D::<Djj& #t$			 	
	    
 	
	sb   BA* B A9B B B*	A63B5A66B9	BB BB BBc                    d| v r	 t        | d         S t        t        | j	                  dd                  S # t        $ r Y /w xY w)u^   이벤트의 epoch(float) 반환. 'epoch' 필드 우선, 없으면 'ts' 문자열 1회 파싱.epochtsrH   )floatrw   _ts_to_epochr>   r?   )es    rD   _get_event_epochr      sM    !|	7$$ AEE$O,--  		s   8 	AA)ra   ledger_pathrO   live_lock_presentr   circuit_marker_pathp0b_flag_pathc               
  !"# |t         }|t        }t        |t        }|t        }t        | d      "t        | d      }t        | d      }	| d|	 !|
t               }t        |      #d!"#fd}
dfd}|	t        vr | |
t        t                    S 	 t        |      D ]X  }t        |j                  dd	            "k(  s"t        |j                  dd	            |k(  sA | |
t        t                    c S  t              D ]\  }t        |j                  d
d	            !k(  s"t        |j                  dd	            t        k(  sE | |
t        t                    c S  	 |#t$        j&                  j)                  t*              }|r | |
t,        t.                    S 	 t1        |      }|j                  dg       }t5        |t6              sg }|j9                         }|D cg c]  }|t;        |      z
  dk  s| }}|D cg c]  }|t;        |      z
  dk  s| }}t=        |      t>        k\  st=        |      t@        k\  r | |
t,        tB                    S |j                  di       }|j                  di       }t5        |tD              si }t5        |tD              si }tG        "      }tI        |j                  "d            }tI        |j                  |d            }|tJ        k\  s	|tL        k\  r | |
t        tN                    S tP        dz  }|D cg c]  }|t;        |      z
  |k  s| }}t=        |      }tS        d |D              }|dkD  r||z  nd}|tT        k\  xs	 |tV        k\  }|r	 t$        j&                  jY                  |      xs d}t%        jZ                  |d       t]        |dd      5 }|j_                  #dz          ddd       t$        j&                  j)                  |      r)	 t]        |dd      5 }|j_                  d       ddd        | |
t`        tb                    S |dz   |"<   |dz   ||<   |je                  #|"dd       ||d<   ||d<   ||d<   	 tg        ||        |
t        tj              }  ||       S # t         $ r  | |
t        t"                    cY S w xY w# t         $ r  | |
t        t2                    cY S w xY wc c}w c c}w c c}w # 1 sw Y   xY w# t         $ r Y .w xY w# 1 sw Y   xY w# t         $ r Y w xY w# t         $ r  | |
t        th                    cY S w xY w) u  spawn 가능 여부를 평가하고 Decision dict를 반환한다.

    차단 순서(첫 차단에서 멈춤, 모든 결정은 spawn_decisions.jsonl에 append):
      1. terminal-only HARD GATE: terminal_state가 TERMINAL_STATES에 없으면 BLOCK
      2. dedup: 동일 key가 이미 ALLOW로 기록돼 있으면 BLOCK
      3. single-flight: live_lock 존재 시 QUEUE
      4. rate-limit: 분당/시간당 초과 시 QUEUE
      5. loop-budget: task/family 누적 ALLOW 초과 시 BLOCK
      6. circuit-breaker: 윈도우 내 spawn/실패율 초과 시 TRIP
      7. 전부 통과 → ALLOW

    반환 구조: dict(decision, reason, task_id, key, ts)
    NrL   head_shaterminal_state|decisionc                $    t        | |      S )N)r   reasonrL   keyr   )r=   )r   r   r   rL   ts_strs     rD   _make_decisionz&evaluate_spawn.<locals>._make_decision   s    
 	
rF   c                     t        |        | S )u+   spawn_decisions.jsonl에 append 후 반환.)r   )dr   s    rD   _record_and_returnz*evaluate_spawn.<locals>._record_and_return  s    +rF   rH   r   r   <   i  per_task
per_familyr   c              3  ^   K   | ]%  }t        |j                  d d            dk(  s"d ' yw)outcomerH   failr   N)r>   r?   ).0r   s     rD   	<genexpr>z!evaluate_spawn.<locals>.<genexpr>J  s'     U1QUU9b5I1Jf1TQUs   #--        rf   Trg   rk   rY   rZ   r   r   r   allow)r   r   rL   r   )r   r>   r   r>   returnr=   )r   r=   r   r=   )6r7   r5   r6   r9   r:   rE   rR   rV   r)   r   r+   r   r>   r?   r,   r   rw   r3   rn   ro   existsr8   r   r-   rd   r1   r<   list	timestampr   lenr"   r#   r.   r=   r!   intr$   r%   r/   r&   sumr'   r(   rp   rq   r\   r   r   r0   appendr}   r2   r*   )$rA   ra   r   rO   r   r   r   r   r   r   r   r   entryry   r   now_tsr   events_1minevents_1hourr   r   family
task_countfamily_count
window_secevents_windowwindow_spawn_count
fail_count	fail_ratecircuit_trippedrz   rb   decision_recr   rL   r   s$        `                           @@@rD   r    r       sF   2 -*."4( y)4Gy*5Hy*:;N Ja'
(C {j#YF
 _,!.8K"LMMS&{3 	SEEIIi,-8EIIj"56(B).@P*QRR	S '~6 	SEEIIeR()S0EIIj"56%?).@P*QRR	S  GGNN+<=!.8P"QRRU!-0 <<"-Ffd#]]_F  &L!2B12E)E)KALKL%N!2B12E)E)MANLN
;<'3|+<+M!.8L"MNN  ||J3H||L"5Jh%j$'
'"Fx||GQ/0Jz~~fa01L]"lo&E!.8L"MNN $b(J &U1&3CA3F*F**TQUMU]+UUUJ5G!5K00QTI 	// 	*))  	77??#67>3DKKt,)3A 'Q&' 77>>-(-w? '1GGI&'
 ".7J"KLL $aHW%)Jv
MM&WQXYZ%HZ'H\#HXV}h/ "%6Ll++{  S!.8P"QRRS  U!.8R"STTU MN, V' ' 		' ' &  V!.8S"TUUVs   ,R$ 
R$ )R$ -R$ 0"R$ R$ ,R$ -S
 0S0S0S5(S5(S:?S:AT S?"T 
T( T*T( ?T8 $ SS
 S-,S-?T	T 	TTT%!T( (	T54T58 UUc                    	 | j                  d      r| dd dz   } t        j                  |       }|j                   |j	                  t
        j                        }|j                         S # t        $ r Y yw xY w)uG   ISO8601 UTC 문자열을 epoch(float)로 변환. 파싱 실패 시 0.0.ZNz+00:00)tzinfor   )	endswithr   fromisoformatr   rv   r   rP   r   rw   )r   rU   s     rD   r   r   z  so    	??3CR[8+F##F+998<<0B||~ s   A)A, ,	A87A8)rH   )rB   r>   rC   r>   r   r>   )rL   r>   r   r>   )r   r   )rU   r   r   r>   )ra   r>   r   r=   )ra   r>   ry   r=   r   None)r   r>   r   r=   r   r   )ro   r>   )r   r=   r   r   )ra   Optional[str]r   r   rO   zOptional[datetime]r   zOptional[bool]r   r   r   r   r   r   r   r=   )r   r>   r   r   )6__doc__
__future__r   r]   rn   rI   rr   r   r   typingr   r"   r#   r$   r%   r&   r'   r(   r)   r   r   r   r   r*   r+   r,   r-   r.   r/   r0   r1   r2   r3   r4   ro   joinr5   r6   r7   r8   r9   r:   __all__rE   r!   rR   rV   rd   r}   r   r   r   r    r   rQ   rF   rD   <module>r      s  
 #  	 	  '      > 	 #) & . * * ) 3 4 1  ''',,~xKhi '',,~xJab '',,~xJ_` '',,~xJbc '',,~xJhi '',,~xJ^_ 4=(&
-
&,8." $(!%"(,$()-#'n, !n, 	n,
 
n, &n, "n, 'n, !n, 
n,brF   