
    -ǚj                    *   U d Z ddlmZ ddlZddlZddlZddlZddlmZmZ ddl	m
Z
mZmZmZmZ ddlmZ ddlmZ ddlmZ dd	lmZmZmZ dd
lmZ ddlmZmZmZ ddl m!Z!  ejD                  e#      Z$dZ%de&d<   	 	 	 	 	 	 	 	 d,dZ'd-dZ(d.dZ)	 	 	 	 	 	 	 	 	 	 d/dZ*	 	 	 	 	 	 	 	 	 	 d/dZ+	 	 	 	 	 	 	 	 	 	 	 	 d0dZ,ddd	 	 	 	 	 	 	 	 	 	 	 d1dZ-d2dZ.	 d3	 	 	 	 	 	 	 	 	 	 	 d4dZ/	 	 	 	 	 	 	 	 d5dZ0	 	 	 	 	 	 	 	 d6dZ1	 	 	 	 	 	 	 	 d7dZ2	 d3	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 d8dZ3dZ4	 	 	 	 	 	 	 	 d9d Z5	 	 	 	 	 	 	 	 d:d!Z6	 	 	 	 	 	 d;d"Z7dd#d#dd$	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 d<d%Z8	 	 	 	 	 	 	 	 d=d&Z9d'Z:d(Z;	 	 d>	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 d?d)Z<	 	 	 	 d@d*Z=	 d3	 	 	 dAd+Z>y)Buq  StepRunner 기반 분석 dispatch 서비스.

Phase 4.1 (architecture-refactor-final/02-final-roadmap.md §Phase 4).
기존 `AnalysisService.run_analysis/reanalyze_scenes` 경로와 `steps.run_all_steps`의
`_run_all_bg` 내부 로직을 단일 dispatch 진입점으로 통합.

공식 경로: Frontend / API / 내부 호출 전부 StepRunner (→ StepCatalog) 경유.
    )annotationsN)datetimetimezone)AnyDictListOptionalTuple)SessionSessionLocal)AppError)ImageCallBudgetinstall_budgetuninstall_budget)submit_background_job)STEP_CATALOGget_active_entriesget_all_downstream_recursive)orchestrate_full_sync)analysisimageall	reanalyzezTuple[str, ...]RUN_ALL_CATEGORY_TOKENSc           
     4   ddl m} ddlm} | j	                  |      j                  |j                  |k(        j                         }| j	                  |      j                  |j                  |k(  |j                  |k(        j                         }|r|j                  n|dd }|r|j                  n|dd }t        j                  t        j                        j                  d      }	| d| d|	 dt!        t#        j$                               dd  }
|||
dS )	u(   Opik thread 그룹핑용 context 생성.r   )ProjectRegistryEpisodeN   z%m%d-%H%M%S_)project_nameepisode_titlerun_tag)app.models.catalogr   app.models.projectr   queryfilteridfirst
project_idnametitler   nowr   utcstrftimestruuiduuid4)dbr+   
episode_idr   r   projectepisoder"   r#   tsr$   s              \/Users/manta/Documents/Projects/TheRoad-I1/backend/app/services/analysis_dispatch_service.pybuild_opik_contextr:   0   s     3* 	!((););z)IJPPR  		

j('*<*<
*J	K	 
 $+7<<
2AL%,GMM*Ra.M	hll	#	,	,]	;Baat1S5Fr5J4KLG$&     c                F   ddl m} | j                  |      j                  |j                  |k(        j                         }|r|j                  si S 	 t        j                  |j                        S # t        j                  $ r t        j                  d|       i cY S w xY w)uG   ProjectSettings.llm_config_json 로드 (없거나 손상 시 빈 dict).r   )ProjectSettingsuE   Malformed llm_config_json for project=%s — falling back to defaults)r&   r=   r'   r(   r+   r*   llm_config_jsonjsonloadsJSONDecodeErrorloggerwarning)r4   r+   r=   pss       r9   load_project_llm_configrE   J   s    2 	!	**j8	9	 
 R''	zz",,-- S	
 	s   A2 2+B B c                "    ddl m}  |||       S )u  text/PDF/checkpoint 중 하나라도 있으면 True (단일 source 위임).

    PDF-only 업로드 후 text 추출 짧은 경우에도 True 반환되어 dispatcher 가
    planning_doc_analysis step 을 cascade 에 포함시킴 (review IMPORTANT 통일).
    r   )project_has_planning_doc)r4   )*app.services.planning_doc_analysis_servicerG   )r4   r+   rG   s      r9   _project_has_planning_docrI   _   s     T#J266r;   c                x    ddl m} | j                   |d      |||d      j                         }|sy|d   dk(  S )u   step_run.status == 'completed' 여부 판정.

    Fix 3.3 prerequisite 검증용. partial은 아직 완료 미달로 처리하여 사용자에게
    명확한 에러를 노출 — silent transitive invalidate 보다 안전한 정책.
    r   )textz\SELECT status FROM step_run WHERE project_id = :pid AND episode_id = :eid AND step_id = :sid)pideidsidF	completed)
sqlalchemyrK   executefetchone)r4   r+   r5   step_id	_sql_textrows         r9   _is_step_completedrV   j   sR     -
**O	
 :g> hj  q6[  r;   c                F    t        | |||      ryddlm}  ||||      dk(  S )u  prerequisite 검증용: step이 완료되었거나 applicability로 skip 가능한지.

    Codex P1-1: image preflight applicability 무시 사고 fix.
        background_mode=='off' (default) 환경에서 background_classify /
        floor_plan_prompt / background_prompt 등은 applicability='if_background_mode'
        → 정적 평가 시 not_applicable. StepRunner는 자동 skip 처리.
        prerequisite 검증이 이를 무시하고 _is_step_completed=False로만 판정하면
        category='image' force가 dispatch.deps_incomplete로 차단되는데, 실제 실행
        시에는 StepRunner가 이를 skip하고 downstream gate를 통과시키므로 차단은 잘못.

    판정 로직:
        1) step_run.status='completed'  → True (이미 완료)
        2) applicability 정적 평가 결과 not_applicable → True (StepRunner skip 가능)
        3) 그 외 → False (prerequisite 미충족 — 사용자에게 명확한 에러 안내)
    Tr   )evaluate_step_applicabilitynot_applicable)rV   app.core.applicabilityrX   )r4   r+   r5   rS   rX   s        r9   _is_step_satisfiedr[      s-    * "j*g>B&w
JGK[[[r;   c           
     6   |dk7  ryg }| D ]N  }t        |      D ]>  }||v r|j                  |      }|s|j                  |k7  s,|j                  ||f       @ P |r6t	        |D 	ch c]  \  }	}|	 c}}	      }
t        dd|d|
 d| dd	      yc c}}	w )
ui  Fix D (cascade fix 2026-05-05): mode=force 시 invalidate_downstream 이
    cross-category step 무효화할 risk 를 시작 시점에 거절.

    StepRunner.run() 의 force 동작이 invalidate_downstream(delete_cp=True) 호출 →
    transitive downstream manifest 모두 delete. batch 안에 없는 cross-category
    step 도 delete 됨 → batch 진행 중 mid-cascade missing 발생 (e.g.
    floor_plan_render(image) force → background_prompt(analysis) manifest delete →
    background_render 시점 fail; scene_detail(analysis) force → scene_image_pipeline
    (image) manifest delete → 후속 image batch 가 stale).

    검증: 각 batch step 의 transitive downstream 중 batch 외 active step 의 category
    가 현재 category 와 다르면 risky. 시작 시점 fail-fast + 사용자에게 `category=all
    &mode=force` 또는 `category=<X>&mode=resume` 안내.

    Args:
        direct: batch 에 포함된 step_id list.
        direct_set: direct 의 set (transitive downstream 검사 시 batch 내 step skip).
        active_by_id: step_id → entry. caller 에서 한 번 build 후 재사용.
        mode: "resume" 또는 "force". "resume" 면 검증 skip.
        category: 현재 batch category ("analysis" / "image"). non-matching downstream 만 risky.

    Raises:
        AppError(dispatch.force_cascade_cross_category) — risky 발견 시.
    forceNz%dispatch.force_cascade_cross_categoryz	category=u6    mode=force 가 cross-category step 무효화 위험: u-   . `category=all&mode=force` 또는 `category=u   &mode=resume` 사용 권장.  codemessagestatus_code)r   getcategoryappendsortedr   )direct
direct_setactive_by_idmoderd   cross_downstreamrN   dsidds_entryr!   cross_stepss              r9   _check_force_cascade_riskro      s    > w.005Dz!#''-H  H, ''d4 6  2BC2Bwq$d2BCD8H< (%%0M 2%J&BD 
 	
 Cs   (B
resumer5   rj   c               @   |dvrt        dd| d      |dvrt        dd|d	d      t        | |      }g }t               D ]K  }|d
k7  r|j                  |k7  r|j                  }|dv r)|dk(  r|s1|j                  |j                         M ddlm |dk7  r|dk(  r|t               D 	ci c]  }	|	j                  |	 }
}	t        |      }g }|D ]p  } |      xs i }|j                  dg       D ]M  }||v r|
j                  |      }|s|j                  dk7  s,t        | |||      r;|j                  ||f       O r |r0t        |D ch c]  \  }}|	 c}}      }t        dd| dd      t        |||
|d       t        |fd      S |t               D 	ci c]  }	|	j                  |	 }}	t        |      }g }|D ]`  } |      xs i }|j                  dg       D ]=  }||v r|j                  |      }|st        | |||      r+|j                  ||f       ? b |r0t        |D ch c]  \  }}|	 c}}      }t        dd| dd      |dk(  r|t        ||d       t        |fd      }|S c c}	w c c}}w c c}	w c c}}w )u  카테고리(analysis/image/all)에 해당하는 실행 가능 step_id 목록.

    applicability(on_demand/disabled/if_planning_doc)를 정적 필터만 수행.
    이외 applicability는 StepRunner 런타임(resolve_applicability)에서 처리.

    Fix 3.3 (M4 해소): category="image" 동작 변경.
    이전: image step의 transitive deps(analysis 포함) 자동 포함 — 사용자가 의도한 "image만 force"
          가 실제로 거의 모든 step force로 폭발 (오늘 13:13 사고 trigger).
    현재: image step만 직접 반환. analysis dep이 미완료/stale이면 prerequisite 검증으로 거부
          (AppError code="dispatch.deps_incomplete"). 사용자가 다음 행동(category=all 또는 해당
          analysis step 먼저 실행)을 명시적으로 선택. silent miss는 prerequisite 에러로 대체.
    prerequisite 검증은 episode_id가 제공된 경우에만 적용 (preview/list 호출은 episode 없음).

    Fix D (cascade fix 2026-05-05): `mode="force"` 시 force-invalidate 가
    cross-category downstream step 을 무효화할 risk 를 시작 시점에 거절.
    StepRunner.run() 의 force 동작이 invalidate_downstream(delete_cp=True) 를
    호출 — 만약 image step 의 transitive downstream 이 active analysis step 을
    포함하면 batch 진행 중 mid-cascade missing 발생 (e.g. floor_plan_render force
    → background_prompt manifest delete → background_render 시점 fail).
    )r   r   r   zdispatch.invalid_categoryzUnknown category: r^   r_   )rp   r]   zdispatch.invalid_modez&mode must be 'resume' or 'force' (got )r   	on_demanddisabledif_planning_docr   get_manifest_dictr   r   
depends_onzdispatch.deps_incompleteuN   analysis step 실행 전 cross-category prerequisite step 미완료입니다: u   . `category=all` 사용 권장.c                <     |       xs i j                  dd      S Norderc   rc   sry   s    r9   <lambda>z+select_steps_for_category.<locals>.<lambda>4  s    ->q-A-GR,L,LWVX,Yr;   )keyuI   image step 실행 전 다음 prerequisite step들이 미완료입니다: uE   . `category=all` 사용 또는 해당 step을 먼저 실행하세요.r]   c                <     |       xs i j                  dd      S r|   r   r   s    r9   r   z+select_steps_for_category.<locals>.<lambda>c  s    ,=a,@,FB+K+KGUW+Xr;   )r   rI   r   rd   applicabilityre   rS   app.core.step_manifestry   setrc   r[   rf   ro   )r4   r+   rd   r5   rj   has_planningrg   entryappleactive_by_id_adirect_set_a
cross_depsrN   metadep	dep_entryr!   dmissingri   rh   incomplete_depsorderedry   s                           @r9   select_steps_for_categoryr      s7   8 33,(
3
 	
 &&(<THAF
 	
 -R<L F#%u8!;"",,$$\emm$ & 97 z!j&<4F4HI4Hqaiil4HNIv;L02J(-388L"5Cl*  . 2 23 7I$  ))Z71"j*cR&--sCj9 6   
!;
1!
!;<3//6i7VX !$  &ndJ f"YZZ .@.BC.B		1.BC[
13C$S)/RDxxb1*$(,,S1	 )"j*cJ#**C:6 2  O<ODAqaO<=G/_`g_h iZ [    w:1 	"&*lD'R V!XYGN_ J "<8 D  =s   1J
J
(J<J
c                $   t        t        d            }|j                  d       t        | |      }g }t	        d      D ]M  }|j
                  |vr|j                  dv r!|j                  dk(  r|s3|j                  |j
                         O |S )uh  씬 재분석(reanalyze-scenes) 대상 step 목록.

    엔티티(entity_* / outlook_*)는 보존하고 씬 관련 step만 force 재실행.
    `scene_segmentation`부터 시작 — baseline `AnalysisService.reanalyze_scenes`가
    `segments.json`을 삭제하여 재추출한 효과를 StepRunner force 모드로 재현.
    (Codex Review Important #2)
    scene_segmentationr   )rd   rt   rw   )r   r   addrI   r   rS   r   re   )r4   r+   
downstreamr   resultr   s         r9   select_scene_reanalysis_stepsr   g  s     12FGHJNN'(,R<L F#Z8==
*";;"33Lemm$ 9 Mr;   c                    ddl m} 	 ddlm} i ||}|j                  |       }	|	st        dd|  d       |	| |||||      S # t        $ r i }Y Ew xY w)	u$   Step ID → StepRunner 인스턴스.r   )STEP_CLASSES)IMAGE_STEP_CLASSESzstep.not_foundzStep class not found:   r_   )rS   r+   r5   r4   project_configopik_context)app.core.stepsr   app.core.steps.image_stepsr   ImportErrorrc   r   )
rS   r+   r5   r4   r   r   
V3_CLASSESr   step_classesclss
             r9   get_step_runnerr     s     : A 8(7J7L


7
#C!,WI6
 	
 %!     s   A AAc                   ddl m} 	 | j                  |      j                  |j                  |k(        j                         }|r2|j                  dk(  r"d|_        |dd |_        | j                          yyy# t        $ rF}t        j                  d|       	 | j                          n# t        $ r Y nw xY wY d}~yY d}~yd}~ww xY w)uO  파이프라인 실패 시 episode.status 복구 — 'analyzing' 잠금 해제.

    Review Critical #1: run_steps_batch 실패 시 status가 'analyzing'에 영구 잠기면
    UI에서 분석 진행 중으로 오인되며 재실행이 차단된다.
    baseline AnalysisService.run_analysis의 except 블록과 동일한 의도.
    r   r   	analyzingerrorN  z$Failed to recover episode status: %sr&   r   r'   r(   r)   r*   statusanalysis_errorcommit	ExceptionrB   r   rollback)r4   r5   error_messager   eprec_excs         r9   "_recover_episode_status_on_failurer     s     +XXg%%gjjJ&>?EEG")){*BI -et 4BIIK +2  ;WE	KKM 		 s<   A-A8 8	CCB)(C)	B52C4B55CCc                   ddl m} 	 | j                  |      j                  |j                  |k(        j                         }|r9|j                  dk(  r)d|_        d|xs d dd |_        | j                          yyy# t        $ rF}t        j                  d	|       	 | j                          n# t        $ r Y nw xY wY d}~yY d}~yd}~ww xY w)
u  정지로 멈췄을 때 `analyzing` 잠금을 푼다 — **`error` 로는 안 적는다.**

    `analyzing` 은 「지금 돌고 있다」는 뜻이라, 착수 검사가 그것을 보고
    재실행을 409로 막는다. 세워 놓고 다시 못 돌리게 되는 것이다.

    ★그렇다고 `error` 도 아니다. 운영자가 멈춘 것이지 깨진 것이 아니다.
     상태는 `stopped` 로 적고, 사유는 `analysis_error` 칸에 남긴다
     (그 칸이 「무슨 일이 있었나」를 담는 유일한 자리다).
    r   r   r   stoppedu   주행 정지:    (사유 없음)Nr   u0   정지 뒤 에피소드 상태 해제 실패: %sr   )r4   r5   reasonr   r   excs         r9   !_release_episode_status_on_cancelr     s     +XXg%%gjjJ&>?EEG")){*!BI"1&2M<M1N OPUQU VBIIK +2  GM	KKM 		 s<   A4A? ?	CC	B0/C	0	B<9C	;B<<C		Cc                
   | dvryddl m} t        |dd      }|dk7  r$|yt        dt	        |            }t        |      S |st        d	d
d      |t	        |      dk  rt        ddd      t        t	        |            S )u|  W20E5 — decide whether the worker thread needs an image-call budget.

    Returns the :class:`ImageCallBudget` instance the worker should install
    for the duration of the run, or ``None`` when no budget is required.

    Behaviour:

    * ``category`` not in ``{"image", "all"}`` → ``None`` (no image steps).
    * ``settings.background_render_reference_mode != "shot_aware_plan"``
      (i.e. legacy / w18j_overlap) → ``None`` unless the caller explicitly
      passes a ``image_call_cap``; in that case the cap is honoured as an
      advisory opt-in (no approval required).
    * Shot-aware image path → fail closed. ``approve_image_generation``
      must be ``True`` and ``image_call_cap`` must be ``>= 1``; otherwise
      :class:`AppError` (``step.image_generation_not_approved`` /
      ``step.image_call_cap_required``) is raised before any worker is
      submitted.
    )r   r   Nr   settings background_render_reference_modelegacyshot_aware_plan)capz"step.image_generation_not_approvedzKW20 shot-aware image phase requires explicit approve_image_generation=true.r^   r_      zstep.image_call_cap_requiredz8W20 shot-aware image phase requires image_call_cap >= 1.)app.core.configr   getattrmaxintr   r   )rd   image_call_capapprove_image_generation	_settingsrender_modecap_ints         r9   evaluate_image_call_budget_gater     s    0 ''5)%GRK'' !a^,-7++ $51 
 	
 ^!4q!8/J
 	
 s>233r;   c           	     "   |t        |       t               }d}	 ddlm}	 d}
d}t	               }|D ]  }ddlm}  ||| |      }|j                  r>t        j                  d||j                                d}d|j                          d	} n`	  |	|      xs i }|j                  d
g       D cg c]	  }||v s| }}|r)t        j                  d||       |j                  |       |j                  dd      r	 t        | ||       t#        || ||||      }|j%                  |      }|j                  dd      }t        j'                  d||       |j                  dd      r	 t        | |||       |dk(  r t        j                  d|       d}
d| d} nR|dk(  rJ|j                  dd      du r t        j                  d|       d}
d| d} nt        j                  d|        |r)t        j                  d&t/        |      t1        |             |rNt        j                  d'|xs d(       t3        |||       dd)|xs d*d+|j5                          |t7                S S |
r1	 t        | ||       dd.d/d+|j5                          |t7                S S t9        |||xs d0       dd|xs d0d+|j5                          |t7                S S c c}w # t        $ r=}t        j                  d||       |j!                          d}
d| d| }Y d}~ ,d}~ww xY w# t        $ r2}t        j                  d||       |j!                          Y d}~d}~ww xY w# t(        $ r}|j*                  dk(  r6t        j                  d ||j,                         d}|j,                  }Y d}~ |j*                  d!k(  r8t        j                  d"||j,                         |j                  |       Y d}~t        j                  d#||j,                         d}
d| d|j,                   }Y d}~ Cd}~wt        $ r-}t        j                  d$||       d}
d| d%| }Y d}~ wd}~ww xY w# t        $ rf}t        j                  d,|       |j!                          t9        ||d-|        ddd-| d+cY d}~|j5                          |t7                S S d}~ww xY w# t        $ r}t        j                  d1|       	 |j!                          n# t        $ r Y nw xY wt9        ||t;        |             ddt;        |      d+cY d}~|j5                          |t7                S S d}~ww xY w# |j5                          |t7                w w xY w)2u  step_ids 순차 실행 — background job 내부에서 호출되는 실제 워커.

    각 step 완료 후 sync projection(체크포인트→DB). 실패/partial 처리 포함.
    실패 시 episode.status를 'analyzing' → 'error'로 복구.

    Returns:
        ``{"ok": bool, "outcome": "completed|cancelled|failed", "reason": str}``

        ★★★2026-09-04 — 종전에는 **아무것도 안 돌려줬다**(`None`). 실패를
         `all_ok=False` 로 접고 episode.status 만 고친 뒤 조용히 끝났다.
         그래서 프로젝트 대기열이 「앞 화가 실패했나」를 알 길이 없어
         **다음 화를 그대로 시작했다** (Codex BLOCK 2026-09-04).
         배경 작업 경로는 이 값을 안 쓰므로 계약은 뒤로 호환된다.

    W20E5 — ``budget`` (optional) is installed on the worker thread for the
    duration of the run and cleared in ``finally`` so it cannot leak to
    other jobs. When ``None``, the existing uncapped behaviour applies.
    Nr   rx   TF)read_cancel_stateu9   [RUN-CANCEL] 배치 정지 — %s 앞에서 멈춘다. %su   주행 정지 요청 (rs   rz   u2   Step %s skipped (cascade) — upstream blocked: %s#requires_projection_sync_before_runu9   Required pre-sync failed before %s: %s — aborting batchzPre-sync failed before z: )rj   r   donezStep %s: %s)rS   uA   Post-step sync failed after %s — DB projection may be stale: %sfailedz!Step %s failed, stopping pipelinezStep z failedpartialallow_partial_downstreamuH   Step %s partial and allow_partial_downstream=False — stopping pipelinez) partial (allow_partial_downstream=False)uC   Step %s partial — continuing (downstream may use partial results)zstep.cancelledu)   주행 정지 — %s 에서 멈춘다: %szgate.blockeduF   Step %s blocked: %s — downstream depending on this will cascade-skipzStep %s error: %szStep %s crashed: %sz
 crashed: u   Pipeline finished with %d blocked step(s): %s. Downstream skipped via cascade. Hint: missing dependencies may be in another category — try category='all'.u   [RUN-CANCEL] 주행이 정지 요청으로 멈췄다 — %s. 다시 돌리려면 POST .../steps/cancel/clear 로 표를 내려라.r   	cancelledu   정지 요청okoutcomer   zFinal sync failed: %szFinal sync failed: rO    zPipeline failedzrun_steps_batch outer error: %s)r   r   r   ry   r   app.core.run_controlr   	requestedrB   rC   describerc   r   r   r   r   r   r   runinfor   r`   ra   lenrf   r   closer   r   r1   )r+   r5   step_idsrun_moder   r   budgetr4   failure_messagery   all_okr   blocked_stepsrN   r   cancel_state	step_metar   blocked_depssync_excrunnerr   r   aer   	outer_excs                             r9   run_steps_batchr     s   6 v	B%)OK< 	 #&%C
 ?,RZHL%%O..0
 !	$:<;P;P;R:SST"Un-c28b	 ,5==r+Ja+JaaS`N`+JaNNL\ "%%c* ==!FN-j*bI )Z^\  2Hf5M37 ==!FN
&-j*bRUV X%LL!DcJ"F(-cU'&:Oy( !}}%?F%Of "'#C5(QR ( NN]C F LL`M"F=$9  NNX4#4
 .b*oNK-@B@ 	
 ? D%j*bA ;"E" 	
 ! /J D3D  H-B1BD 	
 I  b %  W$
 !&,CC58**UD % & +,/ &@  77.. NN#NPSUWU_U_` $I&(jjO77n,NN`RZZ "%%c*0#rzzB$)#b"= 2C=$)#j">	D  	D4h?2
&9($D
 $$7z"BD D& 	
 9	D"  L6	B	KKM 		 	+2z3y>JC	NKK

 L 	
 s}  A4S   M!0	K:K>,M!*S  +M!>KAM!)L#8#M!S  7M!S  M!,A S  *S  -Q. :S  S  M!	L #1LM!S  L  M!#	M,'MM!MM!!	Q+*>P2(S  /AP20S  65P2+S  2Q+>!Q&S  &Q++S  .	S7>S5S6S  SS   	U,)U' TU'	TU'T'U'U,U/ 'U,,U/ /V)r   r   c                   ddl m} ddlm} ddlm} | j                  |      j                  |j                  k(  |j                  k(        j                         j                         }|st        d |d      d      d}|j                  d	k(  rNdd
lm t!        fdt"        D              }|rt        d |d      d      t$        j'                  d       d}|j(                  st        d |d      d      ddlm}	  |	       st        d |d      d      |rdn|j                  }
d	|_        d|_        | j1                          |
S )u  분석 시작 preflight — /analyze와 /steps/run-all 공용 (Codex W1 Critical 수정).

    1. Episode row lock
    2. status != "analyzing" 검증 (이미 실행 중이면 409)
    3. fulltext 존재 검증 (없으면 400)
    4. OPENAI_API_KEY 설정 검증 (없으면 400)
    5. status='analyzing', analysis_error=None set + commit (race window 종료)

    Returns:
        prior_status — dispatch 예외 시 복구에 사용. 호출자가 rollback 책임.
    r   r   )tr   episode.not_foundr   r_   Fr   is_task_runningc           	   3  @   K   | ]  } d  d d|         yw)run_all::N ).0catr5   r   r+   s     r9   	<genexpr>z+preflight_analysis_start.<locals>.<genexpr>/  s1      
. hzl!J<qFG.s   zanalysis.already_running  u   preflight_analysis_start: episode %s status='analyzing' but no active task_registry entry — auto-recovering stale lock from previous abnormal exitTzanalysis.no_textr^   )has_openai_keyzanalysis.openai_key_missingr   N)r   r   app.i18n.loaderr   r&   r   r'   r(   r)   r+   with_for_updater*   r   r   app.core.task_registryr   anyr   rB   rC   fulltextapp.core.openai_keysr   r   r   )r4   r+   r5   r   r   r   r7   stale_recoveredhas_active_jobr   prior_statusr   s    ``        @r9   preflight_analysis_startr  
  sJ     6!* 		

j('*<*<
*J	K				  /;N9O]`aaO~~$
 	; 
.
 
 /45 
 	]	

 .:L8M[^__3.34
 	
 .77>>L GN!GIIKr;   c                    |yddl m} | j                  |      j                  |j                  |k(        j                         }|r||_        | j                          yy)u;   preflight 이후 dispatch 실패 시 episode.status 복구.Nr   r   )r&   r   r'   r(   r)   r*   r   r   )r4   r5   r  r   r   s        r9   rollback_episode_statusr
  R  sS     *	'		!	!'**
":	;	A	A	CB	 	
		 
r;   c                    |yt               }	 t        || |       |j                          y# |j                          w xY w)u  세션이 죽었을 때 **새 세션으로** 상태를 되돌린다 (best-effort).

    실패를 일으킨 세션이 DB 오류로 못 쓰게 됐고 `db.rollback()` 마저 실패한
    경우다. 그 세션으로는 아무것도 못 하므로 새 세션을 연다.
    N)r   r
  r   )r5   r  sessions      r9   _rollback_episode_status_freshr  b  s8     nG\Bs   - ?F)r   r   
run_inlineclaim_held_byc          	        ddl m}	 t        D ]!  }
 |	d|  d| d|
       st        ddd       dd	l m}m d|  |d
u r1 || d|       s"ddl m} t        dd |      xs d dd      s)ddl m}  |      }||k7  rt        dd|xs d dd      d
}	 |t        v rt        || |      }t        |||      }t        ||       }t        || |||      }t        || |      }d|  d| d| }t        j                  t              fd       }| ||||||f}|rM || xs dddd}t!        |j#                  d            |j#                  dd      |j#                  dd      ||dS t%        |||d | d!| "      }|st        ddd      	 d#d'||d(S # t&        $ r d#}	 |j)                          n)# t&        $ r d}t*        j-                  d$|d#%       Y nw xY w	 |rt/        |||       nt1        ||       n(# t&        $ r t*        j3                  d&||d#%       Y nw xY wr         # r	        w w xY ww xY w))u  카테고리 단위 run-all — 단일 진입점.

    ``run_inline`` — 배치를 **이 스레드에서 끝까지** 돌린다. 프로젝트 대기열이
    여러 편을 순서대로 돌릴 때 쓴다. 배경 작업으로 던지면 바로 돌아오므로
    「앞 화가 끝난 뒤 다음 화」를 만들 수 없다.

    ``claim_held_by`` — 호출자가 **프로젝트 자리를 이미 들고 있다**는 표시.
    ★대기열이 화마다 자리를 놓았다 잡으면 그 틈에 단일 실행이 끼어들어
    순서를 깨고 canon·장부를 갱신한다 (Codex BLOCK 2026-09-04). 대기열은
    자리를 **처음부터 끝까지** 들고, 안쪽 화 실행은 그 자리를 물려 쓴다.

    API 엔드포인트(run_all_steps, /analyze) 공용 dispatch 헬퍼.
    category가 analysis/all이면 preflight 포함 (Codex W1 Critical 수정).

    W20E5: when ``category in {"image", "all"}`` and
    ``settings.background_render_reference_mode == "shot_aware_plan"``,
    the caller must explicitly pass ``approve_image_generation=True`` and
    a positive ``image_call_cap`` or the dispatch raises before any
    background job is submitted. The budget is then installed on the
    worker thread so every provider call site reserves against it.
    r   r   r   r   step.already_running!   파이프라인 이미 실행 중r   r_   claim_episode_runrelease_episode_runNholder)episode_run_holderuT   이 프로젝트에서 이미 다른 화가 돌고 있습니다 (먼저 잡은 쪽: ?u(   ). 여러 편은 순서대로 돕니다.zstep.claim_not_heldu8   자리를 물려받았다고 했는데 실제 주인은 u   없음u    이다)rd   r   r   rq   c                 L    	 t        |  r	        S S # r	        w w xY w)uO  배치가 끝나면 에피소드 자리를 놓는다.

            ★놓는 자리가 여기여야 한다 — 아래 `submit_background_job` 은
             제출만 하고 바로 돌아온다. 제출 직후에 놓으면 배치가 도는 동안
             자리가 비어 두 번째 요청이 그대로 들어온다.
            r   )args
_own_claimepisode_keyr  s    r9   _run_then_releasez0dispatch_category_run.<locals>._run_then_release  s3    5&-
 '4 :'4 s    #Fr   u.   run_steps_batch 가 결과를 안 돌려줬다r   r   r   r   r   )r   r   r   stepsjob_keyzRun all z steps for r!  targetr  descriptionTu>   dispatch 실패 정리 중 세션 rollback 실패 (episode=%s)exc_infou   dispatch 실패 정리 중 episode.status 되돌리기 실패 (episode=%s, 되돌릴 값=%s) — 원래 예외를 그대로 올린다startedr   r   r   r!  )r  r   r   r   r  r  r  _ANALYSIS_CATEGORIESr  r   rE   r   r:   	functoolswrapsr   boolrc   r   r   r   rB   rC   r
  r  r   )r+   r5   rd   rj   r4   r   r   r  r  r   running_catr  r  _holderr  r   r   r   r   r!  r  _argsr   r'  session_usabler  r  r  s                            @@@r9   dispatch_category_runr1  t  sA   B 7.Xj\:,a}MN+;  / N ZL)K$&J+:,az :<='''9+'F'M#&N O89 
 	
 =$[1m# *S#/x09	! ! #'Lu++3B
JOL
 1)%=
 1Z@ -
H$
 *"j*EZL*QxjA		)	5 
*	5" 
  (/ LJ4LG w{{401%kk)X>%kk(B7%'; ; ($"8*K
|D	
 +;  d 	 W  ) 	+KKM 	+"NNNPT  +	+	1'JE /z<H 	9
 LLZL4  9	9 #K0 #K0 O)sm   &CF% 8%F% %I1GI#G(%I'G((I,H	H<	"H.+H<-H..H<1I<I		Ic           
        ddl m} t        D ]!  } |d|  d| d|       st        ddd       dd	l m}m}m d|   |d
      st        dd |      xs d dd      	 t        ||       }t        ||       }t        || |      }	d|  d| d}
t        j                  t              fd       }t        |
|| ||d||	fd|       }|st        ddd      	 dd||
dS # t        $ r
          w xY w)uR   씬 재분석(reanalyze-scenes) dispatch — scene_save + downstream force 실행.r   r   r   r   r  r  r   r_   )r  r  r  r   r  uC   이 에피소드는 이미 실행 중입니다 (먼저 잡은 쪽: r  rs   z
:reanalyzec                 B    	 t        |           y #         w xY wNr  )r  r  r  s    r9   r  z4dispatch_scene_reanalysis.<locals>._run_then_releasew  s"    1&#K0#K0s    
r]   zReanalyze scenes for episode r"  Tr'  r(  )r  r   r   r   r  r  r  rE   r   r:   r*  r+  r   r   r   )r+   r5   r4   r   r-  r  r  r   r   r   r!  r  r'  r  r  s                @@r9   dispatch_scene_reanalysisr5  M  st    7.Xj\:,a}MN+;  / 
 ZL)K[='''9+'F'M#&NaQ 
 	
0Z@0Z@)"j*EZL*Z@		)	1 
*	1 ($j(G^ 7
|D
 +;   	 	  K(s   )A8C) )C<queuedblocked_by_previousc           	         ddl m} |st        ddd      |j                  |      j	                  |j
                   k(  |j                  j                  |            j                  |j                        j                         }|D 	ch c]  }	|	j                   }
}	|D 	cg c]	  }	|	|
vs|	 }}	|rt        dd|d	d
  d      |D 	cg c]  }	|	j                   }}	t        |D ch c]  }|j                  |      dkD  s| c}      }|rt        dd| dd      |D 	cg c]  }	|	j                  dk(  s|	 }}	|rt        dd|d   j                   dd      |D 	cg c]  }	|	j                  |	j                  f c}	t        j                  t         j"                        j%                         }|D ]  }	t&        |	_        ||	_         |j+                          d  }d   fd}t-        ||dt/               d        }|sb|D ]?  }	d|	_        t        j                  t         j"                        j%                         |	_        A |j+                          t        ddd      ddD cg c]
  \  }}||d c}}d S c c}	w c c}	w c c}	w c c}w c c}	w c c}	w c c}}w )!u  여러 편을 **화수 순서대로 하나씩** 돌린다.

    ★``image_call_cap``·``approve_image_generation`` — 단일 실행과 **같은
    인자**를 받는다. 없으면 `category="all"` 은 착수 검사에서 선다
    (이미지 생성은 사람이 명시로 켜야 한다). 대기열이라고 그 문을 비켜
    가면 안 된다.

    ★★``image_call_cap`` 은 **화당** 상한이다 — 대기열 전체 합이 아니다.
     화마다 새 예산이 깔린다. N편을 넣으면 최대 `N × cap` 이다. 이 말을
     안 적으면 승인한 수와 실제로 살 수 있는 수가 갈린다.

    ★왜 순차인가 — 앞 화가 끝나야 `entity_canon` 이 서고, 그래야 다음 화가
     앞 화 명부를 보고 같은 신원으로 이어 붙인다(`episode_carry`). 나란히
     돌리면 둘 다 백지에서 시작해 같은 것을 두 신원으로 만든다. 게다가
     `short_id` 장부와 canon 이 프로젝트 범위라 동시 쓰기가 부딪힌다.

    ★다른 프로젝트끼리는 그대로 나란히 돈다 — 자리는 프로젝트마다 하나다.

    ★앞 화가 실패하면 뒤 화는 `blocked_by_previous` 로 **선다.** 자동으로
     넘어가면 앞 화 없이 만든 신원이 그대로 굳는다.
    r   r   zqueue.emptyu#   넣을 에피소드가 없습니다r^   r_   r   u)   이 프로젝트에 없는 에피소드: N   r   r   zqueue.duplicate_episode_numberu   화수가 겹칩니다 uU    — 순서를 정할 수 없습니다. 화수를 고친 뒤 다시 넣어 주세요.r   r  u*   이미 도는 중인 화가 있습니다 (u   화)r   z
run_queue:zqueue:c                 d   ddl m}  ddlm}m} d } ||      sJt
        j                  d       D ].  \  }} |        }	 t        ||t               |j                          0 yd}	 D ]  \  }} |        }	 |#t        ||t               	 |j                          3t        ||d		      }	|	xs i j                  d
      sH|}t
        j                  d||	xs i j                  dd      t        |	xs i j                  dd             |j                           	  ||       y# |j                          w xY w# t        $ rA}
|}t
        j                  d|t        |
d       t        ||t        |
             Y d}
~
vd}
~
ww xY w# |j                          w xY w#  ||       w xY w)uT   앞 화가 **끝난 것을 확인하고** 다음 화. ★배경 작업은 하나다.r   r   r  r   r  uF   프로젝트 대기열 %s: 자리를 못 잡아 시작하지 못했다NT)	r+   r5   rd   rj   r4   r  r  r   r   r   uX   프로젝트 대기열 %s: %d화가 %s 로 끝났다 — 남은 화는 %s 로 둔다: %sr   r  r   r   uQ   프로젝트 대기열 %s: %d화에서 멈춤 — 남은 화는 %s 로 둔다: %sr%  )app.core.databaser   r  r  r  rB   r   _mark_statusEPISODE_STATUS_BLOCKEDr   r1  rc   r   _mark_failed_if_still_queuedr1   )r   r  r  project_keyrM   _nr  
stopped_atnumberoutr   r   rd   claim_ownerr   rj   r   r+   s              r9   _drainz&dispatch_project_queue.<locals>._drain  s   2	A !- ![ALLa#%"R&.$ #/EFMMO # $(
'	-&V&."$!-$Wc3IJ > MMO= 0#-#!)#'{'51IC  I2??40%+
01;V YBOOHc:2SYBOOHb4Q	S" MMOI  'L  ,U MMO8 ! 
I!'JLLk"F,BC!% ! ' 1#s3xHH
I MMO,s[   D.,F% <EF% !A1EF% .E 	F7FFFFF""F% %
F/zRun z episodes in order for )r!  r#  r$  uploadedu;   이 프로젝트의 대기열이 이미 돌고 있습니다Tr6  )r5   episode_number)r   r   episodes)r&   r   r   r'   r(   r+   r)   in_order_byrG  r   rf   countr   r   r.   r   r/   	isoformatEPISODE_STATUS_QUEUED
updated_atr   r   r   )r+   episode_idsrd   rj   r4   r   r   r   rowsr   foundr   _numsn_dupesrunningr.   	queue_keyrE  r'  rM   rD  r   s   ` `` ``              @@r9   dispatch_project_queuerW    s   D +M3X#&( 	( 		""j0'**..2M	N	'((	)		 	   4aQTT4E %8+Q%q+G8$?}M 	
 (,,t!QtE,<1Q!);Q<=F1.vh 7E F	 	
 :$Q!((k"9q$G:'@AZAZ@[[_` 	
 266Aa&&'6G
,,x||
$
.
.
0C(  IIKZL)I:,'K;- ;-z $&3w<.(?
|LNG  A!AH#<<5??AAL  			2\#&( 	( ('.0'.VS! ),qA'.01 1Y !8 -< ; 7l0s<   J)	J!3J!J&1J+J+.J0J0.J5J:
c                   ddl m} 	 | j                  |      j                  |j                  |k(        j                         }||j                  t        k7  ryd|_        |dd |_        t        j                  t        j                        j                         |_        | j                          y# t         $ r+ | j#                          t$        j'                  d|d       Y yw xY w)	u\   착수 전에 터진 화를 `error` 로. ★이미 다른 표가 있으면 안 건드린다.r   r   Nr   i  u   실패 표 기록 실패 (%s)Tr%  )r&   r   r'   r(   r)   r*   r   rM  r   r   r.   r   r/   rL  rN  r   r   r   rB   rC   )r  r5   r   r   rU   s        r9   r>  r>  @  s     +
SmmG$++GJJ*,DEKKM;#**(==
#DS\!hll3==? S6
TRSs   AB1 AB1 11C%$C%c                   ddl m} 	 | j                  |      j                  |j                  |k(        j                         }|t        j                  d||       y||_        |xs1 t        j                  t        j                        j                         |_        | j                          t        j!                  d||       y# t"        $ r, | j%                          t        j'                  d||d       Y yw xY w)	uK   상태 한 칸만 바꾼다. 실패해도 대기열을 멈추지 않는다.r   r   NuB   에피소드 상태를 못 적었다 — 행이 없다 (%s → %s)u   에피소드 상태 %s → %su-   에피소드 상태 기록 실패 (%s → %s)Tr%  )r&   r   r'   r(   r)   r*   rB   r   r   r   r.   r   r/   rL  rN  r   r   r   r   rC   )r  r5   r   r.   r   rU   s         r9   r<  r<  R  s     +&mmG$++GJJ*,DEKKM; LL]#V-
FX\\ : D D F3ZH &F
TZ $ 	 	&&s   AC A(C 2C;:C;)r4   
OrmSessionr+   r1   r5   r1   returnDict[str, str])r4   rZ  r+   r1   r[  Dict[str, Any])r4   rZ  r+   r1   r[  r,  )
r4   rZ  r+   r1   r5   r1   rS   r1   r[  r,  )rg   	List[str]rh   r   ri   r]  rj   r1   rd   r1   r[  None)r4   rZ  r+   r1   rd   r1   r5   Optional[str]rj   r1   r[  r^  )r4   rZ  r+   r1   r[  r^  r4  )rS   r1   r+   r1   r5   r1   r4   rZ  r   r]  r   zOptional[Dict[str, str]])r4   rZ  r5   r1   r   r1   r[  r_  )r4   rZ  r5   r1   r   r1   r[  r_  )rd   r1   r   Optional[int]r   r,  r[  Optional[ImageCallBudget])r+   r1   r5   r1   r   r^  r   r1   r   r]  r   r\  r   rb  r[  r]  )r4   rZ  r+   r1   r5   r1   r[  r`  )r4   rZ  r5   r1   r  r`  r[  r_  )r5   r1   r  r`  r[  r_  )r+   r1   r5   r1   rd   r1   rj   r1   r4   rZ  r   ra  r   r,  r  r,  r  r`  r[  r]  )r+   r1   r5   r1   r4   rZ  r[  r]  )NF)r+   r1   rO  r^  rd   r1   rj   r1   r4   rZ  r   ra  r   r,  r[  r]  )r  rZ  r5   r1   r   r1   r[  r_  )
r  rZ  r5   r1   r   r1   r.   r`  r[  r_  )?__doc__
__future__r   r*  r?   loggingr2   r   r   typingr   r   r   r	   r
   sqlalchemy.ormr   rZ  r;  r   app.core.errorsr   app.core.image_call_budgetr   r   r   app.core.job_managerr   app.core.step_catalogr   r   r   app.services.checkpoint_syncr   	getLogger__name__rB   r   __annotations__r:   rE   rI   rV   r[   ro   r   r   r   r   r   r   r   r)  r  r
  r  r1  r5  rM  r=  rW  r>  r<  r   r;   r9   <module>rp     s|   #     ' 3 3 0 * $ 
 7 
 ?			8	$ ,U  T #144*7!! #!14!?B!	!,\\\ \ 	\
 
\65
5
5
 !5
 	5

 5
 
5
z !%QQQ Q
 Q Q QhD .2!!! ! 		!
 #! +!H #47	4 #-0	:9494 "94 #	94
 94F )-jjj j 	j
 #j !j &j jZ + EEE E 	EP   
	  
2 %)%*#'VVV V 	V
 	V "V #V V !V VrGGG 	G 	G\ !  /  %)%*\1\1\1 \1 	\1
 	\1 "\1 #\1 \1~S),S15S& '+&#&/3&r;   