Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/events/services/triggers.py: 96%

45 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-07 02:04 +0000

1from __future__ import annotations 

2 

3import asyncio 

4from typing import TYPE_CHECKING, NoReturn 

5 

6from docket import Perpetual 

7 

8from prefect.logging import get_logger 

9from prefect.server.events import triggers 

10from prefect.server.services.base import RunInEphemeralServers, Service 

11from prefect.server.services.perpetual_services import perpetual_service 

12from prefect.server.utilities.messaging import Consumer, create_consumer 

13from prefect.server.utilities.messaging._consumer_names import ( 

14 generate_unique_consumer_name, 

15) 

16from prefect.settings.context import get_current_settings 

17from prefect.settings.models.server.services import ServicesBaseSetting 

18 

19if TYPE_CHECKING: 19 ↛ 20line 19 didn't jump to line 20 because the condition on line 19 was never true

20 import logging 

21 

22 

23logger: "logging.Logger" = get_logger(__name__) 

24 

25 

26class ReactiveTriggers(RunInEphemeralServers, Service): 

27 """Evaluates reactive automation triggers""" 

28 

29 consumer_task: asyncio.Task[None] | None = None 

30 

31 @classmethod 

32 def service_settings(cls) -> ServicesBaseSetting: 

33 return get_current_settings().server.services.triggers 

34 

35 async def start(self) -> NoReturn: 

36 assert self.consumer_task is None, "Reactive triggers already started" 

37 consumer_name = generate_unique_consumer_name("reactive-triggers") 

38 logger.info( 

39 f"ReactiveTriggers starting with unique consumer name: {consumer_name}" 

40 ) 

41 self.consumer: Consumer = create_consumer( 

42 "events", 

43 group="reactive-triggers", 

44 name=consumer_name, 

45 read_batch_size=self.service_settings().read_batch_size, 

46 ) 

47 

48 async with triggers.consumer() as handler: 

49 self.consumer_task = asyncio.create_task(self.consumer.run(handler)) 

50 logger.debug("Reactive triggers started") 

51 

52 try: 

53 await self.consumer_task 

54 except asyncio.CancelledError: 

55 pass 

56 

57 async def stop(self) -> None: 

58 assert self.consumer_task is not None, "Reactive triggers not started" 

59 self.consumer_task.cancel() 

60 try: 

61 await self.consumer_task 

62 except asyncio.CancelledError: 

63 pass 

64 finally: 

65 await self.consumer.cleanup() 

66 self.consumer_task = None 

67 logger.debug("Reactive triggers stopped") 

68 

69 

70@perpetual_service( 

71 enabled_getter=lambda: get_current_settings().server.services.triggers.enabled, 

72 run_in_ephemeral=True, 

73) 

74async def evaluate_proactive_triggers_periodic( 

75 perpetual: Perpetual = Perpetual( 

76 automatic=True, 

77 every=get_current_settings().server.events.proactive_granularity, 

78 ), 

79) -> None: 

80 """Evaluate proactive automation triggers on a periodic schedule.""" 

81 await triggers.evaluate_proactive_triggers()