o
    pj,                     @  s  d Z ddlmZ ddlmZmZ ddlmZmZ ddlm	Z	 ddl
mZ ddlmZmZmZmZ ddlmZ dd	lmZmZ d
dddZdZG dd deZeddG dd dZdddd?d"d#Zd@d$d%Zdd&dAd(d)ZdBd+d,ZdCd/d0ZdDd4d5Z dEd7d8Z!dFd=d>Z"dS )Gz-Persistent, multi-node-safe quota accounting.    )annotations)asdict	dataclass)datetime	timedelta)Decimal)Any)casefuncor_select)get_control_plane_session)TenantLimitTenantUsagerequests_per_minuteai_tokens_monthlymonthly_jobs)requests	ai_tokensjobs_resetc                      s   e Zd Zd fddZ  ZS )QuotaExceededdecision'QuotaDecision'returnNonec              	     s2   || _ t d|j d|jdd|jd d S )NzQuota exceeded for z: g/)r   super__init__
limit_nameusedlimit)selfr   	__class__ </var/www/html/flask_server/bridge_platform/quotas/service.pyr      s   
zQuotaExceeded.__init__)r   r   r   r   )__name__
__module____qualname__r   __classcell__r&   r&   r$   r'   r      s    r   T)frozenc                   @  s`   e Zd ZU ded< ded< ded< ded< ded< d	ed
< d	ed< ded< ded< dddZdS )QuotaDecisionstrmetricr    zfloat | Noner"   floatr!   	remainingr   period_start
period_endboolallowedpolicyr   dict[str, Any]c                 C  s(   t | }| j |d< | j |d< |S )Nr2   r3   )r   r2   	isoformatr3   )r#   payloadr&   r&   r'   to_dict-   s   zQuotaDecision.to_dictN)r   r7   )r(   r)   r*   __annotations__r:   r&   r&   r&   r'   r-   !   s   
 r-      N)amountmetadata	tenant_idr.   app_id
str | Noner/   r=   r0   r>   dict[str, Any] | Noner   c                 C  s  |dk rt dt||}t }t }t| || }||	 
 }	t||	\}
}t|| ||
|}t|ttttjdtj| ktj|ktj|ktj|k  p[d}|	ri|	jdurit|	jnd}|	rr|	jpqdnd}|du p|| |kp|dk}|| }t||||r|n||du rdn
td||r|n| |
|||d	}|s|  t||t| ||t t!||
||pi d |"  |W  d   S 1 sw   Y  dS )	zAtomically check and record usage.

    Locking the effective limit row serializes decisions across physical and
    virtual workers when the control plane runs on PostgreSQL/MySQL.
    r   Usage amount cannot be negativeNblock	unlimited        	r/   r    r"   r!   r1   r2   r3   r5   r6   r?   r@   metric_namemetric_valueusage_period_startusage_period_endmetadata_json)#
ValueErrorMETRIC_LIMITSgetr   utcnowr   _effective_limit_querywith_for_updateexecutescalarsfirst_period_usage_startr0   r   r
   coalescesumr   rJ   wherer?   rI   recorded_at
scalar_onelimit_valueoverage_policyr-   maxrollbackr   addr   r.   commit)r?   r@   r/   r=   r>   r    nowsession	limit_rowlimit_recordr2   r3   usage_startr!   r"   r6   r5   	projectedr   r&   r&   r'   consume4   sp    
 
$rj   c                 C  sb  |dk rt dt||}t }t }|t| || 	 
 }t||\}}	t|| |||	}
t|ttttjdtj| ktj|ktj|
ktj|	k  pYd}|rg|jdurgt|jnd}|rp|jpodnd}|du p|| |kp|dk}t|||||du rdntd|| ||	||d	}|st||W  d   S 1 sw   Y  dS )zHCheck capacity atomically without recording usage or calling a provider.r   rC   NrD   rE   rF   rG   )rN   rO   rP   r   rQ   r   rT   rR   rS   rU   rV   rW   rX   r0   r   r
   rY   rZ   r   rJ   r[   r?   rI   r\   r]   r^   r_   r-   r`   r   )r?   r@   r/   r=   r    rd   re   rg   r2   r3   rh   r!   r"   r6   r5   r   r&   r&   r'   ensure_capacityx   sD    $rk   )r>   r   c           
      C  s   t  }t 9}|t| |t||  }t	||\}}	|
t| ||tt|||	|p0i d |  W d   dS 1 sCw   Y  dS )zORecord provider usage that has already occurred; never reject it retroactively.rH   N)r   rQ   r   rT   rR   rO   rP   rU   rV   rW   rb   r   r   r.   rc   )
r?   r@   r/   r=   r>   rd   re   rg   r2   r3   r&   r&   r'   record_observed_usage   s   
"rl   dict[str, dict[str, Any]]c                   s  t  }i }t }|tttj| ktj	tj
  }|D ]  j
dkr+q#t fddt D  j
}t| \}}t|| |||}tj| ktj|ktj|ktj|k g}	 j	rg|	tj	 j	k t|ttttjdj|	  p}d}
 jdurt jnd} j	r j	 d j
 n j
}| j	|
||du rdntd||
 |sdntdt|
| d d	|  |  |dur|
|krш j!pd
d
krdndd	||< q#W d   |S 1 sw   Y  |S )zCReturn current usage and balances for all configured tenant limits.
storage_mbc                 3  s"    | ]\}}| j kr|V  qd S N)r    ).0keyvaluerecordr&   r'   	<genexpr>   s     z!usage_snapshot.<locals>.<genexpr>r   N:rF   d      rD   blocked	available)	r/   r@   r!   r"   r1   
percentager2   r3   status)"r   rQ   r   rT   r   r   r[   r?   order_byr@   r    rU   allnextrO   itemsrW   rX   r   rI   r\   appendr0   r
   rY   rZ   rJ   r]   r^   r`   minroundr8   r_   )r?   rd   resultre   limitsr/   startendrh   filtersr!   r"   rq   r&   rs   r'   usage_snapshot   s^   
$
&&r   reasonr7   c                 C  s   | r|t vst| stdt }t Y}t|d\}}t|	t
tttjdtj| ktj|ktj|ktj|k  pDd}|t| d| t td||t| dd |dd |  W d   n1 ssw   Y  | ||dd	S )
zDReset a tenant counter without deleting its immutable usage history.invalid_usage_resetNr   0   )r   previous_usagerH   rF   )r?   r/   r   r!   )rO   r.   striprN   r   rQ   r   rW   r0   rT   r   r
   rY   rZ   r   rJ   r[   r?   rI   r\   r]   rb   RESET_SUFFIXr   rc   )r?   r/   r   rd   re   r2   r3   previousr&   r&   r'   reset_usage_counter   s4   
r   r2   r   r3   c              	   C  sX   |  tttjtj|ktj| t	 ktj|ktj|k 
 }|r*t||S |S ro   )rT   r   r
   r`   r   r\   r[   r?   rI   r   scalar_one_or_none)re   r?   r/   r2   r3   reset_atr&   r&   r'   rX      s   rX   r    c                 C  sd   |rt tj|ktjd ntjd }tttj| ktj|k|t	tj|kdfdd
dS )Nr   r<   )else_)r   r   r@   is_r   r[   r?   r    r}   r	   r"   )r?   r@   r    scope_filterr&   r&   r'   rR      s   
rR   rd   r"   TenantLimit | Nonetuple[datetime, datetime]c                 C  s   |r"|j r"t|j }t|  }t|||  }||t|d fS t| j| jd}t| j| jdk | jdkr9dn| jd d}||fS )N)secondsr<      )window_secondsint	timestampr   utcfromtimestampr   yearmonth)rd   r"   r   epochr   r   r&   r&   r'   rW   
  s   

,rW   )r?   r.   r@   rA   r/   r.   r=   r0   r>   rB   r   r-   )
r?   r.   r@   rA   r/   r.   r=   r0   r   r-   )r?   r.   r@   rA   r/   r.   r=   r0   r>   rB   r   r   )r?   r.   r   rm   )r?   r.   r/   r.   r   r.   r   r7   )
r?   r.   r/   r.   r2   r   r3   r   r   r   )r?   r.   r@   rA   r    r.   )rd   r   r"   r   r   r   )#__doc__
__future__r   dataclassesr   r   r   r   decimalr   typingr   
sqlalchemyr	   r
   r   r   config.control_planer   bridge_platform.tenants.modelsr   r   rO   r   PermissionErrorr   r-   rj   rk   rl   r   r   rX   rR   rW   r&   r&   r&   r'   <module>   s8    	
D

-


