
    jm                     t	   d Z ddlZddlZddlZddlZddlmZ  ej                  e      Z	ddl
mZmZmZmZ ddlmZ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 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+  eddg      Z,dede-de-de.fdZ/e,ja                  d       ee       ee       ee      fde-de-dede+fd       Z1e,ja                  d       ee       ee       ee      fde-de-de-dede+f
d       Z2e,jg                  d       edd       edd        eddd!"       ed#d$%       ee       ee       ee      fde-d&e-d'e-d(ee4   d)e5de-dede+fd*       Z6e,jo                  d+      e ed,       ed,       ee       ee       ee      fde-d-e4d.e4de-dede+fd/              Z8de-de-defd0Z9 ejt                  d1      Z;e,ja                  d2      e ee       edd3%       ee       ee      fde-de-dee-   de+def
d4              Z<e,jg                  d2      e ee       edd5%       edd6d78       ee       ee      fde-de-dee-   d9ee-   de+defd:              Z=e,jg                  d;      e ee       ed,d<d=8       edd>%       ee       ee      fde-de-d?e-dee-   de+defd@              Z>e,jg                  dA       ed,dB       ed,dB       edCdD       ee       ee       ee      fde-d-e4d.e4d&e-de-dede+fdE       Z?e,ja                  dF      e ee       ee       ee      fde-de-dede+fdG              Z@e,jg                  dH      e eddI%       eddJ%       eddK%       ee       ee       ee      fde-dLe-dee-   dMee-   de-dede+fdN              ZAe,jg                  dO      e ee       ee       ee      fde-de-dede+fdP              ZBe,jg                  dQ      e ed,dR%       ed#dS%       eddT%       ee       ee       ee      fde-de-dMe-dCe5dLe-de-dede+fdU              ZCde-fdVZDe,jg                  dW       edd       ee       ee       ee      fde-de-d&e-de-dede+fdX       ZEd_dYdZde-de-d[e5ddfd\ZFde-de-ddfd]ZGd_d^ZHy)`u@   파이프라인 Step API — 개별 단계 실행/상태 조회.    NPath)AnyDictListOptional)	APIRouterDependsQuery)text)Session)api_endpointget_dbget_current_userverify_project_access)atomic_write_json)AppError)containsget_all_downstream_recursiveget_consumers_ofget_manifest_dictget_ordered_entriesget_resume_sensitive_step_ids)UserAccountz9/api/v1/projects/{project_id}/episodes/{episode_id}/stepssteps)prefixtagsdb
project_id
episode_idreturnc                 "    ddl m}  || ||      S )uG   Opik thread 그룹핑용 context 생성 — dispatch_service로 위임.r   )build_opik_context)&app.services.analysis_dispatch_servicer#   )r   r   r    r#   s       F/Users/manta/Documents/Projects/TheRoad-I1/backend/app/api/v1/steps.py_build_opik_contextr&       s    Ib*j99     current_userc                 &    ddl m} d ||||       iS )u   전체 단계 상태 + DAG 의존성 조회.

    W1-F11: 계산 로직은 `step_readmodel_service.get_all_steps_view`로 이관.
    r   )get_all_steps_viewr   )#app.services.step_readmodel_servicer+   )r    r   r   r)   r+   s        r%   get_all_stepsr-   &   s     G'J
CDDr'   z/{step_id}/resultstep_idc                    ddl }ddlm} ddlm} t        |      st        dd| d       ||j                        |z  d	z  d
z  | z  |z  dz  }|j                         s|dddS 	  |j                  |j                  d            }	||	j                  dd      |	dS # t        $ r}
t        dt        |
      d      d}
~
ww xY w)u8   개별 단계 결과 조회 (체크포인트 데이터).r   Nr   settingszstep.not_foundu   알 수 없는 단계:   codemessagestatus_codecheckpointsepisodeszmanifest.jsonno_data)r.   statusresultzutf-8)encodingr:   unknownzstep.read_errori  )jsonpathlibr   app.core.configr1   _step_containsr   projects_direxistsloads	read_textget	Exceptionstr)r    r.   r   r   r)   r>   r   r1   cp_pathdataexcs              r%   get_step_resultrL   5   s     ('",8OPWy6Ygjkk 	X""#j0
	$	%'1	24;	<>M	N  >>"i4HHRtzz'++W+=>"dhhx.KW[\\ R-s3xSQQRs   (7B   	C)C  Cz/run-allresumez^(resume|force)$)patternanalysisz^(analysis|image|all)$zW20E5: hard cap on provider image-generation calls for this run. Required (>= 1) when category in {image, all} and background_render_reference_mode == 'shot_aware_plan'. Ignored for category=analysis.)gedescriptionFzkW20E5: explicit operator approval to perform real image generation. Required for the shot-aware image path.)rQ   modecategoryimage_call_capapprove_image_generationc           	      ,    ddl m}  ||| |||||      S )u  전체 단계 순차 실행. category: analysis | image | all.

    Phase 4.1: dispatch 로직이 `analysis_dispatch_service`로 이관되어
    `/episodes/{id}/analyze` 경로와 구현을 공유함.

    W20E5: ``image_call_cap`` + ``approve_image_generation`` 두 파라미터는
    shot-aware 이미지 단계의 fail-closed 게이트에 연결된다. legacy
    렌더 모드에서는 둘 다 옵션 (cap을 명시하면 advisory 캡으로 동작).
    r   )dispatch_category_run)rT   rU   )r$   rW   )	r    rR   rS   rT   rU   r   r   r)   rW   s	            r%   run_all_stepsrX   S   s+    F M 
%!9 r'   z/shot_selection/toggle.scene_index
shot_indexc                 B    ddl m}  ||||       j                  ||      S )uG   Shot 선택/해제 토글 — Phase 2.2 ShotSelectionService로 위임.r   )ShotSelectionService)#app.services.shot_selection_servicer\   toggle)r    rY   rZ   r   r   r)   r\   s          r%   toggle_shot_selectionr_      s%     IJ
;BB;PZ[[r'   c                 P    ddl m} t        |j                        | z  dz  dz  |z  S )Nr   r0   r7   r8   )r@   r1   r   rB   )r   r    r1   s      r%   _cp_basera      s+    (%%&3mCjPS]]]r'   z2^manifest_(\d{8}_\d{6})(?:_[a-zA-Z0-9_-]+)?\.json$z
/snapshotsu   특정 단계만 조회c                 B    ddl m}  ||||       j                  |      S )uF   체크포인트 스냅샷 목록 조회 — SnapshotService로 위임.r   SnapshotService)r.   )app.services.snapshot_servicerd   list_versions)r    r   r.   r)   r   rd   s         r%   list_snapshotsrg      s$     >2z:6DDWDUUr'   u-   특정 단계만 스냅샷. 없으면 전체.z^[a-zA-Z0-9_-]{0,32}$u   스냅샷 라벨)rN   rQ   labelc                 D    ddl m}  ||||       j                  ||      S )uV   현재 체크포인트를 수동 스냅샷으로 저장 — SnapshotService로 위임.r   rc   )r.   rh   )re   rd   save_snapshot)r    r   r.   rh   r)   r   rd   s          r%   create_snapshotrk      s'     >2z:6DDW\aDbbr'   z/snapshots/restorez^\d{8}_\d{6}$u"   복원할 버전 (YYYYMMDD_HHMMSS)uB   특정 단계만 복원. 없으면 해당 버전의 전체 단계.versionc                 D    ddl m}  ||||       j                  ||      S )uV   특정 버전의 스냅샷으로 체크포인트 복원 — SnapshotService로 위임.r   rc   )rl   r.   )re   rd   restore)r    r   rl   r.   r)   r   rd   s          r%   restore_snapshotro      s'     >2z:6>>wX_>``r'   z/scene-detail/redo-shot)rP   forcez	^(force)$c                 *    ddl m}  |||| |||      S )u  단건 shot scene_detail re-run — 운영 복구 도구.

    2026-05-11 사용자 권장안 #2 축소형: scene_detail step 의 atomic-fail 대비.
    1 shot 의 LLM 응답이 contract violation 으로 step 전체 abort 된 경우, 본
    endpoint 로 그 shot 만 재실행. cp manifest 의 해당 shot 만 replace +
    checkpoint sync. UI 버튼 X — 운영 호출용.

    Returns:
        service.redo_scene_detail_shot 의 dict 그대로.
    r   )redo_scene_detail_shot)r   r   r    rY   rZ   rR   )&app.services.scene_detail_redo_servicerr   )r    rY   rZ   rR   r   r   r)   rr   s           r%   redo_scene_detail_shot_endpointrt      s'    ( N! r'   z/locksc                    ddl m} ddlm}m}m} |j                  t        d      || d      j                         } |       \  }	}
}g }|D ]M  }|j                  |j                  |j                  |j                  |j                  |j                  d      } |||| |j                  f|j                   |j"                        \  }}|j%                  |j                  |j&                  |j                  |j(                  |j                  r|j                  j+                         nd|j,                  r|j,                  j+                         nd|j                  |j                  |j                  d	|j.                  ||j0                  |j2                  d
       P ddlm}  ||||       }||	|
|d	|j8                  |j;                         ddS )u  지금 누가 무엇을 잡고 있고, **살아있는가**.

    판정은 추측이 아니다 — 같은 프로세스면 등록부를, 같은 호스트면 PID 를
    직접 물어본다 (`app/core/step_lock.py` 머리말의 판정 순서).

    응답의 `verdict` 는 셋이다:
      - `alive`   일하는 중이다. 건드리면 안 된다.
      - `dead`    소유 프로세스가 없다. `release` 로 풀어도 된다.
      - `unknown` 확정 못 한다. 사람이 봐야 한다.
    r   r0   )	LockOwnerjudge_ownerprocess_identityaS  
        SELECT step_id, status, run_id, started_at, heartbeat_at,
               cancel_requested_at, owner_host, owner_pid, owner_boot_id,
               recovery_count, last_recovery_reason
          FROM step_run
         WHERE project_id = :pid AND episode_id = :eid
           AND status = 'running'
         ORDER BY started_at
    )pideid
owner_host	owner_pidowner_boot_idheartbeat_atrun_idkeylease_secondsallow_heartbeat_stealN)hostry   boot_id)r.   r:   r   
started_atr   cancel_requested_atownerverdictverdict_reasonrecovery_countlast_recovery_reason)read_cancel_state)	requesteddescribe)locksthis_process
run_cancel)r@   r1   app.core.step_lockrv   rw   rx   executer   fetchallfrom_rowr|   r}   r~   r   r   r.   step_lock_lease_seconds!step_lock_heartbeat_steal_enabledappendr:   r   	isoformatr   valuer   r   app.core.run_controlr   r   r   )r    r   r   r)   r1   rv   rw   rx   rowsr   ry   r   r   rowr   r   reasonr   cancel_states                      r%   list_step_locksr      s   $ )KK::d  	 Z
02 3;(* 	 *+D#wE"".. ..,,jj$
  &Z5"::"*"L"L	
 	{{jjjj..030@0@  **,d ** ''11304 }},,
 }}$!00$'$<$<+
 	 L 7$RZ@L !%cgF%//$--/
 r'   z/cancelu(   왜 세우는가 — 기록에 남는다uU   특정 스텝만 세운다. 비우면 이 에피소드의 주행 전체를 세운다.u   step_id 를 줄 때는 **필수**. GET /locks 에서 본 run_id 와 다르면 거부한다 — 그 사이 그 스텝이 끝나고 다음 주행이 잡았을 수 있다.r   expected_run_idc           	      z   ddl m} |rt        |      st        dd| d      |st        ddd	      |rd
nd}|j	                  t        d| d      || |t               d|rd|ini       }	|j                          t        |	j                        }
t        j                  d||
|j                  |xs d       d||
|
rddS ddS  |||| |t        |j                              }|j	                  t        d      || t               d       |j                          d|j                  |j                         dS )u?  **협조적 정지** — 멈추라는 표만 찍는다.

    일하는 쪽은 안전 지점(샷 경계, provider 호출 직전)에서 이 표를 읽고
    깨끗이 멈춘다.

    ★진행 중인 provider 호출을 중간에 끊지 않는다. 이미 접수한 이미지는
     결과까지 가져온 뒤 멈춘다 — 끊어 봐야 돈은 이미 나갔고 결과만 잃는다.

    ★`step_id` 를 비우면 **주행 전체**를 세운다. 스텝 하나만 세우면
     배치가 다음 스텝으로 넘어가므로 주행을 멈추려면 이쪽을 써야 한다.
    r   )request_cancelzstep.unknownu   모르는 스텝입니다: r2   r3   zstep.cancel_needs_run_idu   step_id 를 지정할 때는 expected_run_id 도 함께 주세요. GET .../steps/locks 에서 지금 run_id 를 확인하세요.i  zAND run_id = :expected_run_idr(   a  
            UPDATE step_run
               SET cancel_requested_at = CURRENT_TIMESTAMP,
                   updated_at = :now
             WHERE project_id = :pid AND episode_id = :eid AND step_id = :sid
               AND status = 'running'
               z	
        )ry   rz   sidnowr   z/[STEP-CANCEL] step=%s marked=%s by=%s reason=%s-stepu7   정지 표를 찍었다. 안전 지점에서 멈춘다.uJ   지금 running 인 행이 없다 — 이미 끝났거나 시작 전이다.)scoper.   markedr5   )r   requested_byz
        UPDATE step_run
           SET cancel_requested_at = CURRENT_TIMESTAMP, updated_at = :now
         WHERE project_id = :pid AND episode_id = :eid AND status = 'running'
    )ry   rz   r   run)r   r   r5   )r   r   rA   r   r   r   _now_isocommitboolrowcountloggerwarningidrH   r   r   )r    r   r.   r   r   r   r)   r   run_id_clauseupdatedr   states               r%   
cancel_runr   Q  s   B 4g&#5gY?  /T    <K7PR**T '  	#  :	

 8G!?3B
 			g&&'=V\__fm	

   J	
 		
 b	
 		
 
J
C$8E JJt  	 Z

C	E
 IIK __>># r'   z/cancel/clearc                     ddl m}m}  ||||       }|j                  sdddS  |||| |j                        }||rddS ddS )	u  정지 요청을 내린다 — 다시 주행할 수 있게.

    ★이걸 안 부르면 정지 표가 그대로 남아 다음 주행도 즉시 선다.
     `run` 으로 새 claim 을 하면 그 스텝의 표는 지워지지만, 주행 단위 표는
     여기서만 내려간다.
    r   )clear_cancelr   Fu)   걸려 있는 정지 요청이 없었다.)clearedr5   )expected_requested_atu8   정지 요청을 내렸다. 다시 주행할 수 있다.uZ   그 사이 새 정지 요청이 들어왔다 — 내리지 않았다. 다시 확인하라.)r   r   r   r   requested_at)r    r   r   r)   r   r   seenr   s           r%   clear_run_cancelr     so     E
 RZ8D>> -XYY
J
"//G
   G 
 n r'   z/locks/{step_id}/releaseuP   GET /locks 에서 본 run_id. 그 사이 주인이 바뀌었으면 거부된다.u   소유자가 살아있다고 판정돼도 해제한다. 「그 프로세스는 내가 이미 죽였다」를 운영자가 선언하는 경우에만.u%   왜 푸는가 — 기록에 남는다c                    ddl m} ddlm}	m}
m} |j                  t        d      || |d      j                         }|t        d| dd	
      |j                  dk7  rd|j                  d|j                   ddS |j                  |k7  rt        dd| d|j                   dd
      |	j                  |j                  |j                  |j                  |j                   |j                  d      } |||| |f|j"                  |j$                        \  }}||
j&                  ur#|s!t        d| d|j(                   d| dd
      dt+        j,                          }d|j.                   d|j(                   d| d| d |xs | 
}|j                  t        d!      |dd" |t1               || ||d#      }|j3                          t5        |j6                        }t8        j;                  d$|||       ||j(                  ||rd%d'S d&d'S )(uT  죽은 락을 푼다 — `running` → `failed`.

    ★**행을 지우지 않는다.** 죽은 것을 죽었다고 적는 것이다. 다음 `resume`
     이 바로 가져간다.

    ★`expected_run_id` 를 요구하는 이유 — 조회한 뒤 해제하기까지 사이에
     주인이 바뀌었을 수 있다. 그 값을 확인하지 않으면 **새로 들어온 산
     주인을 죽인다.**

    ★기본은 죽음이 **확인될 때만** 통과한다. 살아있거나 확정 못 하면 거부하고,
     운영자가 `force=true` 로 책임을 명시할 때만 통과시킨다.
    r   r0   )rv   OwnerVerdictrw   z
        SELECT status, run_id, started_at, heartbeat_at,
               owner_host, owner_pid, owner_boot_id
          FROM step_run
         WHERE project_id = :pid AND episode_id = :eid AND step_id = :sid
    )ry   rz   r   Nzstep.lock_not_foundu    의 step_run 행이 없다.r2   r3   runningFu   잡혀 있지 않다 (status=u   ) — 할 일이 없다.)releasedr:   r5   zstep.lock_movedu.   그 사이 주인이 바뀌었다 (본 run_id=u   , 지금 run_id=u$   ). 다시 조회하고 판단하라.i  r{   r   zstep.lock_owner_not_deadu:    소유자가 죽었다고 확인되지 않는다 (판정=z: uK   ). 정말 그 프로세스를 죽였다면 force=true 로 다시 부르라.z	released-zmanual release by z
 (verdict=z, force=z, revoked_run_id=z): a  
        UPDATE step_run
           SET status = 'failed',
               run_id = :revoked_run_id,
               error_message = :note,
               last_recovery_reason = :note,
               recovery_count = COALESCE(recovery_count, 0) + 1,
               completed_at = :now,
               updated_at = :now
         WHERE project_id = :pid AND episode_id = :eid AND step_id = :sid
           AND status = 'running'
           AND run_id = :run_id
    i  )noterevoked_run_idr   ry   rz   r   r   u)   [LOCK-RELEASE] step=%s released=%s — %su*   풀었다. 다음 resume 이 가져간다.uE   그 사이 상태가 바뀌어 풀지 않았다. 다시 조회하라.)r   r   r   r5   )r@   r1   r   rv   r   rw   r   r   fetchoner   r:   r   r   r|   r}   r~   r   r   r   DEADr   uuiduuid4r   r   r   r   r   r   r   )r    r.   r   rp   r   r   r   r)   r1   rv   r   rw   r   r   r   r   r   r   r   r   s                       r%   release_step_lockr     s   F )GG
**T  	
 Z
@B
 CK(*  {&i;<
 	

 zzYjj6szzlBZ[
 	

 zz_$"@@Q R!!$,PR 
 	
 nn]]**((**  E *W-66&HH	G^ l'''+) "==/N+; <[[ 
 	
( !/N
\__- .MM?(5' 2)*##^
$	& 	 jj  	 Ud(z!G* IIKG$$%H
NN>SWX==(  9	 	 Y	 	r'   c                  d    ddl m } m} | j                  |j                        j	                         S )uG   `step_run` 의 text 시각 칸에 넣는 값 (기존 계약 그대로).r   datetimetimezone)r   r   r   utcr   r   s     r%   r   r   t  s!    +<<%//11r'   z
/{step_id}c           	          ddl m} t        |||       } |||| |||j                  |      }|j	                  d|       |S )u   개별 단계 실행. mode: resume (기본) | force (재실행). 백그라운드 스레드.

    W1-F11: 실제 orchestration은 `step_execution_service.start_step`로 이관.
    route는 Opik context만 빌드한 뒤 서비스에 위임.
    r   )
start_step)r   r   r    r.   rR   actor_idopik_contextr.   )#app.services.step_execution_servicer   r&   r   
setdefault)	r    r.   rR   r   r   r)   r   r   r;   s	            r%   run_stepr   }  sO     ?&r:zBL!F i)Mr'   Tr   r   c                D    ddl m}  ||| |      j                  |       y)u  T2I appearance count projection — Phase 2.1 EpisodeProjectionService로 이관.

    Legacy wrapper: 기존 호출자(entities.py) 호환 유지용.
    신규 코드는 `app.services.checkpoint_sync.EpisodeProjectionService.sync_t2i_appearance_counts` 직접 사용.
    r   )EpisodeProjectionServicer   N)app.services.checkpoint_syncr   sync_t2i_appearance_counts)r   r    r   r   r   r   s         r%   _sync_t2i_appearance_countsr     s#     FRZ8SS[aSbr'   c                 $    ddl m}  || ||       y)u  분석 체크포인트 → DB 동기화 (Phase 2.1 Service 오케스트레이터로 위임).

    Legacy wrapper: 기존 호출자(image_steps.py, steps.py run_step 등) 호환 유지용.
    실제 로직은 `app.services.checkpoint_sync.orchestrate_full_sync`에 5-way 분해됨.
    r   )orchestrate_full_syncN)r   r   )r   r    r   r   s       r%   _sync_checkpoints_to_dbr     s     C*j"5r'   c                 (    ddl m}  || |||||      S )u   Step ID → StepRunner 인스턴스.

    Legacy wrapper: `app.services.analysis_dispatch_service.get_step_runner`로 위임.
    기존 호출자(test_sync_v3.py, test_pipeline_v3_e2e.py, 내부 run_step) 호환 유지.
    r   )get_step_runner)r$   r   )r.   r   r    r   project_configr   r   s          r%   _get_step_runnerr     s     G7J
BP\]]r'   )N)I__doc__r>   loggingrer   r?   r   	getLogger__name__r   typingr   r   r   r   fastapir	   r
   r   
sqlalchemyr   sqlalchemy.ormr   
OrmSessionapp.api.depsr   r   r   r   app.core.checkpoint_ior   app.core.errorsr   app.core.step_catalogr   rA   r   $catalog_get_all_downstream_recursiver   r   r   r   app.models.catalogr   routerrH   dictr&   rF   r-   rL   postintr   rX   patchr_   ra   compile_SNAP_TS_RErg   rk   ro   rt   r   r   r   r   r   r   r   r   r    r'   r%   <module>r      sm   F   	  			8	$ , , - -  0 V V 4 $  +	U]d\e	f:J :C :S :T : B 34V_ '(8 9	EEE 	E 	E E   34V_ '(8 9RRR R 		R
 R !R: Z h(:;*.FG$)-		% &+B& 34V_ '(8 9/,,
, , SM	, #,* +,, 	-,. /, ,^ &' SzCj34V_ '(8 9
\
\
\ 
\ 	
\
 	
\ 
\  (
\$^ ^# ^$ ^ bjj9
 L 34"45NO '(8 9V_	V	V	V c]	V 		V
 		V  	V \ 34"45de /GUgh '(8 9V_
c
c
c c]
c C=	
c
 
c 	
c  
c !" 34&6Dhi"45yz '(8 9V_
a
a
a 
a c]	
a
 
a 	
a  #
a &' SQ'CA&g{334V_ '(8 9  	
  	  (\ H 34V_ '(8 9	OOO 	O 	O  Od Y (RS"c &+i& 34V_ '(8 9%hhh c]h c]h  !h" 	#h$ %h  hV _ 34V_ '(8 9	    	  	    F '( !^ F (OP34V_ '(8 9'LLL L L  !L" #L$ 	%L& 'L  )L^2# 2 \ h(:;34V_ '(8 9  	
 	  >cei cC cS c^b cnr c6 6 6T 6^r'   