Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/models/block_documents.py: 43%

236 statements  

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

1""" 

2Functions for interacting with block document ORM objects. 

3Intended for internal use by the Prefect REST API. 

4""" 

5 

6from copy import copy 

7from typing import Dict, List, Optional, Sequence, Tuple, TypeVar, Union 

8from uuid import UUID, uuid4 

9 

10import sqlalchemy as sa 

11from sqlalchemy.ext.asyncio import AsyncSession 

12from sqlalchemy.sql import Select 

13 

14import prefect.server.models as models 

15from prefect.server import schemas 

16from prefect.server.database import PrefectDBInterface, db_injector, orm_models 

17from prefect.server.events import clients 

18from prefect.server.events.schemas import lifecycle 

19from prefect.server.schemas.actions import BlockDocumentReferenceCreate 

20from prefect.server.schemas.core import BlockDocument 

21from prefect.server.schemas.filters import BlockSchemaFilter 

22from prefect.server.utilities.database import UUID as UUIDTypeDecorator 

23from prefect.types._datetime import now 

24from prefect.utilities.collections import dict_to_flatdict, flatdict_to_dict 

25from prefect.utilities.names import obfuscate 

26 

27T = TypeVar("T", bound=tuple) 

28 

29 

30async def emit_block_document_created_event(block_document: BlockDocument) -> None: 

31 """Emit an event when a (non-anonymous) block document is created.""" 

32 if block_document.is_anonymous: 

33 return 

34 async with clients.PrefectServerEventsClient() as events_client: 

35 await events_client.emit( 

36 lifecycle.block_document_created_event(block_document, now("UTC")) 

37 ) 

38 

39 

40async def emit_block_document_updated_event(block_document: BlockDocument) -> None: 

41 """Emit an event when a (non-anonymous) block document is updated.""" 

42 if block_document.is_anonymous: 42 ↛ 44line 42 didn't jump to line 44 because the condition on line 42 was always true

43 return 

44 async with clients.PrefectServerEventsClient() as events_client: 

45 await events_client.emit( 

46 lifecycle.block_document_updated_event(block_document, now("UTC")) 

47 ) 

48 

49 

50async def emit_block_document_deleted_event(block_document: BlockDocument) -> None: 

51 """Emit an event when a (non-anonymous) block document is deleted.""" 

52 if block_document.is_anonymous: 

53 return 

54 async with clients.PrefectServerEventsClient() as events_client: 

55 await events_client.emit( 

56 lifecycle.block_document_deleted_event(block_document, now("UTC")) 

57 ) 

58 

59 

60@db_injector 

61async def create_block_document( 

62 db: PrefectDBInterface, 

63 session: AsyncSession, 

64 block_document: schemas.actions.BlockDocumentCreate, 

65) -> BlockDocument: 

66 # lookup block type name and copy to the block document table 

67 block_type = await models.block_types.read_block_type( 

68 session=session, block_type_id=block_document.block_type_id 

69 ) 

70 assert block_type, f"Block type {block_document.block_type_id} not found" 

71 

72 name: Union[str, None] 

73 # anonymous block documents can be given a random name if none is provided 

74 if block_document.is_anonymous and not block_document.name: 

75 name = f"anonymous-{uuid4()}" 

76 else: 

77 name = block_document.name 

78 

79 orm_block = db.BlockDocument( 

80 name=name, 

81 block_schema_id=block_document.block_schema_id, 

82 block_type_id=block_document.block_type_id, 

83 block_type_name=block_type.name, 

84 is_anonymous=block_document.is_anonymous, 

85 ) 

86 

87 ( 

88 block_document_data_without_refs, 

89 block_document_references, 

90 ) = _separate_block_references_from_data(block_document.data) 

91 

92 # encrypt the data and store in block document 

93 await orm_block.encrypt_data(session=session, data=block_document_data_without_refs) 

94 

95 session.add(orm_block) 

96 await session.flush() 

97 

98 # Create a block document reference for each reference in the block document data 

99 for key, reference_block_document_id in block_document_references: 

100 await create_block_document_reference( 

101 session=session, 

102 block_document_reference=BlockDocumentReferenceCreate( 

103 parent_block_document_id=orm_block.id, 

104 reference_block_document_id=reference_block_document_id, 

105 name=key, 

106 ), 

107 ) 

108 

109 # reload the block document in order to load the associated block schema 

110 # relationship 

111 new_block_document = await read_block_document_by_id( 

112 session=session, 

113 block_document_id=orm_block.id, 

114 include_secrets=False, 

115 ) 

116 assert new_block_document 

117 

118 await emit_block_document_created_event(new_block_document) 

119 

120 return new_block_document 

121 

122 

123@db_injector 

124async def block_document_with_unique_values_exists( 

125 db: PrefectDBInterface, session: AsyncSession, block_type_id: UUID, name: str 

126) -> bool: 

127 result = await session.execute( 

128 sa.select(sa.exists(db.BlockDocument)).where( 

129 db.BlockDocument.block_type_id == block_type_id, 

130 db.BlockDocument.name == name, 

131 ) 

132 ) 

133 return bool(result.scalar_one_or_none()) 

134 

135 

136def _separate_block_references_from_data( 

137 block_document_data: Dict, 

138) -> Tuple[Dict, List[Tuple[str, UUID]]]: 

139 """ 

140 Separates block document references from block document data so that a block 

141 document can be saved without references and the corresponding block document 

142 references can be saved. 

143 

144 Args: 

145 block_document_data: Dictionary of block document data passed with request 

146 to create new block document 

147 

148 Returns: 

149 block_document_data_with_out_refs: A copy of the block_document_data supplied 

150 with the block document references removed. 

151 block_document_references: A list of tuples each containing the name of the 

152 field referencing a block document and the ID of the referenced block 

153 document, 

154 """ 

155 block_document_references = [] 

156 block_document_data_without_refs = {} 

157 for key, value in block_document_data.items(): 

158 # Current assumption is that any block references will be stored on a key of 

159 # the block document data and not nested in any other data structures. 

160 if isinstance(value, dict) and "$ref" in value: 160 ↛ 161line 160 didn't jump to line 161 because the condition on line 160 was never true

161 reference_block_document_id = value["$ref"].get("block_document_id") 

162 if reference_block_document_id is None: 

163 raise ValueError( 

164 f"Received block reference without a block_document_id in key {key}" 

165 ) 

166 block_document_references.append((key, reference_block_document_id)) 

167 else: 

168 block_document_data_without_refs[key] = value 

169 return block_document_data_without_refs, block_document_references 

170 

171 

172async def read_block_document_by_id( 

173 session: AsyncSession, 

174 block_document_id: UUID, 

175 include_secrets: bool = False, 

176) -> Union[BlockDocument, None]: 

177 block_documents = await read_block_documents( 

178 session=session, 

179 block_document_filter=schemas.filters.BlockDocumentFilter( 

180 id=dict(any_=[block_document_id]), 

181 # don't apply any anonymous filtering 

182 is_anonymous=None, 

183 ), 

184 include_secrets=include_secrets, 

185 limit=1, 

186 ) 

187 return block_documents[0] if block_documents else None 

188 

189 

190_BLOCK_DOCUMENT_REFERENCE_MAX_DEPTH = 50 

191 

192 

193async def _construct_full_block_document( 

194 db: PrefectDBInterface, 

195 session: AsyncSession, 

196 block_documents_with_references: Sequence[ 

197 Tuple[orm_models.ORMBlockDocument, Optional[str], Optional[UUID]] 

198 ], 

199 parent_block_document: Optional[BlockDocument] = None, 

200 include_secrets: bool = False, 

201 _depth: int = 0, 

202 _max_depth: int = _BLOCK_DOCUMENT_REFERENCE_MAX_DEPTH, 

203) -> Optional[BlockDocument]: 

204 if _depth > _max_depth: 204 ↛ 205line 204 didn't jump to line 205 because the condition on line 204 was never true

205 raise ValueError( 

206 "Block document reference graph exceeds max depth " 

207 f"{_max_depth}; refusing to recurse further." 

208 ) 

209 if len(block_documents_with_references) == 0: 209 ↛ 210line 209 didn't jump to line 210 because the condition on line 209 was never true

210 return None 

211 if parent_block_document is None: 211 ↛ 212line 211 didn't jump to line 212 because the condition on line 211 was never true

212 parent_block_document = copy( 

213 await _find_parent_block_document( 

214 session, 

215 block_documents_with_references, 

216 include_secrets=include_secrets, 

217 ) 

218 ) 

219 

220 if parent_block_document is None: 220 ↛ 221line 220 didn't jump to line 221 because the condition on line 220 was never true

221 raise ValueError("Unable to determine parent block document") 

222 

223 # Recursively walk block document tree and construct the full block 

224 # document data for each child and add it to the parent's block document 

225 # data 

226 for ( 

227 orm_block_document, 

228 name, 

229 parent_block_document_id, 

230 ) in block_documents_with_references: 

231 if parent_block_document_id == parent_block_document.id and name is not None: 231 ↛ 232line 231 didn't jump to line 232 because the condition on line 231 was never true

232 block_document = await BlockDocument.from_orm_model( 

233 session, orm_block_document, include_secrets=include_secrets 

234 ) 

235 full_child_block_document = await _construct_full_block_document( 

236 db, 

237 session, 

238 block_documents_with_references, 

239 parent_block_document=copy(block_document), 

240 include_secrets=include_secrets, 

241 _depth=_depth + 1, 

242 _max_depth=_max_depth, 

243 ) 

244 assert full_child_block_document 

245 parent_block_document.data[name] = full_child_block_document.data 

246 parent_block_document.block_document_references[name] = { 

247 "block_document": { 

248 "id": block_document.id, 

249 "name": block_document.name, 

250 "block_type": block_document.block_type, 

251 "is_anonymous": block_document.is_anonymous, 

252 "block_document_references": ( 

253 full_child_block_document.block_document_references 

254 ), 

255 } 

256 } 

257 

258 return parent_block_document 

259 

260 

261async def _find_parent_block_document( 

262 session: AsyncSession, 

263 block_documents_with_references: Sequence[ 

264 Tuple[orm_models.ORMBlockDocument, Optional[str], Optional[UUID]] 

265 ], 

266 include_secrets: bool = False, 

267) -> Union[BlockDocument, None]: 

268 parent_orm_block_document = next( 

269 ( 

270 block_document 

271 for ( 

272 block_document, 

273 _, 

274 parent_block_document_id, 

275 ) in block_documents_with_references 

276 if parent_block_document_id is None 

277 ), 

278 None, 

279 ) 

280 return ( 

281 await BlockDocument.from_orm_model( 

282 session, 

283 parent_orm_block_document, 

284 include_secrets=include_secrets, 

285 ) 

286 if parent_orm_block_document is not None 

287 else None 

288 ) 

289 

290 

291async def read_block_document_by_name( 

292 session: AsyncSession, 

293 name: str, 

294 block_type_slug: str, 

295 include_secrets: bool = False, 

296) -> Union[BlockDocument, None]: 

297 """ 

298 Read a block document with the given name and block type slug. 

299 """ 

300 block_documents = await read_block_documents( 

301 session=session, 

302 block_document_filter=schemas.filters.BlockDocumentFilter( 

303 name=dict(any_=[name]), 

304 # don't apply any anonymous filtering 

305 is_anonymous=None, 

306 ), 

307 block_type_filter=schemas.filters.BlockTypeFilter( 

308 slug=dict(any_=[block_type_slug]) 

309 ), 

310 include_secrets=include_secrets, 

311 limit=1, 

312 ) 

313 return block_documents[0] if block_documents else None 

314 

315 

316def _apply_block_document_filters( 

317 db: PrefectDBInterface, 

318 query: Select[T], 

319 block_document_filter: Optional[schemas.filters.BlockDocumentFilter] = None, 

320 block_schema_filter: Optional[schemas.filters.BlockSchemaFilter] = None, 

321 block_type_filter: Optional[schemas.filters.BlockTypeFilter] = None, 

322) -> Select[T]: 

323 # if no filter is provided, one is created that excludes anonymous blocks 

324 if block_document_filter is None: 

325 block_document_filter = schemas.filters.BlockDocumentFilter( 

326 is_anonymous=schemas.filters.BlockDocumentFilterIsAnonymous(eq_=False) 

327 ) 

328 

329 # --- Build an initial query that filters for the requested block documents 

330 query = query.where(block_document_filter.as_sql_filter()) 

331 

332 if block_type_filter is not None: 

333 block_type_exists_clause = sa.select(db.BlockType).where( 

334 db.BlockType.id == db.BlockDocument.block_type_id, 

335 block_type_filter.as_sql_filter(), 

336 ) 

337 query = query.where(block_type_exists_clause.exists()) 

338 

339 if block_schema_filter is not None: 

340 block_schema_exists_clause = sa.select(db.BlockSchema).where( 

341 db.BlockSchema.id == db.BlockDocument.block_schema_id, 

342 block_schema_filter.as_sql_filter(), 

343 ) 

344 query = query.where(block_schema_exists_clause.exists()) 

345 

346 return query 

347 

348 

349@db_injector 

350async def read_block_documents( 

351 db: PrefectDBInterface, 

352 session: AsyncSession, 

353 block_document_filter: Optional[schemas.filters.BlockDocumentFilter] = None, 

354 block_type_filter: Optional[schemas.filters.BlockTypeFilter] = None, 

355 block_schema_filter: Optional[schemas.filters.BlockSchemaFilter] = None, 

356 include_secrets: bool = False, 

357 sort: schemas.sorting.BlockDocumentSort = schemas.sorting.BlockDocumentSort.NAME_ASC, 

358 offset: Optional[int] = None, 

359 limit: Optional[int] = None, 

360) -> List[BlockDocument]: 

361 """ 

362 Read block documents with an optional limit and offset 

363 """ 

364 # --- Build an initial query that filters for the requested block documents 

365 filtered_block_documents_query = sa.select(db.BlockDocument.id) 

366 filtered_block_documents_query = _apply_block_document_filters( 

367 db, 

368 query=filtered_block_documents_query, 

369 block_document_filter=block_document_filter, 

370 block_type_filter=block_type_filter, 

371 block_schema_filter=block_schema_filter, 

372 ) 

373 filtered_block_documents_query = filtered_block_documents_query.order_by( 

374 *sort.as_sql_sort() 

375 ) 

376 

377 if offset is not None: 

378 filtered_block_documents_query = filtered_block_documents_query.offset(offset) 

379 

380 if limit is not None: 

381 filtered_block_documents_query = filtered_block_documents_query.limit(limit) 

382 

383 filtered_block_documents_cte = filtered_block_documents_query.cte( 

384 "filtered_block_documents" 

385 ) 

386 

387 # --- Build a recursive query that starts with the filtered block documents 

388 # and iteratively loads all referenced block documents. The query includes 

389 # the ID of each block document as well as the ID of the document that 

390 # references it and name it's referenced by, if applicable. 

391 parent_documents = ( 

392 sa.select( 

393 filtered_block_documents_cte.c.id, 

394 sa.cast(sa.null(), sa.String).label("reference_name"), 

395 sa.cast(sa.null(), UUIDTypeDecorator).label( 

396 "reference_parent_block_document_id" 

397 ), 

398 ) 

399 .select_from(filtered_block_documents_cte) 

400 .cte("all_block_documents", recursive=True) 

401 ) 

402 # recursive part of query 

403 referenced_documents = ( 

404 sa.select( 

405 db.BlockDocumentReference.reference_block_document_id, 

406 db.BlockDocumentReference.name, 

407 db.BlockDocumentReference.parent_block_document_id, 

408 ) 

409 .select_from(parent_documents) 

410 .join( 

411 db.BlockDocumentReference, 

412 db.BlockDocumentReference.parent_block_document_id == parent_documents.c.id, 

413 ) 

414 ) 

415 # union the recursive CTE 

416 all_block_documents_query = parent_documents.union_all(referenced_documents) 

417 

418 # --- Join the recursive query that contains all required document IDs 

419 # back to the BlockDocument table to load info for every document 

420 # and order by name 

421 final_query = ( 

422 sa.select( 

423 db.BlockDocument, 

424 all_block_documents_query.c.reference_name, 

425 all_block_documents_query.c.reference_parent_block_document_id, 

426 ) 

427 .select_from(all_block_documents_query) 

428 .join(db.BlockDocument, db.BlockDocument.id == all_block_documents_query.c.id) 

429 .order_by(*sort.as_sql_sort()) 

430 ) 

431 

432 result = await session.execute( 

433 final_query.execution_options(populate_existing=True) 

434 ) 

435 

436 block_documents_with_references = result.unique().all() 

437 

438 # identify true "parent" documents as those with no reference parent ids 

439 parent_block_document_ids = [ 

440 d[0].id 

441 for d in block_documents_with_references 

442 if d.reference_parent_block_document_id is None 

443 ] 

444 

445 # walk the resulting dataset and hydrate all block documents 

446 fully_constructed_block_documents: List[BlockDocument] = [] 

447 visited_block_document_ids = [] 

448 for root_orm_block_document, _, _ in block_documents_with_references: 

449 if ( 

450 root_orm_block_document.id in parent_block_document_ids 

451 and root_orm_block_document.id not in visited_block_document_ids 

452 ): 

453 root_block_document = await BlockDocument.from_orm_model( 

454 session, root_orm_block_document, include_secrets=include_secrets 

455 ) 

456 constructed = await _construct_full_block_document( 

457 db, 

458 session, 

459 block_documents_with_references, # type: ignore 

460 root_block_document, 

461 include_secrets=include_secrets, 

462 ) 

463 assert constructed 

464 

465 fully_constructed_block_documents.append(constructed) 

466 visited_block_document_ids.append(root_orm_block_document.id) 

467 

468 block_schema_ids = [ 

469 block_document.block_schema_id 

470 for block_document in fully_constructed_block_documents 

471 ] 

472 block_schemas = await models.block_schemas.read_block_schemas( 

473 session=session, 

474 block_schema_filter=BlockSchemaFilter(id=dict(any_=block_schema_ids)), 

475 ) 

476 for block_document in fully_constructed_block_documents: 

477 corresponding_block_schema = next( 

478 block_schema 

479 for block_schema in block_schemas 

480 if block_schema.id == block_document.block_schema_id 

481 ) 

482 block_document.block_schema = corresponding_block_schema 

483 

484 return fully_constructed_block_documents 

485 

486 

487@db_injector 

488async def count_block_documents( 

489 db: PrefectDBInterface, 

490 session: AsyncSession, 

491 block_document_filter: Optional[schemas.filters.BlockDocumentFilter] = None, 

492 block_type_filter: Optional[schemas.filters.BlockTypeFilter] = None, 

493 block_schema_filter: Optional[schemas.filters.BlockSchemaFilter] = None, 

494) -> int: 

495 """ 

496 Count block documents that match the filters. 

497 """ 

498 query = sa.select(sa.func.count()).select_from(db.BlockDocument) 

499 

500 query = _apply_block_document_filters( 

501 db, 

502 query=query, 

503 block_document_filter=block_document_filter, 

504 block_schema_filter=block_schema_filter, 

505 block_type_filter=block_type_filter, 

506 ) 

507 

508 result = await session.execute(query) 

509 return result.scalar() # type: ignore 

510 

511 

512@db_injector 

513async def delete_block_document( 

514 db: PrefectDBInterface, 

515 session: AsyncSession, 

516 block_document_id: UUID, 

517) -> bool: 

518 block_document = await read_block_document_by_id( 

519 session=session, 

520 block_document_id=block_document_id, 

521 include_secrets=False, 

522 ) 

523 if block_document is None: 

524 return False 

525 

526 await emit_block_document_deleted_event(block_document) 

527 

528 query = sa.delete(db.BlockDocument).where(db.BlockDocument.id == block_document_id) 

529 await session.execute(query) 

530 

531 await models.storage_defaults.clear_server_default_result_storage_for_block( 

532 session=session, 

533 block_document_id=block_document_id, 

534 ) 

535 return True 

536 

537 

538@db_injector 

539async def update_block_document( 

540 db: PrefectDBInterface, 

541 session: AsyncSession, 

542 block_document_id: UUID, 

543 block_document: schemas.actions.BlockDocumentUpdate, 

544) -> bool: 

545 merge_existing_data = block_document.merge_existing_data 

546 current_block_document = await session.get(db.BlockDocument, block_document_id) 

547 if not current_block_document: 

548 return False 

549 

550 update_values = block_document.model_dump_for_orm( 

551 exclude_unset=merge_existing_data, 

552 exclude={"merge_existing_data"}, 

553 ) 

554 

555 if "data" in update_values and update_values["data"] is not None: 

556 current_data = await current_block_document.decrypt_data(session=session) 

557 

558 # if a value for a secret field was provided that is identical to the 

559 # obfuscated value of the current secret value, it means someone is 

560 # probably trying to update all of the documents fields without 

561 # realizing they are posting back obfuscated data, so we disregard the update 

562 flat_update_data = dict_to_flatdict(update_values["data"]) 

563 flat_current_data = dict_to_flatdict(current_data) 

564 for secret_field in current_block_document.block_schema.fields.get( 

565 "secret_fields", [] 

566 ): 

567 secret_key = tuple(secret_field.split(".")) 

568 current_secret = flat_current_data.get(secret_key) 

569 if current_secret is not None: 

570 if flat_update_data.get(secret_key) == obfuscate(current_secret): 

571 flat_update_data[secret_key] = current_secret 

572 # Looks for obfuscated values nested under a secret field with a wildcard. 

573 # If any obfuscated values are found, we assume that it shouldn't be update, 

574 # and they are replaced with the current value for that key to avoid losing 

575 # data during update. 

576 elif "*" in secret_key: 

577 wildcard_index = secret_key.index("*") 

578 for data_key in flat_update_data.keys(): 

579 if ( 

580 secret_key[0:wildcard_index] == data_key[0:wildcard_index] 

581 ) and ( 

582 flat_update_data[data_key] 

583 == obfuscate(flat_update_data[data_key]) 

584 ): 

585 flat_update_data[data_key] = flat_current_data[data_key] 

586 

587 update_values["data"] = flatdict_to_dict(flat_update_data) 

588 

589 if merge_existing_data: 

590 # merge the existing data and the new data for partial updates 

591 current_data.update(update_values["data"]) 

592 update_values["data"] = current_data 

593 

594 current_block_document_references = ( 

595 ( 

596 await session.execute( 

597 sa.select(db.BlockDocumentReference).filter_by( 

598 parent_block_document_id=block_document_id 

599 ) 

600 ) 

601 ) 

602 .scalars() 

603 .all() 

604 ) 

605 ( 

606 block_document_data_without_refs, 

607 new_block_document_references, 

608 ) = _separate_block_references_from_data(update_values["data"]) 

609 

610 # encrypt the data and write updated data to the block document 

611 await current_block_document.encrypt_data( 

612 session=session, data=block_document_data_without_refs 

613 ) 

614 

615 # `proposed_block_schema` is always the same as the schema on the client-side 

616 # Block class that is calling `save`, which may or may not be the same schema 

617 # as the one on the saved block document 

618 proposed_block_schema_id = block_document.block_schema_id 

619 

620 # if a new schema is proposed, update the block schema id for the block document 

621 if ( 

622 proposed_block_schema_id is not None 

623 and proposed_block_schema_id != current_block_document.block_schema_id 

624 ): 

625 proposed_block_schema = await session.get( 

626 db.BlockSchema, proposed_block_schema_id 

627 ) 

628 assert proposed_block_schema, ( 

629 f"Block schema {proposed_block_schema_id} not found" 

630 ) 

631 

632 # make sure the proposed schema is of the same block type as the current document 

633 if ( 

634 proposed_block_schema.block_type_id 

635 != current_block_document.block_type_id 

636 ): 

637 raise ValueError( 

638 "Must migrate block document to a block schema of the same block" 

639 " type." 

640 ) 

641 await session.execute( 

642 sa.update(db.BlockDocument) 

643 .where(db.BlockDocument.id == block_document_id) 

644 .values(block_schema_id=proposed_block_schema_id) 

645 ) 

646 

647 unchanged_block_document_references = [] 

648 for name, reference_block_document_id in new_block_document_references: 

649 matching_current_block_document_reference = _find_block_document_reference( 

650 current_block_document_references, 

651 name, 

652 reference_block_document_id, 

653 ) 

654 if matching_current_block_document_reference is None: 

655 await create_block_document_reference( 

656 session=session, 

657 block_document_reference=BlockDocumentReferenceCreate( 

658 parent_block_document_id=block_document_id, 

659 reference_block_document_id=reference_block_document_id, 

660 name=name, 

661 ), 

662 ) 

663 else: 

664 unchanged_block_document_references.append( 

665 matching_current_block_document_reference 

666 ) 

667 

668 for block_document_reference in current_block_document_references: 

669 if block_document_reference not in unchanged_block_document_references: 

670 await delete_block_document_reference( 

671 session, block_document_reference_id=block_document_reference.id 

672 ) 

673 

674 await session.flush() 

675 updated_block_document = await read_block_document_by_id( 

676 session=session, 

677 block_document_id=block_document_id, 

678 include_secrets=False, 

679 ) 

680 if updated_block_document is not None: 

681 await emit_block_document_updated_event(updated_block_document) 

682 

683 return True 

684 

685 

686def _find_block_document_reference( 

687 block_document_references: Sequence[orm_models.BlockDocumentReference], 

688 name: str, 

689 reference_block_document_id: UUID, 

690) -> Optional[orm_models.BlockDocumentReference]: 

691 return next( 

692 ( 

693 block_document_reference 

694 for block_document_reference in block_document_references 

695 if block_document_reference.name == name 

696 and block_document_reference.reference_block_document_id 

697 == reference_block_document_id 

698 ), 

699 None, 

700 ) 

701 

702 

703@db_injector 

704async def _block_document_reference_would_form_cycle( 

705 db: PrefectDBInterface, 

706 session: AsyncSession, 

707 parent_block_document_id: UUID, 

708 reference_block_document_id: UUID, 

709) -> bool: 

710 """Return True if adding a reference edge (parent → reference) would 

711 introduce a cycle into the block_document_reference graph. 

712 

713 A self-reference (parent == reference) is treated as a cycle. For 

714 multi-hop cases, walk forward from `reference_block_document_id` along 

715 its existing reference edges and report a cycle if `parent_block_document_id` 

716 is reachable. 

717 """ 

718 if parent_block_document_id == reference_block_document_id: 

719 return True 

720 

721 visited: set[UUID] = set() 

722 stack: list[UUID] = [reference_block_document_id] 

723 

724 while stack: 

725 current = stack.pop() 

726 if current in visited: 

727 continue 

728 visited.add(current) 

729 

730 if current == parent_block_document_id: 

731 return True 

732 

733 result = await session.execute( 

734 sa.select(db.BlockDocumentReference.reference_block_document_id).where( 

735 db.BlockDocumentReference.parent_block_document_id == current 

736 ) 

737 ) 

738 for child_id in result.scalars().all(): 

739 if child_id not in visited: 

740 stack.append(child_id) 

741 

742 return False 

743 

744 

745@db_injector 

746async def create_block_document_reference( 

747 db: PrefectDBInterface, 

748 session: AsyncSession, 

749 block_document_reference: schemas.actions.BlockDocumentReferenceCreate, 

750) -> Union[orm_models.BlockDocumentReference, None]: 

751 if await _block_document_reference_would_form_cycle( 

752 session=session, 

753 parent_block_document_id=block_document_reference.parent_block_document_id, 

754 reference_block_document_id=block_document_reference.reference_block_document_id, 

755 ): 

756 raise ValueError( 

757 "Cannot create block document reference: it would introduce a " 

758 "cycle into the block_document_reference graph " 

759 f"(parent={block_document_reference.parent_block_document_id}, " 

760 f"reference={block_document_reference.reference_block_document_id})." 

761 ) 

762 

763 insert_stmt = db.queries.insert(db.BlockDocumentReference).values( 

764 **block_document_reference.model_dump_for_orm( 

765 exclude_unset=True, exclude={"created", "updated"} 

766 ) 

767 ) 

768 await session.execute(insert_stmt) 

769 

770 result = await session.execute( 

771 sa.select(db.BlockDocumentReference).where( 

772 db.BlockDocumentReference.id == block_document_reference.id 

773 ) 

774 ) 

775 

776 return result.scalar() 

777 

778 

779@db_injector 

780async def delete_block_document_reference( 

781 db: PrefectDBInterface, 

782 session: AsyncSession, 

783 block_document_reference_id: UUID, 

784) -> bool: 

785 query = sa.delete(db.BlockDocumentReference).where( 

786 db.BlockDocumentReference.id == block_document_reference_id 

787 ) 

788 result = await session.execute(query) 

789 return result.rowcount > 0