
    Oj3                       U d Z ddl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mZ ddlmZ ddlmZ ddlmZmZmZmZmZmZ dd	lmZ dd
lmZ erddlmZmZ ddl m!Z! ddl"m#Z# ddl$m%Z%m&Z&m'Z'm(Z)m*Z*m+Z+ ddl,m-Z.m/Z/m0Z0m1Z1  ejd                  e3      Z4dZ5e6e7d<    G d de      Z8ed        Z9 e
d       G d d             Z:dhZ;dee   de<fdZ= G d d      Z>y)u   StepRunner 베이스 클래스 — 모든 파이프라인 단계의 실행 프레임워크.

각 단계는 이 클래스를 상속하여 _execute()만 구현하면 됨.
게이트, 체크포인트, 무효화, Opik 추적을 자동 처리.
    N)contextmanager)	dataclass)datetimetimezone)Enum)Path)AnyDictListOptionalTupleTYPE_CHECKING)text)Session)CompletionReportCleanupReport)atomic_write_json)AppError)LockKey	LockOwnerOwnerVerdictREGISTRYjudge_ownerprocess_identity)containsget_manifest_dictget_depends_onget_all_downstream_recursive   MAX_RECOVERY_ATTEMPTSc                   (    e Zd ZdZdZdZdZdZdZdZ	y)	ResumeActionuS   ResumeDecision.action — claim 후 실행 분기 결정 (V2 patch I1, spec §2.3).skip
rerun_selfforce_explicitstale_running_recoveryblocknot_applicableN)
__name__
__module____qualname____doc__SKIP
RERUN_SELFFORCE_EXPLICITSTALE_RUNNING_RECOVERYBLOCKNOT_APPLICABLE     J/Users/manta/Documents/Projects/TheRoad-I1/backend/app/core/step_runner.pyr"   r"   =   s$    ]DJ%N5E%Nr4   r"   c              #     K   ddl m} t        | j                         xs i       }t	        |j                  dg       xs g       }|j                  d      } |d| j                   |||      5  d ddd       y# 1 sw Y   yxY ww)u   스텝 하나를 Opik trace 하나로 연다.

    설정이 꺼져 있으면 아무것도 안 한다(바이트 동일). 여는 데 실패해도
    본 스텝은 그대로 돈다 — `open_trace` 가 None 을 줄 뿐이다.
    r   )
open_tracetags	thread_idzstep:)namer8   metadatar9   N)app.modules.llm.opik_tracer7   dictbuild_opik_metadatalistpopgetstep_id)runnerr7   metar8   r9   s        r5   _step_trace_scoperE   G   s      6**,23D$*+D%I	V^^$%	
 	
 
 
s   A-B	/A=4	B	=BB	T)frozenc                       e Zd ZU dZeed<   eed<   dZee   ed<   dZ	ee   ed<   dZ
ee   ed<   dZee   ed<   dZee   ed	<   dZee   ed
<   y)ResumeDecisionuH  판정 결과 + 사유.

    V3 patch B1: STALE_RUNNING_RECOVERY 시 expected_started_at + expected_run_id 동반.
    claim SQL 의 atomic steal 조건에서 SQL cast 없이 (text 비교만) 검증.

    `origin` 은 caller (run()) 가 부수효과 디스패치에 사용:
      - None: first-run 또는 unknown-status (cleanup 없이 execute)
      - "artifact_missing": completed-mismatch (record_recovery + check_exhausted + cleanup)
      - "prior_state": non-completed prior status (force-like log + cleanup)
    후속 task (B11) 에서 CompletionReport.origin 과 통합.
    actionreasonNoriginprior_statusprior_run_idprior_updated_atexpected_started_atexpected_run_id)r)   r*   r+   r,   r"   __annotations__strrK   r   rL   rM   rN   rO   rP   r3   r4   r5   rH   rH   \   sl    
 K FHSM  #'L(3-&"&L(3-&&*hsm*)-#-%)OXc])r4   rH   grounding_modeproject_configreturnc                 "   ddl m}m} t        | xs i       } ||      } ||      }||j	                  dd       n||d<   t        j                  |ddd       }t        j                  |j                  d	            j                         dd
 S )u  project_config의 결정적 해시. P0-3 (Codex H2) — resume 안전성용.

    같은 config면 같은 hash → 체크포인트가 만들어진 시점의 config와
    현재 config가 다르면 mismatch 감지 가능.

    default fallback은 결정적이어야 한다. `default=str`을 쓰면 같은 클래스의
    다른 인스턴스(예: `<X object at 0x10ab40>`)가 매번 다른 hash를 만들어
    silent false-mismatch 회귀를 유발한다. 비-JSON 타입은 클래스명만 기록.

    ★**`grounding_mode` 는 값에 따라 접거나 뺀다** (`_VALUE_NORMALIZED_CONFIG_KEYS`).
    실측: 이 함수는 dict 전체를 해시하므로 그 키를 넣기만 해도 지문이 움직였고
    `shadow_plan` 으로 바꾸면 또 움직였다 — 검색도 산출도 안 바꾸는 관측 스위치가
    하류를 통째로 다시 굽게 만든다 (계약 §12).

    ★**그렇다고 key 째로 빼면 안 된다.** `v2` 까지 숨겨져 **legacy 완료 CP 가 v2
    진입에서 그대로 재사용된다.** 「revision/content hash 가 대신 잡는다」는 전제는
    아직 구현이 없다 (그 값을 만드는 production 코드 0건).

    Returns 16자 md5 prefix.
    r   )fingerprint_valueresolve_grounding_modeNrS   TFc                 4    dt        |       j                   dS )Nz<UNHASHABLE:>)typer)   )os    r5   <lambda>z%compute_config_hash.<locals>.<lambda>   s    La)9)9(:!<r4   )	sort_keysensure_asciidefaultutf-8   )app.core.grounding_moderW   rX   r=   r@   jsondumpshashlibmd5encode	hexdigest)rT   rW   rX   payload	effectivekept	canonicals          r5   compute_config_hashrn      s    * R>'R(G 'w/IY'D|$d+$( !

<	I ;;y''01;;=crBBr4   c                   D   e Zd ZdZ	 	 dedededededee   dee   fd	ZdfdZ	d
e
fdZ ej                  d      ZdZd
e
fdZd
ee   fdZd
ee   fdZedgd       Zded
dfdZdfdZdhdddee   de
d
dfdZded
eeeef      fdZedid       Zdjded
efdZdjded
efd Zd!eeef   d
efd"Z d#ed
efd$Z!d%Z"d&d'd(d(d(d)d)d*d+ed,e
d-e#ed.f   d/e$d0e$d1e$d2ed3ed
e
fd4Z%d'd(d(d(d)d)d5d+ed-e#ed.f   d/e$d0e$d1e$d2ed3ed
dfd6Z&dkd7ed
dfd8Z'd&ddd9d:e
d;ee   d<ee   d
e
fd=Z(d)d&d>d?ed@e
d
dfdAZ)dldBe$dCe$dDe$d
dfdEZ*dfdFZ+ded
dfdGZ,djded
eeef   fdHZ-dIeded
eeef   fdJZ.d
eeef   fdKZ/d
eeef   fdLZ0ddMdNed
eeef   fdOZ1dNed
eeef   fdPZ2dmdQZ3dndRZ4dfdSZ5d7ed
e$fdTZ6dUedVed
dfdWZ7dfdXZ8d
e$fdYZ9d
efdZZ:ed
e;fd[       Z<dmd\Z=d]eeef   d
ee   fd^Z>djded
eeef   fd_Z?d
efd`Z@	 	 dedaeeAe      dbeeeef      d
efdcZBed
efdd       ZCy)o
StepRunneru   파이프라인 단계 실행기.

    서브클래스는 _execute()만 구현.
    run()이 게이트 → 체크포인트 → 실행 → 저장 → 무효화를 자동 처리.
    NrB   
project_id
episode_iddbrT   opik_contextc                 \   t        |      st        d|       || _        || _        || _        || _        |xs i | _        t        |      | _        t        t        j                               | _        |xs i | _        ddlm} t!        |j"                        |z  dz  dz  |z  |z  | _        y )NzUnknown step: r   settingscheckpointsepisodes)_step_contains
ValueErrorrB   rq   rr   rs   rT   r   manifestrR   uuiduuid4run_idrt   app.core.configrw   r   projects_dir_cp_dir)selfrB   rq   rr   rs   rT   rt   rw   s           r5   __init__zStepRunner.__init__   s     g&~gY788$$,2)'2$**,'(.B 	-&&'*4()+568?@ 	r4   rU   c                    g }t        | j                        D ]  }| j                  |      }|s|j                  |       (|d   }|dv r2|dk(  rt	        |      xs i }|j                  dd      du r\|j                  d      }|r5| j                  j                  |      du rt        j                  d||       |j                  | d	       |j                  | d
| d        |r^|D cg c]3  }t	        |j                  d
      d         xs i j                  d|      5 }}t        dddj                  |       d      yc c}w )u  의존 단계 완료 확인. blocked이면 AppError.

        partial cascade contract (problems.md #11 + G4.6 Phase 3 fix iter 2):
          - dep status == 'partial' 이고 dep manifest 의 ``allow_partial_downstream``
            가 ``False`` 면 blocked. 명시 contract 없는 dep (default True) 는
            backward-compat 으로 통과.
          - 즉 emitter (dep step) 의 출력 contract 에 따라 cascade 제어.
          - **Operator override (manifest-driven)**: dep manifest 의
            ``partial_override_config_key`` 가 명시되어 있고 project_config 의
            동일 키 값이 truthy 면, partial 차단을 통과 (운영자가 명시적으로
            partial 진행을 결정한 경우). WARNING log 동시.
        status	completedr(   partialallow_partial_downstreamTFpartial_override_config_keyud   check_gate: dep %s partial — explicit operator override via project_config['%s']=True. proceeding.z(partial,strict)()r   labelzgate.blockedu   선행 단계 미완료: , i  codemessagestatus_codeN)r   rB   _get_step_runappendr   rA   rT   loggerwarningsplitr   join)	r   
blocked_bydepdep_run
dep_statusdep_metaoverride_keyb
dep_labelss	            r5   
check_gatezStepRunner.check_gate   sn    
!$,,/C((-G!!#& *J<<Y&,S17R<< :DAUJ#+<<0M#NL $(;(;(?(?(MQU(UR
 !%%-=&>?Qzl!451 04 ^hi^hYZ,QWWS\!_=CHHRST^hJi#3DIIj4I3JK  is   .8E	c                     ddl m}  ||       S )u   이 단계가 적용 가능한지. False면 not_applicable.

        기본 구현은 `app.core.applicability.resolve_applicability`를 호출.
        특수 로직이 필요한 서브클래스(SceneSplitStep, SceneVerifyStep 등)는 override.
        r   )resolve_applicability)app.core.applicabilityr   )r   r   s     r5   check_applicabilityzStepRunner.check_applicability  s     	A$T**r4   z2^manifest_(\d{8}_\d{6})(?:_[a-zA-Z0-9_-]+)?\.json$z.force_clearedc                 P    | j                   | j                  z  j                         S )um   force_cleared marker 존재 여부 — load_checkpoint이 archive/manifest 모두 무효 처리하는 분기.)r   _FORCE_CLEARED_MARKERexistsr   s    r5   _is_force_clearedzStepRunner._is_force_cleared  s     t999AACCr4   c                    | j                   dz  }| j                         }|j                         r}	 |j                  d      j	                         }|r:t        j                  |      }|r!t        j                  d| j                         y|S t        j                  d| j                         |r!t        j                  d| j                         y| j                         }|y	 t        j                  |j                  d            }	 t        ||       t        j                  d
| j                  |j                         |S # t        $ r+}t        j                  d| j                  |       Y d}~d}~ww xY w# t        $ r6}t        j                  d	| j                  |j                  |       Y d}~yd}~ww xY w# t        $ r7}t        j                  d| j                  |j                  |       Y d}~|S d}~ww xY w)u  체크포인트 manifest.json 로드.

        Primary: manifest.json 직접 로드.
        Fallback (Fix 4 / M5): manifest.json이 빈 파일/없음/손상이면 가장 최근
        archive(`manifest_TIMESTAMP[_runid].json`)로 자동 복원. archive가 살아있는
        한 force 직후 backend kill 같은 사고에서도 dispatcher가 cp를 회복한다.

        F4-1 (Claude IMPORTANT): clear_checkpoint 직후 backend kill로 save 누락된
        시나리오에선 archive(이전 completed)가 force 의도를 silent로 무효화할 수
        있다. clear_checkpoint이 작성한 `.force_cleared` marker를 확인하여 archive
        복구 + manifest.json 복원을 모두 차단 — resume이 stale로 인식하여 재실행.
        save_checkpoint(completed)가 정상 진행되면 marker 자동 제거.

        P2-1 (Codex): archive가 status='running' incremental snapshot이면 복구
        후보에서 제외 — `_find_latest_archive`가 completed/partial archive만 반환.

        archive → manifest 자동 복구는 atomic_write_json으로 기록하므로 재시도
        호출에서 fallback 비용이 1회만 발생.
        manifest.jsonra   encodinguh   Step %s: manifest restored but .force_cleared marker exists — treating as stale (force intent honored)Nz?Checkpoint manifest.json empty for %s. Trying archive fallback.z;Checkpoint load failed for %s: %s. Trying archive fallback.ug   Step %s: .force_cleared marker present — skipping archive fallback (force intent honored, will rerun)z5Archive fallback parse failed for %s (archive=%s): %szKStep %s: manifest.json missing/empty/corrupt. Auto-restored from archive %sz5Archive fallback write failed for %s (archive=%s): %s)r   r   r   	read_textstriprd   loadsr   r   rB   	Exception_find_latest_archiveerrorr:   r   )r   manifest_pathforce_cleared	text_datadataexcarchives          r5   load_checkpointzStepRunner.load_checkpoint  s   ( 6..0 !)33W3EKKM	::i0D$ G LL
  $KULL NN5
  ++-?		::g///ABD	mT2NN] ]  QLL ,  	LLG	 	   	LLG	  	sT   AE E  E $%E: 
7F< 	E7!E22E7:	F9,F44F9<	G<,G77G<c                    | j                   j                         syg }| j                   j                         D ]]  }|j                         s| j                  j                  |j                        }|s<|j                  |j                  d      |f       _ |sy|j                  d d       |D ]y  \  }}	 t        j                  |j                  d            }t!        |t"              r|j%                  d	      nd}||d
v r|c S t        j'                  d||j                         { y# t        $ r+}t        j                  d|j                  |       Y d}~d}~ww xY w)ux  최신 manifest_TIMESTAMP.json archive 파일 경로 반환 (없으면 None).

        파일명 패턴: manifest_YYYYMMDD_HHMMSS[_<runid8>].json
        TIMESTAMP 기준 내림차순으로 정렬 후, status='completed' 또는 'partial' 인
        archive만 후보로 채택 (P2-1 Codex fix).

        P2-1: incremental save_checkpoint(status='running') archive를 복원하면
            step_run.status='completed'로 남고 cp는 running snapshot이 되어
            sync 누락 + downstream silent skip 사고 발생. final 상태 archive만 복구.
            손상된 archive는 skip 후 다음 후보 시도.
        N   c                     | d   S )Nr   r3   )xs    r5   r]   z1StepRunner._find_latest_archive.<locals>.<lambda>  s    AaDr4   T)keyreversera   r   u&   Archive parse failed (skip): %s — %sr   r   r   r(   z*Archive skipped (status=%s, not final): %s)r   r   iterdiris_file_ARCHIVE_PATTERNmatchr:   r   groupsortrd   r   r   r   r   r   
isinstancer=   rA   debug)	r   archivespmtspathrj   r   r   s	            r5   r   zStepRunner._find_latest_archivex  s1    ||""$+-%%'A99;%%++AFF3AQ0 ( .$7 !HB**T^^W^%EF /9$.GW[[*TF~+U!ULL<fdii !   <dii 	s   8%D..	E"7!EE"c                    | j                         sy	 ddl}ddlm}  |j                         j	                  d      }| j
                  d| dz  }|j                         r6ddl}| j
                  d| d |j                         j                  dd  dz  }|j                  t        |       t        |             t        j                  d	|j                         y# t        $ r }t        j                  d
|       Y d}~yd}~ww xY w)u\   기존 manifest.json을 날짜시간 버전으로 복사 보관. 원본은 그대로 유지.Nr   )r   z%Y%m%d_%H%M%S	manifest_z.json_   zArchived checkpoint: %szCheckpoint archive failed: %s)r   shutilr   nowstrftimeparentr}   r~   hexcopy2rR   r   r   r:   r   r   )r   r   r   r   archive_pathr}   r   s          r5   _archive_manifestzStepRunner._archive_manifest  s     ##%	A)((9B(//IbT2GGL""$,33	"Qztzz|GWGWXZYZG[F\\a6bbLL]+S->?LL2L4E4EF 	ANN:C@@	As   CC 	D$C??Dr   c                 X   | j                   dz  }| j                  |d<   | j                  |d<   | j                         |d<   | j	                         |d<   d|vr| j
                  j                  dd      |d<   d|vrt        | j                        |d<   d	|vr| j                  | j                        |d	<   	 | j                  |       t        ||       |j                  d      }|dv r>| j                   | j                   z  }	 |j#                         r|j%                  d       yyy# t        $ rB}t        j                  d
| j                  |       t        d| j                   d|       |d}~ww xY w# t        $ r+}t        j'                  d| j                  |       Y d}~yd}~ww xY w)ug  체크포인트 저장 (원자적). 기존 manifest.json은 복사로 보관 후 덮어쓰기.

        F4-1: status='completed' 또는 'partial' 저장 시 .force_cleared marker 제거 —
        force 명령이 정상 진행 완료된 신호. 'running' incremental save에서는
        marker를 유지하여 backend kill 시 force 의도 보존.
        r   rB   r   resolved_model
updated_atschema_versionr   config_hashproject_config_snapshotz!Checkpoint save failed for %s: %szCheckpoint save failed for : Nr   r   T
missing_okz0Failed to clear .force_cleared marker for %s: %s)r   rB   r   _resolve_model_nowr|   rA   rn   rT   _mask_sensitive_keysr   r   r   r   r   RuntimeErrorr   r   unlinkr   )r   r   r   r   final_statusmarkers         r5   save_checkpointzStepRunner.save_checkpoint  s    6,,YX!%!4!4!6!YY[\ 4'%)]]%6%67G%KD!"$"5d6I6I"JD
 %D0.2.G.GH[H[.\D*+	]""=1mT2 xx)EE\\D$>$>>F==?MMTM2 # F  	]LL<dllCP!<T\\N"SERSY\\	]  FLL# s0   6D' "E5 '	E20=E--E25	F)>!F$$F)c                    | j                   dz  }|j                         r#| j                  |       |j                  d       	 | j                   j	                  dd       | j                   | j
                  z  j                  d       y# t        $ r+}t        j                  d| j                  |       Y d}~yd}~ww xY w)u`  체크포인트 삭제 (기존 파일은 날짜시간 버전으로 보관 후 삭제).

        F4-1: clear는 force trigger의 일부 — 이후 backend kill로 save_checkpoint가
        호출되지 않더라도 force 의도를 보존하기 위해 .force_cleared marker 작성.
        load_checkpoint이 이를 감지하면 archive/manifest 복구를 모두 차단하여
        step을 stale 상태로 유지 → 다음 실행 시 force 동작 보장.

        marker 작성은 best-effort — 실패해도 clear 자체는 진행. 다음 정상 save에서
        marker는 자동 제거된다.
        r   Tr   parentsexist_okr   z0Failed to write .force_cleared marker for %s: %sN)r   r   r   r   mkdirr   touchr   r   r   rB   )r   r   r   s      r5   clear_checkpointzStepRunner.clear_checkpoint  s     6!""=1  D 1	LLtd;\\D666==t=L 	NNBc 	s   AB 	B?!B::B?T)delete_checkpointstarget_step_idr   c                d   |xs | j                   }t        |      }ddlm} t	         ||            }| j                         }ddlm} |D ]u  }	| j                  j                  t        d      | j                  | j                  |	|dt        j                          d      }
|
j                         }|r4t!        |d         j#                  d      rt$        j'                  d|	|       |s|	|v rt$        j)                  d|	|||	       t+        |j,                        | j                  z  d	z  d
z  | j                  z  |	z  dz  }|j/                         s| j1                  |       |j3                  d       	 |j4                  j7                  dd       |j4                  | j8                  z  j;                  d       x |r=| j                  j?                          t$        j)                  dtA        |      |||       yy# t<        $ r"}t$        j'                  d|	|       Y d}~d}~ww xY w)u  downstream step을 stale 처리 + 체크포인트 삭제 (또는 stale 플래그만).

        Args:
            target_step_id: None이면 self.step_id 기준. editorial step이 타 step 체크포인트를
                덮어쓴 경우 해당 step_id를 전달하여 그의 downstream을 무효화.
            delete_checkpoints: True면 체크포인트 파일까지 아카이브+삭제.
                False면 DB step_run만 stale (체크포인트는 덮어쓴 값 그대로 유지, editorial 재실행 용).
        r   )get_consumes_downstreamrv   a  UPDATE step_run SET status = 'stale', updated_at = :now,   run_id = CASE WHEN status = 'running' THEN :revoked                 ELSE run_id END WHERE project_id = :pid AND episode_id = :eid AND step_id = :sid AND status NOT IN ('pending', 'stale') RETURNING run_idzinvalidated-)pideidsidr   revokedu   invalidate_downstream: 도는 중인 하류 %s 를 무효화했다 — 그 주행은 다음 안전 지점에서 멈춘다 (ref=%s)u   Preserved stale checkpoint for %s (consumed by %s) — 파일 보존, DB stale. drift 방지를 위해 %s force 완료 후 %s도 재실행 권장rx   ry   r   Tr   r   r   uC   Failed to mark %s as invalidated (archive 복원 위험 잔존): %sNz:Invalidated %d downstream steps from %s (delete_cp=%s): %s)!rB   r   app.core.step_catalogr   setr   r   rw   rs   executer   rq   rr   r}   r~   fetchonerR   
startswithr   r   infor   r   r   r   r   r   r   r   r   r   commitlen)r   r   r   ref
downstreamr   consumesr   rw   r   res_rowcp_pathr   s                 r5   invalidate_downstreamz StepRunner.invalidate_downstream  s    ,1#6
 	B3C89iik,C ''//$## tsdjjl^'DFGC <<>DDG//?RSVX[ "(?KKuS#s
  ../$//A#$&0137??CEHIKZ[  >>#**73NNdN3
,,TD,I 5567<udu7Ks B GGNNKKLJ&8*  %  125s s   ;AH	H/H**H/c                    | j                   j                  t        d      | j                  | j                  |d      j                         }|y|j                  |j                  |j                  |j                  |j                  |j                  |j                  |j                  |j                  |j                  |j                   |j"                  dS )u2  step_run row 의 모든 결정 필드를 dict 로 반환.

        Block B T0 (plan v2.1.3 / spec V5 S7): atomic claim / stale steal /
        verify path 가 run_id + started_at + recovery_count 의존 — tuple
        index 오해석 위험을 0으로 만든다. tuple fallback 절대 도입 X.
        a~  
            SELECT status, run_id, started_at, completed_count, applicable_count,
                   COALESCE(recovery_count, 0) AS recovery_count, updated_at,
                   owner_host, owner_pid, owner_boot_id, heartbeat_at,
                   cancel_requested_at
            FROM step_run
            WHERE project_id = :pid AND episode_id = :eid AND step_id = :sid
        r   r   r   N)r   r   
started_atcompleted_countapplicable_countrecovery_countr   
owner_host	owner_pidowner_boot_idheartbeat_atcancel_requested_at)rs   r   r   rq   rr   r  r   r   r  r  r  r  r   r  r  r  r  r  )r   rB   rows      r5   r   zStepRunner._get_step_runj  s     ggood $  oodoogNP QYPXPZ 	 ;jjjj.."22 # 4 4!00.. .. ..,,#&#:#:
 	
r4   c                 H    | j                   | j                  | j                  fS )u<   이 러너가 잡는 락의 신원 — 등록부 조회 키.)rq   rr   rB   r   s    r5   lock_keyzStepRunner.lock_key  s     $,,??r4   resumemodec           
         ddl }| j                  | j                        }| j                  |      }|r|j	                  |t        |j                  d      xs d      xs d|j                  d      rt        |j                  d            nd|j                  d      rt        |j                  d            nd      }|S )ui   ★claim 전 행의 상태를 판정에 싣는다 — 판정 본체는 `_evaluate_resume_decision_inner`.r   Nr    r   r   )rL   rM   rN   )dataclassesr   rB   _evaluate_resume_decision_innerreplacerR   rA   )r   r  r  existingdecisions        r5   _evaluate_resume_decisionz$StepRunner._evaluate_resume_decision  s    %%dll377="** h!7!=2>F$=E\\(=Sc(,,x"89Y]EM\\R^E_#hll<&@"Aei	 + lH
 r4   c                    |dk(  rt        t        j                  d      S | j                  | j                        }|st        t        j
                  dd      S |d   }|dk(  r| j                         }|t        t        j
                  d	d
      S | j                  |      }|rI| j                  j                  dd      rt        t        j                  |d      S | j                  |      S | j                         }|j                  rt        t        j                  d      S | j                  j                  dd      r)t        t        j                  d|j                   d      S |j                   dk(  r,t        t        j                  d|j                  dd  d      S |j                   dk(  r,t        t        j                  d|j                  dd  d      S t        t        j
                  d|j                   d
      S |dk(  r| j#                  |      S |dv r t        t        j
                  d| dd      S t        t        j
                  d|d      S )u  resume 판정자 — `(mode, step_run row, cp, verify_completion, project_config)`
        를 관찰해 `ResumeDecision` 반환.

        본 helper 는 자체 raise 안 함. `_safe_verify_completion` (B7) 이
        AppError(step.verify_crashed) raise 시 helper 내부에서 catch 안 함 —
        caller (run() / step_execution_service) 가 fail-fast 처리.

        분기 매핑 (current run() 동작 1:1 보존):
          - mode='force' → FORCE_EXPLICIT
          - first-run (no row) → RERUN_SELF (origin=None — cleanup 없음)
          - status='completed':
              - cp clean + verify pass → SKIP
              - cp None / cp mismatch / verify fail:
                  - strict_resume=True + cp not None → BLOCK
                  - 그 외 → RERUN_SELF (origin='artifact_missing')
          - status in {running, failed, partial, stale, pending}:
              - RERUN_SELF (origin='prior_state')
              - NOTE: running 자동 force 정책 변경은 후속 task (B4 STALE_RUNNING_RECOVERY).
          - 그 외 unknown status → RERUN_SELF (origin=None — current run() 의 fall-through 동작 보존)
            BLOCK 정책은 후속 task (B11) 에서.
        forcezuser-requested force)rI   rJ   zfirst-run (no step_run row)NrI   rJ   rK   r   r   z0checkpoint missing but step_run.status=completedartifact_missingstrict_resumeFzcp clean, verify passedzverify failed: contract_driftz#contract drift detected by verify: r   invariant_driftu2   invariant drift detected by verify (1차 정책): running)failedr   stalepending	cancelledzstatus=u    → force-like recoveryprior_statezunknown status: )rH   r"   r/   r   rB   r.   r   _check_cp_mismatchrT   rA   r1   _evaluate_contract_drift_safe_verify_completionis_completer-   missingrK   _evaluate_running_state)r   r  r"  r   cpcp_mismatchverify_reports          r5   r   z*StepRunner._evaluate_resume_decision_inner  s   , 7?!#22- 
 %%dll3!#..4  (#[ %%'B z%'22M-  11"5K &&**?EB)+11*.  44[AA !88:M((%',,4  ""&&>%'--,]-B-B,CD*  ##'77%'--=(00!457 ,  ##'88%'--**7*?*?*C)DF -  "#..()>)>(?@)  Y//99 KK!#.. (@A$  **%fZ0
 	
r4   r"  c           	         ddl m} |j                  d      }|j                  d      }t        j                  |      }t        || j                  |j                  |j                        \  }}|t        j                  u r t        t        j                  d| ||      S |t        j                  u rt        t        j                  d| d	
      S t!        |j                  d      |j                  d      |j                  d      f      }|r t        t        j                  d| dd	
      S |t        t        j                  d| d
      S 	 t#        j$                  t'        |      j)                  dd            }	|	j0                   |	j)                  t2        j4                        }	|j6                  }
t#        j8                  t2        j4                        |	z
  j;                         }||
k  r&t        t        j                  d|dd|
 d| d	
      S t        t        j                  d|dd|
 d| ||      S # t*        t,        t.        f$ r% t        t        j                  d|d| d
      cY S w xY w)uH  status='running' row 평가 — **죽음을 추측하지 않고 확인한다.**

        2026-08-26 개편. 이전에는 `started_at` 경과 시간만으로 판정해서,
        프로세스를 강제 종료하면 `step_running_timeout_seconds`(기본 3600초)가
        지날 때까지 아무것도 못 했다 (새벽 실측 손실 40분).

        지금은 `step_lock.judge_owner()` 가 소유자 신원으로 먼저 판정하고,
        신원을 못 믿는 자리에서만 시간으로 떨어진다:

        - `OwnerVerdict.DEAD`  → STALE_RUNNING_RECOVERY (바로 풀기)
        - `OwnerVerdict.ALIVE` → BLOCK (진짜로 일하는 중)
        - `OwnerVerdict.UNKNOWN` → 기존 경과 시간 경로 (신원 칸이 없는 구 행)

        ★`started_at` 이 없거나 깨진 것은 더 이상 무조건 BLOCK 이 아니다 —
         신원이 죽음을 말하면 그것이 우선한다. 시간은 신원이 침묵할 때만
         쓰는 마지막 그물이다.

        본 helper 는 자체 raise 안 함 — caller (run()) 가 BLOCK / STALE 디스패치 책임.
        r   rv   r  r   )r   lease_secondsallow_heartbeat_stealu   소유자가 죽었다 — )rI   rJ   rO   rP   u   소유자가 살아있다 — running_healthyr'  r  r  r  u*   소유자 생사를 확정 못 했다 — uu   . GET .../steps/locks 로 확인하고, 죽은 것이 맞으면 POST .../steps/locks/{step_id}/release 로 풀어라.z>running with started_at=None (manual investigation required); running_invalidZz+00:00z&running with started_at parse failed: z; )tzinfozrunning healthy (elapsed=z.0fzs < timeout=zs); zrunning stale (elapsed=zs >= timeout=)r   rw   rA   r   from_rowr   r  step_lock_lease_seconds!step_lock_heartbeat_steal_enabledr   DEADrH   r"   r0   ALIVEr1   anyr   fromisoformatrR   r!  r{   	TypeErrorAttributeErrorrA  r   utcstep_running_timeout_secondsr   total_seconds)r   r"  rw   started_at_raw
run_id_rawownerverdictverdict_reasonhas_identityr  timeout_secondselapseds               r5   r7  z"StepRunner._evaluate_running_state:  s   ( 	-!l3\\(+
 ""8,"-"::"*"L"L	#
 l'''!#::4^4DE$2 *	  l(((!#))77GH(   LL)LL&LL(
 
 !#))@@P QL L )  !!#))T%&( ) 	!//N#++C:J $#++8<<+@J"??<<-
:IIK_$!#))/}LHYY]%&( )  66)'#mOCTTX!"$ !/&
 	
7 I~6 	!#))<^<Nb%&( ) 	s   .H< <6I54I5mismatch_reasonc                     ddl m} d|v r;| j                  |v r-t        t        j
                  d| j                   d| dd      S t        t        j                  |d      S )	u  contract_drift 분류 (Block B B5, plan v2.1.3 / spec V5 §4.5).

        cp_mismatch (schema_version / config_hash) 의 BLOCK 정책. 한정 allowlist
        (`_LEGACY_SCHEMA_BUMP_ALLOWLIST = {entity_t2i}`) 만 schema_version mismatch
        시 RERUN_SELF 허용. config_hash mismatch 는 모든 step BLOCK (사용자 의도
        반영 — project_config 변경은 자동 재실행 시 의도와 다른 결과 위험).

        본 helper 는 자체 raise 안 함 — caller (run()) 가 BLOCK 디스패치 책임.
        r   )_LEGACY_SCHEMA_BUMP_ALLOWLISTzschema_version mismatchzlegacy schema bump allowlist: z (r   r*  r'  )app.core.step_manifestrX  rB   rH   r"   r.   r1   )r   rV  rX  s      r5   r3  z#StepRunner._evaluate_contract_drift  sv     	I %7||<<%'228 G+,A/ ,  %%"#
 	
r4   )r,  r   r   F)r,  r   r  require_ownerexpected_statusr  r  failed_counterror_messageresult_summaryr   r[  r\  .r  r  r]  r^  r_  c                   | j                         }	|dk(  r|	nd}
|dv r|	nd}g }d}|rBt        t        |            D cg c]  }d| 	 }}dj                  d |D              }d| d	}t	        d
| d      }| j
                  j                  |t        t        j                               | j                  | j                  | j                  || j                  | j                         ||||r|dd nd|r|dd nd|
||	dt        ||      D ci c]  \  }}||
 c}}      }|j!                         }| j
                  j#                          |duS c c}w c c}}w )u  step_run upsert. Block C B9 (plan v2.1.3 / spec V5 §4.7, AC-C7):
        require_owner=True 시 ON CONFLICT DO UPDATE 의 WHERE 절에
        `step_run.run_id = :rid` 추가 — 다른 worker 가 steal/overwrite 한 row 를
        조용히 덮어쓰지 못하도록 차단. row 갱신/INSERT 성공 시 True, owner mismatch
        시 False (RETURNING id 의 fetchone is not None 판정).

        `expected_status` — **이 전이가 어느 상태에서 출발해야 하는가.**

        ★상태 조건을 `'running'` 하나로 못박았다가 성공 주행을 망가뜨렸다
         (2026-08-26 Codex PR#4 리뷰 BLOCK-1·2). fan-out 스텝은 마지막 조각을
         끝내며 진행률을 적는데, 그 진행률이 DB 를 곧바로 'completed' 로
         바꿔 놓으면 뒤이은 마무리가 「내 줄이 아니다」로 읽혀 **성공이 예외로
         끝났다.** 마무리가 completed 를 적은 뒤 체크포인트 저장이 실패해
         failed 로 되돌리려 할 때도 같은 자리에서 막혀 **DB 만 성공으로
         남았다.**

         그래서 조건을 없애지 않고 **전이마다 출발 상태를 밝히게** 했다.
         밖에서 status 만 갈아 끼운 줄을 못 덮는다는 원래 성질은 그대로다.
        r,  Nr   r  expr   c              3   &   K   | ]	  }d |   yw):Nr3   ).0ks     r5   	<genexpr>z.StepRunner._update_step_run.<locals>.<genexpr>
  s     7YasGYs   z5WHERE step_run.run_id = :rid AND step_run.status IN (r   aQ  
            INSERT INTO step_run (id, project_id, episode_id, step_id, status, run_id,
                resolved_model, applicable_count, completed_count, failed_count,
                error_message, result_summary, started_at, completed_at, created_at, updated_at)
            VALUES (:id, :pid, :eid, :sid, :status, :rid, :model, :ac, :cc, :fc,
                :err, :rs, :sa, :ca, :now, :now)
            ON CONFLICT (project_id, episode_id, step_id) DO UPDATE SET
                status = :status, run_id = :rid, resolved_model = :model,
                applicable_count = :ac, completed_count = :cc, failed_count = :fc,
                error_message = :err, result_summary = :rs, updated_at = :now,
                started_at = COALESCE(:sa, step_run.started_at),
                completed_at = COALESCE(:ca, step_run.completed_at)
            z"
            RETURNING id
          )idr   r   r   r   ridmodelacccfcerrrssacar   )r   ranger  r   r   rs   r   rR   r}   r~   rq   rr   rB   r   r   zipr  r  )r   r   r[  r\  r  r  r]  r^  r_  r   r  completed_at	_exp_keysowner_clausei_insqlre  vresultr  s                        r5   _update_step_runzStepRunner._update_step_run  s   > iik"i/ST
$(GGsT  "	,1#o2F,GH,Gq3qc,GIH))7Y77CGuAN    N 	   djjl#DOODOO<<6$++((*2B/+8=$'d+9.$'tL	'
 !$I ?@ ?1q!t ?@	'
 	 oo$G I< As   EE)r\  r  r  r]  r^  r_  c          
          | j                  |d||||||      }|s+t        d| j                   d| j                   d| dd      y	)
u  owner-aware update + False 반환 시 step.owner_lost AppError raise.

        V2 patch 추가 d (Codex IMPORTANT #4, plan v2.1.3 §2247): claim 후 transition
        의 모든 path (success / partial / failed / exception handler / cleanup
        failure) 에서 owner check 실패 (False) 는 race lost 신호 — 자동 silent skip
        금지. 모든 transition path 에서 surface 의무.
        TrZ  step.owner_lostz owner check failed (run_id=uY   ) — 다른 worker 가 step_run row 를 갱신했거나 stale steal 당함. transition='u   ' 미적용.  r   N)r|  r   rB   r   )	r   r   r\  r  r  r]  r^  r_  updateds	            r5   _update_step_run_strictz"StepRunner._update_step_run_strict.  sv    $ ''++-%') ( 	
 &||n$@ N##)(,8    r4   rJ   c                 R    | j                  dd       | j                  d|d       y)u  NOT_APPLICABLE 처리 — claim 없이 DB step_run + cp 기록 (Block C, V2 patch I4).

        check_applicability=False 이거나 ResumeAction.NOT_APPLICABLE decision 시 호출.
        claim 안 했으므로 require_owner=False 로 단순 upsert.
        r(   F)r[  r   rJ   N)r|  r   )r   rJ   s     r5   _mark_not_applicablezStepRunner._mark_not_applicableU  s,     	.eD(8FKLr4   allow_stale_stealrO   rP   r  rO   rP   c                   t        t        j                               }| j                         }t	               \  }}}t        d      }	| j                  j                  |	|| j                  | j                  | j                  | j                  |||||||d      }
|
j                         }| j                  j                          |duS )uQ  atomic INSERT ON CONFLICT — race-free running claim (Block C C1).

        AC-C1, C6, C8, C9, C10, C11 (plan v2.1.3 / spec V5 §5.2):
        - first-run (row 없음) → INSERT 경로, RETURNING 한 row 로 success.
        - 기존 row + status != 'running' → DO UPDATE 경로 (failed/partial/stale/...).
        - STALE_RUNNING_RECOVERY (allow_stale_steal=True) → expected-match steal:
          step_run.started_at = :expected_started_at AND run_id = :expected_run_id.

        SQL cast 0 (V3 patch B1): timeout 판정은 _evaluate_running_state() 가 Python
        에서 수행. invalid text started_at row 가 있어도 query 정상 작동.
        fetchone() is not None 판정 (V3 patch I2): rowcount driver 차이 회피.

        Returns: True (claim 성공) / False (다른 worker 이미 잡음 또는 expected mismatch).
        a  
            INSERT INTO step_run (
                id, project_id, episode_id, step_id, run_id, status,
                started_at, created_at, updated_at,
                owner_host, owner_pid, owner_boot_id, heartbeat_at,
                cancel_requested_at
            )
            VALUES (
                :new_step_run_id, :pid, :eid, :sid, :run_id, 'running',
                :now, :now, :now,
                :owner_host, :owner_pid, :owner_boot_id, CURRENT_TIMESTAMP,
                NULL
            )
            ON CONFLICT (project_id, episode_id, step_id) DO UPDATE
              SET run_id = EXCLUDED.run_id,
                  status = 'running',
                  started_at = EXCLUDED.started_at,
                  updated_at = EXCLUDED.updated_at,
                  owner_host = EXCLUDED.owner_host,
                  owner_pid = EXCLUDED.owner_pid,
                  owner_boot_id = EXCLUDED.owner_boot_id,
                  heartbeat_at = EXCLUDED.heartbeat_at,
                  cancel_requested_at = NULL,
                  recovery_count = step_run.recovery_count + CASE
                    WHEN step_run.status = 'running' THEN 1 ELSE 0
                  END
              WHERE step_run.status != 'running'
                 OR (
                   :allow_stale_steal IS TRUE
                   AND step_run.status = 'running'
                   AND step_run.started_at = :expected_started_at
                   AND step_run.run_id = :expected_run_id
                 )
            RETURNING run_id
        )new_step_run_idr   r   r   r   r   r  r  r  r  rO   rP   N)rR   r}   r~   r   r   r   rs   r   rq   rr   rB   r   r  r  )r   r  rO   rP   r  r   r  r  r  ry  r{  r  s               r5   _try_claim_runningzStepRunner._try_claim_running^  s    * djjl+iik/?/A,
I}  " "H .????<<kk$"*!2#6.'
  oo$r4   wherestrictr  r  c                   ddl m} ddlm}  |       }	 |j	                  t        d      | j                  | j                  | j                  d      j                         } ||| j                  | j                  |      }	 |j                          |%|r"t        d| j                   d|xs d	 dd      y|j                   | j                   k7  r<t        d| j                   d| j                    d|j                    d|xs d	 dd      |j"                  dk7  r/t        d| j                   d|j"                   d|xs d	 dd      |j$                  r/t        d| j                   d|j$                   d|xs d	 dd      |j&                  r3t        d| j                   d|j)                          d|xs d	 dd      y# t        $ rt}|j                          |r&t        d| j                   d|xs d	 d
| dd      |t        j                  d| j                  ||       Y d}~|j                          yd}~ww xY w# |j                          w xY w)u/  안전 지점 검사 — **아직 내가 주인인가, 멈추라는 말이 왔나.**

        팬아웃 스텝이 한 단위(샷 하나 등)를 끝낼 때, 그리고 **돈을 쓰기 직전**에
        부른다. 세 가지를 막는다:

        1. **정지 요청** (`cancel_requested_at` 또는 `run_cancel_request`).
           운영자가 `POST .../cancel` 을 불렀다.
        2. **락 풀기** (`status != 'running'`). `release` 가 이 락을 풀기해
           `failed` 로 적었다. `run_id` 는 그대로일 수 있으므로 **상태도 봐야**
           한다 — 안 보면 풀린 뒤에도 옛 worker 가 그냥 통과한다.
        3. **좀비 쓰기** (`run_id` 불일치). 다른 worker 가 락을 가져갔다.
           소유권 검사는 지금까지 **마지막 상태 갱신에만** 있어서, 그 전에
           일어나는 DB commit·파일 저장은 못 막았다.

        Args:
            where: 어디서 불렀는지 (오류 문구에 들어간다).
            strict: DB 를 못 읽었을 때 **멈출지**. 기본 False.

        ★`strict` 를 가르는 이유 — 검사 조회가 실패했을 때 그냥 통과시키면
         DB 장애 동안 좀비가 계속 쓴다. 반대로 항상 멈추면 잠깐의 접속 장애로
         정상 주행이 죽는다. 그래서 **일반 안전 지점은 통과**(fail-open),
         **돈을 쓰거나 결과물을 확정하기 직전은 멈춤**(fail-closed)으로 가른다.

        ★**전용 세션**을 쓴다. 이 검사은 팬아웃 worker 스레드에서도 불리는데,
         SQLAlchemy 세션은 스레드 안전하지 않다. 게다가 PostgreSQL 은 실패한
         문장 하나로 트랜잭션이 abort 되므로, 작업용 세션으로 조회했다가
         실패하면 그 뒤 **그 worker 의 모든 쓰기가 죽는다.**

        Raises:
            AppError(step.cancelled): 정지 요청이 왔다.
            AppError(step.owner_lost): 락을 빼앗겼거나 풀렸다.
            AppError(step.gate_unreadable): strict 인데 검사을 못 읽었다.
        r   )SessionLocal)read_cancel_statez
                SELECT run_id, status, cancel_requested_at
                  FROM step_run
                 WHERE project_id = :pid AND episode_id = :eid AND step_id = :sid
            r  )r  zstep.gate_unreadableu$    소유권 검사을 못 읽었다 (u   안전 지점z): u$   . 돈을 쓰기 전이라 멈춘다.i  r   uD   checkpoint_gate 읽기 실패 (통과시킨다) step=%s where=%s: %sNr~  u    step_run 행이 없다 (u(   ) — 돈을 쓰기 전이라 멈춘다.r  u!    락을 빼앗겼다 (내 run_id=z, DB run_id=u   ) — u   에서 멈춘다.r,  u    락이 풀렸다 (status=step.cancelledu    정지 요청 (u    주행     — )app.core.databaser  app.core.run_controlr  r   r   rq   rr   rB   r  r   rollbackr   r   r   closer   r   r  	requesteddescribe)	r   r  r  r  r  sessionr  
run_cancelr   s	            r5   checkpoint_gatezStepRunner.checkpoint_gate  s   D 	3:."	//$ ( #
 ||	 xz  +$//&J* MMO;*<<. )!4_55]_ !$  ::$&||n$Edkk] S!!$F53KO2LL]_    ::"&||n$>szzl&/00AC    ""%||n$4S5L5L4MV/00AC    %||nHZ-@-@-B,C5/00AC     E  	/<<. )!4_5S >== !$  NNVeS MMO%	$ MMOs+   A+G 	IAH>)I >II Ir   totalr-  c                 .    | j                  d|||       y)u  fan-out 진행률 업데이트 (running 상태에서). claim 후 owner-aware (AC-C7).

        Block C review I4: owner_lost 시 silent skip 대신 strict 로 즉시 raise →
        fan-out 중단 (partial commit orphan 차단). caller 는 outer except 에서 surface.
        r,  )r  r  r]  N)r  )r   r   r  r-  s       r5   update_progresszStepRunner.update_progress4  s#     	$$%"	 	% 	
r4   c                 4    ddl m}  || j                         y)u;  `grounding_mode` 오타를 **첫 지출 전에** 세운다 (fail-closed).

        ``resolve_grounding_mode`` 가 enum 밖 값에 ``AppError`` 를 던진다.
        여기서 부르는 것은 「값을 쓰려고」가 아니라 **「틀린 값이면 아무것도
        시작하지 않으려고」**다.
        r   )rX   N)rc   rX   rT   )r   rX   s     r5   _preflight_grounding_modez$StepRunner._preflight_grounding_modeG  s     	Ct223r4   c                      y)uO  mode 별 실행 가능성 preflight — run() 최선두(모든 mutating 단계
        이전)에서 호출. subclass 가 특정 mode 를 거부해야 할 때 override 해
        raise (예: 표적 씬 슬라이스+force 병용 금지 — 비표적 자산 무효화
        차단, Codex 슬라이스 리뷰 BLOCKING-1). 기본=no-op.Nr3   r   r  s     r5   validate_modezStepRunner.validate_modeR  s    r4   c                 "   | j                          | j                  |       | j                          | j                         s| j	                  d       ddiS | j                  |      }|j                  t        j                  k(  rdddS |j                  t        j                  k(  r|j                  dk(  r| j                   d	|j                   d
}n|j                  dk(  r| j                   d|j                   d}nc|j                  dk(  r| j                   d|j                   d}n7|j                  dv r| j                   d|j                   d}n|j                  }t        d|d      t                t        j                   | j"                  | j$                         	 | j'                  ||      t        j(                  | j"                  | j$                         S # t        j(                  | j"                  | j$                         w xY w)u  단계 실행. mode: resume | force.

        Returns: {"status": "completed|skipped|not_applicable", "result": ...}

        Block C C2 (plan v2.1.3 / spec V5 §S2 + §S6): 4-phase 흐름 —
          (1) gate / applicability (non-mutating)
          (2) ResumeDecision 평가 (non-mutating, mutating side-effects only after claim)
          (3) atomic claim (RERUN_SELF / FORCE_EXPLICIT / STALE_RUNNING_RECOVERY 만)
          (4) origin 별 부수효과 (record_recovery / [RECOVERY] log) + _execute_*
              dispatch — _execute_rerun_self / _execute_force 가 _execute_and_finalize
              로 finalize 까지 위임.
        zcheck_applicability=False)rJ   r   r(   skippedzalready completedr  r)  u%    체크포인트가 stale 입니다 (uC   ). strict_resume=True 이므로 force 모드로 재실행하세요.r*  u    contract drift detected — u.   . force 재실행 또는 수동 진단 필요.r+  u    invariant drift detected — )r>  r?  z step is running (u}   ). 다른 worker 가 진행 중일 가능성 — auto-force 차단됨. 확인 후 mode='force' 로 명시 재실행하세요.zstep.resume_invalidr  r   )r  r  r   r   r  r$  rI   r"   r-   r1   rK   rB   rJ   r   r   LOCK_REGISTRYregisterr  r   _claim_and_run
unregister)r   r  r#  _msgs       r5   runzStepRunner.runX  s   0 	&&( 	4  	 '')%%-H%I.// 11$7??l///'3FGG??l000/1||n % ( ):: 
 $44||n$A(//AR SC C  $55||n$B8??BS TC C  $JJ||n$6x6G HM M   *  	 	t}}dkk:	A&&x6$$T]]DKK@M$$T]]DKK@s   &G" ",Hr#  c           	         |j                   t        j                  k(  }| j                  ||r|j                  nd|r|j
                  nd      }|s5t        d| j                   d|j                   j                   d| dd      |rAt        j                  d	| j                  |j                  |j
                  | j                         |j                   t        j                  k(  r|j                  d
k(  r| j                         }|Wddlm} |j#                  d      }	  ||| j$                        }t        j                  d| j                  |j*                  |       | j-                  |j*                        }
t        j                  d| j                  |j*                  |
       | j/                          d}n |j                   t        j                  k(  r^|j                  dk(  rOt        j                  d| j                  |j0                  |j2                  |j4                  |j*                         d}n|j                   t        j                  k(  ri|j                  dk(  rZ| j-                  |j*                        }
t        j                  d| j                  |
|j*                         | j/                          d}n|j                   t        j6                  k(  rd}|dk(  r=|j                   t        j                  k(  r| j9                         S | j;                         S | j9                         S # t&        $ r}	dt)        |	      i}Y d}	~	d}	~	ww xY w)ue  claim 이후의 실행 — 등록부 등록이 끝난 상태에서만 불린다.

        `run()` 에서 떼어낸 이유는 하나다: 락 등록부의 해제를 `finally` 로
        보장하려면 claim 이후 전체가 한 블록이어야 하는데, 그러면 100줄 넘는
        본문이 통째로 한 칸 들어가 diff 가 안 읽힌다.
        Nr  zstep.already_runninguX    이 이미 다른 worker 에서 실행 중 (또는 stale read 후 갱신됨). decision=z, is_stale_steal=.r  r   zj[STALE_RUNNING] step=%s atomic steal succeeded (expected_started_at=%s, expected_run_id=%s, new_run_id=%s)r(  r   )_diff_project_configr   _diff_errorz3Step %s stale detected: %s. Project config diff: %sz%[RECOVERY] step=%s reason=%s cycle=%dr&  r1  uY   [RECOVERY] step=%s prior_status=%s prior_run=%s prior_updated=%s reason=%s → force-liker*  z8[RECOVERY] step=%s allowlist contract_drift cycle=%d: %s)rI   r"   r0   r  rO   rP   r   rB   valuer   r   r   r.   rK   r   r   r  rA   rT   r   rR   rJ   _record_recovery_check_recovery_exhaustedrL   rM   rN   r/   _execute_rerun_self_execute_force)r   r#  r  is_stale_stealclaim_okr8  r  cp_snapshotconfig_diffdiff_exc	new_counts              r5   r  zStepRunner._claim_and_run  s    #//\-P-PP**,@N < <TX8FH44D + 

 +||n %CCK??CXCXBY Z&&4%5Q8    NNN,,(( ??l555(//M_:_ %%'B~G ff%>?A"6{DDWDW"XK ILL(//; --hoo>INN7hooy **,D__ 7 77HOO}<\ NNkh33X5J5JHLeLegogvgv D__ 7 77HOOO_<_ --hoo>INNJi **,D__ ; ;;D 7? ,"9"99//11&&(( ''))U ! A#0#h-"@KAs   L* *	M3MMc                 &    | j                  d      S )u  auto-recovery 경로 — 자기 step 만 rerun (Block B B12, AC-B6).

        force 와의 차이:
        - cleanup_artifacts() 호출 X
        - invalidate_downstream() 호출 X
        - clear_checkpoint() 호출 X
        - downstream cp 보존 (auto-recovery 가 downstream 영향 X)

        finalize 는 _execute_and_finalize() 공통 helper. _execute(mode='resume')
        호출 — force-like 격상 X.
        r  execute_mode)_execute_and_finalizer   s    r5   r  zStepRunner._execute_rerun_self8  s     ))x)@@r4   c                 t   	 | j                         }|j                  s|j                  rAt
        j                  d| j                  |j                  |j                  |j                         | j                          | j                          | j                  d	
      S # t        $ r}	 | j                  dd|        nR# t        $ rF}|j                  dk7  r t
        j                  d| j                  |j                  |       Y d}~nd}~ww xY wt
        j                  d| j                  |        d}~ww xY w)uR  사용자 명시 force — cleanup + invalidate_downstream cascade (Block B B12, AC-B6).

        AC-B6 의 대척 경로 — force 는 의도적으로 downstream invalidate. cleanup_
        artifacts 예외 시 status='failed' 영속 + raise (silent fallback 금지).
        finalize 는 _execute_and_finalize() 공통 helper.
        r-  zcleanup_artifacts crashed: r^  r~  z2Step %s force-cleanup owner_lost: %s (original=%s)Nz%Step %s cleanup_artifacts crashed: %sz-[CLEANUP] step=%s rows=%d files=%d targets=%sr&  r  )cleanup_artifactsr   r  r   r   r   r   rB   r   r   deleted_db_rowsdeleted_filestargetsr  r   r  )r   cleanup_reportr   	owner_excs       r5   r  zStepRunner._execute_forceF  s%   	!335N, ))^-I-INN?..,,&& 	""$))w)??A  	,,$?u"E -   >>%66HLL)"3"3S  LL7c )	s;   B 	D7'B>=D2>	D<DD2D%D22D7r  r  c                T    ddl m}  j                  d       t        j	                  d j
                   j                   j                                 | j                                ddl	m
} d fd}t               5   || j                  |      cddd       S # 1 sw Y   yxY w)	u   rerun_self / force 두 경로 공통 finalize (Block B B12, plan v2.1.3 §S2).

        7 단계 (V5 S2 + P3):
          (1) _update_step_run("running") + opik context 설정
          (2) _execute(mode=execute_mode) + _last_execute_result attribute
          (3) final_status 계산 (failed/partial/completed)
          (4) exit verify (final_status='completed' 시) — `step.verify_crashed`
              AppError raise 시 status='failed' 영속 후 re-raise (running leak 차단).
              그 외 verify is_complete=False → final_status='partial' 격상.
          (5) _update_step_run(final_status, ...)
          (6) save_checkpoint({"status": final_status, **result}) + recovery
              counter reset (final_status='completed' 시)
          (7) modifies_checkpoints cascade policy (1-pass default)

        예외 처리: _execute / finalize 단계의 unexpected exception → status='failed'
        영속 + re-raise. status='running' leak 차단 (다음 resume 에서 stale 분기 회피).
        r   set_opik_contextr,  z%Step %s started (run_id=%s, model=%s))run_with_stop_checkNc                  ,     j                  dd       y )Nu   유료 이미지 호출 직전Tr  )r  r   s   r5   _paid_call_gatez9StepRunner._execute_and_finalize.<locals>._paid_call_gate  s      'GPT Ur4   rU   N)app.modules.llm.llm_clientr  r  r   r  rB   r   r   r>   app.core.image_call_budgetr  rE   _execute_and_finalize_inner)r   r  r  r  r  s   `    r5   r  z StepRunner._execute_and_finalizeq  s    $ 	@ 	$$Y/3LL$++t':':'<	
 	1134" 	C	V t$&00 %$$s    BB'c                    ddl m} d}	 | j                  |      }|| _        |j	                  dd      }|j	                  dd      }|j	                  dd      }|dk(  rd	n|dkD  rd
nd}|d	k(  rJ	 | j                         }	|	j                  s-t        j                  d| j                  |	j                          d
}| j                  ||||t#        j$                  |j'                         D ci c]  \  }}|dk7  st)        |t*              r||! c}}dt,              dd        | j/                  d|i|       |d	k(  r| j1                          | j2                  j	                  d      xs g }| j2                  j	                  dd      }|dv r@|r>|r|D ]  }| j5                  |d        n!t        j7                  d | j                  |       t        j9                  d!| j                  ||||       ||d" |d       S # t        $ r}
t        j                  d| j                  t        |
dd      |
j                         	 | j                  d|||dt        |
dd       d|
j                          d} # t        $ rG}|j                  dk7  r t        j                  d| j                  |j                         Y d}~d} d}~ww xY wd}
~
ww xY wc c}}w # t:        $ rr}ddl}t)        |t              rt        |dd#      d$k(  rt        j                  d%| j                  |j                         	 | j                  d&|j                  '        # t        $ rE}|j                  dk7  r t        j                  d(| j                  |j                         Y d}~ d}~ww xY w|s{	 | j                  dt-        |      | j>                  )       nR# t        $ rF}|j                  dk7  r t        j                  d*| j                  |j                  |       Y d}~nd}~ww xY wt        j                  d+| j                  ||jA                                 d}~ww xY w#  |d       w xY w),uG  `_execute_and_finalize` 의 본문 — 옮기기만 했다(내용 변경 0).

        스텝 경계에 Opik trace 를 겹치려면 본문을 with 로 감싸야 하는데,
        그러면 100줄 넘는 블록이 통째로 한 칸 들어가 diff 가 안 읽힌다.
        그래서 본문을 메서드로 떼어 내고 바깥에서 감쌌다.

        ★`set_opik_context` local import 는 함께 옮긴다 — finally 절이 그것을
        쓰는데, 바깥 함수의 지역 이름은 여기서 안 보인다(호출 시점까지
        안 드러나는 결함이다).
        r   r  F)r  r  r   r  r]  r   r   r-  u4   [VERIFY-EXIT] step=%s exit verify crashed: %s — %sr   ?zexit verify crashed: r  )r  r  r]  r^  r~  z3[VERIFY-EXIT] step=%s owner_lost on failed-mark: %sNTu,   [VERIFY-EXIT] step=%s failed → partial: %sr   )r_   r`   rg  )r  r  r]  r_  r   modifies_checkpointsinvalidate_downstream_on_edit)r   r   )r   r   z>Step %s: editorial cascade skipped (policy=1-pass, targets=%s)z'Step %s %s (completed=%d/%d, failed=%d))r   r{  r  r  u   [CANCELLED] step=%s — %sr0  r  u(   Step %s cancel 기록 중 owner_lost: %s)r^  r\  z3Step %s exception path owner_lost: %s (original=%s)z,Step %s _execute_and_finalize crashed: %s
%s)!r  r  _execute_last_execute_resultrA   r4  r   r   r   rB   getattrr   r  r   r   r5  r6  rd   re   itemsr   bytesrR   r   _reset_recovery_counterr|   r  r   r  r   	traceback_TERMINAL_ROLLBACK_FROM
format_exc)r   r  r  failed_recordedr{  r   r  r-  r   exit_report
verify_excr  re  rz  mods
cascade_on
target_sidr   r  s                      r5   r  z&StepRunner._execute_and_finalize_inner  sh    	@  [	#]]]5F(.D%

#4a8IJJ115EZZ2F  &{#,q=ih  {*"&">">"@K< #..NNFk&9&9 $-L (( )!&##zz&,llnandaVJWXZ_L`QTna!&  4  ) 	   (L!CF!CD{*,,. ==$$%;<BD**+JERJ77D&*
22+5/4 3  '+ LLXd
 KK9lIuf +f=| T"O   LLN
FC8"**	44$,5-2)/"7#*:vs#C"DE*J\J\I]!_ 5 	" '+O $ $>>->>!Q LL)*;*;  '+O)V b@  :	 3)C,0@@;T\\3;;W
00#3;; 1      ~~)::NNBi&7&7   # 00 C(,(D(D 1     ~~)::NNMi&7&7  LL?c9#7#7#9 u:	x T"s   A$K /G> ?A*K )K7KKC(K >	K7K?1I30K3	K<;J>7K>KKK	K 
QAQ.MQ	N;NQNQ 'OQ	P<PQP4QQQ 
Qc                 &    ddl m}  |dg di       S )u  결과물 무결성 검증. 베이스는 항상 complete — override 안 한 step은 검증 없음.

        Resume entry + execute exit 양방향에서 호출됨. is_complete=False 반환 시:
          - entry: mode='force' 격상 (auto-rerun)
          - exit: status='partial' 마킹

        Override 시 자기 step 결과물(DB row + 파일 stat)만 검증한다. cross-step
        dep verify는 후속 spec.
        r   r   Tclean)r5  r6  severityr;   )app.core.integrity_reportr   )r   r   s     r5   verify_completionzStepRunner.verify_completion\  s     	?D"wY[\\r4   c                 &    ddl m}  |ddg g       S )u@  force/auto-rerun 시 stale artifact 정리. 베이스는 noop.

        호출 site: run() force 분기 (Task 12 구현 예정)에서 _execute() 직전 1회 호출.
        resume entry verify가 is_complete=False 반환 → mode='force' 격상 → cleanup → execute 흐름.

        Override 작성 기준 (매우 보수적):
        - 전체 force = delete-and-rebuild가 항상 안전한 step만
        - 부분 재생성을 endpoint로 지원하는 step은 절대 override 금지
        - 이미지 step은 거의 모두 default noop 유지 (사용자 caveat)
        r   )r   )r  r  r  r  )r  r   )r   r   s     r5   r  zStepRunner.cleanup_artifactsi  s     	<QaUWXXr4   c                     | j                         }|t        k\  r6t        d| j                   dt         d| d| j	                          dd      y)	u  누적 recovery_count >= MAX_RECOVERY_ATTEMPTS면 step.recovery_exhausted raise.

        Audit A4 (DRY): 기존 dual gate(resume 분기 record 전 + force 분기 진입 직전)를
        한 곳으로 통합. 호출 site는 `_record_recovery` **직후** — post-increment된
        새 카운트로 즉시 검사하여 force 분기 진입 자체를 차단한다.

        message에는 module 상수 MAX_RECOVERY_ATTEMPTS를 명시 — silent recovery로
        묻히지 않게 사용자에게 한계값을 노출하고, 회귀 테스트의 "3" 포함 검사도
        보존한다(누적 카운트와 한계가 모두 표시됨).
        zstep.recovery_exhaustedu    auto-recovery 한계(u   회) 도달 — 누적 u   회 실패. 마지막 사유: u   . 수동 진단 필요.r  r   N)_get_recovery_countr    r   rB   _get_last_recovery_reason)r   counts     r5   r  z$StepRunner._check_recovery_exhaustedy  sn     ((*)).||n$:;P:Q R#W$B4CaCaCcBd e,,    *r4   c           
      H   | j                         }| j                  j                  t        d      |dd | j                  | j
                  | j                  || j                  d      }| j                  j                          | j                  |d       | j                         S )u)  recovery_count++ + last_recovery_reason 적재. 새 카운터값 반환.

        호출 site: entry verify 실패 시. mode='force' 격상 직전.

        ★한 줄도 안 바뀌면 **내 토큰이 이미 버려진 것**이다 (2026-08-26 Codex
         재리뷰 BLOCK-4). 종전에는 `run_id` 조건만 걸고 결과를 안 봐서, 토큰이
         버려진 뒤에도 조용히 0줄로 끝나고 바로 아래 `_get_recovery_count` 가
         **새 주인의 줄**을 읽었다. 남의 카운터로 내 한계를 판단한 셈이다.
        a4  
            UPDATE step_run SET
                recovery_count = recovery_count + 1,
                last_recovery_reason = :reason,
                updated_at = :now
            WHERE project_id = :pid AND episode_id = :eid AND step_id = :sid
              AND run_id = :rid AND status = 'running'
        Nrg  )rJ   r   r   r   r   ri  u   recovery_count 증가)r   rs   r   r   rq   rr   rB   r   r  _require_owned_rowr  )r   rJ   r   r	  s       r5   r  zStepRunner._record_recovery  s     iikggood $   toodll3kk#	$ 	 	%<=''))r4   r{  whatc           	      ~    t        |dd      }|dk(  r+t        d| j                   d| d| j                   dd	      y)
u   방금 UPDATE 가 **내 줄**을 건드렸는지 확인한다.

        0줄이면 내 토큰이 버림됐다는 뜻 — 조용히 넘어가면 그 뒤 판단이 전부
        남의 줄을 근거로 이뤄진다. 여기서 멈춘다.
        rowcountNr   r~   u    실패 (run_id=uw   ) — 이 주행의 토큰이 버림됐다. 다른 주행이 이 스텝을 가져갔거나 정지·풀기로 놓았다.r  r   )r  r   rB   r   )r   r{  r  r  s       r5   r  zStepRunner._require_owned_row  sX     6:t4q=&||nAdV+;DKK= I3 3    r4   c           	      "   | j                         }| j                  j                  t        d      | j                  | j
                  | j                  || j                  d      }| j                  j                          | j                  |d       y)u  status='completed' 정상 진행 시 카운터 reset (다음 sweep 깨끗한 상태).

        ★부르는 자리가 `_update_step_run_strict(final_status)` **뒤**이고
         `final_status == "completed"` 일 때뿐이다. 그러니 이 시점의 DB
         status 는 'completed' 여야 한다 — 아니면 그 사이에 누가 끼어든 것이다
         (2026-08-26 Codex 2차 재리뷰).
        z
            UPDATE step_run SET recovery_count = 0, last_recovery_reason = NULL, updated_at = :now
            WHERE project_id = :pid AND episode_id = :eid AND step_id = :sid
              AND run_id = :rid AND status = 'completed'
        )r   r   r   r   ri  u   recovery_count 초기화N)
r   rs   r   r   rq   rr   rB   r   r  r  )r   r   r	  s      r5   r  z"StepRunner._reset_recovery_counter  ss     iikggood $  oodooll3t{{D	E 	%?@r4   c                     | j                   j                  t        d      | j                  | j                  | j
                  | j                  d      j                         }|r|d   S dS )NzxSELECT recovery_count FROM step_run WHERE project_id = :pid AND episode_id = :eid AND step_id = :sid   AND run_id = :ridr   r   r   ri  r   rs   r   r   rq   rr   rB   r   r  r   r  s     r5   r  zStepRunner._get_recovery_count  sb     ggood"
 ??4??<<5	6
 7?hj 	 s1v#!#r4   c                     | j                   j                  t        d      | j                  | j                  | j
                  | j                  d      j                         }|r
|d   r|d   S dS )Nz~SELECT last_recovery_reason FROM step_run WHERE project_id = :pid AND episode_id = :eid AND step_id = :sid   AND run_id = :ridr  r   u   (없음)r  r  s     r5   r  z$StepRunner._get_last_recovery_reason  sf    ggood"
 ??4??<<5	6
 7?hj 	 #a&A9j9r4   c                     d}| si S i }| j                         D ]M  \  }t        fd|D              rd|<    t        |t              rt        j                  |      |<   I||<   O |S )u   project_config 저장 전 민감 키 마스킹. api_key/password/secret/token/credential 패턴.

        nested dict는 재귀. None → 빈 dict.
        D2 fix — save_checkpoint이 project_config_snapshot 기록 시 사용.
        )api_keypasswordsecrettoken
credentialc              3   T   K   | ]  }|t              j                         v  ! y wN)rR   lower)rd  sre  s     r5   rf  z2StepRunner._mask_sensitive_keys.<locals>.<genexpr>  s!     :	11A&	s   %(z***)r  rG  r   r=   rp   r   )config	SENSITIVEmaskedrz  re  s       @r5   r   zStepRunner._mask_sensitive_keys  sq     M	ILLNDAq:	::!q	At$&;;A>q	q	 # r4   c                    ddl m} 	 | j                         S # t        $ rz}t	        |dd      dk(  r t
        j                  d| j                  |j                  t	        |dd              |dd	|j                   gd
dt	        |dd      id      cY d}~S d}~wt        $ rO}t
        j                  d| j                  |d       t        ddt        |      j                   d|       |d}~ww xY w)u  verify_completion 실행 + 예외 분류 (Block B B7, plan v2.1.3 / spec V5 §4.3).

        - AppError(code='step.verify_crashed') → propagate (B8 verify_completion 이
          card recompute 등에서 미리 격상한 케이스 — 이중 wrap 금지).
        - 그 외 AppError (loader contract / step.contract_violation 등) →
          CompletionReport(origin='contract_drift', is_complete=False). caller
          (run() / ResumeDecision helper) 가 BLOCK 처리 (B11 후속).
        - 그 외 Exception (KeyError / ValueError / AttributeError 등) →
          AppError(step.verify_crashed) raise (자동 force 금지, fail-fast).
          silent recovery 사고 패턴 차단.
        r   r  r   Nzstep.verify_crashedz/verify_completion AppError for %s: %s (code=%s)r  Fz
contract: r6  app_error_coder*  )r5  r6  r  r;   rK   z$verify_completion crashed for %s: %sT)exc_infozverify_completion crashed: r   )r   r   )r  r   r  r   r  r   r   rB   r   r   r   r[   r)   )r   r   r   s      r5   r4  z"StepRunner._safe_verify_completion  s     	?	))++ 	sFD)-BBNNAckk73+D $!%ckk]34"*GC,FG'   	LL6cTX   *5d3i6H6H5IC5Q 	s(    	C1A/BC1C1"A
C,,C1r8  c                 4   | j                   j                  dd      }|j                  d      }|j                  d      }|||k7  rd| d| S |It        | dd      }t        |      r
 |       }d}nt	        | j
                        }d	}||k7  rd
| d| d| S y)u|  cp의 schema_version + config_hash와 현재 manifest/project_config 비교.

        Returns: mismatch reason str (None이면 일치).
        run() resume 분기의 P0-3 mismatch 검증을 헬퍼로 이관.
        legacy 체크포인트(cp_schema/cp_hash 누락)는 None 반환 — 호환 보장.

        config_hash 비교 (Codex iter root cause):
          save_checkpoint (line 349-350) 는 step 이 _execute() 에서 명시 반환한
          'config_hash' 가 있으면 그것을 보존하고, 없으면 compute_config_hash(
          project_config) 로 fallback. 따라서 mismatch check 도 step 의 자체
          hash method (`_config_hash()`) 가 있으면 그것으로 계산해야 일관.
          이전 결함: _execute 가 SHA256(prompt_version+model+...) hash 저장 →
          check 는 MD5(project_config) hash 비교 → 영구 false-positive mismatch.
        r   r   r   Nu)   schema_version mismatch: 체크포인트=u	   , 현재=_config_hashz
step-localrT   zconfig_hash mismatch (z): cp=z
, current=)r|   rA   r  callablern   rT   )r   r8  current_schema	cp_schemacp_hashlocal_hash_fncurrent_hash	hash_kinds           r5   r2  zStepRunner._check_cp_mismatch"  s     **+;Q?FF+,	&&' Y.%@>ykSaRbcc $D.$?M&  -(	243F3FG,	,&,YK 8!*\N<
 r4   c                 4    t        d| j                   d      )u   서브클래스에서 구현. 실제 LLM 호출 수행.

        Returns: {
            "completed_count": N,
            "applicable_count": N,
            "failed_count": N,
            "data": {...},  # 단계별 결과 데이터
        }
        zStep z must implement _execute())NotImplementedErrorrB   r  s     r5   r  zStepRunner._executeP  s     "E$,,7Q"RSSr4   c                 H    ddl m}  || j                  | j                        S )u&   현재 단계의 모델 별칭 반환.r   )r   )r  r   rB   rT   )r   r   s     r5   r   zStepRunner._resolve_model^  s    =dllD,?,?@@r4   
extra_tagsextra_metadatac                    | j                   | j                  d}| j                  }|j                  d      r?|j                  dd       d|j                  dd       d| j                   |d<   |d   |d<   | j                  g}|j                  d      r|d   |d	<   |j                  |d          |j                  d      r|d   |d
<   |j                  |d          |r|j                  |       ||d<   |r"|j                         D ]  \  }}||vs|||<    	 ddlm	} t        |dd      rddlm}	m}
 |j                  dd       |j                  dd       |j                  d	d       |j                  d
d        |
|j                  dd      |j                  dd      | j                  xs d      |d<   |j                  d      r|d   |d<   |j                  d      r|d   |d<   |j                  d      r|d   |d<    |	| j                        }|xs g D ]  }d| }||vs|j                  |        ||d<   |S # t        $ r!}t         j#                  d|       Y d}~|S d}~ww xY w)u  Opik metadata 구성 — thread 그룹핑 + 프로젝트/에피소드 정보.

        extra_tags: Opik 태그에 추가할 낮은 카디널리티 문자열만 사용 (예: "retry", "gpt_fallback").
                    동적 ID(scene_index, shot_index 등)는 태그가 아니라 extra_metadata에 넣어야
                    Opik 인덱스 카디널리티 폭발을 방지한다.
        extra_metadata: 메타데이터 dict에 merge. 필터링은 metadata 기반으로 가능.
        )rq   rr   run_tagproject_namer  z > episode_title
trace_name
session_idproject_name_tagepisode_title_tagr8   r   rv   opik_trace_v2_enabledF)build_axis_tagsepisode_thread_idN)r  r  rr   r9   pipeline_project_name)stepzstatus:u-   build_opik_metadata v2 실패 (non-fatal): %s)rq   rr   rt   rA   rB   r   extendr  r   rw   r  r<   r!  r"  r@   r   r   r   )r   r  r  rD   ctxr8   re  rz  rw   r!  r"  axisttagr   s                  r5   r>   zStepRunner.build_opik_metadatac  s   " //// 
 779$'GGNB$?#@CGGO]_D`Caadeieqeqdr!sD!$YD~77>"'*>':D#$KKN+,77?#(+O(<D$%KKO,-KK
#V&,,.1D=DG /
,	O0x!8%@8
 t,t,+T2,d3$5!$!<"%''/2">#4"%[! 77>*474GD0177?+,/,@D)779%&))nDO 'DLL9$**A#A3-C$C( +  $V   	OLLH#NN	Os   DI +I 	I/I**I/c                  d    t        j                  t        j                        j	                         S r  )r   r   r   rK  	isoformatr3   r4   r5   r   zStepRunner._now  s    ||HLL)3355r4   )NNr  )r   r   rU   Nr  )rU   r   )r  )r  )r   )rU   r   )rU   r   )Dr)   r*   r+   r,   rR   
OrmSessionr   r
   r   r   boolr   recompiler   r   r   r   r   r   staticmethodr   r   r   r  r	   r   propertyr  rH   r$  r   r7  r3  r  r   intr|  r  r  r  r  r  r  r  r  r  r  r  r  r  r  r  r  r  r  r  r  r  r=   r   r4  r2  r  r   r   r>   r   r3   r4   r5   rp   rp      sG    *.'+

 
 	

 
 !
 tn
<.`+T + "rzz= -D4 DY$ Yv)htn )V A A$3D 3T 3j8_hl _HSM _ae _qu _F 
S  
Xd38n-E  
D @ @.c  G
C G
~ G
R|
S#X |
> |
|
 
 
D B $+7  ! PP 	P
 sCxP P P P P P 
Pl ,8  ! %% sCx	%
 % % % % % 
%NM3 M M #(-1)-R  R &c]	R
 "#R 
Rh /1 @ @$ @4 @D
 
S 
# 
d 
&	4H# H$ HbA bA4S> bAHx*~ x*S x*T#s(^ x*xAT#s(^ A)@S#X )@V <D 9S 9SRUX 9vl# l#S#X l#`]Y .*s *s *8 C D $A$	$S 	$:3 :   &'R,T#s(^ , ,\
TS 
TS#X 
TA A +/37VT#Y'V !c3h0V 
	Vp 6# 6 6r4   rp   )?r,   rf   rd   loggingosr.  r}   
contextlibr   r  r   r   r   enumr   pathlibr   typingr	   r
   r   r   r   r   
sqlalchemyr   sqlalchemy.ormr   r,  r  r   r   app.core.checkpoint_ior   app.core.errorsr   app.core.step_lockr   r   r   r   r  r   r   rY  r   rz   r   r   r   	getLoggerr)   r   r    r2  rQ   r"   rE   rH   _VALUE_NORMALIZED_CONFIG_KEYSrR   rn   rp   r3   r4   r5   <module>r@     s       	 	  % ! '   B B  0I 4 $   
		8	$  s &4 &  ( $* * *H "2 2 )C )C3 )CXN 6 N 6r4   