+
    qj>G                        R t ^ RIt^ RIt^ RIt^ RIt^ RIHtHt ^ RIHt ^ RI	H
t
 ^ RIHt ^ RIHt ^ RIHt ^ RIHtHt ^ R	IHt ^ R
IHtHt ^ RIHtHtHt ]P:                  ! ]4      t ! R R] 4      t!R R lt"R R lt#R^ RRR^ /R R llt$RRRRR^ RR/R R llt%R R lt&R*R R llt'R+R  R! llt(^t)R,R" R# llt*R$RR%R/R& R' llt+R-R( R) llt,R# ).ur  Recolha do registo comercial: histórico por NIPC e expansão pelo grafo.

Duas vias, com objectivos diferentes e que se complementam:

* **Por NIPC** (aqui). O `publications_for_nif` do IRN é ilimitado no tempo,
  portanto dá o **histórico completo** de uma entidade — incluindo a constituição
  de 2008 que nenhum varrimento dos últimos anos encontraria. É a via das
  entidades que nos interessam: as monitorizadas e, em BFS, cada empresa
  descoberta como titular ou participada.
* **Nacional, dia a dia** (`registry_sweep.py`). A única forma de responder à
  pergunta inversa — "em que outras empresas é que esta pessoa é sócia" — porque
  o IRN **não** permite pesquisar pelo NIF de uma pessoa: com o NIF de um
  indivíduo devolve zero, medido.

As pessoas nunca se expandem por NIPC, por essa mesma razão. O BFS segue apenas
sócios que sejam entidades.
N)datedatetime)Any)text)AsyncSession)settings)RateLimiter)IrnVersionChangedMjIrnClient)html_to_text)	parse_actwants_content)apply_parseenqueue_entity
upsert_actc                       ] tR t^&tRtRtR# )RateLimiteduD   O IRN respondeu 429/403. Para o pool inteiro, não só quem apanhou. N)__name__
__module____qualname____firstlineno____doc____static_attributes__r       app/services/registry_sync.pyr   r   &   s    Nr   r   c                L    V ^8  d   QhR\         R,          R\        R,          /# )   valueNreturn)strr   )formats   "r   __annotate__r"   *   s"      sTz dTk r   c                     V '       g   R #  \         P                  ! V R,          R4      P                  4       #   \         d     R # i ; i)N:N
   Nz%Y-%m-%d)r   strptimer   
ValueError)r   s   &r   _parse_dater'   *   s?      sZ8==?? s   +8 AAc                    V ^8  d   QhR\         R\        R\        R\        R\        R\        \        \
        3,          R\        \        \        3,          R,          /# )	r   sessionclientlimiteract_idnipcpubr   N)r   r
   r   r    dictr   int)r!   s   "r   r"   r"   3   sc     0 000 0
 0 0 
c3h0 
#s(^d0r   c               R  "   VP                  ^4      G Rj  xL
  VP                  V4      G Rj  xL
 pV'       d   \        V4      MRpV'       g(   V P                  \	        R4      RV/4      G Rj  xL
  R# \        VP                  R4      4      pV P                  \	        R4      RVRVR\        P                  ! VP                  4       4      P                  4       /4      G Rj  xL
  \        WuP                  R4      4      p	\        WWIVR	7      G Rj  xL
 #  EL L L L; L5i)
uC   Abre o conteúdo de um acto e grava nós e arestas. 3 pedidos HTTP.Nz
                UPDATE registry_acts
                   SET content_status = 'empty', content_attempts = content_attempts + 1,
                       content_fetched_at = now()
                 WHERE id = CAST(:aid AS UUID)
                aidDatea  
            UPDATE registry_acts SET
                body = :body,
                body_sha256 = :sha,
                body_source = 'irn',
                body_truncated = FALSE,
                content_status = 'fetched',
                content_fetched_at = now(),
                content_attempts = content_attempts + 1,
                parse_error = NULL
            WHERE id = CAST(:aid AS UUID)
            bodyshaAct_Fact)r,   r-   parsedact_date)acquirecontent_forr   executer   r'   gethashlibsha256encode	hexdigestr   r   )
r)   r*   r+   r,   r-   r.   htmlr4   r8   r7   s
   &&&$$$    r   _fetch_and_parserB   3   s     //!
##C((D!%<4Doo FO

 
	
 
	
 3776?+H
//	
 
eW^^DKKM-J-T-T-VW  " tWWZ01FT8  I (
	
$sa   D'DD'DD'D' D'2D!3A4D''D#(/D'D%D'D'!D'#D'%D'priorityfetch_contentTdepthc                    V ^8  d   QhR\         R\        R\        R\        R\        R\
        R\        R\        \        \        3,          /# )	r   r)   r*   r+   r-   rC   rD   rE   r   )r   r
   r   r    r0   boolr/   )r!   s   "r   r"   r"   f   s`        	    
#s(^r   c          
        "   VP                  4       G Rj  xL
  VP                  V4      G Rj  xL
 p\        WVWWER7      G Rj  xL
 #  L4 L L5i)u	  Histórico completo de um NIPC, com conteúdo dos actos que trazem pessoas.

Só se abre o conteúdo do que ainda não temos: a listagem é uma chamada, o
conteúdo são três. Prestações de contas — mais de metade do volume de uma
entidade — nunca se abrem.
N)r*   r+   rC   rD   )r9   publications_for_nifingest_publications)r)   r*   r+   r-   rC   rD   rE   pubss   &&&&$$$ r   	sync_nipcrL   f   sT       //
,,T22D$t   2s1   AAAAAAAAAr*   r+   c                    V ^8  d   QhR\         R\        R\        \        \        \        3,          ,          R\
        R,          R\        R,          R\        R\        R	\        \        \        3,          /# )
r   r)   r-   rK   r*   Nr+   rC   rD   r   )	r   r    listr/   r   r
   r   r0   rG   )r!   s   "r   r"   r"      s     Y YY
Y tCH~
Y
 $Y 4Y Y Y 
#s(^Yr   c                 "   R\        V4      R^ R^ R^ R^ /pV'       g(   V P                  \        R4      RV/4      G Rj  xL
  V# V P                  \        R	4      RV/4      G Rj  xL
 P                  4        Uu0 uF  pV^ ,          '       g   K  V^ ,          kK  	  p	pV EFc  p
\	        V
P                  R
4      4      p\        V VVV
P                  R4      \        V
P                  R4      4      V
P                  R4      V
V
P                  R4      V
P                  R4      VR7
      G Rj  xL
 w  rV'       g   K  VR;;,          \        V4      ,          uu&   V'       d	   Ve   Vf   K  W9   d   K  \        V
P                  R4      4      '       g   K  \        WWLWR7      G Rj  xL
 pVR;;,          ^,          uu&   V'       g   EK+  VR;;,          VR,          ,          uu&   VR;;,          VR,          ,          uu&   EKf  	  V P                  \        R4      RV/4      G Rj  xL
  V#  EL ELu upi  EL L L5i)u  Grava a listagem de um NIPC. Não vai à rede se `fetch_content` for falso.

Está separada do `sync_nipc` porque a listagem passou a poder chegar de duas
origens: o CT, que a vai buscar ele próprio, e o trabalhador remoto, que a
entrega já buscada. O que é escasso no IRN é o IP, não o processador — mas o
que escreve na base de dados tem de ser um só, senão duas versões da mesma
regra divergem em silêncio na primeira alteração. É a mesma razão pela qual o
parser nunca saiu daqui.
listadosnovos	conteudosquotascargosz]UPDATE registry_entities SET history_fetched_at = now(), updated_at = now() WHERE nif = :nipcr-   NzTSELECT irn_publication_id FROM registry_acts WHERE nipc = :nipc AND body IS NOT NULLIdr6   r3   EntityDistrictCouncil)	r-   irn_idact_typer8   entity_namelistingdistrictcouncilrC   )r,   r-   r.   officersz
            UPDATE registry_entities SET history_fetched_at = now(), updated_at = now()
             WHERE nif = :nipc
            )lenr;   r   allr    r<   r   r'   r0   r   rB   )r)   r-   rK   r*   r+   rC   rD   statsrknown_bodiesr.   pub_idr,   is_newcountss   &&&$$$$        r   rJ   rJ      s:    & TGQQ!XWXYE oo8 TN
 	
 	
 
 //?   #%A Q44 	!   SWWT]#)WWZ( 1)WWZ(GGI& 
 
 g#f+% '/!SWWZ011'W$
 
 	ka6(Ovh//O(Ovj11OA D //	
 
   LE	
 
0
st   >I H2&I'H5(I<H8H8BI-H=.A;I)I *IAI+I,I5I8I IIc                0    V ^8  d   QhR\         R\        /# )r   r)   r   r   r0   )r!   s   "r   r"   r"      s      , 3 r   c                   "   V P                  \        R4      4      G Rj  xL
 P                  4       p^ pV F<  w  r4T\        \	        WP                  4       RV 2^ R7      G Rj  xL
 4      ,          pK>  	  V#  LX L5i)ul  Põe na fila **todas** as empresas monitorizadas, os cinco tipos.

O `run_mj_irn` corre só sobre internal/competitor/client, e é por isso que as
7 de análise têm `registry_fetched_at` a NULL — zero informação de registo — e
que as 17 `related` (descobertas como sócias) nunca são sincronizadas, o que
trava o grafo no primeiro hop. Aqui entram as 205.
aY  
                SELECT nif, monitoring_type FROM companies
                 WHERE active = TRUE AND length(trim(nif)) = 9
                 ORDER BY CASE monitoring_type
                            WHEN 'internal' THEN 0 WHEN 'client' THEN 1
                            WHEN 'competitor' THEN 2 WHEN 'analysis' THEN 3 ELSE 4 END
                Nz
monitored:)nifreasonrE   )r;   r   ra   r0   r   strip)r)   rowsaddedrk   mtypes   &    r   seed_monitoredrq      s      oo

 
	
 
ce 	 E
>'yy{ZX]W^K_ghiijj L
	
 js"   A=A9AA=$A;%A=;A=c                <    V ^8  d   QhR\         R\        R\        /# )r   r)   rE   r   ri   )r!   s   "r   r"   r"      s!        l  3  s  r   c                   "   V P                  \        R4      RV/4      G Rj  xL
 pVP                  ;'       g    ^ #  L5i)u  Põe na fila **todas** as entidades colectivas do país sem histórico próprio.

A listagem diária do IRN não é um índice completo: numa amostra de 30
empresas, **28 tinham actos que nunca vimos** — 379 no total, 83 deles de
estrutura. A JOIN THE MOMENT tinha 4 actos no nosso índice e 30 na fonte.

Entra a `depth = 9`, muito acima dos 0–3 da vizinhança, porque o
`take_from_queue` ordena por `depth, requested_at`: as monitorizadas e os
seus sócios continuam a ser servidos primeiro e o resto do país vem atrás
sem lhes passar à frente. O `DO NOTHING` preserva quem já lá esteja mais
perto.

Uma só instrução, não 507 mil chamadas ao `enqueue_entity`.
a  
            INSERT INTO registry_entity_fetch (nif, depth, reason)
            SELECT e.nif, :depth, 'nacional'
              FROM registry_entities e
             WHERE e.merged_into IS NULL
               AND e.nif IS NOT NULL
               AND e.kind IN ('company', 'public', 'foreign')
               AND e.history_fetched_at IS NULL
            ON CONFLICT (nif) DO NOTHING
            rE   N)r;   r   rowcount)r)   rE   results   && r   seed_all_companiesrv      sJ      ??		
 
% F ??as   !A >A A c                J    V ^8  d   QhR\         R\        R,          R\        /# )r   r)   	max_depthNr   ri   )r!   s   "r   r"   r"     s&     2 2 2t 2_b 2r   c                @  "   Ve   TM\         P                  pV P                  \        R4      RV/4      G Rj  xL
 P	                  4       p^ pV FD  w  rEpT\        \        WP                  4       R\        V4      VR7      G Rj  xL
 4      ,          pKF  	  V#  L` L5i)u"  Alarga a fila às empresas ligadas ao que já temos.

Só entidades: um sócio pessoa não se expande porque o IRN não pesquisa por
NIF de pessoa. A profundidade de uma entrada nunca aumenta — um NIPC pedido
como cliente não passa a hop 2 por aparecer também como sócio de um sócio.
NaH  
                WITH known AS (
                    SELECT f.nif, f.depth FROM registry_entity_fetch f
                )
                SELECT DISTINCT other.nif, min(k.depth) + 1 AS depth, min(k.nif) AS seed
                  FROM registry_edges e
                  JOIN registry_entities subj ON subj.id = e.subject_id
                  JOIN registry_entities other ON other.id = e.holder_id
                  JOIN known k ON k.nif = subj.nif
                 WHERE e.is_current
                   AND other.nif IS NOT NULL
                   AND other.kind IN ('company','public','foreign')
                   AND other.history_fetched_at IS NULL
                   AND k.depth < :maxd
                 GROUP BY other.nif
                UNION
                SELECT DISTINCT subj.nif, min(k.depth) + 1, min(k.nif)
                  FROM registry_edges e
                  JOIN registry_entities holder ON holder.id = e.holder_id
                  JOIN registry_entities subj ON subj.id = e.subject_id
                  JOIN known k ON k.nif = holder.nif
                 WHERE e.is_current
                   AND subj.nif IS NOT NULL
                   AND subj.kind IN ('company','public','foreign')
                   AND subj.history_fetched_at IS NULL
                   AND k.depth < :maxd
                 GROUP BY subj.nif
                maxdzsocio-empresa)rk   rl   rE   seed_nif)r   REGISTRY_MAX_DEPTHr;   r   ra   r0   r   rm   )r)   rx   rn   ro   rk   rE   seeds   &&     r   enqueue_corporate_neighboursr~     s      '2	8S8SIoo< Y? 
  	
B 
ceE 	F E D YY[E
]a 
 	
 ! LS 	
Js"   7BBABBBBc          	      t    V ^8  d   QhR\         R\        R\        \        \        \
        3,          ,          /# )r   r)   limitr   )r   r0   rN   r/   r    r   )r!   s   "r   r"   r"   Q  s4     %P %P< %P %PT$sTWx.EY %Pr   c           
       "   V P                  \        R\         R24      RV/4      G Rj  xL
 P                  4       pV Uu. uF.  pRV^ ,          P	                  4       RV^,          RV^,          /NK0  	  up#  LMu upi 5i)u  Tira trabalho da fila, do mais perto do nosso universo para o mais longe.

`FOR UPDATE SKIP LOCKED` para dois trabalhadores nunca pegarem no mesmo NIPC.

**E recupera reclamações penduradas.** Enquanto só o CT consumia esta fila,
em processo único, um restart levava tudo consigo e não havia nada preso.
Com o trabalhador remoto pela frente — noutra máquina, noutra linha — um
processo que morra deixaria entidades em `running` para sempre, invisíveis à
fila. Com 507 mil na fila, ninguém daria por elas.
aX  
                WITH picked AS (
                    SELECT nif FROM registry_entity_fetch
                     WHERE attempts < 3
                       AND (status IN ('pending','error')
                            OR (status = 'running'
                                AND claimed_at
                                    < now() - interval 'a   minutes'))
                     ORDER BY depth, requested_at
                     FOR UPDATE SKIP LOCKED
                     LIMIT :lim
                )
                UPDATE registry_entity_fetch f
                   SET status = 'running', attempts = f.attempts + 1,
                       claimed_at = now()
                  FROM picked p
                 WHERE f.nif = p.nif
                RETURNING f.nif, f.depth, f.reason
                limNrk   rE   rl   )r;   r   STALE_CLAIM_MINUTESra   rm   )r)   r   rn   rc   s   &&  r   take_from_queuer   Q  s      oo9 :M8M N* EN-
 	
0 
ce3 	4 KOO$QUAaDJJL'1Q41Q4@$OO3	
2 Ps!   )B A9B 4A;6B ;B rb   errorc                    V ^8  d   QhR\         R\        R\        R\        \        \        3,          R,          R\        R,          RR/# )r   r)   rk   statusrb   Nr   r   )r   r    r/   r0   )r!   s   "r   r"   r"   y  sO       #03<@cNT<Q: 
r   c                   "   T P                  \        R 4      RTRTRTRT;'       g    / P                  R4      RT;'       g    / P                  R4      /4      G Rj  xL
  R#  L5i)	ao  
            UPDATE registry_entity_fetch SET
                status = :status,
                acts_seen = COALESCE(:seen, acts_seen),
                acts_new = COALESCE(:new, acts_new),
                last_error = :err,
                history_fetched_at = CASE WHEN :status = 'ok' THEN now() ELSE history_fetched_at END
             WHERE nif = :nif
            rk   r   errseenrP   newrQ   N)r;   r   r<   )r)   rk   r   rb   r   s   &&$$$r   finish_queue_itemr   y  sj      //
	
 3&%U[[b%%j1EKKR$$W-	
  s   A A&A&A$A&c                    V ^8  d   QhR\         R\        R\         R,          R\        R\        \        \         3,          /# )r   	max_itemsmax_secondsworkersNrD   r   )r0   floatrG   r/   r    )r!   s   "r   r"   r"     sM     M MMM 4ZM 	M
 
#s(^Mr   c                  a aaaaaaa	a
"   ^ RI Ho T;'       g    \        P                  p\	        4       o\
        P                  ! 4       V,           oR^ R^ R^ R^ R^ R^ R^ /o	\        P                  ! 4       o\        P                  ! 4       oR	 VVVVVV VV	3R
 llo
\        P                  ! V
3R l\        V4       4       !  G Rj  xL
  SP                  S	R&   S	#  L5i)un  Consome a fila com um pool de trabalhadores e um tecto de ritmo comum.

O orçamento é por tempo de parede **e** por contagem, o que chegar primeiro:
a contagem sozinha mente, porque o custo de um NIPC varia entre uma chamada
(uma entidade com dois actos) e trinta (uma com histórico longo).

Com `fetch_content=False` só se lista, o que muda a ordem de grandeza: uma
chamada por empresa em vez de uma mais três por cada acto que interesse.
Medido, 102 empresas/min — o país inteiro em ~83 horas em vez de meses. Os
actos ficam na fila de conteúdo e são abertos depois, pela escada de
prioridades que já existe.
)AsyncSessionLocal	entidadesrP   rQ   rR   rS   rT   errosc                (    V ^8  d   QhR\         RR/# )r   indexr   N)r0   )r!   s   "r   r"   #process_queue.<locals>.__annotate__  s     -- --C --D --r   c                 ^	  <"   \        4       ;_uu_4       GR j  xL
 pSP                  4       '       Eg   \        P                  ! 4       S	8  Ed   S;_uu_4       GR j  xL
  SR,          S8  d$    R R R 4      GR j  xL
  R R R 4      GR j  xL
  R # R R R 4      GR j  xL
  S! 4       ;_uu_4       GR j  xL
 p\	        V^R7      G R j  xL
 pVP                  4       G R j  xL
  R R R 4      GR j  xL
  X'       g    R R R 4      GR j  xL
  R # V^ ,          p S! 4       ;_uu_4       GR j  xL
 p\        Y!SVR,          VR,          ^ 8X  d   ^ M^
VR,          S
R7      G R j  xL
 p\        W$R,          RVR7      G R j  xL
  VP                  4       G R j  xL
  R R R 4      GR j  xL
  S;_uu_4       GR j  xL
  SR;;,          ^,          uu&   R F(  pSV;;,          XP                  V^ 4      ,          uu&   K*  	  R R R 4      GR j  xL
  EK  R R R 4      GR j  xL
  R #  EL EL EL EL EL  + GR j  xL 
 '       g   i     EL; i EL EL EL{ ELn  + GR j  xL 
 '       g   i     EL; i ELu ELT EL  EL L L  + GR j  xL 
 '       g   i     L; i L L  + GR j  xL 
 '       g   i     EK  ; i  \         d    S! 4       ;_uu_4       GR j  xL 
 p\        Y$R,          RR	R
7      G R j  xL 
  TP                  4       G R j  xL 
  R R R 4      GR j  xL 
  M  + GR j  xL 
 '       g   i     M; iSP                  4        h \         Ed   p\        P                  RTR,          T4       S! 4       ;_uu_4       GR j  xL 
 p\        Y$R,          R\        T4      R,          R
7      G R j  xL 
  TP                  4       G R j  xL 
  R R R 4      GR j  xL 
  M  + GR j  xL 
 '       g   i     M; iS;_uu_4       GR j  xL 
  SR;;,          ^,          uu&   R R R 4      GR j  xL 
   R p?EKD    + GR j  xL 
 '       g   i      R p?EKc  ; iR p?ii ; i ELf  + GR j  xL 
 '       g   i     R # ; i5i)Nr   )r   rk   rE   )rC   rE   rD   ok)r   rb   pendingr	   )r   r   zregisto nif=%s falhou: %sr   :Ni  Nr   )rP   rQ   rR   rS   rT   )r
   is_settime	monotonicr   commitrL   r   r<   r	   set	Exceptionloggerwarningr    )r   r*   r)   batchitemrb   keyer   deadlinerD   r+   lockr   stoptotalss   &       r   workerprocess_queue.<locals>.worker  s    ===Fkkmm(88(C44k*i7  4 !==4 -...'"1'"CCE!..*** /.  !== Qx!-0222g&/#Wd5k*.w-1*<Q""&w-*7	' ! 0eTY^___%nn...  32  $tt{+q0+#YC"3K599S!+<<K $Z  $tt+ !== !444 /C* /... !  3! `.  3222  $ttt ) 
  1222g/#%[J]   &nn...	  322222
 HHJ  -NN#>UQO0222g/#%[At   &nn...	  322222
  $ttw1,  $tttttt-I !===sh  R-H.R-R)RH1 R#H=	3R>H4?RR-H7R-R H:!R8I9R<I&	II&	%I &I&	*R5I#6RRR-JR-	RK2J3K63J	)J	*J	JJ	JJ	#K.J/KJ/KAJ3		KJ1KRR-'R(R-1R4R7R-:R=II
II
RI&	 I&	#R&J ,I/-
J 8J :	RR-K	J	J	J	KJ,J
J,%J,'	K1K3K9J<:
KKKRKR
/K20R
4L>LL>&L)'L>,R
7L:
8R
>MM
MM"R
5R
61R'N*(R,)POP.O1/P4R?P
 RPP
PPR0P31R5Q#RQ
RR#R)Q,*
R5R7R;RRR

RR-R*	R
R*	"R*	$	R-c              3   4   <"   T F  pS! V4      x  K  	  R # 5iNr   ).0ir   s   & r   	<genexpr> process_queue.<locals>.<genexpr>  s     =n6!99ns   Npedidos)app.dbr   r   REGISTRY_WORKERSr   r   r   asyncioLockEventgatherrangeacquired)r   r   r   rD   r   r   r+   r   r   r   r   s   f&&f@@@@@@@r   process_queuer     s     $ )2222GmG~~+-H1j!Waa8Q4F<<>D==?D-- --^ ..=eGn=
>>>((F9M ?s   B;CCC)	   r   )   )2   g      @NT)-r   r   r=   loggingr   r   r   typingr   
sqlalchemyr   sqlalchemy.ext.asyncior   
app.configr   app.scrapers.baser   app.scrapers.mj_irnr	   r
   "app.services.mj_company_enrichmentr    app.services.registry_act_parserr   r   app.services.registry_ingestr   r   r   	getLoggerr   r   RuntimeErrorr   r'   rB   rL   rJ   rq   rv   r~   r   r   r   r   r   r   r   <module>r      s   "     #   /  ) > ; E P P			8	$O, O0f   2Y
 "&Y #'Y Y Yx6 B2n  %PPTX0M Mr   