Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/services/base.py: 88%
76 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 abc
4import asyncio
5import inspect
6from abc import ABC, abstractmethod
7from contextlib import asynccontextmanager
8from logging import Logger
9from types import ModuleType
10from typing import AsyncGenerator, NoReturn, Sequence
12from typing_extensions import Self
14from prefect.logging.loggers import get_logger
15from prefect.settings.models.root import canonical_environment_prefix
16from prefect.settings.models.server.services import ServicesBaseSetting
18logger: Logger = get_logger(__name__)
21def _known_service_modules() -> list[ModuleType]:
22 """Get list of Prefect server modules containing Service subclasses"""
23 from prefect.server.events import stream
24 from prefect.server.events.services import (
25 actions,
26 event_logger,
27 event_persister,
28 triggers,
29 )
30 from prefect.server.logs import stream as logs_stream
31 from prefect.server.services import (
32 task_run_recorder,
33 )
35 return [
36 # Orchestration services
37 task_run_recorder,
38 # Events services
39 event_logger,
40 event_persister,
41 triggers,
42 actions,
43 stream,
44 # Logs services
45 logs_stream,
46 ]
49class Service(ABC):
50 name: str
51 logger: Logger
53 @classmethod
54 @abstractmethod
55 def service_settings(cls) -> ServicesBaseSetting:
56 """The Prefect setting that controls whether the service is enabled"""
57 ...
59 @classmethod
60 def environment_variable_name(cls) -> str:
61 return canonical_environment_prefix(cls.service_settings()) + "ENABLED"
63 @classmethod
64 def enabled(cls) -> bool:
65 """Whether the service is enabled"""
66 return cls.service_settings().enabled
68 @classmethod
69 def all_services(cls) -> Sequence[type[Self]]:
70 """Get list of all service classes"""
71 discovered: list[type[Self]] = []
72 for module in _known_service_modules():
73 for _, obj in inspect.getmembers(module):
74 if (
75 inspect.isclass(obj)
76 and issubclass(obj, cls)
77 and not inspect.isabstract(obj)
78 ):
79 discovered.append(obj)
80 return discovered
82 @classmethod
83 def enabled_services(cls) -> list[type[Self]]:
84 """Get list of enabled service classes"""
85 return [svc for svc in cls.all_services() if svc.enabled()]
87 @classmethod
88 @asynccontextmanager
89 async def running(cls) -> AsyncGenerator[None, None]:
90 """A context manager that runs enabled services on entry and stops them on
91 exit."""
92 service_tasks: dict[Service, asyncio.Task[None]] = {}
93 for service_class in cls.enabled_services():
94 service = service_class()
95 service_tasks[service] = asyncio.create_task(service.start())
97 try:
98 yield
99 finally:
100 await asyncio.gather(*[service.stop() for service in service_tasks])
101 await asyncio.gather(*service_tasks.values(), return_exceptions=True)
103 @classmethod
104 async def run_services(cls) -> NoReturn:
105 """Run enabled services until cancelled."""
106 async with cls.running():
107 heat_death_of_the_universe = asyncio.get_running_loop().create_future()
108 try:
109 await heat_death_of_the_universe
110 except asyncio.CancelledError:
111 logger.info("Received cancellation, stopping services...")
113 @abstractmethod
114 async def start(self) -> NoReturn:
115 """Start running the service, which may run indefinitely"""
116 ...
118 @abstractmethod
119 async def stop(self) -> None:
120 """Stop the service"""
121 ...
123 def __init__(self):
124 self.name = self.__class__.__name__
125 self.logger = get_logger(f"server.services.{self.name.lower()}")
128class RunInEphemeralServers(Service, abc.ABC):
129 """
130 A marker class for services that should run even when running an ephemeral server
131 """
133 pass
136class RunInWebservers(Service, abc.ABC):
137 """
138 A marker class for services that should run when running a webserver
139 """
141 pass