This is an automated email from the ASF dual-hosted git repository.
bbovenzi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new 938a9d1f06f Add reparse button to the Dag import errors modal (#72896)
938a9d1f06f is described below
commit 938a9d1f06faf2a763e552c34c956be5ab843b99
Author: Pierre Jeambrun <[email protected]>
AuthorDate: Wed Sep 23 19:37:33 2026 +0200
Add reparse button to the Dag import errors modal (#72896)
* Add reparse button to the Dag import errors modal
The import errors modal lists files that failed to parse but gave no way to
retry them; users had to wait for the next scheduled parse or trigger a reparse
elsewhere after fixing the file. Surfacing a reparse action next to each failed
file lets them retry from the same place they see the error.
* Add release note for gating reparse of files with no registered Dag
Reparsing a file with no registered Dag now returns 403 instead of 404 when
the caller lacks the new REPARSE_ALL permission -- an observable behaviour
change that warrants a release note.
* Check reparse permission before probing file existence
The two-step check on a file with no registered Dag returned 404 for a file
Airflow had never heard of and 403 for one with an import error, so an
unauthorized caller could tell them apart. Auth runs first, and the
existence
probe now serves only its own purpose (an authorized caller against a stale
file). The dashboard modal also invalidates the import-errors and DAG-stats
queries on a successful reparse, so a fixed error disappears from the list
and the resolved Dag shows up in the stats without a manual refresh.
* Refresh the Dags table after a successful reparse too
* Drop the invalidation comment
* Drop the reparse-mutation query invalidation
Reparse is asynchronous: the endpoint inserts a DagPriorityParsingRequest
row
and returns, and the dag processor picks it up on its next cycle.
Invalidating
the modal, dashboard, and Dags-table queries the instant the mutation
resolves
refetches against a DB the processor has not touched yet, so every refetch
returns the same rows and only marks the queries fresh again. The toast is
the
honest signal that the reparse has been queued.
* Poll the import-errors modal after a reparse is queued
Reparse enqueues a DagPriorityParsingRequest and returns; the dag processor
picks it up on its next cycle. Nothing tells the modal when that happens, so
a user watching the list has no way to see the row disappear without closing
and reopening the modal. Polling on the same auto-refresh cadence the
dashboard
cards use closes that gap. Polling only starts once a reparse has been
issued
in this modal session, so a modal that is opened and never mutated hits the
endpoint no more than it did before, and it stops when the modal is closed.
---
airflow-core/newsfragments/72896.significant.rst | 16 +++++
.../core_api/datamodels/import_error.py | 11 ++-
.../core_api/openapi/v2-rest-api-generated.yaml | 7 ++
.../src/airflow/api_fastapi/core_api/security.py | 29 +++++++-
.../airflow/ui/openapi-gen/requests/schemas.gen.ts | 8 ++-
.../airflow/ui/openapi-gen/requests/types.gen.ts | 4 ++
.../pages/Dashboard/Stats/DagImportErrorsModal.tsx | 69 ++++++++++++++++--
.../core_api/routes/public/test_dag_parsing.py | 83 +++++++++++++++++++++-
.../core_api/routes/public/test_import_error.py | 22 +++++-
.../unit/api_fastapi/core_api/test_security.py | 62 ++++++++++++++++
.../src/airflowctl/api/datamodels/generated.py | 7 ++
.../tests/airflow_ctl/api/test_operations.py | 1 +
12 files changed, 310 insertions(+), 9 deletions(-)
diff --git a/airflow-core/newsfragments/72896.significant.rst
b/airflow-core/newsfragments/72896.significant.rst
new file mode 100644
index 00000000000..44230aa6850
--- /dev/null
+++ b/airflow-core/newsfragments/72896.significant.rst
@@ -0,0 +1,16 @@
+Reparsing a file with no registered Dag now requires the ``REPARSE_ALL`` view
+
+A file that failed to import before defining any Dag has no per-Dag key to
authorize a reparse
+against, so ``PUT /parseDagFile/{file_token}`` previously returned ``404`` for
it. Reparsing such a
+file is now gated on a dedicated ``AccessView.REPARSE_ALL`` -- a write action
of its own, distinct
+from the permission to view import errors -- and the Dag import errors modal
gains a reparse action
+for these files.
+
+**Behaviour changes:**
+
+- ``PUT /parseDagFile/{file_token}`` for a file with no registered Dag
authorizes against
+ ``REPARSE_ALL`` and returns ``403`` when the caller lacks it, instead of
``404``. A file that has a
+ registered Dag keeps its existing per-Dag ``can_edit`` check.
+- With the simple auth manager the view is granted to the ``ADMIN`` role only.
With the FAB auth
+ manager it maps to a new ``All Reparses`` resource, which ships in the
``Admin`` role; custom roles
+ that need to reparse such files must be granted ``All Reparses.can_read``
explicitly.
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/import_error.py
b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/import_error.py
index 084434cadbf..d070df98fdb 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/import_error.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/import_error.py
@@ -19,9 +19,10 @@ from __future__ import annotations
from collections.abc import Iterable
from datetime import datetime
-from pydantic import Field
+from pydantic import Field, computed_field
from airflow.api_fastapi.core_api.base import BaseModel
+from airflow.api_fastapi.core_api.datamodels.dags import
_get_file_token_serializer
class ImportErrorResponse(BaseModel):
@@ -33,6 +34,14 @@ class ImportErrorResponse(BaseModel):
bundle_name: str | None
stacktrace: str = Field(alias="stack_trace")
+ @computed_field # type: ignore[prop-decorator]
+ @property
+ def file_token(self) -> str:
+ """Return a signed token identifying the file, used to request its
reparse."""
+ return _get_file_token_serializer().dumps(
+ {"bundle_name": self.bundle_name, "relative_fileloc":
self.filename}
+ )
+
class ImportErrorCollectionResponse(BaseModel):
"""Import Error Collection Response."""
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
index 203a8e3ab54..f8ebe110b55 100644
---
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
+++
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
@@ -15290,6 +15290,12 @@ components:
stack_trace:
type: string
title: Stack Trace
+ file_token:
+ type: string
+ title: File Token
+ description: Return a signed token identifying the file, used to
request
+ its reparse.
+ readOnly: true
type: object
required:
- import_error_id
@@ -15297,6 +15303,7 @@ components:
- filename
- bundle_name
- stack_trace
+ - file_token
title: ImportErrorResponse
description: Import Error Response.
JobCollectionResponse:
diff --git a/airflow-core/src/airflow/api_fastapi/core_api/security.py
b/airflow-core/src/airflow/api_fastapi/core_api/security.py
index 3217fccc33f..61c7444bde2 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/security.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/security.py
@@ -77,6 +77,7 @@ from airflow.models.dag import DagModel, DagRun, DagTag
from airflow.models.dag_version import DagVersion
from airflow.models.dagbundle import DagBundleModel
from airflow.models.dagwarning import DagWarning
+from airflow.models.errors import ParseImportError
from airflow.models.log import Log
from airflow.models.taskinstance import TaskInstance as TI
from airflow.models.team import Team
@@ -256,7 +257,33 @@ def requires_access_dag_from_file_token(
)
)
if not dag_ids:
- raise HTTPException(status.HTTP_404_NOT_FOUND, "File not found")
+ # A file with an import error has no registered Dag to authorize
per-Dag against, so
+ # reparsing it is gated on the dedicated ``REPARSE_ALL``
permission -- admin-by-default,
+ # scoped to the file's team via its bundle -- rather than on the
permission to view
+ # import errors. Reparse is an action, so it must not ride on
being able to see the error.
+ # The auth check runs before the existence check so an
unauthorized caller cannot tell a
+ # file with an import error apart from one Airflow has never heard
of.
+ team_name = (
+ DagBundleModel.get_team_name(payload["bundle_name"],
session=session)
+ if payload["bundle_name"]
+ else None
+ )
+ if not get_auth_manager().authorize_view(
+ access_view=AccessView.REPARSE_ALL, user=user,
team_name=team_name
+ ):
+ raise HTTPException(
+ status.HTTP_403_FORBIDDEN,
+ "You do not have permission to reparse files with no
registered Dag",
+ )
+ has_import_error = session.scalar(
+ select(ParseImportError.id).where(
+ ParseImportError.bundle_name == payload["bundle_name"],
+ ParseImportError.filename == payload["relative_fileloc"],
+ )
+ )
+ if has_import_error is None:
+ raise HTTPException(status.HTTP_404_NOT_FOUND, "File not
found")
+ return
dag_id_to_team = DagModel.get_dag_id_to_team_name_mapping(dag_ids,
session=session)
requests: list[IsAuthorizedDagRequest] = [
diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
index 2aca8e5aeab..61bbeeb57e1 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
@@ -5968,10 +5968,16 @@ export const $ImportErrorResponse = {
stack_trace: {
type: 'string',
title: 'Stack Trace'
+ },
+ file_token: {
+ type: 'string',
+ title: 'File Token',
+ description: 'Return a signed token identifying the file, used to
request its reparse.',
+ readOnly: true
}
},
type: 'object',
- required: ['import_error_id', 'timestamp', 'filename', 'bundle_name',
'stack_trace'],
+ required: ['import_error_id', 'timestamp', 'filename', 'bundle_name',
'stack_trace', 'file_token'],
title: 'ImportErrorResponse',
description: 'Import Error Response.'
} as const;
diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
index e175c7203aa..f2cc5a075d7 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
@@ -1607,6 +1607,10 @@ export type ImportErrorResponse = {
filename: string;
bundle_name: string | null;
stack_trace: string;
+ /**
+ * Return a signed token identifying the file, used to request its reparse.
+ */
+ readonly file_token: string;
};
/**
diff --git
a/airflow-core/src/airflow/ui/src/pages/Dashboard/Stats/DagImportErrorsModal.tsx
b/airflow-core/src/airflow/ui/src/pages/Dashboard/Stats/DagImportErrorsModal.tsx
index 2fe98868f51..4fe139d1a6b 100644
---
a/airflow-core/src/airflow/ui/src/pages/Dashboard/Stats/DagImportErrorsModal.tsx
+++
b/airflow-core/src/airflow/ui/src/pages/Dashboard/Stats/DagImportErrorsModal.tsx
@@ -20,16 +20,27 @@ import { useState } from "react";
import { Box, ClipboardRoot, Heading, HStack, Text } from "@chakra-ui/react";
import { useTranslation } from "react-i18next";
+import { AiOutlineFileSync } from "react-icons/ai";
import { LuFileWarning } from "react-icons/lu";
import { PiFilePy } from "react-icons/pi";
-import { useImportErrorServiceGetImportErrors } from "openapi/queries";
+import { useDagParsingServiceReparseDagFile,
useImportErrorServiceGetImportErrors } from "openapi/queries";
-import { Accordion, ClipboardIconButton, Modal, Pagination } from
"src/system-components";
+import {
+ Accordion,
+ ClipboardIconButton,
+ IconButton,
+ Modal,
+ Pagination,
+ toaster,
+} from "src/system-components";
import { SearchBar } from "src/components/SearchBar";
import Time from "src/components/Time";
+import { useConfig } from "src/queries/useConfig";
+import { createErrorToaster } from "src/utils";
+
type ImportDAGErrorModalProps = {
readonly onClose: () => void;
readonly open: boolean;
@@ -37,9 +48,47 @@ type ImportDAGErrorModalProps = {
const PAGE_LIMIT = 15;
+const ReparseButton = ({
+ fileToken,
+ onReparsed,
+}: {
+ readonly fileToken: string;
+ readonly onReparsed: () => void;
+}) => {
+ const { t: translate } = useTranslation(["components", "dag"]);
+
+ const { isPending, mutate } = useDagParsingServiceReparseDagFile({
+ onError: (error) => createErrorToaster(error, { titleKey:
"dag:parse.toaster.error.title" }, translate),
+ onSuccess: () => {
+ // Reparse is queued for the DAG processor, so the modal cannot know when
+ // it lands; polling starts on the first reparse and stops on close.
+ onReparsed();
+ toaster.create({
+ description: translate("dag:parse.toaster.success.description"),
+ title: translate("dag:parse.toaster.success.title"),
+ type: "success",
+ });
+ },
+ });
+
+ return (
+ <IconButton
+ data-testid="reparse-import-error"
+ label={translate("components:reparseDag")}
+ loading={isPending}
+ onClick={() => mutate({ fileToken })}
+ variant="outline"
+ >
+ <AiOutlineFileSync />
+ </IconButton>
+ );
+};
+
export const DagImportErrorsModal = ({ onClose, open }:
ImportDAGErrorModalProps) => {
const [page, setPage] = useState(1);
const [searchQuery, setSearchQuery] = useState("");
+ const [pollAfterReparse, setPollAfterReparse] = useState(false);
+ const autoRefreshInterval = useConfig("auto_refresh_interval") as number |
undefined;
const { data } = useImportErrorServiceGetImportErrors(
{
@@ -48,7 +97,14 @@ export const DagImportErrorsModal = ({ onClose, open }:
ImportDAGErrorModalProps
offset: PAGE_LIMIT * (page - 1),
},
undefined,
- { enabled: open },
+ {
+ enabled: open,
+ // Reparse is a queued request the dag processor picks up on its next
+ // cycle, so once a reparse has been issued the modal polls until it is
+ // closed. autoRefreshInterval mirrors what the dashboard cards use.
+ refetchInterval:
+ pollAfterReparse && autoRefreshInterval !== undefined ?
autoRefreshInterval * 1000 : false,
+ },
);
const { t: translate } = useTranslation(["dashboard", "components"]);
@@ -56,6 +112,7 @@ export const DagImportErrorsModal = ({ onClose, open }:
ImportDAGErrorModalProps
const onOpenChange = () => {
setSearchQuery("");
setPage(1);
+ setPollAfterReparse(false);
onClose();
};
@@ -121,7 +178,11 @@ export const DagImportErrorsModal = ({ onClose, open }:
ImportDAGErrorModalProps
{importError.filename}
</HStack>
</Accordion.ItemTrigger>
- <Box alignItems="center" display="flex" flexShrink={0} pr={2}>
+ <Box alignItems="center" display="flex" flexShrink={0} gap={1}
pr={2}>
+ <ReparseButton
+ fileToken={importError.file_token}
+ onReparsed={() => setPollAfterReparse(true)}
+ />
<ClipboardRoot value={importError.filename}>
<ClipboardIconButton variant="outline" />
</ClipboardRoot>
diff --git
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_parsing.py
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_parsing.py
index b60b7011a9e..7f645f81cc3 100644
---
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_parsing.py
+++
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_parsing.py
@@ -16,13 +16,23 @@
# under the License.
from __future__ import annotations
+from unittest import mock
+
import pytest
+from fastapi.testclient import TestClient
from sqlalchemy import select
+from airflow.api_fastapi.auth.managers.simple.user import SimpleAuthManagerUser
from airflow.models.dagbag import DagPriorityParsingRequest, DBDagBag
+from airflow.models.errors import ParseImportError
from tests_common.test_utils.api_fastapi import _check_last_log
-from tests_common.test_utils.db import clear_db_dag_parsing_requests,
clear_db_logs, parse_and_sync_to_db
+from tests_common.test_utils.db import (
+ clear_db_dag_parsing_requests,
+ clear_db_import_errors,
+ clear_db_logs,
+ parse_and_sync_to_db,
+)
from tests_common.test_utils.paths import AIRFLOW_CORE_SOURCES_PATH
pytestmark = pytest.mark.db_test
@@ -33,10 +43,26 @@ NOT_READABLE_DAG_ID = "latest_only_with_trigger"
TEST_MULTIPLE_DAGS_ID = "asset_produces_1"
[email protected]
+def dag_reader_test_client(test_client):
+ """A caller who may read the Dags (and import errors) but not edit them:
viewer is below the role edits require."""
+ auth_manager = test_client.app.state.auth_manager
+ token = auth_manager._get_token_signer().generate(
+ auth_manager.serialize_user(SimpleAuthManagerUser(username="reader",
role="viewer"))
+ )
+ with mock.patch("airflow.models.revoked_token.RevokedToken.is_revoked",
return_value=False):
+ yield TestClient(
+ test_client.app,
+ headers={"Authorization": f"Bearer {token}"},
+ base_url=str(test_client.base_url),
+ )
+
+
class TestDagParsingEndpoint:
@staticmethod
def clear_db():
clear_db_dag_parsing_requests()
+ clear_db_import_errors()
@pytest.fixture(autouse=True)
def setup(self, session) -> None:
@@ -94,3 +120,58 @@ class TestDagParsingEndpoint:
parsing_requests =
session.scalars(select(DagPriorityParsingRequest)).all()
assert parsing_requests == []
+
+ def test_reparse_import_error_file(self, url_safe_serializer, session,
test_client):
+ # A file with an import error has no registered Dags, but reparse must
still be allowed
+ # so the user can retry after fixing the file.
+ session.add(ParseImportError(bundle_name="some_bundle",
filename="dags/broken.py", stacktrace="boom"))
+ session.commit()
+ token = url_safe_serializer.dumps(
+ {"bundle_name": "some_bundle", "relative_fileloc":
"dags/broken.py"}
+ )
+
+ response = test_client.put(f"/parseDagFile/{token}",
headers={"Accept": "application/json"})
+
+ assert response.status_code == 201
+ parsing_requests =
session.scalars(select(DagPriorityParsingRequest)).all()
+ assert len(parsing_requests) == 1
+ assert parsing_requests[0].bundle_name == "some_bundle"
+ assert parsing_requests[0].relative_fileloc == "dags/broken.py"
+
+ def test_reparse_import_error_file_forbidden(
+ self, url_safe_serializer, session, unauthorized_test_client
+ ):
+ session.add(ParseImportError(bundle_name="some_bundle",
filename="dags/broken.py", stacktrace="boom"))
+ session.commit()
+ token = url_safe_serializer.dumps(
+ {"bundle_name": "some_bundle", "relative_fileloc":
"dags/broken.py"}
+ )
+
+ response = unauthorized_test_client.put(
+ f"/parseDagFile/{token}", headers={"Accept": "application/json"}
+ )
+
+ assert response.status_code == 403
+ assert session.scalars(select(DagPriorityParsingRequest)).all() == []
+
+ def
test_reparse_import_error_file_forbidden_for_basic_import_errors_viewer(
+ self, url_safe_serializer, session, dag_reader_test_client
+ ):
+ # Reparsing a file with no registered Dag requires the dedicated
REPARSE_ALL permission
+ # (admin-by-default), so a caller who can view the import-errors list
(basic IMPORT_ERRORS)
+ # but lacks REPARSE_ALL must not be able to reparse it.
+ session.add(ParseImportError(bundle_name="some_bundle",
filename="dags/broken.py", stacktrace="boom"))
+ session.commit()
+ token = url_safe_serializer.dumps(
+ {"bundle_name": "some_bundle", "relative_fileloc":
"dags/broken.py"}
+ )
+
+ response = dag_reader_test_client.put(
+ f"/parseDagFile/{token}", headers={"Accept": "application/json"}
+ )
+
+ assert response.status_code == 403
+ assert (
+ response.json()["detail"] == "You do not have permission to
reparse files with no registered Dag"
+ )
+ assert session.scalars(select(DagPriorityParsingRequest)).all() == []
diff --git
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_import_error.py
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_import_error.py
index e703d98b2ad..1bde7d2e1a3 100644
---
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_import_error.py
+++
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_import_error.py
@@ -208,6 +208,7 @@ class TestGetImportError:
test_client,
permitted_dag_model_all,
import_errors,
+ url_safe_serializer,
):
import_error: ParseImportError | None = (
import_errors[prepared_import_error_idx] if
prepared_import_error_idx is not None else None
@@ -219,7 +220,14 @@ class TestGetImportError:
if expected_status_code != 200:
return
- expected_body.update({"import_error_id": import_error_id})
+ expected_body.update(
+ {
+ "import_error_id": import_error_id,
+ "file_token": url_safe_serializer.dumps(
+ {"bundle_name": import_error.bundle_name,
"relative_fileloc": import_error.filename}
+ ),
+ }
+ )
assert response.json() == expected_body
def test_should_raises_401_unauthenticated(self,
unauthenticated_test_client, import_errors):
@@ -255,6 +263,7 @@ class TestGetImportError:
permitted_dag_model_all,
not_permitted_dag_model,
import_errors,
+ url_safe_serializer,
):
import_error_id = import_errors[0].id
set_mock_auth_manager__get_authorized_dag_ids(mock_get_auth_manager,
permitted_dag_model_all)
@@ -268,6 +277,9 @@ class TestGetImportError:
"filename": FILENAME1,
"stack_trace": "REDACTED - you do not have read permission on all
Dags in the file",
"bundle_name": BUNDLE_NAME,
+ "file_token": url_safe_serializer.dumps(
+ {"bundle_name": BUNDLE_NAME, "relative_fileloc": FILENAME1}
+ ),
}
@pytest.mark.parametrize(
@@ -285,6 +297,7 @@ class TestGetImportError:
import_errors,
can_view_all_import_errors,
expected_status_code,
+ url_safe_serializer,
):
"""A file with no registered Dag has no per-Dag key to authorize on, so
visibility is gated on the dedicated ``IMPORT_ERRORS_ALL`` view:
callers
@@ -305,6 +318,9 @@ class TestGetImportError:
"filename": FILENAME1,
"stack_trace": STACKTRACE1,
"bundle_name": BUNDLE_NAME,
+ "file_token": url_safe_serializer.dumps(
+ {"bundle_name": BUNDLE_NAME, "relative_fileloc": FILENAME1}
+ ),
}
# The unregistered-file view is what gates access, scoped to the file's
# team (None for the un-teamed "testing" bundle).
@@ -746,6 +762,7 @@ class TestGetImportErrors:
expected_stack_trace,
permitted_dag_model_all,
import_errors,
+ url_safe_serializer,
):
dag_id1 = "dag_id1"
mock_get_dag_id_to_team_name_mapping.return_value = {dag_id1: team}
@@ -770,6 +787,9 @@ class TestGetImportErrors:
"filename": FILENAME1,
"stack_trace": expected_stack_trace,
"bundle_name": BUNDLE_NAME,
+ "file_token": url_safe_serializer.dumps(
+ {"bundle_name": BUNDLE_NAME, "relative_fileloc":
FILENAME1}
+ ),
}
],
}
diff --git a/airflow-core/tests/unit/api_fastapi/core_api/test_security.py
b/airflow-core/tests/unit/api_fastapi/core_api/test_security.py
index df7a8abfb1a..c5ca94e5b26 100644
--- a/airflow-core/tests/unit/api_fastapi/core_api/test_security.py
+++ b/airflow-core/tests/unit/api_fastapi/core_api/test_security.py
@@ -21,6 +21,7 @@ from unittest.mock import AsyncMock, Mock, patch
import pytest
from fastapi import HTTPException, Request
+from itsdangerous import URLSafeSerializer
from jwt import ExpiredSignatureError, InvalidTokenError
from sqlalchemy import select
from sqlalchemy.orm import Session
@@ -55,6 +56,7 @@ from airflow.api_fastapi.core_api.security import (
requires_access_connection,
requires_access_connection_bulk,
requires_access_dag,
+ requires_access_dag_from_file_token,
requires_access_event_log,
requires_access_pool,
requires_access_pool_bulk,
@@ -365,6 +367,66 @@ class TestFastApiSecurity:
user=user,
)
+ @patch.object(DagBundleModel, "get_team_name")
+ @patch("airflow.api_fastapi.core_api.security.get_auth_manager")
+ def
test_requires_access_dag_from_file_token_no_dag_reparse_scoped_to_files_team(
+ self, mock_get_auth_manager, mock_get_team_name
+ ):
+ # Reparsing a file with no registered Dag is authorized on the
dedicated REPARSE_ALL
+ # permission, scoped to the file's own team (from its bundle), so a
caller cannot reparse
+ # another team's file.
+ auth_manager = Mock()
+ auth_manager.authorize_view.return_value = True
+ mock_get_auth_manager.return_value = auth_manager
+ mock_get_team_name.return_value = "team_b"
+
+ secret_key = "secret"
+ request = Mock()
+ request.app.state.secret_key = secret_key
+ token = URLSafeSerializer(secret_key).dumps(
+ {"bundle_name": "team_b_bundle", "relative_fileloc":
"dags/broken.py"}
+ )
+ session = Mock()
+ session.scalars.return_value = [] # no registered Dags
+ session.scalar.return_value = 1 # an import error row exists
+ user = Mock()
+
+ requires_access_dag_from_file_token("PUT")(token, request, user,
session)
+
+ mock_get_team_name.assert_called_once_with("team_b_bundle",
session=session)
+ auth_manager.authorize_view.assert_called_once_with(
+ access_view=AccessView.REPARSE_ALL, user=user, team_name="team_b"
+ )
+
+ @patch.object(DagBundleModel, "get_team_name")
+ @patch("airflow.api_fastapi.core_api.security.get_auth_manager")
+ def test_requires_access_dag_from_file_token_no_dag_reparse_forbidden(
+ self, mock_get_auth_manager, mock_get_team_name
+ ):
+ # Without REPARSE_ALL on the file's team, reparse of a no-Dag file is
denied with a
+ # message that names the no-registered-Dag case.
+ auth_manager = Mock()
+ auth_manager.authorize_view.return_value = False
+ mock_get_auth_manager.return_value = auth_manager
+ mock_get_team_name.return_value = None
+
+ secret_key = "secret"
+ request = Mock()
+ request.app.state.secret_key = secret_key
+ token = URLSafeSerializer(secret_key).dumps(
+ {"bundle_name": "some_bundle", "relative_fileloc":
"dags/broken.py"}
+ )
+ session = Mock()
+ session.scalars.return_value = []
+ session.scalar.return_value = 1
+ user = Mock()
+
+ with pytest.raises(HTTPException) as exc_info:
+ requires_access_dag_from_file_token("PUT")(token, request, user,
session)
+
+ assert exc_info.value.status_code == 403
+ assert exc_info.value.detail == "You do not have permission to reparse
files with no registered Dag"
+
@pytest.mark.db_test
@pytest.mark.asyncio
@patch.object(DagModel, "get_team_name")
diff --git a/airflow-ctl/src/airflowctl/api/datamodels/generated.py
b/airflow-ctl/src/airflowctl/api/datamodels/generated.py
index 23329487c61..d93dd23dec3 100644
--- a/airflow-ctl/src/airflowctl/api/datamodels/generated.py
+++ b/airflow-ctl/src/airflowctl/api/datamodels/generated.py
@@ -920,6 +920,13 @@ class ImportErrorResponse(BaseModel):
filename: Annotated[str, Field(title="Filename")]
bundle_name: Annotated[str | None, Field(title="Bundle Name")]
stack_trace: Annotated[str, Field(title="Stack Trace")]
+ file_token: Annotated[
+ str,
+ Field(
+ description="Return a signed token identifying the file, used to
request its reparse.",
+ title="File Token",
+ ),
+ ]
class JobResponse(BaseModel):
diff --git a/airflow-ctl/tests/airflow_ctl/api/test_operations.py
b/airflow-ctl/tests/airflow_ctl/api/test_operations.py
index f6a9802af26..f1eaf2e89ef 100644
--- a/airflow-ctl/tests/airflow_ctl/api/test_operations.py
+++ b/airflow-ctl/tests/airflow_ctl/api/test_operations.py
@@ -1153,6 +1153,7 @@ class TestDagOperations:
filename="filename",
bundle_name="bundle_name",
stack_trace="stack_trace",
+ file_token="file_token",
)
import_error_collection_response = ImportErrorCollectionResponse(