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

pierrejeambrun 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 0b6c38502ab Add runtime DagRun note updates (#69403)
0b6c38502ab is described below

commit 0b6c38502abb6a152ab6a6334910f171c75c4e5f
Author: Hemkumar Chheda <[email protected]>
AuthorDate: Wed Sep 30 22:06:11 2026 +0530

    Add runtime DagRun note updates (#69403)
    
    * Add runtime DagRun note updates
    
    Runtime code needs a supported Task SDK path for updating the current Dag 
run note without writing directly to the metadata database. Keeping the write 
scoped to the authenticated task instance preserves the Execution API isolation 
boundary while supporting callback use cases.
    
    * Mark DagRun note updates as worker-only
    
    The triggerer message-type guard needs to record worker-only runtime 
operations explicitly. Dag run note updates require a task execution context 
and should not become part of the triggerer supervisor protocol.
    
    * Reuse patch_dag_run_note for runtime DagRun note updates
    
    The runtime note endpoint reimplemented the set/update/clear logic from
    patch_dag_run_note and diverged on clearing: a runtime clear left a
    dag_run_note row with NULL content, while the UI/public API deletes it.
    Make patch_dag_run_note's user attribution optional (user_id) and reuse
    it from the execution API so both paths behave the same, including
    removing the note row on clear.
    
    Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
    
    * Preserve runtime DagRun notes on null update payloads
    
    * Register the runtime DagRun note endpoint on the in-progress API version
    
    The endpoint was registered under Version("2026-06-30"), which was the
    in-progress execution API version when this branch was cut. That version has
    since been released and 2026-10-30 is the in-progress one now, so the 
released
    2026-06-30 spec was advertising an endpoint that release never shipped.
    
    Move the VersionChange to v2026_10_30.py and add the previous-version test 
the
    versioning guide asks for, which fails if the change sits on a released
    version. Also document that a null note leaves an existing note untouched
    while an empty note removes it, since the signature alone reads the other 
way.
    
    Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
    
    * Address review feedback on runtime DagRun note updates
    
    Skip the supervisor round-trip when the note is None, log when a runtime
    write replaces an attributed note, and rename the attribution test to say
    what it proves.
    
    ---------
    
    Co-authored-by: Claude Opus 4.8 (1M context) <[email protected]>
---
 airflow-core/docs/public-airflow-interface.rst     |   5 +-
 .../api_fastapi/core_api/routes/public/dag_run.py  |   2 +-
 .../core_api/services/public/dag_run.py            |  17 ++-
 .../execution_api/datamodels/taskinstance.py       |   6 +
 .../execution_api/routes/task_instances.py         |  55 +++++++++
 .../api_fastapi/execution_api/versions/__init__.py |   2 +
 .../execution_api/versions/v2026_10_30.py          |  10 ++
 .../versions/head/test_task_instances.py           | 125 +++++++++++++++++++++
 .../versions/v2026_10_30/test_dag_run_notes.py     |  43 +++++++
 .../tests/unit/dag_processing/test_processor.py    |   1 +
 airflow-core/tests/unit/jobs/test_triggerer_job.py |   1 +
 task-sdk/src/airflow/sdk/api/client.py             |  12 ++
 .../src/airflow/sdk/api/datamodels/_generated.py   |  15 +++
 task-sdk/src/airflow/sdk/execution_time/comms.py   |   7 ++
 .../airflow/sdk/execution_time/request_handlers.py |   9 ++
 .../airflow/sdk/execution_time/schema/schema.json  |  32 ++++++
 .../execution_time/schema/versions/v2026_10_30.py  |   9 ++
 .../src/airflow/sdk/execution_time/supervisor.py   |   8 ++
 .../src/airflow/sdk/execution_time/task_runner.py  |  12 ++
 task-sdk/src/airflow/sdk/types.py                  |   2 +
 task-sdk/tests/task_sdk/api/test_client.py         |  15 +++
 .../task_sdk/execution_time/test_supervisor.py     |  15 +++
 .../task_sdk/execution_time/test_task_runner.py    |  20 ++++
 ts-sdk/src/generated/supervisor.ts                 |  50 +++++----
 24 files changed, 446 insertions(+), 27 deletions(-)

diff --git a/airflow-core/docs/public-airflow-interface.rst 
b/airflow-core/docs/public-airflow-interface.rst
index 756ea2e8bc0..b2180aeca87 100644
--- a/airflow-core/docs/public-airflow-interface.rst
+++ b/airflow-core/docs/public-airflow-interface.rst
@@ -556,7 +556,10 @@ of the following alternatives:
 * **Task Context**: Use :func:`~airflow.sdk.get_current_context` to access 
task instance
   information and methods like 
:meth:`~airflow.sdk.types.RuntimeTaskInstanceProtocol.get_dr_count`,
   :meth:`~airflow.sdk.types.RuntimeTaskInstanceProtocol.get_dagrun_state`, and
-  :meth:`~airflow.sdk.types.RuntimeTaskInstanceProtocol.get_task_states`.
+  :meth:`~airflow.sdk.types.RuntimeTaskInstanceProtocol.get_task_states`. To 
update the note for
+  the current Dag run at runtime, use
+  :meth:`~airflow.sdk.types.RuntimeTaskInstanceProtocol.update_dagrun_note` 
instead of writing
+  directly to ``DagRun.note`` in the metadata database.
 
 * **REST API**: Use the :doc:`Stable REST API <stable-rest-api-ref>` for 
programmatic
   access to Airflow metadata.
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 3345dff44ef..2515344c50d 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
@@ -239,7 +239,7 @@ def patch_dag_run(
     if "note" in data:
         updated_dag_run = session.get(DagRun, dag_run.id)
         if updated_dag_run is not None:
-            patch_dag_run_note(dag_run=updated_dag_run, note=data["note"], 
user=user)
+            patch_dag_run_note(dag_run=updated_dag_run, note=data["note"], 
user_id=user.get_id())
     if "state" in data and patch_body.state is not None:
         patch_dag_run_state(dag=dag, dag_run=dag_run, state=patch_body.state, 
session=session)
 
diff --git 
a/airflow-core/src/airflow/api_fastapi/core_api/services/public/dag_run.py 
b/airflow-core/src/airflow/api_fastapi/core_api/services/public/dag_run.py
index 9f5e65501c0..ebd149e05dd 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/services/public/dag_run.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/services/public/dag_run.py
@@ -156,7 +156,7 @@ def perform_clear_dag_run(
     if not dag_run_cleared:
         raise HTTPException(status.HTTP_404_NOT_FOUND, "Dag run not found 
after clearing")
     if note is not None:
-        patch_dag_run_note(dag_run=dag_run_cleared, note=note, user=user)
+        patch_dag_run_note(dag_run=dag_run_cleared, note=note, 
user_id=user.get_id())
     return dag_run_cleared
 
 
@@ -228,15 +228,20 @@ def patch_dag_run_state(
             log.exception("error calling listener")
 
 
-def patch_dag_run_note(*, dag_run: DagRun, note: str | None, user: BaseUser) 
-> None:
-    """Set, update, or clear a Dag Run's note. An empty note removes it so the 
run is left without a note."""
+def patch_dag_run_note(*, dag_run: DagRun, note: str | None, user_id: str | 
None) -> None:
+    """
+    Set, update, or clear a Dag Run's note. An empty note removes it so the 
run is left without a note.
+
+    ``user_id`` is the author to attribute the note to, or ``None`` for an 
unattributed note
+    (e.g. a note written from task runtime, which has no acting user).
+    """
     if note == "":
         dag_run.dag_run_note = None
     elif dag_run.dag_run_note is None:
-        dag_run.note = (note, user.get_id())
+        dag_run.note = (note, user_id)
     else:
         dag_run.dag_run_note.content = note
-        dag_run.dag_run_note.user_id = user.get_id()
+        dag_run.dag_run_note.user_id = user_id
 
 
 @attrs.define
@@ -426,7 +431,7 @@ class BulkDagRunService(BulkService[BulkDAGRunBody]):
                     dag = get_dag_for_run(self.dag_bag, dag_run, 
session=self.session)
                     patch_dag_run_state(dag=dag, dag_run=dag_run, 
state=entity.state, session=self.session)
                 if entity.note is not None:
-                    patch_dag_run_note(dag_run=dag_run, note=entity.note, 
user=self.user)
+                    patch_dag_run_note(dag_run=dag_run, note=entity.note, 
user_id=self.user.get_id())
                 results.success.append(f"{dag_id}.{run_id}")
         except HTTPException as e:
             results.errors.append({"error": f"{e.detail}", "status_code": 
e.status_code})
diff --git 
a/airflow-core/src/airflow/api_fastapi/execution_api/datamodels/taskinstance.py 
b/airflow-core/src/airflow/api_fastapi/execution_api/datamodels/taskinstance.py
index f5153b5de46..b1b4301f360 100644
--- 
a/airflow-core/src/airflow/api_fastapi/execution_api/datamodels/taskinstance.py
+++ 
b/airflow-core/src/airflow/api_fastapi/execution_api/datamodels/taskinstance.py
@@ -67,6 +67,12 @@ class TIEnterRunningPayload(StrictBaseModel):
     """When the task started executing"""
 
 
+class DagRunNoteUpdatePayload(StrictBaseModel):
+    """Schema for updating the DagRun note associated with a task instance."""
+
+    note: str | None = Field(None, max_length=1000)
+
+
 # Create an enum to give a nice name in the generated datamodels
 class TerminalStateNonSuccess(str, Enum):
     """TaskInstance states that can be reported without extra information."""
diff --git 
a/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py 
b/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py
index 1731bae3fe2..4ef48dce0bf 100644
--- 
a/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py
+++ 
b/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py
@@ -52,8 +52,10 @@ from airflow.api_fastapi.common.db.dags import 
eager_load_teams
 from airflow.api_fastapi.common.types import UtcDateTime
 from airflow.api_fastapi.compat import HTTP_422_UNPROCESSABLE_CONTENT
 from airflow.api_fastapi.core_api.openapi.exceptions import 
create_openapi_http_exception_doc
+from airflow.api_fastapi.core_api.services.public.dag_run import 
patch_dag_run_note
 from airflow.api_fastapi.execution_api.datamodels.task_arg_binding import 
get_arg_bindings_adapter
 from airflow.api_fastapi.execution_api.datamodels.taskinstance import (
+    DagRunNoteUpdatePayload,
     InactiveAssetsResponse,
     PreviousTIResponse,
     PrevSuccessfulDagRunResponse,
@@ -968,6 +970,59 @@ def _raise_ti_not_in_live_table(task_instance_id: UUID, *, 
archived_in_history:
     )
 
 
+@ti_id_router.patch(
+    "/{task_instance_id}/dag-run-note",
+    status_code=status.HTTP_204_NO_CONTENT,
+    responses=create_openapi_http_exception_doc(
+        [
+            (status.HTTP_404_NOT_FOUND, "Task Instance not found"),
+            (HTTP_422_UNPROCESSABLE_CONTENT, "Invalid payload for the DagRun 
note update"),
+        ]
+    ),
+)
+def update_dag_run_note(
+    task_instance_id: UUID,
+    body: DagRunNoteUpdatePayload,
+    session: SessionDep,
+) -> None:
+    """
+    Update the note for the DagRun associated with this task instance.
+
+    An empty note removes the existing note, matching the public API. A null 
note is a
+    no-op so runtime callers can leave a user-authored note untouched.
+    """
+    bind_contextvars(ti_id=str(task_instance_id))
+
+    dag_run = session.scalar(
+        select(DR)
+        .join(TI, and_(TI.dag_id == DR.dag_id, TI.run_id == DR.run_id))
+        .options(joinedload(DR.dag_run_note))
+        .where(TI.id == task_instance_id)
+    )
+    if dag_run is None:
+        raise HTTPException(
+            status_code=status.HTTP_404_NOT_FOUND,
+            detail={"reason": "not_found", "message": "Task Instance not 
found"},
+        )
+
+    if body.note is None:
+        return
+
+    # Runtime notes have no acting user, so they are stored unattributed. 
Carrying over the
+    # previous author would credit them with content they did not write, so 
log the drop
+    # instead of keeping it.
+    if dag_run.dag_run_note is not None and dag_run.dag_run_note.user_id is 
not None:
+        log.info(
+            "Replacing an attributed DagRun note from task runtime; the note 
becomes unattributed",
+            dag_id=dag_run.dag_id,
+            run_id=dag_run.run_id,
+            previous_user_id=dag_run.dag_run_note.user_id,
+        )
+
+    # Reuse the public API note logic so both editing paths stay consistent.
+    patch_dag_run_note(dag_run=dag_run, note=body.note, user_id=None)
+
+
 @ti_id_router.put(
     "/{task_instance_id}/heartbeat",
     status_code=status.HTTP_204_NO_CONTENT,
diff --git 
a/airflow-core/src/airflow/api_fastapi/execution_api/versions/__init__.py 
b/airflow-core/src/airflow/api_fastapi/execution_api/versions/__init__.py
index 6df38592984..3f5050b19b1 100644
--- a/airflow-core/src/airflow/api_fastapi/execution_api/versions/__init__.py
+++ b/airflow-core/src/airflow/api_fastapi/execution_api/versions/__init__.py
@@ -55,6 +55,7 @@ from airflow.api_fastapi.execution_api.versions.v2026_06_30 
import (
 from airflow.api_fastapi.execution_api.versions.v2026_10_30 import (
     AddArgBindingsToTIRunContext,
     AddCallbackRunEndpoint,
+    AddDagRunNoteUpdateEndpoint,
     AddMultiTeamToTIRunContext,
     AddStoppedTaskReport,
     AddTerminalStateRetryReasonField,
@@ -67,6 +68,7 @@ bundle = VersionBundle(
         "2026-10-30",
         AddArgBindingsToTIRunContext,
         AddCallbackRunEndpoint,
+        AddDagRunNoteUpdateEndpoint,
         AddTerminalStateRetryReasonField,
         AddMultiTeamToTIRunContext,
         AddStoppedTaskReport,
diff --git 
a/airflow-core/src/airflow/api_fastapi/execution_api/versions/v2026_10_30.py 
b/airflow-core/src/airflow/api_fastapi/execution_api/versions/v2026_10_30.py
index 3c1dc379580..050018f45a3 100644
--- a/airflow-core/src/airflow/api_fastapi/execution_api/versions/v2026_10_30.py
+++ b/airflow-core/src/airflow/api_fastapi/execution_api/versions/v2026_10_30.py
@@ -77,6 +77,16 @@ class AddCallbackRunEndpoint(VersionChange):
     )
 
 
+class AddDagRunNoteUpdateEndpoint(VersionChange):
+    """Add endpoint for updating a DagRun note from task runtime code."""
+
+    description = __doc__
+
+    instructions_to_migrate_to_previous_version = (
+        endpoint("/task-instances/{task_instance_id}/dag-run-note", 
["PATCH"]).didnt_exist,
+    )
+
+
 class AddTerminalStateRetryReasonField(VersionChange):
     """Add the `retry_reason` field to TITerminalStatePayload for failed 
retry-policy decisions."""
 
diff --git 
a/airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py
 
b/airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py
index c447149ac2c..9d77f666bed 100644
--- 
a/airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py
+++ 
b/airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py
@@ -4972,6 +4972,131 @@ class TestTIPatchRenderedMapIndex:
         assert response.status_code == 422
 
 
+class TestTIDagRunNoteUpdate:
+    def setup_method(self):
+        clear_db_runs()
+
+    def teardown_method(self):
+        clear_db_runs()
+
+    def test_create_dag_run_note(self, client, session, create_task_instance):
+        ti = create_task_instance(
+            task_id="test_create_dag_run_note",
+            state=State.RUNNING,
+            session=session,
+        )
+        ti_id = ti.id
+        session.commit()
+
+        response = client.patch(
+            f"/execution/task-instances/{ti_id}/dag-run-note",
+            json={"note": "Created from task runtime"},
+        )
+
+        assert response.status_code == 204
+        assert response.text == ""
+
+        session.expire_all()
+        dag_run = session.get(TaskInstance, ti_id).dag_run
+        assert dag_run.note == "Created from task runtime"
+        assert dag_run.dag_run_note.user_id is None
+
+    def test_runtime_update_of_user_note_becomes_unattributed(self, client, 
session, create_task_instance):
+        """Runtime rewrites the content, so the previous author is not carried 
over."""
+        ti = create_task_instance(
+            task_id="test_update_dag_run_note",
+            state=State.RUNNING,
+            session=session,
+        )
+        ti.dag_run.note = ("Created from UI", "user_id")
+        ti_id = ti.id
+        session.commit()
+
+        response = client.patch(
+            f"/execution/task-instances/{ti_id}/dag-run-note",
+            json={"note": "Updated from task runtime"},
+        )
+
+        assert response.status_code == 204
+        assert response.text == ""
+
+        session.expire_all()
+        dag_run = session.get(TaskInstance, ti_id).dag_run
+        assert dag_run.note == "Updated from task runtime"
+        assert dag_run.dag_run_note.user_id is None
+
+    def test_clear_dag_run_note(self, client, session, create_task_instance):
+        ti = create_task_instance(
+            task_id="test_clear_dag_run_note",
+            state=State.RUNNING,
+            session=session,
+        )
+        ti.dag_run.note = ("Will be cleared", "user_id")
+        ti_id = ti.id
+        session.commit()
+
+        response = client.patch(
+            f"/execution/task-instances/{ti_id}/dag-run-note",
+            json={"note": ""},
+        )
+
+        assert response.status_code == 204
+        assert response.text == ""
+
+        session.expire_all()
+        dag_run = session.get(TaskInstance, ti_id).dag_run
+        # Clearing removes the note row entirely, same as the UI/public API 
path.
+        assert dag_run.note is None
+        assert dag_run.dag_run_note is None
+
+    def test_null_note_leaves_existing_dag_run_note_unchanged(self, client, 
session, create_task_instance):
+        ti = create_task_instance(
+            task_id="test_null_note_leaves_existing_dag_run_note_unchanged",
+            state=State.RUNNING,
+            session=session,
+        )
+        ti.dag_run.note = ("Keep existing note", "user_id")
+        ti_id = ti.id
+        session.commit()
+
+        response = client.patch(
+            f"/execution/task-instances/{ti_id}/dag-run-note",
+            json={"note": None},
+        )
+
+        assert response.status_code == 204
+        assert response.text == ""
+
+        session.expire_all()
+        dag_run = session.get(TaskInstance, ti_id).dag_run
+        assert dag_run.note == "Keep existing note"
+        assert dag_run.dag_run_note.user_id == "user_id"
+
+    def test_update_dag_run_note_task_instance_not_found(self, client, 
session):
+        response = client.patch(
+            f"/execution/task-instances/{uuid4()}/dag-run-note",
+            json={"note": "Does not matter"},
+        )
+
+        assert response.status_code == 404
+
+    def test_update_dag_run_note_rejects_too_long_note(self, client, session, 
create_task_instance):
+        ti = create_task_instance(
+            task_id="test_update_dag_run_note_rejects_too_long_note",
+            state=State.RUNNING,
+            session=session,
+        )
+        ti_id = ti.id
+        session.commit()
+
+        response = client.patch(
+            f"/execution/task-instances/{ti_id}/dag-run-note",
+            json={"note": "x" * 1001},
+        )
+
+        assert response.status_code == 422
+
+
 @pytest.mark.usefixtures("_use_real_jwt_bearer")
 class TestTokenTypeValidation:
     """Test token scope enforcement (workload vs execution)."""
diff --git 
a/airflow-core/tests/unit/api_fastapi/execution_api/versions/v2026_10_30/test_dag_run_notes.py
 
b/airflow-core/tests/unit/api_fastapi/execution_api/versions/v2026_10_30/test_dag_run_notes.py
new file mode 100644
index 00000000000..2702a9ad508
--- /dev/null
+++ 
b/airflow-core/tests/unit/api_fastapi/execution_api/versions/v2026_10_30/test_dag_run_notes.py
@@ -0,0 +1,43 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+from __future__ import annotations
+
+import pytest
+
+pytestmark = pytest.mark.db_test
+
+MISSING_TI_NOTE_URL = 
"/execution/task-instances/00000000-0000-0000-0000-000000000000/dag-run-note"
+
+
+class TestUpdateDagRunNoteEndpointVersioning:
+    """The task-instances/{task_instance_id}/dag-run-note endpoint didn't 
exist before 2026-10-30."""
+
+    def test_old_version_returns_404(self, client):
+        """Before 2026-10-30 the route is absent, so routing itself 404s (no 
endpoint-shaped detail)."""
+        client.headers["Airflow-API-Version"] = "2026-06-30"
+
+        response = client.patch(MISSING_TI_NOTE_URL, json={"note": "from 
runtime"})
+
+        assert response.status_code == 404
+        assert response.json() == {"detail": "Not Found"}
+
+    def test_head_version_routes_to_endpoint(self, client):
+        """At head the route exists: the same request reaches the endpoint's 
own 404 handling."""
+        response = client.patch(MISSING_TI_NOTE_URL, json={"note": "from 
runtime"})
+
+        assert response.status_code == 404
+        assert response.json()["detail"]["reason"] == "not_found"
diff --git a/airflow-core/tests/unit/dag_processing/test_processor.py 
b/airflow-core/tests/unit/dag_processing/test_processor.py
index f3966f3916a..57cce320a80 100644
--- a/airflow-core/tests/unit/dag_processing/test_processor.py
+++ b/airflow-core/tests/unit/dag_processing/test_processor.py
@@ -2283,6 +2283,7 @@ class TestDagProcessingMessageTypes:
             "DeleteAssetStateStoreByUri",
             "ClearAssetStateStoreByName",
             "ClearAssetStateStoreByUri",
+            "UpdateDagRunNote",
         }
 
         in_task_runner_but_not_in_dag_processing_process = {
diff --git a/airflow-core/tests/unit/jobs/test_triggerer_job.py 
b/airflow-core/tests/unit/jobs/test_triggerer_job.py
index 3336bedd3b0..ff9d8806cf4 100644
--- a/airflow-core/tests/unit/jobs/test_triggerer_job.py
+++ b/airflow-core/tests/unit/jobs/test_triggerer_job.py
@@ -3071,6 +3071,7 @@ class TestTriggererMessageTypes:
             "SetTaskStateStore",
             "DeleteTaskStateStore",
             "ClearTaskStateStore",
+            "UpdateDagRunNote",
         }
 
         in_task_but_not_in_trigger_runner = {
diff --git a/task-sdk/src/airflow/sdk/api/client.py 
b/task-sdk/src/airflow/sdk/api/client.py
index 88c3d017f16..b9877e7b12d 100644
--- a/task-sdk/src/airflow/sdk/api/client.py
+++ b/task-sdk/src/airflow/sdk/api/client.py
@@ -56,12 +56,14 @@ from airflow.sdk.api.datamodels._generated import (
     ConnectionTestState,
     DagResponse,
     DagRun,
+    DagRunNoteUpdatePayload,
     DagRunStateResponse,
     DagRunType,
     HITLDetailRequest,
     HITLDetailResponse,
     HITLUser,
     InactiveAssetsResponse,
+    Note,
     PrevSuccessfulDagRunResponse,
     ResultMessage,
     TaskBreadcrumbsResponse,
@@ -354,6 +356,16 @@ class TaskInstanceOperations:
         body = TISkippedDownstreamTasksStatePayload(tasks=msg.tasks)
         self.client.patch(f"task-instances/{id}/skip-downstream", 
content=body.model_dump_json())
 
+    def update_dagrun_note(self, id: uuid.UUID, note: str | None) -> 
OKResponse:
+        """
+        Update the note for the DagRun associated with this task instance.
+
+        An empty note removes it, and ``None`` leaves any existing note 
untouched.
+        """
+        body = DagRunNoteUpdatePayload(note=Note(note) if note is not None 
else None)
+        self.client.patch(f"task-instances/{id}/dag-run-note", 
content=body.model_dump_json())
+        return OKResponse(ok=True)
+
     def set_rtif(self, id: uuid.UUID, body: dict[str, str]) -> OKResponse:
         """Set Rendered Task Instance Fields via the API server."""
         self.client.put(f"task-instances/{id}/rtif", json=body)
diff --git a/task-sdk/src/airflow/sdk/api/datamodels/_generated.py 
b/task-sdk/src/airflow/sdk/api/datamodels/_generated.py
index 9f221cf0f3f..96fd94e355f 100644
--- a/task-sdk/src/airflow/sdk/api/datamodels/_generated.py
+++ b/task-sdk/src/airflow/sdk/api/datamodels/_generated.py
@@ -143,6 +143,21 @@ class DagRunAssetReference(BaseModel):
     partition_key: Annotated[str | None, Field(title="Partition Key")]
 
 
+class Note(RootModel[str]):
+    root: Annotated[str, Field(max_length=1000, title="Note")]
+
+
+class DagRunNoteUpdatePayload(BaseModel):
+    """
+    Schema for updating the DagRun note associated with a task instance.
+    """
+
+    model_config = ConfigDict(
+        extra="forbid",
+    )
+    note: Annotated[Note | None, Field(title="Note")] = None
+
+
 class DagRunState(str, Enum):
     """
     All possible states that a DagRun can be in.
diff --git a/task-sdk/src/airflow/sdk/execution_time/comms.py 
b/task-sdk/src/airflow/sdk/execution_time/comms.py
index 3ac67112c53..ff151a01777 100644
--- a/task-sdk/src/airflow/sdk/execution_time/comms.py
+++ b/task-sdk/src/airflow/sdk/execution_time/comms.py
@@ -1143,6 +1143,12 @@ class GetDagRunState(BaseModel):
     type: Literal["GetDagRunState"] = "GetDagRunState"
 
 
+class UpdateDagRunNote(BaseModel):
+    ti_id: UUID
+    note: str | None
+    type: Literal["UpdateDagRunNote"] = "UpdateDagRunNote"
+
+
 class GetPreviousDagRun(BaseModel):
     dag_id: str
     logical_date: AwareDatetime
@@ -1332,6 +1338,7 @@ ToSupervisor = Annotated[
     | ValidateInletsAndOutlets
     | TaskState
     | TriggerDagRun
+    | UpdateDagRunNote
     | DeleteVariable
     | ResendLoggingFD
     | CreateHITLDetailPayload
diff --git a/task-sdk/src/airflow/sdk/execution_time/request_handlers.py 
b/task-sdk/src/airflow/sdk/execution_time/request_handlers.py
index 902e19447da..ed0e547463a 100644
--- a/task-sdk/src/airflow/sdk/execution_time/request_handlers.py
+++ b/task-sdk/src/airflow/sdk/execution_time/request_handlers.py
@@ -74,6 +74,7 @@ from airflow.sdk.execution_time.comms import (
     SetAssetStateStoreByUri,
     SetXCom,
     TaskStatesResult,
+    UpdateDagRunNote,
     VariableKeysResult,
     VariableResult,
     XComResult,
@@ -219,6 +220,14 @@ def handle_get_dag_run_state(client: Client, msg: 
GetDagRunState) -> tuple[BaseM
     return dr_resp, {}
 
 
+def handle_update_dag_run_note(
+    client: Client, msg: UpdateDagRunNote
+) -> tuple[BaseModel | None, dict[str, bool]]:
+    """Update the note for the DagRun associated with the current task 
instance."""
+    resp = client.task_instances.update_dagrun_note(msg.ti_id, msg.note)
+    return resp, {}
+
+
 def handle_get_previous_dag_run(
     client: Client, msg: GetPreviousDagRun
 ) -> tuple[BaseModel | None, dict[str, bool]]:
diff --git a/task-sdk/src/airflow/sdk/execution_time/schema/schema.json 
b/task-sdk/src/airflow/sdk/execution_time/schema/schema.json
index 4072443e52c..4777a686c98 100644
--- a/task-sdk/src/airflow/sdk/execution_time/schema/schema.json
+++ b/task-sdk/src/airflow/sdk/execution_time/schema/schema.json
@@ -4380,6 +4380,38 @@
       "title": "TriggerDagRun",
       "type": "object"
     },
+    "UpdateDagRunNote": {
+      "properties": {
+        "ti_id": {
+          "format": "uuid",
+          "title": "Ti Id",
+          "type": "string"
+        },
+        "note": {
+          "anyOf": [
+            {
+              "type": "string"
+            },
+            {
+              "type": "null"
+            }
+          ],
+          "title": "Note"
+        },
+        "type": {
+          "const": "UpdateDagRunNote",
+          "default": "UpdateDagRunNote",
+          "title": "Type",
+          "type": "string"
+        }
+      },
+      "required": [
+        "ti_id",
+        "note"
+      ],
+      "title": "UpdateDagRunNote",
+      "type": "object"
+    },
     "UpdateHITLDetail": {
       "description": "Update the response content part of an existing 
Human-in-the-loop response.",
       "properties": {
diff --git 
a/task-sdk/src/airflow/sdk/execution_time/schema/versions/v2026_10_30.py 
b/task-sdk/src/airflow/sdk/execution_time/schema/versions/v2026_10_30.py
index e6f6e920d45..af54445542b 100644
--- a/task-sdk/src/airflow/sdk/execution_time/schema/versions/v2026_10_30.py
+++ b/task-sdk/src/airflow/sdk/execution_time/schema/versions/v2026_10_30.py
@@ -15,6 +15,15 @@
 # specific language governing permissions and limitations
 # under the License.
 
+"""
+Supervisor schema version 2026-10-30.
+
+A brand-new message body needs no field-level migration instructions here: a 
lang-SDK
+pinned to an older version simply never sends it, so there is nothing to strip 
on the
+way down. Only *changes* to bodies that already existed at an earlier dated 
version
+require a ``VersionChange`` entry below.
+"""
+
 from __future__ import annotations
 
 from cadwyn import VersionChange, schema
diff --git a/task-sdk/src/airflow/sdk/execution_time/supervisor.py 
b/task-sdk/src/airflow/sdk/execution_time/supervisor.py
index 312766cd756..f2373b3bb83 100644
--- a/task-sdk/src/airflow/sdk/execution_time/supervisor.py
+++ b/task-sdk/src/airflow/sdk/execution_time/supervisor.py
@@ -141,6 +141,7 @@ from airflow.sdk.execution_time.comms import (
     TaskStateStoreResult,
     ToSupervisor,
     TriggerDagRun,
+    UpdateDagRunNote,
     ValidateInletsAndOutlets,
     _RequestFrame,
     _ResponseFrame,
@@ -166,6 +167,7 @@ from airflow.sdk.execution_time.request_handlers import (
     handle_mask_secret,
     handle_put_variable,
     handle_set_xcom,
+    handle_update_dag_run_note,
 )
 from airflow.sdk.execution_time.schema import get_schema_version_migrator, 
resolve_body_class
 
@@ -2132,6 +2134,11 @@ class ActivitySubprocess(WatchedSubprocess):
         )
         return resp, {}
 
+    def _handle_update_dag_run_note(
+        self, msg: UpdateDagRunNote, log: FilteringBoundLogger, req_id: int
+    ) -> RequestResult:
+        return handle_update_dag_run_note(self.client, msg)
+
     def _handle_get_dag_run(self, msg: GetDagRun, log: FilteringBoundLogger, 
req_id: int) -> RequestResult:
         dr_resp = self.client.dag_runs.get_detail(msg.dag_id, msg.run_id)
         resp = DagRunResult.from_api_response(dr_resp)
@@ -2344,6 +2351,7 @@ class ActivitySubprocess(WatchedSubprocess):
                 register_request_method(SucceedTask, _handle_finished_task),
                 register_request_method(TaskState, _handle_task_state),
                 register_request_method(TriggerDagRun, 
_handle_trigger_dag_run),
+                register_request_method(UpdateDagRunNote, 
_handle_update_dag_run_note),
                 register_request_method(ValidateInletsAndOutlets, 
_handle_validate_inlets_and_outlets),
             ]
         ),
diff --git a/task-sdk/src/airflow/sdk/execution_time/task_runner.py 
b/task-sdk/src/airflow/sdk/execution_time/task_runner.py
index 4e2b5bc7354..5dc3de05e13 100644
--- a/task-sdk/src/airflow/sdk/execution_time/task_runner.py
+++ b/task-sdk/src/airflow/sdk/execution_time/task_runner.py
@@ -124,6 +124,7 @@ from airflow.sdk.execution_time.comms import (
     ToSupervisor,
     ToTask,
     TriggerDagRun,
+    UpdateDagRunNote,
     ValidateInletsAndOutlets,
 )
 from airflow.sdk.execution_time.context import (
@@ -701,6 +702,17 @@ class RuntimeTaskInstance(TaskInstance):
 
         return response.dag_run
 
+    def update_dagrun_note(self, note: str | None) -> None:
+        """
+        Update the note for this task instance's DagRun.
+
+        A string sets or replaces the note and an empty string removes it. 
``None`` is a
+        no-op, so an existing user-authored note is left untouched.
+        """
+        if note is None:
+            return
+        SUPERVISOR_COMMS.send(msg=UpdateDagRunNote(ti_id=self.id, note=note))
+
     def get_previous_ti(
         self,
         state: TaskInstanceState | None = None,
diff --git a/task-sdk/src/airflow/sdk/types.py 
b/task-sdk/src/airflow/sdk/types.py
index efa3c9f6d94..c6d1a01bbd4 100644
--- a/task-sdk/src/airflow/sdk/types.py
+++ b/task-sdk/src/airflow/sdk/types.py
@@ -181,6 +181,8 @@ class RuntimeTaskInstanceProtocol(Protocol):
 
     def get_previous_dagrun(self, state: str | None = None) -> DagRunProtocol 
| None: ...
 
+    def update_dagrun_note(self, note: str | None) -> None: ...
+
     def get_previous_ti(
         self,
         state: TaskInstanceState | None = None,
diff --git a/task-sdk/tests/task_sdk/api/test_client.py 
b/task-sdk/tests/task_sdk/api/test_client.py
index a61e417d740..6a80ebb2a96 100644
--- a/task-sdk/tests/task_sdk/api/test_client.py
+++ b/task-sdk/tests/task_sdk/api/test_client.py
@@ -482,6 +482,21 @@ class TestTaskInstanceOperations:
 
         assert len(responses) == 1
 
+    def test_task_instance_update_dagrun_note(self):
+        ti_id = uuid6.uuid7()
+
+        def handle_request(request: httpx.Request) -> httpx.Response:
+            if request.url.path == f"/task-instances/{ti_id}/dag-run-note":
+                assert json.loads(request.read()) == {"note": "Updated from 
task runtime"}
+                return httpx.Response(status_code=204)
+            return httpx.Response(status_code=400, json={"detail": "Bad 
Request"})
+
+        client = make_client(transport=httpx.MockTransport(handle_request))
+
+        response = client.task_instances.update_dagrun_note(ti_id, "Updated 
from task runtime")
+
+        assert response == OKResponse(ok=True)
+
     @pytest.mark.parametrize("queues_enabled", [False, True])
     def test_task_instance_defer(self, queues_enabled: bool):
         # Simulate a successful response from the server that defers a task
diff --git a/task-sdk/tests/task_sdk/execution_time/test_supervisor.py 
b/task-sdk/tests/task_sdk/execution_time/test_supervisor.py
index a18cd809e88..12daa4173d5 100644
--- a/task-sdk/tests/task_sdk/execution_time/test_supervisor.py
+++ b/task-sdk/tests/task_sdk/execution_time/test_supervisor.py
@@ -38,6 +38,7 @@ from time import sleep
 from typing import TYPE_CHECKING, Any, get_args, get_type_hints
 from unittest import mock
 from unittest.mock import MagicMock, patch
+from uuid import UUID
 
 import httpx
 import msgspec
@@ -150,6 +151,7 @@ from airflow.sdk.execution_time.comms import (
     TICount,
     ToSupervisor,
     TriggerDagRun,
+    UpdateDagRunNote,
     UpdateHITLDetail,
     ValidateInletsAndOutlets,
     VariableKeysResult,
@@ -2932,6 +2934,19 @@ REQUEST_TEST_CASES = [
         ),
         test_id="get_dag_run_state",
     ),
+    RequestTestCase(
+        message=UpdateDagRunNote(
+            ti_id=UUID("9c230b40-da03-451d-8bd7-be30471be383"),
+            note="Updated from task runtime",
+        ),
+        expected_body={"ok": True, "type": "OKResponse"},
+        client_mock=ClientMock(
+            method_path="task_instances.update_dagrun_note",
+            args=(UUID("9c230b40-da03-451d-8bd7-be30471be383"), "Updated from 
task runtime"),
+            response=OKResponse(ok=True),
+        ),
+        test_id="update_dag_run_note",
+    ),
     RequestTestCase(
         message=GetPreviousDagRun(
             dag_id="test_dag",
diff --git a/task-sdk/tests/task_sdk/execution_time/test_task_runner.py 
b/task-sdk/tests/task_sdk/execution_time/test_task_runner.py
index 42e7c045449..ce9b1187c5e 100644
--- a/task-sdk/tests/task_sdk/execution_time/test_task_runner.py
+++ b/task-sdk/tests/task_sdk/execution_time/test_task_runner.py
@@ -155,6 +155,7 @@ from airflow.sdk.execution_time.comms import (
     TaskStatesResult,
     TICount,
     TriggerDagRun,
+    UpdateDagRunNote,
     ValidateInletsAndOutlets,
     VariableResult,
     XComResult,
@@ -3343,6 +3344,25 @@ class TestRuntimeTaskInstance:
         assert dr.run_id == "prev_run"
         assert dr.state == "success"
 
+    def test_update_dagrun_note(self, create_runtime_ti, 
mock_supervisor_comms):
+        task = BaseOperator(task_id="hello")
+        runtime_ti = create_runtime_ti(task=task)
+
+        runtime_ti.update_dagrun_note("Updated from task runtime")
+
+        mock_supervisor_comms.send.assert_called_once_with(
+            msg=UpdateDagRunNote(ti_id=runtime_ti.id, note="Updated from task 
runtime")
+        )
+
+    def test_update_dagrun_note_none_skips_request(self, create_runtime_ti, 
mock_supervisor_comms):
+        """A null note is a server-side no-op, so don't spend a round-trip on 
it."""
+        task = BaseOperator(task_id="hello")
+        runtime_ti = create_runtime_ti(task=task)
+
+        runtime_ti.update_dagrun_note(None)
+
+        mock_supervisor_comms.send.assert_not_called()
+
     def test_get_previous_dagrun_with_state(self, create_runtime_ti, 
mock_supervisor_comms):
         """Test that get_previous_dagrun sends the correct request with state 
filter."""
 
diff --git a/ts-sdk/src/generated/supervisor.ts 
b/ts-sdk/src/generated/supervisor.ts
index 1f37c7f3af6..1f9fe56f892 100644
--- a/ts-sdk/src/generated/supervisor.ts
+++ b/ts-sdk/src/generated/supervisor.ts
@@ -613,6 +613,9 @@ export type DagId22 = string;
 export type DagRunId = string;
 export type Type81 = "TriggerDagRun";
 export type TiId9 = string;
+export type Note3 = string | null;
+export type Type82 = "UpdateDagRunNote";
+export type TiId10 = string;
 /**
  * @minItems 1
  */
@@ -620,22 +623,22 @@ export type ChosenOptions = [string, ...string[]];
 export type ParamsInput = {
   [k: string]: unknown;
 } | null;
-export type Type82 = "UpdateHITLDetail";
-export type TiId10 = string;
-export type Type83 = "ValidateInletsAndOutlets";
+export type Type83 = "UpdateHITLDetail";
+export type TiId11 = string;
+export type Type84 = "ValidateInletsAndOutlets";
 export type Keys = string[];
 export type TotalEntries = number;
-export type Type84 = "VariableKeysResult";
+export type Type85 = "VariableKeysResult";
 export type Key19 = string;
 export type Value2 = string | null;
-export type Type85 = "VariableResult";
+export type Type86 = "VariableResult";
 export type Len = number;
-export type Type86 = "XComCountResponse";
+export type Type87 = "XComCountResponse";
 export type Key20 = string;
-export type Type87 = "XComResult";
-export type Type88 = "XComSequenceIndexResult";
+export type Type88 = "XComResult";
+export type Type89 = "XComSequenceIndexResult";
 export type Root = JsonValue[];
-export type Type89 = "XComSequenceSliceResult";
+export type Type90 = "XComSequenceSliceResult";
 
 export interface SupervisorWireSchema {}
 /**
@@ -1866,6 +1869,15 @@ export interface TriggerDagRun {
   run_id: DagRunId;
   type?: Type81;
 }
+/**
+ * This interface was referenced by `SupervisorWireSchema`'s JSON-Schema
+ * via the `definition` "UpdateDagRunNote".
+ */
+export interface UpdateDagRunNote {
+  ti_id: TiId9;
+  note: Note3;
+  type?: Type82;
+}
 /**
  * Update the response content part of an existing Human-in-the-loop response.
  *
@@ -1873,18 +1885,18 @@ export interface TriggerDagRun {
  * via the `definition` "UpdateHITLDetail".
  */
 export interface UpdateHITLDetail {
-  ti_id: TiId9;
+  ti_id: TiId10;
   chosen_options: ChosenOptions;
   params_input?: ParamsInput;
-  type?: Type82;
+  type?: Type83;
 }
 /**
  * This interface was referenced by `SupervisorWireSchema`'s JSON-Schema
  * via the `definition` "ValidateInletsAndOutlets".
  */
 export interface ValidateInletsAndOutlets {
-  ti_id: TiId10;
-  type?: Type83;
+  ti_id: TiId11;
+  type?: Type84;
 }
 /**
  * This interface was referenced by `SupervisorWireSchema`'s JSON-Schema
@@ -1893,7 +1905,7 @@ export interface ValidateInletsAndOutlets {
 export interface VariableKeysResult {
   keys: Keys;
   total_entries: TotalEntries;
-  type?: Type84;
+  type?: Type85;
 }
 /**
  * This interface was referenced by `SupervisorWireSchema`'s JSON-Schema
@@ -1902,7 +1914,7 @@ export interface VariableKeysResult {
 export interface VariableResult {
   key: Key19;
   value: Value2;
-  type?: Type85;
+  type?: Type86;
 }
 /**
  * This interface was referenced by `SupervisorWireSchema`'s JSON-Schema
@@ -1910,7 +1922,7 @@ export interface VariableResult {
  */
 export interface XComCountResponse {
   len: Len;
-  type?: Type86;
+  type?: Type87;
 }
 /**
  * Response to ReadXCom request.
@@ -1921,7 +1933,7 @@ export interface XComCountResponse {
 export interface XComResult {
   key: Key20;
   value: JsonValue | null;
-  type?: Type87;
+  type?: Type88;
 }
 /**
  * This interface was referenced by `SupervisorWireSchema`'s JSON-Schema
@@ -1929,7 +1941,7 @@ export interface XComResult {
  */
 export interface XComSequenceIndexResult {
   root: JsonValue;
-  type?: Type88;
+  type?: Type89;
 }
 /**
  * This interface was referenced by `SupervisorWireSchema`'s JSON-Schema
@@ -1937,7 +1949,7 @@ export interface XComSequenceIndexResult {
  */
 export interface XComSequenceSliceResult {
   root: Root;
-  type?: Type89;
+  type?: Type90;
 }
 
 /** Cadwyn schema version this SDK was generated against.

Reply via email to