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
« 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.
18from __future__ import annotations
20from typing import Annotated
22from fastapi import APIRouter, HTTPException, Query, status
23from sqlalchemy import select
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
29router = APIRouter(
30 responses={
31 status.HTTP_404_NOT_FOUND: {"description": "Asset not found"},
32 status.HTTP_401_UNAUTHORIZED: {"description": "Unauthorized"},
33 },
34)
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")
46 return AssetResponse.model_validate(asset)
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")
58 return AssetResponse.model_validate(asset)
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)]
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 )