
    AHj&                        d dl mZ d dlZd dlmZ d dlmZ d dlmZ ddl	m
Z
mZ erddlmZ dd	Zddd
ZdddZddZdddZy)    )annotationsN)TYPE_CHECKING)Redis)Pipeline   )DuplicateSchedulerErrorSchedulerNotFound)CronSchedulerc                      y)z0Get the Redis key for the CronScheduler registryzrq:cron_schedulers r       R/root/tools/cai/cai_env/lib/python3.12/site-packages/rq/cron_scheduler_registry.pyget_registry_keyr      s    r   c                    ||n| j                   }t               }t        j                         }|j                  || j                  |id      }|dk(  rt        d| j                   d      y)a9  Register a CronScheduler in the registry with current timestamp as score

    Args:
        cron_scheduler: CronScheduler instance to register
        pipeline: Redis pipeline to use. If None, uses cron_scheduler.connection

    Raises:
        DuplicateSchedulerError: If the scheduler is already registered
    NT)nxr   CronScheduler 'z' is already registered)
connectionr   timezaddnamer   )cron_schedulerpipeliner   registry_keyscoreadded_counts         r   registerr      sv     &1~7P7PJ#%L IIKE //,1D1De0LQU/VKa%8K8K7LLc&dee r   c                    ||n| j                   }t               }|j                  || j                        }|st	        d| j                   d      y)a  Remove a CronScheduler from the registry

    Args:
        cron_scheduler: CronScheduler instance to unregister
        pipeline: Redis pipeline to use. If None, uses cron_scheduler.connection

    Raises:
        SchedulerNotFound: If the scheduler is not found in the registry
    Nr   z' not found in registry)r   r   zremr   r	   )r   r   r   r   results        r   
unregisterr    +   sZ     &1~7P7PJ#%L __\>+>+>?F/.2E2E1FF] ^__ r   c                    t               }| j                  |dd      }|D cg c]%  }t        |t              r|j	                  d      n|' c}S c c}w )zGet all registered CronScheduler names from the registry

    Args:
        connection: Redis connection to use

    Returns:
        List of CronScheduler names (strings) sorted by registration time (oldest first)
    r   zutf-8)r   zrange
isinstancebytesdecode)r   r   keyskeys       r   get_keysr)   >   sQ     $%L \1b1D OSSs:c5#9CJJwsBSSSs   *Ac                f    t        j                          |z
  }| j                  t               d|      S )a^  Remove stale CronScheduler entries from the registry

    Removes schedulers that haven't sent a heartbeat in more than `threshold` seconds.

    Args:
        connection: Redis connection to use
        threshold: Number of seconds after which a scheduler is considered stale (default: 120)

    Returns:
        Number of stale entries removed
    r   )r   zremrangebyscorer   )r   	thresholdcutoff_times      r   cleanupr.   Q   s/     ))+	)K &&'7'91kJJr   )returnstr)N)r   r
   r   zPipeline | Noner/   None)r   r   r/   z	list[str])x   )r   r   r,   intr/   r3   )
__future__r   r   typingr   redisr   redis.clientr   
exceptionsr   r	   cronr
   r   r   r    r)   r.   r   r   r   <module>r:      s:    "     ! B# 
f.`&T&Kr   