diff --git a/apps/server/src/telemetry/AnalyticsService.test.ts b/apps/server/src/telemetry/AnalyticsService.test.ts index c135467efffa..e2966306329d 100644 --- a/apps/server/src/telemetry/AnalyticsService.test.ts +++ b/apps/server/src/telemetry/AnalyticsService.test.ts @@ -4,6 +4,10 @@ import { assert, it } from "@effect/vitest"; import * as ConfigProvider from "effect/ConfigProvider"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; +import * as Schema from "effect/Schema"; +import * as TestClock from "effect/testing/TestClock"; +import * as HttpClient from "effect/unstable/http/HttpClient"; +import * as HttpClientError from "effect/unstable/http/HttpClientError"; import * as HttpServer from "effect/unstable/http/HttpServer"; import * as HttpServerRequest from "effect/unstable/http/HttpServerRequest"; import * as HttpServerResponse from "effect/unstable/http/HttpServerResponse"; @@ -46,7 +50,86 @@ interface RecordedBatchBody { }>; } +const SentBatch = Schema.fromJsonString( + Schema.Struct({ + batch: Schema.Array(Schema.Struct({ uuid: Schema.String })), + }), +); + +/** + * HTTP client that reads each batch, then fails as if the connection dropped + * before the response arrived. PostHog stores these batches, so the server + * must not send them forever. + */ +const acceptThenFailClient = (batches: Array>) => + Layer.succeed( + HttpClient.HttpClient, + HttpClient.make((request) => + Effect.gen(function* () { + if (request.body._tag === "Uint8Array") { + const body = yield* Schema.decodeEffect(SentBatch)( + new TextDecoder().decode(request.body.body), + ).pipe(Effect.orDie); + batches.push(body.batch); + } + return yield* new HttpClientError.HttpClientError({ + reason: new HttpClientError.TransportError({ request, cause: "connection reset" }), + }); + }), + ), + ); + +it("retryDelayMs doubles from 2s and stays under the 5 minute cap", () => { + assert.equal(AnalyticsService.retryDelayMs(1, 0), 1_000); + assert.equal(AnalyticsService.retryDelayMs(2, 0.999_999), 4_000); + assert.equal(AnalyticsService.retryDelayMs(30, 0), 150_000); + assert.equal(AnalyticsService.retryDelayMs(30, 0.999_999), 300_000); +}); + it.layer(NodeServices.layer)("AnalyticsService test", (it) => { + it.effect("a batch that keeps failing is retried with backoff, then dropped", () => + Effect.gen(function* () { + const batches: Array> = []; + const runtimeLayer = AnalyticsService.layer.pipe( + Layer.provideMerge( + ServerConfig.ServerConfig.layerTest(process.cwd(), { prefix: "t3-telemetry-retry-" }), + ), + Layer.provide( + ConfigProvider.layer( + ConfigProvider.fromUnknown({ + T3CODE_TELEMETRY_ENABLED: true, + T3CODE_POSTHOG_KEY: "phc_test_key", + T3CODE_POSTHOG_HOST: "http://localhost", + }), + ), + ), + Layer.provide( + Layer.mergeAll( + Layer.succeed(HostProcessPlatform, "win32"), + Layer.succeed(HostProcessArchitecture, "x64"), + acceptThenFailClient(batches), + ), + ), + ); + + yield* Effect.gen(function* () { + const analytics = yield* AnalyticsService.AnalyticsService; + for (let index = 0; index < 20; index += 1) { + yield* analytics.record("test.retry", { index }); + } + // Before the fix this loop sent the batch about once a second. + for (let second = 0; second < 600; second += 1) { + yield* TestClock.adjust("1 second"); + } + }).pipe(Effect.provide(runtimeLayer)); + + assert.equal(batches.length, 5); + const uuids = batches.map((batch) => batch.map((event) => event.uuid).join(",")); + assert.equal(new Set(uuids).size, 1, "every retry carries the same uuids"); + assert.equal(new Set(batches[0]?.map((event) => event.uuid)).size, 20); + }), + ); + it.effect("flush drains all buffered events across multiple batches", () => Effect.gen(function* () { const capturedRequests: Array = []; diff --git a/apps/server/src/telemetry/AnalyticsService.ts b/apps/server/src/telemetry/AnalyticsService.ts index d663a8f321e9..4d1db714eb96 100644 --- a/apps/server/src/telemetry/AnalyticsService.ts +++ b/apps/server/src/telemetry/AnalyticsService.ts @@ -2,19 +2,27 @@ * Anonymous PostHog telemetry service. * * Persists an installation-scoped anonymous identifier, buffers events in - * memory, and flushes batches over Effect's HTTP client. + * memory, and flushes batches over Effect's HTTP client. A failed batch is + * retried with backoff and dropped after a few tries. Each event carries a + * uuid, so PostHog can tell a retried copy from a new event. * * @module AnalyticsService */ import { HostProcessArchitecture, HostProcessPlatform } from "@t3tools/shared/hostProcess"; import type { ClientOs } from "@t3tools/contracts"; +import * as Clock from "effect/Clock"; import * as Config from "effect/Config"; import * as Context from "effect/Context"; +import * as Crypto from "effect/Crypto"; import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; +import * as Exit from "effect/Exit"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; +import * as Random from "effect/Random"; import * as Ref from "effect/Ref"; +import * as Result from "effect/Result"; +import * as Semaphore from "effect/Semaphore"; import * as HttpClient from "effect/unstable/http/HttpClient"; import * as HttpClientRequest from "effect/unstable/http/HttpClientRequest"; import * as HttpClientResponse from "effect/unstable/http/HttpClientResponse"; @@ -24,11 +32,40 @@ import * as ServerConfig from "../config.ts"; import { getTelemetryIdentifier } from "./Identify.ts"; interface BufferedAnalyticsEvent { + readonly uuid: string; readonly event: string; readonly properties?: Readonly>; readonly capturedAt: string; } +interface DeliveryState { + /** Batch that failed last. It is sent again before newer events. */ + readonly failedBatch: ReadonlyArray; + /** Failed sends of `failedBatch`. */ + readonly batchAttempts: number; + /** Failed sends since the last success. Sets the backoff delay. */ + readonly failures: number; + /** The background flush does not send before this time (epoch ms). */ + readonly retryAt: number; +} + +const FLUSH_INTERVAL_MS = 1_000; +// A hung send would hold the flush lock, and with it the shutdown flush. +const SEND_TIMEOUT = "10 seconds"; +const MAX_BATCH_ATTEMPTS = 5; +const RETRY_BASE_DELAY_MS = 2_000; +const RETRY_MAX_DELAY_MS = 300_000; + +/** + * Delay before the next send after `failures` consecutive failed sends. The + * ceiling doubles from 2s up to 5 minutes, and the delay is a random point in + * its upper half. `random` is in [0, 1). + */ +export function retryDelayMs(failures: number, random: number): number { + const ceiling = Math.min(RETRY_MAX_DELAY_MS, RETRY_BASE_DELAY_MS * 2 ** (failures - 1)); + return Math.round(ceiling / 2 + (ceiling / 2) * random); +} + const TelemetryEnvConfig = Config.all({ posthogKey: Config.String("T3CODE_POSTHOG_KEY").pipe( Config.withDefault("phc_XOWci4oZP4VvLiEyrFqkFjP4CZn55mjYYBMREK5Wd6m"), @@ -88,17 +125,31 @@ export const make = Effect.gen(function* () { const httpClient = yield* HttpClient.HttpClient; const serverConfig = yield* ServerConfig.ServerConfig; const identifier = yield* getTelemetryIdentifier; + const crypto = yield* Crypto.Crypto; const bufferRef = yield* Ref.make>([]); + const deliveryRef = yield* Ref.make({ + failedBatch: [], + batchAttempts: 0, + failures: 0, + retryAt: 0, + }); + // The background flush and the shutdown flush must not send the same batch at once. + const flushLock = yield* Semaphore.make(1); const clientType = serverConfig.mode === "desktop" ? "desktop-app" : "cli-web-client"; const hostPlatform = yield* HostProcessPlatform; const hostArchitecture = yield* HostProcessArchitecture; - const enqueueBufferedEvent = (event: string, properties?: Readonly>) => + const enqueueBufferedEvent = ( + uuid: string, + event: string, + properties?: Readonly>, + ) => Effect.flatMap(DateTime.now, (now) => Ref.modify(bufferRef, (current) => { const appended = [ ...current, { + uuid, event, ...(properties ? { properties } : {}), capturedAt: DateTime.formatIso(now), @@ -128,6 +179,7 @@ export const make = Effect.gen(function* () { const payload = { api_key: telemetryConfig.posthogKey, batch: events.map((event) => ({ + uuid: event.uuid, event: event.event, distinct_id: identifier, properties: { @@ -152,39 +204,69 @@ export const make = Effect.gen(function* () { HttpClientRequest.bodyJson(payload), Effect.flatMap(httpClient.execute), Effect.flatMap(HttpClientResponse.filterStatusOk), + Effect.timeout(SEND_TIMEOUT), ); }); + const takeBatch = Ref.modify(bufferRef, (current) => { + const nextBatch = current.slice(0, telemetryConfig.flushBatchSize); + return [nextBatch, current.slice(nextBatch.length)] as const; + }); + + // Sends batches until the buffer is empty or a send fails. A failed batch is + // kept for the next flush, and dropped after MAX_BATCH_ATTEMPTS failed sends. const flush: AnalyticsService["Service"]["flush"] = Effect.gen(function* () { while (true) { - const batch = yield* Ref.modify(bufferRef, (current) => { - if (current.length === 0) { - return [[] as ReadonlyArray, current] as const; - } - const nextBatch = current.slice(0, telemetryConfig.flushBatchSize); - const remaining = current.slice(nextBatch.length); - return [nextBatch, remaining] as const; - }); - + const delivery = yield* Ref.get(deliveryRef); + const batch = delivery.failedBatch.length > 0 ? delivery.failedBatch : yield* takeBatch; if (batch.length === 0) { return; } - yield* sendBatch(batch).pipe( - Effect.catch((error) => - Ref.update(bufferRef, (current) => [...batch, ...current]).pipe( - Effect.flatMap(() => Effect.fail(error)), - ), - ), - ); + const sent = yield* Effect.result(sendBatch(batch)); + if (Result.isSuccess(sent)) { + yield* Ref.set(deliveryRef, { failedBatch: [], batchAttempts: 0, failures: 0, retryAt: 0 }); + continue; + } + + const failures = delivery.failures + 1; + const batchAttempts = delivery.batchAttempts + 1; + const retryAt = (yield* Clock.currentTimeMillis) + retryDelayMs(failures, yield* Random.next); + if (batchAttempts < MAX_BATCH_ATTEMPTS) { + yield* Ref.set(deliveryRef, { failedBatch: batch, batchAttempts, failures, retryAt }); + yield* Effect.logDebug("Failed to send telemetry batch; will retry", { + attempt: batchAttempts, + cause: sent.failure, + }); + return; + } + yield* Ref.set(deliveryRef, { failedBatch: [], batchAttempts: 0, failures, retryAt }); + yield* Effect.logWarning("Dropped telemetry batch after repeated send failures", { + events: batch.length, + attempts: batchAttempts, + cause: sent.failure, + }); + return; } - }).pipe(Effect.catch((cause) => Effect.logError("Failed to flush telemetry", { cause }))); + }).pipe(flushLock.withPermit); + + const flushWhenDue = Effect.gen(function* () { + const { retryAt } = yield* Ref.get(deliveryRef); + if ((yield* Clock.currentTimeMillis) >= retryAt) { + yield* flush; + } + }); const record: AnalyticsService["Service"]["record"] = Effect.fn("AnalyticsService.record")( function* (event, properties) { if (!telemetryConfig.enabled || !identifier) return; - const enqueueResult = yield* enqueueBufferedEvent(event, properties); + // Telemetry is best effort: an event without a uuid is not sent. The + // Node implementation throws (a defect) rather than failing, so catch both. + const uuid = yield* Effect.exit(crypto.randomUUIDv7); + if (Exit.isFailure(uuid)) return; + + const enqueueResult = yield* enqueueBufferedEvent(uuid.value, event, properties); if (enqueueResult.dropped) { yield* Effect.logDebug("analytics buffer full; dropping oldest event", { size: enqueueResult.size, @@ -194,7 +276,7 @@ export const make = Effect.gen(function* () { }, ); - yield* Effect.forever(Effect.sleep(1000).pipe(Effect.flatMap(() => flush)), { + yield* Effect.forever(Effect.sleep(FLUSH_INTERVAL_MS).pipe(Effect.flatMap(() => flushWhenDue)), { disableYield: true, }).pipe(Effect.forkScoped); diff --git a/docs/internals/connection-runtime.md b/docs/internals/connection-runtime.md index 301e1e0ffecb..35c8b0a73658 100644 --- a/docs/internals/connection-runtime.md +++ b/docs/internals/connection-runtime.md @@ -10,16 +10,20 @@ several views need the same environment. The [supervisor](../../packages/client-runtime/src/connection/supervisor.ts) owns transport retry policy; resolving an endpoint and opening an RPC session are single -attempts. Transient failures retry with capped backoff. Offline states and -authentication failures wait for a wakeup instead of spending attempts on -unchanged conditions. +attempts. Transient failures retry with jittered exponential backoff, capped at +five minutes, that resets only after a connection stays up. Without jitter, every +client of a restarted server reconnects in the same second; with a short cap, a +client that can never connect retries all day. Offline states and authentication +failures wait for a wakeup instead of spending attempts on unchanged conditions. -Foregrounding needs different treatment depending on the connection's state. -It wakes a retry immediately, leaves an ordinary in-flight attempt alone, and -probes an established session before replacing it. A long mobile background -suspension forces replacement because the OS can kill a socket without reporting -closure. Treating every foreground event as a reconnect delays healthy attempts; -treating every resume as harmless leaves suspended sockets stuck. +Foregrounding, an explicit retry, and an offline report probe the established +session, and only a failed probe reconnects. Offline reports are often wrong, for +example for a loopback server. A long mobile background suspension is the one +exception: it replaces the session at once, because the OS can kill a socket +without reporting closure, and a probe would hold a dead socket in "Resuming" +until it times out. That fresh attempt runs even while the network reports +offline. Foregrounding also wakes a pending retry immediately and +leaves an ordinary in-flight attempt alone. The [registry](../../packages/client-runtime/src/connection/registry.ts) scopes connections by environment. An involuntary disconnect retains the registration diff --git a/docs/internals/product-analytics.md b/docs/internals/product-analytics.md index 0e88610ada99..e9123640a541 100644 --- a/docs/internals/product-analytics.md +++ b/docs/internals/product-analytics.md @@ -41,6 +41,13 @@ Unknown counts stay absent. Partial usage contains valid observed counts but cannot establish a whole-turn total. Keep these distinctions when changing token normalization or building reports. +## Delivery + +A send can fail after PostHog has stored the batch, so every retry is a copy. +[Delivery](../../apps/server/src/telemetry/AnalyticsService.ts) gives each event a +uuid when it is recorded, backs off after a failed send, and drops a batch after a +few tries. Without these limits, one stuck batch was sent every second for days. + ## Collection boundary Keep analytics payloads to product metadata and normalized measurements. Do not diff --git a/packages/client-runtime/src/connection/registry.test.ts b/packages/client-runtime/src/connection/registry.test.ts index 998cbe479f2a..7a4d581d3c8e 100644 --- a/packages/client-runtime/src/connection/registry.test.ts +++ b/packages/client-runtime/src/connection/registry.test.ts @@ -1443,6 +1443,46 @@ describe("EnvironmentRegistry", () => { }), ); + it.effect("keeps one session per environment across concurrent registrations and retries", () => + Effect.gen(function* () { + const harness = yield* makeHarness([]); + + yield* Effect.gen(function* () { + const registry = yield* EnvironmentRegistry.EnvironmentRegistry; + const registration = new PrimaryConnectionRegistration({ target: TARGET }); + yield* Effect.all( + Array.from({ length: 5 }, () => registry.registerPlatform(registration)), + { concurrency: "unbounded", discard: true }, + ); + yield* awaitConnectionState( + registry, + TARGET.environmentId, + (state) => state.phase === "connected", + ); + + // Platform polls and explicit retries reach a healthy connection at once. + yield* Effect.all( + [ + ...Array.from({ length: 5 }, () => registry.registerPlatform(registration)), + ...Array.from({ length: 5 }, () => registry.retryNow(TARGET.environmentId)), + registry.reconcilePlatform([registration]), + ], + { concurrency: "unbounded", discard: true }, + ); + for (let attempt = 0; attempt < 100; attempt += 1) { + yield* Effect.yieldNow; + } + + expect(yield* Ref.get(harness.sessions)).toHaveLength(1); + expect(yield* Ref.get(harness.releasedSessions)).toBe(0); + expect(yield* registry.state(TARGET.environmentId)).toMatchObject({ + phase: "connected", + generation: 1, + }); + }).pipe(Effect.provide(harness.layer), Effect.scoped); + }), + ); + it.effect("retains a healthy runtime when the platform repeats an identical registration", () => Effect.gen(function* () { const harness = yield* makeHarness([]); diff --git a/packages/client-runtime/src/connection/supervisor.test.ts b/packages/client-runtime/src/connection/supervisor.test.ts index d64863b9ffcf..1d518f017d6d 100644 --- a/packages/client-runtime/src/connection/supervisor.test.ts +++ b/packages/client-runtime/src/connection/supervisor.test.ts @@ -5,6 +5,7 @@ import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; +import * as Random from "effect/Random"; import * as Ref from "effect/Ref"; import * as Stream from "effect/Stream"; import * as SubscriptionRef from "effect/SubscriptionRef"; @@ -186,6 +187,11 @@ const makeHarness = Effect.fn("TestConnectionHarness.make")(function* (options?: }); const dependencies = Layer.mergeAll( + // Jitter at its maximum, so each retry waits exactly its ceiling: 2s, 4s, 8s... + Layer.succeed(Random.Random, { + nextDoubleUnsafe: () => 1 - Number.EPSILON, + nextIntUnsafe: () => 0, + }), Layer.succeed(Connectivity.Connectivity, connectivity), Layer.succeed( ConnectionWakeups.ConnectionWakeups, @@ -225,6 +231,17 @@ const makeHarness = Effect.fn("TestConnectionHarness.make")(function* (options?: }; }); +describe("retryDelayMs", () => { + it("doubles from 2 seconds to a 5 minute cap, jittered within the upper half of each step", () => { + const ceilings = [2_000, 4_000, 8_000, 16_000, 32_000, 64_000, 128_000, 256_000, 300_000]; + for (const [failureCount, ceiling] of [...ceilings, 300_000].entries()) { + expect(EnvironmentSupervisor.retryDelayMs(failureCount, 0)).toBe(ceiling / 2); + expect(EnvironmentSupervisor.retryDelayMs(failureCount, 0.5)).toBe((ceiling * 3) / 4); + expect(EnvironmentSupervisor.retryDelayMs(failureCount, 1 - Number.EPSILON)).toBe(ceiling); + } + }); +}); + describe("EnvironmentSupervisor", () => { it.effect("exports each relay setup as a standalone linked trace that ends at readiness", () => Effect.gen(function* () { @@ -353,7 +370,7 @@ describe("EnvironmentSupervisor", () => { }), ); - it.effect("retries forever with exponential backoff capped at sixteen seconds", () => + it.effect("retries forever with exponential backoff capped at five minutes", () => Effect.gen(function* () { const harness = yield* makeHarness({ prepare: () => Effect.fail(transient()), @@ -368,15 +385,20 @@ describe("EnvironmentSupervisor", () => { ); expect(yield* Ref.get(harness.prepareCount)).toBe(1); - for (const [index, delay] of [3_000, 4_000, 8_000, 16_000, 16_000, 16_000].entries()) { - yield* TestClock.adjust(delay); + const delays = [ + 2_000, 4_000, 8_000, 16_000, 32_000, 64_000, 128_000, 256_000, 300_000, 300_000, + ]; + for (const [index, delay] of delays.entries()) { + yield* TestClock.adjust(delay - 1); + expect(yield* Ref.get(harness.prepareCount)).toBe(index + 1); + yield* TestClock.adjust(1); yield* eventuallyState( supervisor.state, (state) => state.phase === "backoff" && state.attempt === index + 2, ); } - expect(yield* Ref.get(harness.prepareCount)).toBe(7); + expect(yield* Ref.get(harness.prepareCount)).toBe(delays.length + 1); }).pipe(Effect.provide(TestClock.layer())), ); @@ -578,7 +600,7 @@ describe("EnvironmentSupervisor", () => { ); expect(yield* Ref.get(harness.prepareCount)).toBe(3); - yield* TestClock.adjust("2999 millis"); + yield* TestClock.adjust("1999 millis"); expect(yield* Ref.get(harness.prepareCount)).toBe(3); yield* TestClock.adjust("1 milli"); yield* eventuallyState( @@ -641,9 +663,12 @@ describe("EnvironmentSupervisor", () => { }).pipe(Effect.provide(TestClock.layer())), ); - it.effect("releases a live session while offline and starts a new generation when online", () => + it.effect("keeps a session that still answers when the network reports offline", () => Effect.gen(function* () { - const harness = yield* makeHarness(); + const probeCount = yield* Ref.make(0); + const harness = yield* makeHarness({ + probe: () => Ref.update(probeCount, (count) => count + 1), + }); const supervisor = yield* EnvironmentSupervisor.make(TARGET_ENTRY, { initiallyDesired: true, }).pipe(Effect.provide(harness.dependencies)); @@ -652,21 +677,160 @@ describe("EnvironmentSupervisor", () => { supervisor.state, (state) => state.phase === "connected" && state.generation === 1, ); + // A loopback server, or a flap shorter than the probe, keeps working. yield* harness.setNetworkStatus("offline"); - yield* awaitState(supervisor.state, (state) => state.phase === "offline"); + for (let attempt = 0; attempt < 100; attempt += 1) { + if ((yield* Ref.get(probeCount)) > 0) break; + yield* Effect.yieldNow; + } + yield* harness.setNetworkStatus("online"); + yield* Effect.yieldNow; - expect(yield* Ref.get(harness.releaseCount)).toBe(1); - expect(Option.isNone(yield* SubscriptionRef.get(supervisor.session))).toBe(true); + expect(yield* Ref.get(probeCount)).toBe(1); + expect(yield* Ref.get(harness.sessionCount)).toBe(1); + expect(yield* Ref.get(harness.releaseCount)).toBe(0); + expect(yield* SubscriptionRef.get(supervisor.state)).toMatchObject({ + phase: "connected", + generation: 1, + }); + }), + ); + + it.effect("replaces the session on a long resume while the network reports offline", () => + Effect.gen(function* () { + const probeCount = yield* Ref.make(0); + const harness = yield* makeHarness({ + probe: () => Ref.update(probeCount, (count) => count + 1), + }); + const supervisor = yield* EnvironmentSupervisor.make(TARGET_ENTRY, { + initiallyDesired: true, + }).pipe(Effect.provide(harness.dependencies)); - yield* harness.setNetworkStatus("online"); yield* awaitState( + supervisor.state, + (state) => state.phase === "connected" && state.generation === 1, + ); + // A wrong offline report: the probe answers, so the session stays. + yield* harness.setNetworkStatus("offline"); + for (let attempt = 0; attempt < 100; attempt += 1) { + if ((yield* Ref.get(probeCount)) > 0) break; + yield* Effect.yieldNow; + } + expect(yield* Ref.get(harness.sessionCount)).toBe(1); + + // The replacement connects although the network still reports offline. + yield* harness.wake("application-active-reconnect"); + const replaced = yield* awaitState( supervisor.state, (state) => state.phase === "connected" && state.generation === 2, ); + + expect(replaced.attempt).toBe(1); + expect(yield* Ref.get(probeCount)).toBe(1); expect(yield* Ref.get(harness.sessionCount)).toBe(2); + expect(yield* Ref.get(harness.releaseCount)).toBe(1); + }), + ); + + it.effect( + "releases a session that stops answering while offline and reconnects when online", + () => + Effect.gen(function* () { + const harness = yield* makeHarness({ + probe: (attempt) => (attempt === 1 ? Effect.never : Effect.void), + }); + const supervisor = yield* EnvironmentSupervisor.make(TARGET_ENTRY, { + initiallyDesired: true, + }).pipe(Effect.provide(harness.dependencies)); + + yield* awaitState( + supervisor.state, + (state) => state.phase === "connected" && state.generation === 1, + ); + yield* harness.setNetworkStatus("offline"); + yield* TestClock.adjust("3 seconds"); + yield* awaitState(supervisor.state, (state) => state.phase === "offline"); + + expect(yield* Ref.get(harness.releaseCount)).toBe(1); + expect(Option.isNone(yield* SubscriptionRef.get(supervisor.session))).toBe(true); + + yield* harness.setNetworkStatus("online"); + yield* awaitState( + supervisor.state, + (state) => state.phase === "connected" && state.generation === 2, + ); + expect(yield* Ref.get(harness.sessionCount)).toBe(2); + }).pipe(Effect.provide(TestClock.layer())), + ); + + it.effect("probes instead of replacing a healthy session on an explicit retry", () => + Effect.gen(function* () { + const probeCount = yield* Ref.make(0); + const harness = yield* makeHarness({ + probe: () => Ref.update(probeCount, (count) => count + 1), + }); + const supervisor = yield* EnvironmentSupervisor.make(TARGET_ENTRY, { + initiallyDesired: true, + }).pipe(Effect.provide(harness.dependencies)); + + yield* awaitState(supervisor.state, (state) => state.phase === "connected"); + yield* supervisor.retryNow; + for (let attempt = 0; attempt < 100; attempt += 1) { + if ((yield* Ref.get(probeCount)) > 0) break; + yield* Effect.yieldNow; + } + + expect(yield* Ref.get(probeCount)).toBe(1); + expect(yield* Ref.get(harness.sessionCount)).toBe(1); + expect(yield* Ref.get(harness.releaseCount)).toBe(0); }), ); + it.effect("keeps the backoff ladder after an explicit retry finds a healthy session", () => + Effect.gen(function* () { + const probeCount = yield* Ref.make(0); + const harness = yield* makeHarness({ + probe: () => Ref.update(probeCount, (count) => count + 1), + }); + const supervisor = yield* EnvironmentSupervisor.make(TARGET_ENTRY, { + initiallyDesired: true, + }).pipe(Effect.provide(harness.dependencies)); + + yield* awaitState(supervisor.state, (state) => state.phase === "connected"); + yield* harness.closeLatestSession(); + yield* awaitState( + supervisor.state, + (state) => state.phase === "backoff" && state.attempt === 1, + ); + yield* TestClock.adjust("2 seconds"); + yield* awaitState( + supervisor.state, + (state) => state.phase === "connected" && state.generation === 2, + ); + + yield* supervisor.retryNow; + for (let attempt = 0; attempt < 100; attempt += 1) { + if ((yield* Ref.get(probeCount)) > 0) break; + yield* Effect.yieldNow; + } + expect(yield* Ref.get(probeCount)).toBe(1); + + // The flapping session keeps climbing the ladder: the answered retry does + // not reset it after this unrelated close. + yield* harness.closeLatestSession(); + yield* awaitState( + supervisor.state, + (state) => state.phase === "backoff" && state.attempt === 2, + ); + yield* TestClock.adjust("4 seconds"); + const reconnected = yield* awaitState( + supervisor.state, + (state) => state.phase === "connected" && state.generation === 3, + ); + expect(reconnected.attempt).toBe(3); + }).pipe(Effect.provide(TestClock.layer())), + ); + it.effect("retries a blocked connection when platform credentials change", () => Effect.gen(function* () { const harness = yield* makeHarness({ @@ -844,7 +1008,7 @@ describe("EnvironmentSupervisor", () => { supervisor.state, (state) => state.phase === "backoff" && state.attempt === 1, ); - yield* TestClock.adjust("3 seconds"); + yield* TestClock.adjust("2 seconds"); yield* awaitState( supervisor.state, (state) => state.phase === "connected" && state.generation === 2 && state.attempt === 2, @@ -964,6 +1128,35 @@ describe("EnvironmentSupervisor", () => { }), ); + it.effect("reconnects immediately when the session closes during a resume probe", () => + Effect.gen(function* () { + const probeStarted = yield* Deferred.make(); + const harness = yield* makeHarness({ + probe: (attempt) => + attempt === 1 + ? Deferred.succeed(probeStarted, undefined).pipe(Effect.andThen(Effect.never)) + : Effect.void, + }); + const supervisor = yield* EnvironmentSupervisor.make(TARGET_ENTRY, { + initiallyDesired: true, + }).pipe(Effect.provide(harness.dependencies)); + + yield* awaitState(supervisor.state, (state) => state.phase === "connected"); + yield* harness.wake("application-active-probe"); + yield* Deferred.await(probeStarted); + // The OS reports the suspended socket's close before the probe answers. + yield* harness.closeLatestSession(); + + // No TestClock advance: the unanswered probe skips the first backoff rung. + const reconnected = yield* awaitState( + supervisor.state, + (state) => state.phase === "connected" && state.generation === 2, + ); + expect(reconnected.attempt).toBe(1); + expect(yield* Ref.get(harness.sessionCount)).toBe(2); + }).pipe(Effect.provide(TestClock.layer())), + ); + it.effect("reconnects immediately when the foreground liveness probe fails", () => Effect.gen(function* () { const allowReconnect = yield* Deferred.make(); @@ -1020,7 +1213,7 @@ describe("EnvironmentSupervisor", () => { supervisor.state, (state) => state.phase === "backoff" && state.attempt === 1, ); - yield* TestClock.adjust("2999 millis"); + yield* TestClock.adjust("1999 millis"); expect(yield* Ref.get(harness.prepareCount)).toBe(2); yield* TestClock.adjust("1 milli"); yield* eventuallyState( @@ -1055,6 +1248,32 @@ describe("EnvironmentSupervisor", () => { }).pipe(Effect.provide(TestClock.layer())), ); + it.effect("an explicit retry shortens a stalled desktop foreground probe", () => + Effect.gen(function* () { + const harness = yield* makeHarness({ + probe: (attempt) => (attempt === 1 ? Effect.never : Effect.void), + }); + const supervisor = yield* EnvironmentSupervisor.make(TARGET_ENTRY, { + initiallyDesired: true, + }).pipe(Effect.provide(harness.dependencies)); + + yield* awaitState(supervisor.state, (state) => state.phase === "connected"); + yield* harness.wake("application-active"); + yield* TestClock.adjust("5 seconds"); + yield* supervisor.retryNow; + // The retry's 3 second limit applies, not the 10 seconds left of the 15. + yield* TestClock.adjust("2999 millis"); + expect(yield* Ref.get(harness.sessionCount)).toBe(1); + yield* TestClock.adjust("1 milli"); + yield* awaitState( + supervisor.state, + (state) => state.phase === "connected" && state.generation === 2 && state.attempt === 1, + ); + + expect(yield* Ref.get(harness.sessionCount)).toBe(2); + }).pipe(Effect.provide(TestClock.layer())), + ); + it.effect("quickly times out a stalled mobile foreground liveness probe", () => Effect.gen(function* () { const harness = yield* makeHarness({ diff --git a/packages/client-runtime/src/connection/supervisor.ts b/packages/client-runtime/src/connection/supervisor.ts index b75413bbb7de..a8e1a0831f90 100644 --- a/packages/client-runtime/src/connection/supervisor.ts +++ b/packages/client-runtime/src/connection/supervisor.ts @@ -2,11 +2,13 @@ import { withRelayClientTracing } from "@t3tools/shared/relayTracing"; import * as Cause from "effect/Cause"; import * as Clock from "effect/Clock"; import * as Context from "effect/Context"; +import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; import * as Fiber from "effect/Fiber"; import * as Option from "effect/Option"; import * as Queue from "effect/Queue"; +import * as Random from "effect/Random"; import * as Ref from "effect/Ref"; import * as Scope from "effect/Scope"; import * as Stream from "effect/Stream"; @@ -29,10 +31,13 @@ import { safeErrorLogAttributes } from "../errors/safeLog.ts"; import { NETWORK_BLOCKING_HINT } from "../errors/network.ts"; import * as ConnectionWakeups from "./wakeups.ts"; -const RETRY_DELAYS_MS = [3_000, 4_000, 8_000, 16_000] as const; +const RETRY_BASE_DELAY_MS = 1_000; +const RETRY_MAX_DELAY_MS = 300_000; const CONNECTION_ESTABLISHMENT_TIMEOUT = "15 seconds"; const CONNECTION_PROBE_TIMEOUT = "15 seconds"; -const MOBILE_CONNECTION_PROBE_TIMEOUT = "3 seconds"; +// Mobile resumes, explicit retries, and offline events want a fast answer: +// the user is waiting, or the network may be gone. +const QUICK_CONNECTION_PROBE_TIMEOUT = "3 seconds"; const BACKOFF_RESET_AFTER_MS = 30_000; interface SupervisorIntent { @@ -101,8 +106,19 @@ export interface EnvironmentSupervisorOptions { readonly initiallyDesired?: boolean; } -function retryDelayMs(failureCount: number): number { - return RETRY_DELAYS_MS[Math.min(failureCount, RETRY_DELAYS_MS.length - 1)] ?? 16_000; +/** + * Delay before the next attempt after `failureCount` consecutive failures + * (0 for the first retry). The ceiling doubles from 2s up to 5 minutes, and + * the delay is a random point in its upper half: never quicker than half the + * ceiling, and spread out so clients that lost the same server do not all + * reconnect in the same second. `random` is in [0, 1). + * + * The long cap only applies to a connection that keeps failing. Returning to + * the app, the network coming back, and an explicit retry all skip the wait. + */ +export function retryDelayMs(failureCount: number, random: number): number { + const ceiling = Math.min(RETRY_MAX_DELAY_MS, RETRY_BASE_DELAY_MS * 2 ** (failureCount + 1)); + return Math.round(ceiling / 2 + (ceiling / 2) * random); } function annotateTarget(target: ConnectionTarget) { @@ -235,10 +251,11 @@ export const make = Effect.fn("EnvironmentSupervisor.make")(function* ( const intent = yield* Ref.make(initialIntent); const signals = yield* Queue.unbounded(); const resetRetryState = yield* Ref.make(false); - // Set when a foreground wake probe fails or times out: the user is actively - // returning to the app on a dead transport, so the follow-up reconnect skips - // the first backoff rung instead of sleeping. - const wakeProbeFailed = yield* Ref.make(false); + // Set while a probe of the live session is running, and kept when it fails + // or times out: something asked whether the connection still works and it + // closed or failed before answering, so the follow-up reconnect skips the + // first backoff rung instead of sleeping. + const probeUnanswered = yield* Ref.make(false); const state = yield* SubscriptionRef.make( !initialIntent.desired ? availableState(initialIntent, 0) @@ -392,96 +409,118 @@ export const make = Effect.fn("EnvironmentSupervisor.make")(function* ( } }); + // Signals that end a connected lease whatever its health: "reset" ends it + // and restarts the retry ladder, "end" ends it, undefined keeps it. + const connectedLeaseEnd = Effect.fnUntraced(function* (next: SupervisorSignal) { + if (next._tag === "DisconnectRequested") { + return "end" as const; + } + if (next._tag !== "Wakeup") { + return undefined; + } + if (next.reason === "application-active-reconnect") { + // Mobile operating systems often kill a suspended socket without a close + // event. A probe would show a dead socket as "Resuming" until it times + // out, so a long background resume replaces the session at once. + return "reset" as const; + } + if (next.reason === "credentials-changed" && target._tag === "RelayConnectionTarget") { + yield* logManagedRelayAccountChange; + return "end" as const; + } + return undefined; + }); + + // How long a signal waits for the live session to answer a probe, or + // undefined when the signal does not question the connection. + const probeTimeoutFor = (next: SupervisorSignal): Duration.Input | undefined => { + switch (next._tag) { + case "RetryRequested": + return QUICK_CONNECTION_PROBE_TIMEOUT; + case "NetworkChanged": + return next.network === "offline" ? QUICK_CONNECTION_PROBE_TIMEOUT : undefined; + case "Wakeup": + if (next.reason === "application-active") { + return CONNECTION_PROBE_TIMEOUT; + } + return next.reason === "application-active-probe" + ? QUICK_CONNECTION_PROBE_TIMEOUT + : undefined; + case "ConnectRequested": + case "DisconnectRequested": + return undefined; + } + }; + + // Holds a connected lease until it must end, and returns whether to restart + // the retry ladder. Returning to the app, an explicit retry, and the network + // reporting offline all probe the live session instead of replacing it, so a + // healthy socket is not torn down (the offline report is often wrong, for + // example for a loopback server). Only a long mobile resume replaces the + // session without a probe. A failed probe fails this effect, and the + // supervisor reconnects. const monitorConnectedLease = Effect.fnUntraced(function* ( lease: ConnectionDriver.EnvironmentConnectionLease, ) { + // A probe answers an explicit retry here, so the retry must not also reset + // the backoff of a later, unrelated failure. + const takeSignal = Queue.take(signals).pipe( + Effect.tap((next) => + next._tag === "RetryRequested" ? Ref.set(resetRetryState, false) : Effect.void, + ), + ); for (;;) { - const next = yield* Queue.take(signals); - switch (next._tag) { - case "DisconnectRequested": - case "RetryRequested": - return false; - case "NetworkChanged": - if (next.network === "offline") { - return false; - } - break; - case "Wakeup": - if (next.reason === "credentials-changed" && target._tag === "RelayConnectionTarget") { - yield* logManagedRelayAccountChange; - return false; - } - if (next.reason === "application-active-reconnect") { - // Mobile operating systems commonly suspend sockets without - // delivering a close event. A long background resume deliberately - // replaces that lease and starts a fresh attempt without backoff. - return true; - } - if (next.reason === "application-active" || next.reason === "application-active-probe") { - const probe = yield* lease.session.probe.pipe( - Effect.timeoutOrElse({ - duration: - next.reason === "application-active-probe" - ? MOBILE_CONNECTION_PROBE_TIMEOUT - : CONNECTION_PROBE_TIMEOUT, - orElse: () => - Effect.fail( - new ConnectionTransientError({ - reason: "timeout", - detail: `${target.label} did not respond to a connection health check.`, - }), - ), - }), - Effect.forkChild, - ); - for (;;) { - const probeEvent = yield* Effect.raceFirst( - Fiber.await(probe).pipe( - Effect.map((exit) => ({ _tag: "ProbeCompleted" as const, exit })), - ), - Queue.take(signals).pipe( - Effect.map((signal) => ({ _tag: "Signal" as const, signal })), - ), - ); - if (probeEvent._tag === "ProbeCompleted") { - if (Exit.isFailure(probeEvent.exit)) { - yield* Ref.set(wakeProbeFailed, true); - } - yield* probeEvent.exit; - break; - } - switch (probeEvent.signal._tag) { - case "DisconnectRequested": - case "RetryRequested": - yield* Fiber.interrupt(probe); - return false; - case "NetworkChanged": - if (probeEvent.signal.network === "offline") { - yield* Fiber.interrupt(probe); - return false; - } - break; - case "Wakeup": - if (probeEvent.signal.reason === "application-active-reconnect") { - yield* Fiber.interrupt(probe); - return true; - } - if ( - probeEvent.signal.reason === "credentials-changed" && - target._tag === "RelayConnectionTarget" - ) { - yield* Fiber.interrupt(probe); - return false; - } - break; - case "ConnectRequested": - break; - } - } + const next = yield* takeSignal; + const end = yield* connectedLeaseEnd(next); + if (end !== undefined) { + return end === "reset"; + } + const probeTimeout = probeTimeoutFor(next); + if (probeTimeout === undefined) { + continue; + } + yield* Ref.set(probeUnanswered, true); + const probe = yield* Effect.forkChild(lease.session.probe); + // Monotonic nanoseconds, so a wall-clock correction cannot move the deadline. + let deadline = (yield* Clock.monotonicTimeNanos) + Duration.toNanosUnsafe(probeTimeout); + for (;;) { + const remaining = deadline - (yield* Clock.monotonicTimeNanos); + const probeEvent = yield* Effect.raceAllFirst([ + Fiber.await(probe).pipe( + Effect.map((exit) => ({ _tag: "ProbeCompleted" as const, exit })), + ), + takeSignal.pipe(Effect.map((signal) => ({ _tag: "Signal" as const, signal }))), + Effect.sleep(Duration.nanos(remaining > 0n ? remaining : 0n)).pipe( + Effect.as({ _tag: "TimedOut" as const }), + ), + ]); + if (probeEvent._tag === "TimedOut") { + yield* Fiber.interrupt(probe); + return yield* new ConnectionTransientError({ + reason: "timeout", + detail: `${target.label} did not respond to a connection health check.`, + }); + } + if (probeEvent._tag === "ProbeCompleted") { + if (Exit.isSuccess(probeEvent.exit)) { + yield* Ref.set(probeUnanswered, false); } + yield* probeEvent.exit; break; - case "ConnectRequested": - break; + } + const endDuringProbe = yield* connectedLeaseEnd(probeEvent.signal); + if (endDuringProbe !== undefined) { + yield* Fiber.interrupt(probe); + return endDuringProbe === "reset"; + } + // A retry or an offline report during a desktop foreground probe wants + // its quicker answer, so it shortens the running probe. + const signalTimeout = probeTimeoutFor(probeEvent.signal); + if (signalTimeout !== undefined) { + const signalDeadline = + (yield* Clock.monotonicTimeNanos) + Duration.toNanosUnsafe(signalTimeout); + if (signalDeadline < deadline) deadline = signalDeadline; + } } } }); @@ -491,6 +530,7 @@ export const make = Effect.fn("EnvironmentSupervisor.make")(function* ( generation: number, lastFailure: ConnectionAttemptError | null, pendingRetry: Option.Option, + ignoreOffline: boolean, ) { yield* SubscriptionRef.set(prepared, Option.none()); const establishment = yield* Effect.raceAllFirst([ @@ -556,7 +596,7 @@ export const make = Effect.fn("EnvironmentSupervisor.make")(function* ( const active = establishment.exit.value; const currentIntent = yield* Ref.get(intent); - if (!currentIntent.desired || currentIntent.network === "offline") { + if (!currentIntent.desired || (currentIntent.network === "offline" && !ignoreOffline)) { return { _tag: "Interrupted", established: false, @@ -641,6 +681,10 @@ export const make = Effect.fn("EnvironmentSupervisor.make")(function* ( failureCount = 0; pendingRetry = Option.none(); }; + // Set after a long resume ends an attempt or a session. The fresh attempt + // runs even while the network reports offline: the report is often wrong, + // and the replaced session must not leave the client offline. + let replacing = false; for (;;) { if (yield* Ref.getAndSet(resetRetryState, false)) { @@ -657,7 +701,7 @@ export const make = Effect.fn("EnvironmentSupervisor.make")(function* ( yield* waitForSignal; continue; } - if (currentIntent.network === "offline") { + if (currentIntent.network === "offline" && !replacing) { yield* clearLease; yield* setState(offlineState(currentIntent, generation, failureCount + 1, latestFailure)); const applicationActivated = yield* waitForSignal; @@ -670,11 +714,12 @@ export const make = Effect.fn("EnvironmentSupervisor.make")(function* ( const attempt = failureCount + 1; const nextGeneration = generation + 1; const outcome: AttemptOutcome = yield* Effect.scoped( - runAttempt(attempt, nextGeneration, latestFailure, pendingRetry), + runAttempt(attempt, nextGeneration, latestFailure, pendingRetry, replacing), ); + replacing = false; // Consumed on every iteration so a stale marker can never leak into a // later, unrelated failure. - const failedWakeProbe = yield* Ref.getAndSet(wakeProbeFailed, false); + const failedProbe = yield* Ref.getAndSet(probeUnanswered, false); if (outcome.established) { generation = nextGeneration; if (outcome.stable) { @@ -685,6 +730,7 @@ export const make = Effect.fn("EnvironmentSupervisor.make")(function* ( if (outcome._tag === "Interrupted") { if (outcome.resetRetry) { resetRetryLadder(); + replacing = true; } continue; } @@ -711,18 +757,19 @@ export const make = Effect.fn("EnvironmentSupervisor.make")(function* ( continue; } - if (failedWakeProbe) { - // The wake probe found a dead transport while the user is returning to - // the app, so reconnect immediately instead of sleeping the first - // backoff rung. Only this first attempt skips the ladder; if it fails - // too, normal backoff resumes. + if (failedProbe) { + // A probe found a dead transport, or the transport closed while a probe + // waited for an answer (the user returned to the app, asked to retry, + // or the network changed), so reconnect immediately instead + // of sleeping the first backoff rung. Only this first attempt skips the + // ladder; if it fails too, normal backoff resumes. resetRetryLadder(); yield* setState(connectingState(yield* Ref.get(intent), generation, 1, error)); continue; } failureCount += 1; - const delayMs = retryDelayMs(failureCount - 1); + const delayMs = retryDelayMs(failureCount - 1, yield* Random.next); pendingRetry = Option.map(attemptSpan, (previousAttempt) => ({ previousAttempt, failureCount, diff --git a/packages/client-runtime/src/connection/wakeups.ts b/packages/client-runtime/src/connection/wakeups.ts index 721b941645c0..5f19924c47a1 100644 --- a/packages/client-runtime/src/connection/wakeups.ts +++ b/packages/client-runtime/src/connection/wakeups.ts @@ -16,6 +16,7 @@ export function isApplicationActiveWakeup(reason: ConnectionWakeup): boolean { ); } +// A long resume replaces the session, and the new session subscribes on its own. export function shouldResubscribeAfterWakeup(reason: ConnectionWakeup): boolean { return reason === "application-active" || reason === "application-active-probe"; } diff --git a/packages/client-runtime/src/state/server.ts b/packages/client-runtime/src/state/server.ts index f17fa59adcc9..0a0e3ca377aa 100644 --- a/packages/client-runtime/src/state/server.ts +++ b/packages/client-runtime/src/state/server.ts @@ -179,9 +179,9 @@ export function validateServerUpdateReadyEvent( * Keeps reconnect attempts ~1s apart for the whole update restart. * * A restart takes the server down for ~15 seconds, but the supervisor's normal - * backoff ladder (1/2/4/8/16s) assumes an unexpected failure and lands attempts - * at ~3, 5, 9, 17 and 33 seconds — so a 15-second restart is observed as a - * 33-second "Resuming". Nudging on every backoff entry (not just the first) + * backoff assumes an unexpected failure and doubles its delay after each failed + * attempt, so a 15-second restart can be observed as a ~30-second "Resuming". + * Nudging on every backoff entry (not just the first) * holds the retry cadence flat until the server answers again. The sleep before * each nudge is the pacer: a connection that fails instantly re-enters backoff * immediately and would otherwise spin a tight retry loop.