This is an automated email from the ASF dual-hosted git repository.

vincbeck pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/main by this push:
     new a0d563aa4b3 UI: Add team column and filter to the audit log (#71916)
a0d563aa4b3 is described below

commit a0d563aa4b391177548fa86500d34625268525d5
Author: Vincent <[email protected]>
AuthorDate: Tue Aug 25 08:44:12 2026 -0400

    UI: Add team column and filter to the audit log (#71916)
    
    When multi-team mode is enabled, operators reviewing the audit log need to
    know which team owns the resource an event was recorded against, and to 
scope
    the log to a single team when auditing that team's activity.
    
    Some audited actions own no Dag to resolve a team through — a triggerer 
started
    for one team, a team-scoped pool or variable, the team commands themselves 
— so
    the team is recorded on the event itself and read back in preference to the 
one
    derived from the Dag. Without it, those events would show no team at all 
and,
    worse, would silently disappear from the log whenever a team filter was 
applied.
    
    The column and the filter render only when the `multi_team` configuration is
    enabled, and no team data is loaded when it is off, so single-team 
deployments
    are unaffected.
---
 airflow-core/docs/migrations-ref.rst               |   4 +-
 .../api_fastapi/core_api/datamodels/event_logs.py  |   1 +
 .../core_api/openapi/v2-rest-api-generated.yaml    |  13 ++
 .../core_api/routes/public/event_logs.py           |   7 ++
 .../src/airflow/api_fastapi/logging/decorators.py  |  46 ++++++-
 .../src/airflow/jobs/scheduler_job_runner.py       |  11 +-
 .../versions/0132_3_4_0_add_team_name_to_log.py    |  50 ++++++++
 airflow-core/src/airflow/models/log.py             |  48 +++++++-
 .../src/airflow/ui/openapi-gen/queries/common.ts   |   5 +-
 .../ui/openapi-gen/queries/ensureQueryData.ts      |   6 +-
 .../src/airflow/ui/openapi-gen/queries/prefetch.ts |   6 +-
 .../src/airflow/ui/openapi-gen/queries/queries.ts  |   6 +-
 .../src/airflow/ui/openapi-gen/queries/suspense.ts |   6 +-
 .../airflow/ui/openapi-gen/requests/schemas.gen.ts |  11 ++
 .../ui/openapi-gen/requests/services.gen.ts        |   4 +-
 .../airflow/ui/openapi-gen/requests/types.gen.ts   |   2 +
 .../src/airflow/ui/src/pages/Events/Events.tsx     |  22 +++-
 .../airflow/ui/src/pages/Events/EventsFilters.tsx  |   6 +
 airflow-core/src/airflow/utils/cli.py              |   8 +-
 .../src/airflow/utils/cli_action_loggers.py        |   5 +-
 airflow-core/src/airflow/utils/db.py               |   2 +-
 .../core_api/routes/public/test_event_logs.py      |  55 +++++++++
 .../unit/api_fastapi/logging/test_decorators.py    | 137 +++++++++++++++++++++
 airflow-core/tests/unit/jobs/test_scheduler_job.py |  46 ++++++-
 airflow-core/tests/unit/models/test_log.py         |  88 +++++++++++++
 airflow-core/tests/unit/utils/test_cli_util.py     |  58 +++++++++
 .../src/airflowctl/api/datamodels/generated.py     |   1 +
 27 files changed, 632 insertions(+), 22 deletions(-)

diff --git a/airflow-core/docs/migrations-ref.rst 
b/airflow-core/docs/migrations-ref.rst
index cc18c48f493..8e03f37ef54 100644
--- a/airflow-core/docs/migrations-ref.rst
+++ b/airflow-core/docs/migrations-ref.rst
@@ -39,7 +39,9 @@ Here's the list of all the Database Migrations that are 
executed via when you ru
 
+-------------------------+------------------+-------------------+--------------------------------------------------------------+
 | Revision ID             | Revises ID       | Airflow Version   | Description 
                                                 |
 
+=========================+==================+===================+==============================================================+
-| ``c7f0a5d2e9b4`` (head) | ``76c46545c91e`` | ``3.4.0``         | Lower case 
team names.                                       |
+| ``8d3f1a6b2c47`` (head) | ``c7f0a5d2e9b4`` | ``3.4.0``         | Add 
team_name to log.                                        |
++-------------------------+------------------+-------------------+--------------------------------------------------------------+
+| ``c7f0a5d2e9b4``        | ``76c46545c91e`` | ``3.4.0``         | Lower case 
team names.                                       |
 
+-------------------------+------------------+-------------------+--------------------------------------------------------------+
 | ``76c46545c91e``        | ``3c525f44bea8`` | ``3.4.0``         | Add new 
index for trigger.                                   |
 
+-------------------------+------------------+-------------------+--------------------------------------------------------------+
diff --git 
a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/event_logs.py 
b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/event_logs.py
index 5666e46de68..0c67909add7 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/event_logs.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/event_logs.py
@@ -46,6 +46,7 @@ class EventLogResponse(BaseModel):
     task_display_name: str | None = Field(
         validation_alias=AliasPath("task_instance", "task_display_name"), 
default=None
     )
+    team_name: str | None = None
 
 
 class EventLogCollectionResponse(BaseModel):
diff --git 
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
 
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
index ce8920b702d..bd5a6ce52b0 100644
--- 
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
+++ 
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
@@ -4771,6 +4771,14 @@ paths:
           title: Event Prefix Pattern
         description: Case-sensitive, index-friendly prefix match. See 
"Filtering with
           pattern parameters".
+      - name: teams
+        in: query
+        required: false
+        schema:
+          type: array
+          items:
+            type: string
+          title: Teams
       responses:
         '200':
           description: Successful Response
@@ -14319,6 +14327,11 @@ components:
           - type: string
           - type: 'null'
           title: Task Display Name
+        team_name:
+          anyOf:
+          - type: string
+          - type: 'null'
+          title: Team Name
       type: object
       required:
       - event_log_id
diff --git 
a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/event_logs.py 
b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/event_logs.py
index eb0c4009814..524b66b8275 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/event_logs.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/event_logs.py
@@ -171,6 +171,12 @@ def get_event_logs(
         _PrefixSearchParam,
         Depends(prefix_search_param_factory(Log.event, 
"event_prefix_pattern")),
     ],
+    teams: Annotated[
+        FilterParam[list[str]],
+        Depends(
+            filter_param_factory(Log.team_name, list[str], 
FilterOptionEnum.IN, "teams", default_factory=list)
+        ),
+    ],
     readable_event_logs_filter: ReadableEventLogsFilterDep,
 ) -> EventLogCollectionResponse:
     """Get all Event Logs."""
@@ -214,6 +220,7 @@ def get_event_logs(
             owner_display_name_prefix_pattern,
             event_pattern,
             event_prefix_pattern,
+            teams,
             # Permission
             readable_event_logs_filter,
         ],
diff --git a/airflow-core/src/airflow/api_fastapi/logging/decorators.py 
b/airflow-core/src/airflow/api_fastapi/logging/decorators.py
index 22acc54b130..3ed9dd86c79 100644
--- a/airflow-core/src/airflow/api_fastapi/logging/decorators.py
+++ b/airflow-core/src/airflow/api_fastapi/logging/decorators.py
@@ -20,18 +20,32 @@ import itertools
 import json
 import logging
 from datetime import datetime
+from typing import TYPE_CHECKING
 
 import pendulum
 from fastapi import Request
 from pendulum.parsing.exceptions import ParserError
+from sqlalchemy import select
 
 from airflow._shared.secrets_masker import secrets_masker
 from airflow.api_fastapi.common.db.common import SessionDep
 from airflow.api_fastapi.core_api.security import GetUserDep
-from airflow.models import Log
+from airflow.configuration import conf
+from airflow.models import Connection, Log, Pool, Variable
+from airflow.models.team import find_invalid_team_names
+
+if TYPE_CHECKING:
+    from sqlalchemy.orm import Session
 
 logger = logging.getLogger(__name__)
 
+# Request parameter identifying a team-scoped resource, and the columns to 
read its team from.
+_TEAM_SCOPED_RESOURCES = {
+    "pool_name": (Pool.team_name, Pool.pool),
+    "variable_key": (Variable.team_name, Variable.key),
+    "connection_id": (Connection.team_name, Connection.conn_id),
+}
+
 
 def _sanitize_for_stdlib_log(value: str) -> str:
     """
@@ -138,6 +152,35 @@ def _mask_variable_entity(extra_fields):
     return result
 
 
+def _resolve_team_name(params: dict, *, session: Session) -> str | None:
+    """
+    Return the team the audited action belongs to, for the resources that own 
no Dag.
+
+    A Dag-scoped event has its team stamped from ``dag_id`` when the row is 
inserted, and a request
+    that names a team carries it directly. What is left is an action on a 
team-scoped resource that
+    names no team -- a deletion, or a patch that does not touch ``team_name`` 
-- where the team can
+    only come from the resource being acted on. It is read here rather than in 
the routes because
+    this dependency runs before the endpoint, so the row is still there to 
read even when the action
+    is about to delete it.
+    """
+    team_name = params.get("team_name")
+    if isinstance(team_name, str):
+        # The endpoint's own validation rejects a name that is too long or 
malformed, but this row
+        # is committed before that runs, so recording it would fail the insert 
on a backend that
+        # enforces the column width. The value stays visible in ``extra`` 
either way.
+        return None if find_invalid_team_names([team_name]) else team_name
+    if params.get("dag_id"):
+        # Left to the insert-time hook on ``Log``, which covers every writer 
of an audit row rather
+        # than only this one, and resolves a Dag's team through its bundle 
instead of a column.
+        return None
+    if not conf.getboolean("core", "multi_team"):
+        return None
+    for param, (team_column, resource_column) in 
_TEAM_SCOPED_RESOURCES.items():
+        if (resource_id := params.get(param)) is not None:
+            return session.scalar(select(team_column).where(resource_column == 
resource_id))
+    return None
+
+
 def action_logging(event: str | None = None):
     async def log_action(
         request: Request,
@@ -216,6 +259,7 @@ def action_logging(event: str | None = None):
             task_id=params.get("task_id"),
             dag_id=params.get("dag_id"),
             run_id=params.get("run_id") or params.get("dag_run_id"),
+            team_name=_resolve_team_name(params, session=session),
         )
 
         if "logical_date" in request.query_params:
diff --git a/airflow-core/src/airflow/jobs/scheduler_job_runner.py 
b/airflow-core/src/airflow/jobs/scheduler_job_runner.py
index 6b08dab3f02..f8b2df4aac2 100644
--- a/airflow-core/src/airflow/jobs/scheduler_job_runner.py
+++ b/airflow-core/src/airflow/jobs/scheduler_job_runner.py
@@ -104,6 +104,7 @@ from airflow.models.dagbag import CachedDBDagBag, DBDagBag
 from airflow.models.dagbundle import DagBundleModel
 from airflow.models.dagrun import DagRun
 from airflow.models.dagwarning import DagWarning, DagWarningType
+from airflow.models.log import resolve_team_name
 from airflow.models.pool import normalize_pool_name_for_stats
 from airflow.models.serialized_dag import SerializedDagModel
 from airflow.models.taskinstance import TaskInstance
@@ -1320,7 +1321,15 @@ class SchedulerJobRunner(BaseJobRunner, LoggingMixin):
 
     @staticmethod
     def _process_task_event_logs(log_records: deque[Log], session: Session):
-        objects = (log_records.popleft() for _ in range(len(log_records)))
+        objects = [log_records.popleft() for _ in range(len(log_records))]
+        # A bulk insert skips the ORM hook that would stamp this from 
``dag_id``. One flush carries
+        # a record per task event, so resolve each Dag once rather than once 
per record.
+        teams_by_dag_id = {
+            dag_id: resolve_team_name(dag_id, session=session)
+            for dag_id in {log_record.dag_id for log_record in objects}
+        }
+        for log_record in objects:
+            log_record.team_name = teams_by_dag_id[log_record.dag_id]
         session.bulk_save_objects(objects=objects, preserve_order=False)
 
     @staticmethod
diff --git 
a/airflow-core/src/airflow/migrations/versions/0132_3_4_0_add_team_name_to_log.py
 
b/airflow-core/src/airflow/migrations/versions/0132_3_4_0_add_team_name_to_log.py
new file mode 100644
index 00000000000..f7ed523b52b
--- /dev/null
+++ 
b/airflow-core/src/airflow/migrations/versions/0132_3_4_0_add_team_name_to_log.py
@@ -0,0 +1,50 @@
+#
+# 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.
+
+"""
+Add team_name to log.
+
+Revision ID: 8d3f1a6b2c47
+Revises: c7f0a5d2e9b4
+Create Date: 2026-08-20 15:02:11.409771
+
+"""
+
+from __future__ import annotations
+
+import sqlalchemy as sa
+from alembic import op
+
+# revision identifiers, used by Alembic.
+revision = "8d3f1a6b2c47"
+down_revision = "c7f0a5d2e9b4"
+branch_labels = None
+depends_on = None
+airflow_version = "3.4.0"
+
+
+def upgrade():
+    """Add team_name to log."""
+    op.add_column("log", sa.Column("team_name", sa.String(length=50), 
nullable=True))
+    op.create_index("idx_log_team_name", "log", ["team_name"], unique=False)
+
+
+def downgrade():
+    """Unapply Add team_name to log."""
+    op.drop_index("idx_log_team_name", table_name="log")
+    op.drop_column("log", "team_name")
diff --git a/airflow-core/src/airflow/models/log.py 
b/airflow-core/src/airflow/models/log.py
index 8eb1d4bbcd2..c966694ca05 100644
--- a/airflow-core/src/airflow/models/log.py
+++ b/airflow-core/src/airflow/models/log.py
@@ -20,10 +20,11 @@ from __future__ import annotations
 from datetime import datetime
 from typing import TYPE_CHECKING
 
-from sqlalchemy import Index, Integer, String, Text
-from sqlalchemy.orm import Mapped, mapped_column, relationship
+from sqlalchemy import Index, Integer, String, Text, event
+from sqlalchemy.orm import Mapped, Session, mapped_column, relationship
 
 from airflow._shared.timezones import timezone
+from airflow.configuration import conf
 from airflow.models.base import Base, StringID
 from airflow.utils.sqlalchemy import UtcDateTime
 
@@ -50,6 +51,14 @@ class Log(Base):
     owner_display_name: Mapped[str | None] = mapped_column(String(500), 
nullable=True)
     extra: Mapped[str | None] = mapped_column(Text, nullable=True)
     try_number: Mapped[int | None] = mapped_column(Integer, nullable=True)
+    # Team this event belongs to, recorded when the row is written: either 
passed in by a caller
+    # acting on a team-scoped resource that owns no Dag (a triggerer started 
with ``--team-name``,
+    # a team-scoped pool, the ``teams`` commands), or resolved from ``dag_id`` 
on insert. Stamped
+    # rather than resolved on read so an event keeps the team that owned the 
resource at the time,
+    # even if the Dag's bundle is later reassigned. Deliberately not a foreign 
key to ``team.name``,
+    # for the same reason ``dag_id`` is not one to ``dag.dag_id``: deleting a 
team must not rewrite
+    # or delete the audit trail of what was done to it.
+    team_name: Mapped[str | None] = mapped_column(String(50), nullable=True)
 
     dag_model: Mapped[DagModel | None] = relationship(
         "DagModel",
@@ -70,6 +79,7 @@ class Log(Base):
         Index("idx_log_dttm", dttm),
         Index("idx_log_event", event),
         Index("idx_log_task_instance", dag_id, task_id, run_id, map_index, 
try_number),
+        Index("idx_log_team_name", team_name),
     )
 
     def __init__(
@@ -111,9 +121,43 @@ class Log(Base):
             self.map_index = kwargs["map_index"]
         if "try_number" in kwargs:
             self.try_number = kwargs["try_number"]
+        if "team_name" in kwargs:
+            self.team_name = kwargs["team_name"]
 
         self.owner = owner or task_owner
         self.owner_display_name = owner_display_name or None
 
     def __str__(self) -> str:
         return f"Log({self.event}, {self.task_id}, {self.owner}, 
{self.owner_display_name}, {self.extra})"
+
+
+def resolve_team_name(dag_id: str | None, *, session: Session) -> str | None:
+    """
+    Return the team owning ``dag_id``, or ``None`` outside multi-team mode and 
for no Dag.
+
+    Callers that insert :class:`Log` rows without the ORM must stamp 
``team_name`` with this;
+    the ``before_insert`` hook below does it for every other write path.
+    """
+    if not dag_id or not conf.getboolean("core", "multi_team"):
+        return None
+
+    from airflow.models.dag import DagModel
+
+    return DagModel.get_team_name(dag_id, session=session)
+
+
[email protected]_for(Log, "before_insert")
+def _stamp_team_name(mapper, connection, target: Log) -> None:
+    """
+    Record the team owning the event's Dag, unless the caller already named a 
team.
+
+    Done here rather than at each of the many places that log an event, so 
that none of them --
+    including ones added later -- can leave a Dag event unattributed. Bulk 
inserts bypass ORM
+    events, so those few paths call :func:`resolve_team_name` themselves.
+    """
+    if target.team_name or not target.dag_id:
+        return
+    # ``Session(bind=connection)`` runs the lookup inside the flush's own 
transaction. Cheap in
+    # practice: ``get_team_name`` is cached per Dag, so a warm cache emits no 
query at all.
+    with Session(bind=connection) as session:
+        target.team_name = resolve_team_name(target.dag_id, session=session)
diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts 
b/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts
index f11d2d808d1..7d86b648d4a 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts
@@ -419,7 +419,7 @@ export const UseEventLogServiceGetEventLogKeyFn = ({ 
eventLogId }: {
 export type EventLogServiceGetEventLogsDefaultResponse = 
Awaited<ReturnType<typeof EventLogService.getEventLogs>>;
 export type EventLogServiceGetEventLogsQueryResult<TData = 
EventLogServiceGetEventLogsDefaultResponse, TError = unknown> = 
UseQueryResult<TData, TError>;
 export const useEventLogServiceGetEventLogsKey = "EventLogServiceGetEventLogs";
-export const UseEventLogServiceGetEventLogsKeyFn = ({ after, before, dagId, 
dagIdPattern, dagIdPrefixPattern, event, eventPattern, eventPrefixPattern, 
excludedEvents, includedEvents, limit, mapIndex, offset, orderBy, owner, 
ownerDisplayNamePattern, ownerDisplayNamePrefixPattern, ownerPattern, 
ownerPrefixPattern, runId, runIdPattern, runIdPrefixPattern, taskId, 
taskIdPattern, taskIdPrefixPattern, tryNumber }: {
+export const UseEventLogServiceGetEventLogsKeyFn = ({ after, before, dagId, 
dagIdPattern, dagIdPrefixPattern, event, eventPattern, eventPrefixPattern, 
excludedEvents, includedEvents, limit, mapIndex, offset, orderBy, owner, 
ownerDisplayNamePattern, ownerDisplayNamePrefixPattern, ownerPattern, 
ownerPrefixPattern, runId, runIdPattern, runIdPrefixPattern, taskId, 
taskIdPattern, taskIdPrefixPattern, teams, tryNumber }: {
   after?: string;
   before?: string;
   dagId?: string;
@@ -445,8 +445,9 @@ export const UseEventLogServiceGetEventLogsKeyFn = ({ 
after, before, dagId, dagI
   taskId?: string;
   taskIdPattern?: string;
   taskIdPrefixPattern?: string;
+  teams?: string[];
   tryNumber?: number;
-} = {}, queryKey?: Array<unknown>) => [useEventLogServiceGetEventLogsKey, 
...(queryKey ?? [{ after, before, dagId, dagIdPattern, dagIdPrefixPattern, 
event, eventPattern, eventPrefixPattern, excludedEvents, includedEvents, limit, 
mapIndex, offset, orderBy, owner, ownerDisplayNamePattern, 
ownerDisplayNamePrefixPattern, ownerPattern, ownerPrefixPattern, runId, 
runIdPattern, runIdPrefixPattern, taskId, taskIdPattern, taskIdPrefixPattern, 
tryNumber }])];
+} = {}, queryKey?: Array<unknown>) => [useEventLogServiceGetEventLogsKey, 
...(queryKey ?? [{ after, before, dagId, dagIdPattern, dagIdPrefixPattern, 
event, eventPattern, eventPrefixPattern, excludedEvents, includedEvents, limit, 
mapIndex, offset, orderBy, owner, ownerDisplayNamePattern, 
ownerDisplayNamePrefixPattern, ownerPattern, ownerPrefixPattern, runId, 
runIdPattern, runIdPrefixPattern, taskId, taskIdPattern, taskIdPrefixPattern, 
teams, tryNumber }])];
 export type ExtraLinksServiceGetExtraLinksDefaultResponse = 
Awaited<ReturnType<typeof ExtraLinksService.getExtraLinks>>;
 export type ExtraLinksServiceGetExtraLinksQueryResult<TData = 
ExtraLinksServiceGetExtraLinksDefaultResponse, TError = unknown> = 
UseQueryResult<TData, TError>;
 export const useExtraLinksServiceGetExtraLinksKey = 
"ExtraLinksServiceGetExtraLinks";
diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts 
b/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts
index 8f7ed508d3a..40fb19af300 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts
@@ -825,10 +825,11 @@ export const ensureUseEventLogServiceGetEventLogData = 
(queryClient: QueryClient
 * @param data.ownerPrefixPattern Case-sensitive, index-friendly prefix match. 
See "Filtering with pattern parameters".
 * @param data.ownerDisplayNamePrefixPattern Case-sensitive, index-friendly 
prefix match. See "Filtering with pattern parameters".
 * @param data.eventPrefixPattern Case-sensitive, index-friendly prefix match. 
See "Filtering with pattern parameters".
+* @param data.teams
 * @returns EventLogCollectionResponse Successful Response
 * @throws ApiError
 */
-export const ensureUseEventLogServiceGetEventLogsData = (queryClient: 
QueryClient, { after, before, dagId, dagIdPattern, dagIdPrefixPattern, event, 
eventPattern, eventPrefixPattern, excludedEvents, includedEvents, limit, 
mapIndex, offset, orderBy, owner, ownerDisplayNamePattern, 
ownerDisplayNamePrefixPattern, ownerPattern, ownerPrefixPattern, runId, 
runIdPattern, runIdPrefixPattern, taskId, taskIdPattern, taskIdPrefixPattern, 
tryNumber }: {
+export const ensureUseEventLogServiceGetEventLogsData = (queryClient: 
QueryClient, { after, before, dagId, dagIdPattern, dagIdPrefixPattern, event, 
eventPattern, eventPrefixPattern, excludedEvents, includedEvents, limit, 
mapIndex, offset, orderBy, owner, ownerDisplayNamePattern, 
ownerDisplayNamePrefixPattern, ownerPattern, ownerPrefixPattern, runId, 
runIdPattern, runIdPrefixPattern, taskId, taskIdPattern, taskIdPrefixPattern, 
teams, tryNumber }: {
   after?: string;
   before?: string;
   dagId?: string;
@@ -854,8 +855,9 @@ export const ensureUseEventLogServiceGetEventLogsData = 
(queryClient: QueryClien
   taskId?: string;
   taskIdPattern?: string;
   taskIdPrefixPattern?: string;
+  teams?: string[];
   tryNumber?: number;
-} = {}) => queryClient.ensureQueryData({ queryKey: 
Common.UseEventLogServiceGetEventLogsKeyFn({ after, before, dagId, 
dagIdPattern, dagIdPrefixPattern, event, eventPattern, eventPrefixPattern, 
excludedEvents, includedEvents, limit, mapIndex, offset, orderBy, owner, 
ownerDisplayNamePattern, ownerDisplayNamePrefixPattern, ownerPattern, 
ownerPrefixPattern, runId, runIdPattern, runIdPrefixPattern, taskId, 
taskIdPattern, taskIdPrefixPattern, tryNumber }), queryFn: () => 
EventLogService.getEve [...]
+} = {}) => queryClient.ensureQueryData({ queryKey: 
Common.UseEventLogServiceGetEventLogsKeyFn({ after, before, dagId, 
dagIdPattern, dagIdPrefixPattern, event, eventPattern, eventPrefixPattern, 
excludedEvents, includedEvents, limit, mapIndex, offset, orderBy, owner, 
ownerDisplayNamePattern, ownerDisplayNamePrefixPattern, ownerPattern, 
ownerPrefixPattern, runId, runIdPattern, runIdPrefixPattern, taskId, 
taskIdPattern, taskIdPrefixPattern, teams, tryNumber }), queryFn: () => 
EventLogService [...]
 /**
 * Get Extra Links
 * Get extra links for task instance.
diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts 
b/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts
index fb50dcc1247..431b7480b0d 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts
@@ -825,10 +825,11 @@ export const prefetchUseEventLogServiceGetEventLog = 
(queryClient: QueryClient,
 * @param data.ownerPrefixPattern Case-sensitive, index-friendly prefix match. 
See "Filtering with pattern parameters".
 * @param data.ownerDisplayNamePrefixPattern Case-sensitive, index-friendly 
prefix match. See "Filtering with pattern parameters".
 * @param data.eventPrefixPattern Case-sensitive, index-friendly prefix match. 
See "Filtering with pattern parameters".
+* @param data.teams
 * @returns EventLogCollectionResponse Successful Response
 * @throws ApiError
 */
-export const prefetchUseEventLogServiceGetEventLogs = (queryClient: 
QueryClient, { after, before, dagId, dagIdPattern, dagIdPrefixPattern, event, 
eventPattern, eventPrefixPattern, excludedEvents, includedEvents, limit, 
mapIndex, offset, orderBy, owner, ownerDisplayNamePattern, 
ownerDisplayNamePrefixPattern, ownerPattern, ownerPrefixPattern, runId, 
runIdPattern, runIdPrefixPattern, taskId, taskIdPattern, taskIdPrefixPattern, 
tryNumber }: {
+export const prefetchUseEventLogServiceGetEventLogs = (queryClient: 
QueryClient, { after, before, dagId, dagIdPattern, dagIdPrefixPattern, event, 
eventPattern, eventPrefixPattern, excludedEvents, includedEvents, limit, 
mapIndex, offset, orderBy, owner, ownerDisplayNamePattern, 
ownerDisplayNamePrefixPattern, ownerPattern, ownerPrefixPattern, runId, 
runIdPattern, runIdPrefixPattern, taskId, taskIdPattern, taskIdPrefixPattern, 
teams, tryNumber }: {
   after?: string;
   before?: string;
   dagId?: string;
@@ -854,8 +855,9 @@ export const prefetchUseEventLogServiceGetEventLogs = 
(queryClient: QueryClient,
   taskId?: string;
   taskIdPattern?: string;
   taskIdPrefixPattern?: string;
+  teams?: string[];
   tryNumber?: number;
-} = {}) => queryClient.prefetchQuery({ queryKey: 
Common.UseEventLogServiceGetEventLogsKeyFn({ after, before, dagId, 
dagIdPattern, dagIdPrefixPattern, event, eventPattern, eventPrefixPattern, 
excludedEvents, includedEvents, limit, mapIndex, offset, orderBy, owner, 
ownerDisplayNamePattern, ownerDisplayNamePrefixPattern, ownerPattern, 
ownerPrefixPattern, runId, runIdPattern, runIdPrefixPattern, taskId, 
taskIdPattern, taskIdPrefixPattern, tryNumber }), queryFn: () => 
EventLogService.getEvent [...]
+} = {}) => queryClient.prefetchQuery({ queryKey: 
Common.UseEventLogServiceGetEventLogsKeyFn({ after, before, dagId, 
dagIdPattern, dagIdPrefixPattern, event, eventPattern, eventPrefixPattern, 
excludedEvents, includedEvents, limit, mapIndex, offset, orderBy, owner, 
ownerDisplayNamePattern, ownerDisplayNamePrefixPattern, ownerPattern, 
ownerPrefixPattern, runId, runIdPattern, runIdPrefixPattern, taskId, 
taskIdPattern, taskIdPrefixPattern, teams, tryNumber }), queryFn: () => 
EventLogService.g [...]
 /**
 * Get Extra Links
 * Get extra links for task instance.
diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts 
b/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
index f1d64923f48..bd447ebeed2 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
@@ -825,10 +825,11 @@ export const useEventLogServiceGetEventLog = <TData = 
Common.EventLogServiceGetE
 * @param data.ownerPrefixPattern Case-sensitive, index-friendly prefix match. 
See "Filtering with pattern parameters".
 * @param data.ownerDisplayNamePrefixPattern Case-sensitive, index-friendly 
prefix match. See "Filtering with pattern parameters".
 * @param data.eventPrefixPattern Case-sensitive, index-friendly prefix match. 
See "Filtering with pattern parameters".
+* @param data.teams
 * @returns EventLogCollectionResponse Successful Response
 * @throws ApiError
 */
-export const useEventLogServiceGetEventLogs = <TData = 
Common.EventLogServiceGetEventLogsDefaultResponse, TError = unknown, TQueryKey 
extends Array<unknown> = unknown[]>({ after, before, dagId, dagIdPattern, 
dagIdPrefixPattern, event, eventPattern, eventPrefixPattern, excludedEvents, 
includedEvents, limit, mapIndex, offset, orderBy, owner, 
ownerDisplayNamePattern, ownerDisplayNamePrefixPattern, ownerPattern, 
ownerPrefixPattern, runId, runIdPattern, runIdPrefixPattern, taskId, 
taskIdPatte [...]
+export const useEventLogServiceGetEventLogs = <TData = 
Common.EventLogServiceGetEventLogsDefaultResponse, TError = unknown, TQueryKey 
extends Array<unknown> = unknown[]>({ after, before, dagId, dagIdPattern, 
dagIdPrefixPattern, event, eventPattern, eventPrefixPattern, excludedEvents, 
includedEvents, limit, mapIndex, offset, orderBy, owner, 
ownerDisplayNamePattern, ownerDisplayNamePrefixPattern, ownerPattern, 
ownerPrefixPattern, runId, runIdPattern, runIdPrefixPattern, taskId, 
taskIdPatte [...]
   after?: string;
   before?: string;
   dagId?: string;
@@ -854,8 +855,9 @@ export const useEventLogServiceGetEventLogs = <TData = 
Common.EventLogServiceGet
   taskId?: string;
   taskIdPattern?: string;
   taskIdPrefixPattern?: string;
+  teams?: string[];
   tryNumber?: number;
-} = {}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>, 
"queryKey" | "queryFn">) => useQuery<TData, TError>({ queryKey: 
Common.UseEventLogServiceGetEventLogsKeyFn({ after, before, dagId, 
dagIdPattern, dagIdPrefixPattern, event, eventPattern, eventPrefixPattern, 
excludedEvents, includedEvents, limit, mapIndex, offset, orderBy, owner, 
ownerDisplayNamePattern, ownerDisplayNamePrefixPattern, ownerPattern, 
ownerPrefixPattern, runId, runIdPattern, runIdPrefixPattern, taskI [...]
+} = {}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>, 
"queryKey" | "queryFn">) => useQuery<TData, TError>({ queryKey: 
Common.UseEventLogServiceGetEventLogsKeyFn({ after, before, dagId, 
dagIdPattern, dagIdPrefixPattern, event, eventPattern, eventPrefixPattern, 
excludedEvents, includedEvents, limit, mapIndex, offset, orderBy, owner, 
ownerDisplayNamePattern, ownerDisplayNamePrefixPattern, ownerPattern, 
ownerPrefixPattern, runId, runIdPattern, runIdPrefixPattern, taskI [...]
 /**
 * Get Extra Links
 * Get extra links for task instance.
diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts 
b/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts
index 1c0bb7d0033..27f4bdce056 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts
@@ -825,10 +825,11 @@ export const useEventLogServiceGetEventLogSuspense = 
<TData = Common.EventLogSer
 * @param data.ownerPrefixPattern Case-sensitive, index-friendly prefix match. 
See "Filtering with pattern parameters".
 * @param data.ownerDisplayNamePrefixPattern Case-sensitive, index-friendly 
prefix match. See "Filtering with pattern parameters".
 * @param data.eventPrefixPattern Case-sensitive, index-friendly prefix match. 
See "Filtering with pattern parameters".
+* @param data.teams
 * @returns EventLogCollectionResponse Successful Response
 * @throws ApiError
 */
-export const useEventLogServiceGetEventLogsSuspense = <TData = 
Common.EventLogServiceGetEventLogsDefaultResponse, TError = unknown, TQueryKey 
extends Array<unknown> = unknown[]>({ after, before, dagId, dagIdPattern, 
dagIdPrefixPattern, event, eventPattern, eventPrefixPattern, excludedEvents, 
includedEvents, limit, mapIndex, offset, orderBy, owner, 
ownerDisplayNamePattern, ownerDisplayNamePrefixPattern, ownerPattern, 
ownerPrefixPattern, runId, runIdPattern, runIdPrefixPattern, taskId, tas [...]
+export const useEventLogServiceGetEventLogsSuspense = <TData = 
Common.EventLogServiceGetEventLogsDefaultResponse, TError = unknown, TQueryKey 
extends Array<unknown> = unknown[]>({ after, before, dagId, dagIdPattern, 
dagIdPrefixPattern, event, eventPattern, eventPrefixPattern, excludedEvents, 
includedEvents, limit, mapIndex, offset, orderBy, owner, 
ownerDisplayNamePattern, ownerDisplayNamePrefixPattern, ownerPattern, 
ownerPrefixPattern, runId, runIdPattern, runIdPrefixPattern, taskId, tas [...]
   after?: string;
   before?: string;
   dagId?: string;
@@ -854,8 +855,9 @@ export const useEventLogServiceGetEventLogsSuspense = 
<TData = Common.EventLogSe
   taskId?: string;
   taskIdPattern?: string;
   taskIdPrefixPattern?: string;
+  teams?: string[];
   tryNumber?: number;
-} = {}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>, 
"queryKey" | "queryFn">) => useSuspenseQuery<TData, TError>({ queryKey: 
Common.UseEventLogServiceGetEventLogsKeyFn({ after, before, dagId, 
dagIdPattern, dagIdPrefixPattern, event, eventPattern, eventPrefixPattern, 
excludedEvents, includedEvents, limit, mapIndex, offset, orderBy, owner, 
ownerDisplayNamePattern, ownerDisplayNamePrefixPattern, ownerPattern, 
ownerPrefixPattern, runId, runIdPattern, runIdPrefixPatter [...]
+} = {}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>, 
"queryKey" | "queryFn">) => useSuspenseQuery<TData, TError>({ queryKey: 
Common.UseEventLogServiceGetEventLogsKeyFn({ after, before, dagId, 
dagIdPattern, dagIdPrefixPattern, event, eventPattern, eventPrefixPattern, 
excludedEvents, includedEvents, limit, mapIndex, offset, orderBy, owner, 
ownerDisplayNamePattern, ownerDisplayNamePrefixPattern, ownerPattern, 
ownerPrefixPattern, runId, runIdPattern, runIdPrefixPatter [...]
 /**
 * Get Extra Links
 * Get extra links for task instance.
diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts 
b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
index 824a7f4d4b7..ca7eee96aee 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
@@ -4959,6 +4959,17 @@ export const $EventLogResponse = {
                 }
             ],
             title: 'Task Display Name'
+        },
+        team_name: {
+            anyOf: [
+                {
+                    type: 'string'
+                },
+                {
+                    type: 'null'
+                }
+            ],
+            title: 'Team Name'
         }
     },
     type: 'object',
diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts 
b/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts
index 18f475eb78d..9b8bc8f5b07 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts
@@ -2170,6 +2170,7 @@ export class EventLogService {
      * @param data.ownerPrefixPattern Case-sensitive, index-friendly prefix 
match. See "Filtering with pattern parameters".
      * @param data.ownerDisplayNamePrefixPattern Case-sensitive, 
index-friendly prefix match. See "Filtering with pattern parameters".
      * @param data.eventPrefixPattern Case-sensitive, index-friendly prefix 
match. See "Filtering with pattern parameters".
+     * @param data.teams
      * @returns EventLogCollectionResponse Successful Response
      * @throws ApiError
      */
@@ -2203,7 +2204,8 @@ export class EventLogService {
                 run_id_prefix_pattern: data.runIdPrefixPattern,
                 owner_prefix_pattern: data.ownerPrefixPattern,
                 owner_display_name_prefix_pattern: 
data.ownerDisplayNamePrefixPattern,
-                event_prefix_pattern: data.eventPrefixPattern
+                event_prefix_pattern: data.eventPrefixPattern,
+                teams: data.teams
             },
             errors: {
                 401: 'Unauthorized',
diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts 
b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
index c060e08aa19..91336b9c42d 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
@@ -1295,6 +1295,7 @@ export type EventLogResponse = {
     extra: string | null;
     dag_display_name?: string | null;
     task_display_name?: string | null;
+    team_name?: string | null;
 };
 
 /**
@@ -3735,6 +3736,7 @@ export type GetEventLogsData = {
      * Case-sensitive, index-friendly prefix match. See "Filtering with 
pattern parameters".
      */
     taskIdPrefixPattern?: string | null;
+    teams?: Array<(string)>;
     tryNumber?: number | null;
 };
 
diff --git a/airflow-core/src/airflow/ui/src/pages/Events/Events.tsx 
b/airflow-core/src/airflow/ui/src/pages/Events/Events.tsx
index 2f5abef3f40..38ae6ed0d04 100644
--- a/airflow-core/src/airflow/ui/src/pages/Events/Events.tsx
+++ b/airflow-core/src/airflow/ui/src/pages/Events/Events.tsx
@@ -32,19 +32,21 @@ import RenderedJsonField from 
"src/components/RenderedJsonField";
 import Time from "src/components/Time";
 import { SearchParamsKeys, type SearchParamsKeysType } from 
"src/constants/searchParams";
 import { useAdvancedSearchArg } from "src/hooks/useAdvancedSearch";
+import { useConfig } from "src/queries/useConfig";
 import { useDocumentTitle } from "src/utils";
 
 import { EventsFilters } from "./EventsFilters";
 
 type EventsColumn = {
   dagId?: string;
+  multiTeam: boolean;
   open?: boolean;
   runId?: string;
   taskId?: string;
 };
 
 const eventsColumn = (
-  { dagId, open, runId, taskId }: EventsColumn,
+  { dagId, multiTeam, open, runId, taskId }: EventsColumn,
   translate: (key: string) => string,
 ): Array<ColumnDef<EventLogResponse>> => [
   {
@@ -73,6 +75,18 @@ const eventsColumn = (
       skeletonWidth: 10,
     },
   },
+  ...(multiTeam
+    ? [
+        {
+          accessorKey: "team_name",
+          enableSorting: false,
+          header: translate("common:dagDetails.team"),
+          meta: {
+            skeletonWidth: 10,
+          },
+        },
+      ]
+    : []),
   {
     accessorKey: "extra",
     cell: ({ row: { original } }) => {
@@ -156,6 +170,7 @@ const {
   MAP_INDEX: MAP_INDEX_PARAM,
   RUN_ID: RUN_ID_PARAM,
   TASK_ID: TASK_ID_PARAM,
+  TEAMS: TEAMS_PARAM,
   TRY_NUMBER: TRY_NUMBER_PARAM,
   USER: USER_PARAM,
 }: SearchParamsKeysType = SearchParamsKeys;
@@ -163,6 +178,7 @@ const {
 export const Events = () => {
   const { t: translate } = useTranslation(["browse", "common"]);
   const { dagId, runId, taskId } = useParams();
+  const multiTeamEnabled = Boolean(useConfig("multi_team"));
 
   // Only the standalone audit-log page owns the tab title; nested tabs 
inherit their parent page's title.
   useDocumentTitle(dagId === undefined ? translate("common:browse.auditLog") : 
undefined);
@@ -182,6 +198,7 @@ export const Events = () => {
   const taskIdFilter = searchParams.get(TASK_ID_PARAM);
   const tryNumberFilter = searchParams.get(TRY_NUMBER_PARAM);
   const userFilter = searchParams.get(USER_PARAM);
+  const teams = searchParams.getAll(TEAMS_PARAM);
 
   const orderBy = sort ? [`${sort.desc ? "-" : ""}${sort.id}`] : ["-when"];
   // Convert string filters to appropriate types for API
@@ -240,13 +257,14 @@ export const Events = () => {
       ...runIdArg,
       taskId: taskId ?? undefined,
       ...taskIdArg,
+      teams: teams.length > 0 ? teams : undefined,
       tryNumber: tryNumberNumber,
     },
     undefined,
   );
 
   const eventLogs = data?.event_logs ?? [];
-  const columns = eventsColumn({ dagId, open, runId, taskId }, translate);
+  const columns = eventsColumn({ dagId, multiTeam: multiTeamEnabled, open, 
runId, taskId }, translate);
 
   return (
     <>
diff --git a/airflow-core/src/airflow/ui/src/pages/Events/EventsFilters.tsx 
b/airflow-core/src/airflow/ui/src/pages/Events/EventsFilters.tsx
index 1352fa38f8f..61c9cb58855 100644
--- a/airflow-core/src/airflow/ui/src/pages/Events/EventsFilters.tsx
+++ b/airflow-core/src/airflow/ui/src/pages/Events/EventsFilters.tsx
@@ -18,6 +18,7 @@
  */
 import { FilterBar } from "src/components/FilterBar";
 import { SearchParamsKeys } from "src/constants/searchParams";
+import { useConfig } from "src/queries/useConfig";
 import { useFiltersHandler, type FilterableSearchParamsKeys } from "src/utils";
 
 type EventsFiltersProps = {
@@ -27,6 +28,7 @@ type EventsFiltersProps = {
 };
 
 export const EventsFilters = ({ urlDagId, urlRunId, urlTaskId }: 
EventsFiltersProps) => {
+  const multiTeamEnabled = Boolean(useConfig("multi_team"));
   const searchParamKeys: Array<FilterableSearchParamsKeys> = [
     SearchParamsKeys.EVENT_DATE_RANGE,
     SearchParamsKeys.EVENT_TYPE,
@@ -35,6 +37,10 @@ export const EventsFilters = ({ urlDagId, urlRunId, 
urlTaskId }: EventsFiltersPr
     SearchParamsKeys.TRY_NUMBER,
   ];
 
+  if (multiTeamEnabled) {
+    searchParamKeys.push(SearchParamsKeys.TEAMS);
+  }
+
   // Only add Dag ID filter if not in URL context
   if (urlDagId === undefined) {
     searchParamKeys.push(SearchParamsKeys.DAG_ID);
diff --git a/airflow-core/src/airflow/utils/cli.py 
b/airflow-core/src/airflow/utils/cli.py
index d0105e86efa..e3491bf91f1 100644
--- a/airflow-core/src/airflow/utils/cli.py
+++ b/airflow-core/src/airflow/utils/cli.py
@@ -77,6 +77,7 @@ def action_cli(func=None, check_db=True):
             dag_id : dag id (optional)
             task_id : task_id (optional)
             logical_date : logical date (optional)
+            team_name : team the command is scoped to (optional)
             error : exception instance if there's an exception
 
         :param f: function instance
@@ -142,7 +143,7 @@ def _build_metrics(func_name, namespace):
 
     It assumes that function arguments is from airflow.bin.cli module's 
function
     and has Namespace instance where it optionally contains "dag_id", 
"task_id",
-    and "logical_date".
+    "logical_date", and a team name.
 
     :param func_name: name of function
     :param namespace: Namespace instance from argparse
@@ -222,6 +223,11 @@ def _build_metrics(func_name, namespace):
     metrics["dag_id"] = tmp_dic.get("dag_id")
     metrics["task_id"] = tmp_dic.get("task_id")
     metrics["logical_date"] = tmp_dic.get("logical_date")
+    # The ``teams`` commands take the team they act on as a positional 
``name``; everything else
+    # scoped to a team (``triggerer``, ``pools set``) takes it as 
``--team-name``.
+    metrics["team_name"] = tmp_dic.get("team_name") or (
+        tmp_dic.get("name") if func_name.startswith("team_") else None
+    )
     metrics["host_name"] = socket.gethostname()
 
     return metrics
diff --git a/airflow-core/src/airflow/utils/cli_action_loggers.py 
b/airflow-core/src/airflow/utils/cli_action_loggers.py
index 632430e8a93..f1f5877c72a 100644
--- a/airflow-core/src/airflow/utils/cli_action_loggers.py
+++ b/airflow-core/src/airflow/utils/cli_action_loggers.py
@@ -112,6 +112,7 @@ def default_action_log(
     logical_date,
     host_name,
     full_command,
+    team_name=None,
     *,
     session: Session = NEW_SESSION,
     **_,
@@ -125,7 +126,7 @@ def default_action_log(
     from sqlalchemy.exc import OperationalError, ProgrammingError
 
     from airflow._shared.timezones import timezone
-    from airflow.models.log import Log
+    from airflow.models.log import Log, resolve_team_name
 
     try:
         # Use bulk_insert_mappings here to avoid importing all models (which 
using the classes does) early
@@ -141,6 +142,8 @@ def default_action_log(
                     "task_id": task_id,
                     "dag_id": dag_id,
                     "logical_date": logical_date,
+                    # A bulk insert skips the ORM hook that would stamp this 
from ``dag_id``.
+                    "team_name": team_name or resolve_team_name(dag_id, 
session=session),
                     "dttm": timezone.utcnow(),
                 }
             ],
diff --git a/airflow-core/src/airflow/utils/db.py 
b/airflow-core/src/airflow/utils/db.py
index 19f96a2ea6e..71ec2d2004b 100644
--- a/airflow-core/src/airflow/utils/db.py
+++ b/airflow-core/src/airflow/utils/db.py
@@ -117,7 +117,7 @@ _REVISION_HEADS_MAP: dict[str, str] = {
     "3.1.8": "509b94a1042d",
     "3.2.0": "1d6611b6ab7c",
     "3.3.0": "d2f4e1b3c5a7",
-    "3.4.0": "c7f0a5d2e9b4",
+    "3.4.0": "8d3f1a6b2c47",
 }
 
 # Prefix used to identify tables holding data moved during migration.
diff --git 
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_event_logs.py 
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_event_logs.py
index 75fd6229d5d..240c95a46a2 100644
--- 
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_event_logs.py
+++ 
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_event_logs.py
@@ -54,6 +54,8 @@ EVENT_WITH_TASK_INSTANCE = "EVENT_WITH_TASK_INSTANCE"
 EVENT_WITH_OWNER_AND_TASK_INSTANCE = "EVENT_WITH_OWNER_AND_TASK_INSTANCE"
 EVENT_WITHOUT_DTTM = "EVENT_WITHOUT_DTTM"
 EVENT_NON_EXISTED_ID = 9999
+TEAM_EVENT = "TEAM_EVENT"
+TEAM_NAME = "TEST_TEAM"
 
 
 class TestEventLogsEndpoint:
@@ -192,10 +194,21 @@ class TestGetEventLog(TestEventLogsEndpoint):
             "owner": expected_body.get("owner"),
             "owner_display_name": expected_body.get("owner_display_name"),
             "extra": expected_body.get("extra"),
+            "team_name": None,
         }
 
         assert response.json() == expected_json
 
+    def test_get_event_log_returns_the_recorded_team(self, test_client, 
session):
+        event_log = Log(event="cli_triggerer", team_name=TEAM_NAME)
+        session.add(event_log)
+        session.commit()
+
+        response = test_client.get(f"/eventLogs/{event_log.id}")
+
+        assert response.status_code == 200
+        assert response.json()["team_name"] == TEAM_NAME
+
     def test_should_raises_401_unauthenticated(self, 
unauthenticated_test_client, setup):
         event_log_id = setup[EVENT_NORMAL].id
         response = 
unauthenticated_test_client.get(f"/eventLogs/{event_log_id}")
@@ -450,6 +463,48 @@ class TestGetEventLogs(TestEventLogsEndpoint):
         assert event_log["owner"] == OWNER_AIRFLOW
         assert event_log["owner_display_name"] == OWNER_AIRFLOW
 
+    def test_get_event_logs_returns_the_recorded_team(self, test_client, 
session):
+        session.add(Log(event=TEAM_EVENT, dag_id=DAG_ID, team_name=TEAM_NAME))
+        session.commit()
+
+        with assert_queries_count(4):
+            response = test_client.get("/eventLogs")
+
+        assert response.status_code == 200
+        teams_by_event = {
+            event_log["event"]: event_log["team_name"] for event_log in 
response.json()["event_logs"]
+        }
+        assert teams_by_event == {
+            EVENT_NORMAL: None,
+            EVENT_WITH_OWNER: None,
+            TASK_INSTANCE_EVENT: None,
+            EVENT_WITH_OWNER_AND_TASK_INSTANCE: None,
+            TEAM_EVENT: TEAM_NAME,
+        }
+
+    def test_get_event_logs_filtered_by_team(self, test_client, session):
+        session.add_all(
+            [
+                Log(event=TEAM_EVENT, dag_id=DAG_ID, team_name=TEAM_NAME),
+                Log(event="cli_triggerer", team_name=TEAM_NAME),
+                Log(event="cli_triggerer", team_name="other-team"),
+            ]
+        )
+        session.commit()
+
+        with assert_queries_count(4):
+            response = test_client.get("/eventLogs", params={"teams": 
[TEAM_NAME]})
+
+        assert response.status_code == 200
+        body = response.json()
+        assert body["total_entries"] == 2
+        assert {event_log["event"] for event_log in body["event_logs"]} == 
{TEAM_EVENT, "cli_triggerer"}
+
+        # A team no event was recorded for returns nothing.
+        response = test_client.get("/eventLogs", params={"teams": 
["nonexistent-team"]})
+        assert response.status_code == 200
+        assert response.json()["total_entries"] == 0
+
     def test_get_event_logs_filters_by_owner_display_name_pattern(self, 
test_client):
         response = test_client.get("/eventLogs", 
params={"owner_display_name_pattern": "est Own"})
 
diff --git a/airflow-core/tests/unit/api_fastapi/logging/test_decorators.py 
b/airflow-core/tests/unit/api_fastapi/logging/test_decorators.py
index ea356b87036..e1b97ebcd2d 100644
--- a/airflow-core/tests/unit/api_fastapi/logging/test_decorators.py
+++ b/airflow-core/tests/unit/api_fastapi/logging/test_decorators.py
@@ -22,6 +22,7 @@ from unittest.mock import MagicMock
 
 import pytest
 from fastapi import Request
+from sqlalchemy import select
 from sqlalchemy.orm import Session
 
 from airflow.api_fastapi.auth.managers.models.base_user import BaseUser
@@ -31,6 +32,21 @@ from airflow.api_fastapi.logging.decorators import (
     _sanitize_for_stdlib_log,
     action_logging,
 )
+from airflow.models import Connection, Log, Pool, Variable
+from airflow.models.dag import DagModel, clear_team_name_cache
+from airflow.models.dagbundle import DagBundleModel
+from airflow.models.team import Team
+
+from tests_common.test_utils.config import conf_vars
+from tests_common.test_utils.db import (
+    clear_db_connections,
+    clear_db_dag_bundles,
+    clear_db_dags,
+    clear_db_logs,
+    clear_db_pools,
+    clear_db_teams,
+    clear_db_variables,
+)
 
 
 class TestSanitizeForStdlibLog:
@@ -303,3 +319,124 @@ class TestActionLoggingUserFields:
         (logged,) = session.add.call_args.args
         assert logged.owner == "jdoe"
         assert logged.owner_display_name == "Jane Doe"
+
+
+class TestActionLoggingTeamName:
+    """A team-scoped resource that owns no Dag records its team on the Log row 
itself, since
+    there is no ``dag_id`` to resolve one through."""
+
+    @pytest.mark.parametrize(
+        ("query_string", "expected_team_name"),
+        [
+            pytest.param(b"team_name=payments", "payments", id="team-scoped"),
+            pytest.param(b"", None, id="not-team-scoped"),
+        ],
+    )
+    def test_team_name_is_taken_from_the_request(self, query_string, 
expected_team_name):
+        request = Request({"type": "http", "method": "POST", "headers": [], 
"query_string": query_string})
+        session = MagicMock(spec=Session)
+
+        asyncio.run(action_logging(event="post_pool")(request=request, 
session=session, user=None))
+
+        (logged,) = session.add.call_args.args
+        assert logged.team_name == expected_team_name
+
+    @pytest.mark.parametrize(
+        "query_string",
+        [
+            pytest.param(b"team_name=" + b"a" * 60, 
id="longer-than-the-column"),
+            pytest.param(b"team_name=Payments", 
id="not-a-name-a-team-can-have"),
+        ],
+    )
+    def test_a_name_no_team_can_have_is_not_recorded(self, query_string):
+        """The row is committed before the endpoint validates the request, so 
a name the column
+        cannot hold would fail the insert instead of the request."""
+        request = Request({"type": "http", "method": "POST", "headers": [], 
"query_string": query_string})
+        session = MagicMock(spec=Session)
+
+        asyncio.run(action_logging(event="post_pool")(request=request, 
session=session, user=None))
+
+        (logged,) = session.add.call_args.args
+        assert logged.team_name is None
+
+
[email protected]_test
+class TestActionLoggingResourceTeamName:
+    """An action that names no team -- a deletion, or a patch leaving 
``team_name`` out -- takes the
+    team from the resource it acts on, which this dependency can still read 
because it runs before
+    the endpoint."""
+
+    def teardown_method(self):
+        clear_db_logs()
+        clear_db_pools()
+        clear_db_variables()
+        clear_db_connections()
+        clear_db_dags()
+        clear_db_dag_bundles()
+        clear_db_teams()
+
+    @staticmethod
+    def _create_resources_owned_by_a_team(session, team_name="payments"):
+        session.add(Team(name=team_name))
+        session.flush()
+        session.add(Pool(pool="team_pool", slots=1, include_deferred=False, 
team_name=team_name))
+        session.add(Variable(key="team_var", val="something", 
team_name=team_name))
+        session.add(Connection(conn_id="team_conn", conn_type="http", 
team_name=team_name))
+        session.commit()
+
+    @staticmethod
+    def _log_action(session, path_params):
+        request = Request(
+            {
+                "type": "http",
+                "method": "DELETE",
+                "headers": [],
+                "query_string": b"",
+                "path_params": path_params,
+            }
+        )
+        asyncio.run(action_logging(event="delete_resource")(request=request, 
session=session, user=None))
+        return session.scalar(select(Log).order_by(Log.id.desc()))
+
+    @pytest.mark.parametrize(
+        "path_params",
+        [
+            pytest.param({"pool_name": "team_pool"}, id="pool"),
+            pytest.param({"variable_key": "team_var"}, id="variable"),
+            pytest.param({"connection_id": "team_conn"}, id="connection"),
+        ],
+    )
+    @conf_vars({("core", "multi_team"): "True"})
+    def test_team_comes_from_the_resource_being_acted_on(self, session, 
path_params):
+        self._create_resources_owned_by_a_team(session)
+
+        assert self._log_action(session, path_params).team_name == "payments"
+
+    @conf_vars({("core", "multi_team"): "True"})
+    def test_no_team_is_recorded_for_a_resource_owned_by_none(self, session):
+        session.add(Pool(pool="team_pool", slots=1, include_deferred=False))
+        session.commit()
+
+        assert self._log_action(session, {"pool_name": "team_pool"}).team_name 
is None
+
+    def test_the_resource_is_not_read_when_multi_team_is_off(self, session):
+        self._create_resources_owned_by_a_team(session)
+
+        assert self._log_action(session, {"pool_name": "team_pool"}).team_name 
is None
+
+    @conf_vars({("core", "multi_team"): "True"})
+    def test_a_dag_scoped_event_keeps_the_team_of_its_dag(self, session):
+        """The Dag's team is stamped when the row is inserted, so a resource 
named alongside it
+        must not take precedence."""
+        self._create_resources_owned_by_a_team(session)
+        bundle = DagBundleModel(name="team-bundle")
+        bundle.teams.append(Team(name="infra"))
+        session.add(bundle)
+        session.flush()
+        session.add(DagModel(dag_id="dag_owned_by_infra", 
bundle_name="team-bundle", is_stale=False))
+        session.commit()
+        clear_team_name_cache()
+
+        log = self._log_action(session, {"dag_id": "dag_owned_by_infra", 
"pool_name": "team_pool"})
+
+        assert log.team_name == "infra"
diff --git a/airflow-core/tests/unit/jobs/test_scheduler_job.py 
b/airflow-core/tests/unit/jobs/test_scheduler_job.py
index 9f14dee0b78..a7b332e763f 100644
--- a/airflow-core/tests/unit/jobs/test_scheduler_job.py
+++ b/airflow-core/tests/unit/jobs/test_scheduler_job.py
@@ -92,7 +92,7 @@ from airflow.models.db_callback_request import 
DbCallbackRequest
 from airflow.models.deadline import Deadline
 from airflow.models.deadline_alert import DeadlineAlert
 from airflow.models.hitl import HITLDetail
-from airflow.models.log import Log
+from airflow.models.log import Log, resolve_team_name
 from airflow.models.pool import Pool
 from airflow.models.serialized_dag import SerializedDagModel
 from airflow.models.taskinstance import TaskInstance
@@ -9976,6 +9976,50 @@ class TestSchedulerJob:
         call_args = mock_listener_manager.hook.on_dag_run_success.call_args
         assert call_args.kwargs["dag_run"]._team_name == "testing"
 
+    @conf_vars({("core", "multi_team"): "true"})
+    def test_process_task_event_logs_records_the_team_owning_the_dag(self, 
dag_maker, session, team_bundle):
+        with dag_maker(dag_id="test_task_event_log_team", 
bundle_name="testing", session=session):
+            EmptyOperator(task_id="test_task")
+        session.commit()
+
+        self.job_runner = SchedulerJobRunner(Job(), executors=[MagicMock()])
+        self.job_runner._process_task_event_logs(
+            deque([Log(event="test_task_event_log_team_event", 
dag_id="test_task_event_log_team")]), session
+        )
+
+        log = session.scalar(select(Log).where(Log.event == 
"test_task_event_log_team_event"))
+        assert log.team_name == "testing"
+
+    @conf_vars({("core", "multi_team"): "true"})
+    @mock.patch("airflow.jobs.scheduler_job_runner.resolve_team_name", 
side_effect=resolve_team_name)
+    def test_process_task_event_logs_resolves_each_dag_once(
+        self, mock_resolve_team_name, dag_maker, session, team_bundle
+    ):
+        for dag_id in ("test_task_event_log_dag_1", 
"test_task_event_log_dag_2"):
+            with dag_maker(dag_id=dag_id, bundle_name="testing", 
session=session):
+                EmptyOperator(task_id="test_task")
+        session.commit()
+
+        self.job_runner = SchedulerJobRunner(Job(), executors=[MagicMock()])
+        self.job_runner._process_task_event_logs(
+            deque(
+                Log(event="test_task_event_log_dedupe", dag_id=dag_id)
+                for dag_id in (
+                    "test_task_event_log_dag_1",
+                    "test_task_event_log_dag_2",
+                    "test_task_event_log_dag_1",
+                    "test_task_event_log_dag_2",
+                    "test_task_event_log_dag_1",
+                )
+            ),
+            session,
+        )
+
+        assert mock_resolve_team_name.call_count == 2
+        logs = session.scalars(select(Log).where(Log.event == 
"test_task_event_log_dedupe")).all()
+        assert len(logs) == 5
+        assert {log.team_name for log in logs} == {"testing"}
+
     @mock.patch("airflow.models.Deadline.handle_miss")
     def test_process_expired_deadlines(self, mock_handle_miss, session, 
dag_maker):
         """Verify all expired and unhandled deadlines (and only those) are 
processed by the scheduler."""
diff --git a/airflow-core/tests/unit/models/test_log.py 
b/airflow-core/tests/unit/models/test_log.py
index 24961e3b00c..cd579ce39b6 100644
--- a/airflow-core/tests/unit/models/test_log.py
+++ b/airflow-core/tests/unit/models/test_log.py
@@ -22,10 +22,21 @@ from sqlalchemy import select
 from sqlalchemy.exc import InvalidRequestError
 from sqlalchemy.orm import joinedload
 
+from airflow.models.dag import DagModel, clear_team_name_cache
+from airflow.models.dagbundle import DagBundleModel
 from airflow.models.log import Log
+from airflow.models.team import Team
 from airflow.operators.empty import EmptyOperator
 from airflow.utils.state import TaskInstanceState
 
+from tests_common.test_utils.config import conf_vars
+from tests_common.test_utils.db import (
+    clear_db_dag_bundles,
+    clear_db_dags,
+    clear_db_logs,
+    clear_db_teams,
+)
+
 pytestmark = pytest.mark.db_test
 
 
@@ -104,3 +115,80 @@ class TestLogTaskInstanceReproduction:
         assert loaded_log2.task_instance is not None
         assert loaded_log2.task_instance.dag_id == "dag_2"
         assert loaded_log2.task_instance.run_id == ti2.run_id
+
+
+DAG_IN_TEAM = "dag_owned_by_a_team"
+
+
+class TestLogTeamName:
+    def teardown_method(self):
+        clear_db_logs()
+        clear_db_dags()
+        clear_db_dag_bundles()
+        clear_db_teams()
+
+    @staticmethod
+    def _create_dag_owned_by_team(session, team_name: str, *, 
dag_id=DAG_IN_TEAM, bundle_name="team-bundle"):
+        bundle = DagBundleModel(name=bundle_name)
+        bundle.teams.append(Team(name=team_name))
+        session.add(bundle)
+        session.flush()
+        session.add(DagModel(dag_id=dag_id, bundle_name=bundle_name, 
is_stale=False))
+        session.commit()
+        clear_team_name_cache()
+
+    @conf_vars({("core", "multi_team"): "True"})
+    def test_team_owning_the_dag_is_recorded(self, session):
+        self._create_dag_owned_by_team(session, "payments")
+
+        log = Log(event="test_event", dag_id=DAG_IN_TEAM)
+        session.add(log)
+        session.commit()
+
+        assert log.team_name == "payments"
+
+    @conf_vars({("core", "multi_team"): "True"})
+    def test_team_named_by_the_caller_is_kept(self, session):
+        self._create_dag_owned_by_team(session, "payments")
+
+        log = Log(event="test_event", dag_id=DAG_IN_TEAM, team_name="infra")
+        session.add(log)
+        session.commit()
+
+        assert log.team_name == "infra"
+
+    @conf_vars({("core", "multi_team"): "True"})
+    def test_recorded_team_outlives_the_dag_changing_teams(self, session):
+        self._create_dag_owned_by_team(session, "payments")
+        log = Log(event="test_event", dag_id=DAG_IN_TEAM)
+        session.add(log)
+        session.commit()
+
+        other_bundle = DagBundleModel(name="other-team-bundle")
+        other_bundle.teams.append(Team(name="infra"))
+        session.add(other_bundle)
+        session.flush()
+        session.scalar(
+            select(DagModel).where(DagModel.dag_id == DAG_IN_TEAM)
+        ).bundle_name = "other-team-bundle"
+        session.commit()
+        clear_team_name_cache()
+
+        assert log.team_name == "payments"
+
+    @conf_vars({("core", "multi_team"): "True"})
+    def test_no_team_is_recorded_for_an_event_owning_no_dag(self, session):
+        log = Log(event="test_event")
+        session.add(log)
+        session.commit()
+
+        assert log.team_name is None
+
+    def test_no_team_is_recorded_when_multi_team_is_off(self, session):
+        self._create_dag_owned_by_team(session, "payments")
+
+        log = Log(event="test_event", dag_id=DAG_IN_TEAM)
+        session.add(log)
+        session.commit()
+
+        assert log.team_name is None
diff --git a/airflow-core/tests/unit/utils/test_cli_util.py 
b/airflow-core/tests/unit/utils/test_cli_util.py
index 9cc8a7350c4..3539f8b3e52 100644
--- a/airflow-core/tests/unit/utils/test_cli_util.py
+++ b/airflow-core/tests/unit/utils/test_cli_util.py
@@ -33,10 +33,15 @@ import airflow
 from airflow import settings
 from airflow._shared.timezones import timezone
 from airflow.exceptions import AirflowException
+from airflow.models.dag import DagModel
+from airflow.models.dagbundle import DagBundleModel
 from airflow.models.log import Log
+from airflow.models.team import Team
 from airflow.utils import cli, cli_action_loggers
 from airflow.utils.cli import _search_for_dag_file
 
+from tests_common.test_utils.config import conf_vars
+
 # Mark entire module as db_test because ``action_cli`` wrapper still could use 
DB on callbacks:
 # - ``cli_action_loggers.on_pre_execution``
 # - ``cli_action_loggers.on_post_execution``
@@ -62,6 +67,18 @@ class TestCliUtil:
         assert metrics.get("start_datetime") <= timezone.utcnow()
         assert metrics.get("full_command")
 
+    @pytest.mark.parametrize(
+        ("func_name", "namespace", "expected_team_name"),
+        [
+            pytest.param("triggerer", Namespace(team_name="payments"), 
"payments", id="team-name-option"),
+            pytest.param("team_create", Namespace(name="payments"), 
"payments", id="teams-positional"),
+            pytest.param("team_list", Namespace(output="table"), None, 
id="no-team"),
+            pytest.param("dag_list", Namespace(name="not-a-team"), None, 
id="name-of-something-else"),
+        ],
+    )
+    def test_metrics_build_team_name(self, func_name, namespace, 
expected_team_name):
+        assert cli._build_metrics(func_name, namespace).get("team_name") == 
expected_team_name
+
     def test_fail_function(self):
         """
         Actual function is failing and fail needs to be propagated.
@@ -185,6 +202,47 @@ class TestCliUtil:
         command = ast.literal_eval(command)
         assert command == expected_command
 
+    def test_action_log_records_team_name(self, session):
+        namespace = Namespace(team_name="payments")
+        with (
+            mock.patch.object(sys, "argv", ["airflow", "triggerer", 
"--team-name", "payments"]),
+            mock.patch("airflow.utils.session.create_session") as 
mock_create_session,
+        ):
+            metrics = cli._build_metrics("triggerer", namespace)
+            mock_create_session.return_value = session.begin_nested()
+            mock_create_session.return_value.bulk_insert_mappings = 
session.bulk_insert_mappings
+            cli_action_loggers.default_action_log(**metrics)
+
+            log = session.scalar(select(Log).order_by(Log.dttm.desc()))
+
+        assert log.event == "cli_triggerer"
+        assert log.dag_id is None
+        assert log.team_name == "payments"
+
+    @conf_vars({("core", "multi_team"): "True"})
+    def test_action_log_records_the_team_owning_the_dag(self, session):
+        bundle = DagBundleModel(name="team-bundle")
+        bundle.teams.append(Team(name="payments"))
+        session.add(bundle)
+        session.flush()
+        session.add(DagModel(dag_id="dag_owned_by_a_team", 
bundle_name="team-bundle", is_stale=False))
+        session.flush()
+
+        namespace = Namespace(dag_id="dag_owned_by_a_team", subcommand="pause")
+        with (
+            mock.patch.object(sys, "argv", ["airflow", "dags", "pause", 
"dag_owned_by_a_team"]),
+            mock.patch("airflow.utils.session.create_session") as 
mock_create_session,
+        ):
+            metrics = cli._build_metrics("dag_pause", namespace)
+            mock_create_session.return_value = session.begin_nested()
+            mock_create_session.return_value.bulk_insert_mappings = 
session.bulk_insert_mappings
+            mock_create_session.return_value.scalar = session.scalar
+            cli_action_loggers.default_action_log(**metrics)
+
+            log = session.scalar(select(Log).where(Log.dag_id == 
"dag_owned_by_a_team"))
+
+        assert log.team_name == "payments"
+
     def test_setup_locations_relative_pid_path(self):
         relative_pid_path = "fake.pid"
         pid_full_path = os.path.join(os.getcwd(), relative_pid_path)
diff --git a/airflow-ctl/src/airflowctl/api/datamodels/generated.py 
b/airflow-ctl/src/airflowctl/api/datamodels/generated.py
index a1383e643e1..b5ef5665ac0 100644
--- a/airflow-ctl/src/airflowctl/api/datamodels/generated.py
+++ b/airflow-ctl/src/airflowctl/api/datamodels/generated.py
@@ -686,6 +686,7 @@ class EventLogResponse(BaseModel):
     extra: Annotated[str | None, Field(title="Extra")]
     dag_display_name: Annotated[str | None, Field(title="Dag Display Name")] = 
None
     task_display_name: Annotated[str | None, Field(title="Task Display Name")] 
= None
+    team_name: Annotated[str | None, Field(title="Team Name")] = None
 
 
 class ExternalLogUrlResponse(BaseModel):

Reply via email to