Skip to content
Closed
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
54 changes: 54 additions & 0 deletions src/deploy/functions/release/fabricator.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -344,6 +344,24 @@
await fab.createV1Function(ep1, new scraper.SourceTokenScraper());
expect(gcf.setInvokerCreate).to.not.have.been.called;
});

it("does not deadlock when createFunction receives a retryable 409 on the first attempt", async () => {
const retryingFab = new fabricator.Fabricator({
...ctorArgs,
functionExecutor: new executor.QueueExecutor({ retries: 1, backoff: 10, maxBackoff: 10 }),
});

const err409 = new Error("unable to queue the operation");
(err409 as any).status = 409;

Check warning on line 355 in src/deploy/functions/release/fabricator.spec.ts

View workflow job for this annotation

GitHub Actions / lint (24)

Unexpected any. Specify a different type

Check warning on line 355 in src/deploy/functions/release/fabricator.spec.ts

View workflow job for this annotation

GitHub Actions / lint (24)

Unsafe member access .status on an `any` value
gcf.createFunction.onFirstCall().rejects(err409);
gcf.createFunction.resolves({ name: "op", type: "create", done: false });
poller.pollOperation.resolves();

const ep = endpoint({ scheduleTrigger: {} });
await expect(retryingFab.createV1Function(ep, new scraper.SourceTokenScraper())).to.eventually
.be.undefined;
expect(gcf.createFunction).to.have.been.calledTwice;
});
});

describe("updateV1Function", () => {
Expand Down Expand Up @@ -432,6 +450,24 @@
await fab.updateV1Function(ep, new scraper.SourceTokenScraper());
expect(gcf.setInvokerUpdate).to.not.have.been.called;
});

it("does not deadlock when updateFunction receives a retryable 409 on the first attempt", async () => {
const retryingFab = new fabricator.Fabricator({
...ctorArgs,
functionExecutor: new executor.QueueExecutor({ retries: 1, backoff: 10, maxBackoff: 10 }),
});

const err409 = new Error("unable to queue the operation");
(err409 as any).status = 409;

Check warning on line 461 in src/deploy/functions/release/fabricator.spec.ts

View workflow job for this annotation

GitHub Actions / lint (24)

Unexpected any. Specify a different type

Check warning on line 461 in src/deploy/functions/release/fabricator.spec.ts

View workflow job for this annotation

GitHub Actions / lint (24)

Unsafe member access .status on an `any` value
gcf.updateFunction.onFirstCall().rejects(err409);
gcf.updateFunction.resolves({ name: "op", type: "update", done: false });
poller.pollOperation.resolves();

const ep = endpoint({ httpsTrigger: {} });
await expect(retryingFab.updateV1Function(ep, new scraper.SourceTokenScraper())).to.eventually
.be.undefined;
expect(gcf.updateFunction).to.have.been.calledTwice;
});
});

describe("deleteV1Function", () => {
Expand All @@ -452,7 +488,7 @@
it("handles topics that already exist", async () => {
pubsub.createTopic.callsFake(() => {
const err = new Error("Already exists");
(err as any).status = 409;

Check warning on line 491 in src/deploy/functions/release/fabricator.spec.ts

View workflow job for this annotation

GitHub Actions / lint (24)

Unexpected any. Specify a different type

Check warning on line 491 in src/deploy/functions/release/fabricator.spec.ts

View workflow job for this annotation

GitHub Actions / lint (24)

Unsafe member access .status on an `any` value
return Promise.reject(err);
});
gcfv2.createFunction.resolves({ name: "op", done: false });
Expand Down Expand Up @@ -527,7 +563,7 @@
eventarc.createChannel.callsFake(({ name }) => {
expect(name).to.equal("channel");
const err = new Error("Already exists");
(err as any).status = 409;

Check warning on line 566 in src/deploy/functions/release/fabricator.spec.ts

View workflow job for this annotation

GitHub Actions / lint (24)

Unexpected any. Specify a different type

Check warning on line 566 in src/deploy/functions/release/fabricator.spec.ts

View workflow job for this annotation

GitHub Actions / lint (24)

Unsafe member access .status on an `any` value
return Promise.reject(err);
});
gcfv2.createFunction.resolves({ name: "op", done: false });
Expand Down Expand Up @@ -593,7 +629,7 @@
eventarc.getChannel.resolves(undefined);
eventarc.createChannel.callsFake(() => {
const err = new Error("🤷‍♂️");
(err as any).status = 400;

Check warning on line 632 in src/deploy/functions/release/fabricator.spec.ts

View workflow job for this annotation

GitHub Actions / lint (24)

Unexpected any. Specify a different type

Check warning on line 632 in src/deploy/functions/release/fabricator.spec.ts

View workflow job for this annotation

GitHub Actions / lint (24)

Unsafe member access .status on an `any` value
return Promise.reject(err);
});

Expand Down Expand Up @@ -942,6 +978,24 @@
);
});

it("does not deadlock when updateFunction receives a retryable 409 on the first attempt", async () => {
const retryingFab = new fabricator.Fabricator({
...ctorArgs,
functionExecutor: new executor.QueueExecutor({ retries: 1, backoff: 10, maxBackoff: 10 }),
});

const err409 = new Error("unable to queue the operation");
(err409 as any).status = 409;
gcfv2.updateFunction.onFirstCall().rejects(err409);
gcfv2.updateFunction.resolves({ name: "op", done: false });
poller.pollOperation.resolves({ serviceConfig: { service: "service" } });

const ep = endpoint({ httpsTrigger: {} }, { platform: "gcfv2" });
await expect(retryingFab.updateV2Function(ep, new scraper.SourceTokenScraper())).to.eventually
.be.undefined;
expect(gcfv2.updateFunction).to.have.been.calledTwice;
});

it("throws on set invoker failure", async () => {
gcfv2.updateFunction.resolves({ name: "op", done: false });
poller.pollOperation.resolves({ serviceConfig: { service: "service" } });
Expand Down
20 changes: 19 additions & 1 deletion src/deploy/functions/release/sourceTokenScraper.ts
Original file line number Diff line number Diff line change
@@ -1,11 +1,17 @@
import { FirebaseError } from "../../../error";
import { assertExhaustive } from "../../../functional";
import { logger } from "../../../logger";
import { timeoutFallback } from "../../../timeout";

// How long a concurrent getToken() call will wait for the first deploy's poller
// to produce a token before giving up and proceeding without one.
const SOURCE_TOKEN_FETCH_TIMEOUT_MS = 25 * 60 * 1000; // 25 minutes, matches fabricator.ts masterTimeout

type TokenFetchState = "NONE" | "FETCHING" | "VALID";
interface TokenFetchResult {
token?: string;
aborted: boolean;
timedOut?: boolean;
}

/**
Expand Down Expand Up @@ -54,8 +60,20 @@ export class SourceTokenScraper {
this.fetchState = "FETCHING";
return undefined;
} else if (this.fetchState === "FETCHING") {
const tokenResult = await this.promise;
// Race the token promise against a deadline so a lost token producer
// (e.g. a failed deploy that didn't call abort()) degrades to a
// token-less deploy rather than hanging the process forever.
const tokenResult: TokenFetchResult = await timeoutFallback(
this.promise,
{ aborted: true, timedOut: true },
SOURCE_TOKEN_FETCH_TIMEOUT_MS,
);
if (tokenResult.aborted) {
if (tokenResult.timedOut) {
logger.warn(
"Timed out waiting for a source token. Proceeding without one, which may slow the deploy.",
);
}
return undefined;
}
return tokenResult.token;
Expand Down
15 changes: 11 additions & 4 deletions src/timeout.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,10 +6,17 @@ export async function timeoutFallback<T, V>(
value: V,
timeoutMillis = 2000,
): Promise<T | V> {
return Promise.race([
promise,
new Promise<V>((resolve) => setTimeout(() => resolve(value), timeoutMillis)),
]);
let timer: NodeJS.Timeout | undefined;
try {
return await Promise.race([
promise,
new Promise<V>((resolve) => {
timer = setTimeout(() => resolve(value), timeoutMillis);
}),
]);
} finally {
clearTimeout(timer);
}
}

export async function timeoutError<T>(
Expand Down
Loading