potiuk commented on code in PR #74204:
URL: https://github.com/apache/airflow/pull/74204#discussion_r4184907581
##########
providers/elasticsearch/src/airflow/providers/elasticsearch/log/es_task_handler.py:
##########
@@ -226,6 +230,29 @@ def _render_log_id(log_id_template: str, ti: TaskInstance
| TaskInstanceKey, try
)
+def _get_ti_id_fields(ti: TaskInstance | TaskInstanceKey) -> dict[str, str]:
+ # Airflow 2 task instances have no id.
+ return {"ti_id": str(ti_id)} if (ti_id := getattr(ti, "id", None)) else {}
+
+
+def _build_log_query(log_id: str, ti: RuntimeTI) -> list[dict[str, Any]]:
+ log_id_match = {"match_phrase": {"log_id": log_id}}
+ # Before 3.4 a cleared task instance gets a new id, which can differ from
the id its logs were written under.
+ if not AIRFLOW_V_3_4_PLUS:
+ return [log_id_match]
+ return [
+ {
+ "bool": {
+ "should": [
+ {"match_phrase": {"ti_id": str(ti.id)}},
Review Comment:
Adding one more narrow case to this thread. Before 3.4, a try in
`UP_FOR_RETRY` is stored as `(B, N)` (see #73554), while its logs carry the
running UUID A. If a deployment upgrades to 3.4 while a task is in that state,
viewing try N looks up `ti_id=B`, and because those lines do carry `ti_id`, the
`log_id` fallback skips them. The logs reappear once the retry runs and the try
is archived with A, so this is temporary, unlike the case above. Approving on
my side either way; I'll leave this thread for you two to settle.
---
Drafted-by: Claude Code (Opus 5.5); reviewed by @potiuk before posting
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]