jason810496 commented on code in PR #74071:
URL: https://github.com/apache/airflow/pull/74071#discussion_r4192784822


##########
ts-sdk/src/sdk/dag.ts:
##########
@@ -401,10 +405,15 @@ export type TaskFactory<TArgs extends object | void = 
void, TReturn = unknown> =
  */
 export type TaskOptions = TaskSpec;
 
-/** Per-task record a Dag retains: the reference, the handler, and its spec. */
+/** Per-task record a Dag retains: the reference, what runs it, and its spec.
+ *
+ *  Exactly one of `fn` and `trigger` is set: a task has either a handler, or
+ *  the options of a `triggerDagRun(...)` the runtime carries out itself. */
 export interface TaskRecord {
   readonly task: TaskRef;
-  readonly fn: TaskFunction;
+  readonly fn?: TaskFunction;
+  /** What a `triggerDagRun(...)` task triggers, when it has no handler. */
+  readonly trigger?: TriggerDagRunTask;

Review Comment:
   `fn` and `trigger` are two independently-optional fields with only the 
comment above stating "exactly one is set", not a discriminated union. This is 
the root cause behind the `#wrapDecider` bug (a stray `trigger` surviving 
alongside a new `fn`): a shape like `{ kind: "handler", fn } | { kind: 
"trigger", trigger }` would make that spread a compile error instead of a 
silent runtime corruption.
   



##########
ts-sdk/src/coordinator/serde.ts:
##########
@@ -83,6 +84,31 @@ const TASK_MODULE = "airflow.sdk.coordinators.node";
  */
 const TASK_LANGUAGE = "typescript";
 
+/**

Review Comment:
   Stray doubled JSDoc opener, this line should be removed, the next line is 
already a complete one-line comment.
   
   ```suggestion
   ```
   



##########
ts-sdk/src/coordinator/serde.ts:
##########
@@ -83,6 +84,31 @@ const TASK_MODULE = "airflow.sdk.coordinators.node";
  */
 const TASK_LANGUAGE = "typescript";
 
+/**
+/** The Dags this one triggers, for the UI dependency graph. */
+function serializeDagDependencies(dag: Dag): SerializedValue {
+  const dependencies: SerializedValue[] = [];
+  for (const [taskId, record] of getDagTaskRecords(dag)) {
+    if (!record.trigger) continue;
+    dependencies.push({
+      source: dag.dagId,
+      target: record.trigger.dagId,
+      label: taskId,
+      dependency_type: "trigger",
+      dependency_id: taskId,
+    });
+  }
+  return dependencies;
+}
+
+/**

Review Comment:
   Same doubled JSDoc opener as above, before `TRIGGER_DAG_RUN_FIELDS`.
   
   ```suggestion
   ```
   



##########
ts-sdk/src/coordinator/serde.ts:
##########
@@ -196,9 +224,15 @@ function serializeTask(
     is_stub: true,
   };
   const label = `task "${taskId}" of Dag "${dagId}"`;
+  if (record.trigger) {
+    if (record.trigger.conf !== undefined) {
+      toPlainJson(record.trigger.conf, `conf of ${label}`);
+    }
+    Object.assign(data, structuredClone(TRIGGER_DAG_RUN_FIELDS));
+  }

Review Comment:
   This stamps `_operator_name: TriggerDagRunOperator` whenever 
`record.trigger` is set, with no check that `record.canSkipDownstream` (set by 
the `dag.if`/`dag.switch` path) isn't also set on the same record. If the 
`#wrapDecider` bug on the Dag side ever produces a record with both fields set, 
this serializes self-contradictory operator metadata with no error. A defensive 
check here would catch that case even if the root cause elsewhere isn't fixed.
   
   ```suggestion
     if (record.trigger) {
       if (record.canSkipDownstream) {
         throw new Error(`${label} has both a trigger and canSkipDownstream 
set`);
       }
       if (record.trigger.conf !== undefined) {
         toPlainJson(record.trigger.conf, `conf of ${label}`);
       }
       Object.assign(data, structuredClone(TRIGGER_DAG_RUN_FIELDS));
     }
   ```
   



##########
ts-sdk/src/coordinator/trigger-runner.ts:
##########
@@ -0,0 +1,234 @@
+/*!
+ * 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.
+ */
+
+// Mirrors `TriggerDagRunOperator.execute` and `execute_complete` in the 
standard provider.
+
+import { setTimeout as sleep } from "node:timers/promises";
+import type { CoordinatorClient } from "./client.js";
+import type { LogChannel } from "./log-channel.js";
+import type {
+  RuntimeDeferTask,
+  RuntimeRetryTask,
+  RuntimeSucceedTask,
+  RuntimeTaskState,
+  StartupDetails,
+} from "./protocol.js";
+import { getBooleanEnv, type TriggerDagRunTask } from 
"../sdk/trigger-dag-run.js";
+
+export type TriggerOutcome =
+  RuntimeSucceedTask | RuntimeRetryTask | RuntimeTaskState | RuntimeDeferTask;
+
+/** A failure the task's retries apply to, as `_handle_current_task_failed` 
decides. */
+export type FailTask = (message: string) => RuntimeRetryTask | 
RuntimeTaskState;
+
+const DAG_STATE_TRIGGER = 
"airflow.providers.standard.triggers.external_task.DagStateTrigger";
+/** `TriggerDagRunLink().xcom_key`: the XCom the "Triggered DAG" extra link 
reads. */
+const LINK_XCOM_KEY = "_link_TriggerDagRunLink";
+const RUN_ID_XCOM_KEY = "trigger_run_id";
+/** `TRIGGER_FAIL_REPR`: the `next_method` a failed or timed-out trigger 
resumes with. */
+const TRIGGER_FAIL = "__fail__";
+const EXECUTE_COMPLETE = "execute_complete";
+
+export async function runTriggerDagRun(
+  details: StartupDetails,
+  trigger: TriggerDagRunTask,
+  client: CoordinatorClient,
+  logs: LogChannel,
+  signal: AbortSignal,
+  fail: FailTask,
+): Promise<TriggerOutcome> {
+  const nextMethod = details.ti_context.next_method;
+  if (nextMethod) return resume(nextMethod, details.ti_context.next_kwargs, 
trigger, logs, fail);
+
+  const logicalDate = new Date();
+  const runId = trigger.runId ?? `manual__${pythonIsoformat(logicalDate)}`;
+
+  if (trigger.failWhenDagIsPaused && (await 
client.isDagPaused(trigger.dagId))) {
+    return fail(`Dag ${trigger.dagId} is paused`);
+  }
+
+  await client.setXCom({ key: LINK_XCOM_KEY, value: dagRunUrl(trigger.dagId, 
runId) });
+
+  logs.info("Triggering Dag Run.", { trigger_dag_id: trigger.dagId });
+  const triggered = await client.triggerDagRun({
+    dag_id: trigger.dagId,
+    run_id: runId,
+    logical_date: logicalDate.toISOString(),
+    run_after: null,
+    conf: (trigger.conf as Record<string, unknown> | undefined) ?? null,
+    reset_dag_run: trigger.resetDagRun,
+    note: trigger.note ?? null,
+  });

Review Comment:
   This `setXCom` (the Triggered-DAG link) and the `triggerDagRun` call below 
it are independent RPCs, the XCom value depends only on `trigger.dagId` and the 
locally generated `runId`, not on the trigger call's result, but they run as 
two sequential awaits.
   
   ```suggestion
     logs.info("Triggering Dag Run.", { trigger_dag_id: trigger.dagId });
     const [, triggered] = await Promise.all([
       client.setXCom({ key: LINK_XCOM_KEY, value: dagRunUrl(trigger.dagId, 
runId) }),
       client.triggerDagRun({
         dag_id: trigger.dagId,
         run_id: runId,
         logical_date: logicalDate.toISOString(),
         run_after: null,
         conf: (trigger.conf as Record<string, unknown> | undefined) ?? null,
         reset_dag_run: trigger.resetDagRun,
         note: trigger.note ?? null,
       }),
     ]);
   ```
   



##########
ts-sdk/src/coordinator/client.ts:
##########
@@ -202,6 +214,46 @@ export function createCoordinatorClient(
       await rpc("SetXCom", null, msg, () => undefined, "throw");
     },
 
+    // ---- Dag runs ----
+
+    async triggerDagRun(msg: Omit<TriggerDagRun, "type">) {

Review Comment:
   `triggerDagRun` bypasses the shared `rpc()` helper every other method on 
this client uses (`getVariable`, `setXCom`, `getDagRunState`, etc.), 
hand-rolling its own `comm.request` + `parseFrameError` + logging + throw. A 
future change to logging format or error-message wording in `rpc()` is picked 
up automatically everywhere except here.
   



##########
ts-sdk/src/coordinator/trigger-runner.ts:
##########
@@ -0,0 +1,234 @@
+/*!
+ * 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.
+ */
+
+// Mirrors `TriggerDagRunOperator.execute` and `execute_complete` in the 
standard provider.
+
+import { setTimeout as sleep } from "node:timers/promises";
+import type { CoordinatorClient } from "./client.js";
+import type { LogChannel } from "./log-channel.js";
+import type {
+  RuntimeDeferTask,
+  RuntimeRetryTask,
+  RuntimeSucceedTask,
+  RuntimeTaskState,
+  StartupDetails,
+} from "./protocol.js";
+import { getBooleanEnv, type TriggerDagRunTask } from 
"../sdk/trigger-dag-run.js";
+
+export type TriggerOutcome =
+  RuntimeSucceedTask | RuntimeRetryTask | RuntimeTaskState | RuntimeDeferTask;
+
+/** A failure the task's retries apply to, as `_handle_current_task_failed` 
decides. */
+export type FailTask = (message: string) => RuntimeRetryTask | 
RuntimeTaskState;
+
+const DAG_STATE_TRIGGER = 
"airflow.providers.standard.triggers.external_task.DagStateTrigger";
+/** `TriggerDagRunLink().xcom_key`: the XCom the "Triggered DAG" extra link 
reads. */
+const LINK_XCOM_KEY = "_link_TriggerDagRunLink";
+const RUN_ID_XCOM_KEY = "trigger_run_id";
+/** `TRIGGER_FAIL_REPR`: the `next_method` a failed or timed-out trigger 
resumes with. */
+const TRIGGER_FAIL = "__fail__";
+const EXECUTE_COMPLETE = "execute_complete";
+
+export async function runTriggerDagRun(
+  details: StartupDetails,
+  trigger: TriggerDagRunTask,
+  client: CoordinatorClient,
+  logs: LogChannel,
+  signal: AbortSignal,
+  fail: FailTask,
+): Promise<TriggerOutcome> {
+  const nextMethod = details.ti_context.next_method;
+  if (nextMethod) return resume(nextMethod, details.ti_context.next_kwargs, 
trigger, logs, fail);
+
+  const logicalDate = new Date();
+  const runId = trigger.runId ?? `manual__${pythonIsoformat(logicalDate)}`;
+
+  if (trigger.failWhenDagIsPaused && (await 
client.isDagPaused(trigger.dagId))) {
+    return fail(`Dag ${trigger.dagId} is paused`);
+  }
+
+  await client.setXCom({ key: LINK_XCOM_KEY, value: dagRunUrl(trigger.dagId, 
runId) });
+
+  logs.info("Triggering Dag Run.", { trigger_dag_id: trigger.dagId });
+  const triggered = await client.triggerDagRun({
+    dag_id: trigger.dagId,
+    run_id: runId,
+    logical_date: logicalDate.toISOString(),
+    run_after: null,
+    conf: (trigger.conf as Record<string, unknown> | undefined) ?? null,
+    reset_dag_run: trigger.resetDagRun,
+    note: trigger.note ?? null,
+  });
+  if (triggered === "already_exists") {
+    if (trigger.skipWhenAlreadyExists) {
+      logs.info(
+        "Dag Run already exists, skipping task as skip_when_already_exists is 
set to True.",
+        { dag_id: trigger.dagId },
+      );
+      return { type: "TaskState", state: "skipped", end_date: new 
Date().toISOString() };
+    }
+    logs.error("Dag Run already exists, marking task as failed.", { dag_id: 
trigger.dagId });
+    return { type: "TaskState", state: "failed", end_date: new 
Date().toISOString() };
+  }
+  logs.info("Dag Run triggered successfully.", { trigger_dag_id: trigger.dagId 
});
+  await client.setXCom({ key: RUN_ID_XCOM_KEY, value: runId });
+
+  if (!trigger.waitForCompletion) {
+    if (trigger.deferrable) {
+      logs.info(
+        "Ignoring deferrable=True because wait_for_completion=False. " +
+          "Task will complete immediately without waiting for the triggered 
DAG run.",
+        { trigger_dag_id: trigger.dagId },
+      );
+    }
+    return succeeded();
+  }
+
+  if (trigger.deferrable) {
+    logs.info("Pausing task as DEFERRED.", { trigger_dag_id: trigger.dagId, 
run_id: runId });
+    return {
+      type: "DeferTask",
+      state: "deferred",
+      classpath: DAG_STATE_TRIGGER,
+      // `DagStateTrigger.serialize()`, key for key.
+      trigger_kwargs: {
+        dag_id: trigger.dagId,
+        states: [...trigger.allowedStates, ...trigger.failedStates],
+        poll_interval: trigger.pokeInterval,
+        run_ids: [runId],
+        execution_dates: null,
+      },
+      trigger_timeout: null,
+      // `_defer_task` hands the trigger the task's queue only when triggerer 
queues are enabled.
+      queue: getBooleanEnv("AIRFLOW__TRIGGERER__QUEUES_ENABLED", false)
+        ? (details.ti.queue ?? null)
+        : null,
+      next_method: EXECUTE_COMPLETE,
+      next_kwargs: {},
+    };
+  }
+
+  while (true) {
+    logs.info("Waiting for dag run to complete execution in allowed state.", {
+      dag_id: trigger.dagId,
+      run_id: runId,
+      allowed_state: trigger.allowedStates,
+    });
+    await sleep(trigger.pokeInterval * 1000, undefined, { signal });
+    const state = await client.getDagRunState(trigger.dagId, runId);
+    if (includes(trigger.failedStates, state)) {
+      logs.error("DagRun finished with failed state.", { dag_id: 
trigger.dagId, state });
+      return fail(`${trigger.dagId} failed with failed state ${state}`);
+    }
+    if (includes(trigger.allowedStates, state)) {
+      logs.info("DagRun finished with allowed state.", { dag_id: 
trigger.dagId, state });
+      return succeeded();
+    }
+    logs.debug("DagRun not yet in allowed or failed state.", { dag_id: 
trigger.dagId, state });
+  }
+}
+
+/** `BaseOperator.resume_execution` for this task: `__fail__` or 
`execute_complete`. */
+function resume(
+  nextMethod: string,
+  nextKwargs: unknown,
+  trigger: TriggerDagRunTask,
+  logs: LogChannel,
+  fail: FailTask,
+): TriggerOutcome {
+  const kwargs = isRecord(nextKwargs) ? nextKwargs : {};
+  if (nextMethod === TRIGGER_FAIL) {
+    const traceback = kwargs["traceback"];
+    if (Array.isArray(traceback)) logs.error(`Trigger 
failed:\n${traceback.join("\n")}`);
+    return fail(String(kwargs["error"] ?? "Unknown"));
+  }
+  if (nextMethod !== EXECUTE_COMPLETE) {
+    return fail(`Task cannot resume with next_method "${nextMethod}"`);
+  }
+  const eventData = decodeEvent(kwargs["event"]);
+  const runIds = eventData?.["run_ids"];
+  if (eventData === undefined || !Array.isArray(runIds)) {
+    return fail(`Task resumed with an event it cannot read: 
${JSON.stringify(kwargs["event"])}`);
+  }
+  const failedRunIds: string[] = [];
+  for (const runId of runIds) {
+    const state = eventData[String(runId)];
+    if (includes(trigger.failedStates, state)) {
+      failedRunIds.push(String(runId));
+      continue;
+    }
+    if (includes(trigger.allowedStates, state)) {
+      logs.info("Triggered Dag run finished with allowed state.", {
+        dag_id: trigger.dagId,
+        state,
+        run_id: runId,
+      });
+    }
+  }
+  if (failedRunIds.length > 0) {
+    return fail(
+      `${trigger.dagId} failed with failed states 
${JSON.stringify(trigger.failedStates)} ` +
+        `for run_ids ${JSON.stringify(failedRunIds)}`,
+    );
+  }
+  return succeeded();
+}
+
+/**
+ * The payload `DagStateTrigger` fired with: `(classpath, data)`. The triggerer
+ * stores it with serde, which encodes a tuple as `{__classname__, __data__}`.
+ */
+function decodeEvent(event: unknown): Record<string, unknown> | undefined {
+  const pair =
+    isRecord(event) && event["__classname__"] === "builtins.tuple" ? 
event["__data__"] : event;
+  if (!Array.isArray(pair) || pair.length !== 2 || !isRecord(pair[1])) return 
undefined;
+  return pair[1];
+}
+
+function succeeded(): RuntimeSucceedTask {
+  return {
+    type: "SucceedTask",
+    end_date: new Date().toISOString(),
+    task_outlets: [],
+    outlet_events: [],
+  };
+}
+
+/** `datetime.isoformat()` of a UTC instant, which is how Python spells a run 
ID's date. */
+export function pythonIsoformat(date: Date): string {
+  const iso = date.toISOString();
+  const seconds = iso.slice(0, 19);
+  const millis = date.getUTCMilliseconds();
+  const fraction = millis === 0 ? "" : `.${String(millis * 1000).padStart(6, 
"0")}`;
+  return `${seconds}${fraction}+00:00`;
+}
+
+/** `build_airflow_dagrun_url`, on `[api] base_url` from the environment, or 
"/" when unset. */
+export function dagRunUrl(dagId: string, runId: string): string {
+  const base = process.env["AIRFLOW__API__BASE_URL"] || "/";
+  return `${base.replace(/\/+$/, "")}/dags/${dagId}/runs/${runId}`;
+}
+
+function includes(states: readonly string[], state: unknown): boolean {
+  return typeof state === "string" && states.includes(state);
+}
+
+function isRecord(value: unknown): value is Record<string, unknown> {

Review Comment:
   This reimplements the already-exported `isPlainRecord` from `sdk/dag.ts`. 
`sdk/trigger-dag-run.ts` also inlines the same object-shape check twice more, 
four slightly different versions of the same predicate now exist across the 
codebase, and a future fix to one (e.g. handling `Object.create(null)`) won't 
propagate to the rest.
   



##########
ts-sdk/src/sdk/trigger-dag-run.ts:
##########
@@ -0,0 +1,194 @@
+/*!
+ * 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.
+ */
+
+import { brand, hasBrand } from "./brand.js";
+import type { JsonValue } from "./client-types.js";
+
+/** A Dag run state, as `allowedStates` and `failedStates` name one. */
+export type DagRunState = "queued" | "running" | "success" | "failed";
+
+const DAG_RUN_STATES: ReadonlySet<string> = new Set<DagRunState>([
+  "queued",
+  "running",
+  "success",
+  "failed",
+]);
+
+/** The `TriggerDagRunOperator` options, named as TypeScript spells them. */
+export interface TriggerDagRunSpec {
+  /** Identifier of the Dag to trigger. */
+  readonly dagId: string;
+  /** Run ID for the triggered run; generated from the trigger time when 
unset. */
+  readonly runId?: string;
+  /** Configuration the triggered run is started with. */
+  readonly conf?: Readonly<Record<string, JsonValue>>;
+  /** Clear an existing run with the same ID instead of failing. */
+  readonly resetDagRun?: boolean;
+  /** Hold this task open until the triggered run finishes. */
+  readonly waitForCompletion?: boolean;
+  /** Seconds between checks while waiting. Defaults to 60. */
+  readonly pokeInterval?: number;
+  /** Run states that count as success when waiting. Defaults to 
`["success"]`. */
+  readonly allowedStates?: readonly DagRunState[];
+  /** Run states that count as failure when waiting. Defaults to `["failed"]`. 
*/
+  readonly failedStates?: readonly DagRunState[];
+  /** Skip rather than fail when the run already exists. */
+  readonly skipWhenAlreadyExists?: boolean;
+  /** Fail rather than trigger when the target Dag is paused. */
+  readonly failWhenDagIsPaused?: boolean;
+  /** Note recorded against the triggered run. */
+  readonly note?: string;
+  /**
+   * While waiting, free the worker slot and defer to `DagStateTrigger`, which
+   * the Python triggerer runs. Defaults to `[operators] default_deferrable`, 
or
+   * to false when that is unset.
+   */
+  readonly deferrable?: boolean;
+}
+
+/** Internal: a trigger task's options with `TriggerDagRunOperator`'s defaults 
applied. */
+export interface TriggerDagRunTask {
+  readonly dagId: string;
+  readonly runId: string | undefined;
+  readonly conf: Readonly<Record<string, JsonValue>> | undefined;
+  readonly resetDagRun: boolean;
+  readonly waitForCompletion: boolean;
+  readonly pokeInterval: number;
+  readonly allowedStates: readonly DagRunState[];
+  readonly failedStates: readonly DagRunState[];
+  readonly skipWhenAlreadyExists: boolean;
+  readonly failWhenDagIsPaused: boolean;
+  readonly note: string | undefined;
+  readonly deferrable: boolean;
+}
+
+/** Internal: whether `value` is a trigger task built by any copy of this 
package. */
+export function isTriggerDagRunTask(value: unknown): value is 
TriggerDagRunTask {
+  return hasBrand(value, "TriggerDagRunTask");
+}
+
+const OPTION_NAMES: ReadonlySet<string> = new Set<keyof TriggerDagRunSpec>([
+  "dagId",
+  "runId",
+  "conf",
+  "resetDagRun",
+  "waitForCompletion",
+  "pokeInterval",
+  "allowedStates",
+  "failedStates",
+  "skipWhenAlreadyExists",
+  "failWhenDagIsPaused",
+  "note",
+  "deferrable",
+]);
+
+function checkStates(name: string, states: unknown): DagRunState[] | undefined 
{
+  if (states === undefined) return undefined;
+  if (!Array.isArray(states)) {
+    throw new Error(`triggerDagRun(...) option "${name}" must be an array of 
Dag run states`);
+  }
+  for (const state of states) {
+    if (typeof state !== "string" || !DAG_RUN_STATES.has(state)) {
+      throw new Error(
+        `triggerDagRun(...) option "${name}" holds ${JSON.stringify(state)}, 
which is not a Dag ` +
+          `run state; use one of ${[...DAG_RUN_STATES].join(", ")}`,
+      );
+    }
+  }
+  return [...(states as DagRunState[])];
+}
+
+function checkType(name: string, value: unknown, type: "string" | "boolean"): 
void {
+  if (value !== undefined && typeof value !== type) {
+    throw new Error(`triggerDagRun(...) option "${name}" must be a ${type}`);
+  }
+}
+
+/**

Review Comment:
   Same doubled JSDoc opener, before `getBooleanEnv`.
   
   ```suggestion
   ```
   



##########
ts-sdk/src/sdk/trigger-dag-run.ts:
##########
@@ -0,0 +1,194 @@
+/*!
+ * 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.
+ */
+
+import { brand, hasBrand } from "./brand.js";
+import type { JsonValue } from "./client-types.js";
+
+/** A Dag run state, as `allowedStates` and `failedStates` name one. */
+export type DagRunState = "queued" | "running" | "success" | "failed";
+
+const DAG_RUN_STATES: ReadonlySet<string> = new Set<DagRunState>([
+  "queued",
+  "running",
+  "success",
+  "failed",
+]);
+
+/** The `TriggerDagRunOperator` options, named as TypeScript spells them. */
+export interface TriggerDagRunSpec {
+  /** Identifier of the Dag to trigger. */
+  readonly dagId: string;
+  /** Run ID for the triggered run; generated from the trigger time when 
unset. */
+  readonly runId?: string;
+  /** Configuration the triggered run is started with. */
+  readonly conf?: Readonly<Record<string, JsonValue>>;
+  /** Clear an existing run with the same ID instead of failing. */
+  readonly resetDagRun?: boolean;
+  /** Hold this task open until the triggered run finishes. */
+  readonly waitForCompletion?: boolean;
+  /** Seconds between checks while waiting. Defaults to 60. */
+  readonly pokeInterval?: number;
+  /** Run states that count as success when waiting. Defaults to 
`["success"]`. */
+  readonly allowedStates?: readonly DagRunState[];
+  /** Run states that count as failure when waiting. Defaults to `["failed"]`. 
*/
+  readonly failedStates?: readonly DagRunState[];
+  /** Skip rather than fail when the run already exists. */
+  readonly skipWhenAlreadyExists?: boolean;
+  /** Fail rather than trigger when the target Dag is paused. */
+  readonly failWhenDagIsPaused?: boolean;
+  /** Note recorded against the triggered run. */
+  readonly note?: string;
+  /**
+   * While waiting, free the worker slot and defer to `DagStateTrigger`, which
+   * the Python triggerer runs. Defaults to `[operators] default_deferrable`, 
or
+   * to false when that is unset.
+   */
+  readonly deferrable?: boolean;
+}
+
+/** Internal: a trigger task's options with `TriggerDagRunOperator`'s defaults 
applied. */
+export interface TriggerDagRunTask {
+  readonly dagId: string;
+  readonly runId: string | undefined;
+  readonly conf: Readonly<Record<string, JsonValue>> | undefined;
+  readonly resetDagRun: boolean;
+  readonly waitForCompletion: boolean;
+  readonly pokeInterval: number;
+  readonly allowedStates: readonly DagRunState[];
+  readonly failedStates: readonly DagRunState[];
+  readonly skipWhenAlreadyExists: boolean;
+  readonly failWhenDagIsPaused: boolean;
+  readonly note: string | undefined;
+  readonly deferrable: boolean;
+}
+
+/** Internal: whether `value` is a trigger task built by any copy of this 
package. */
+export function isTriggerDagRunTask(value: unknown): value is 
TriggerDagRunTask {
+  return hasBrand(value, "TriggerDagRunTask");
+}
+
+const OPTION_NAMES: ReadonlySet<string> = new Set<keyof TriggerDagRunSpec>([
+  "dagId",
+  "runId",
+  "conf",
+  "resetDagRun",
+  "waitForCompletion",
+  "pokeInterval",
+  "allowedStates",
+  "failedStates",
+  "skipWhenAlreadyExists",
+  "failWhenDagIsPaused",
+  "note",
+  "deferrable",
+]);
+
+function checkStates(name: string, states: unknown): DagRunState[] | undefined 
{
+  if (states === undefined) return undefined;
+  if (!Array.isArray(states)) {
+    throw new Error(`triggerDagRun(...) option "${name}" must be an array of 
Dag run states`);
+  }
+  for (const state of states) {
+    if (typeof state !== "string" || !DAG_RUN_STATES.has(state)) {
+      throw new Error(
+        `triggerDagRun(...) option "${name}" holds ${JSON.stringify(state)}, 
which is not a Dag ` +
+          `run state; use one of ${[...DAG_RUN_STATES].join(", ")}`,
+      );
+    }
+  }
+  return [...(states as DagRunState[])];
+}
+
+function checkType(name: string, value: unknown, type: "string" | "boolean"): 
void {
+  if (value !== undefined && typeof value !== type) {
+    throw new Error(`triggerDagRun(...) option "${name}" must be a ${type}`);
+  }
+}
+
+/**
+/** Internal: a boolean Airflow option from the environment, read as 
`conf.getboolean` does. */
+export function getBooleanEnv(name: string, fallback: boolean): boolean {
+  const raw = process.env[name];
+  if (raw === undefined) return fallback;
+  const value = raw.trim().toLowerCase();
+  if (value === "t" || value === "true" || value === "1") return true;
+  if (value === "f" || value === "false" || value === "0") return false;
+  throw new Error(`${name} is ${JSON.stringify(raw)}, which is not a boolean; 
use true or false`);
+}
+
+/**

Review Comment:
   Same artifact, here two consecutive `/**` lines open the `triggerDagRun` doc 
block, one of them should go.
   
   ```suggestion
   ```
   



##########
ts-sdk/src/coordinator/serde.ts:
##########
@@ -140,12 +166,13 @@ export function serializeDag(
       serializeTask(
         dag.dagId,
         taskId,
-        withDagQueue(record.spec, dag.spec.queue),
+        record,
         graph.downstreamTaskIds.get(taskId),
         inputs.get(taskId),
+        dag.spec.queue,
       ),
     ),
-    dag_dependencies: [],
+    dag_dependencies: serializeDagDependencies(dag),

Review Comment:
   Not a blocking one, can be follow-up to improve the overall serialization 
overhead after this series.
   As the change scope might be a be large for this one.
   
   ---
   
   `serializeDagDependencies` does its own full loop over every task record 
right after the `tasks:` mapping a few lines below already iterated the same 
records. For Dags with many tasks this walks the task list twice, 
`dag_dependencies` could instead be collected as a side effect of the existing 
`tasks.map()` pass.
   



-- 
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]

Reply via email to