diff --git a/knip.json b/knip.json index c57c8be24..14df12e42 100644 --- a/knip.json +++ b/knip.json @@ -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", diff --git a/packages/appkit/package.json b/packages/appkit/package.json index a72554fcd..d081b45f0 100644 --- a/packages/appkit/package.json +++ b/packages/appkit/package.json @@ -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", diff --git a/packages/appkit/src/connectors/jobs/client.ts b/packages/appkit/src/connectors/jobs/client.ts index a00474d21..ab048a202 100644 --- a/packages/appkit/src/connectors/jobs/client.ts +++ b/packages/appkit/src/connectors/jobs/client.ts @@ -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` 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 = { + pipelineTaskParameters: "parameters", +}; + +function isPlainObject(value: unknown): value is Record { + 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; @@ -55,7 +135,11 @@ export class JobsConnector { signal?: AbortSignal, ): Promise { 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) }; }); } @@ -65,7 +149,17 @@ export class JobsConnector { signal?: AbortSignal, ): Promise { 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 }; }); } @@ -75,7 +169,14 @@ export class JobsConnector { signal?: AbortSignal, ): Promise { 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; }); } @@ -85,10 +186,11 @@ export class JobsConnector { signal?: AbortSignal, ): Promise { 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; }); } @@ -98,9 +200,9 @@ export class JobsConnector { signal?: AbortSignal, ): Promise { await this._callApi("cancelRun", async () => { - return workspaceClient.jobs.cancelRun( - request, - this._createContext(signal), + await workspaceClient.jobs.cancelRun( + { runId: BigInt(request.run_id) }, + { signal }, ); }); } @@ -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; @@ -130,7 +237,14 @@ export class JobsConnector { signal?: AbortSignal, ): Promise { 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; }); } @@ -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 ( @@ -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 }); - }, - }, - }); - } } diff --git a/packages/appkit/src/connectors/jobs/tests/client.test.ts b/packages/appkit/src/connectors/jobs/tests/client.test.ts new file mode 100644 index 000000000..542175284 --- /dev/null +++ b/packages/appkit/src/connectors/jobs/tests/client.test.ts @@ -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" }); + }); +}); diff --git a/packages/appkit/src/plugins/jobs/tests/plugin.test.ts b/packages/appkit/src/plugins/jobs/tests/plugin.test.ts index 84b44ed0c..a9fc1a33d 100644 --- a/packages/appkit/src/plugins/jobs/tests/plugin.test.ts +++ b/packages/appkit/src/plugins/jobs/tests/plugin.test.ts @@ -24,29 +24,22 @@ const { mockClient, jobsApi } = await vi.hoisted(async () => { const mockClient = createMockWorkspaceClient(); - // Facade accessors are typed against the legacy SDK, so `.mockResolvedValue` + // Facade accessors are typed against the SDK, so `.mockResolvedValue` // on them would not typecheck. `getMock` is the typed handle; it mints // idempotently, so these are the very functions the plugin will call. const jobsApi = { runNow: getMock(mockClient, "jobs.runNow"), - submit: getMock(mockClient, "jobs.submit"), + submitRun: getMock(mockClient, "jobs.submitRun"), getRun: getMock(mockClient, "jobs.getRun"), getRunOutput: getMock(mockClient, "jobs.getRunOutput"), cancelRun: getMock(mockClient, "jobs.cancelRun"), - listRuns: getMock(mockClient, "jobs.listRuns"), - get: getMock(mockClient, "jobs.get"), + listRunsIter: getMock(mockClient, "jobs.listRunsIter"), + getJob: getMock(mockClient, "jobs.getJob"), }; return { mockClient, jobsApi }; }); -vi.mock("../../../workspace-client", async (importOriginal) => { - const actual = - await importOriginal(); - // Only `Context` — the client itself is injected through ServiceContext. - return { ...actual, Context: vi.fn() }; -}); - // Boots AppKit's real in-memory cache (no cache-module mock needed). useTestCache(); @@ -264,7 +257,7 @@ describe("JobsPlugin", () => { test("runNow passes configured job_id to connector", async () => { process.env.DATABRICKS_JOB_ETL = "123"; - jobsApi.runNow.mockResolvedValue({ run_id: 42 }); + jobsApi.runNow.mockResolvedValue({ runId: 42n }); const plugin = new JobsPlugin({}); const exported = plugin.exports(); @@ -273,7 +266,7 @@ describe("JobsPlugin", () => { await handle.runNow(); expect(jobsApi.runNow).toHaveBeenCalledWith( - expect.objectContaining({ job_id: 123 }), + expect.objectContaining({ jobId: 123n }), expect.anything(), ); }); @@ -281,7 +274,7 @@ describe("JobsPlugin", () => { test("runNow merges user params with configured job_id (no taskType)", async () => { process.env.DATABRICKS_JOB_ETL = "123"; - jobsApi.runNow.mockResolvedValue({ run_id: 42 }); + jobsApi.runNow.mockResolvedValue({ runId: 42n }); const plugin = new JobsPlugin({}); const exported = plugin.exports(); @@ -293,8 +286,8 @@ describe("JobsPlugin", () => { expect(jobsApi.runNow).toHaveBeenCalledWith( expect.objectContaining({ - job_id: 123, - notebook_params: { key: "value" }, + jobId: 123n, + notebookParams: { key: "value" }, }), expect.anything(), ); @@ -323,7 +316,7 @@ describe("JobsPlugin", () => { test("runNow maps validated params to SDK fields when taskType is set", async () => { process.env.DATABRICKS_JOB_ETL = "123"; - jobsApi.runNow.mockResolvedValue({ run_id: 42 }); + jobsApi.runNow.mockResolvedValue({ runId: 42n }); const plugin = new JobsPlugin({ jobs: { @@ -339,8 +332,8 @@ describe("JobsPlugin", () => { expect(jobsApi.runNow).toHaveBeenCalledWith( expect.objectContaining({ - job_id: 123, - notebook_params: { key: "value" }, + jobId: 123n, + notebookParams: { key: "value" }, }), expect.anything(), ); @@ -349,7 +342,7 @@ describe("JobsPlugin", () => { test("runNow skips validation when no schema is configured", async () => { process.env.DATABRICKS_JOB_ETL = "123"; - jobsApi.runNow.mockResolvedValue({ run_id: 42 }); + jobsApi.runNow.mockResolvedValue({ runId: 42n }); const plugin = new JobsPlugin({}); const handle = plugin.exports()("etl"); @@ -363,8 +356,8 @@ describe("JobsPlugin", () => { process.env.DATABRICKS_JOB_ETL = "123"; jobsApi.getRun.mockResolvedValue({ - run_id: 1, - state: { life_cycle_state: "TERMINATED" }, + runId: 1n, + state: { lifeCycleState: "TERMINATED" }, }); const plugin = new JobsPlugin({}); @@ -389,7 +382,7 @@ describe("JobsPlugin", () => { test("getJob wraps call in execute", async () => { process.env.DATABRICKS_JOB_ETL = "123"; - jobsApi.get.mockResolvedValue({ job_id: 123 }); + jobsApi.getJob.mockResolvedValue({ jobId: 123n }); const plugin = new JobsPlugin({}); const executeSpy = vi.spyOn(plugin as any, "execute"); @@ -413,7 +406,7 @@ describe("JobsPlugin", () => { test("listRuns clamps caller-supplied limit before calling the SDK", async () => { process.env.DATABRICKS_JOB_ETL = "123"; - jobsApi.listRuns.mockReturnValue((async function* () {})()); + jobsApi.listRunsIter.mockReturnValue((async function* () {})()); const plugin = new JobsPlugin({}); const handle = plugin.exports()("etl"); @@ -421,7 +414,7 @@ describe("JobsPlugin", () => { await handle.listRuns({ limit: 10000 }); // SDK should receive the clamped limit, not the caller-supplied 10000. - expect(jobsApi.listRuns).toHaveBeenCalledWith( + expect(jobsApi.listRunsIter).toHaveBeenCalledWith( expect.objectContaining({ limit: 100 }), expect.anything(), ); @@ -431,7 +424,7 @@ describe("JobsPlugin", () => { process.env.DATABRICKS_JOB_ETL = "123"; // Pre-flight getRun verifies the run belongs to the configured jobId. - jobsApi.getRun.mockResolvedValue({ run_id: 1, job_id: 123 }); + jobsApi.getRun.mockResolvedValue({ runId: 1n, jobId: 123n }); jobsApi.cancelRun.mockResolvedValue(undefined); const plugin = new JobsPlugin({}); @@ -452,15 +445,15 @@ describe("JobsPlugin", () => { test("runAndWait yields status updates and terminates on TERMINATED", async () => { process.env.DATABRICKS_JOB_ETL = "123"; - jobsApi.runNow.mockResolvedValue({ run_id: 42 }); + jobsApi.runNow.mockResolvedValue({ runId: 42n }); jobsApi.getRun .mockResolvedValueOnce({ - run_id: 42, - state: { life_cycle_state: "RUNNING" }, + runId: 42n, + state: { lifeCycleState: "RUNNING" }, }) .mockResolvedValueOnce({ - run_id: 42, - state: { life_cycle_state: "TERMINATED" }, + runId: 42n, + state: { lifeCycleState: "TERMINATED" }, }); const plugin = new JobsPlugin({ pollIntervalMs: 10 }); @@ -544,7 +537,7 @@ describe("JobsPlugin", () => { test("listRuns returns error result on execute failure", async () => { process.env.DATABRICKS_JOB_ETL = "123"; - jobsApi.listRuns.mockImplementation(() => { + jobsApi.listRunsIter.mockImplementation(() => { throw new Error("Auth failure"); }); @@ -587,7 +580,7 @@ describe("JobsPlugin", () => { test("successful operations return ok result with data", async () => { process.env.DATABRICKS_JOB_ETL = "123"; - jobsApi.runNow.mockResolvedValue({ run_id: 42 }); + jobsApi.runNow.mockResolvedValue({ runId: 42n }); const plugin = new JobsPlugin({}); const handle = plugin.exports()("etl"); @@ -604,7 +597,7 @@ describe("JobsPlugin", () => { test("getRun returns 404 when run.job_id does not match configured jobId", async () => { process.env.DATABRICKS_JOB_ETL = "123"; - jobsApi.getRun.mockResolvedValue({ run_id: 99, job_id: 456 }); + jobsApi.getRun.mockResolvedValue({ runId: 99n, jobId: 456n }); const plugin = new JobsPlugin({}); const handle = plugin.exports()("etl"); @@ -617,7 +610,7 @@ describe("JobsPlugin", () => { test("getRunOutput returns 404 when run belongs to another job", async () => { process.env.DATABRICKS_JOB_ETL = "123"; - jobsApi.getRun.mockResolvedValue({ run_id: 99, job_id: 456 }); + jobsApi.getRun.mockResolvedValue({ runId: 99n, jobId: 456n }); jobsApi.getRunOutput.mockResolvedValue({ logs: "nope" }); const plugin = new JobsPlugin({}); @@ -633,7 +626,7 @@ describe("JobsPlugin", () => { test("cancelRun returns 404 when run belongs to another job", async () => { process.env.DATABRICKS_JOB_ETL = "123"; - jobsApi.getRun.mockResolvedValue({ run_id: 99, job_id: 456 }); + jobsApi.getRun.mockResolvedValue({ runId: 99n, jobId: 456n }); jobsApi.cancelRun.mockResolvedValue(undefined); const plugin = new JobsPlugin({}); @@ -649,9 +642,9 @@ describe("JobsPlugin", () => { process.env.DATABRICKS_JOB_ETL = "123"; jobsApi.getRun.mockResolvedValue({ - run_id: 42, - job_id: 123, - state: { life_cycle_state: "TERMINATED" }, + runId: 42n, + jobId: 123n, + state: { lifeCycleState: "TERMINATED" }, }); const plugin = new JobsPlugin({}); @@ -662,17 +655,11 @@ describe("JobsPlugin", () => { if (result.ok) expect(result.data.run_id).toBe(42); }); - test("connector's cancellation token reflects signal state live", async () => { - process.env.DATABRICKS_JOB_ETL = "123"; - - const { Context } = await import("../../../workspace-client"); - const mockContext = Context as unknown as ReturnType; - mockContext.mockClear(); - + test("connector forwards the abort signal as CallOptions", async () => { const { JobsConnector } = await import("../../../connectors/jobs"); const connector = new JobsConnector({}); - jobsApi.get.mockResolvedValue({ job_id: 123 }); + jobsApi.getJob.mockResolvedValue({ jobId: 123n }); const controller = new AbortController(); await connector.getJob( @@ -681,12 +668,10 @@ describe("JobsPlugin", () => { controller.signal, ); - const ctorArg = mockContext.mock.calls.at(-1)?.[0] as { - cancellationToken: { isCancellationRequested: boolean }; - }; - expect(ctorArg.cancellationToken.isCancellationRequested).toBe(false); - controller.abort(); - expect(ctorArg.cancellationToken.isCancellationRequested).toBe(true); + expect(jobsApi.getJob).toHaveBeenCalledWith( + { jobId: 123n }, + { signal: controller.signal }, + ); }); }); @@ -694,10 +679,10 @@ describe("JobsPlugin", () => { test("runAndWait stops polling when signal is aborted", async () => { process.env.DATABRICKS_JOB_ETL = "123"; - jobsApi.runNow.mockResolvedValue({ run_id: 42 }); + jobsApi.runNow.mockResolvedValue({ runId: 42n }); jobsApi.getRun.mockResolvedValue({ - run_id: 42, - state: { life_cycle_state: "RUNNING" }, + runId: 42n, + state: { lifeCycleState: "RUNNING" }, }); const plugin = new JobsPlugin({ pollIntervalMs: 10 }); @@ -805,14 +790,14 @@ describe("JobsPlugin", () => { process.env.DATABRICKS_JOB_ETL = "100"; process.env.DATABRICKS_JOB_ML = "200"; - jobsApi.runNow.mockResolvedValue({ run_id: 1 }); + jobsApi.runNow.mockResolvedValue({ runId: 1n }); const plugin = new JobsPlugin({}); const exported = plugin.exports(); await exported("etl").runNow(); expect(jobsApi.runNow).toHaveBeenCalledWith( - expect.objectContaining({ job_id: 100 }), + expect.objectContaining({ jobId: 100n }), expect.anything(), ); @@ -820,7 +805,7 @@ describe("JobsPlugin", () => { await exported("ml").runNow(); expect(jobsApi.runNow).toHaveBeenCalledWith( - expect.objectContaining({ job_id: 200 }), + expect.objectContaining({ jobId: 200n }), expect.anything(), ); }); @@ -1060,7 +1045,7 @@ describe("injectRoutes", () => { test("returns runId on successful non-streaming run", async () => { process.env.DATABRICKS_JOB_ETL = "123"; - jobsApi.runNow.mockResolvedValue({ run_id: 42 }); + jobsApi.runNow.mockResolvedValue({ runId: 42n }); const plugin = new JobsPlugin({}); const routeSpy = vi.spyOn(plugin as any, "route"); @@ -1179,10 +1164,10 @@ describe("injectRoutes", () => { process.env.DATABRICKS_JOB_ETL = "123"; const mockRuns = [ - { run_id: 1, state: { life_cycle_state: "TERMINATED" } }, - { run_id: 2, state: { life_cycle_state: "RUNNING" } }, + { runId: 1n, state: { lifeCycleState: "TERMINATED" } }, + { runId: 2n, state: { lifeCycleState: "RUNNING" } }, ]; - jobsApi.listRuns.mockReturnValue( + jobsApi.listRunsIter.mockReturnValue( (async function* () { for (const run of mockRuns) yield run; })(), @@ -1212,15 +1197,19 @@ describe("injectRoutes", () => { await handler(mockReq, mockRes); + // The modular SDK's camelCase/bigint models go out in the legacy wire shape. expect(mockRes.json).toHaveBeenCalledWith({ - runs: mockRuns, + runs: [ + { run_id: 1, state: { life_cycle_state: "TERMINATED" } }, + { run_id: 2, state: { life_cycle_state: "RUNNING" } }, + ], }); }); test("passes limit query param to listRuns", async () => { process.env.DATABRICKS_JOB_ETL = "123"; - jobsApi.listRuns.mockReturnValue((async function* () {})()); + jobsApi.listRunsIter.mockReturnValue((async function* () {})()); const plugin = new JobsPlugin({}); const routeSpy = vi.spyOn(plugin as any, "route"); @@ -1247,7 +1236,7 @@ describe("injectRoutes", () => { await handler(mockReq, mockRes); // Verify the connector was called with limit 5 - expect(jobsApi.listRuns).toHaveBeenCalledWith( + expect(jobsApi.listRunsIter).toHaveBeenCalledWith( expect.objectContaining({ limit: 5 }), expect.anything(), ); @@ -1259,9 +1248,9 @@ describe("injectRoutes", () => { process.env.DATABRICKS_JOB_ETL = "123"; const mockRun = { - run_id: 42, - job_id: 123, - state: { life_cycle_state: "TERMINATED" }, + runId: 42n, + jobId: 123n, + state: { lifeCycleState: "TERMINATED" }, }; jobsApi.getRun.mockResolvedValue(mockRun); @@ -1289,7 +1278,11 @@ describe("injectRoutes", () => { await handler(mockReq, mockRes); - expect(mockRes.json).toHaveBeenCalledWith(mockRun); + expect(mockRes.json).toHaveBeenCalledWith({ + run_id: 42, + job_id: 123, + state: { life_cycle_state: "TERMINATED" }, + }); }); test("returns 400 for invalid runId", async () => { @@ -1330,7 +1323,7 @@ describe("injectRoutes", () => { process.env.DATABRICKS_JOB_ETL = "123"; // Run exists upstream but is owned by job 456, not the configured 123. - jobsApi.getRun.mockResolvedValue({ run_id: 99, job_id: 456 }); + jobsApi.getRun.mockResolvedValue({ runId: 99n, jobId: 456n }); const plugin = new JobsPlugin({}); const routeSpy = vi.spyOn(plugin as any, "route"); @@ -1369,10 +1362,10 @@ describe("injectRoutes", () => { process.env.DATABRICKS_JOB_ETL = "123"; const mockRun = { - run_id: 42, - state: { life_cycle_state: "TERMINATED" }, + runId: 42n, + state: { lifeCycleState: "TERMINATED" }, }; - jobsApi.listRuns.mockReturnValue( + jobsApi.listRunsIter.mockReturnValue( (async function* () { yield mockRun; })(), @@ -1404,14 +1397,14 @@ describe("injectRoutes", () => { expect(mockRes.json).toHaveBeenCalledWith({ status: "TERMINATED", - run: mockRun, + run: { run_id: 42, state: { life_cycle_state: "TERMINATED" } }, }); }); test("returns null status when no runs exist", async () => { process.env.DATABRICKS_JOB_ETL = "123"; - jobsApi.listRuns.mockReturnValue((async function* () {})()); + jobsApi.listRunsIter.mockReturnValue((async function* () {})()); const plugin = new JobsPlugin({}); const routeSpy = vi.spyOn(plugin as any, "route"); @@ -1448,7 +1441,7 @@ describe("injectRoutes", () => { test("cancels run and returns 204", async () => { process.env.DATABRICKS_JOB_ETL = "123"; - jobsApi.getRun.mockResolvedValue({ run_id: 42, job_id: 123 }); + jobsApi.getRun.mockResolvedValue({ runId: 42n, jobId: 123n }); jobsApi.cancelRun.mockResolvedValue(undefined); const plugin = new JobsPlugin({}); @@ -1519,7 +1512,7 @@ describe("injectRoutes", () => { process.env.DATABRICKS_JOB_ETL = "123"; // Pre-flight getRun reports a run owned by a different job. - jobsApi.getRun.mockResolvedValue({ run_id: 99, job_id: 456 }); + jobsApi.getRun.mockResolvedValue({ runId: 99n, jobId: 456n }); jobsApi.cancelRun.mockResolvedValue(undefined); const plugin = new JobsPlugin({}); @@ -1701,7 +1694,7 @@ describe("injectRoutes", () => { test("allows exactly MAX_UNVALIDATED_PARAM_KEYS (50) keys without schema", async () => { process.env.DATABRICKS_JOB_ETL = "123"; - jobsApi.runNow.mockResolvedValue({ run_id: 42 }); + jobsApi.runNow.mockResolvedValue({ runId: 42n }); const plugin = new JobsPlugin({ jobs: { etl: { taskType: "notebook" } }, @@ -1745,7 +1738,7 @@ describe("injectRoutes", () => { test("allows undefined params", async () => { process.env.DATABRICKS_JOB_ETL = "123"; - jobsApi.runNow.mockResolvedValue({ run_id: 42 }); + jobsApi.runNow.mockResolvedValue({ runId: 42n }); const plugin = new JobsPlugin({}); const routeSpy = vi.spyOn(plugin as any, "route"); @@ -1826,7 +1819,7 @@ describe("injectRoutes", () => { test("GET /:jobKey/runs returns upstream status on failure", async () => { process.env.DATABRICKS_JOB_ETL = "123"; - jobsApi.listRuns.mockImplementation(() => { + jobsApi.listRunsIter.mockImplementation(() => { throw createApiError({ statusCode: 401, message: "Unauthorized", @@ -1865,7 +1858,7 @@ describe("injectRoutes", () => { process.env.DATABRICKS_JOB_ETL = "123"; // Pre-flight succeeds so we reach the actual cancel call. - jobsApi.getRun.mockResolvedValue({ run_id: 42, job_id: 123 }); + jobsApi.getRun.mockResolvedValue({ runId: 42n, jobId: 123n }); const error = new Error("Forbidden"); (error as any).statusCode = 403; jobsApi.cancelRun.mockRejectedValue(error); diff --git a/packages/shared/package.json b/packages/shared/package.json index 6002e4693..01938030b 100644 --- a/packages/shared/package.json +++ b/packages/shared/package.json @@ -52,6 +52,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", diff --git a/packages/shared/src/workspace-client/client.ts b/packages/shared/src/workspace-client/client.ts index 1552392e5..92ca56b11 100644 --- a/packages/shared/src/workspace-client/client.ts +++ b/packages/shared/src/workspace-client/client.ts @@ -20,6 +20,8 @@ import { type ScimClient, buildGenieClient, type GenieClient, + buildJobsClient, + type JobsClient, type StatementExecutionClient, type WarehousesClient, type WorkspaceAuth, @@ -35,6 +37,7 @@ export class AppKitWorkspaceClient implements WorkspaceClient { #auth?: WorkspaceAuth; #currentUser?: ScimClient; #genie?: GenieClient; + #jobs?: JobsClient; constructor(opts: WorkspaceClientOptions) { this.#opts = opts; @@ -60,8 +63,12 @@ export class AppKitWorkspaceClient implements WorkspaceClient { return this.#genie; } - get jobs() { - return this.#getLegacy().jobs; + // Migrated to the modular SDK — built lazily, independent of the legacy client. + get jobs(): JobsClient { + if (!this.#jobs) { + this.#jobs = buildJobsClient(this.#opts); + } + return this.#jobs; } // Migrated to the modular SDK — built lazily, independent of the legacy client. diff --git a/packages/shared/src/workspace-client/modular.ts b/packages/shared/src/workspace-client/modular.ts index c4baade6b..43ca41374 100644 --- a/packages/shared/src/workspace-client/modular.ts +++ b/packages/shared/src/workspace-client/modular.ts @@ -8,9 +8,9 @@ * * Migrated services are built here as per-service clients; the facade delegates * their accessors to these instead of the legacy monolithic client. Currently - * `warehouses` and `statementExecution` are migrated, plus the auth + raw-request - * seam ({@link buildWorkspaceAuth}); every other service still routes through - * `legacy.ts`. + * `warehouses`, `statementExecution`, `currentUser` (SCIM), `genie` and `jobs` + * are migrated, plus the auth + raw-request seam ({@link buildWorkspaceAuth}); + * every other service still routes through `legacy.ts`. * * NOTE: statementExecution relies on a pinned pnpm patch * (`patches/@databricks__sdk-statementexecution@0.46.0.patch`) that restores the @@ -42,6 +42,7 @@ import { } from "@databricks/sdk-core/http"; import { resolve } from "@databricks/sdk-core/profiles"; import { GenieClient } from "@databricks/sdk-genie/v1"; +import { JobsClient } from "@databricks/sdk-jobs/v2"; import type { ClientOptions } from "@databricks/sdk-options/client"; import { ScimClient } from "@databricks/sdk-scim/v1"; import { StatementExecutionClient } from "@databricks/sdk-statementexecution/v1"; @@ -333,8 +334,14 @@ export function buildGenieClient(opts: WorkspaceClientOptions): GenieClient { return new GenieClient(mapToClientOptions(opts)); } +/** Build a modular Jobs (API 2.2) client from wrapper options. */ +export function buildJobsClient(opts: WorkspaceClientOptions): JobsClient { + return new JobsClient(mapToClientOptions(opts)); +} + // ── Client type re-exports (for the facade accessor types) ─────────────── export type { GenieClient } from "@databricks/sdk-genie/v1"; +export type { JobsClient } from "@databricks/sdk-jobs/v2"; export type { ScimClient } from "@databricks/sdk-scim/v1"; export type { StatementExecutionClient } from "@databricks/sdk-statementexecution/v1"; export type { WarehousesClient } from "@databricks/sdk-warehouses/v1"; @@ -363,6 +370,13 @@ export type { GenieGetMessageQueryResultResponse, GenieMessage, } from "@databricks/sdk-genie/v1"; +export type { + GetJobRequest, + GetRunRequest, + ListRunsRequest, + RunNowRequest, + SubmitRunRequest, +} from "@databricks/sdk-jobs/v2"; export type { EndpointHealth, EndpointInfo, diff --git a/packages/shared/src/workspace-client/types.ts b/packages/shared/src/workspace-client/types.ts index 75b8b80c0..352b1780a 100644 --- a/packages/shared/src/workspace-client/types.ts +++ b/packages/shared/src/workspace-client/types.ts @@ -20,6 +20,7 @@ import type { WorkspaceAuth, ScimClient, GenieClient, + JobsClient, } from "./modular"; // Legacy SDK type namespaces for un-migrated services, re-exported so AppKit @@ -30,8 +31,11 @@ import type { // as `sql.EndpointInfo[]`. It is not on the modular `listWarehouses` because // that request has no `skip_cannot_use` filter, so it could pick a warehouse // the caller can't use. Statement + warehouse service types come from `./modular`. +// `jobs` stays as the wire-shape (snake_case, `number` IDs) type of the jobs +// plugin's public API; appkit's jobs connector translates the modular client's +// camelCase/`bigint` models back to it. export type { files, jobs, serving, sql } from "@databricks/sdk-experimental"; -// Modular SDK client + model types (warehouses, statementExecution, genie). +// Modular SDK client + model types (warehouses, statementExecution, genie, jobs). export type * from "./modular"; /** @@ -52,8 +56,8 @@ export interface WorkspaceClient extends WorkspaceAuth { /** Genie (modular SDK). */ readonly genie: GenieClient; - /** Jobs. */ - readonly jobs: LegacyWorkspaceClient["jobs"]; + /** Jobs (modular SDK, Jobs API 2.2). */ + readonly jobs: JobsClient; /** Statement Execution (modular SDK). */ readonly statementExecution: StatementExecutionClient; diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index fd69b375f..c1b0f795e 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -292,6 +292,9 @@ importers: '@databricks/sdk-genie': specifier: 0.54.0 version: 0.54.0(patch_hash=664f3e13b155eb1b5499486a4f6fc1ab9a6a502e2c6b200d06749bd25dd9ae54) + '@databricks/sdk-jobs': + specifier: 0.57.0 + version: 0.57.0 '@databricks/sdk-options': specifier: 0.51.0 version: 0.51.0 @@ -628,6 +631,9 @@ importers: '@databricks/sdk-genie': specifier: 0.54.0 version: 0.54.0(patch_hash=664f3e13b155eb1b5499486a4f6fc1ab9a6a502e2c6b200d06749bd25dd9ae54) + '@databricks/sdk-jobs': + specifier: 0.57.0 + version: 0.57.0 '@databricks/sdk-options': specifier: 0.51.0 version: 0.51.0 @@ -2025,6 +2031,10 @@ packages: resolution: {integrity: sha512-xT+dyuXocgwPJ5VA1Kr5/dSdKzzR9hRAkEOybeO2i+v891kuuFzVd/M9OM31PtO2/YZp95vdFH5Sb+jVGtm4nw==} engines: {node: '>=22.0.0'} + '@databricks/sdk-jobs@0.57.0': + resolution: {integrity: sha512-LUBLNRhJq0OtUUU6+2T/zRzZMEd4UL0lVm7qTxLIfuDPDNDPy/EHGs3rQ1cXt1Fv0t3Rik7QT/SueZixSsH/Bg==} + engines: {node: '>=22.0.0'} + '@databricks/sdk-options@0.51.0': resolution: {integrity: sha512-p5uBh64Y1onnvwlEwfRsmKmGRI23Y4Ly8fZiCgNUwVDe4Z9d0ctI0mo1v4gnDqlmBaBml3TZfqn5BJBROCXsEw==} engines: {node: '>=22.0.0'} @@ -14346,6 +14356,15 @@ snapshots: json-bigint: 1.0.0 zod: 4.3.6 + '@databricks/sdk-jobs@0.57.0': + dependencies: + '@databricks/sdk-auth': 0.51.0 + '@databricks/sdk-core': 0.51.0 + '@databricks/sdk-options': 0.51.0 + '@js-temporal/polyfill': 0.5.1 + json-bigint: 1.0.0 + zod: 4.3.6 + '@databricks/sdk-options@0.51.0': dependencies: '@databricks/sdk-auth': 0.51.0