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
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 02:04 +0000
1from __future__ import annotations
3import asyncio
4from typing import TYPE_CHECKING, NoReturn
6from docket import Perpetual
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
19if TYPE_CHECKING: 19 ↛ 20line 19 didn't jump to line 20 because the condition on line 19 was never true
20 import logging
23logger: "logging.Logger" = get_logger(__name__)
26class ReactiveTriggers(RunInEphemeralServers, Service):
27 """Evaluates reactive automation triggers"""
29 consumer_task: asyncio.Task[None] | None = None
31 @classmethod
32 def service_settings(cls) -> ServicesBaseSetting:
33 return get_current_settings().server.services.triggers
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 )
48 async with triggers.consumer() as handler:
49 self.consumer_task = asyncio.create_task(self.consumer.run(handler))
50 logger.debug("Reactive triggers started")
52 try:
53 await self.consumer_task
54 except asyncio.CancelledError:
55 pass
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")
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()