Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/execution_api/routes/assets.py: 77%

24 statements  

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

1# Licensed to the Apache Software Foundation (ASF) under one 

2# or more contributor license agreements. See the NOTICE file 

3# distributed with this work for additional information 

4# regarding copyright ownership. The ASF licenses this file 

5# to you under the Apache License, Version 2.0 (the 

6# "License"); you may not use this file except in compliance 

7# with the License. You may obtain a copy of the License at 

8# 

9# http://www.apache.org/licenses/LICENSE-2.0 

10# 

11# Unless required by applicable law or agreed to in writing, 

12# software distributed under the License is distributed on an 

13# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY 

14# KIND, either express or implied. See the License for the 

15# specific language governing permissions and limitations 

16# under the License. 

17 

18from __future__ import annotations 

19 

20from typing import Annotated 

21 

22from fastapi import APIRouter, HTTPException, Query, status 

23from sqlalchemy import select 

24 

25from airflow.api_fastapi.common.db.common import SessionDep 

26from airflow.api_fastapi.execution_api.datamodels.asset import AssetResponse 

27from airflow.models.asset import AssetModel, expand_alias_to_assets 

28 

29router = APIRouter( 

30 responses={ 

31 status.HTTP_404_NOT_FOUND: {"description": "Asset not found"}, 

32 status.HTTP_401_UNAUTHORIZED: {"description": "Unauthorized"}, 

33 }, 

34) 

35 

36 

37@router.get("/by-name") 

38def get_asset_by_name( 

39 name: Annotated[str, Query(description="The name of the Asset")], 

40 session: SessionDep, 

41) -> AssetResponse: 

42 """Get an Airflow Asset by `name`.""" 

43 asset = session.scalar(select(AssetModel).where(AssetModel.name == name, AssetModel.active.has())) 

44 _raise_if_not_found(asset, f"Asset with name {name} not found") 

45 

46 return AssetResponse.model_validate(asset) 

47 

48 

49@router.get("/by-uri") 

50def get_asset_by_uri( 

51 uri: Annotated[str, Query(description="The URI of the Asset")], 

52 session: SessionDep, 

53) -> AssetResponse: 

54 """Get an Airflow Asset by `uri`.""" 

55 asset = session.scalar(select(AssetModel).where(AssetModel.uri == uri, AssetModel.active.has())) 

56 _raise_if_not_found(asset, f"Asset with URI {uri} not found") 

57 

58 return AssetResponse.model_validate(asset) 

59 

60 

61@router.get("/by-alias") 

62def get_assets_by_alias( 

63 alias_name: Annotated[str, Query(description="The name of the AssetAlias")], 

64 session: SessionDep, 

65) -> list[AssetResponse]: 

66 """Get all Airflow Assets resolved from an AssetAlias by `alias_name`.""" 

67 return [AssetResponse.model_validate(a) for a in expand_alias_to_assets(alias_name, session=session)] 

68 

69 

70def _raise_if_not_found(asset, msg): 

71 if asset is None: 71 ↛ 72line 71 didn't jump to line 72 because the condition on line 71 was never true

72 raise HTTPException( 

73 status.HTTP_404_NOT_FOUND, 

74 detail={ 

75 "reason": "not_found", 

76 "message": msg, 

77 }, 

78 )