Skip to content
Draft
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
1 change: 1 addition & 0 deletions knip.json
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@
"@databricks/sdk-core",
"@databricks/sdk-experimental",
"@databricks/sdk-genie",
"@databricks/sdk-jobs",
"@databricks/sdk-options",
"@databricks/sdk-scim",
"@databricks/sdk-statementexecution",
Expand Down
1 change: 1 addition & 0 deletions packages/appkit/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,7 @@
"@databricks/sdk-core": "0.51.0",
"@databricks/sdk-experimental": "0.17.0",
"@databricks/sdk-genie": "0.54.0",
"@databricks/sdk-jobs": "0.57.0",
"@databricks/sdk-options": "0.51.0",
"@databricks/sdk-scim": "0.51.0",
"@databricks/sdk-statementexecution": "0.52.0",
Expand Down
176 changes: 143 additions & 33 deletions packages/appkit/src/connectors/jobs/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,15 +9,95 @@ import {
SpanStatusCode,
TelemetryManager,
} from "../../telemetry";
import {
Context,
type jobs,
type WorkspaceClient,
import type {
GetJobRequest,
GetRunRequest,
jobs,
ListRunsRequest,
RunNowRequest,
SubmitRunRequest,
WorkspaceClient,
} from "../../workspace-client";
import type { JobsConnectorConfig } from "./types";

const logger = createLogger("connectors:jobs");

/**
* `Record<string, string>` fields of the Jobs model. Their keys are user data
* (notebook params, tags, Spark conf), so they are copied verbatim instead of
* re-cased.
*/
const MAP_FIELDS = new Set([
"artifactsHeaders",
"baseParameters",
"customTags",
"filters",
"jobParameters",
"namedParameters",
"notebookBaseParameters",
"notebookParams",
"parameters",
"pipelineTaskParameters",
"pythonNamedParams",
"sparkConf",
"sparkEnvVars",
"sqlParams",
"tags",
"variables",
"violations",
]);

/** The one model field whose camelCase name doesn't round-trip to its wire key. */
const WIRE_KEY_OVERRIDES: Record<string, string> = {
pipelineTaskParameters: "parameters",
};

function isPlainObject(value: unknown): value is Record<string, unknown> {
return value !== null && typeof value === "object" && !Array.isArray(value);
}

/**
* Modular model → the legacy wire shape the plugin's public API (HTTP JSON, SSE,
* cache) has always exposed: snake_case keys and `number` int64s. `bigint` would
* make `JSON.stringify` throw. `Number()` loses precision past 2^53, exactly as
* the legacy SDK's plain `JSON.parse` did.
*/
function toWire(value: unknown, verbatimKeys = false): unknown {
if (typeof value === "bigint") return Number(value);
if (Array.isArray(value)) return value.map((v) => toWire(v));
if (!isPlainObject(value)) return value;
return Object.fromEntries(
Object.entries(value).map(([key, v]) =>
verbatimKeys
? [key, v]
: [
WIRE_KEY_OVERRIDES[key] ??
key.replace(/[A-Z]/g, (c) => `_${c.toLowerCase()}`),
toWire(v, MAP_FIELDS.has(key)),
],
),
);
}

/**
* Legacy snake_case request → modular camelCase request. int64 fields (`job_id`,
* `run_id`) still need an explicit `BigInt()` at the call site: the SDK's
* marshal schemas reject a `number` there.
*/
function fromWire(value: unknown, verbatimKeys = false): unknown {
if (Array.isArray(value)) return value.map((v) => fromWire(v));
if (!isPlainObject(value)) return value;
return Object.fromEntries(
Object.entries(value).map(([key, v]) => {
if (verbatimKeys) return [key, v];
const camel = key.replace(/_([a-z0-9])/g, (_, c: string) =>
c.toUpperCase(),
);
return [camel, fromWire(v, MAP_FIELDS.has(camel))];
}),
);
}

export class JobsConnector {
private readonly name = "jobs";
private readonly config: JobsConnectorConfig;
Expand Down Expand Up @@ -55,7 +135,11 @@ export class JobsConnector {
signal?: AbortSignal,
): Promise<jobs.SubmitRunResponse> {
return this._callApi("submit", async () => {
return workspaceClient.jobs.submit(request, this._createContext(signal));
const waiter = await workspaceClient.jobs.submitRun(
fromWire(request) as SubmitRunRequest,
{ signal },
);
return { run_id: Number(waiter.runId) };
});
}

Expand All @@ -65,7 +149,17 @@ export class JobsConnector {
signal?: AbortSignal,
): Promise<jobs.RunNowResponse> {
return this._callApi("runNow", async () => {
return workspaceClient.jobs.runNow(request, this._createContext(signal));
const waiter = await workspaceClient.jobs.runNow(
{
...(fromWire(request) as RunNowRequest),
jobId: BigInt(request.job_id),
},
{ signal },
);
// The waiter only exposes runId; the API documents number_in_job as
// "set to the same value as run_id".
const runId = Number(waiter.runId);
return { run_id: runId, number_in_job: runId };
});
}

Expand All @@ -75,7 +169,14 @@ export class JobsConnector {
signal?: AbortSignal,
): Promise<jobs.Run> {
return this._callApi("getRun", async () => {
return workspaceClient.jobs.getRun(request, this._createContext(signal));
const run = await workspaceClient.jobs.getRun(
{
...(fromWire(request) as GetRunRequest),
runId: BigInt(request.run_id),
},
{ signal },
);
return toWire(run) as jobs.Run;
});
}

Expand All @@ -85,10 +186,11 @@ export class JobsConnector {
signal?: AbortSignal,
): Promise<jobs.RunOutput> {
return this._callApi("getRunOutput", async () => {
return workspaceClient.jobs.getRunOutput(
request,
this._createContext(signal),
const output = await workspaceClient.jobs.getRunOutput(
{ runId: BigInt(request.run_id) },
{ signal },
);
return toWire(output) as jobs.RunOutput;
});
}

Expand All @@ -98,9 +200,9 @@ export class JobsConnector {
signal?: AbortSignal,
): Promise<void> {
await this._callApi("cancelRun", async () => {
return workspaceClient.jobs.cancelRun(
request,
this._createContext(signal),
await workspaceClient.jobs.cancelRun(
{ runId: BigInt(request.run_id) },
{ signal },
);
});
}
Expand All @@ -113,11 +215,16 @@ export class JobsConnector {
return this._callApi("listRuns", async () => {
const runs: jobs.BaseRun[] = [];
const limit = Math.max(1, Math.min(request.limit ?? 100, 100));
for await (const run of workspaceClient.jobs.listRuns(
{ ...request, limit },
this._createContext(signal),
for await (const run of workspaceClient.jobs.listRunsIter(
{
...(fromWire(request) as ListRunsRequest),
jobId:
request.job_id === undefined ? undefined : BigInt(request.job_id),
limit,
},
{ signal },
)) {
runs.push(run);
runs.push(toWire(run) as jobs.BaseRun);
if (runs.length >= limit) break;
}
return runs;
Expand All @@ -130,7 +237,14 @@ export class JobsConnector {
signal?: AbortSignal,
): Promise<jobs.Job> {
return this._callApi("getJob", async () => {
return workspaceClient.jobs.get(request, this._createContext(signal));
const job = await workspaceClient.jobs.getJob(
{
...(fromWire(request) as GetJobRequest),
jobId: BigInt(request.job_id),
},
{ signal },
);
return toWire(job) as jobs.Job;
});
}

Expand Down Expand Up @@ -164,6 +278,17 @@ export class JobsConnector {
if (error instanceof AppKitError) {
throw error;
}
// The modular SDK's ApiError exposes the HTTP status as
// `httpStatusCode` (-1 when not an HTTP error); Plugin.execute()
// maps on `statusCode`.
if (
error instanceof Error &&
"httpStatusCode" in error &&
typeof error.httpStatusCode === "number" &&
error.httpStatusCode > 0
) {
throw Object.assign(error, { statusCode: error.httpStatusCode });
}
// Preserve SDK ApiError (and any error with a numeric statusCode)
// so Plugin.execute() can map it to the correct HTTP status.
if (
Expand Down Expand Up @@ -198,19 +323,4 @@ export class JobsConnector {
{ name: this.name, includePrefix: true },
);
}

private _createContext(signal?: AbortSignal) {
return new Context({
cancellationToken: {
// Getter — evaluated on every read so SDK code paths that poll
// (rather than subscribe) observe cancellation live.
get isCancellationRequested() {
return signal?.aborted ?? false;
},
onCancellationRequested: (cb: () => void) => {
signal?.addEventListener("abort", cb, { once: true });
},
},
});
}
}
74 changes: 74 additions & 0 deletions packages/appkit/src/connectors/jobs/tests/client.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,74 @@
import { describe, expect, test, vi } from "vitest";

import { JobsConnector } from "../client";

// The modular Jobs SDK returns camelCase models with `bigint` int64s, which
// `JSON.stringify` rejects. The connector must hand the plugin the legacy wire
// shape (snake_case, `number`) its public HTTP/SSE/cache surface has always had.
describe("JobsConnector wire-shape translation", () => {
test("getRun returns a JSON-serializable snake_case run with numeric IDs", async () => {
const getRun = vi.fn().mockResolvedValue({
runId: 9007199254740991n,
jobId: 123n,
startTime: 1700000000000n,
state: { lifeCycleState: "TERMINATED" },
overridingParameters: { notebookParams: { myParam: "x" } },
tasks: [
{
taskKey: "t1",
pipelineTask: { pipelineTaskParameters: { fullRefresh: "true" } },
},
],
});
const client = { jobs: { getRun } } as never;

const run = await new JobsConnector({}).getRun(client, { run_id: 42 });

expect(getRun).toHaveBeenCalledWith({ runId: 42n }, { signal: undefined });
expect(JSON.parse(JSON.stringify(run))).toEqual({
run_id: 9007199254740991,
job_id: 123,
start_time: 1700000000000,
state: { life_cycle_state: "TERMINATED" },
// Map keys are user data: copied verbatim, not re-cased.
overriding_parameters: { notebook_params: { myParam: "x" } },
tasks: [
{
task_key: "t1",
pipeline_task: { parameters: { fullRefresh: "true" } },
},
],
});
});

test("runNow camelCases the request without touching param keys", async () => {
const runNow = vi.fn().mockResolvedValue({ runId: 7n });
const client = { jobs: { runNow } } as never;

const result = await new JobsConnector({}).runNow(client, {
job_id: 123,
notebook_params: { my_param: "v" },
});

expect(runNow).toHaveBeenCalledWith(
{ jobId: 123n, notebookParams: { my_param: "v" } },
{ signal: undefined },
);
expect(result).toEqual({ run_id: 7, number_in_job: 7 });
});

test("surfaces the modular ApiError's httpStatusCode as statusCode", async () => {
// Shaped like sdk-core's ApiError: the status is a getter, not `statusCode`.
class ModularApiError extends Error {
get httpStatusCode() {
return 404;
}
}
const getRun = vi.fn().mockRejectedValue(new ModularApiError("not found"));
const client = { jobs: { getRun } } as never;

await expect(
new JobsConnector({}).getRun(client, { run_id: 1 }),
).rejects.toMatchObject({ statusCode: 404, message: "not found" });
});
});
Loading
Loading