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
1 change: 1 addition & 0 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@
"@typescript-eslint/parser": "^5.33.0",
"eslint": "^8.21.0",
"eslint-config-prettier": "^8.5.0",
"ts-mockito": "^2.6.1",
"typescript": "^4.7.4"
}
}
1 change: 1 addition & 0 deletions packages/databricks-sdk-js/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,7 @@
"tmp-promise": "^3.0.3",
"ts-loader": "^9.3.1",
"ts-mocha": "^10.0.0",
"ts-mockito": "^2.6.1",
"ts-node": "^10.9.1",
"typescript": "^4.7.4",
"uuid": "^8.3.2"
Expand Down
33 changes: 20 additions & 13 deletions packages/databricks-sdk-js/src/retries/retries.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,19 +9,23 @@ export interface RetriableResult<T> {
error?: unknown;
}

const maxJitter = new Time(750, TimeUnits.milliseconds);
const minJitter = new Time(50, TimeUnits.milliseconds);
const maxWaitTime = new Time(10, TimeUnits.seconds);
const defaultTimeout = new Time(10, TimeUnits.minutes);
export class RetryConfigs {
static maxJitter = new Time(750, TimeUnits.milliseconds);
static minJitter = new Time(50, TimeUnits.milliseconds);
static maxWaitTime = new Time(10, TimeUnits.seconds);
static defaultTimeout = new Time(10, TimeUnits.minutes);

function waitTime(attempt: number) {
const jitter = maxJitter
.sub(minJitter)
.multiply(Math.random())
.add(minJitter);
const timeout = new Time(attempt, TimeUnits.seconds).add(jitter);
static waitTime(attempt: number) {
const jitter = RetryConfigs.maxJitter
.sub(RetryConfigs.minJitter)
.multiply(Math.random())
.add(RetryConfigs.minJitter);
const timeout = new Time(attempt, TimeUnits.seconds).add(jitter);

return timeout.gt(maxWaitTime) ? maxWaitTime : timeout;
return timeout.gt(RetryConfigs.maxWaitTime)
? RetryConfigs.maxWaitTime
: timeout;
}
}

interface RetryArgs<T> {
Expand All @@ -30,7 +34,7 @@ interface RetryArgs<T> {
}

export default async function retry<T>({
timeout = defaultTimeout,
timeout = RetryConfigs.defaultTimeout,
fn,
}: RetryArgs<T>): Promise<T> {
let attempt = 1;
Expand Down Expand Up @@ -60,7 +64,10 @@ export default async function retry<T>({
}

await new Promise((resolve) =>
setTimeout(resolve, waitTime(attempt).toMillSeconds().value)
setTimeout(
resolve,
RetryConfigs.waitTime(attempt).toMillSeconds().value
)
);

attempt += 1;
Expand Down
47 changes: 5 additions & 42 deletions packages/databricks-sdk-js/src/services/Cluster.integ.ts
Original file line number Diff line number Diff line change
@@ -1,10 +1,13 @@
/* eslint-disable @typescript-eslint/naming-convention */

import {Cluster, ClusterService} from "..";
import {Cluster} from "..";
import assert from "assert";

import chai from "chai";
import chaiAsPromised from "chai-as-promised";
import {IntegrationTestSetup} from "../test/IntegrationTestSetup";

chai.use(chaiAsPromised);
Comment thread
kartikgupta-db marked this conversation as resolved.

describe(__filename, function () {
let integSetup: IntegrationTestSetup;

Expand Down Expand Up @@ -36,44 +39,4 @@ describe(__filename, function () {
assert(clusterA.id);
assert.equal(clusterA.id, clusterB?.id);
});

// TODO: run tests changing the state of cluster in a seperate job, on a single node
it.skip("calling start on a non terminated state should not throw an error", async () => {
await integSetup.cluster.stop();
await new ClusterService(integSetup.client).start({
cluster_id: integSetup.cluster.id,
});

await integSetup.cluster.refresh();
assert(integSetup.cluster.state !== "RUNNING");

await integSetup.cluster.start();
});

it.skip("should terminate cluster", async () => {
integSetup.cluster.start();
assert.equal(integSetup.cluster.state, "RUNNING");

await integSetup.cluster.stop();
assert.equal(integSetup.cluster.state, "TERMINATED");

integSetup.cluster.start();
});

it.skip("should terminate non running clusters", async () => {
await integSetup.cluster.stop();
assert.equal(integSetup.cluster.state, "TERMINATED");

await new ClusterService(integSetup.client).start({
cluster_id: integSetup.cluster.id,
});
await integSetup.cluster.refresh();
assert.notEqual(integSetup.cluster.state, "RUNNING");

await integSetup.cluster.stop();
assert.equal(integSetup.cluster.state, "TERMINATED");

await integSetup.cluster.start();
assert.equal(integSetup.cluster.state, "RUNNING");
});
});
242 changes: 242 additions & 0 deletions packages/databricks-sdk-js/src/services/Cluster.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,242 @@
/* eslint-disable @typescript-eslint/naming-convention */

import {ApiClient, Cluster} from "..";
import chai, {assert} from "chai";
import {mock, when, instance, deepEqual, verify, anything} from "ts-mockito";
import chaiAsPromised from "chai-as-promised";
import Time, {TimeUnits} from "../retries/Time";
import getMockTestCluster from "../test/fixtures/ClusterFixtures";
import {GetClusterResponse} from "../apis/cluster";
import TokenFixture from "../test/fixtures/TokenFixture";
import {RetryConfigs} from "../retries/retries";

chai.use(chaiAsPromised);

describe(__filename, function () {
this.timeout(new Time(10, TimeUnits.minutes).toMillSeconds().value);

let mockedClient: ApiClient;
let mockedCluster: Cluster;
let testClusterDetails: GetClusterResponse;

beforeEach(async () => {
({mockedCluster, mockedClient, testClusterDetails} =
await getMockTestCluster());

RetryConfigs.waitTime = (attempt) => {
return new Time(0, TimeUnits.milliseconds);
};
});

it("calling start on a non terminated state should not throw an error", async () => {
when(
mockedClient.request(
"/api/2.0/clusters/get",
"GET",
deepEqual({
cluster_id: testClusterDetails.cluster_id,
})
)
).thenResolve(
{
...testClusterDetails,
state: "PENDING",
},
{
...testClusterDetails,
state: "PENDING",
},
Comment thread
kartikgupta-db marked this conversation as resolved.
{
...testClusterDetails,
state: "PENDING",
},
{
...testClusterDetails,
state: "PENDING",
},
{
...testClusterDetails,
state: "PENDING",
},
{
...testClusterDetails,
state: "RUNNING",
}
);

await mockedCluster.refresh();
assert.notEqual(mockedCluster.state, "RUNNING");

await mockedCluster.start();
assert.equal(mockedCluster.state, "RUNNING");

verify(
mockedClient.request(
"/api/2.0/clusters/get",
anything(),
anything()
)
).times(6);

verify(
mockedClient.request(
"/api/2.0/clusters/start",
anything(),
anything()
)
).never();
});

it("should terminate cluster", async () => {
when(
mockedClient.request(
"/api/2.0/clusters/get",
"GET",
deepEqual({
cluster_id: testClusterDetails.cluster_id,
})
)
).thenResolve(
{
...testClusterDetails,
state: "RUNNING",
},
{
...testClusterDetails,
state: "TERMINATING",
},
{
...testClusterDetails,
state: "TERMINATED",
}
);

when(
mockedClient.request(
"/api/2.0/clusters/delete",
"POST",
deepEqual({
cluster_id: testClusterDetails.cluster_id,
})
)
).thenResolve({});

assert.equal(mockedCluster.state, "RUNNING");

await mockedCluster.stop();
assert.equal(mockedCluster.state, "TERMINATED");

verify(
mockedClient.request(
"/api/2.0/clusters/get",
anything(),
anything()
)
).times(3);

verify(
mockedClient.request(
"/api/2.0/clusters/delete",
anything(),
anything()
)
).once();
});

it("should terminate non running clusters", async () => {
when(
mockedClient.request(
"/api/2.0/clusters/get",
"GET",
deepEqual({
cluster_id: testClusterDetails.cluster_id,
})
)
).thenResolve(
{
...testClusterDetails,
state: "PENDING",
},
{
...testClusterDetails,
state: "TERMINATING",
},
{
...testClusterDetails,
state: "TERMINATED",
}
);

await mockedCluster.refresh();
assert.notEqual(mockedCluster.state, "RUNNING");

await mockedCluster.stop();
assert.equal(mockedCluster.state, "TERMINATED");
Comment thread
kartikgupta-db marked this conversation as resolved.

verify(
mockedClient.request(
"/api/2.0/clusters/get",
anything(),
anything()
)
).times(3);

verify(
mockedClient.request(
"/api/2.0/clusters/delete",
anything(),
anything()
)
).once();
});

it("should cancel cluster start", async () => {
const whenMockGetCluster = when(
mockedClient.request(
"/api/2.0/clusters/get",
"GET",
deepEqual({
cluster_id: testClusterDetails.cluster_id,
})
)
);

whenMockGetCluster.thenResolve({
...testClusterDetails,
state: "PENDING",
});

when(
mockedClient.request(
"/api/2.0/clusters/delete",
"POST",
deepEqual({
cluster_id: testClusterDetails.cluster_id,
})
)
).thenCall(() => {
whenMockGetCluster.thenResolve(
{
...testClusterDetails,
state: "TERMINATING",
},
{
...testClusterDetails,
state: "TERMINATED",
}
);
return {};
});

const token = mock(TokenFixture);
when(token.isCancellationRequested).thenReturn(false, false, true);
//mocked cluster is initially in running state, this gets it to pending state
await mockedCluster.refresh();

assert.equal(mockedCluster.state, "PENDING");
await mockedCluster.start(instance(token));

verify(token.isCancellationRequested).thrice();
Comment thread
kartikgupta-db marked this conversation as resolved.
assert.equal(mockedCluster.state, "TERMINATED");
});
});
Loading