Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/models/variables.py: 73%
92 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 typing import Optional, Sequence
2from uuid import UUID
4import sqlalchemy as sa
5from sqlalchemy.ext.asyncio import AsyncSession
7from prefect.server.database import PrefectDBInterface, db_injector, orm_models
8from prefect.server.events import clients
9from prefect.server.events.schemas import lifecycle
10from prefect.server.schemas import filters, sorting
11from prefect.server.schemas.actions import VariableCreate, VariableUpdate
12from prefect.types._datetime import now
15async def emit_variable_created_event(variable: orm_models.Variable) -> None:
16 """Emit an event when a variable is created."""
17 async with clients.PrefectServerEventsClient() as events_client:
18 await events_client.emit(lifecycle.variable_created_event(variable, now("UTC")))
21async def emit_variable_updated_event(variable: orm_models.Variable) -> None:
22 """Emit an event when a variable is updated."""
23 async with clients.PrefectServerEventsClient() as events_client:
24 await events_client.emit(lifecycle.variable_updated_event(variable, now("UTC")))
27async def emit_variable_deleted_event(variable: orm_models.Variable) -> None:
28 """Emit an event when a variable is deleted."""
29 async with clients.PrefectServerEventsClient() as events_client:
30 await events_client.emit(lifecycle.variable_deleted_event(variable, now("UTC")))
33@db_injector
34async def create_variable(
35 db: PrefectDBInterface, session: AsyncSession, variable: VariableCreate
36) -> orm_models.Variable:
37 """
38 Create a variable
40 Args:
41 session: async database session
42 variable: variable to create
44 Returns:
45 orm_models.Variable
46 """
47 model = db.Variable(**variable.model_dump())
48 session.add(model)
49 await session.flush()
51 await emit_variable_created_event(model)
53 return model
56@db_injector
57async def read_variable(
58 db: PrefectDBInterface, session: AsyncSession, variable_id: UUID
59) -> Optional[orm_models.Variable]:
60 """
61 Reads a variable by id.
62 """
64 query = sa.select(db.Variable).where(db.Variable.id == variable_id)
66 result = await session.execute(query)
67 return result.scalar()
70@db_injector
71async def read_variable_by_name(
72 db: PrefectDBInterface, session: AsyncSession, name: str
73) -> Optional[orm_models.Variable]:
74 """
75 Reads a variable by name.
76 """
78 query = sa.select(db.Variable).where(db.Variable.name == name)
80 result = await session.execute(query)
81 return result.scalar()
84@db_injector
85async def read_variables(
86 db: PrefectDBInterface,
87 session: AsyncSession,
88 variable_filter: Optional[filters.VariableFilter] = None,
89 sort: sorting.VariableSort = sorting.VariableSort.NAME_ASC,
90 offset: Optional[int] = None,
91 limit: Optional[int] = None,
92) -> Sequence[orm_models.Variable]:
93 """
94 Read variables, applying filers.
95 """
96 query = sa.select(db.Variable).order_by(*sort.as_sql_sort())
98 if variable_filter:
99 query = query.where(variable_filter.as_sql_filter())
101 if offset is not None:
102 query = query.offset(offset)
103 if limit is not None:
104 query = query.limit(limit)
106 result = await session.execute(query)
107 return result.scalars().unique().all()
110@db_injector
111async def count_variables(
112 db: PrefectDBInterface,
113 session: AsyncSession,
114 variable_filter: Optional[filters.VariableFilter] = None,
115) -> int:
116 """
117 Count variables, applying filters.
118 """
120 query = sa.select(sa.func.count()).select_from(db.Variable)
122 if variable_filter:
123 query = query.where(variable_filter.as_sql_filter())
125 result = await session.execute(query)
126 return result.scalar_one()
129@db_injector
130async def update_variable(
131 db: PrefectDBInterface,
132 session: AsyncSession,
133 variable_id: UUID,
134 variable: VariableUpdate,
135) -> bool:
136 """
137 Updates a variable by id.
138 """
139 existing = await read_variable(session, variable_id)
140 if existing is None:
141 return False
143 query = (
144 sa.update(db.Variable)
145 .where(db.Variable.id == variable_id)
146 .values(**variable.model_dump_for_orm(exclude_unset=True))
147 )
148 await session.execute(query)
150 await session.refresh(existing)
151 await emit_variable_updated_event(existing)
152 return True
155@db_injector
156async def update_variable_by_name(
157 db: PrefectDBInterface, session: AsyncSession, name: str, variable: VariableUpdate
158) -> bool:
159 """
160 Updates a variable by name.
161 """
162 existing = await read_variable_by_name(session, name)
163 if existing is None:
164 return False
166 query = (
167 sa.update(db.Variable)
168 .where(db.Variable.id == existing.id)
169 .values(**variable.model_dump_for_orm(exclude_unset=True))
170 )
171 await session.execute(query)
173 await session.refresh(existing)
174 await emit_variable_updated_event(existing)
175 return True
178@db_injector
179async def delete_variable(
180 db: PrefectDBInterface, session: AsyncSession, variable_id: UUID
181) -> bool:
182 """
183 Delete a variable by id.
184 """
185 existing = await read_variable(session, variable_id)
186 if existing is None:
187 return False
189 await emit_variable_deleted_event(existing)
191 query = sa.delete(db.Variable).where(db.Variable.id == variable_id)
192 await session.execute(query)
193 return True
196@db_injector
197async def delete_variable_by_name(
198 db: PrefectDBInterface, session: AsyncSession, name: str
199) -> bool:
200 """
201 Delete a variable by name.
202 """
203 existing = await read_variable_by_name(session, name)
204 if existing is None:
205 return False
207 await emit_variable_deleted_event(existing)
209 query = sa.delete(db.Variable).where(db.Variable.id == existing.id)
210 await session.execute(query)
211 return True