Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion airflow/api_fastapi/core_api/datamodels/job.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -31,7 +31,7 @@ class JobResponse(BaseModel):
start_date: datetime | None
end_date: datetime | None
latest_heartbeat: datetime | None
executor_class: datetime | None
executor_class: str | None
hostname: str | None
unixname: str | None

Expand Down
1 change: 0 additions & 1 deletion airflow/api_fastapi/core_api/openapi/v1-generated.yaml
Original file line numberDiff line numberDiff line change
Expand Up@@ -10051,7 +10051,6 @@ components:
executor_class:
anyOf:
- type: string
format: date-time
- type: 'null'
title: Executor Class
hostname:
Expand Down
16 changes: 16 additions & 0 deletions airflow/cli/api/__init__.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,16 @@
# 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.
264 changes: 264 additions & 0 deletions airflow/cli/api/client.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,264 @@
# 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

import contextlib
import json
import os
import sys
from functools import wraps
from typing import TYPE_CHECKING, Any, Callable, TypeVar, cast

import httpx
import keyring
import structlog
from platformdirs import user_config_path
from uuid6 import uuid7

from airflow.cli.api.operations import (
AssetsOperations,
BackfillsOperations,
ConfigOperations,
ConnectionsOperations,
DagOperations,
DagRunOperations,
JobsOperations,
PoolsOperations,
ProvidersOperations,
ServerResponseError,
VariablesOperations,
VersionOperations,
)
from airflow.exceptions import AirflowNotFoundException
from airflow.typing_compat import ParamSpec
from airflow.version import version

if TYPE_CHECKING:
# # methodtools doesn't have typestubs, so give a stub
def lru_cache(maxsize: int | None = 128):
def wrapper(f):
return f

return wrapper
else:
from methodtools import lru_cache

log = structlog.get_logger(logger_name=__name__)

__all__ = [
"Client",
"Credentials",
]

PS = ParamSpec("PS")
RT = TypeVar("RT")


def add_correlation_id(request: httpx.Request):
request.headers["correlation-id"] = str(uuid7())


def get_json_error(response: httpx.Response):
"""Raise a ServerResponseError if we can extract error info from the error."""
err = ServerResponseError.from_response(response)
if err:
log.warning("Server error ", extra=dict(err.response.json()))
raise err


def raise_on_4xx_5xx(response: httpx.Response):
return get_json_error(response) or response.raise_for_status()


# Credentials for the API
class Credentials:
"""Credentials for the API."""

api_url: str | None
api_token: str | None
api_environment: str

def __init__(
self,
api_url: str | None = None,
api_token: str | None = None,
api_environment: str = "production",
Comment thread
jedcunningham marked this conversation as resolved.
):
self.api_url = api_url
self.api_token = api_token
self.api_environment = os.getenv("AIRFLOW_CLI_ENVIRONMENT") or api_environment

@property
def input_cli_config_file(self) -> str:
"""Generate path for the CLI config file."""
return f"{self.api_environment}.json"

def save(self):
"""Save the credentials to keyring and URL to disk as a file."""
default_config_dir = user_config_path("airflow", "Apache Software Foundation")
if not os.path.exists(default_config_dir):
os.makedirs(default_config_dir)
with open(os.path.join(default_config_dir, self.input_cli_config_file), "w") as f:
json.dump({"api_url": self.api_url}, f)
keyring.set_password("airflow-cli", f"api_token-{self.api_environment}", self.api_token)

def load(self) -> Credentials:
"""Load the credentials from keyring and URL from disk file."""
default_config_dir = user_config_path("airflow", "Apache Software Foundation")
if os.path.exists(default_config_dir):
with open(os.path.join(default_config_dir, self.input_cli_config_file)) as f:
credentials = json.load(f)
self.api_url = credentials["api_url"]
self.api_token = keyring.get_password("airflow-cli", f"api_token-{self.api_environment}")
return self
else:
raise AirflowNotFoundException(f"No credentials found in {default_config_dir}")


class BearerAuth(httpx.Auth):
def __init__(self, token: str):
self.token: str = token

def auth_flow(self, request: httpx.Request):
if self.token:
request.headers["Authorization"] = "Bearer " + self.token
yield request


class Client(httpx.Client):
"""Client for the Airflow REST API."""

def __init__(self, *, base_url: str, token: str, **kwargs: Any):
auth = BearerAuth(token)
print(f"token: {token}")
kwargs["base_url"] = f"{base_url}/public"
pyver = f"{'.'.join(map(str, sys.version_info[:3]))}"
super().__init__(
auth=auth,
headers={"user-agent": f"apache-airflow-cli/{version} (Python/{pyver})"},
event_hooks={"response": [raise_on_4xx_5xx], "request": [add_correlation_id]},
**kwargs,
)

@lru_cache() # type: ignore[misc]
@property
def assets(self):
"""Operations related to assets."""
return AssetsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def backfills(self):
"""Operations related to backfills."""
return BackfillsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def configs(self):
"""Operations related to configs."""
return ConfigOperations(self)

@lru_cache() # type: ignore[misc]
@property
def connections(self):
"""Operations related to connections."""
return ConnectionsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def dags(self):
"""Operations related to DAGs."""
return DagOperations(self)

@lru_cache() # type: ignore[misc]
@property
def dag_runs(self):
"""Operations related to DAG runs."""
return DagRunOperations(self)

@lru_cache() # type: ignore[misc]
@property
def jobs(self):
"""Operations related to jobs."""
return JobsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def pools(self):
"""Operations related to pools."""
return PoolsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def providers(self):
"""Operations related to providers."""
return ProvidersOperations(self)

@lru_cache() # type: ignore[misc]
@property
def variables(self):
"""Operations related to variables."""
return VariablesOperations(self)

@lru_cache() # type: ignore[misc]
@property
def version(self):
"""Get the version of the server."""
return VersionOperations(self)


# API Client Decorator for CLI Actions
@contextlib.contextmanager
def get_client():
"""Get CLI API client."""
cli_api_client = None
try:
credentials = Credentials().load()
limits = httpx.Limits(max_keepalive_connections=1, max_connections=1)
cli_api_client = Client(base_url=credentials.api_url, limits=limits, token=credentials.api_token)
yield cli_api_client
except AirflowNotFoundException as e:
raise e
Comment thread
jedcunningham marked this conversation as resolved.
finally:
if cli_api_client:
cli_api_client.close()


def provide_api_client(func: Callable[PS, RT]) -> Callable[PS, RT]:
"""
Provide a CLI API Client to the decorated function.

CLI API Client shouldn't be passed to the function when this wrapper is used
if the purpose is not mocking or testing.
If you want to reuse a CLI API Client or run the function as part of
API call, you pass it to the function, if not this wrapper
will create one and close it for you.
"""

@wraps(func)
def wrapper(*args, **kwargs) -> RT:
if "cli_api_client" not in kwargs:
with get_client() as cli_api_client:
return func(*args, cli_api_client=cli_api_client, **kwargs)
# The CLI API Client should be only passed for Mocking and Testing
return func(*args, **kwargs)

return wrapper


NEW_CLI_API_CLIENT: Client = cast(Client, None)
16 changes: 16 additions & 0 deletions airflow/cli/api/datamodels/__init__.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,16 @@
# 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.
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Add copy buttons to all
 blocks
(function() {
function addCopyButtons() {
document.querySelectorAll('pre code').forEach(function(codeBlock) {
if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;
codeBlock.parentElement.setAttribute('data-copy-added', 'true');
var btn = document.createElement('button');
btn.textContent = 'Copy';
btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';
btn.onmouseover = function() { this.style.opacity = '1'; };
btn.onmouseout = function() { this.style.opacity = '0.7'; };
btn.onclick = function() {
navigator.clipboard.writeText(codeBlock.textContent).then(function() {
btn.textContent = 'Copied!';
setTimeout(function() { btn.textContent = 'Copy'; }, 1500);
});
};
codeBlock.parentElement.style.position = 'relative';
codeBlock.parentElement.appendChild(btn);
});
}
addCopyButtons();
// Re-run on dynamic content
var observer = new MutationObserver(addCopyButtons);
observer.observe(document.body, { childList: true, subtree: true });
})();
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion airflow/api_fastapi/core_api/datamodels/job.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -31,7 +31,7 @@ class JobResponse(BaseModel):
start_date: datetime | None
end_date: datetime | None
latest_heartbeat: datetime | None
executor_class: datetime | None
executor_class: str | None
hostname: str | None
unixname: str | None

Expand Down
1 change: 0 additions & 1 deletion airflow/api_fastapi/core_api/openapi/v1-generated.yaml
Original file line numberDiff line numberDiff line change
Expand Up@@ -10051,7 +10051,6 @@ components:
executor_class:
anyOf:
- type: string
format: date-time
- type: 'null'
title: Executor Class
hostname:
Expand Down
16 changes: 16 additions & 0 deletions airflow/cli/api/__init__.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,16 @@
# 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.
264 changes: 264 additions & 0 deletions airflow/cli/api/client.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,264 @@
# 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

import contextlib
import json
import os
import sys
from functools import wraps
from typing import TYPE_CHECKING, Any, Callable, TypeVar, cast

import httpx
import keyring
import structlog
from platformdirs import user_config_path
from uuid6 import uuid7

from airflow.cli.api.operations import (
AssetsOperations,
BackfillsOperations,
ConfigOperations,
ConnectionsOperations,
DagOperations,
DagRunOperations,
JobsOperations,
PoolsOperations,
ProvidersOperations,
ServerResponseError,
VariablesOperations,
VersionOperations,
)
from airflow.exceptions import AirflowNotFoundException
from airflow.typing_compat import ParamSpec
from airflow.version import version

if TYPE_CHECKING:
# # methodtools doesn't have typestubs, so give a stub
def lru_cache(maxsize: int | None = 128):
def wrapper(f):
return f

return wrapper
else:
from methodtools import lru_cache

log = structlog.get_logger(logger_name=__name__)

__all__ = [
"Client",
"Credentials",
]

PS = ParamSpec("PS")
RT = TypeVar("RT")


def add_correlation_id(request: httpx.Request):
request.headers["correlation-id"] = str(uuid7())


def get_json_error(response: httpx.Response):
"""Raise a ServerResponseError if we can extract error info from the error."""
err = ServerResponseError.from_response(response)
if err:
log.warning("Server error ", extra=dict(err.response.json()))
raise err


def raise_on_4xx_5xx(response: httpx.Response):
return get_json_error(response) or response.raise_for_status()


# Credentials for the API
class Credentials:
"""Credentials for the API."""

api_url: str | None
api_token: str | None
api_environment: str

def __init__(
self,
api_url: str | None = None,
api_token: str | None = None,
api_environment: str = "production",
Comment thread
jedcunningham marked this conversation as resolved.
):
self.api_url = api_url
self.api_token = api_token
self.api_environment = os.getenv("AIRFLOW_CLI_ENVIRONMENT") or api_environment

@property
def input_cli_config_file(self) -> str:
"""Generate path for the CLI config file."""
return f"{self.api_environment}.json"

def save(self):
"""Save the credentials to keyring and URL to disk as a file."""
default_config_dir = user_config_path("airflow", "Apache Software Foundation")
if not os.path.exists(default_config_dir):
os.makedirs(default_config_dir)
with open(os.path.join(default_config_dir, self.input_cli_config_file), "w") as f:
json.dump({"api_url": self.api_url}, f)
keyring.set_password("airflow-cli", f"api_token-{self.api_environment}", self.api_token)

def load(self) -> Credentials:
"""Load the credentials from keyring and URL from disk file."""
default_config_dir = user_config_path("airflow", "Apache Software Foundation")
if os.path.exists(default_config_dir):
with open(os.path.join(default_config_dir, self.input_cli_config_file)) as f:
credentials = json.load(f)
self.api_url = credentials["api_url"]
self.api_token = keyring.get_password("airflow-cli", f"api_token-{self.api_environment}")
return self
else:
raise AirflowNotFoundException(f"No credentials found in {default_config_dir}")


class BearerAuth(httpx.Auth):
def __init__(self, token: str):
self.token: str = token

def auth_flow(self, request: httpx.Request):
if self.token:
request.headers["Authorization"] = "Bearer " + self.token
yield request


class Client(httpx.Client):
"""Client for the Airflow REST API."""

def __init__(self, *, base_url: str, token: str, **kwargs: Any):
auth = BearerAuth(token)
print(f"token: {token}")
kwargs["base_url"] = f"{base_url}/public"
pyver = f"{'.'.join(map(str, sys.version_info[:3]))}"
super().__init__(
auth=auth,
headers={"user-agent": f"apache-airflow-cli/{version} (Python/{pyver})"},
event_hooks={"response": [raise_on_4xx_5xx], "request": [add_correlation_id]},
**kwargs,
)

@lru_cache() # type: ignore[misc]
@property
def assets(self):
"""Operations related to assets."""
return AssetsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def backfills(self):
"""Operations related to backfills."""
return BackfillsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def configs(self):
"""Operations related to configs."""
return ConfigOperations(self)

@lru_cache() # type: ignore[misc]
@property
def connections(self):
"""Operations related to connections."""
return ConnectionsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def dags(self):
"""Operations related to DAGs."""
return DagOperations(self)

@lru_cache() # type: ignore[misc]
@property
def dag_runs(self):
"""Operations related to DAG runs."""
return DagRunOperations(self)

@lru_cache() # type: ignore[misc]
@property
def jobs(self):
"""Operations related to jobs."""
return JobsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def pools(self):
"""Operations related to pools."""
return PoolsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def providers(self):
"""Operations related to providers."""
return ProvidersOperations(self)

@lru_cache() # type: ignore[misc]
@property
def variables(self):
"""Operations related to variables."""
return VariablesOperations(self)

@lru_cache() # type: ignore[misc]
@property
def version(self):
"""Get the version of the server."""
return VersionOperations(self)


# API Client Decorator for CLI Actions
@contextlib.contextmanager
def get_client():
"""Get CLI API client."""
cli_api_client = None
try:
credentials = Credentials().load()
limits = httpx.Limits(max_keepalive_connections=1, max_connections=1)
cli_api_client = Client(base_url=credentials.api_url, limits=limits, token=credentials.api_token)
yield cli_api_client
except AirflowNotFoundException as e:
raise e
Comment thread
jedcunningham marked this conversation as resolved.
finally:
if cli_api_client:
cli_api_client.close()


def provide_api_client(func: Callable[PS, RT]) -> Callable[PS, RT]:
"""
Provide a CLI API Client to the decorated function.

CLI API Client shouldn't be passed to the function when this wrapper is used
if the purpose is not mocking or testing.
If you want to reuse a CLI API Client or run the function as part of
API call, you pass it to the function, if not this wrapper
will create one and close it for you.
"""

@wraps(func)
def wrapper(*args, **kwargs) -> RT:
if "cli_api_client" not in kwargs:
with get_client() as cli_api_client:
return func(*args, cli_api_client=cli_api_client, **kwargs)
# The CLI API Client should be only passed for Mocking and Testing
return func(*args, **kwargs)

return wrapper


NEW_CLI_API_CLIENT: Client = cast(Client, None)
16 changes: 16 additions & 0 deletions airflow/cli/api/datamodels/__init__.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,16 @@
# 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.
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Force GitHub README to respect dark mode (function() { var style = document.createElement('style'); style.textContent = ' .markdown-body { color-scheme: dark light; } .markdown-body pre { background: #161b22 !important; } .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; } .markdown-body table th, .markdown-body table td { border-color: #30363d !important; } .markdown-body img { background: #0d1117; } .markdown-body blockquote { border-left-color: #8b949e; } .markdown-body hr { border-color: #30363d; } '; document.head.appendChild(style); })(); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion airflow/api_fastapi/core_api/datamodels/job.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -31,7 +31,7 @@ class JobResponse(BaseModel):
start_date: datetime | None
end_date: datetime | None
latest_heartbeat: datetime | None
executor_class: datetime | None
executor_class: str | None
hostname: str | None
unixname: str | None

Expand Down
1 change: 0 additions & 1 deletion airflow/api_fastapi/core_api/openapi/v1-generated.yaml
Original file line numberDiff line numberDiff line change
Expand Up@@ -10051,7 +10051,6 @@ components:
executor_class:
anyOf:
- type: string
format: date-time
- type: 'null'
title: Executor Class
hostname:
Expand Down
16 changes: 16 additions & 0 deletions airflow/cli/api/__init__.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,16 @@
# 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.
264 changes: 264 additions & 0 deletions airflow/cli/api/client.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,264 @@
# 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

import contextlib
import json
import os
import sys
from functools import wraps
from typing import TYPE_CHECKING, Any, Callable, TypeVar, cast

import httpx
import keyring
import structlog
from platformdirs import user_config_path
from uuid6 import uuid7

from airflow.cli.api.operations import (
AssetsOperations,
BackfillsOperations,
ConfigOperations,
ConnectionsOperations,
DagOperations,
DagRunOperations,
JobsOperations,
PoolsOperations,
ProvidersOperations,
ServerResponseError,
VariablesOperations,
VersionOperations,
)
from airflow.exceptions import AirflowNotFoundException
from airflow.typing_compat import ParamSpec
from airflow.version import version

if TYPE_CHECKING:
# # methodtools doesn't have typestubs, so give a stub
def lru_cache(maxsize: int | None = 128):
def wrapper(f):
return f

return wrapper
else:
from methodtools import lru_cache

log = structlog.get_logger(logger_name=__name__)

__all__ = [
"Client",
"Credentials",
]

PS = ParamSpec("PS")
RT = TypeVar("RT")


def add_correlation_id(request: httpx.Request):
request.headers["correlation-id"] = str(uuid7())


def get_json_error(response: httpx.Response):
"""Raise a ServerResponseError if we can extract error info from the error."""
err = ServerResponseError.from_response(response)
if err:
log.warning("Server error ", extra=dict(err.response.json()))
raise err


def raise_on_4xx_5xx(response: httpx.Response):
return get_json_error(response) or response.raise_for_status()


# Credentials for the API
class Credentials:
"""Credentials for the API."""

api_url: str | None
api_token: str | None
api_environment: str

def __init__(
self,
api_url: str | None = None,
api_token: str | None = None,
api_environment: str = "production",
Comment thread
jedcunningham marked this conversation as resolved.
):
self.api_url = api_url
self.api_token = api_token
self.api_environment = os.getenv("AIRFLOW_CLI_ENVIRONMENT") or api_environment

@property
def input_cli_config_file(self) -> str:
"""Generate path for the CLI config file."""
return f"{self.api_environment}.json"

def save(self):
"""Save the credentials to keyring and URL to disk as a file."""
default_config_dir = user_config_path("airflow", "Apache Software Foundation")
if not os.path.exists(default_config_dir):
os.makedirs(default_config_dir)
with open(os.path.join(default_config_dir, self.input_cli_config_file), "w") as f:
json.dump({"api_url": self.api_url}, f)
keyring.set_password("airflow-cli", f"api_token-{self.api_environment}", self.api_token)

def load(self) -> Credentials:
"""Load the credentials from keyring and URL from disk file."""
default_config_dir = user_config_path("airflow", "Apache Software Foundation")
if os.path.exists(default_config_dir):
with open(os.path.join(default_config_dir, self.input_cli_config_file)) as f:
credentials = json.load(f)
self.api_url = credentials["api_url"]
self.api_token = keyring.get_password("airflow-cli", f"api_token-{self.api_environment}")
return self
else:
raise AirflowNotFoundException(f"No credentials found in {default_config_dir}")


class BearerAuth(httpx.Auth):
def __init__(self, token: str):
self.token: str = token

def auth_flow(self, request: httpx.Request):
if self.token:
request.headers["Authorization"] = "Bearer " + self.token
yield request


class Client(httpx.Client):
"""Client for the Airflow REST API."""

def __init__(self, *, base_url: str, token: str, **kwargs: Any):
auth = BearerAuth(token)
print(f"token: {token}")
kwargs["base_url"] = f"{base_url}/public"
pyver = f"{'.'.join(map(str, sys.version_info[:3]))}"
super().__init__(
auth=auth,
headers={"user-agent": f"apache-airflow-cli/{version} (Python/{pyver})"},
event_hooks={"response": [raise_on_4xx_5xx], "request": [add_correlation_id]},
**kwargs,
)

@lru_cache() # type: ignore[misc]
@property
def assets(self):
"""Operations related to assets."""
return AssetsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def backfills(self):
"""Operations related to backfills."""
return BackfillsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def configs(self):
"""Operations related to configs."""
return ConfigOperations(self)

@lru_cache() # type: ignore[misc]
@property
def connections(self):
"""Operations related to connections."""
return ConnectionsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def dags(self):
"""Operations related to DAGs."""
return DagOperations(self)

@lru_cache() # type: ignore[misc]
@property
def dag_runs(self):
"""Operations related to DAG runs."""
return DagRunOperations(self)

@lru_cache() # type: ignore[misc]
@property
def jobs(self):
"""Operations related to jobs."""
return JobsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def pools(self):
"""Operations related to pools."""
return PoolsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def providers(self):
"""Operations related to providers."""
return ProvidersOperations(self)

@lru_cache() # type: ignore[misc]
@property
def variables(self):
"""Operations related to variables."""
return VariablesOperations(self)

@lru_cache() # type: ignore[misc]
@property
def version(self):
"""Get the version of the server."""
return VersionOperations(self)


# API Client Decorator for CLI Actions
@contextlib.contextmanager
def get_client():
"""Get CLI API client."""
cli_api_client = None
try:
credentials = Credentials().load()
limits = httpx.Limits(max_keepalive_connections=1, max_connections=1)
cli_api_client = Client(base_url=credentials.api_url, limits=limits, token=credentials.api_token)
yield cli_api_client
except AirflowNotFoundException as e:
raise e
Comment thread
jedcunningham marked this conversation as resolved.
finally:
if cli_api_client:
cli_api_client.close()


def provide_api_client(func: Callable[PS, RT]) -> Callable[PS, RT]:
"""
Provide a CLI API Client to the decorated function.

CLI API Client shouldn't be passed to the function when this wrapper is used
if the purpose is not mocking or testing.
If you want to reuse a CLI API Client or run the function as part of
API call, you pass it to the function, if not this wrapper
will create one and close it for you.
"""

@wraps(func)
def wrapper(*args, **kwargs) -> RT:
if "cli_api_client" not in kwargs:
with get_client() as cli_api_client:
return func(*args, cli_api_client=cli_api_client, **kwargs)
# The CLI API Client should be only passed for Mocking and Testing
return func(*args, **kwargs)

return wrapper


NEW_CLI_API_CLIENT: Client = cast(Client, None)
16 changes: 16 additions & 0 deletions airflow/cli/api/datamodels/__init__.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,16 @@
# 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.
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Highlight search terms from Google/DuckDuckGo/Bing referrer (function() { var ref = document.referrer; var terms = []; if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) { var url = new URL(ref); var q = url.searchParams.get('q') || url.searchParams.get('p'); if (q) { terms = q.split(/\s+/).filter(function(t) { return t.length > 2; }); } } if (terms.length === 0) return; var style = document.createElement('style'); style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }'; document.head.appendChild(style); function highlight(node) { if (node.nodeType === 3) { // text node var text = node.textContent; var found = false; terms.forEach(function(term) { var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\]\\]/g, '\\') + ')', 'gi'); if (regex.test(text)) { found = true; var frag = document.createDocumentFragment(); var parts = text.split(regex); parts.forEach(function(part, i) { if (i % 2 === 0) { frag.appendChild(document.createTextNode(part)); } else { var span = document.createElement('span'); span.className = 'userscript-highlight'; span.textContent = part; frag.appendChild(span); } }); node.parentNode.replaceChild(frag, node); } }); } else if (node.nodeType === 1 && node.childNodes) { // element var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT']; if (!skipTags.includes(node.tagName)) { Array.from(node.childNodes).forEach(highlight); } } } highlight(document.body); // Re-highlight on dynamic content var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1 || node.nodeType === 3) highlight(node); }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion airflow/api_fastapi/core_api/datamodels/job.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -31,7 +31,7 @@ class JobResponse(BaseModel):
start_date: datetime | None
end_date: datetime | None
latest_heartbeat: datetime | None
executor_class: datetime | None
executor_class: str | None
hostname: str | None
unixname: str | None

Expand Down
1 change: 0 additions & 1 deletion airflow/api_fastapi/core_api/openapi/v1-generated.yaml
Original file line numberDiff line numberDiff line change
Expand Up@@ -10051,7 +10051,6 @@ components:
executor_class:
anyOf:
- type: string
format: date-time
- type: 'null'
title: Executor Class
hostname:
Expand Down
16 changes: 16 additions & 0 deletions airflow/cli/api/__init__.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,16 @@
# 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.
264 changes: 264 additions & 0 deletions airflow/cli/api/client.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,264 @@
# 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

import contextlib
import json
import os
import sys
from functools import wraps
from typing import TYPE_CHECKING, Any, Callable, TypeVar, cast

import httpx
import keyring
import structlog
from platformdirs import user_config_path
from uuid6 import uuid7

from airflow.cli.api.operations import (
AssetsOperations,
BackfillsOperations,
ConfigOperations,
ConnectionsOperations,
DagOperations,
DagRunOperations,
JobsOperations,
PoolsOperations,
ProvidersOperations,
ServerResponseError,
VariablesOperations,
VersionOperations,
)
from airflow.exceptions import AirflowNotFoundException
from airflow.typing_compat import ParamSpec
from airflow.version import version

if TYPE_CHECKING:
# # methodtools doesn't have typestubs, so give a stub
def lru_cache(maxsize: int | None = 128):
def wrapper(f):
return f

return wrapper
else:
from methodtools import lru_cache

log = structlog.get_logger(logger_name=__name__)

__all__ = [
"Client",
"Credentials",
]

PS = ParamSpec("PS")
RT = TypeVar("RT")


def add_correlation_id(request: httpx.Request):
request.headers["correlation-id"] = str(uuid7())


def get_json_error(response: httpx.Response):
"""Raise a ServerResponseError if we can extract error info from the error."""
err = ServerResponseError.from_response(response)
if err:
log.warning("Server error ", extra=dict(err.response.json()))
raise err


def raise_on_4xx_5xx(response: httpx.Response):
return get_json_error(response) or response.raise_for_status()


# Credentials for the API
class Credentials:
"""Credentials for the API."""

api_url: str | None
api_token: str | None
api_environment: str

def __init__(
self,
api_url: str | None = None,
api_token: str | None = None,
api_environment: str = "production",
Comment thread
jedcunningham marked this conversation as resolved.
):
self.api_url = api_url
self.api_token = api_token
self.api_environment = os.getenv("AIRFLOW_CLI_ENVIRONMENT") or api_environment

@property
def input_cli_config_file(self) -> str:
"""Generate path for the CLI config file."""
return f"{self.api_environment}.json"

def save(self):
"""Save the credentials to keyring and URL to disk as a file."""
default_config_dir = user_config_path("airflow", "Apache Software Foundation")
if not os.path.exists(default_config_dir):
os.makedirs(default_config_dir)
with open(os.path.join(default_config_dir, self.input_cli_config_file), "w") as f:
json.dump({"api_url": self.api_url}, f)
keyring.set_password("airflow-cli", f"api_token-{self.api_environment}", self.api_token)

def load(self) -> Credentials:
"""Load the credentials from keyring and URL from disk file."""
default_config_dir = user_config_path("airflow", "Apache Software Foundation")
if os.path.exists(default_config_dir):
with open(os.path.join(default_config_dir, self.input_cli_config_file)) as f:
credentials = json.load(f)
self.api_url = credentials["api_url"]
self.api_token = keyring.get_password("airflow-cli", f"api_token-{self.api_environment}")
return self
else:
raise AirflowNotFoundException(f"No credentials found in {default_config_dir}")


class BearerAuth(httpx.Auth):
def __init__(self, token: str):
self.token: str = token

def auth_flow(self, request: httpx.Request):
if self.token:
request.headers["Authorization"] = "Bearer " + self.token
yield request


class Client(httpx.Client):
"""Client for the Airflow REST API."""

def __init__(self, *, base_url: str, token: str, **kwargs: Any):
auth = BearerAuth(token)
print(f"token: {token}")
kwargs["base_url"] = f"{base_url}/public"
pyver = f"{'.'.join(map(str, sys.version_info[:3]))}"
super().__init__(
auth=auth,
headers={"user-agent": f"apache-airflow-cli/{version} (Python/{pyver})"},
event_hooks={"response": [raise_on_4xx_5xx], "request": [add_correlation_id]},
**kwargs,
)

@lru_cache() # type: ignore[misc]
@property
def assets(self):
"""Operations related to assets."""
return AssetsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def backfills(self):
"""Operations related to backfills."""
return BackfillsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def configs(self):
"""Operations related to configs."""
return ConfigOperations(self)

@lru_cache() # type: ignore[misc]
@property
def connections(self):
"""Operations related to connections."""
return ConnectionsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def dags(self):
"""Operations related to DAGs."""
return DagOperations(self)

@lru_cache() # type: ignore[misc]
@property
def dag_runs(self):
"""Operations related to DAG runs."""
return DagRunOperations(self)

@lru_cache() # type: ignore[misc]
@property
def jobs(self):
"""Operations related to jobs."""
return JobsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def pools(self):
"""Operations related to pools."""
return PoolsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def providers(self):
"""Operations related to providers."""
return ProvidersOperations(self)

@lru_cache() # type: ignore[misc]
@property
def variables(self):
"""Operations related to variables."""
return VariablesOperations(self)

@lru_cache() # type: ignore[misc]
@property
def version(self):
"""Get the version of the server."""
return VersionOperations(self)


# API Client Decorator for CLI Actions
@contextlib.contextmanager
def get_client():
"""Get CLI API client."""
cli_api_client = None
try:
credentials = Credentials().load()
limits = httpx.Limits(max_keepalive_connections=1, max_connections=1)
cli_api_client = Client(base_url=credentials.api_url, limits=limits, token=credentials.api_token)
yield cli_api_client
except AirflowNotFoundException as e:
raise e
Comment thread
jedcunningham marked this conversation as resolved.
finally:
if cli_api_client:
cli_api_client.close()


def provide_api_client(func: Callable[PS, RT]) -> Callable[PS, RT]:
"""
Provide a CLI API Client to the decorated function.

CLI API Client shouldn't be passed to the function when this wrapper is used
if the purpose is not mocking or testing.
If you want to reuse a CLI API Client or run the function as part of
API call, you pass it to the function, if not this wrapper
will create one and close it for you.
"""

@wraps(func)
def wrapper(*args, **kwargs) -> RT:
if "cli_api_client" not in kwargs:
with get_client() as cli_api_client:
return func(*args, cli_api_client=cli_api_client, **kwargs)
# The CLI API Client should be only passed for Mocking and Testing
return func(*args, **kwargs)

return wrapper


NEW_CLI_API_CLIENT: Client = cast(Client, None)
16 changes: 16 additions & 0 deletions airflow/cli/api/datamodels/__init__.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,16 @@
# 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.
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Strip utm_, fbclid, gclid, etc. from all links on page (function() { var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content', 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid', 'ref', 'ref_src', 'source', 'medium', 'campaign']; function cleanUrl(url) { try { var u = new URL(url, window.location.origin); var changed = false; trackingParams.forEach(function(p) { if (u.searchParams.has(p)) { u.searchParams.delete(p); changed = true; } }); return changed ? u.toString() : url; } catch (e) { return url; } } function cleanLinks() { document.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } cleanLinks(); var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1) { if (node.tagName === 'A') cleanLinks(); node.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion airflow/api_fastapi/core_api/datamodels/job.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -31,7 +31,7 @@ class JobResponse(BaseModel):
start_date: datetime | None
end_date: datetime | None
latest_heartbeat: datetime | None
executor_class: datetime | None
executor_class: str | None
hostname: str | None
unixname: str | None

Expand Down
1 change: 0 additions & 1 deletion airflow/api_fastapi/core_api/openapi/v1-generated.yaml
Original file line numberDiff line numberDiff line change
Expand Up@@ -10051,7 +10051,6 @@ components:
executor_class:
anyOf:
- type: string
format: date-time
- type: 'null'
title: Executor Class
hostname:
Expand Down
16 changes: 16 additions & 0 deletions airflow/cli/api/__init__.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,16 @@
# 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.
264 changes: 264 additions & 0 deletions airflow/cli/api/client.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,264 @@
# 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

import contextlib
import json
import os
import sys
from functools import wraps
from typing import TYPE_CHECKING, Any, Callable, TypeVar, cast

import httpx
import keyring
import structlog
from platformdirs import user_config_path
from uuid6 import uuid7

from airflow.cli.api.operations import (
AssetsOperations,
BackfillsOperations,
ConfigOperations,
ConnectionsOperations,
DagOperations,
DagRunOperations,
JobsOperations,
PoolsOperations,
ProvidersOperations,
ServerResponseError,
VariablesOperations,
VersionOperations,
)
from airflow.exceptions import AirflowNotFoundException
from airflow.typing_compat import ParamSpec
from airflow.version import version

if TYPE_CHECKING:
# # methodtools doesn't have typestubs, so give a stub
def lru_cache(maxsize: int | None = 128):
def wrapper(f):
return f

return wrapper
else:
from methodtools import lru_cache

log = structlog.get_logger(logger_name=__name__)

__all__ = [
"Client",
"Credentials",
]

PS = ParamSpec("PS")
RT = TypeVar("RT")


def add_correlation_id(request: httpx.Request):
request.headers["correlation-id"] = str(uuid7())


def get_json_error(response: httpx.Response):
"""Raise a ServerResponseError if we can extract error info from the error."""
err = ServerResponseError.from_response(response)
if err:
log.warning("Server error ", extra=dict(err.response.json()))
raise err


def raise_on_4xx_5xx(response: httpx.Response):
return get_json_error(response) or response.raise_for_status()


# Credentials for the API
class Credentials:
"""Credentials for the API."""

api_url: str | None
api_token: str | None
api_environment: str

def __init__(
self,
api_url: str | None = None,
api_token: str | None = None,
api_environment: str = "production",
Comment thread
jedcunningham marked this conversation as resolved.
):
self.api_url = api_url
self.api_token = api_token
self.api_environment = os.getenv("AIRFLOW_CLI_ENVIRONMENT") or api_environment

@property
def input_cli_config_file(self) -> str:
"""Generate path for the CLI config file."""
return f"{self.api_environment}.json"

def save(self):
"""Save the credentials to keyring and URL to disk as a file."""
default_config_dir = user_config_path("airflow", "Apache Software Foundation")
if not os.path.exists(default_config_dir):
os.makedirs(default_config_dir)
with open(os.path.join(default_config_dir, self.input_cli_config_file), "w") as f:
json.dump({"api_url": self.api_url}, f)
keyring.set_password("airflow-cli", f"api_token-{self.api_environment}", self.api_token)

def load(self) -> Credentials:
"""Load the credentials from keyring and URL from disk file."""
default_config_dir = user_config_path("airflow", "Apache Software Foundation")
if os.path.exists(default_config_dir):
with open(os.path.join(default_config_dir, self.input_cli_config_file)) as f:
credentials = json.load(f)
self.api_url = credentials["api_url"]
self.api_token = keyring.get_password("airflow-cli", f"api_token-{self.api_environment}")
return self
else:
raise AirflowNotFoundException(f"No credentials found in {default_config_dir}")


class BearerAuth(httpx.Auth):
def __init__(self, token: str):
self.token: str = token

def auth_flow(self, request: httpx.Request):
if self.token:
request.headers["Authorization"] = "Bearer " + self.token
yield request


class Client(httpx.Client):
"""Client for the Airflow REST API."""

def __init__(self, *, base_url: str, token: str, **kwargs: Any):
auth = BearerAuth(token)
print(f"token: {token}")
kwargs["base_url"] = f"{base_url}/public"
pyver = f"{'.'.join(map(str, sys.version_info[:3]))}"
super().__init__(
auth=auth,
headers={"user-agent": f"apache-airflow-cli/{version} (Python/{pyver})"},
event_hooks={"response": [raise_on_4xx_5xx], "request": [add_correlation_id]},
**kwargs,
)

@lru_cache() # type: ignore[misc]
@property
def assets(self):
"""Operations related to assets."""
return AssetsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def backfills(self):
"""Operations related to backfills."""
return BackfillsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def configs(self):
"""Operations related to configs."""
return ConfigOperations(self)

@lru_cache() # type: ignore[misc]
@property
def connections(self):
"""Operations related to connections."""
return ConnectionsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def dags(self):
"""Operations related to DAGs."""
return DagOperations(self)

@lru_cache() # type: ignore[misc]
@property
def dag_runs(self):
"""Operations related to DAG runs."""
return DagRunOperations(self)

@lru_cache() # type: ignore[misc]
@property
def jobs(self):
"""Operations related to jobs."""
return JobsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def pools(self):
"""Operations related to pools."""
return PoolsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def providers(self):
"""Operations related to providers."""
return ProvidersOperations(self)

@lru_cache() # type: ignore[misc]
@property
def variables(self):
"""Operations related to variables."""
return VariablesOperations(self)

@lru_cache() # type: ignore[misc]
@property
def version(self):
"""Get the version of the server."""
return VersionOperations(self)


# API Client Decorator for CLI Actions
@contextlib.contextmanager
def get_client():
"""Get CLI API client."""
cli_api_client = None
try:
credentials = Credentials().load()
limits = httpx.Limits(max_keepalive_connections=1, max_connections=1)
cli_api_client = Client(base_url=credentials.api_url, limits=limits, token=credentials.api_token)
yield cli_api_client
except AirflowNotFoundException as e:
raise e
Comment thread
jedcunningham marked this conversation as resolved.
finally:
if cli_api_client:
cli_api_client.close()


def provide_api_client(func: Callable[PS, RT]) -> Callable[PS, RT]:
"""
Provide a CLI API Client to the decorated function.

CLI API Client shouldn't be passed to the function when this wrapper is used
if the purpose is not mocking or testing.
If you want to reuse a CLI API Client or run the function as part of
API call, you pass it to the function, if not this wrapper
will create one and close it for you.
"""

@wraps(func)
def wrapper(*args, **kwargs) -> RT:
if "cli_api_client" not in kwargs:
with get_client() as cli_api_client:
return func(*args, cli_api_client=cli_api_client, **kwargs)
# The CLI API Client should be only passed for Mocking and Testing
return func(*args, **kwargs)

return wrapper


NEW_CLI_API_CLIENT: Client = cast(Client, None)
16 changes: 16 additions & 0 deletions airflow/cli/api/datamodels/__init__.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,16 @@
# 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.
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Auto-enable theater mode on YouTube (function() { function tryTheater() { var btn = document.querySelector('button[aria-label="Theater mode"], ytd-player #player button[title="Theater mode"]'); if (btn && !btn.classList.contains('activated')) { btn.click(); } } // Try immediately tryTheater(); // Try after navigation (SPA) var lastUrl = location.href; setInterval(function() { if (location.href !== lastUrl) { lastUrl = location.href; setTimeout(tryTheater, 500); } }, 1000); // Also try on player load var observer = new MutationObserver(tryTheater); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion airflow/api_fastapi/core_api/datamodels/job.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -31,7 +31,7 @@ class JobResponse(BaseModel):
start_date: datetime | None
end_date: datetime | None
latest_heartbeat: datetime | None
executor_class: datetime | None
executor_class: str | None
hostname: str | None
unixname: str | None

Expand Down
1 change: 0 additions & 1 deletion airflow/api_fastapi/core_api/openapi/v1-generated.yaml
Original file line numberDiff line numberDiff line change
Expand Up@@ -10051,7 +10051,6 @@ components:
executor_class:
anyOf:
- type: string
format: date-time
- type: 'null'
title: Executor Class
hostname:
Expand Down
16 changes: 16 additions & 0 deletions airflow/cli/api/__init__.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,16 @@
# 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.
264 changes: 264 additions & 0 deletions airflow/cli/api/client.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,264 @@
# 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

import contextlib
import json
import os
import sys
from functools import wraps
from typing import TYPE_CHECKING, Any, Callable, TypeVar, cast

import httpx
import keyring
import structlog
from platformdirs import user_config_path
from uuid6 import uuid7

from airflow.cli.api.operations import (
AssetsOperations,
BackfillsOperations,
ConfigOperations,
ConnectionsOperations,
DagOperations,
DagRunOperations,
JobsOperations,
PoolsOperations,
ProvidersOperations,
ServerResponseError,
VariablesOperations,
VersionOperations,
)
from airflow.exceptions import AirflowNotFoundException
from airflow.typing_compat import ParamSpec
from airflow.version import version

if TYPE_CHECKING:
# # methodtools doesn't have typestubs, so give a stub
def lru_cache(maxsize: int | None = 128):
def wrapper(f):
return f

return wrapper
else:
from methodtools import lru_cache

log = structlog.get_logger(logger_name=__name__)

__all__ = [
"Client",
"Credentials",
]

PS = ParamSpec("PS")
RT = TypeVar("RT")


def add_correlation_id(request: httpx.Request):
request.headers["correlation-id"] = str(uuid7())


def get_json_error(response: httpx.Response):
"""Raise a ServerResponseError if we can extract error info from the error."""
err = ServerResponseError.from_response(response)
if err:
log.warning("Server error ", extra=dict(err.response.json()))
raise err


def raise_on_4xx_5xx(response: httpx.Response):
return get_json_error(response) or response.raise_for_status()


# Credentials for the API
class Credentials:
"""Credentials for the API."""

api_url: str | None
api_token: str | None
api_environment: str

def __init__(
self,
api_url: str | None = None,
api_token: str | None = None,
api_environment: str = "production",
Comment thread
jedcunningham marked this conversation as resolved.
):
self.api_url = api_url
self.api_token = api_token
self.api_environment = os.getenv("AIRFLOW_CLI_ENVIRONMENT") or api_environment

@property
def input_cli_config_file(self) -> str:
"""Generate path for the CLI config file."""
return f"{self.api_environment}.json"

def save(self):
"""Save the credentials to keyring and URL to disk as a file."""
default_config_dir = user_config_path("airflow", "Apache Software Foundation")
if not os.path.exists(default_config_dir):
os.makedirs(default_config_dir)
with open(os.path.join(default_config_dir, self.input_cli_config_file), "w") as f:
json.dump({"api_url": self.api_url}, f)
keyring.set_password("airflow-cli", f"api_token-{self.api_environment}", self.api_token)

def load(self) -> Credentials:
"""Load the credentials from keyring and URL from disk file."""
default_config_dir = user_config_path("airflow", "Apache Software Foundation")
if os.path.exists(default_config_dir):
with open(os.path.join(default_config_dir, self.input_cli_config_file)) as f:
credentials = json.load(f)
self.api_url = credentials["api_url"]
self.api_token = keyring.get_password("airflow-cli", f"api_token-{self.api_environment}")
return self
else:
raise AirflowNotFoundException(f"No credentials found in {default_config_dir}")


class BearerAuth(httpx.Auth):
def __init__(self, token: str):
self.token: str = token

def auth_flow(self, request: httpx.Request):
if self.token:
request.headers["Authorization"] = "Bearer " + self.token
yield request


class Client(httpx.Client):
"""Client for the Airflow REST API."""

def __init__(self, *, base_url: str, token: str, **kwargs: Any):
auth = BearerAuth(token)
print(f"token: {token}")
kwargs["base_url"] = f"{base_url}/public"
pyver = f"{'.'.join(map(str, sys.version_info[:3]))}"
super().__init__(
auth=auth,
headers={"user-agent": f"apache-airflow-cli/{version} (Python/{pyver})"},
event_hooks={"response": [raise_on_4xx_5xx], "request": [add_correlation_id]},
**kwargs,
)

@lru_cache() # type: ignore[misc]
@property
def assets(self):
"""Operations related to assets."""
return AssetsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def backfills(self):
"""Operations related to backfills."""
return BackfillsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def configs(self):
"""Operations related to configs."""
return ConfigOperations(self)

@lru_cache() # type: ignore[misc]
@property
def connections(self):
"""Operations related to connections."""
return ConnectionsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def dags(self):
"""Operations related to DAGs."""
return DagOperations(self)

@lru_cache() # type: ignore[misc]
@property
def dag_runs(self):
"""Operations related to DAG runs."""
return DagRunOperations(self)

@lru_cache() # type: ignore[misc]
@property
def jobs(self):
"""Operations related to jobs."""
return JobsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def pools(self):
"""Operations related to pools."""
return PoolsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def providers(self):
"""Operations related to providers."""
return ProvidersOperations(self)

@lru_cache() # type: ignore[misc]
@property
def variables(self):
"""Operations related to variables."""
return VariablesOperations(self)

@lru_cache() # type: ignore[misc]
@property
def version(self):
"""Get the version of the server."""
return VersionOperations(self)


# API Client Decorator for CLI Actions
@contextlib.contextmanager
def get_client():
"""Get CLI API client."""
cli_api_client = None
try:
credentials = Credentials().load()
limits = httpx.Limits(max_keepalive_connections=1, max_connections=1)
cli_api_client = Client(base_url=credentials.api_url, limits=limits, token=credentials.api_token)
yield cli_api_client
except AirflowNotFoundException as e:
raise e
Comment thread
jedcunningham marked this conversation as resolved.
finally:
if cli_api_client:
cli_api_client.close()


def provide_api_client(func: Callable[PS, RT]) -> Callable[PS, RT]:
"""
Provide a CLI API Client to the decorated function.

CLI API Client shouldn't be passed to the function when this wrapper is used
if the purpose is not mocking or testing.
If you want to reuse a CLI API Client or run the function as part of
API call, you pass it to the function, if not this wrapper
will create one and close it for you.
"""

@wraps(func)
def wrapper(*args, **kwargs) -> RT:
if "cli_api_client" not in kwargs:
with get_client() as cli_api_client:
return func(*args, cli_api_client=cli_api_client, **kwargs)
# The CLI API Client should be only passed for Mocking and Testing
return func(*args, **kwargs)

return wrapper


NEW_CLI_API_CLIENT: Client = cast(Client, None)
16 changes: 16 additions & 0 deletions airflow/cli/api/datamodels/__init__.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,16 @@
# 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.
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Remove or un-stick sticky/fixed headers that block content (function() { function unstick() { document.querySelectorAll('header, nav, [role="banner"], .header, .navbar, .sticky, .fixed-top, [style*="position: fixed"], [style*="position:sticky"]').forEach(function(el) { if (el.style.position === 'fixed' || el.style.position === 'sticky' || getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') { el.style.position = 'static'; el.style.top = 'auto'; el.style.zIndex = 'auto'; } }); } unstick(); var observer = new MutationObserver(unstick); observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] }); })(); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion airflow/api_fastapi/core_api/datamodels/job.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -31,7 +31,7 @@ class JobResponse(BaseModel):
start_date: datetime | None
end_date: datetime | None
latest_heartbeat: datetime | None
executor_class: datetime | None
executor_class: str | None
hostname: str | None
unixname: str | None

Expand Down
1 change: 0 additions & 1 deletion airflow/api_fastapi/core_api/openapi/v1-generated.yaml
Original file line numberDiff line numberDiff line change
Expand Up@@ -10051,7 +10051,6 @@ components:
executor_class:
anyOf:
- type: string
format: date-time
- type: 'null'
title: Executor Class
hostname:
Expand Down
16 changes: 16 additions & 0 deletions airflow/cli/api/__init__.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,16 @@
# 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.
264 changes: 264 additions & 0 deletions airflow/cli/api/client.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,264 @@
# 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

import contextlib
import json
import os
import sys
from functools import wraps
from typing import TYPE_CHECKING, Any, Callable, TypeVar, cast

import httpx
import keyring
import structlog
from platformdirs import user_config_path
from uuid6 import uuid7

from airflow.cli.api.operations import (
AssetsOperations,
BackfillsOperations,
ConfigOperations,
ConnectionsOperations,
DagOperations,
DagRunOperations,
JobsOperations,
PoolsOperations,
ProvidersOperations,
ServerResponseError,
VariablesOperations,
VersionOperations,
)
from airflow.exceptions import AirflowNotFoundException
from airflow.typing_compat import ParamSpec
from airflow.version import version

if TYPE_CHECKING:
# # methodtools doesn't have typestubs, so give a stub
def lru_cache(maxsize: int | None = 128):
def wrapper(f):
return f

return wrapper
else:
from methodtools import lru_cache

log = structlog.get_logger(logger_name=__name__)

__all__ = [
"Client",
"Credentials",
]

PS = ParamSpec("PS")
RT = TypeVar("RT")


def add_correlation_id(request: httpx.Request):
request.headers["correlation-id"] = str(uuid7())


def get_json_error(response: httpx.Response):
"""Raise a ServerResponseError if we can extract error info from the error."""
err = ServerResponseError.from_response(response)
if err:
log.warning("Server error ", extra=dict(err.response.json()))
raise err


def raise_on_4xx_5xx(response: httpx.Response):
return get_json_error(response) or response.raise_for_status()


# Credentials for the API
class Credentials:
"""Credentials for the API."""

api_url: str | None
api_token: str | None
api_environment: str

def __init__(
self,
api_url: str | None = None,
api_token: str | None = None,
api_environment: str = "production",
Comment thread
jedcunningham marked this conversation as resolved.
):
self.api_url = api_url
self.api_token = api_token
self.api_environment = os.getenv("AIRFLOW_CLI_ENVIRONMENT") or api_environment

@property
def input_cli_config_file(self) -> str:
"""Generate path for the CLI config file."""
return f"{self.api_environment}.json"

def save(self):
"""Save the credentials to keyring and URL to disk as a file."""
default_config_dir = user_config_path("airflow", "Apache Software Foundation")
if not os.path.exists(default_config_dir):
os.makedirs(default_config_dir)
with open(os.path.join(default_config_dir, self.input_cli_config_file), "w") as f:
json.dump({"api_url": self.api_url}, f)
keyring.set_password("airflow-cli", f"api_token-{self.api_environment}", self.api_token)

def load(self) -> Credentials:
"""Load the credentials from keyring and URL from disk file."""
default_config_dir = user_config_path("airflow", "Apache Software Foundation")
if os.path.exists(default_config_dir):
with open(os.path.join(default_config_dir, self.input_cli_config_file)) as f:
credentials = json.load(f)
self.api_url = credentials["api_url"]
self.api_token = keyring.get_password("airflow-cli", f"api_token-{self.api_environment}")
return self
else:
raise AirflowNotFoundException(f"No credentials found in {default_config_dir}")


class BearerAuth(httpx.Auth):
def __init__(self, token: str):
self.token: str = token

def auth_flow(self, request: httpx.Request):
if self.token:
request.headers["Authorization"] = "Bearer " + self.token
yield request


class Client(httpx.Client):
"""Client for the Airflow REST API."""

def __init__(self, *, base_url: str, token: str, **kwargs: Any):
auth = BearerAuth(token)
print(f"token: {token}")
kwargs["base_url"] = f"{base_url}/public"
pyver = f"{'.'.join(map(str, sys.version_info[:3]))}"
super().__init__(
auth=auth,
headers={"user-agent": f"apache-airflow-cli/{version} (Python/{pyver})"},
event_hooks={"response": [raise_on_4xx_5xx], "request": [add_correlation_id]},
**kwargs,
)

@lru_cache() # type: ignore[misc]
@property
def assets(self):
"""Operations related to assets."""
return AssetsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def backfills(self):
"""Operations related to backfills."""
return BackfillsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def configs(self):
"""Operations related to configs."""
return ConfigOperations(self)

@lru_cache() # type: ignore[misc]
@property
def connections(self):
"""Operations related to connections."""
return ConnectionsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def dags(self):
"""Operations related to DAGs."""
return DagOperations(self)

@lru_cache() # type: ignore[misc]
@property
def dag_runs(self):
"""Operations related to DAG runs."""
return DagRunOperations(self)

@lru_cache() # type: ignore[misc]
@property
def jobs(self):
"""Operations related to jobs."""
return JobsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def pools(self):
"""Operations related to pools."""
return PoolsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def providers(self):
"""Operations related to providers."""
return ProvidersOperations(self)

@lru_cache() # type: ignore[misc]
@property
def variables(self):
"""Operations related to variables."""
return VariablesOperations(self)

@lru_cache() # type: ignore[misc]
@property
def version(self):
"""Get the version of the server."""
return VersionOperations(self)


# API Client Decorator for CLI Actions
@contextlib.contextmanager
def get_client():
"""Get CLI API client."""
cli_api_client = None
try:
credentials = Credentials().load()
limits = httpx.Limits(max_keepalive_connections=1, max_connections=1)
cli_api_client = Client(base_url=credentials.api_url, limits=limits, token=credentials.api_token)
yield cli_api_client
except AirflowNotFoundException as e:
raise e
Comment thread
jedcunningham marked this conversation as resolved.
finally:
if cli_api_client:
cli_api_client.close()


def provide_api_client(func: Callable[PS, RT]) -> Callable[PS, RT]:
"""
Provide a CLI API Client to the decorated function.

CLI API Client shouldn't be passed to the function when this wrapper is used
if the purpose is not mocking or testing.
If you want to reuse a CLI API Client or run the function as part of
API call, you pass it to the function, if not this wrapper
will create one and close it for you.
"""

@wraps(func)
def wrapper(*args, **kwargs) -> RT:
if "cli_api_client" not in kwargs:
with get_client() as cli_api_client:
return func(*args, cli_api_client=cli_api_client, **kwargs)
# The CLI API Client should be only passed for Mocking and Testing
return func(*args, **kwargs)

return wrapper


NEW_CLI_API_CLIENT: Client = cast(Client, None)
16 changes: 16 additions & 0 deletions airflow/cli/api/datamodels/__init__.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,16 @@
# 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.
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Universal Dark Mode - works on any site (function() { var enabled = true; function applyDarkMode() { if (!enabled) return; // Create style element if it doesn't exist var style = document.getElementById('universal-dark-mode-style'); if (!style) { style = document.createElement('style'); style.id = 'universal-dark-mode-style'; document.head.appendChild(style); } // Dark mode CSS - inverts colors but preserves images/video style.textContent = ' /* Invert everything except media */ html { filter: invert(1) hue-rotate(180deg) !important; background: #1a1a2e !important; } /* Restore images, videos, iframes, canvas */ img, video, iframe, canvas, svg, picture, [style*="background-image"] { filter: invert(1) hue-rotate(180deg) !important; } /* Preserve specific elements that should not be inverted */ .no-dark-mode, .no-dark-mode *, [data-theme="light"], [data-theme="light"], .ace_editor, .ace_editor *, .CodeMirror, .CodeMirror *, .monaco-editor, .monaco-editor *, .markdown-body pre, .markdown-body pre *, .highlight, .highlight *, pre code, pre code * { filter: none !important; } /* Fix common UI elements */ .modal, .popup, .dropdown-menu, .tooltip, .popover { filter: invert(1) hue-rotate(180deg) !important; background: #2d2d44 !important; border-color: #444 !important; } /* Scrollbars */ ::-webkit-scrollbar { background: #1a1a2e !important; } ::-webkit-scrollbar-thumb { background: #444 !important; } ::-webkit-scrollbar-thumb:hover { background: #555 !important; } /* Selection */ ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; } ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; } '; } function removeDarkMode() { var style = document.getElementById('universal-dark-mode-style'); if (style) style.remove(); } // Toggle with Alt+Shift+D document.addEventListener('keydown', function(e) { if (e.altKey && e.shiftKey && e.key === 'D') { e.preventDefault(); enabled = !enabled; if (enabled) { applyDarkMode(); console.log('[Universal Dark Mode] Enabled'); } else { removeDarkMode(); console.log('[Universal Dark Mode] Disabled'); } } }); // Apply on load applyDarkMode(); // Re-apply on dynamic content var observer = new MutationObserver(function(mutations) { if (enabled && !document.getElementById('universal-dark-mode-style')) { applyDarkMode(); } }); observer.observe(document.head, { childList: true }); console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle'); })(); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion airflow/api_fastapi/core_api/datamodels/job.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -31,7 +31,7 @@ class JobResponse(BaseModel):
start_date: datetime | None
end_date: datetime | None
latest_heartbeat: datetime | None
executor_class: datetime | None
executor_class: str | None
hostname: str | None
unixname: str | None

Expand Down
1 change: 0 additions & 1 deletion airflow/api_fastapi/core_api/openapi/v1-generated.yaml
Original file line numberDiff line numberDiff line change
Expand Up@@ -10051,7 +10051,6 @@ components:
executor_class:
anyOf:
- type: string
format: date-time
- type: 'null'
title: Executor Class
hostname:
Expand Down
16 changes: 16 additions & 0 deletions airflow/cli/api/__init__.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,16 @@
# 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.
264 changes: 264 additions & 0 deletions airflow/cli/api/client.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,264 @@
# 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

import contextlib
import json
import os
import sys
from functools import wraps
from typing import TYPE_CHECKING, Any, Callable, TypeVar, cast

import httpx
import keyring
import structlog
from platformdirs import user_config_path
from uuid6 import uuid7

from airflow.cli.api.operations import (
AssetsOperations,
BackfillsOperations,
ConfigOperations,
ConnectionsOperations,
DagOperations,
DagRunOperations,
JobsOperations,
PoolsOperations,
ProvidersOperations,
ServerResponseError,
VariablesOperations,
VersionOperations,
)
from airflow.exceptions import AirflowNotFoundException
from airflow.typing_compat import ParamSpec
from airflow.version import version

if TYPE_CHECKING:
# # methodtools doesn't have typestubs, so give a stub
def lru_cache(maxsize: int | None = 128):
def wrapper(f):
return f

return wrapper
else:
from methodtools import lru_cache

log = structlog.get_logger(logger_name=__name__)

__all__ = [
"Client",
"Credentials",
]

PS = ParamSpec("PS")
RT = TypeVar("RT")


def add_correlation_id(request: httpx.Request):
request.headers["correlation-id"] = str(uuid7())


def get_json_error(response: httpx.Response):
"""Raise a ServerResponseError if we can extract error info from the error."""
err = ServerResponseError.from_response(response)
if err:
log.warning("Server error ", extra=dict(err.response.json()))
raise err


def raise_on_4xx_5xx(response: httpx.Response):
return get_json_error(response) or response.raise_for_status()


# Credentials for the API
class Credentials:
"""Credentials for the API."""

api_url: str | None
api_token: str | None
api_environment: str

def __init__(
self,
api_url: str | None = None,
api_token: str | None = None,
api_environment: str = "production",
Comment thread
jedcunningham marked this conversation as resolved.
):
self.api_url = api_url
self.api_token = api_token
self.api_environment = os.getenv("AIRFLOW_CLI_ENVIRONMENT") or api_environment

@property
def input_cli_config_file(self) -> str:
"""Generate path for the CLI config file."""
return f"{self.api_environment}.json"

def save(self):
"""Save the credentials to keyring and URL to disk as a file."""
default_config_dir = user_config_path("airflow", "Apache Software Foundation")
if not os.path.exists(default_config_dir):
os.makedirs(default_config_dir)
with open(os.path.join(default_config_dir, self.input_cli_config_file), "w") as f:
json.dump({"api_url": self.api_url}, f)
keyring.set_password("airflow-cli", f"api_token-{self.api_environment}", self.api_token)

def load(self) -> Credentials:
"""Load the credentials from keyring and URL from disk file."""
default_config_dir = user_config_path("airflow", "Apache Software Foundation")
if os.path.exists(default_config_dir):
with open(os.path.join(default_config_dir, self.input_cli_config_file)) as f:
credentials = json.load(f)
self.api_url = credentials["api_url"]
self.api_token = keyring.get_password("airflow-cli", f"api_token-{self.api_environment}")
return self
else:
raise AirflowNotFoundException(f"No credentials found in {default_config_dir}")


class BearerAuth(httpx.Auth):
def __init__(self, token: str):
self.token: str = token

def auth_flow(self, request: httpx.Request):
if self.token:
request.headers["Authorization"] = "Bearer " + self.token
yield request


class Client(httpx.Client):
"""Client for the Airflow REST API."""

def __init__(self, *, base_url: str, token: str, **kwargs: Any):
auth = BearerAuth(token)
print(f"token: {token}")
kwargs["base_url"] = f"{base_url}/public"
pyver = f"{'.'.join(map(str, sys.version_info[:3]))}"
super().__init__(
auth=auth,
headers={"user-agent": f"apache-airflow-cli/{version} (Python/{pyver})"},
event_hooks={"response": [raise_on_4xx_5xx], "request": [add_correlation_id]},
**kwargs,
)

@lru_cache() # type: ignore[misc]
@property
def assets(self):
"""Operations related to assets."""
return AssetsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def backfills(self):
"""Operations related to backfills."""
return BackfillsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def configs(self):
"""Operations related to configs."""
return ConfigOperations(self)

@lru_cache() # type: ignore[misc]
@property
def connections(self):
"""Operations related to connections."""
return ConnectionsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def dags(self):
"""Operations related to DAGs."""
return DagOperations(self)

@lru_cache() # type: ignore[misc]
@property
def dag_runs(self):
"""Operations related to DAG runs."""
return DagRunOperations(self)

@lru_cache() # type: ignore[misc]
@property
def jobs(self):
"""Operations related to jobs."""
return JobsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def pools(self):
"""Operations related to pools."""
return PoolsOperations(self)

@lru_cache() # type: ignore[misc]
@property
def providers(self):
"""Operations related to providers."""
return ProvidersOperations(self)

@lru_cache() # type: ignore[misc]
@property
def variables(self):
"""Operations related to variables."""
return VariablesOperations(self)

@lru_cache() # type: ignore[misc]
@property
def version(self):
"""Get the version of the server."""
return VersionOperations(self)


# API Client Decorator for CLI Actions
@contextlib.contextmanager
def get_client():
"""Get CLI API client."""
cli_api_client = None
try:
credentials = Credentials().load()
limits = httpx.Limits(max_keepalive_connections=1, max_connections=1)
cli_api_client = Client(base_url=credentials.api_url, limits=limits, token=credentials.api_token)
yield cli_api_client
except AirflowNotFoundException as e:
raise e
Comment thread
jedcunningham marked this conversation as resolved.
finally:
if cli_api_client:
cli_api_client.close()


def provide_api_client(func: Callable[PS, RT]) -> Callable[PS, RT]:
"""
Provide a CLI API Client to the decorated function.

CLI API Client shouldn't be passed to the function when this wrapper is used
if the purpose is not mocking or testing.
If you want to reuse a CLI API Client or run the function as part of
API call, you pass it to the function, if not this wrapper
will create one and close it for you.
"""

@wraps(func)
def wrapper(*args, **kwargs) -> RT:
if "cli_api_client" not in kwargs:
with get_client() as cli_api_client:
return func(*args, cli_api_client=cli_api_client, **kwargs)
# The CLI API Client should be only passed for Mocking and Testing
return func(*args, **kwargs)

return wrapper


NEW_CLI_API_CLIENT: Client = cast(Client, None)
16 changes: 16 additions & 0 deletions airflow/cli/api/datamodels/__init__.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,16 @@
# 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.
Loading