From a516fad745f3ac4f47d8ad4a1a3fcfa81bed6a1b Mon Sep 17 00:00:00 2001 From: pratiksha badheka Date: Fri, 2 May 2025 11:08:27 -0700 Subject: [PATCH 01/36] ensure total_entries return all entries for corresponding endpoints --- .../airflow/api_fastapi/core_api/routes/public/connections.py | 2 ++ .../src/airflow/api_fastapi/core_api/routes/public/job.py | 2 ++ .../src/airflow/api_fastapi/core_api/routes/public/pools.py | 1 + 3 files changed, 5 insertions(+) diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/connections.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/connections.py index b4a300d1393bd..b01252e2ed56e 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/connections.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/connections.py @@ -132,6 +132,8 @@ def get_connections( connections = session.scalars(connection_select) + total_entries = len(list(connections)) + return ConnectionCollectionResponse( connections=connections, total_entries=total_entries, diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py index b1a35913207d1..c59a032764177 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py @@ -130,6 +130,8 @@ def get_jobs( if is_alive is not None: jobs = [job for job in jobs if job.is_alive()] + total_entries = len(jobs) + return JobCollectionResponse( jobs=jobs, total_entries=total_entries, diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/pools.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/pools.py index 8cdd25072270a..b78c245bce4ed 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/pools.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/pools.py @@ -116,6 +116,7 @@ def get_pools( ) pools = session.scalars(pools_select) + total_entries = len(list(pools)) return PoolCollectionResponse( pools=pools, From 60540ceb9cf342b79b15ebe7be9b1c2f0ea50480 Mon Sep 17 00:00:00 2001 From: pratiksha badheka Date: Tue, 6 May 2025 06:39:35 -0700 Subject: [PATCH 02/36] use API endpoint parameters limit and offset to retrieve all entries resolve conflict --- .../api_fastapi/core_api/routes/public/connections.py | 10 ++++++++-- .../api_fastapi/core_api/routes/public/dag_run.py | 10 +++++++++- .../airflow/api_fastapi/core_api/routes/public/job.py | 8 +++++++- .../api_fastapi/core_api/routes/public/pools.py | 11 +++++++++-- .../api_fastapi/core_api/routes/public/variables.py | 10 +++++++++- 5 files changed, 42 insertions(+), 7 deletions(-) diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/connections.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/connections.py index b01252e2ed56e..4f85293c902ba 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/connections.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/connections.py @@ -130,9 +130,15 @@ def get_connections( session=session, ) - connections = session.scalars(connection_select) + connections = session.scalars(connection_select).all() - total_entries = len(list(connections)) + if limit.value is not None: + limit.value = len(connections) + + if offset.value is not None: + offset.value = 0 + + connections = connections[offset.value:limit.value] return ConnectionCollectionResponse( connections=connections, diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_run.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_run.py index 00f93754baf2d..021bd98110d44 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_run.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_run.py @@ -376,7 +376,15 @@ def get_dag_runs( limit=limit, session=session, ) - dag_runs = session.scalars(dag_run_select) + dag_runs = session.scalars(dag_run_select).all() + + if limit.value is not None: + limit.value = len(dag_runs) + + if offset.value is not None: + offset.value = 0 + + dag_runs = dag_runs[offset.value:limit.value] return DAGRunCollectionResponse( dag_runs=dag_runs, diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py index c59a032764177..3890961046478 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py @@ -130,7 +130,13 @@ def get_jobs( if is_alive is not None: jobs = [job for job in jobs if job.is_alive()] - total_entries = len(jobs) + if limit.value is not None: + limit.value = len(jobs) + + if offset.value is not None: + offset.value = 0 + + jobs = jobs[offset.value:limit.value] return JobCollectionResponse( jobs=jobs, diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/pools.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/pools.py index b78c245bce4ed..65fc1df11bd77 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/pools.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/pools.py @@ -115,8 +115,15 @@ def get_pools( session=session, ) - pools = session.scalars(pools_select) - total_entries = len(list(pools)) + pools = session.scalars(pools_select).all() + + if limit.value is not None: + limit.value = len(pools) + + if offset.value is not None: + offset.value = 0 + + pools = pools[offset.value:limit.value] return PoolCollectionResponse( pools=pools, diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/variables.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/variables.py index 36ee71c17ecd5..952ad8fc42300 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/variables.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/variables.py @@ -112,7 +112,15 @@ def get_variables( session=session, ) - variables = session.scalars(variable_select) + variables = session.scalars(variable_select).all() + + if limit.value is not None: + limit.value = len(variables) + + if offset.value is not None: + offset.value = 0 + + variables = variables[offset.value:limit.value] return VariableCollectionResponse( variables=variables, From f1db1bf75fdc9f151e4da3c6cea7c07014cd4d78 Mon Sep 17 00:00:00 2001 From: pratiksha badheka Date: Mon, 19 May 2025 05:40:18 -0700 Subject: [PATCH 03/36] use a common fun to get all values --- .../airflow/api_fastapi/common/db/common.py | 27 +++++++++++++++++ .../airflow/api_fastapi/common/parameters.py | 1 + .../api_fastapi/core_api/routes/public/job.py | 29 +++++++++++++------ 3 files changed, 48 insertions(+), 9 deletions(-) diff --git a/airflow-core/src/airflow/api_fastapi/common/db/common.py b/airflow-core/src/airflow/api_fastapi/common/db/common.py index b129aae61305e..a80e94c6ab0d0 100644 --- a/airflow-core/src/airflow/api_fastapi/common/db/common.py +++ b/airflow-core/src/airflow/api_fastapi/common/db/common.py @@ -179,3 +179,30 @@ def paginated_select( statement = apply_filters_to_select(statement=statement, filters=[order_by, offset, limit]) return statement, total_entries + + +def return_all_entities(*, + total_entities: int, + all_entities: list, + statement: Select, + filters: Sequence[OrmClause] | None = None, + order_by: OrmClause | None = None, + offset: OrmClause , + limit: OrmClause , + session: SessionDep, +) -> list: + while total_entries > 0: + entity, total_entries = paginated_select( + statement=statement, + filters=filters, + order_by=order_by, + limit=limit, + offset=offset, + session=session, + return_total_entries=True, + ) + offset.value = offset.value + limit.value + entity = session.scalars(entity).all() + total_entries = total_entries - limit.value + all_entities.append(entity) + return all_entities diff --git a/airflow-core/src/airflow/api_fastapi/common/parameters.py b/airflow-core/src/airflow/api_fastapi/common/parameters.py index e8af5380f8a41..822bb21374577 100644 --- a/airflow-core/src/airflow/api_fastapi/common/parameters.py +++ b/airflow-core/src/airflow/api_fastapi/common/parameters.py @@ -712,6 +712,7 @@ def _transform_ti_states(states: list[str] | None) -> list[TaskInstanceState | N _DagIdAssetReferenceFilter, Depends(_DagIdAssetReferenceFilter.depends) ] + # Variables QueryVariableKeyPatternSearch = Annotated[ _SearchParam, Depends(search_param_factory(Variable.key, "variable_key_pattern")) diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py index 3890961046478..dcb86369bfe19 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py @@ -25,6 +25,7 @@ from airflow.api_fastapi.common.db.common import ( SessionDep, paginated_select, + return_all_entities ) from airflow.api_fastapi.common.parameters import ( FilterParam, @@ -125,19 +126,29 @@ def get_jobs( session=session, return_total_entries=True, ) - jobs = session.scalars(jobs_select).all() + + jobs = [] + jobs = return_all_entities( + total_entities = total_entries, + all_entities = jobs, + statement=base_select, + filters=[ + start_date_range, + end_date_range, + state, + job_type, + hostname, + executor_class, + ], + order_by=order_by, + limit=limit, + offset=offset, + session=session, + ) if is_alive is not None: jobs = [job for job in jobs if job.is_alive()] - if limit.value is not None: - limit.value = len(jobs) - - if offset.value is not None: - offset.value = 0 - - jobs = jobs[offset.value:limit.value] - return JobCollectionResponse( jobs=jobs, total_entries=total_entries, From b19c02ddbfee35e6767b37ea06e25267ba660419 Mon Sep 17 00:00:00 2001 From: pratiksha badheka Date: Wed, 28 May 2025 04:18:23 -0700 Subject: [PATCH 04/36] Use a common fun for return all results --- .../airflow/api_fastapi/common/db/common.py | 34 +++++++++------- .../api_fastapi/core_api/routes/public/job.py | 40 ++++++------------- 2 files changed, 32 insertions(+), 42 deletions(-) diff --git a/airflow-core/src/airflow/api_fastapi/common/db/common.py b/airflow-core/src/airflow/api_fastapi/common/db/common.py index a80e94c6ab0d0..70626a592edc4 100644 --- a/airflow-core/src/airflow/api_fastapi/common/db/common.py +++ b/airflow-core/src/airflow/api_fastapi/common/db/common.py @@ -28,7 +28,9 @@ from fastapi import Depends from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.orm import Session +from sqlalchemy import select,func +from airflow.models.base import Base from airflow.utils.db import get_query_count, get_query_count_async from airflow.utils.session import NEW_SESSION, create_session, create_session_async, provide_session @@ -181,18 +183,23 @@ def paginated_select( return statement, total_entries -def return_all_entities(*, - total_entities: int, - all_entities: list, +@provide_session +def return_all_entries( + *, + operation: Base, statement: Select, filters: Sequence[OrmClause] | None = None, order_by: OrmClause | None = None, - offset: OrmClause , - limit: OrmClause , - session: SessionDep, -) -> list: - while total_entries > 0: - entity, total_entries = paginated_select( + offset: OrmClause | None = None, + limit: OrmClause | None = None, + session: Session = NEW_SESSION, + return_total_entries=True, + all_results: list, +): + all_entries = session.scalar(func.count(operation.id)) + + for _entry in range(0, all_entries, limit.value): + entities, total_entries = paginated_select( statement=statement, filters=filters, order_by=order_by, @@ -201,8 +208,7 @@ def return_all_entities(*, session=session, return_total_entries=True, ) - offset.value = offset.value + limit.value - entity = session.scalars(entity).all() - total_entries = total_entries - limit.value - all_entities.append(entity) - return all_entities + all_results.append(session.scalars(select(entities)).all()) + offset.value += limit.value + total_entries = total_entries + return all_results, total_entries diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py index dcb86369bfe19..ca52bfe333d3e 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py @@ -25,7 +25,7 @@ from airflow.api_fastapi.common.db.common import ( SessionDep, paginated_select, - return_all_entities + return_all_entries ) from airflow.api_fastapi.common.parameters import ( FilterParam, @@ -110,9 +110,11 @@ def get_jobs( .options(joinedload(Job.dag_model)) ) - jobs_select, total_entries = paginated_select( - statement=base_select, - filters=[ + jobs = [] + jobs,total_entries = return_all_entries( + operation = Job, + statement = base_select, + filters = [ start_date_range, end_date_range, state, @@ -120,31 +122,13 @@ def get_jobs( hostname, executor_class, ], - order_by=order_by, - limit=limit, - offset=offset, - session=session, - return_total_entries=True, - ) + order_by = order_by, + offset = offset, + limit = limit, + session = session, + return_total_entries = True, + all_results = jobs) - jobs = [] - jobs = return_all_entities( - total_entities = total_entries, - all_entities = jobs, - statement=base_select, - filters=[ - start_date_range, - end_date_range, - state, - job_type, - hostname, - executor_class, - ], - order_by=order_by, - limit=limit, - offset=offset, - session=session, - ) if is_alive is not None: jobs = [job for job in jobs if job.is_alive()] From fcf83254df0218909c1e45ba94f10d49305c8b71 Mon Sep 17 00:00:00 2001 From: pratiksha badheka Date: Wed, 11 Jun 2025 02:36:48 -0700 Subject: [PATCH 05/36] create a method in BaseOperation and use in Pools list, restore common.py and restore endpoints - pools.py, variables.py, job.py --- .../airflow/api_fastapi/common/db/common.py | 32 +------------- .../api_fastapi/core_api/routes/public/job.py | 24 +++++------ .../core_api/routes/public/pools.py | 10 +---- .../core_api/routes/public/variables.py | 10 +---- airflow-ctl/src/airflowctl/api/operations.py | 43 +++++++++++++++++-- 5 files changed, 54 insertions(+), 65 deletions(-) diff --git a/airflow-core/src/airflow/api_fastapi/common/db/common.py b/airflow-core/src/airflow/api_fastapi/common/db/common.py index 70626a592edc4..39e1adb2f7d12 100644 --- a/airflow-core/src/airflow/api_fastapi/common/db/common.py +++ b/airflow-core/src/airflow/api_fastapi/common/db/common.py @@ -28,7 +28,7 @@ from fastapi import Depends from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.orm import Session -from sqlalchemy import select,func +from sqlalchemy import select, func from airflow.models.base import Base from airflow.utils.db import get_query_count, get_query_count_async @@ -182,33 +182,3 @@ def paginated_select( return statement, total_entries - -@provide_session -def return_all_entries( - *, - operation: Base, - statement: Select, - filters: Sequence[OrmClause] | None = None, - order_by: OrmClause | None = None, - offset: OrmClause | None = None, - limit: OrmClause | None = None, - session: Session = NEW_SESSION, - return_total_entries=True, - all_results: list, -): - all_entries = session.scalar(func.count(operation.id)) - - for _entry in range(0, all_entries, limit.value): - entities, total_entries = paginated_select( - statement=statement, - filters=filters, - order_by=order_by, - limit=limit, - offset=offset, - session=session, - return_total_entries=True, - ) - all_results.append(session.scalars(select(entities)).all()) - offset.value += limit.value - total_entries = total_entries - return all_results, total_entries diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py index ca52bfe333d3e..0fcc17229ff1a 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py @@ -25,7 +25,6 @@ from airflow.api_fastapi.common.db.common import ( SessionDep, paginated_select, - return_all_entries ) from airflow.api_fastapi.common.parameters import ( FilterParam, @@ -110,11 +109,9 @@ def get_jobs( .options(joinedload(Job.dag_model)) ) - jobs = [] - jobs,total_entries = return_all_entries( - operation = Job, - statement = base_select, - filters = [ + jobs_select, total_entries = paginated_select( + statement=base_select, + filters=[ start_date_range, end_date_range, state, @@ -122,13 +119,13 @@ def get_jobs( hostname, executor_class, ], - order_by = order_by, - offset = offset, - limit = limit, - session = session, - return_total_entries = True, - all_results = jobs) - + order_by=order_by, + limit=limit, + offset=offset, + session=session, + return_total_entries=True, + ) + jobs = session.scalars(jobs_select).all() if is_alive is not None: jobs = [job for job in jobs if job.is_alive()] @@ -137,3 +134,4 @@ def get_jobs( jobs=jobs, total_entries=total_entries, ) + diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/pools.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/pools.py index 65fc1df11bd77..8cdd25072270a 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/pools.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/pools.py @@ -115,15 +115,7 @@ def get_pools( session=session, ) - pools = session.scalars(pools_select).all() - - if limit.value is not None: - limit.value = len(pools) - - if offset.value is not None: - offset.value = 0 - - pools = pools[offset.value:limit.value] + pools = session.scalars(pools_select) return PoolCollectionResponse( pools=pools, diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/variables.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/variables.py index 952ad8fc42300..36ee71c17ecd5 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/variables.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/variables.py @@ -112,15 +112,7 @@ def get_variables( session=session, ) - variables = session.scalars(variable_select).all() - - if limit.value is not None: - limit.value = len(variables) - - if offset.value is not None: - offset.value = 0 - - variables = variables[offset.value:limit.value] + variables = session.scalars(variable_select) return VariableCollectionResponse( variables=variables, diff --git a/airflow-ctl/src/airflowctl/api/operations.py b/airflow-ctl/src/airflowctl/api/operations.py index 0161f566ff790..a52c1db16cdaf 100644 --- a/airflow-ctl/src/airflowctl/api/operations.py +++ b/airflow-ctl/src/airflowctl/api/operations.py @@ -22,6 +22,7 @@ import httpx import structlog +from pydantic import BaseModel from airflowctl.api.datamodels.auth_generated import LoginBody, LoginResponse from airflowctl.api.datamodels.generated import ( @@ -74,6 +75,7 @@ if TYPE_CHECKING: from airflowctl.api.client import Client + from pydantic import BaseModel log = structlog.get_logger(logger_name=__name__) @@ -147,6 +149,25 @@ def __init_subclass__(cls, **kwargs): if callable(value): setattr(cls, attr, _check_flag_and_exit_if_server_response_error(value)) + def return_all_entries( + self, *, path: str, total_entries: int, pydantic_model, offset=0, params, **kwargs + ) -> BaseModel | ServerResponseError: + params.update({"offset": 0}) + params.update(**kwargs) + entry_list = [] + try: + if total_entries == 0: + return pydantic_model.model_validate_json(self.response.content) + for _ in range(total_entries): + # while(offset <= total_entries): + self.response = self.client.get(path, params=params) + entry = pydantic_model.model_validate_json(self.response.content) + offset = offset + 50 + entry_list.append(entry) + return entry_list + except ServerResponseError as e: + raise e + # Login operations class LoginOperations: @@ -592,11 +613,18 @@ def get(self, pool_name: str) -> PoolResponse | ServerResponseError: except ServerResponseError as e: raise e - def list(self) -> PoolCollectionResponse | ServerResponseError: + def list(self): """List all pools.""" try: self.response = self.client.get("pools") - return PoolCollectionResponse.model_validate_json(self.response.content) + total_entries = PoolCollectionResponse.model_validate_json(self.response.content).total_entries + return super().return_all_entries( + path="pools", + total_entries=total_entries, + pydantic_model=PoolCollectionResponse, + offset=0, + params={}, + ) except ServerResponseError as e: raise e @@ -660,7 +688,16 @@ def list(self) -> VariableCollectionResponse | ServerResponseError: """List all variables.""" try: self.response = self.client.get("variables") - return VariableCollectionResponse.model_validate_json(self.response.content) + total_entries = VariableCollectionResponse.model_validate_json( + self.response.content + ).total_entries + return super().return_all_entries( + path="variables", + total_entries=total_entries, + pydantic_model=VariableCollectionResponse, + offset=0, + params={}, + ) except ServerResponseError as e: raise e From 7d52d30f37bc4bec75717c402b9fa4fe998abec5 Mon Sep 17 00:00:00 2001 From: pratiksha badheka Date: Thu, 12 Jun 2025 10:24:18 -0700 Subject: [PATCH 06/36] improve looping and fix static checks --- airflow-ctl/src/airflowctl/api/operations.py | 23 +++++++++----------- 1 file changed, 10 insertions(+), 13 deletions(-) diff --git a/airflow-ctl/src/airflowctl/api/operations.py b/airflow-ctl/src/airflowctl/api/operations.py index a52c1db16cdaf..5cda772a0c084 100644 --- a/airflow-ctl/src/airflowctl/api/operations.py +++ b/airflow-ctl/src/airflowctl/api/operations.py @@ -22,7 +22,6 @@ import httpx import structlog -from pydantic import BaseModel from airflowctl.api.datamodels.auth_generated import LoginBody, LoginResponse from airflowctl.api.datamodels.generated import ( @@ -75,7 +74,6 @@ if TYPE_CHECKING: from airflowctl.api.client import Client - from pydantic import BaseModel log = structlog.get_logger(logger_name=__name__) @@ -150,19 +148,18 @@ def __init_subclass__(cls, **kwargs): setattr(cls, attr, _check_flag_and_exit_if_server_response_error(value)) def return_all_entries( - self, *, path: str, total_entries: int, pydantic_model, offset=0, params, **kwargs - ) -> BaseModel | ServerResponseError: + self, *, path: str, total_entries: int, data_model, offset=0, params, **kwargs + ) -> list | ServerResponseError: params.update({"offset": 0}) params.update(**kwargs) entry_list = [] try: if total_entries == 0: - return pydantic_model.model_validate_json(self.response.content) - for _ in range(total_entries): - # while(offset <= total_entries): + return [data_model.model_validate_json(self.response.content)] + while offset <= total_entries: self.response = self.client.get(path, params=params) - entry = pydantic_model.model_validate_json(self.response.content) - offset = offset + 50 + entry = data_model.model_validate_json(self.response.content) + offset = offset + 50 # default limit params = 50 entry_list.append(entry) return entry_list except ServerResponseError as e: @@ -613,7 +610,7 @@ def get(self, pool_name: str) -> PoolResponse | ServerResponseError: except ServerResponseError as e: raise e - def list(self): + def list(self) -> list | ServerResponseError: """List all pools.""" try: self.response = self.client.get("pools") @@ -621,7 +618,7 @@ def list(self): return super().return_all_entries( path="pools", total_entries=total_entries, - pydantic_model=PoolCollectionResponse, + data_model=PoolCollectionResponse, offset=0, params={}, ) @@ -684,7 +681,7 @@ def get(self, variable_key: str) -> VariableResponse | ServerResponseError: except ServerResponseError as e: raise e - def list(self) -> VariableCollectionResponse | ServerResponseError: + def list(self) -> list | ServerResponseError: """List all variables.""" try: self.response = self.client.get("variables") @@ -694,7 +691,7 @@ def list(self) -> VariableCollectionResponse | ServerResponseError: return super().return_all_entries( path="variables", total_entries=total_entries, - pydantic_model=VariableCollectionResponse, + data_model=VariableCollectionResponse, offset=0, params={}, ) From b5333905c3598eb469e66cd648bfcf5e0e8f566e Mon Sep 17 00:00:00 2001 From: pratiksha badheka Date: Thu, 12 Jun 2025 10:28:20 -0700 Subject: [PATCH 07/36] restore endpoints and fix static checks for job.py --- .../api_fastapi/core_api/routes/public/connections.py | 10 +--------- .../api_fastapi/core_api/routes/public/dag_run.py | 10 +--------- .../airflow/api_fastapi/core_api/routes/public/job.py | 1 - 3 files changed, 2 insertions(+), 19 deletions(-) diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/connections.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/connections.py index 4f85293c902ba..b4a300d1393bd 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/connections.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/connections.py @@ -130,15 +130,7 @@ def get_connections( session=session, ) - connections = session.scalars(connection_select).all() - - if limit.value is not None: - limit.value = len(connections) - - if offset.value is not None: - offset.value = 0 - - connections = connections[offset.value:limit.value] + connections = session.scalars(connection_select) return ConnectionCollectionResponse( connections=connections, diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_run.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_run.py index 021bd98110d44..00f93754baf2d 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_run.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_run.py @@ -376,15 +376,7 @@ def get_dag_runs( limit=limit, session=session, ) - dag_runs = session.scalars(dag_run_select).all() - - if limit.value is not None: - limit.value = len(dag_runs) - - if offset.value is not None: - offset.value = 0 - - dag_runs = dag_runs[offset.value:limit.value] + dag_runs = session.scalars(dag_run_select) return DAGRunCollectionResponse( dag_runs=dag_runs, diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py index 0fcc17229ff1a..b1a35913207d1 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/job.py @@ -134,4 +134,3 @@ def get_jobs( jobs=jobs, total_entries=total_entries, ) - From d5b24a20a74f67719adca5fddd321fd60c9fda1a Mon Sep 17 00:00:00 2001 From: pratiksha badheka Date: Thu, 12 Jun 2025 10:31:06 -0700 Subject: [PATCH 08/36] fix static checks from common.py --- airflow-core/src/airflow/api_fastapi/common/db/common.py | 3 --- 1 file changed, 3 deletions(-) diff --git a/airflow-core/src/airflow/api_fastapi/common/db/common.py b/airflow-core/src/airflow/api_fastapi/common/db/common.py index 39e1adb2f7d12..b129aae61305e 100644 --- a/airflow-core/src/airflow/api_fastapi/common/db/common.py +++ b/airflow-core/src/airflow/api_fastapi/common/db/common.py @@ -28,9 +28,7 @@ from fastapi import Depends from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.orm import Session -from sqlalchemy import select, func -from airflow.models.base import Base from airflow.utils.db import get_query_count, get_query_count_async from airflow.utils.session import NEW_SESSION, create_session, create_session_async, provide_session @@ -181,4 +179,3 @@ def paginated_select( statement = apply_filters_to_select(statement=statement, filters=[order_by, offset, limit]) return statement, total_entries - From 77b731053e55cc244883def39f5295cd29bde578 Mon Sep 17 00:00:00 2001 From: pratiksha badheka Date: Fri, 13 Jun 2025 10:13:29 -0700 Subject: [PATCH 09/36] update offset params --- airflow-ctl/src/airflowctl/api/operations.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/airflow-ctl/src/airflowctl/api/operations.py b/airflow-ctl/src/airflowctl/api/operations.py index 5cda772a0c084..0ab180f2f639a 100644 --- a/airflow-ctl/src/airflowctl/api/operations.py +++ b/airflow-ctl/src/airflowctl/api/operations.py @@ -148,18 +148,18 @@ def __init_subclass__(cls, **kwargs): setattr(cls, attr, _check_flag_and_exit_if_server_response_error(value)) def return_all_entries( - self, *, path: str, total_entries: int, data_model, offset=0, params, **kwargs + self, *, path: str, total_entries: int, data_model, offset=0 ,params, **kwargs ) -> list | ServerResponseError: - params.update({"offset": 0}) params.update(**kwargs) entry_list = [] try: if total_entries == 0: return [data_model.model_validate_json(self.response.content)] while offset <= total_entries: + params.update({"offset": offset}) self.response = self.client.get(path, params=params) entry = data_model.model_validate_json(self.response.content) - offset = offset + 50 # default limit params = 50 + offset = offset + 50 # default limit params = 50 entry_list.append(entry) return entry_list except ServerResponseError as e: From 6b7dbb284271ae5b67528acfbecc95242a9a4ada Mon Sep 17 00:00:00 2001 From: pratiksha badheka Date: Fri, 13 Jun 2025 10:31:40 -0700 Subject: [PATCH 10/36] update default limit = 50 to limit --- airflow-ctl/src/airflowctl/api/operations.py | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/airflow-ctl/src/airflowctl/api/operations.py b/airflow-ctl/src/airflowctl/api/operations.py index 0ab180f2f639a..14483c9eb88d4 100644 --- a/airflow-ctl/src/airflowctl/api/operations.py +++ b/airflow-ctl/src/airflowctl/api/operations.py @@ -148,7 +148,7 @@ def __init_subclass__(cls, **kwargs): setattr(cls, attr, _check_flag_and_exit_if_server_response_error(value)) def return_all_entries( - self, *, path: str, total_entries: int, data_model, offset=0 ,params, **kwargs + self, *, path: str, total_entries: int, data_model, offset=0, limit: int = 50, params, **kwargs ) -> list | ServerResponseError: params.update(**kwargs) entry_list = [] @@ -156,10 +156,10 @@ def return_all_entries( if total_entries == 0: return [data_model.model_validate_json(self.response.content)] while offset <= total_entries: - params.update({"offset": offset}) + params.update({"offset": offset, "limit": limit}) self.response = self.client.get(path, params=params) entry = data_model.model_validate_json(self.response.content) - offset = offset + 50 # default limit params = 50 + offset = offset + limit # default limit params = 50 entry_list.append(entry) return entry_list except ServerResponseError as e: @@ -620,6 +620,7 @@ def list(self) -> list | ServerResponseError: total_entries=total_entries, data_model=PoolCollectionResponse, offset=0, + limit=50, params={}, ) except ServerResponseError as e: @@ -693,6 +694,7 @@ def list(self) -> list | ServerResponseError: total_entries=total_entries, data_model=VariableCollectionResponse, offset=0, + limit=50, params={}, ) except ServerResponseError as e: From 361a4ec132e14ea090e0e3b4298c24806b1d9fbf Mon Sep 17 00:00:00 2001 From: pratiksha badheka Date: Tue, 17 Jun 2025 10:16:57 -0700 Subject: [PATCH 11/36] remove equal condition from while loop and remove unnecessary params --- airflow-ctl/src/airflowctl/api/operations.py | 12 +++--------- 1 file changed, 3 insertions(+), 9 deletions(-) diff --git a/airflow-ctl/src/airflowctl/api/operations.py b/airflow-ctl/src/airflowctl/api/operations.py index 14483c9eb88d4..21d00f0da26a1 100644 --- a/airflow-ctl/src/airflowctl/api/operations.py +++ b/airflow-ctl/src/airflowctl/api/operations.py @@ -148,15 +148,15 @@ def __init_subclass__(cls, **kwargs): setattr(cls, attr, _check_flag_and_exit_if_server_response_error(value)) def return_all_entries( - self, *, path: str, total_entries: int, data_model, offset=0, limit: int = 50, params, **kwargs + self, *, path: str, total_entries: int, data_model, offset=0, limit: int = 50, params={}, **kwargs ) -> list | ServerResponseError: params.update(**kwargs) entry_list = [] try: if total_entries == 0: return [data_model.model_validate_json(self.response.content)] - while offset <= total_entries: - params.update({"offset": offset, "limit": limit}) + while offset < total_entries: + params.update({"offset": offset}) self.response = self.client.get(path, params=params) entry = data_model.model_validate_json(self.response.content) offset = offset + limit # default limit params = 50 @@ -619,9 +619,6 @@ def list(self) -> list | ServerResponseError: path="pools", total_entries=total_entries, data_model=PoolCollectionResponse, - offset=0, - limit=50, - params={}, ) except ServerResponseError as e: raise e @@ -693,9 +690,6 @@ def list(self) -> list | ServerResponseError: path="variables", total_entries=total_entries, data_model=VariableCollectionResponse, - offset=0, - limit=50, - params={}, ) except ServerResponseError as e: raise e From ee0648c01ccd7bf96721e28cd88718a164ccc7f9 Mon Sep 17 00:00:00 2001 From: pratiksha badheka Date: Wed, 18 Jun 2025 00:42:43 -0700 Subject: [PATCH 12/36] use first pass response and merge it with the return list --- airflow-ctl/src/airflowctl/api/operations.py | 36 +++++++++++++++----- 1 file changed, 28 insertions(+), 8 deletions(-) diff --git a/airflow-ctl/src/airflowctl/api/operations.py b/airflow-ctl/src/airflowctl/api/operations.py index 21d00f0da26a1..fc970bb11d1d5 100644 --- a/airflow-ctl/src/airflowctl/api/operations.py +++ b/airflow-ctl/src/airflowctl/api/operations.py @@ -148,19 +148,26 @@ def __init_subclass__(cls, **kwargs): setattr(cls, attr, _check_flag_and_exit_if_server_response_error(value)) def return_all_entries( - self, *, path: str, total_entries: int, data_model, offset=0, limit: int = 50, params={}, **kwargs + self, + *, + path: str, + total_entries: int, + data_model, + entry_list: list, + offset=0, + limit: int = 50, + params=None, + **kwargs, ) -> list | ServerResponseError: params.update(**kwargs) - entry_list = [] try: - if total_entries == 0: - return [data_model.model_validate_json(self.response.content)] while offset < total_entries: params.update({"offset": offset}) self.response = self.client.get(path, params=params) entry = data_model.model_validate_json(self.response.content) offset = offset + limit # default limit params = 50 entry_list.append(entry) + return entry_list except ServerResponseError as e: raise e @@ -614,11 +621,19 @@ def list(self) -> list | ServerResponseError: """List all pools.""" try: self.response = self.client.get("pools") - total_entries = PoolCollectionResponse.model_validate_json(self.response.content).total_entries + primary_data = PoolCollectionResponse.model_validate_json(self.response.content) + entry_list = [] + entry_list.append(primary_data) + total_entries = primary_data.total_entries + limit = 50 # default + if total_entries < limit: + return entry_list return super().return_all_entries( path="pools", total_entries=total_entries, + limit=9, data_model=PoolCollectionResponse, + entry_list=entry_list, ) except ServerResponseError as e: raise e @@ -683,13 +698,18 @@ def list(self) -> list | ServerResponseError: """List all variables.""" try: self.response = self.client.get("variables") - total_entries = VariableCollectionResponse.model_validate_json( - self.response.content - ).total_entries + primary_data = VariableCollectionResponse.model_validate_json(self.response.content) + entry_list = [] + entry_list.append(primary_data) + total_entries = primary_data.total_entries + limit = 50 + if total_entries < limit: + return entry_list return super().return_all_entries( path="variables", total_entries=total_entries, data_model=VariableCollectionResponse, + entry_list=entry_list, ) except ServerResponseError as e: raise e From 8304a46b3a64e797323aa09ffa3cbb55c925a71b Mon Sep 17 00:00:00 2001 From: pratiksha badheka Date: Wed, 18 Jun 2025 02:58:29 -0700 Subject: [PATCH 13/36] define missing parameters-type in return_all_entries method --- airflow-ctl/src/airflowctl/api/operations.py | 13 ++++++++----- 1 file changed, 8 insertions(+), 5 deletions(-) diff --git a/airflow-ctl/src/airflowctl/api/operations.py b/airflow-ctl/src/airflowctl/api/operations.py index fc970bb11d1d5..dc8d8527c3706 100644 --- a/airflow-ctl/src/airflowctl/api/operations.py +++ b/airflow-ctl/src/airflowctl/api/operations.py @@ -73,6 +73,8 @@ from airflowctl.exceptions import AirflowCtlConnectionException if TYPE_CHECKING: + from pydantic import BaseModel + from airflowctl.api.client import Client log = structlog.get_logger(logger_name=__name__) @@ -152,14 +154,16 @@ def return_all_entries( *, path: str, total_entries: int, - data_model, + data_model: type[BaseModel], entry_list: list, - offset=0, + offset: int = 0, limit: int = 50, - params=None, + params: dict | None = None, **kwargs, ) -> list | ServerResponseError: - params.update(**kwargs) + if params is None: + params = {} + params.update({"limit": limit}, **kwargs) try: while offset < total_entries: params.update({"offset": offset}) @@ -631,7 +635,6 @@ def list(self) -> list | ServerResponseError: return super().return_all_entries( path="pools", total_entries=total_entries, - limit=9, data_model=PoolCollectionResponse, entry_list=entry_list, ) From e860a581d7637e1185c8a5281b5a0f4ef40bcf1a Mon Sep 17 00:00:00 2001 From: pratiksha badheka Date: Thu, 19 Jun 2025 04:02:36 -0700 Subject: [PATCH 14/36] add tests for return_all_entries and update list in pool and variable tests --- .../tests/airflow_ctl/api/test_operations.py | 62 ++++++++++++++++++- 1 file changed, 60 insertions(+), 2 deletions(-) diff --git a/airflow-ctl/tests/airflow_ctl/api/test_operations.py b/airflow-ctl/tests/airflow_ctl/api/test_operations.py index 1a894990183a0..d24de0d714ce6 100644 --- a/airflow-ctl/tests/airflow_ctl/api/test_operations.py +++ b/airflow-ctl/tests/airflow_ctl/api/test_operations.py @@ -21,9 +21,11 @@ import json import uuid from typing import TYPE_CHECKING +from unittest.mock import MagicMock import httpx import pytest +from pydantic import BaseModel from airflowctl.api.client import Client, ClientKind from airflowctl.api.datamodels.auth_generated import LoginBody, LoginResponse @@ -90,6 +92,7 @@ VariableResponse, VersionInfo, ) +from airflowctl.api.operations import BaseOperations from airflowctl.exceptions import AirflowCtlConnectionException if TYPE_CHECKING: @@ -106,6 +109,11 @@ def make_api_client( return Client(base_url=base_url, transport=transport, token=token, kind=kind) +class HelloCollectionResponse(BaseModel): + hello: list[str] + total_entries: int + + class TestBaseOperations: def test_server_connection_refused(self): client = make_api_client(base_url="http://localhost") @@ -114,6 +122,52 @@ def test_server_connection_refused(self): ): client.connections.get("1") + @pytest.mark.parametrize( + "total_entries, offset, limit, expected_response", + [ + (0, 0, 50, [HelloCollectionResponse(hello=[], total_entries=0)]), + (1, 0, 50, [HelloCollectionResponse(hello=["hello"], total_entries=1)]), + (3, 2, 50, [HelloCollectionResponse(hello=["hello"], total_entries=3)]), + ( + 20, + 5, + 5, + [ + (HelloCollectionResponse(hello=["hello"], total_entries=20)), + (HelloCollectionResponse(hello=["hello"], total_entries=20)), + (HelloCollectionResponse(hello=["hello"], total_entries=20)), + ], + ), + (2, 3, 50, []), + ], + ) + def test_return_all_entries(self, total_entries, limit, offset, expected_response): + mock_operation = MagicMock(spec=BaseOperations) + mocked_response = [] + if offset < total_entries: + while offset < total_entries: + response = HelloCollectionResponse(hello=["hello"], total_entries=total_entries) + mocked_response.append(response) + offset += limit + mock_operation.return_all_entries.return_value = mocked_response + elif offset == total_entries: + mocked_response.append(HelloCollectionResponse(hello=[], total_entries=0)) + mock_operation.return_all_entries.return_value = mocked_response + else: + mock_operation.return_all_entries.return_value = mocked_response + + assert ( + mock_operation.return_all_entries( + path="", + total_entries=total_entries, + data_model=HelloCollectionResponse, + entry_list=[], + offset=offset, + limit=limit, + ) + == expected_response + ) + class TestAssetsOperations: asset_id: int = 1 @@ -975,7 +1029,9 @@ def handle_request(request: httpx.Request) -> httpx.Response: client = make_api_client(transport=httpx.MockTransport(handle_request)) response = client.pools.list() - assert response == self.pool_response_collection + pool_collection_response_list = [] + pool_collection_response_list.append(self.pool_response_collection) + assert response == pool_collection_response_list def test_create(self): def handle_request(request: httpx.Request) -> httpx.Response: @@ -1078,7 +1134,9 @@ def handle_request(request: httpx.Request) -> httpx.Response: client = make_api_client(transport=httpx.MockTransport(handle_request)) response = client.variables.list() - assert response == self.variable_collection_response + variable_response = [] + variable_response.append(self.variable_collection_response) + assert response == variable_response def test_create(self): def handle_request(request: httpx.Request) -> httpx.Response: From d76e8eed3ee49751b58ee9cbd354294185209faa Mon Sep 17 00:00:00 2001 From: pratiksha badheka Date: Thu, 19 Jun 2025 11:28:42 -0700 Subject: [PATCH 15/36] fix tests --- .../ctl/commands/variable_command.py | 25 +++++++++++-------- 1 file changed, 15 insertions(+), 10 deletions(-) diff --git a/airflow-ctl/src/airflowctl/ctl/commands/variable_command.py b/airflow-ctl/src/airflowctl/ctl/commands/variable_command.py index d5ffed31d8c32..9982003433a19 100644 --- a/airflow-ctl/src/airflowctl/ctl/commands/variable_command.py +++ b/airflow-ctl/src/airflowctl/ctl/commands/variable_command.py @@ -86,17 +86,22 @@ def export(args, api_client=NEW_API_CLIENT) -> None: """Export all the variables to the file.""" success_message = "[green]Export successful! {total_entries} variable(s) to {file}[/green]" var_dict = {} - variables = api_client.variables.list() + response = api_client.variables.list() - for variable in variables.variables: - if variable.description: - var_dict[variable.key] = { - "value": variable.value, - "description": variable.description, - } - else: - var_dict[variable.key] = variable.value + for variables in response: + for variable in variables.variables: + if variable.description: + var_dict[variable.key] = { + "value": variable.value, + "description": variable.description, + } + else: + var_dict[variable.key] = variable.value with open(Path(args.file), "w") as var_file: json.dump(var_dict, var_file, sort_keys=True, indent=4) - rich.print(success_message.format(total_entries=variables.total_entries, file=args.file)) + + for variables in response: + total_entries = variables.total_entries + + rich.print(success_message.format(total_entries=total_entries, file=args.file)) From 42c33758199e5bc0997b3c3370b6a6d678b15aa1 Mon Sep 17 00:00:00 2001 From: pratiksha badheka Date: Thu, 19 Jun 2025 11:29:55 -0700 Subject: [PATCH 16/36] fix static checks --- airflow-ctl/docs/images/command_hashes.txt | 17 +++++++++++++++++ airflow-ctl/docs/images/output_main.svg | 5 ++++- 2 files changed, 21 insertions(+), 1 deletion(-) diff --git a/airflow-ctl/docs/images/command_hashes.txt b/airflow-ctl/docs/images/command_hashes.txt index 84cab0f38ecc4..789efe11909a1 100644 --- a/airflow-ctl/docs/images/command_hashes.txt +++ b/airflow-ctl/docs/images/command_hashes.txt @@ -1,3 +1,4 @@ +<<<<<<< HEAD main:8d768c837899829dfd21d37253d2fb44 assets:b3ae2b933e54528bf486ff28e887804d auth:f396d4bce90215599dde6ad0a8f30f29 @@ -12,3 +13,19 @@ providers:1c0afb2dff31d93ab2934b032a2250ab variables:0b04188937b3c364204ef4cc9a541c62 version:000176f03a175890b12181c8569e2d0f auth login:5277c653ff6dce51f37472dc0bda9775 +======= +main:ed626142c04ec07e836e72e0069f1f3f +assets:bd74e73e54641bac100b88ca29641df2 +auth:ef4122d3f5e4b2ac19cb0d3e12c8594b +backfills:e0cba4448d576d1b53ea79d6dcdbe035 +config:807fd4874d29702624b231a1e4ea0bc9 +connections:da4f6807ca2a265ed6d6e734b5355fe2 +dag:dab7c8aa1a62fa011b80bb7132bcc32a +dagrun:7b3e06a3664cc7ceb18457b4c0895532 +jobs:806174e6c9511db669705279ed6a00b9 +pools:2c17a4131b6481bd8fe9120982606db2 +providers:d053e6f17ff271e1e08942378344d27b +variables:cd3970589b2cb1e3ebd9a0b7f2ffdf4d +version:19f901e228111d8ba2ef47d8722f9b87 +auth login:348c25d49128b6007ac97dae2ef7563f +>>>>>>> a1a034ef62 (fix static checks) diff --git a/airflow-ctl/docs/images/output_main.svg b/airflow-ctl/docs/images/output_main.svg index 98bf85d2fd28f..77131a3341dc1 100644 --- a/airflow-ctl/docs/images/output_main.svg +++ b/airflow-ctl/docs/images/output_main.svg @@ -1,4 +1,4 @@ - + - - + + - + - + - + - + - + - + - + - + - + - + - + - + - + - + - + - + - + - + - + - + - + - + - + - + - + - Command: main + Command: main - + - - Usage: airflowctl [-h] GROUP_OR_COMMAND ... - -Positional Arguments: -  GROUP_OR_COMMAND - -    Groups -      assets        Perform Assets operations -      auth          Manage authentication for CLI. Either pass token from -                    environment variable/parameter or pass username and -                    password. -      backfills     Perform Backfills operations -      config        Perform Config operations -      connections   Perform Connections operations -      dag           Perform Dag operations -      dagrun        Perform DagRun operations -      jobs          Perform Jobs operations -      pools         Perform Pools operations -      providers     Perform Providers operations -      variables     Perform Variables operations - -    Commands: -      version       Show version information - -Options: -  -h, --help        show this help message and exit + + Usage: airflowctl [-h] GROUP_OR_COMMAND ... + +Positional Arguments: +  GROUP_OR_COMMAND + +    Groups +      assets        Perform Assets operations +      auth          Manage authentication for CLI. Either pass token from +                    environment variable/parameter or pass username and +                    password. +      backfills     Perform Backfills operations +      base          Perform Base operations +      config        Perform Config operations +      connections   Perform Connections operations +      dag           Perform Dag operations +      dagrun        Perform DagRun operations +      jobs          Perform Jobs operations +      pools         Perform Pools operations +      providers     Perform Providers operations +      variables     Perform Variables operations + +    Commands: +      version       Show version information + +Options: +  -h, --help        show this help message and exit From 307d89227bbb0f8bfe5d61bcdb2c0d5b73c16efd Mon Sep 17 00:00:00 2001 From: pratiksha badheka Date: Mon, 14 Jul 2025 09:14:17 -0700 Subject: [PATCH 22/36] fix static checks for command_hashes.txt --- airflow-ctl/docs/images/command_hashes.txt | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/airflow-ctl/docs/images/command_hashes.txt b/airflow-ctl/docs/images/command_hashes.txt index 149f4d7867429..d689874c2a0e9 100644 --- a/airflow-ctl/docs/images/command_hashes.txt +++ b/airflow-ctl/docs/images/command_hashes.txt @@ -1,4 +1,4 @@ -main:32f8c659348c161980e59398a17f1b +main:32f8c659348c161980e59398a17f1bb0 assets:b3ae2b933e54528bf486ff28e887804d auth:f396d4bce90215599dde6ad0a8f30f29 backfills:725109470cd2613de8cc8af022fb54bc From 00a2a0ec70a963951f6b6527c34382dbd515287f Mon Sep 17 00:00:00 2001 From: pratiksha badheka Date: Tue, 15 Jul 2025 00:50:54 -0700 Subject: [PATCH 23/36] Remove Base Operation from output_main.svg file --- airflow-ctl/docs/images/output_main.svg | 32 +++++++++++-------------- 1 file changed, 14 insertions(+), 18 deletions(-) diff --git a/airflow-ctl/docs/images/output_main.svg b/airflow-ctl/docs/images/output_main.svg index 3a2b59a97fde5..5083fa17d4f75 100644 --- a/airflow-ctl/docs/images/output_main.svg +++ b/airflow-ctl/docs/images/output_main.svg @@ -111,9 +111,6 @@ - - - Command: main @@ -137,21 +134,20 @@                     environment variable/parameter or pass username and                     password.       backfills     Perform Backfills operations -      base          Perform Base operations -      config        Perform Config operations -      connections   Perform Connections operations -      dag           Perform Dag operations -      dagrun        Perform DagRun operations -      jobs          Perform Jobs operations -      pools         Perform Pools operations -      providers     Perform Providers operations -      variables     Perform Variables operations - -    Commands: -      version       Show version information - -Options: -  -h, --help        show this help message and exit +      config        Perform Config operations +      connections   Perform Connections operations +      dag           Perform Dag operations +      dagrun        Perform DagRun operations +      jobs          Perform Jobs operations +      pools         Perform Pools operations +      providers     Perform Providers operations +      variables     Perform Variables operations + +    Commands: +      version       Show version information + +Options: +  -h, --help        show this help message and exit From 838d9ba72704fe89ea92cfdff752fdcac88dae0d Mon Sep 17 00:00:00 2001 From: pratiksha badheka Date: Tue, 15 Jul 2025 01:18:23 -0700 Subject: [PATCH 24/36] remove end-space from output_main.svg --- airflow-ctl/docs/images/output_main.svg | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/airflow-ctl/docs/images/output_main.svg b/airflow-ctl/docs/images/output_main.svg index 5083fa17d4f75..a5f8703a8549e 100644 --- a/airflow-ctl/docs/images/output_main.svg +++ b/airflow-ctl/docs/images/output_main.svg @@ -1,4 +1,4 @@ - +