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

1from typing import Optional, Sequence 

2from uuid import UUID 

3 

4import sqlalchemy as sa 

5from sqlalchemy.ext.asyncio import AsyncSession 

6 

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 

13 

14 

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"))) 

19 

20 

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"))) 

25 

26 

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"))) 

31 

32 

33@db_injector 

34async def create_variable( 

35 db: PrefectDBInterface, session: AsyncSession, variable: VariableCreate 

36) -> orm_models.Variable: 

37 """ 

38 Create a variable 

39 

40 Args: 

41 session: async database session 

42 variable: variable to create 

43 

44 Returns: 

45 orm_models.Variable 

46 """ 

47 model = db.Variable(**variable.model_dump()) 

48 session.add(model) 

49 await session.flush() 

50 

51 await emit_variable_created_event(model) 

52 

53 return model 

54 

55 

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 """ 

63 

64 query = sa.select(db.Variable).where(db.Variable.id == variable_id) 

65 

66 result = await session.execute(query) 

67 return result.scalar() 

68 

69 

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 """ 

77 

78 query = sa.select(db.Variable).where(db.Variable.name == name) 

79 

80 result = await session.execute(query) 

81 return result.scalar() 

82 

83 

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()) 

97 

98 if variable_filter: 

99 query = query.where(variable_filter.as_sql_filter()) 

100 

101 if offset is not None: 

102 query = query.offset(offset) 

103 if limit is not None: 

104 query = query.limit(limit) 

105 

106 result = await session.execute(query) 

107 return result.scalars().unique().all() 

108 

109 

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 """ 

119 

120 query = sa.select(sa.func.count()).select_from(db.Variable) 

121 

122 if variable_filter: 

123 query = query.where(variable_filter.as_sql_filter()) 

124 

125 result = await session.execute(query) 

126 return result.scalar_one() 

127 

128 

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 

142 

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) 

149 

150 await session.refresh(existing) 

151 await emit_variable_updated_event(existing) 

152 return True 

153 

154 

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 

165 

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) 

172 

173 await session.refresh(existing) 

174 await emit_variable_updated_event(existing) 

175 return True 

176 

177 

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 

188 

189 await emit_variable_deleted_event(existing) 

190 

191 query = sa.delete(db.Variable).where(db.Variable.id == variable_id) 

192 await session.execute(query) 

193 return True 

194 

195 

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 

206 

207 await emit_variable_deleted_event(existing) 

208 

209 query = sa.delete(db.Variable).where(db.Variable.id == existing.id) 

210 await session.execute(query) 

211 return True