ashb commented on code in PR #74370:
URL: https://github.com/apache/airflow/pull/74370#discussion_r4207177464


##########
airflow-core/src/airflow/api_fastapi/execution_api/datamodels/token.py:
##########
@@ -17,31 +17,60 @@
 
 from __future__ import annotations
 
-from typing import Literal
+from typing import Annotated, Literal
 from uuid import UUID
 
-from pydantic import ConfigDict
+from pydantic import ConfigDict, Field, model_validator
 
 from airflow.api_fastapi.core_api.base import BaseModel
+from airflow.typing_compat import Self
 
-TokenScope = Literal["execution", "workload", "callback"]
+TokenScope = Literal[
+    "execution", "workload", "callback", "dag_processor_session", 
"dag_processor", "dag_parse"

Review Comment:
   What's the difference between dag_processor_session and dag_processor?



##########
airflow-core/src/airflow/api_fastapi/execution_api/dag_processor_tokens.py:
##########
@@ -0,0 +1,65 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""
+Issue Dag processor session tokens.
+
+Only trusted provisioning runs this: it needs the Execution API signing key, 
which a Dag processor
+must never hold, because the processor runs the code it parses.
+"""
+
+from __future__ import annotations
+
+import contextlib
+import os
+import tempfile
+from typing import TYPE_CHECKING
+
+from airflow.api_fastapi.execution_api.app import _jwt_generator
+
+if TYPE_CHECKING:
+    from collections.abc import Collection
+    from pathlib import Path
+    from uuid import UUID
+
+
+def generate_dag_processor_session_token(
+    *, session_id: UUID, bundle_names: Collection[str], valid_for: float
+) -> str:
+    """Return a ``dag_processor_session`` token, which a Dag processor 
exchanges for a Job token."""
+    return _jwt_generator().generate(
+        extras={
+            "sub": str(session_id),
+            "scope": "dag_processor_session",
+            "dag_bundles": sorted(bundle_names),
+        },
+        valid_for=valid_for,
+    )
+
+
+def write_token_file(path: Path, token: str) -> None:

Review Comment:
   This feels slightly out of place in here.



##########
airflow-core/src/airflow/api_fastapi/execution_api/datamodels/token.py:
##########
@@ -17,31 +17,60 @@
 
 from __future__ import annotations
 
-from typing import Literal
+from typing import Annotated, Literal
 from uuid import UUID
 
-from pydantic import ConfigDict
+from pydantic import ConfigDict, Field, model_validator
 
 from airflow.api_fastapi.core_api.base import BaseModel
+from airflow.typing_compat import Self
 
-TokenScope = Literal["execution", "workload", "callback"]
+TokenScope = Literal[
+    "execution", "workload", "callback", "dag_processor_session", 
"dag_processor", "dag_parse"
+]
 
 
-class TIClaims(BaseModel):
+class ExecutionClaims(BaseModel):
     """
-    Validated JWT claims for a task identity token.
+    Validated JWT claims for an Execution API principal.
 
-    Only fields used by the Execution API (sub, scope) are explicitly typed.
-    JWTValidator already validates exp/iat/nbf/aud/etc. Extra claims are 
allowed.
+    JWTValidator validates exp/iat/nbf/aud before these claims are constructed.
+    Extra claims are allowed for compatibility with existing task tokens.
     """
 
     model_config = ConfigDict(extra="allow")
 
     scope: TokenScope = "execution"
+    exp: float | None = None
+    dag_bundles: frozenset[Annotated[str, Field(min_length=1)]] | None = None
+    """Dag bundles a Dag processor token may act for."""
+    job_id: int | None = None
+    """Job a ``dag_processor`` token was issued for when that Job 
registered."""
+    session_id: UUID | None = None
+    """Processor session that owns a parsing attempt, whose own identity is 
the token subject."""
+    relative_fileloc: str | None = Field(default=None, min_length=1, 
max_length=2000)
+    """Bundle-relative file being parsed; an archive is one file under the 
current processor model."""
 
+    @model_validator(mode="after")
+    def validate_dag_processor_claims(self) -> Self:
+        if self.scope in ("dag_processor_session", "dag_processor", 
"dag_parse") and not self.dag_bundles:
+            raise ValueError(f"A {self.scope} token must grant at least one 
Dag bundle")
+        if self.scope in ("dag_processor", "dag_parse") and self.job_id is 
None:
+            raise ValueError(f"A {self.scope} token must name the Job it was 
issued for")
+        if self.scope == "dag_parse" and (
+            len(self.dag_bundles or ()) != 1 or self.session_id is None or 
self.relative_fileloc is None
+        ):
+            raise ValueError("A dag_parse token must name one bundle, its 
file, and its processor session")
+        return self
 
-class TIToken(BaseModel):
-    """Task Identity Token."""
+
+class ExecutionToken(BaseModel):
+    """Authenticated task, callback, processor session, or parsing-attempt 
identity."""
 
     id: UUID
-    claims: TIClaims
+    claims: ExecutionClaims
+
+
+# Preserve imports for task-specific consumers while shared endpoints use the 
general principal.
+TIClaims = ExecutionClaims

Review Comment:
   I'm not sure about this -- I think it might be nicer to distinqiush the 
claims we have for TI ExecutorTask work loads form dag parsing ones.



##########
airflow-core/src/airflow/api_fastapi/execution_api/routes/jobs.py:
##########
@@ -0,0 +1,274 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+
+from __future__ import annotations
+
+from typing import TYPE_CHECKING
+
+import svcs
+from cadwyn import VersionedAPIRouter
+from fastapi import HTTPException, Security, status
+from sqlalchemy import insert, select, update
+from sqlalchemy.exc import IntegrityError
+
+from airflow._shared.timezones import timezone
+from airflow.api_fastapi.auth.tokens import JWTGenerator
+from airflow.api_fastapi.common.db.common import SessionDep
+from airflow.api_fastapi.core_api.openapi.exceptions import 
create_openapi_http_exception_doc
+from airflow.api_fastapi.execution_api.datamodels.job import (
+    DagParseTokenBody,
+    DagParseTokenResponse,
+    JobCompleteBody,
+    JobHeartbeatResponse,
+    JobRegisterBody,
+    JobRegisterResponse,
+)
+from airflow.api_fastapi.execution_api.datamodels.token import ExecutionToken
+from airflow.api_fastapi.execution_api.deps import DepContainer
+from airflow.api_fastapi.execution_api.security import (
+    JOB_UNCHECKED_SCOPE,
+    CurrentExecutionToken,
+    ExecutionAPIRoute,
+    require_auth,
+)
+from airflow.configuration import conf
+from airflow.jobs.dag_processor_job_runner import DagProcessorJobRunner
+from airflow.jobs.job import Job, JobState
+from airflow.models.dagbundle import DagBundleModel
+from airflow.models.team import JobTeam
+
+if TYPE_CHECKING:
+    from uuid import UUID
+
+    from sqlalchemy.orm import Session
+
+router = VersionedAPIRouter(route_class=ExecutionAPIRoute)
+
+_JOB_NOT_FOUND = create_openapi_http_exception_doc(
+    [(status.HTTP_404_NOT_FOUND, "Job not found for this token")]
+)
+
+
+def _issue_job_token(

Review Comment:
   This appears to be the second version of this sort of helper. The other is 
in airflow-core/src/airflow/api_fastapi/execution_api/dag_processor_tokens.py



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to