Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
46 changes: 46 additions & 0 deletions packages/clickzetta-sdk/src/abort.ts
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))
Comment on lines +4 to +8

Copy link
Copy Markdown
Contributor

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 work settles, so a work that never settles leaks it.

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))
})

Cleanup hangs off work, not off the race outcome. abortable exists precisely to escape promises that may never resolve — client.tokens.get() against a stalled portal is the motivating case, and test/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 to signal for the signal's lifetime along with the retained work and 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 work fixes it:

Suggested change
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))
return new Promise<T>((resolve, reject) => {
const abort = () => reject(signal.reason)
const done = () => signal.removeEventListener("abort", abort)
if (signal.aborted) abort()
else signal.addEventListener("abort", abort, { once: true })
work.then(
(value) => {
done()
resolve(value)
},
(error) => {
done()
reject(error)
},
)
})

})
}

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)
},
}
}
48 changes: 29 additions & 19 deletions packages/clickzetta-sdk/src/client.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import { abortable, abortAfter, delay } from "./abort.js"
import { ClickZettaApiError, type ApiResponse } from "./types/api.js"
import type { Credential, RequestContext, TokenSource } from "./types/index.js"
import { currentTraceparent } from "./traceparent.js"
Expand Down Expand Up @@ -28,6 +29,9 @@ export interface ClientOptions {
customHeaders?: Record<string, string>
traceparent?: string
timeout?: number
/** Cancellation covers requests, credential waits and retry backoff. */
signal?: AbortSignal
maxRetries?: number
/** Non-auth metadata some request bodies embed — see {@link RequestContext}. */
context?: RequestContext
}
Expand All @@ -41,10 +45,6 @@ export function retryDelayMs(attempt: number): number {
return base + Math.random() * 500
}

function sleep(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms))
}

/**
* Generate a request id matching the Python connector format
* (`pysdk-v{version}-{uuid12}`, client.py:292). The server uses this for
Expand Down Expand Up @@ -72,12 +72,12 @@ function buildHeaders(opts: ClientOptions, credential: Credential): Record<strin
return mergeHeaders(
{
"Content-Type": "application/json",
"Accept": "application/json, text/plain, */*",
Accept: "application/json, text/plain, */*",
"User-Agent": `tssdk/${SDK_VERSION}`,
// client.py:293 — trace id header, required by the gateway for correlation
"requestId": requestId,
requestId: requestId,
"X-Request-ID": requestId,
"traceparent": opts.traceparent ?? currentTraceparent(),
traceparent: opts.traceparent ?? currentTraceparent(),
...(instanceName ? { instanceName } : {}),
},
// The credential's own headers sit under the caller's: a Cookie belongs to
Expand Down Expand Up @@ -105,7 +105,8 @@ async function doRequest<T>(
parseWrapper: boolean,
): Promise<T> {
const url = `${opts.baseUrl}${path}`
let credential = await opts.tokens.get()
opts.signal?.throwIfAborted()
let credential = await abortable(opts.tokens.get(), opts.signal)
let headers = buildHeaders(opts, credential)
// Credential a rotation just produced, to be used verbatim by the next attempt.
// `TokenSource.rotate` is contracted to RETURN the replacement, not to make
Expand All @@ -123,21 +124,23 @@ async function doRequest<T>(
let authExhausted = false

let lastError: Error | undefined
for (let attempt = 0; attempt <= MAX_RETRIES; attempt++) {
for (let attempt = 0; attempt <= (opts.maxRetries ?? MAX_RETRIES); attempt++) {
const deadline = abortAfter(opts.timeout ?? DEFAULT_TIMEOUT_MS, opts.signal)
try {
if (attempt > 0) {
// Prefer the rotated credential; otherwise re-resolve, because a retry
// after a multi-second backoff must not resend one that expired while
// we waited. `get()` is cache-backed, so re-resolving costs nothing.
credential = rotatedCredential ?? await opts.tokens.get()
opts.signal?.throwIfAborted()
credential = rotatedCredential ?? (await abortable(opts.tokens.get(), opts.signal))
rotatedCredential = undefined
headers = buildHeaders(opts, credential)
}
const resp = await fetch(url, {
method,
headers,
body: body !== undefined ? JSON.stringify(body) : undefined,
signal: AbortSignal.timeout(opts.timeout ?? DEFAULT_TIMEOUT_MS),
signal: deadline.signal,
})
const text = await resp.text()
if (!resp.ok) {
Expand All @@ -150,9 +153,10 @@ async function doRequest<T>(
if (resp.status === AUTH_EXPIRED_STATUS) {
// Rotation is offered once; a source that cannot rotate (or a second
// rejection) makes this 401 the final answer for this identity.
const fresh = rotated || attempt >= MAX_RETRIES
? undefined
: await opts.tokens.rotate(credential)
const fresh =
rotated || attempt >= (opts.maxRetries ?? MAX_RETRIES)
? undefined
: await abortable(opts.tokens.rotate(credential), opts.signal)
if (!fresh) {
authExhausted = true
throw apiErr
Expand All @@ -169,16 +173,22 @@ async function doRequest<T>(
throw new ClickZettaApiError("PARSE_ERROR", `Invalid JSON response: ${text.slice(0, 200)}`, 0)
}
} catch (err) {
opts.signal?.throwIfAborted()
lastError = err instanceof Error ? err : new Error(String(err))
if (err instanceof ClickZettaApiError && (NON_RETRYABLE_STATUS.has(err.statusCode ?? 0) || err.code === "PARSE_ERROR")) {
if (
err instanceof ClickZettaApiError &&
(NON_RETRYABLE_STATUS.has(err.statusCode ?? 0) || err.code === "PARSE_ERROR")
) {
throw err
}
if (authExhausted) throw err
if (TERMINAL_ERROR_CODES.has(String((err as { code?: unknown }).code ?? ""))) throw err
if (attempt < MAX_RETRIES) {
await sleep(retryDelayMs(attempt))
if (attempt < (opts.maxRetries ?? MAX_RETRIES)) {
await delay(retryDelayMs(attempt), opts.signal)
continue
}
} finally {
deadline.dispose()
}
}
// Signal to callers: parseWrapper is unused on the error path but
Expand All @@ -193,7 +203,7 @@ export async function request<T>(
body?: unknown,
method: string = "POST",
): Promise<ApiResponse<T>> {
return doRequest<ApiResponse<T>>(options, path, body, method, true)
return abortable(doRequest<ApiResponse<T>>(options, path, body, method, true), options.signal)
}

/**
Expand All @@ -207,5 +217,5 @@ export async function requestRaw<T = unknown>(
body?: unknown,
method: string = "POST",
): Promise<T> {
return doRequest<T>(options, path, body, method, false)
return abortable(doRequest<T>(options, path, body, method, false), options.signal)
}
2 changes: 2 additions & 0 deletions packages/clickzetta-sdk/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -47,3 +47,5 @@ export const threadsafety = 2
/** DB-API 2.0 parameter style (dbapi.py:29). */
export const paramstyle = "qmark"
export { czStruct } from "./sql/converter.js"

export { abortable, abortAfter } from "./abort.js"
84 changes: 71 additions & 13 deletions packages/clickzetta-sdk/src/sql/cancel.ts
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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

MEDIUM (confidence: medium) — the new raw.code allowlist changes cz-cli job cancel output and will reject some success responses.

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:

  1. code: null throws. null !== undefined passes the guard and is not in the allowlist, so a gateway that serialises an absent code as null (common for JSON envelopes that don't use proto3 omission) yields ClickZettaApiError("null", "Cancellation rejected"). The respStatus check directly above is careful to treat an absent status as success; this one is not. test/fixtures/cancellation.ts:23 covers { code: 0, data: {} } but not null.

  2. Changed behavior for an existing command. packages/cz-cli/src/commands/job.ts:399 calls cancelJob and, on success, prints { job_id, cancelled: true }. With this validation a coordinator rejection now throws, so that call site emits JOB_CANCEL_ERROR with exit 1 where it previously reported cancelled: true and exit 0. That is almost certainly the better behavior, but it is an output-shape change for cz-cli job cancel that scripts may parse, it is unrelated to the PR's stated purpose, and nothing in the test suite covers the job cancel command path.

(The request → requestRaw switch itself is inert — doRequest's parseWrapper is unused, void parseWrapper at client.ts:196 — so the only behavior delta here is this validation.)

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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LOW (confidence: high) — the deadline overwrites reason with a wrong attribution, and the diagnostic record is the only place it is ever read.

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, client.signal aborts, cancelJob rejects with the TimeoutError DOMException, and this handler records "Cancellation request failed" — indistinguishable from a genuine transport failure. The fallback reason initialised above ("not confirmed before the cleanup deadline") is the accurate one in that case but gets clobbered on the final iteration.

That string is what lands in ~/.clickzetta/sql-cleanup.jsonl via createSqlSupervisor's onWarning, which per specs/sql-cancellation.md is the sole record of an unconfirmed cleanup — so the one artifact meant for post-mortem diagnosis systematically mislabels timeouts as request failures. Skipping the assignment when client.signal.aborted would preserve the real cause.

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 ?? "")

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LOW (confidence: medium) — CANCELLING and SUCCEEDED never confirm, so the common cancel path burns the full budget and logs a false warning.

["SUCCEED", "FAILED", "CANCELLED"].includes(status?.state ?? "")

Two states the SDK demonstrably sees are missing:

  • CANCELLING — poll.ts:206-208 maps case "CANCELLED": case "CANCELLING": to JobStatus.CANCELLED, i.e. the SDK already knows the coordinator reports this intermediate state. Since a successful cancel transitions RUNNING → CANCELLING → CANCELLED, a job that is still cancelling at the deadline returns confirmed: false and writes an unconfirmed-cleanup record to sql-cleanup.jsonl — even though cancellation is observably in progress. With the child budget at 1.5 s and a 250 ms poll interval, that will not be rare, and it makes the diagnostic log noisy enough to stop being useful.
  • SUCCEEDED — poll.ts:197-200 handles case "SUCCEEDED": case "SUCCEED":. Here only the short spelling confirms.

The short-spelling-only convention does match poll.ts:10's TERMINAL_STATES, session.ts:497 and volume.ts:593, so this is consistent with existing code rather than newly wrong — but toJobStatus is the function that actually normalises wire states, and it accepts both. Reusing a shared helper (or at least TERMINAL_STATES) instead of a fourth inline copy of the list would keep these from drifting; adding CANCELLING as "confirmed enough for cleanup" is the part with a visible payoff.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The 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 SUCCEEDED plus the lowercase forms. Two places in-tree say those reach us:

  • packages/clickzetta-sdk/src/sql/poll.ts:198-200 maps both "SUCCEEDED" and "SUCCEED" to JobStatus.SUCCEEDED.
  • packages/cz-cli/src/commands/job-profile.ts:179 treats ["SUCCEED", "SUCCEEDED", "FAILED", "CANCELLED", "succeed", "failed", "cancelled"] as terminal.

Failure scenario: a job finishes and the deployment reports state: "SUCCEEDED". cancelJobAndWait never matches, so the while (!client.signal.aborted) loop keeps issuing cancelJob + getJobResultRaw every 250 ms for the entire budget — ~12 round trips at the 1.5 s child budget, ~40 at the 5 s supervisor budget — and then returns { confirmed: false }. Downstream that becomes a false SQL job <id>: cancellation unconfirmed; parent cleanup or server timeout must recover it. on stderr (sql-lifecycle.ts:34-36) or a bogus record in ~/.clickzetta/sql-cleanup.jsonl, for a job that completed normally.

Smaller correct change: reuse one terminal-state predicate rather than a third inline literal — poll.ts already has TERMINAL_STATES (line 10, same narrow set) and job-profile.ts:179 has the wide one. Exporting the wide predicate from the SDK and calling it here (and ideally from poll.ts and exec.ts:178) removes the drift instead of adding to it.

) {
return { confirmed: true, state: status.state }
}
}
await delay(250, client.signal).catch(() => {})
Comment on lines +51 to +73

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The 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. createSqlSupervisor admits up to 256 connections (cleanup-scope.ts:85) and close() fans out cleanup across all of them concurrently, so a worst-case shutdown is on the order of 10⁴ requests to the gateway inside 5 seconds — from a client that is trying to exit.

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 SUCCEEDED state-matching gap I flagged on line 68: that bug is what turns a job that finished cleanly into a full-budget retry loop rather than a single confirmed round trip.

}
return { confirmed: false, reason }
} finally {
deadline.dispose()
}
}
Loading
Loading