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
8 changes: 6 additions & 2 deletions docs/operations/release.md
Original file line number Diff line number Diff line change
Expand Up @@ -210,8 +210,12 @@ A deleted legacy tunnel keeps its allocation, so its hostname is kept. When the
3. Run the legacy steps of the disposable-host canary below.
4. Before enabling, confirm the web and mobile builds that show the "update T3 Code on that computer"
message are live. Without them, a user whose older host lost its tunnel only sees it as offline.
5. Set the legacy mode to `enabled`. One sweep deletes at most 100 tunnels, so a backlog of
20,000 takes about 17 hours. Watch `deletedLegacy`, `failed`, and `truncated`.
5. Set the legacy mode to `enabled`. One sweep deletes at most 100 tunnels, four at a time, and
stops starting new deletions after 90 seconds. A backlog of 20,000 takes about 17 hours if each
sweep finishes its 100. Watch `deletedLegacy`, `attempted`, `failed`, and `truncated`. An
`attempted` well under 100 with `truncated` set means the sweep stopped early: either the time
budget ran out or Cloudflare rate-limited a deletion. The counters don't say which; the relay
logs a warning with the Cloudflare error for each failed deletion.

### Disposable-host canary

Expand Down
102 changes: 101 additions & 1 deletion infra/relay/src/environments/ManagedEndpointReaper.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,10 @@ function harness(input?: {
readonly cleanupMode?: RelayConfiguration.ManagedEndpointCleanupMode;
readonly legacyCleanupMode?: RelayConfiguration.ManagedEndpointCleanupMode;
readonly legacyTunnelGraceMinutes?: number;
/** Simulated time each release takes, advanced on the test clock. */
readonly releaseDelayMs?: number;
/** Simulated time each Cloudflare list request takes. */
readonly listDelayMs?: number;
}) {
const listRequests: ManagedEndpointProvider.ManagedEndpointTunnelListRequest[] = [];
const deleted: string[] = [];
Expand Down Expand Up @@ -119,7 +123,10 @@ function harness(input?: {
return Effect.succeed(found);
}),
list: (request) =>
Effect.sync(() => {
Effect.gen(function* () {
if (input?.listDelayMs !== undefined) {
yield* TestClock.adjust(input.listDelayMs);
}
listRequests.push(request);
const matching = remaining.filter((entry) => entry.status === request.status);
const start = ((request.page ?? 1) - 1) * (request.perPage ?? 100);
Expand Down Expand Up @@ -191,6 +198,9 @@ function harness(input?: {
release: (request) =>
Effect.gen(function* () {
releases.push(request);
if (input?.releaseDelayMs !== undefined) {
yield* TestClock.adjust(input.releaseDelayMs);
}
if (request.expectedTunnelId === input?.skipTunnelId) {
return false;
}
Expand Down Expand Up @@ -689,6 +699,96 @@ describe("ManagedEndpointReaper", () => {
}).pipe(Effect.provide(state.layer));
});

it.effect("starts no new deletion after a rate limit, even with deletions in flight", () => {
// Enough candidates to fill every concurrent slot several times over.
const entries = Array.from({ length: 12 }, (_, index) =>
tunnel({
id: index === 0 ? "limited" : `ok-${index}`,
suffix: index.toString(16).padStart(16, "0"),
status: "down",
timestamp: "2026-08-25T11:00:00.000Z",
}),
);
const state = harness({
tunnels: entries,
allocations: recoverableOwners(entries),
rateLimitedTunnelId: "limited",
});

return Effect.gen(function* () {
yield* TestClock.setTime(NOW_MILLIS);
const reaper = yield* ManagedEndpointReaper.ManagedEndpointReaper;
const result = yield* reaper.sweep;
expect(result.truncated).toBe(true);
expect(result.failed).toBe(1);
// Releases already running may finish, but none start afterwards.
expect(result.attempted).toBeLessThanOrEqual(
ManagedEndpointReaper.MANAGED_ENDPOINT_SWEEP_DELETE_CONCURRENCY,
);
expect(state.releases.length).toBe(result.attempted);
}).pipe(Effect.provide(state.layer));
});

it.effect("stops starting deletions when the sweep's time budget runs out", () => {
const entries = Array.from({ length: 40 }, (_, index) =>
tunnel({
id: `slow-${index}`,
suffix: index.toString(16).padStart(16, "0"),
status: "down",
timestamp: "2026-08-25T11:00:00.000Z",
}),
);
const state = harness({
tunnels: entries,
allocations: recoverableOwners(entries),
// Ten seconds each: the 90-second budget allows about 9 rounds.
releaseDelayMs: 10_000,
});

return Effect.gen(function* () {
yield* TestClock.setTime(NOW_MILLIS);
const reaper = yield* ManagedEndpointReaper.ManagedEndpointReaper;
const result = yield* reaper.sweep;
expect(result.truncated).toBe(true);
expect(result.attempted).toBeGreaterThan(0);
expect(result.attempted).toBeLessThan(entries.length);
expect(result.deleted).toBe(result.attempted);
}).pipe(Effect.provide(state.layer));
});

it.effect("counts listing time against the deletion budget", () => {
const entries = Array.from({ length: 40 }, (_, index) =>
tunnel({
id: `slow-${index}`,
suffix: index.toString(16).padStart(16, "0"),
status: "down",
timestamp: "2026-08-25T11:00:00.000Z",
}),
);
const sweepWith = (listDelayMs: number) =>
Effect.gen(function* () {
yield* TestClock.setTime(NOW_MILLIS);
const reaper = yield* ManagedEndpointReaper.ManagedEndpointReaper;
return (yield* reaper.sweep).attempted;
}).pipe(
Effect.provide(
harness({
tunnels: entries,
allocations: recoverableOwners(entries),
releaseDelayMs: 10_000,
listDelayMs,
}).layer,
),
);

return Effect.gen(function* () {
const fastListing = yield* sweepWith(0);
// Two list requests of 20 seconds each use 40 of the 90-second budget.
const slowListing = yield* sweepWith(20_000);
expect(slowListing).toBeLessThan(fastListing);
});
});

it.effect("continues past a page of older hosts to find recoverable tunnels", () => {
const entries = Array.from({ length: 101 }, (_, index) =>
tunnel({
Expand Down
114 changes: 79 additions & 35 deletions infra/relay/src/environments/ManagedEndpointReaper.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import * as Clock from "effect/Clock";
import * as Context from "effect/Context";
import * as DateTime from "effect/DateTime";
import * as Effect from "effect/Effect";
Expand Down Expand Up @@ -27,6 +28,16 @@ const MANAGED_ENDPOINT_LEGACY_AGE_BUCKET_DAYS = [7, 30, 90] as const;
// returning host that has updated recovers the tunnel at the same hostname;
// one that has not sees the client's "update T3 Code" message instead.
const MANAGED_ENDPOINT_LEGACY_GRACE_PERIOD_DAYS = 7;
// Deletions run a few at a time: each one is a row-locked database
// transaction plus two Cloudflare calls, so one at a time cannot finish a
// full attempt budget inside the cron's timeout. Each holds a Hyperdrive
// connection while it runs; keep this well under the 20-connection origin
// limit that request handlers share.
export const MANAGED_ENDPOINT_SWEEP_DELETE_CONCURRENCY = 4;
// Stop starting deletions this long after the sweep starts, leaving room under
// the cron's two-minute timeout to finish in-flight ones and record the
// counters. Counted from sweep start so slow listing eats into it.
const MANAGED_ENDPOINT_SWEEP_DELETE_BUDGET_MS = 90_000;

export interface ManagedEndpointSweepResult {
readonly mode: RelayConfiguration.ManagedEndpointCleanupMode;
Expand Down Expand Up @@ -180,6 +191,7 @@ export const make = Effect.gen(function* () {
if ((mode === "off" && legacyMode === "off") || !namespace) {
return emptyResult(mode, legacyMode);
}
const sweepStartedAtMillis = yield* Clock.currentTimeMillis;
// The override exists for the disposable canary stage; prod always
// waits the full grace period.
const legacyGraceMinutes =
Expand Down Expand Up @@ -275,6 +287,16 @@ export const make = Effect.gen(function* () {
let attempted = 0;
let deleted = 0;
let wouldDelete = 0;
const candidates: Array<{
readonly owner: ManagedEndpointAllocations.ManagedEndpointTunnelAllocation;
readonly tunnel: ManagedEndpointProvider.ManagedEndpointTunnel & {
readonly id: string;
readonly name: string;
};
readonly status: "down" | "inactive";
readonly legacy: boolean;
readonly inactiveBefore: string;
}> = [];
let wouldDeleteLegacy = 0;
let deletedLegacy = 0;
let skippedLegacy = 0;
Expand Down Expand Up @@ -330,46 +352,68 @@ export const make = Effect.gen(function* () {
wouldDelete += 1;
if (mode === "dry-run") continue;
}
if (attempted >= MANAGED_ENDPOINT_SWEEP_ATTEMPT_LIMIT) {
if (candidates.length >= MANAGED_ENDPOINT_SWEEP_ATTEMPT_LIMIT) {
truncated = true;
break;
}
attempted += 1;
// The release re-reads the tunnel and deletes only if it is still in
// this status and inactive since before this tunnel's cutoff.
const result = yield* provider
.release({
userId: owner.userId,
environmentId: owner.environmentId,
expectedTunnelId: tunnel.id,
expectedInactiveBefore: DateTime.formatIso(legacy ? legacyCutoff : cutoff),
expectedStatus: status,
})
.pipe(Effect.result);
if (result._tag === "Failure") {
failed += 1;
yield* Effect.logWarning("Failed to delete an inactive managed tunnel", {
tunnelId: tunnel.id,
tunnelName: tunnel.name,
legacy,
cause: result.failure,
});
if (isRateLimited(result.failure)) {
truncated = true;
break;
}
} else if (result.success) {
deleted += 1;
if (legacy) deletedLegacy += 1;
yield* Effect.logInfo("Deleted an inactive managed tunnel", {
tunnelId: tunnel.id,
tunnelName: tunnel.name,
status,
legacy,
});
}
candidates.push({
owner,
tunnel,
status,
legacy,
// The release re-reads the tunnel and deletes only if it is still in
// this status and inactive since before this cutoff.
inactiveBefore: DateTime.formatIso(legacy ? legacyCutoff : cutoff),
});
}

const deleteDeadline = sweepStartedAtMillis + MANAGED_ENDPOINT_SWEEP_DELETE_BUDGET_MS;
let stopDeleting = false;
yield* Effect.forEach(
candidates,
(candidate) =>
Effect.gen(function* () {
if (stopDeleting || (yield* Clock.currentTimeMillis) >= deleteDeadline) {
stopDeleting = true;
truncated = true;
return;
}
attempted += 1;
const result = yield* provider
.release({
userId: candidate.owner.userId,
environmentId: candidate.owner.environmentId,
expectedTunnelId: candidate.tunnel.id,
expectedInactiveBefore: candidate.inactiveBefore,
expectedStatus: candidate.status,
})
.pipe(Effect.result);
if (result._tag === "Failure") {
failed += 1;
yield* Effect.logWarning("Failed to delete an inactive managed tunnel", {
tunnelId: candidate.tunnel.id,
tunnelName: candidate.tunnel.name,
legacy: candidate.legacy,
cause: result.failure,
});
if (isRateLimited(result.failure)) {
stopDeleting = true;
truncated = true;
}
} else if (result.success) {
deleted += 1;
if (candidate.legacy) deletedLegacy += 1;
yield* Effect.logInfo("Deleted an inactive managed tunnel", {
tunnelId: candidate.tunnel.id,
tunnelName: candidate.tunnel.name,
status: candidate.status,
legacy: candidate.legacy,
});
}
}),
{ concurrency: MANAGED_ENDPOINT_SWEEP_DELETE_CONCURRENCY, discard: true },
);

return {
mode,
legacyMode,
Expand Down
Loading