ó
    ÒQœjì.  ã                  ó¤  • % S r SSKJr  SSKrSSKrSSKrSSKrSSKrSSKJ	r	  SSK
JrJrJrJrJrJrJrJr  SSKJr  SSKJr  \	" \5      R1                  S5      r\R4                  " 5        S	\R6                  " 5        3rS
rSrSrSr Sr!Sr"S\#S'   1 Skr$Sr%SS jr&SS jr'S S jr(S!S jr)S"S jr*S#S jr+S$S jr, " S S5      r-        S%S jr.g)&aY  Job records in Postgres.

An upload job and a batch run each live in the memory of the task running
them, as they always have; this module is the durable copy. The running
task writes every stage transition, every finished batch row, and a
heartbeat every few seconds. Any task can then answer a status poll or
serve a run's rows from the tables, and a task that dies leaves a record
whose heartbeat stops: a job with no heartbeat for ``STALE_SECONDS`` is
marked failed with that reason the next time anything reads it, instead of
sitting at its last stage forever.

The tables are created by ``ensure_tables`` at startup from ``jobs.sql``,
which is idempotent, so a running database gains them without a migration
step. Nothing here is imported by the batch runner; the runner takes a store
object with the same method names, and tests run it without one.
é    )ÚannotationsN)ÚPath)ÚAnyÚCallableÚDictÚIterableÚListÚLiteralStringÚOptionalÚcast)Úlogger)Údbzjobs.sqlÚ:éZ   é
   éÈ   )ÚdoneÚfailedz!the task running this job stopped)Újob_idÚfilenameÚtitleÚfolderÚfiling_typeÚyearÚperiodÚversionÚtrackÚreplaceÚ
replace_idÚreplace_keyÚdoc_keyÚcontent_hashÚdoc_idÚ
source_uriÚstageÚstagesÚerrorÚ
error_infoÚlogÚ
scan_pagesÚscan_pages_reasonztuple[str, ...]Ú_UPLOAD_COLUMNS>   r)   r&   r(   iÙ=o c                 óH  • [        [        [        R                  5       5      n [        R
                  " 5        oR                  5          UR                  S[        45        UR                  U 5        S S S 5        S S S 5        g ! , (       d  f       N= f! , (       d  f       g = f)Nz SELECT pg_advisory_xact_lock(%s))	r   r
   ÚJOBS_SQLÚ	read_textr   ÚconnectÚtransactionÚexecuteÚ_TABLES_LOCK)ÚsqlÚconns     Ú3/home/mande/repo/quber/src/quber/playground/jobs.pyÚensure_tablesr7   N   s]   € Ü
Œ}œh×0Ñ0Ó2Ó
3€CÜ	�ŠŒ˜×/Ñ/Õ1Ø�‰Ð7¼,¸ÔIØ�‰�SÔ÷  2�ˆ×1Õ1ú��ús#   ¸BÁ)BÁ1BÂ
B	ÂBÂ
B!c                ó\  • / n[          Hl  nU R                  U5      nUS:X  a  [        U=(       d    / 5      [        * S nU[        ;   a  Ub  [
        R                  " U5      OSnUR                  U5        Mn     SR                  [         5      nSR                  S/[        [         5      -  5      nSR                  S [         SS  5       5      nSU-   S-   U-   S	-   U-   S
-   n[        R                  " 5        nUR                  U/ UQ[        P75        SSS5        g! , (       d  f       g= f)zGWrite the job as it stands, log tail included, and stamp the heartbeat.r)   Nú, z%sc              3  ó.   #   • U  H  o S U 3v •  M     g7f)z = EXCLUDED.N© )Ú.0Úcs     r6   Ú	<genexpr>Úsave_upload.<locals>.<genexpr>d   s   é € Ð&ZÒFYÀ¨¨L¸¸Õ'<ÒFYùs   ‚é   z(INSERT INTO ade_playground.upload_jobs (z(, owner, heartbeat, updated_at) VALUES (z7, %s, now(), now()) ON CONFLICT (job_id) DO UPDATE SET z?, owner = EXCLUDED.owner, heartbeat = now(), updated_at = now())r,   ÚgetÚlistÚLOG_TAILÚ_JSON_COLUMNSÚjsonÚdumpsÚappendÚjoinÚlenr   r0   r2   ÚOWNER)	ÚjobÚvaluesÚcolumnÚvalueÚcolumnsÚplaceholdersÚupdatesÚqueryr5   s	            r6   Úsave_uploadrS   X   s   € à€Fß!ˆØ—‘˜“ˆØ�U‹?Ü˜Ÿ "Ó%¤x i jÐ1ˆEØ”]Ó"Ø).Ñ):”D—J’J˜uÔ%ÀˆEØ�‰�eÖñ "ð "ŸY™Y¤Ó7€GØ"&§)¡)¨T¨F´S¼Ó5IÑ,IÓ"J€LØ!ŸY™YÑ&ZÄoÐVWÐVXÑFYÓ&ZÓZ€Gà2Ø
ñ	à
4ñ	5ð ñ	ð Dñ		Dð
 ñ	ð Lñ	Lð 
ô 
�ŠŒ˜Ø�‰�UÐ,˜fÐ,¤eÑ,Ô-÷ 
�Žús   Ã8DÄ
D+c                ó
  • [        [        [        U SS95      n[         HC  nUR	                  U5      n[        U[        5      (       d  M+  [        R                  " U5      X'   ME     UR	                  S5      =(       d    / US'   U$ )NF)Ústrictr)   )	ÚdictÚzipr,   rD   rA   Ú
isinstanceÚstrrE   Úloads)ÚrowrK   rM   rN   s       r6   Ú_upload_from_rowr\   r   sf   € Ü
Œs”? C°Ñ6Ó
7€CßˆØ—‘˜“ˆÜ�eœS×!Ó!ÜŸ*š* UÓ+ˆC‹Kñ  ð —‘˜“×% 2€Cˆ�JØ€Jó    c                 ó4  • [         R                  " 5        n U R                  S[        [	        [
        5      [        45      R                  5       nSSS5        W H!  u  n[        R                  " SU[        5        M#     [        U5      $ ! , (       d  f       N@= f)zHFail every unfinished job whose heartbeat has stopped. Returns how many.a  UPDATE ade_playground.upload_jobs
               SET stage = 'failed', error = %s, updated_at = now()
               WHERE stage <> ALL(%s)
                 AND (heartbeat IS NULL OR heartbeat < now() - make_interval(secs => %s))
               RETURNING job_idNzupload job {} failed: {})r   r0   r2   ÚSTOPPED_REASONrB   ÚTERMINAL_STAGESÚSTALE_SECONDSÚfetchallr   ÚwarningrI   )r5   Úrowsr   s      r6   Úmark_stale_uploadsre   |   ss   € ä	�ŠŒ˜Ø�|‰|ð#ô
 œT¤/Ó2´MÐBó
÷ ‰(‹*ð 	÷ 
ó ‰	ˆÜ�ŠÐ1°6¼>ÖJñ äˆt‹9Ð÷ 
�ús   –9B	Â	
Bc                ó  • [        5         SR                  [        5      n[        R                  " 5        nUR                  SU-   S-   U 45      R                  5       nS S S 5        W(       a  [        U5      $ S $ ! , (       d  f       N"= f)Nr9   úSELECT z2 FROM ade_playground.upload_jobs WHERE job_id = %s)re   rH   r,   r   r0   r2   Úfetchoner\   )r   rO   r5   r[   s       r6   Úload_uploadri   Œ   sp   € ÜÔØ!ŸY™Y¤Ó7€GÜ	�ŠŒ˜Ø�l‰lØ˜ÑÐ"VÑVØˆIó
÷ ‰(‹*ð 	÷ 
ö
 %(Ô˜CÓ Ð1¨TÐ1÷ 
�ús   µ(A9Á9
Bc                 óH  • [        5         SR                  [        5      n [        R                  " 5        nUR                  SU -   S-   [        [        5      45      R                  5       nSSS5        W Vs/ s H  n[        U5      PM     sn$ ! , (       d  f       N*= fs  snf )z+Every upload job still running on any task.r9   rg   z7 FROM ade_playground.upload_jobs WHERE stage <> ALL(%s)N)
re   rH   r,   r   r0   r2   rB   r`   rb   r\   )rO   r5   rd   Úrs       r6   Úlive_uploadsrl   —   s‚   € äÔØ!ŸY™Y¤Ó7€GÜ	�ŠŒ˜Ø�|‰|Ø˜ÑÐ"[Ñ[Ü”/Ó"Ð$ó
÷ ‰(‹*ð 	÷ 
ñ
 *.Ó.ª AÔ˜QÖ©Ñ.Ð.÷ 
�üò
 /s   µ5BÁ6BÂ
Bc                ó®   • [        U 5      nU(       d  g [        R                  " 5        nUR                  SU45        S S S 5        g ! , (       d  f       g = f)NzNUPDATE ade_playground.upload_jobs SET heartbeat = now() WHERE job_id = ANY(%s)©rB   r   r0   r2   )Újob_idsÚidsr5   s      r6   Úheartbeat_uploadsrq   £   s:   € Ü
ˆw‹-€CÞØÜ	�ŠŒ˜Ø�‰Ø\ØˆFô	
÷ 
�Žúó   ©AÁ
Ac                  óv   • \ rS rSrSrSS jr            SS jrSS jrSS jrSS jr	SS jr
SS	 jrS
rg)Ú
BatchStoreé±   zHThe batch runner's durable side. Method names are the runner's contract.c           
     ó¾   • [         R                  " 5        nUR                  SXU[        R                  " U5      [
        45        S S S 5        g ! , (       d  f       g = f)NzÂINSERT INTO ade_playground.batch_runs (job_id, doc_id, want, questions, owner, heartbeat)
                   VALUES (%s, %s, %s, %s, %s, now())
                   ON CONFLICT (job_id) DO NOTHING)r   r0   r2   rE   rF   rJ   )Úselfr   r#   ÚwantÚ	questionsr5   s         r6   Úsave_runÚBatchStore.save_run´   s?   € Ü�ZŠZŒ\˜TØ�L‰Lð6ð  ¤t§z¢z°)Ó'<¼eÐDô	÷ �\Ž\ús   –/AÁ
Ac                óø   • [         R                  " 5        nUR                  SXX4Ub  [        R                  " U5      OS 45        UR                  SU(       a  SOSU45        S S S 5        g ! , (       d  f       g = f)Na  INSERT INTO ade_playground.batch_rows (job_id, index, question, error, response)
                   VALUES (%s, %s, %s, %s, %s)
                   ON CONFLICT (job_id, index) DO UPDATE
                   SET error = EXCLUDED.error, response = EXCLUDED.response, finished_at = now()zrUPDATE ade_playground.batch_runs SET failed = failed + %s, heartbeat = now(), updated_at = now() WHERE job_id = %sr@   r   )r   r0   r2   rE   rF   )rw   r   ÚindexÚquestionr'   Úresponser5   s          r6   Úsave_rowÚBatchStore.save_row½   sh   € ô �ZŠZŒ\˜TØ�L‰Lðdð  ÈÑI]´·²¸HÔ1EÐcgÐhôð �L‰Lð EÞ‘  FÐ+ô÷ �\Ž\ús   –AA+Á+
A9c                óˆ   • [         R                  " 5        nUR                  SU45        S S S 5        g ! , (       d  f       g = f)NzZUPDATE ade_playground.batch_runs SET running = false, updated_at = now() WHERE job_id = %s)r   r0   r2   )rw   r   r5   s      r6   Ú
finish_runÚBatchStore.finish_runÍ   s+   € Ü�ZŠZŒ\˜TØ�L‰LØlØ�	ô÷ �\Ž\ús	   –3³
Ac           
     ó¼  • [         R                  " 5        nUR                  S[        45      R	                  5       nU Hø  u  p4[        U[        5      (       a  [        R                  " U5      OUnUR                  SU45      R	                  5        Vs1 s H  nUS   iM
     nn[        [        U5      5       Vs/ s H  owU;  d  M
  UPM     nnU H  n	UR                  SX9XI   [        45        M      UR                  S[        U5      U45        [        R                  " SU[        [        U5      5        Mú     SSS5        gs  snf s  snf ! , (       d  f       g= f)z€A run whose task stopped: every unanswered question becomes an error
row saying so, and the run is closed, so its counts add up.z¡SELECT job_id, questions FROM ade_playground.batch_runs
                   WHERE running AND (heartbeat IS NULL OR heartbeat < now() - make_interval(secs => %s))z=SELECT index FROM ade_playground.batch_rows WHERE job_id = %sr   z�INSERT INTO ade_playground.batch_rows (job_id, index, question, error)
                           VALUES (%s, %s, %s, %s) ON CONFLICT DO NOTHINGzpUPDATE ade_playground.batch_runs SET running = false, failed = failed + %s, updated_at = now() WHERE job_id = %sz'batch run {} closed: {} ({} unanswered)N)r   r0   r2   ra   rb   rX   rY   rE   rZ   ÚrangerI   r_   r   rc   )
rw   r5   Ústaler   ry   rk   r   ÚiÚmissingr}   s
             r6   Úmark_stale_runsÚBatchStore.mark_stale_runsÔ   sO  € ô �ZŠZŒ\˜TØ—L‘LðmäÐ ó÷ ‰h‹jð	 ó
 &+Ñ!�Ü5?À	Ì3×5OÑ5OœDŸJšJ yÔ1ÐU^�	ð "Ÿ\™\ØWÐZ`ÐYbóç‘h“jð!óò!˜ð �a”Dñ!ð ð ô ',¬C°	«NÔ&;ÓMÒ&; È¹}Ÿ1Ñ&;�ÐMÛ$�EØ—L‘LðMà¨	Ñ(8¼.ÐIöñ %ð —‘ð GÜ˜“\ 6Ð*ôô —’Ø=¸vÄ~ÔWZÐ[bÓWcöñ' &+÷ ˆ\ùòùò N÷ �\ús1   –A=EÂEÂ"EÂ:	EÃEÃA-EÅ
EÅ
Ec                óf  • U R                  5         [        R                  " 5        nUR                  SU45      R	                  5       nUc
   SSS5        gUR                  SU45      R                  5       nSSS5        [        WS   [        5      (       a  [        R                  " US   5      OUS   nUUS   US   UUS   US   W Vs/ s HG  nUS   US   US   [        US   [        5      (       a  [        R                  " US   5      OUS   S	.PMI     snS
.$ ! , (       d  f       N­= fs  snf )z,The run and its rows as plain data, or None.z`SELECT doc_id, want, questions, failed, running FROM ade_playground.batch_runs WHERE job_id = %sNzgSELECT index, question, error, response FROM ade_playground.batch_rows WHERE job_id = %s ORDER BY indexé   r   r@   é   é   )r}   r~   r'   r   )r   r#   rx   ry   r   Úrunningrd   )
rŠ   r   r0   r2   rh   rb   rX   rY   rE   rZ   )rw   r   r5   Úrunrd   ry   rk   s          r6   Úload_runÚBatchStore.load_runô   s6  € à×ÑÔÜ�ZŠZŒ\˜TØ—,‘,ØrØ�	ó÷ ‰h‹jð ð ‰{Ø÷ ˆ\ð —<‘<ØyØ�	ó÷ ‰h‹jð ÷ ô +5°S¸±V¼S×*AÑ*A”D—J’J˜s 1™vÔ&ÀsÈ1Ávˆ	àØ˜!‘fØ˜‘FØ"Ø˜!‘fØ˜1‘vñ óò �Að ˜q™TØ ! !¡Ø˜q™TÜ4>¸qÀ¹tÄS×4IÑ4I¤§
¢
¨1¨Q©4Ô 0ÈqÐQRÉtô	ñ ññ
ð 	
÷ �\üò&s   ¦&DÁ!DÃ
AD.Ä
D+c                óÐ   • U R                  5         [        R                  " 5        nUR                  SU45      R	                  5       nS S S 5        US L$ ! , (       d  f       WS L$ = f)NzMSELECT 1 FROM ade_playground.batch_runs WHERE doc_id = %s AND running LIMIT 1)rŠ   r   r0   r2   rh   )rw   r#   r5   r[   s       r6   Úrunning_forÚBatchStore.running_for  s_   € Ø×ÑÔÜ�ZŠZŒ\˜TØ—,‘,Ø_ÐbhÐajóç‰h‹jð ÷ ð ˜$ˆÐ÷	 Œ\ð ˜$ˆÐús   ¦"AÁ
A%c                ó®   • [        U5      nU(       d  g [        R                  " 5        nUR                  SU45        S S S 5        g ! , (       d  f       g = f)NzMUPDATE ade_playground.batch_runs SET heartbeat = now() WHERE job_id = ANY(%s)rn   )rw   ro   rp   r5   s       r6   Ú	heartbeatÚBatchStore.heartbeat  s;   € Ü�7‹mˆÞØÜ�ZŠZŒ\˜TØ�L‰LØ_ÐbeÐagô÷ �\Ž\úrr   r;   N)
r   rY   r#   Úintrx   rY   ry   z	List[str]ÚreturnÚNone)r   rY   r}   rš   r~   rY   r'   zOptional[str]r   úOptional[Dict[str, Any]]r›   rœ   )r   rY   r›   rœ   ©r›   rœ   ©r   rY   r›   r�   )r#   rš   r›   Úbool©ro   zIterable[str]r›   rœ   )Ú__name__Ú
__module__Ú__qualname__Ú__firstlineno__Ú__doc__rz   r€   rƒ   rŠ   r’   r•   r˜   Ú__static_attributes__r;   r]   r6   rt   rt   ±   sZ   † ÙRôðØðØ"%ðØ14ðØ=JðØVnðà	ôô ôô@
ôB÷r]   rt   c                ól   ^ ^^• SUUU 4S jjn[         R                  " USSS9nUR                  5         U$ )z¶Beat, every ``HEARTBEAT_SECONDS``, for every job this process is still
running. The callables name them at each beat, so a job that finished
between beats is simply not beaten again.c                 óè   >•  [         R                  " [        5         [        T" 5       5        TR	                  T" 5       5        MC  ! [
         a!  n [        R                  " SU 5         S n A N(S n A ff = f)Nzheartbeat failed: {})ÚtimeÚsleepÚHEARTBEAT_SECONDSrq   r˜   Ú	Exceptionr   rc   )ÚexcÚrun_idsÚstoreÚ
upload_idss    €€€r6   ÚbeatÚstart_heartbeat.<locals>.beat3  sY   ø€ ØÜ�JŠJÔ(Ô)ð<Ü!¡*£,Ô/Ø—‘¡£	Ô*ñ	 øô
 ó <Ü—’Ð5°s×;Ñ;ûð<ús   ž&A Á
A1ÁA,Á,A1zjob-heartbeatT)ÚtargetÚnameÚdaemonrž   )Ú	threadingÚThreadÚstart)r±   r¯   r°   r²   Úthreads   ```  r6   Ústart_heartbeatr»   *  s1   ú€ ÷<ñ <ô ×Ò T°ÈÑM€FØ
‡L�L„NØ€Mr]   rž   )rK   úDict[str, Any]r›   rœ   )r[   ztuple[Any, ...]r›   r¼   )r›   rš   rŸ   )r›   zList[Dict[str, Any]]r¡   )r±   úCallable[[], Iterable[str]]r¯   r½   r°   rt   r›   zthreading.Thread)/r¦   Ú
__future__r   rE   ÚosÚsocketr·   rª   Úpathlibr   Útypingr   r   r   r   r	   r
   r   r   Úlogurur   Úquber.playgroundr   Ú__file__Ú	with_namer.   ÚgethostnameÚgetpidrJ   ra   r¬   rC   r`   r_   r,   Ú__annotations__rD   r3   r7   rS   r\   re   ri   rl   rq   rt   r»   r;   r]   r6   Ú<module>rÊ      sõ   ðòõ" #ã Û 	Û Û Û Ý ß U× UÓ Uå å á�‹>×#Ñ# JÓ/€ð ×ÒÓÐ
   "§)¢)£+ Ð/€à€ØÐ Ø€Ø$€Ø4€ð$€�ó ò2 0€ð €ôô.ô4ôô 2ô	/ô
÷sñ sðrØ+ðà(ðð ðð õ	r]   