jason810496 commented on code in PR #73618:
URL: https://github.com/apache/airflow/pull/73618#discussion_r4193282743
##########
ts-sdk/src/coordinator/client.ts:
##########
@@ -84,9 +100,83 @@ export interface CoordinatorClient extends TaskClient {
getXComEntry(opts: GetXComOpts): Promise<XComEntry>;
}
+const VARIABLE_ABSENT_POLICY: AbsentRowPolicy = {
+ codes: ["VARIABLE_NOT_FOUND"],
+ apiServer404: true,
+};
+const XCOM_ABSENT_POLICY: AbsentRowPolicy = { codes: ["XCOM_NOT_FOUND"],
apiServer404: true };
+const CONNECTION_ABSENT_POLICY: AbsentRowPolicy = {
+ codes: ["CONNECTION_NOT_FOUND"],
+ apiServer404: true,
+};
+const TASK_STATE_STORE_ABSENT_POLICY: AbsentRowPolicy = {
+ codes: ["TASK_STORE_NOT_FOUND"],
+ apiServer404: false,
+};
+
+// Matches Airflow's `[state_store] default_retention_days` config default.
+const FALLBACK_RETENTION_DAYS = 30;
+const RETENTION_DAYS_ENV_VAR = "AIRFLOW__STATE_STORE__DEFAULT_RETENTION_DAYS";
+const MS_PER_DAY = 24 * 60 * 60 * 1000;
+
+function assertKey(key: string): void {
+ if (typeof key !== "string") {
+ throw new TypeError(`task state store key must be a string, got ${typeof
key}`);
+ }
+ if (key === "") {
+ throw new RangeError("task state store key must not be empty");
+ }
+}
+
+// Python's `datetime` (and the wire format it parses) cannot represent a year
+// past 9999, and the supervisor silently drops a frame it cannot decode
+// rather than replying with an error, so the task would hang instead of
+// failing fast. Route both retention paths through this check.
+function toExpiresAt(ms: number): string {
+ const at = new Date(Date.now() + ms);
+ if (Number.isNaN(at.getTime()) || at.getUTCFullYear() > 9999) {
+ throw new RangeError(
+ `retention of ${ms}ms overflows the wire timestamp; use NEVER_EXPIRE
instead`,
+ );
+ }
+ return at.toISOString();
+}
+
+/** Resolve `retentionMs` (or the deployment default) to a wire `expires_at`.
*/
+function resolveExpiresAt(retentionMs: number | undefined): string | null {
+ if (retentionMs === undefined) {
+ return resolveDefaultExpiresAt();
+ }
+ if (retentionMs === NEVER_EXPIRE) {
+ return null;
+ }
+ if (!Number.isFinite(retentionMs) || retentionMs < 0) {
+ throw new RangeError(`retentionMs must be >= 0 or NEVER_EXPIRE, got
${retentionMs}`);
+ }
+ return toExpiresAt(retentionMs);
+}
+
+function resolveDefaultExpiresAt(): string | null {
+ const raw = process.env[RETENTION_DAYS_ENV_VAR];
+ if (raw === undefined || raw === "") {
+ return toExpiresAt(FALLBACK_RETENTION_DAYS * MS_PER_DAY);
+ }
Review Comment:
Ditto as https://github.com/apache/airflow/pull/73420/changes#r4193124778,
let's remove the hardcoded fallback.
--
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]