Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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
7 changes: 7 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,13 @@ All notable changes to this extension will be documented in this file.

## Unreleased

### Fixed

- Flink statement results no longer stop loading when Confluent Cloud returns a temporary error
right after a statement is submitted. The Results Viewer now retries before giving up, honoring
the `Retry-After` delay when Confluent Cloud sends one, instead of showing "Failed to load
results."

## 2.3.1

### Fixed
Expand Down
22 changes: 22 additions & 0 deletions src/errors.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import {
getNestedErrorChain,
hasErrorCause,
isResponseErrorWithStatus,
isTransientResponseError,
logError,
} from "./errors";
import { Logger } from "./logging";
Expand Down Expand Up @@ -169,6 +170,27 @@ describe("errors.ts isResponseErrorWithStatus()", () => {
});
});

describe("errors.ts isTransientResponseError()", () => {
it("should return false for not-a-response-error", () => {
const error = new Error("test");
assert.strictEqual(isTransientResponseError(error), false);
});

for (const status of [429, 500, 502, 503, 504]) {
it(`should return true for a ${status} response error`, () => {
const error = createResponseError(status, "Transient", "test");
assert.strictEqual(isTransientResponseError(error), true);
});
}

for (const status of [400, 401, 403, 404, 409]) {
it(`should return false for a ${status} response error`, () => {
const error = createResponseError(status, "Client Error", "test");
assert.strictEqual(isTransientResponseError(error), false);
});
}
});

describe("errors.ts extractResponseBody()", () => {
it("should return the response body as JSON if it is valid JSON", async () => {
const embeddedObject = { message: "test" };
Expand Down
8 changes: 8 additions & 0 deletions src/errors.ts
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,14 @@ export function isResponseErrorWithStatus(
return isResponseError(error) && error.response.status === statusCode;
}

/** HTTP statuses where repeating the same request after a short delay may succeed. */
const TRANSIENT_RESPONSE_STATUSES = [429, 500, 502, 503, 504];

/** Was this a response error whose status suggests the request is worth retrying? */
export function isTransientResponseError(error: unknown): error is AnyResponseError {
return isResponseError(error) && TRANSIENT_RESPONSE_STATUSES.includes(error.response.status);
}

/**
* If error is a response error, try to decode its response body
* from JSON and return the resulting object.
Expand Down
164 changes: 145 additions & 19 deletions src/flinkSql/flinkStatementResultsManager.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,20 +5,33 @@ import type { FlinkStatementResultsManagerTestContext } from "../../tests/create
import { createTestResultsManagerContext } from "../../tests/createResultsManager";
import { eventually } from "../../tests/eventually";
import { loadFixtureFromFile } from "../../tests/fixtures/utils";
import { createResponseError } from "../../tests/unit/testUtils";
import {
createResponseError,
createSingleUseResponseError,
ResponseErrorSource,
} from "../../tests/unit/testUtils";
import type { GetSqlv1StatementResult200Response } from "../clients/flinkSql";
import {
GetSqlv1StatementResult200ResponseApiVersionEnum,
GetSqlv1StatementResult200ResponseKindEnum,
} from "../clients/flinkSql";
import * as messageUtils from "../documentProviders/message";
import { transientBackoffWindow } from "./flinkStatementResultsManager";
import { FlinkStatement, Phase } from "../models/flinkStatement";
import type { WebviewStorage } from "../webview/comms/comms";
import type {
FlinkStatementResultsViewModel,
ResultsViewerStorageState,
} from "../webview/flink-statement-results";

/** A successful results response carrying no rows. */
const EMPTY_RESULTS_RESPONSE: GetSqlv1StatementResult200Response = {
api_version: GetSqlv1StatementResult200ResponseApiVersionEnum.SqlV1,
kind: GetSqlv1StatementResult200ResponseKindEnum.StatementResult,
metadata: {},
results: { data: [] },
};

function createMockStatement(): FlinkStatement {
const fakeFlinkStatement = loadFixtureFromFile(
"flink-statement-results-processing/fake-flink-statement.json",
Expand Down Expand Up @@ -401,14 +414,9 @@ describe("FlinkStatementResultsViewModel and FlinkStatementResultsManager", () =
ctx.flinkSqlStatementResultsApi.getSqlv1StatementResult
.onSecondCall()
.rejects(createResponseError(409, "Conflict", "{}"));
ctx.flinkSqlStatementResultsApi.getSqlv1StatementResult.onThirdCall().resolves({
api_version: GetSqlv1StatementResult200ResponseApiVersionEnum.SqlV1,
kind: GetSqlv1StatementResult200ResponseKindEnum.StatementResult,
metadata: {},
results: {
data: [],
},
});
ctx.flinkSqlStatementResultsApi.getSqlv1StatementResult
.onThirdCall()
.resolves(EMPTY_RESULTS_RESPONSE);

// Trigger a fetch
const fetchPromise = ctx.manager.fetchResults();
Expand Down Expand Up @@ -442,9 +450,8 @@ describe("FlinkStatementResultsViewModel and FlinkStatementResultsManager", () =
assert.ok(ctx.manager["_latestError"]());
});

it("should not retry on non-409 errors during fetch", async () => {
// Mock the getSqlv1StatementResult to fail with 500
const responseError = createResponseError(500, "Internal Server Error", "{}");
it("should not retry on errors that are neither 409 nor transient during fetch", async () => {
const responseError = createResponseError(403, "Forbidden", "{}");
ctx.flinkSqlStatementResultsApi.getSqlv1StatementResult.rejects(responseError);

// Trigger a fetch
Expand All @@ -455,11 +462,91 @@ describe("FlinkStatementResultsViewModel and FlinkStatementResultsManager", () =

await fetchPromise;

assert.equal(ctx.flinkSqlStatementResultsApi.getSqlv1StatementResult.callCount, 1);
sinon.assert.calledOnce(ctx.flinkSqlStatementResultsApi.getSqlv1StatementResult);
// Verify error state is set
assert.ok(ctx.manager["_latestError"]());
});

it("should retry get statement results on transient errors", async () => {
// CCloud briefly can't resolve a just-created statement, answering 429 or 5xx before the
// results endpoint starts working
ctx.flinkSqlStatementResultsApi.getSqlv1StatementResult
.onFirstCall()
.rejects(createResponseError(429, "Too Many Requests", "{}"));
ctx.flinkSqlStatementResultsApi.getSqlv1StatementResult
.onSecondCall()
.rejects(createResponseError(500, "Internal Server Error", "{}"));
ctx.flinkSqlStatementResultsApi.getSqlv1StatementResult
.onThirdCall()
.resolves(EMPTY_RESULTS_RESPONSE);

const fetchPromise = ctx.manager.fetchResults();

// backoff doubles from 500ms and is jittered, so tick past the two maximums
await clock.tickAsync(500 + 1000);

await fetchPromise;

sinon.assert.calledThrice(ctx.flinkSqlStatementResultsApi.getSqlv1StatementResult);
assert.equal(ctx.manager["_latestError"](), null);
});

it("should wait for the server's Retry-After rather than the exponential curve", async () => {
ctx.flinkSqlStatementResultsApi.getSqlv1StatementResult.onFirstCall().rejects(
createResponseError(429, "Too Many Requests", "{}", ResponseErrorSource.Sidecar, {
"retry-after": "2",
}),
);
ctx.flinkSqlStatementResultsApi.getSqlv1StatementResult
.onSecondCall()
.resolves(EMPTY_RESULTS_RESPONSE);

const fetchPromise = ctx.manager.fetchResults();

// the exponential curve would have retried by now, but the server asked for 2s
await clock.tickAsync(1000);
sinon.assert.calledOnce(ctx.flinkSqlStatementResultsApi.getSqlv1StatementResult);

// 2s as requested, plus up to one base delay of jitter on top
await clock.tickAsync(1500);
await fetchPromise;

sinon.assert.calledTwice(ctx.flinkSqlStatementResultsApi.getSqlv1StatementResult);
assert.equal(ctx.manager["_latestError"](), null);
});

it("should complete the stream after exhausting transient retries", async () => {
ctx.flinkSqlStatementResultsApi.getSqlv1StatementResult.rejects(
createResponseError(500, "Internal Server Error", "{}"),
);

const fetchPromise = ctx.manager.fetchResults();

// 4 transient retries at up to 500/1000/2000/4000ms
await clock.tickAsync(7500);

await fetchPromise;

// 1 initial attempt + MAX_TRANSIENT_RETRIES
sinon.assert.callCount(ctx.flinkSqlStatementResultsApi.getSqlv1StatementResult, 5);
assert.ok(ctx.manager["_latestError"]());
assert.equal(ctx.manager["_state"](), "completed");
});

it("should leave the error response body readable for logging", async () => {
// a real single-use Response, so reading the body without cloning would be observable
const responseError = createSingleUseResponseError(
400,
"Bad Request",
'{"errors":[{"code":"cr_failed_get_stmt_name"}]}',
);
ctx.flinkSqlStatementResultsApi.getSqlv1StatementResult.rejects(responseError);

await ctx.manager.fetchResults();

assert.strictEqual(responseError.response.bodyUsed, false);
});

it("should only allow one instance of fetchResults to run at a time", async () => {
// Create a promise that we can resolve manually to simulate a slow API call
let resolveRequest: (value: GetSqlv1StatementResult200Response) => void;
Expand All @@ -486,12 +573,7 @@ describe("FlinkStatementResultsViewModel and FlinkStatementResultsManager", () =
assert.equal(ctx.flinkSqlStatementResultsApi.getSqlv1StatementResult.callCount, 1);

// Resolve the API call
resolveRequest!({
api_version: GetSqlv1StatementResult200ResponseApiVersionEnum.SqlV1,
kind: GetSqlv1StatementResult200ResponseKindEnum.StatementResult,
metadata: {},
results: { data: [] },
});
resolveRequest!(EMPTY_RESULTS_RESPONSE);

// Wait for all calls to complete
await Promise.all(fetchPromises);
Expand Down Expand Up @@ -768,3 +850,47 @@ describe("FlinkStatementResultsViewModel only", () => {
}
});
});

describe("flinkStatementResultsManager.ts transientBackoffWindow()", () => {
function responseWithHeaders(headers: Record<string, string>): Response {
return new Response("{}", { status: 429, headers });
}

it("should honor a Retry-After the server sends, never shortening it", () => {
const { minMs, maxMs } = transientBackoffWindow(responseWithHeaders({ "retry-after": "2" }), 0);

assert.equal(minMs, 2000);
assert.ok(maxMs > minMs, "jitter should only extend a server-requested delay");
});

it("should cap an outsized Retry-After", () => {
const { minMs } = transientBackoffWindow(responseWithHeaders({ "retry-after": "600" }), 0);

assert.equal(minMs, 8000);
});

it("should fall back to an exponential curve without a Retry-After", () => {
const windows = [0, 1, 2, 3].map((attempt) =>
transientBackoffWindow(responseWithHeaders({}), attempt),
);

assert.deepEqual(
windows.map((w) => w.maxMs),
[500, 1000, 2000, 4000],
);
assert.deepEqual(
windows.map((w) => w.minMs),
[250, 500, 1000, 2000],
);
});

it("should fall back to the curve for a non-numeric Retry-After", () => {
// RFC 7231 also permits an HTTP-date, which we don't parse
const { maxMs } = transientBackoffWindow(
responseWithHeaders({ "retry-after": "Wed, 21 Oct 2015 07:28:00 GMT" }),
0,
);

assert.equal(maxMs, 500);
});
});
Loading