Repository navigation
fix(sql): cancel remote queries after subprocess timeout #115
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
f49ff94
7574edf
2536ed9
3dc09d6
13f661b
ad49723
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,46 @@ | ||
| /** Race non-fetch boundaries (credential sources and backoff) against cancellation. */ | ||
| export function abortable<T>(work: Promise<T>, signal?: AbortSignal): Promise<T> { | ||
| if (!signal) return work | ||
| return new Promise<T>((resolve, reject) => { | ||
| const abort = () => reject(signal.reason) | ||
| if (signal.aborted) abort() | ||
| else signal.addEventListener("abort", abort, { once: true }) | ||
| work.then(resolve, reject).finally(() => signal.removeEventListener("abort", abort)) | ||
| }) | ||
| } | ||
|
|
||
| export async function delay(ms: number, signal?: AbortSignal): Promise<void> { | ||
| signal?.throwIfAborted() | ||
| return new Promise((resolve, reject) => { | ||
| const finish = () => { | ||
| signal?.removeEventListener("abort", abort) | ||
| resolve() | ||
| } | ||
| const timer = setTimeout(finish, ms) | ||
| const abort = () => { | ||
| clearTimeout(timer) | ||
| signal?.removeEventListener("abort", abort) | ||
| reject(signal?.reason) | ||
| } | ||
| signal?.addEventListener("abort", abort, { once: true }) | ||
| }) | ||
| } | ||
|
|
||
| /** Explicitly own and dispose deadline timers, including across repeated requests. */ | ||
| export function abortAfter(timeoutMs: number, parent?: AbortSignal) { | ||
| const controller = new AbortController() | ||
| const abort = () => controller.abort(parent?.reason) | ||
| if (parent?.aborted) abort() | ||
| else parent?.addEventListener("abort", abort, { once: true }) | ||
| const timer = setTimeout( | ||
| () => controller.abort(new DOMException("The operation timed out", "TimeoutError")), | ||
| timeoutMs, | ||
| ) | ||
| return { | ||
| signal: controller.signal, | ||
| dispose() { | ||
| clearTimeout(timer) | ||
| parent?.removeEventListener("abort", abort) | ||
| }, | ||
| } | ||
| } | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -1,21 +1,79 @@ | ||
| import { request, type ClientOptions } from "../client.js" | ||
| import { requestRaw, type ClientOptions } from "../client.js" | ||
| import { ClickZettaApiError } from "../types/api.js" | ||
| import { abortAfter, delay } from "../abort.js" | ||
| import { isRetryableErrorCode } from "./errors.js" | ||
| import { getJobResultRaw } from "./job-info.js" | ||
| import type { JobID } from "./types.js" | ||
|
|
||
| export async function cancelJob( | ||
| opts: ClientOptions, | ||
| jobId: JobID, | ||
| ): Promise<unknown> { | ||
| const body = { | ||
| export async function cancelJob(opts: ClientOptions, jobId: JobID): Promise<unknown> { | ||
| const response = await requestRaw<unknown>(opts, "/lh/cancelJob", { | ||
| account: { user_id: 0 }, | ||
| job_id: { | ||
| id: jobId.id, | ||
| workspace: jobId.workspace, | ||
| instance_id: jobId.instanceId, | ||
| }, | ||
| job_id: { id: jobId.id, workspace: jobId.workspace, instance_id: jobId.instanceId }, | ||
| user_agent: "", | ||
| force: false, | ||
| }) | ||
| // Coordinator protobuf JSON uses respStatus; some gateways preserve snake_case. | ||
| // Proto3 omits empty fields individually, so an absent status (alone or beside | ||
| // other fields) is success; only a populated error status rejects. | ||
| if (!response || typeof response !== "object" || Array.isArray(response)) { | ||
| throw new ClickZettaApiError("INVALID_CANCEL_RESPONSE", "Invalid cancellation response") | ||
| } | ||
| const raw = response as Record<string, unknown> | ||
| const value = raw.respStatus ?? raw.resp_status | ||
| if (value !== undefined && (!value || typeof value !== "object" || Array.isArray(value))) { | ||
| throw new ClickZettaApiError("INVALID_CANCEL_RESPONSE", "Invalid cancellation status") | ||
| } | ||
| const status = value as Record<string, unknown> | undefined | ||
| const code = status?.errorCode ?? status?.error_code | ||
| const message = status?.errorMsg ?? status?.error_msg | ||
| if (code || message) throw new ClickZettaApiError(String(code || "CANCEL_FAILED"), String(message || code)) | ||
| if (raw.code !== undefined && ![0, "0", 200, "200", "SUCCESS"].includes(raw.code as string | number)) { | ||
| throw new ClickZettaApiError(String(raw.code), String(raw.message ?? raw.msg ?? "Cancellation rejected")) | ||
| } | ||
|
Comment on lines
+30
to
+32
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. MEDIUM (confidence: medium) — the new if (raw.code !== undefined && ![0, "0", 200, "200", "SUCCESS"].includes(raw.code as string | number)) {
throw new ClickZettaApiError(String(raw.code), String(raw.message ?? raw.msg ?? "Cancellation rejected"))
}Two concerns:
(The |
||
| return response | ||
| } | ||
|
|
||
| const resp = await request<unknown>(opts, "/lh/cancelJob", body) | ||
| return resp | ||
| export type CancellationResult = { confirmed: true; state: string } | { confirmed: false; reason: string } | ||
|
|
||
| /** Independent total budget, including credential resolution, requests and retries. */ | ||
| export async function cancelJobAndWait( | ||
| opts: ClientOptions, | ||
| jobId: JobID, | ||
| timeoutMs = 5000, | ||
| ): Promise<CancellationResult> { | ||
| const deadline = abortAfter(timeoutMs) | ||
| const client = { ...opts, signal: deadline.signal, maxRetries: 0 } | ||
| try { | ||
| let reason = "Cancellation was not confirmed before the cleanup deadline" | ||
| while (!client.signal.aborted) { | ||
| // Repeat cancellation to cover a submit that becomes visible after the first cancel. | ||
| // A missing job is not proof of termination: submission may still be in flight. | ||
| await cancelJob(client, jobId).catch((error: unknown) => { | ||
| reason = | ||
| error instanceof ClickZettaApiError ? `Cancellation rejected (${error.code})` : "Cancellation request failed" | ||
| }) | ||
|
Comment on lines
+51
to
+54
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. LOW (confidence: high) — the deadline overwrites await cancelJob(client, jobId).catch((error: unknown) => {
reason =
error instanceof ClickZettaApiError ? `Cancellation rejected (${error.code})` : "Cancellation request failed"
})When the 5 s / 1.5 s budget expires mid-request, That string is what lands in |
||
| const raw = await getJobResultRaw(client, jobId).catch(() => undefined) | ||
| if (raw && typeof raw === "object" && "status" in raw) { | ||
| const response = raw as { | ||
| status?: { state?: string; errorCode?: string } | ||
| respStatus?: { errorCode?: string } | ||
| resp_status?: { error_code?: string } | ||
| } | ||
| const status = response.status | ||
| if ( | ||
| status?.state && | ||
| !response.respStatus?.errorCode && | ||
| !response.resp_status?.error_code && | ||
| !isRetryableErrorCode(status?.errorCode) && | ||
| ["SUCCEED", "FAILED", "CANCELLED"].includes(status?.state ?? "") | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. LOW (confidence: medium) — ["SUCCEED", "FAILED", "CANCELLED"].includes(status?.state ?? "")Two states the SDK demonstrably sees are missing:
The short-spelling-only convention does match
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. MEDIUM — terminal-state set is narrower than the states this repo already knows the gateway sends (confidence: medium-high) ["SUCCEED", "FAILED", "CANCELLED"].includes(status?.state ?? "")This is the only place confirmation can come from, and it omits
Failure scenario: a job finishes and the deployment reports Smaller correct change: reuse one terminal-state predicate rather than a third inline literal — |
||
| ) { | ||
| return { confirmed: true, state: status.state } | ||
| } | ||
| } | ||
| await delay(250, client.signal).catch(() => {}) | ||
|
Comment on lines
+51
to
+73
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. LOW — unconditional re-cancellation every 250 ms makes the request volume scale with the number of orphans (confidence: medium) await cancelJob(client, jobId).catch((error: unknown) => { ... })
const raw = await getJobResultRaw(client, jobId).catch(() => undefined)
...
await delay(250, client.signal).catch(() => {})Each iteration issues two requests, so one unconfirmed job costs ~12 requests at the child's 1.5 s budget and ~40 at the supervisor's 5 s default. The comment justifies the repeat as covering "a submit that becomes visible after the first cancel", which is a real race, but it only needs to be covered for as long as submission could still be in flight — not for the whole budget at full rate. Cancelling once, polling, and re-cancelling only after the state is still non-terminal for a couple of polls would keep the property and cut the volume several-fold. Backing off the 250 ms interval would help too. This compounds with the |
||
| } | ||
| return { confirmed: false, reason } | ||
| } finally { | ||
| deadline.dispose() | ||
| } | ||
| } | ||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
LOW (confidence: high) — the abort listener is only removed when
worksettles, so aworkthat never settles leaks it.Cleanup hangs off
work, not off the race outcome.abortableexists precisely to escape promises that may never resolve —client.tokens.get()against a stalled portal is the motivating case, andtest/fixtures/cancellation.ts:78(get: () => new Promise(() => {})) builds exactly that. In that scenario the outer promise rejects on abort, the caller moves on, and the listener stays attached tosignalfor the signal's lifetime along with the retainedworkand its closure.It's bounded in practice — the long-lived signals here are per-query (
trackSqlJob's deadline) or per-cleanup (cancelJobAndWait's 5 s budget), so they are discarded shortly after. The one place it could accumulate is a long agent session where each supervised job leaks one listener on its own signal, which then goes away with the signal. So: real but low impact.Attaching the removal to the race result instead of to
workfixes it: