
    Oj(+                    d   U d Z ddlmZ ddlZddlZddlZddlmZ ddlmZm	Z	 ddl
mZ ddlmZ  ej                  e      ZdZ ed	
       G d d             Zdddd	 	 	 	 	 	 	 	 	 	 	 ddZddd	 	 	 	 	 	 	 	 	 ddZddd	 	 	 	 	 	 	 	 	 ddZ G d d      Z eh d      Zded<   d dZg dZy)!uM  주행을 **멈추는 길** — 프로세스를 죽이지 않고 세우는 방법.

## 왜 있나

이 시스템에는 주행을 멈출 방법이 프로세스를 죽이는 것밖에 없었다. 그런데
프로세스를 죽이면 `step_run` 에 `running` 행이 남고, 락은 소유자의 죽음을
확인할 길이 없어 경과 시간(기본 3600초)을 기다렸다. **멈추는 길이 없어서
락이 망가지는 순환**이었다.

`step_run.cancel_requested_at` 은 **스텝 하나**를 세운다. 그것만으로는
부족하다 — `run_steps_batch` 는 `for sid in step_ids:` 로 도는지라
(`analysis_dispatch_service.py`), 지금 도는 스텝이 멈춰도 **다음 스텝으로
넘어간다.** 그래서 주행 전체를 덮는 표식이 따로 필요하다.

## 무엇을 하나

`run_cancel_request` 한 행이 「이 에피소드의 주행을 세워라」를 뜻한다.
읽는 자리는 셋이다:

1. **스텝 사이** — `run_steps_batch` 가 다음 `sid` 로 넘어가기 전
2. **단위 사이** — 팬아웃 스텝이 샷 하나를 끝낼 때 (`checkpoint_gate`)
3. **돈 쓰기 직전** — worker 가 provider 를 부르기 전

3번이 중요하다. 팬아웃은 작업을 executor 에 **한꺼번에 제출**하므로,
대기열에 이미 들어간 작업은 부모가 멈춰도 시작된다. 그 작업들이 유료
호출 앞에서 스스로 물러나야 돈이 안 나간다.

★멈춤은 **협조적**이다. 진행 중인 provider 호출을 중간에 끊지 않는다 —
 이미 산 것은 결과까지 가져온다. 끊어 봐야 돈은 이미 나갔고 결과만 잃는다.
    )annotationsN)	dataclass)datetimetimezone)Optional)texta  
CREATE TABLE IF NOT EXISTS run_cancel_request (
    id text PRIMARY KEY,
    project_id text NOT NULL,
    episode_id text NOT NULL,
    scope text NOT NULL DEFAULT 'episode',
    requested_at timestamptz NOT NULL DEFAULT CURRENT_TIMESTAMP,
    requested_by text,
    reason text,
    cleared_at timestamptz,
    CONSTRAINT run_cancel_request_scope_key UNIQUE(project_id, episode_id, scope)
)
T)frozenc                  N    e Zd ZU dZded<   dZded<   dZded<   dZded	<   dd
Zy)CancelStateu:   지금 이 에피소드에 정지 요청이 걸려 있나.bool	requestedNOptional[datetime]requested_atzOptional[str]requested_byreasonc                    | j                   sy| j                  xs d}| j                  r| j                  j                         nd}| j                  xs d}d| d| d| S )Nu   정지 요청 없음?u   (사유 없음)u   정지 요청 — u    가 u    에, 사유: )r   r   r   	isoformatr   )selfwhowhenwhys       J/Users/manta/Documents/Projects/TheRoad-I1/backend/app/core/run_control.pydescribezCancelState.describeF   sb    ~~)&3040A0At  **,skk..#C5dV>#GG    )returnstr)	__name__
__module____qualname____doc____annotations__r   r   r   r    r   r   r   r   =   s/    DO'+L$+"&L-& FM Hr   r    episode)r   r   scopec          
        | j                  t        d      t        t        j                               ||||xs d|xs dd       | j                          t        j                  d|||xs d|xs d       t        | |||      S )u   정지를 요청한다. 이미 걸려 있으면 사유만 갱신한다 (멱등).

    ★시각은 DB 시계로 찍는다 — 읽는 쪽과 같은 시계여야 한다.
    a  
        INSERT INTO run_cancel_request (
            id, project_id, episode_id, scope,
            requested_at, requested_by, reason, cleared_at
        )
        VALUES (
            :id, :pid, :eid, :scope,
            CURRENT_TIMESTAMP, :by, :reason, NULL
        )
        ON CONFLICT (project_id, episode_id, scope) DO UPDATE
          SET requested_at = CURRENT_TIMESTAMP,
              requested_by = EXCLUDED.requested_by,
              reason = EXCLUDED.reason,
              cleared_at = NULL
    N)idpideidr&   byr   u@   [RUN-CANCEL] 정지 요청 project=%s episode=%s by=%s reason=%sr   -)r&   )	executer   r   uuiduuid4commitloggerwarningread_cancel_state)db
project_id
episode_idr   r   r&   s         r   request_cancelr7   O   s     JJt  	 $**,#t>T, IIK
NNJJ 3V]s RZuEEr   )r&   expected_requested_atc                   |||d}d}|d}||d<   | j                  t        d| d      |      }| j                          t        |j                        }|rt
        j                  d||       |S )u  정지 요청을 내린다 — 다시 주행할 수 있게 한다.

    ★행을 지우지 않고 `cleared_at` 을 찍는다. 언제 세웠고 언제 풀었는지가
     남아야 나중에 읽을 수 있다.

    Args:
        expected_requested_at: 본 요청의 시각. 주면 **그 요청일 때만** 내린다.

    ★`expected_requested_at` 이 왜 있나 — 조회하고 내리기까지 사이에 누가
     **새 정지**를 요청했을 수 있다. 확인하지 않으면 내가 본 옛 요청을 푼다는
     것이 남의 새 요청까지 지운다. 「멈춰라」가 조용히 사라지는 것은
     「멈추지 못하는 것」과 같다.

    Returns: 실제로 내린 것이 있으면 True.
    r)   r*   r&   r$   zAND requested_at = :expected_atexpected_atz
        UPDATE run_cancel_request
           SET cleared_at = CURRENT_TIMESTAMP
         WHERE project_id = :pid AND episode_id = :eid AND scope = :scope
           AND cleared_at IS NULL
           z
    u7   [RUN-CANCEL] 정지 요청 해제 project=%s episode=%s)r-   r   r0   r   rowcountr1   info)	r4   r5   r6   r&   r8   params
cas_clauseresultcleareds	            r   clear_cancelrB   z   s    . &j5IFJ(6
 5}ZZ "
 <  	 F IIK6??#GE
	
 Nr   F)r&   strictc               J   	 | j                  t        d      |||d      j                         }|t        d      S t        d|j                  |j                  |j                        S # t        $ r/}|r t        j                  d|       t        d      cY d}~S d}~ww xY w)	u  지금 정지 요청이 걸려 있나.

    Args:
        strict: 읽기 실패를 **올릴지**. 기본 False.

    ★기본값은 읽기가 실패하면 **걸려 있지 않다**고 답한다. 표가 아직 없거나
     DB 가 잠깐 끊긴 것으로 정상 주행이 멈추는 쪽이 더 나쁘기 때문이다.

    ★그러나 **돈을 쓰기 직전**에는 이 관용이 거짓말이 된다. 호출자가
     「fail-closed 로 확인했다」고 믿는데 여기서 조용히 「정지 없음」을
     돌려주면, DB 장애 동안 정지가 통째로 안 들리면서 지출은 계속된다.
     그 자리에서는 `strict=True` 로 불러 실패를 올린다.
    z
            SELECT requested_at, requested_by, reason
              FROM run_cancel_request
             WHERE project_id = :pid AND episode_id = :eid AND scope = :scope
               AND cleared_at IS NULL
        r:   u9   정지 요청 조회 실패 (없는 것으로 본다): %sF)r   NT)r   r   r   r   )
r-   r   fetchone	Exceptionr1   debugr   r   r   r   )r4   r5   r6   r&   rC   rowexcs          r   r3   r3      s    *,jj  
 !eDF
 GOhj 	 {U++%%%%zz	   ,PRUVU++	,s   -A* *	B"3$BB"B"c                  R    e Zd ZdZdd	 	 	 	 	 	 	 ddZddddZddddZddd	Zy
)CancellationTokenu  worker 가 유료 호출 직전에 물어보는 표.

    팬아웃은 작업을 한꺼번에 제출하므로, 부모가 멈춰도 대기열의 작업은
    시작된다. 그 작업들이 **돈을 쓰기 전에** 이 표를 보고 물러나야 한다.

    ★DB 를 매번 때리지 않는다. `ttl_seconds` 동안 답을 재사용한다 —
     샷 하나가 수십 초 걸리므로 몇 초 지연은 무해하고, 수십 개 worker 가
     동시에 조회하는 것은 유해하다.

    ★세션을 공유하지 않는다. 스레드마다 다른 세션이 필요하므로
     `session_factory` 를 받아 조회할 때마다 열고 닫는다.
    g      @)ttl_secondsc                   || _         || _        || _        || _        t	        j
                         | _        d | _        d | _        y N)	_session_factory_project_id_episode_id_ttl	threadingLock_lock_cached
_cached_at)r   session_factoryr5   r6   rL   s        r   __init__zCancellationToken.__init__   s?     !0%%	^^%
.2.2r   Fforcec                  t        j                  t        j                        }| j                  5  | xrJ | j
                  d uxr: | j                  d uxr* || j                  z
  j                         | j                  k  }|r!| j
                  | j
                  cd d d        S d d d        | j                         }	 t        || j                  | j                        }|j                          | j                  5  || _        || _        d d d        |S # 1 sw Y   pxY w# |j                          w xY w# 1 sw Y   |S xY wrN   )r   nowr   utcrU   rV   rW   total_secondsrR   rO   r3   rP   rQ   close)r   r[   r]   freshsessionstates         r   rc   zCancellationToken.state   s   ll8<<(ZZ	 HLL,HOO4/H 4??*99;diiG	  1|| ZZ '')	%gt/?/?AQAQREMMOZZ DL!DO  ' Z MMO s$   A)D<!D 9D3DD03D=c               :    | j                  |      j                  S )NrZ   )rc   r   )r   r[   s     r   is_cancelledzCancellationToken.is_cancelled	  s    zzz&000r   c                    | j                         }|j                  syddlm}  |d|xs d d|j	                          d      )	ub   멈추라는 말이 왔으면 즉시 멈춘다.

        Raises: AppError(step.cancelled)
        Nr   )AppErrorstep.cancelledu   주행u    정지 — i  )codemessagestatus_code)rc   r   app.core.errorsrg   r   )r   whererc   rg   s       r   raise_if_cancelledz$CancellationToken.raise_if_cancelled  sK    
 

,!()enn6F5GH
 	
r   N)r5   r   r6   r   rL   floatr   None)r[   r   r   r   )r[   r   r   r   )r$   )rm   r   r   rp   )r   r   r    r!   rY   rc   re   rn   r#   r   r   rK   rK      sR    & !3 3 	3 3 
3  &+ . -2 1
r   rK   >   step.owner_loststep.gate_unreadablerh   	frozensetABORT_CODESc                (    t        | dd      t        v S )uc   이 예외가 **주행을 세워야 하는 것**인가. ★code 로 본다 — 문구가 아니라.ri   r$   )getattrrt   )rI   s    r   is_abortrw   *  s    3#{22r   )CREATE_TABLE_SQLr   rK   r7   rB   r3   rt   rw   )r5   r   r6   r   r   r   r   r   r&   r   r   r   )
r5   r   r6   r   r&   r   r8   r   r   r   )
r5   r   r6   r   r&   r   rC   r   r   r   )rI   BaseExceptionr   r   )r!   
__future__r   loggingrS   r.   dataclassesr   r   r   typingr   
sqlalchemyr   	getLoggerr   r1   rx   r   r7   rB   r3   rK   rs   rt   r"   rw   __all__r#   r   r   <module>r      s^  > #    ! '  			8	$  $H H H, (F(F (F
 (F (F (F (F` 04++ +
 + .+ 
+f )) )
 ) ) )XF
 F
^ # $ Y 3
	r   