Ë
    ÿ�j®"  ã                   óT  — d 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dZd	Zd
efd„Zded
efd„Zded
efd„Zdd„Zdededededed
dfd„Zded
ee   fd„Zdeded
edz  fd„Zdededed
dfd„Zdededed
dfd„Zded
efd„Zdeded
dfd„Z y)a{  Storage for every product ingestion source, keyed by (kind, external_ref).

Lives in the master database rather than a per-tenant schema for two reasons.
Shopify webhooks arrive keyed only by shop domain with no tenant, so per-tenant
storage would mean scanning every schema to route one. And the UNIQUE
constraint enforces one-shop-one-tenant globally, which per-schema cannot.
é    N)ÚFernet)ÚJsonÚRealDictCursor)Úsettings)Úget_master_db_connectionÚshopifyÚhttp_apiÚcrawlÚreturnc                  óP   — t        t        j                  j                  «       «      S )N)r   r   ÚSOURCE_CREDENTIALS_KEYÚencode© ó    ú;/var/www/html/strategist-ai/app/services/catalog/sources.pyÚ_fernetr      s   € Ü”(×1Ñ1×8Ñ8Ó:Ó;Ð;r   Úcredsc                 ó–   — t        «       j                  t        j                  | d¬«      j	                  «       «      j                  «       S )aO  Encrypt a credentials object.

    A dict rather than a bare string, so the same envelope holds a Shopify
    access token or an API key plus client secret without a schema change.
    sort_keys makes the plaintext stable, which matters for tests but not
    for security -- Fernet's IV already makes ciphertext non-deterministic.
    T)Ú	sort_keys)r   ÚencryptÚjsonÚdumpsr   Údecode)r   s    r   Úencrypt_credentialsr      s4   € ô ‹9×ÑœTŸZ™Z¨¸Ô>×EÑEÓGÓH×OÑOÓQÐQr   Úblobc                 ó’   — t        j                  t        «       j                  | j	                  «       «      j                  «       «      S )zMDecrypt a credentials blob. Raises InvalidToken if tampered or wrongly keyed.)r   Úloadsr   Údecryptr   r   )r   s    r   Údecrypt_credentialsr   *   s-   € ä�:‰:”g“i×'Ñ'¨¯©«Ó6×=Ñ=Ó?Ó@Ð@r   c                  ó~  — t        «       } 	 | j                  «       5 }|j                  d«       |j                  d«       d d d «       | j                  «        	 | j                  «        y # 1 sw Y   Œ+xY w# t        $ r) | j                  «        t        j                  dd¬«       ‚ w xY w# | j                  «        w xY w)Naû  
                CREATE TABLE IF NOT EXISTS product_sources (
                    id                    SERIAL PRIMARY KEY,
                    tenant_id             TEXT NOT NULL,
                    kind                  TEXT NOT NULL,
                    external_ref          TEXT NOT NULL,
                    config                JSONB NOT NULL DEFAULT '{}',
                    credentials_encrypted TEXT NOT NULL,
                    status                TEXT NOT NULL DEFAULT 'active',
                    connected_at          TIMESTAMPTZ NOT NULL DEFAULT now(),
                    disconnected_at       TIMESTAMPTZ,
                    last_synced_at        TIMESTAMPTZ,
                    UNIQUE (kind, external_ref)
                )
            zTCREATE INDEX IF NOT EXISTS product_sources_tenant_idx ON product_sources (tenant_id)z&Could not ensure product_sources tableT©Úexc_info©	r   ÚcursorÚexecuteÚcommitÚ	ExceptionÚrollbackÚloggerÚerrorÚclose)ÚconnÚcurs     r   Úensure_tabler.   /   s¢   € Ü#Ó%€DðØ�[‰[‹]ð 	˜cØ�K‰Kð ô ð �K‰Kð1ô÷!	ð( 	�‰�ð 	�
‰
�÷5	ð 	ûô* ò Ø�‰ŒÜ�‰Ð=ÈˆÔMØðûð
 	�
‰
�ús-   ŒA5 œ#A)¿A5 Á)A2Á.A5 Á52B'Â'B* Â*B<Ú	tenant_idÚkindÚexternal_refÚconfigÚcredentialsc                 óÔ  — t        «        t        «       }	 |j                  «       5 }|j                  d| ||t	        |«      t        |«      f«       ddd«       |j                  «        t        j                  d||| «       	 |j                  «        y# 1 sw Y   ŒCxY w# t        $ r+ |j                  «        t        j                  d|| d¬«       ‚ w xY w# |j                  «        w xY w)zÚInsert a source, or update it in place on (kind, external_ref) conflict.

    The ON CONFLICT target is what makes a reinstall update the existing row
    instead of creating a second row under a different tenant.
    a­  
                INSERT INTO product_sources
                    (tenant_id, kind, external_ref, config,
                     credentials_encrypted, status, connected_at, disconnected_at)
                VALUES (%s, %s, %s, %s, %s, 'active', now(), NULL)
                ON CONFLICT (kind, external_ref) DO UPDATE SET
                    tenant_id             = EXCLUDED.tenant_id,
                    config                = EXCLUDED.config,
                    credentials_encrypted = EXCLUDED.credentials_encrypted,
                    status                = 'active',
                    connected_at          = now(),
                    disconnected_at       = NULL
            Nz!Stored %s source %s for tenant %sz'Could not store %s source for tenant %sTr!   )r.   r   r$   r%   r   r   r&   r)   Úinfor'   r(   r*   r+   )r/   r0   r1   r2   r3   r,   r-   s          r   Úupsert_sourcer6   O   sÏ   € ô „NÜ#Ó%€DðØ�[‰[‹]ð 	5˜cØ�K‰Kð ð ˜T <´°f³Ü% kÓ2ð4ô5÷	5ð 	�‰ŒÜ�‰Ð7¸¸|ÈYÕWð 	�
‰
�÷/	5ð 	5ûô" ò Ø�‰ŒÜ�‰Ð>ÀÀiØ"ð 	ô 	$àð	ûð 	�
‰
�ús.   –B ¦*BÁ0B ÂBÂB Â4CÃC ÃC'c                 óp  — t        «        t        «       }	 |j                  t        ¬«      5 }|j	                  d| f«       |j                  «       D �cg c]  }t        |«      ‘Œ c}cddd«       |j                  «        S c c}w # 1 sw Y   nxY w	 |j                  «        y# |j                  «        w xY w)z=Every active source for a tenant. Never includes credentials.)Úcursor_factorya  
                SELECT kind, external_ref, config, status,
                       connected_at, last_synced_at
                FROM product_sources
                WHERE tenant_id = %s AND status = 'active'
                ORDER BY connected_at DESC
            N)r.   r   r$   r   r%   ÚfetchallÚdictr+   )r/   r,   r-   Úrs       r   Úget_sourcesr<   s   sš   € ä„NÜ#Ó%€DðØ�[‰[¬ˆ[Ó7ð 	5¸3Ø�K‰Kð ð �ôð &)§\¡\£^Ö4 ”D˜•GÒ4÷	5ð 	5ð 	�
‰
�ùò 5÷	5ð 	5úð 	5ð 	�
‰
�øˆ�
‰
�ús4   –B# ¬&BÁB Á$BÁ&	B# Â BÂBÂ
B# Â#B5c                 óR  — t        «        t        «       }	 |j                  «       5 }|j                  d| |f«       |j	                  «       }|rt        |d   «      ndcddd«       |j                  «        S # 1 sw Y   nxY w	 |j                  «        y# |j                  «        w xY w)zsDecrypted credentials for one active source, or None if not found.

    Callers must not log the return value.
    zmSELECT credentials_encrypted FROM product_sources WHERE kind = %s AND external_ref = %s AND status = 'active'r   N)r.   r   r$   r%   Úfetchoner   r+   )r0   r1   r,   r-   Úrows        r   Úget_credentialsr@   …   sž   € ô
 „NÜ#Ó%€Dð
Ø�[‰[‹]ð 	@˜cØ�K‰KðNà�|Ð$ôð
 —,‘,“.ˆCÙ25Ô& s¨1¡vÔ.¸4÷	@ð 	@ð 	�
‰
�÷	@ð 	@úð 	@ð 	�
‰
�øˆ�
‰
�ús"   –B ¦6A6Á	B Á6A?Á;B ÂB&c                 óz  — t        «       }	 |j                  «       5 }|j                  dt        |«      | |f«       ddd«       |j	                  «        	 |j                  «        y# 1 sw Y   Œ+xY w# t
        $ r+ |j                  «        t        j                  d| |d¬«       ‚ w xY w# |j                  «        w xY w)zePersist a config change (e.g. first-sync field inference) without
    touching credentials or status.zLUPDATE product_sources SET config = %s WHERE kind = %s AND external_ref = %sNz!Could not update config for %s:%sTr!   )
r   r$   r%   r   r&   r'   r(   r)   r*   r+   )r0   r1   r2   r,   r-   s        r   Úupdate_source_configrB   ™   s©   € ô $Ó%€DðØ�[‰[‹]ð 	˜cØ�K‰Kð8ä�f“˜t \Ð2ô÷	ð 	�‰�ð 	�
‰
�÷	ð 	ûô ò Ø�‰ŒÜ�‰Ð8¸$ÀØ"ð 	ô 	$àð	ûð 	�
‰
�ús-   ŒA1 œA%»A1 Á%A.Á*A1 Á14B%Â%B( Â(B:Ústatusc                 óh  — t        «       }	 |j                  «       5 }|j                  d|| |f«       ddd«       |j                  «        	 |j                  «        y# 1 sw Y   Œ+xY w# t        $ r+ |j                  «        t        j                  d| |d¬«       ‚ w xY w# |j                  «        w xY w)zˆA source whose credentials stopped working must stop being retried
    silently; the dashboard needs to show that it needs reconnecting.zLUPDATE product_sources SET status = %s WHERE kind = %s AND external_ref = %sNzCould not set status for %s:%sTr!   r#   )r0   r1   rC   r,   r-   s        r   Úmark_source_statusrE   ®   s¥   € ô $Ó%€DðØ�[‰[‹]ð 	˜cØ�K‰Kð8à˜˜|Ð,ô÷	ð 	�‰�ð 	�
‰
�÷	ð 	ûô ò Ø�‰ŒÜ�‰Ð5°t¸\Ø"ð 	ô 	$àð	ûð 	�
‰
�úó-   ŒA( œA²A( ÁA%Á!A( Á(4BÂB ÂB1c                 ó¼  — t        «        t        «       }	 |j                  «       5 }|j                  d| f«       |j                  }ddd«       |j                  «        t        j                  d| «       ||j                  «        S # 1 sw Y   ŒBxY w# t        $ r* |j                  «        t        j                  d| d¬«       ‚ w xY w# |j                  «        w xY w)a  Revoke every active source a tenant has connected. Pairs with wiping
    the whole catalogue: without this, a source stays connected after its
    products are deleted, and the next build silently re-crawls it and
    repopulates exactly what was just removed.zqUPDATE product_sources SET status = 'revoked', disconnected_at = now() WHERE tenant_id = %s AND status = 'active'Nz'Disconnected %d source(s) for tenant %sz*Could not disconnect sources for tenant %sTr!   )r.   r   r$   r%   Úrowcountr&   r)   r5   r+   r'   r(   r*   )r/   r,   r-   Úcounts       r   Údisconnect_all_sourcesrJ   Ã   s½   € ô
 „NÜ#Ó%€DðØ�[‰[‹]ð 	!˜cØ�K‰Kð=à�ôð
 —L‘LˆE÷	!ð 	�‰ŒÜ�‰Ð=¸uÀiÔPØð 	�
‰
�÷	!ð 	!ûô ò Ø�‰ŒÜ�‰ÐAÀ9ÐW[ˆÔ\Øðûð
 	�
‰
�ús.   –B ¦ BÁ0B ÂBÂB Â3CÃC	 Ã	Cc                 óh  — t        «       }	 |j                  «       5 }|j                  d|| |f«       d d d «       |j                  «        	 |j                  «        y # 1 sw Y   Œ+xY w# t        $ r+ |j                  «        t        j                  d| |d¬«       ‚ w xY w# |j                  «        w xY w)NzTUPDATE product_sources SET last_synced_at = %s WHERE kind = %s AND external_ref = %sz&Could not set last_synced_at for %s:%sTr!   r#   )r0   r1   Úwhenr,   r-   s        r   Útouch_last_syncedrM   Ý   s£   € Ü#Ó%€DðØ�[‰[‹]ð 	˜cØ�K‰Kð8à�t˜\Ð*ô÷	ð 	�‰�ð 	�
‰
�÷	ð 	ûô ò Ø�‰ŒÜ�‰Ð=¸tÀ\Ø"ð 	ô 	$àð	ûð 	�
‰
�úrF   )r   N)!Ú__doc__r   ÚloggingÚcryptography.fernetr   Úpsycopg2.extrasr   r   Úapp.core.configr   Úapp.services.infra.databaser   Ú	getLoggerÚ__name__r)   ÚKIND_SHOPIFYÚKIND_HTTP_APIÚ
KIND_CRAWLr   r:   Ústrr   r   r.   r6   Úlistr<   r@   rB   rE   ÚintrJ   rM   r   r   r   ú<module>r\      sV  ðñó Û å &ß 0å $Ý @à	ˆ×	Ñ	˜8Ó	$€à€Ø€ð €
ð<�ó <ðR˜tð R¨ó RðA˜cð A dó Aó
ð@!˜Sð !¨ð !¸3ð !Øð!Ø.2ð!Ø7;ó!ðH˜3ð  4¨¡:ó ð$˜#ð ¨Sð °T¸D±[ó ð(˜sð °#ð ¸tð Èó ð*˜Sð °ð ¸Sð ÀTó ð* cð ¨có ð4˜Cð ¨sð ¸Tô r   