
    PmjF                       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 ddl	m
Z
 ddl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 ddlmZ ddlmZmZmZ dZdZdZdZdZ  ej!        e"          Z#ddZ$ddZ% G d d          Z&dS )z?Durable aggregation and local export for Hermes shared metrics.    )annotationsN)Iterator)contextmanager)datetime	timedeltatimezone)Path)Any)	write_txn)get_hermes_home)atomic_json_write   )COUNTER_METRICSMODEL_CALL_METRICcounter_dimensions_are_validzhermes.shared_metrics.v11   i     returnr   c                 >    t          j        t          j                  S N)r   nowr   utc     M/home/thesage/.hermes/hermes-agent/hermes_cli/observability/shared_metrics.py_utc_nowr   #   s    <%%%r   valuestrc                    |                      t          j                                                                      dd          S )Nz+00:00Z)
astimezoner   r   	isoformatreplace)r   s    r   
_isoformatr%   '   s4    HL))3355==hLLLr   c                     e Zd ZdZ	 	 d3d4dZd5dZd6dZd7dZd7dZd7dZ	d8dZ
eedd9d            Zed:d            Zed:d            Zd;dZed<d"            Zd=d#Zd>d$Zd;d%Zd?d'Zd@d*ZedAd.            Zd7d/Zdd0dBd2ZdS )CSharedMetricsStorezAPersist allowlisted counters and export immutable delta packages.Ndatabase_pathPath | Noneoutbox_directoryr   Nonec                ,   t                      dz  dz  }|p|dz  | _        |p|dz  | _        |                     | j        j                   |                     | j                   |                     | j                   |                                  d S )N	telemetryshared_metricszmetrics.sqlite3outbox)r   r(   r*   _ensure_private_directoryparent_ensure_private_file_ensure_schema)selfr(   r*   roots       r   __init__zSharedMetricsStore.__init__.   s    
   ;.1AA*Fd5F.F 0 CD8O&&t'9'@AAA&&t'<===!!$"4555r   
dimensionsdict[str, str]hermes_versionr   c                >    |                      t          ||           dS )zBIncrement the terminal model-call counter for the current UTC day.N)record_counterr   )r4   r7   r9   s      r   record_model_callz$SharedMetricsStore.record_model_call;   s#     	-z>JJJJJr   metric_namec                   |t           vrt          d|           t          ||          st          d|           t          j        |dd          }t                                                                                      }|                                 5 }|	                    d|||pd|f           ddd           dS # 1 swxY w Y   dS )	z:Increment one allowlisted counter for the current UTC day.zUnsupported shared metric: *Unsupported dimensions for shared metric: T,:	sort_keys
separatorsa!  
                INSERT INTO counter_aggregates(
                    period_start,
                    metric_name,
                    hermes_version,
                    dimensions_json,
                    value,
                    packaged_value
                ) VALUES (?, ?, ?, ?, 1, 0)
                ON CONFLICT(
                    period_start,
                    metric_name,
                    hermes_version,
                    dimensions_json
                )
                DO UPDATE SET value = value + 1
                unknownN)
r   
ValueErrorr   jsondumpsr   dater#   _connectionexecute)r4   r=   r7   r9   dimensions_jsonperiod_start
connections          r   r;   z!SharedMetricsStore.record_counterC   s8    o--H;HHIII+KDD 	YW+WWXXX*!
 
 

  zz((2244 	:$ !"/i#	%  	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	s   CC	C	
list[Path]c                    |                                  }t          |          D ]}|                                  n|                                 S )zDCommit one pending delta package, then atomically export the outbox.)_pending_period_countrange_create_package_export_and_prune)r4   pending_periods_s      r   create_and_export_packagez,SharedMetricsStore.create_and_export_packageo   sX    4466'' 	 	A##%%- .%%'''r   c                R    |                                   |                                 S )zCCreate pending packages at most once per UTC day, then export them.)_create_pending_packages_if_duerU   )r4   s    r    create_and_export_package_if_duez3SharedMetricsStore.create_and_export_package_if_duew   s&    ,,...%%'''r   c                    |                                  }	 |                                  n,# t          $ r t                              dd           Y nw xY w|S )Nz.Unable to prune expired shared-metrics historyTexc_info)_export_pending_packages_prune_expired_history	Exceptionloggerwarning)r4   exporteds     r   rU   z$SharedMetricsStore._export_and_prune|   sz    0022	'')))) 	 	 	NN@      	
 s   + &AAlist[dict[str, Any]]c                    |                                  5 }|                    d                                          }ddd           n# 1 swxY w Y   d |D             S )zBReturn cumulative counters for focused tests and local inspection.an  
                SELECT
                    period_start,
                    metric_name,
                    hermes_version,
                    dimensions_json,
                    value,
                    packaged_value
                FROM counter_aggregates
                ORDER BY period_start, hermes_version, metric_name, dimensions_json
                Nc           	         g | ]A}|d          |d         |d         t          j        |d                   |d         |d         dBS )rN   r=   r9   rM   r   packaged_value)rN   r=   r9   r7   r   rh   )rH   loads).0rows     r   
<listcomp>z7SharedMetricsStore.counter_snapshot.<locals>.<listcomp>   sm     

 

 

  !$N 3"=1"%&6"7"j->)?@@W"%&6"7 

 

 

r   )rK   rL   fetchall)r4   rO   rowss      r   counter_snapshotz#SharedMetricsStore.counter_snapshot   s     	:%%
  hjj 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	

 

 

 

 

 
	
   (A		AAbusy_timeout_msrr   intIterator[sqlite3.Connection]c             #  0  K   t          j        | j        |dz            }	 t           j        |_        |                    d|            |5  |V  d d d            n# 1 swxY w Y   |                                 d S # |                                 w xY w)Ni  )timeoutzPRAGMA busy_timeout=)sqlite3connectr(   Rowrow_factoryrL   close)r4   rr   rO   s      r   rK   zSharedMetricsStore._connection   s       _#d*
 
 

	%,[J"GoGGHHH ! !    ! ! ! ! ! ! ! ! ! ! ! ! ! ! ! Js/   +A? AA? A""A? %A"&A? ?Bpathr	   c                    |                      ddd           	 |                     d           d S # t          $ r Y d S w xY w)NTi  )parentsexist_okmode)mkdirchmodOSErrorr|   s    r   r0   z,SharedMetricsStore._ensure_private_directory   sY    

4$U
;;;	JJu 	 	 	DD	s   1 
??c                    |                      dd           	 |                     d           d S # t          $ r Y d S w xY w)N  T)r   r   )touchr   r   r   s    r   r2   z'SharedMetricsStore._ensure_private_file   sW    


---	JJu 	 	 	DD	s   0 
>>c                    |                      t                    5 }t          |          5  |                     |           d d d            n# 1 swxY w Y   d d d            d S # 1 swxY w Y   d S )Nrq   )rK   _SCHEMA_BUSY_TIMEOUT_MSr   _ensure_schema_in_transaction)r4   rO   s     r   r3   z!SharedMetricsStore._ensure_schema   s    .EFF 	?*:&& ? ?22:>>>? ? ? ? ? ? ? ? ? ? ? ? ? ? ?	? 	? 	? 	? 	? 	? 	? 	? 	? 	? 	? 	? 	? 	? 	? 	? 	? 	?s4   A&AA&A	A&A	A&&A*-A*rO   sqlite3.Connectionc                z   |                      d           |                      d                                          }|6t          |d                   t          k    rt	          d|d                    |                      d           |                      d           |                      dt          f           d S )Nz
            CREATE TABLE IF NOT EXISTS telemetry_state (
                key TEXT PRIMARY KEY,
                value TEXT NOT NULL
            )
            z>SELECT value FROM telemetry_state WHERE key = 'schema_version'r   z1Unsupported shared-metrics store schema version: a)  
            CREATE TABLE IF NOT EXISTS counter_aggregates (
                period_start TEXT NOT NULL,
                metric_name TEXT NOT NULL,
                hermes_version TEXT NOT NULL,
                dimensions_json TEXT NOT NULL,
                value INTEGER NOT NULL,
                packaged_value INTEGER NOT NULL DEFAULT 0,
                PRIMARY KEY (
                    period_start,
                    metric_name,
                    hermes_version,
                    dimensions_json
                )
            )
            aM  
            CREATE TABLE IF NOT EXISTS package_outbox (
                package_id TEXT PRIMARY KEY,
                period_start TEXT NOT NULL,
                period_end TEXT NOT NULL,
                payload_json TEXT NOT NULL,
                created_at TEXT NOT NULL,
                exported_at TEXT
            )
            zt
            INSERT OR IGNORE INTO telemetry_state(key, value)
            VALUES ('schema_version', ?)
            )rL   fetchoner   _STORE_SCHEMA_VERSIONRuntimeError)rO   
schema_rows     r   r   z0SharedMetricsStore._ensure_schema_in_transaction   s    	
 	
 	
  ''L
 

(** 	 !c*W*=&>&>BW&W&W)g&) )   		
 	
 	
$ 			
 	
 	
 	 #$	
 	
 	
 	
 	
r   c                   |                     d                                          }|t          |d                   S t          t          j                              }|                     d|f           |                     d                                          }|t          d          t          |d                   S )Nz:SELECT value FROM telemetry_state WHERE key = 'install_id'r   zJINSERT OR IGNORE INTO telemetry_state(key, value) VALUES ('install_id', ?)z4Unable to create the shared-metrics install identity)rL   r   r   uuiduuid4r   )r4   rO   rk   	candidates       r   _install_idzSharedMetricsStore._install_id  s      H
 

(** 	 ?s7|$$$
%%	XL	
 	
 	
   H
 

(** 	 ;UVVV3w<   r   c                    |                                  5 }|                    d                                          }d d d            n# 1 swxY w Y   |t          |d                   ndS )Na9  
                SELECT COUNT(*) AS period_count
                FROM (
                    SELECT period_start, hermes_version
                    FROM counter_aggregates
                    WHERE value > packaged_value
                    GROUP BY period_start, hermes_version
                )
                period_countr   )rK   rL   r   rs   )r4   rO   rk   s      r   rR   z(SharedMetricsStore._pending_period_count  s     	:$$
 
 hjj 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 ,/?s3~&'''Arp   c                   t                      }|                                 5 }t          |          5  |                    d|                                                                f                                          }|	 d d d            d d d            d S |                     ||          	 |                     ||          d d d            n# 1 swxY w Y   d d d            d S # 1 swxY w Y   d S )Nz
                    SELECT 1
                    FROM package_outbox
                    WHERE substr(created_at, 1, 10) >= ?
                    LIMIT 1
                    )r   rK   r   rL   rJ   r#   r   _create_package_in_transaction)r4   r   rO   package_created_todays       r   rZ   z2SharedMetricsStore._create_pending_packages_if_due$  s   jj 	::&&   )3(:(: XXZZ))++-) ) (** & )4     	 	 	 	 	 	 	 	 99*cJJV 99*cJJV              	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	sA   C-ACC--C	C-C	C-C	C--C14C1dict[str, Any] | Nonec                   t                      }|                                 5 }t          |          5  |                     ||          cd d d            cd d d            S # 1 swxY w Y   	 d d d            d S # 1 swxY w Y   d S r   )r   rK   r   r   )r4   r   rO   s      r   rT   z"SharedMetricsStore._create_package8  sN   jj 	L::&& L L:::sKKL L L L L L L	L 	L 	L 	L 	L 	L 	L 	LL L L L L L L L L	L 	L 	L 	L 	L 	L 	L 	L 	L 	L 	L 	L 	L 	L 	L 	L 	L 	Ls4   A;A"	A;"A&	&A;)A&	*A;;A?A?r   r   c           	     `    |                     d                                          }||d         nd }|sd S |                     d||d         f                                          }t          j        t          |                                        t          j                  }|t          d          z   }t          t          j                              }t          |                     |          t          |          t          |          t          |          d|d         i fd|D             d	}	t          j        |	d
d          }
|                     d||	d         |	d         |
|	d         f           |D ].}|                     d||d         |d         |d         f           /|	S )Nz
                SELECT period_start, hermes_version
                FROM counter_aggregates
                WHERE value > packaged_value
                ORDER BY period_start, hermes_version
                LIMIT 1
                rN   a7  
                SELECT metric_name, dimensions_json, value, packaged_value
                FROM counter_aggregates
                WHERE period_start = ?
                  AND hermes_version = ?
                  AND value > packaged_value
                ORDER BY metric_name, dimensions_json
                r9   )tzinfor   daysc                :    g | ]}                     |          S r   )_package_metric)rj   rk   r4   s     r   rl   zESharedMetricsStore._create_package_in_transaction.<locals>.<listcomp>h  s'    BBBc,,S11BBBr   )schema_version
package_id
install_idrN   
period_endgenerated_atresourcemetricsTr@   rC   a	  
                INSERT INTO package_outbox(
                    package_id,
                    period_start,
                    period_end,
                    payload_json,
                    created_at
                ) VALUES (?, ?, ?, ?, ?)
                r   r   a"  
                    UPDATE counter_aggregates
                    SET packaged_value = value
                    WHERE period_start = ?
                      AND metric_name = ?
                      AND hermes_version = ?
                      AND dimensions_json = ?
                    r=   rM   )rL   r   rm   r   fromisoformatr   r$   r   r   r   r   r   _PACKAGE_SCHEMA_VERSIONr   r%   rH   rI   )r4   rO   r   
period_rowperiod_valuern   rN   r   r   payloadpayload_jsonrk   s   `           r   r   z1SharedMetricsStore._create_package_in_transaction>  s   
  ''
 
 (** 	 6@5Kz.11QU 	4!! :&678

 

 (** 	  -c,.?.?@@HH< I 
 
 "I1$5$5$55
&&
5$**:66&|44$Z00&sOO):6F+GHBBBBTBBB	
 	
 z!
 
 

 	 '%'	
 	
 	
$  	 	C !&/0)*	     r   rk   sqlite3.Rowdict[str, Any]c                    t          | d                   }t          j        | d                   }t          |t                    rt          ||          st          d|           |d|| d         | d         z
  dS )Nr=   rM   r?   counterr   rh   )nametyper7   r   )r   rH   ri   
isinstancedictr   rG   )rk   r=   r7   s      r   r   z"SharedMetricsStore._package_metric  s    #m,--Z$5 677
*d++ 	Y3O4
 4
 	Y W+WWXXX$\C(8$99	
 
 	
r   c           	     :   |                                  5 }|                    d                                          }d d d            n# 1 swxY w Y   g }|D ]}t          |d                   }| j        | dz  }t          |t          j        |d                   ddd           |                                  5 }|                    d	t          t                                |f           d d d            n# 1 swxY w Y   |
                    |           |S )
Nz
                SELECT package_id, payload_json
                FROM package_outbox
                WHERE exported_at IS NULL
                ORDER BY created_at, package_id
                r   .jsonr      Tr   )indentrD   r   z
                    UPDATE package_outbox
                    SET exported_at = ?
                    WHERE package_id = ? AND exported_at IS NULL
                    )rK   rL   rm   r   r*   r   rH   ri   r%   r   append)r4   rO   rn   rd   rk   r   r|   s          r   r_   z+SharedMetricsStore._export_pending_packages  s    	:%%  hjj 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	  " 	" 	"CS.//J(j+?+?+??D
3~.//    !!## z""
  

++Z8                 OOD!!!!s#   (A		AA<2C::C>	C>	)r   datetime | Nonec               x   |pt                      t          t                    z
  }t          |          }|                                                                }|                                 5 }|                    d|f                                          }ddd           n# 1 swxY w Y   g }|D ]|}t          |d                   }		 | j
        |	 dz                      d           n-# t          $ r  t                              d|	d	           Y cw xY w|                    |	           }|                                 5 }t!          |          5  |D ]}	|                    d
|	|f           |                    d|f           ddd           n# 1 swxY w Y   ddd           dS # 1 swxY w Y   dS )zARemove exported local history after the bounded retention window.r   z
                SELECT package_id
                FROM package_outbox
                WHERE exported_at IS NOT NULL
                  AND exported_at < ?
                ORDER BY exported_at, package_id
                Nr   r   T)
missing_okz1Unable to prune expired shared-metrics package %sr]   z
                        DELETE FROM package_outbox
                        WHERE package_id = ?
                          AND exported_at IS NOT NULL
                          AND exported_at < ?
                        a  
                    DELETE FROM counter_aggregates
                    WHERE period_start < ?
                      AND value = packaged_value
                      AND NOT EXISTS (
                          SELECT 1
                          FROM package_outbox
                          WHERE exported_at IS NULL
                            AND substr(package_outbox.period_start, 1, 10)
                                = counter_aggregates.period_start
                      )
                    )r   r   _LOCAL_HISTORY_RETENTION_DAYSr%   rJ   r#   rK   rL   rm   r   r*   unlinkr   rb   rc   r   r   )
r4   r   cutoffcutoff_timestampcutoff_periodrO   rn   removable_package_idsrk   r   s
             r   r`   z)SharedMetricsStore._prune_expired_history  s   #y.(
 (
 (
 
 &f--//11 
	:%% "#	 	 hjj 
	 
	 
	 
	 
	 
	 
	 
	 
	 
	 
	 
	 
	 
	 
	 ,. 	5 	5CS.//J
&J)=)=)==EE# F        G!    
  "((4444 	::&&  "7 	 	J&& $%56    "" #$                	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	sZ   0*B&&B*-B*!C//'DDF/5FF/F	F/F	F//F36F3)NN)r(   r)   r*   r)   r   r+   )r7   r8   r9   r   r   r+   )r=   r   r7   r8   r9   r   r   r+   )r   rP   )r   re   )rr   rs   r   rt   )r|   r	   r   r+   )r   r+   )rO   r   r   r+   )rO   r   r   r   )r   rs   )r   r   )rO   r   r   r   r   r   )rk   r   r   r   )r   r   r   r+   )__name__
__module____qualname____doc__r6   r<   r;   rX   r[   rU   ro   r   _BUSY_TIMEOUT_MSrK   staticmethodr0   r2   r3   r   r   rR   rZ   rT   r   r   r_   r`   r   r   r   r'   r'   +   s'       KK &*(,    K K K K* * * *X( ( ( (( ( ( (
	 	 	 	
 
 
 
8   0     ^"    \    \? ? ? ? 5
 5
 5
 \5
n! ! ! !$B B B B   (L L L LT T T Tl 
 
 
 \
       D @D = = = = = = = =r   r'   )r   r   )r   r   r   r   )'r   
__future__r   rH   loggingrw   r   collections.abcr   
contextlibr   r   r   r   pathlibr	   typingr
   hermes_cli.sqlite_utilr   hermes_constantsr   utilsr   shared_metrics_contractr   r   r   r   r   r   r   r   	getLoggerr   rb   r   r%   r'   r   r   r   <module>r      s   E E " " " " " "     $ $ $ $ $ $ % % % % % % 2 2 2 2 2 2 2 2 2 2             , , , , , , , , , , , , # # # # # #          5     " 		8	$	$& & & &M M M MW W W W W W W W W Wr   