mirror of
https://github.com/VibedByKaKi/t3-code-android-nightly.git
synced 2026-10-09 03:41:17 +02:00
1039 lines
36 KiB
TypeScript
1039 lines
36 KiB
TypeScript
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";
|
|
import * as SubscriptionRef from "effect/SubscriptionRef";
|
|
import * as Tracer from "effect/Tracer";
|
|
|
|
import type { ConnectionCatalogEntry, ConnectionRoute } from "./catalog.ts";
|
|
import * as Connectivity from "./connectivity.ts";
|
|
import * as ConnectionDriver from "./driver.ts";
|
|
import {
|
|
type ConnectionAttemptError,
|
|
type ConnectionTarget,
|
|
ConnectionTransientError,
|
|
type NetworkStatus,
|
|
type PreparedConnection,
|
|
type SupervisorConnectionState,
|
|
} from "./model.ts";
|
|
import * as RpcSession from "../rpc/session.ts";
|
|
import { safeErrorLogAttributes } from "../errors/safeLog.ts";
|
|
import { NETWORK_BLOCKING_HINT } from "../errors/network.ts";
|
|
import * as ConnectionWakeups from "./wakeups.ts";
|
|
import { connectionRouteId, connectionRoutes, entryWithRoutes } from "./routes.ts";
|
|
|
|
const RETRY_BASE_DELAY_MS = 1_000;
|
|
const RETRY_MAX_DELAY_MS = 300_000;
|
|
const CONNECTION_ESTABLISHMENT_TIMEOUT = "15 seconds";
|
|
const establishmentTimeout = Duration.fromInputUnsafe(CONNECTION_ESTABLISHMENT_TIMEOUT);
|
|
const CONNECTION_PROBE_TIMEOUT = "15 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;
|
|
// While connected over a fallback route, how often to look for a better one.
|
|
// Network changes and returning to the app also trigger a check.
|
|
const BETTER_ROUTE_CHECK_INTERVAL = "60 seconds";
|
|
// A better route that answered the check but then failed to connect is not
|
|
// tried again for this long, so a flaky LAN cannot bounce the connection.
|
|
const BETTER_ROUTE_COOLDOWN_MS = 300_000;
|
|
|
|
interface SupervisorIntent {
|
|
readonly desired: boolean;
|
|
readonly network: NetworkStatus;
|
|
}
|
|
|
|
type SupervisorSignal =
|
|
| { readonly _tag: "ConnectRequested" }
|
|
| { readonly _tag: "DisconnectRequested" }
|
|
| { readonly _tag: "RetryRequested" }
|
|
| { readonly _tag: "NetworkChanged"; readonly network: NetworkStatus }
|
|
| { readonly _tag: "Wakeup"; readonly reason: ConnectionWakeups.ConnectionWakeup }
|
|
| { readonly _tag: "BetterRouteAvailable"; readonly routeId: string };
|
|
|
|
interface PendingRetryTrace {
|
|
readonly previousAttempt: Tracer.Span;
|
|
readonly failureCount: number;
|
|
readonly delayMs: number;
|
|
readonly reason: ConnectionAttemptError["reason"];
|
|
}
|
|
|
|
interface TracedAttemptFailure {
|
|
readonly error: ConnectionAttemptError;
|
|
readonly attemptSpan: Option.Option<Tracer.Span>;
|
|
}
|
|
|
|
type AttemptOutcome =
|
|
| {
|
|
readonly _tag: "Interrupted";
|
|
readonly established: boolean;
|
|
readonly stable: boolean;
|
|
readonly resetRetry: boolean;
|
|
}
|
|
| {
|
|
readonly _tag: "Failure";
|
|
readonly established: boolean;
|
|
readonly stable: boolean;
|
|
readonly failure: TracedAttemptFailure;
|
|
};
|
|
|
|
type EstablishmentEvent =
|
|
| {
|
|
readonly _tag: "Completed";
|
|
readonly exit: Exit.Exit<
|
|
{
|
|
readonly attemptSpan: Option.Option<Tracer.Span>;
|
|
readonly lease: ConnectionDriver.EnvironmentConnectionLease;
|
|
},
|
|
TracedAttemptFailure
|
|
>;
|
|
}
|
|
| { readonly _tag: "Interrupted"; readonly resetRetry: boolean }
|
|
| { readonly _tag: "TimedOut" };
|
|
|
|
function exitUnlessInterrupted<A, E, R>(
|
|
effect: Effect.Effect<A, E, R>,
|
|
): Effect.Effect<Exit.Exit<A, E>, never, R> {
|
|
return Effect.matchCauseEffect(effect, {
|
|
onFailure: (cause) =>
|
|
Cause.hasInterrupts(cause) ? Effect.interrupt : Effect.succeed(Exit.failCause(cause)),
|
|
onSuccess: (value) => Effect.succeed(Exit.succeed(value)),
|
|
});
|
|
}
|
|
|
|
export interface EnvironmentSupervisorOptions {
|
|
readonly initiallyDesired?: boolean;
|
|
/**
|
|
* Saves the direct addresses the server reports once a session is ready and
|
|
* returns the entry with its updated routes, which later attempts use. The
|
|
* live session is left alone: it already works.
|
|
*/
|
|
readonly learnRoutes?: (input: {
|
|
readonly activeRoute: ConnectionRoute;
|
|
readonly reported: ReadonlyArray<{ readonly httpBaseUrl: string }>;
|
|
}) => Effect.Effect<Option.Option<ConnectionCatalogEntry>>;
|
|
}
|
|
|
|
/**
|
|
* 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) {
|
|
return Effect.annotateCurrentSpan({
|
|
"environment.id": target.environmentId,
|
|
"environment.label": target.label,
|
|
"environment.target.kind": target._tag,
|
|
});
|
|
}
|
|
|
|
function availableState(intent: SupervisorIntent, generation: number): SupervisorConnectionState {
|
|
return {
|
|
desired: false,
|
|
network: intent.network,
|
|
phase: "available",
|
|
stage: null,
|
|
attempt: 0,
|
|
generation,
|
|
lastFailure: null,
|
|
retryAt: null,
|
|
};
|
|
}
|
|
|
|
function offlineState(
|
|
intent: SupervisorIntent,
|
|
generation: number,
|
|
attempt: number,
|
|
lastFailure: ConnectionAttemptError | null,
|
|
): SupervisorConnectionState {
|
|
return {
|
|
desired: true,
|
|
network: intent.network,
|
|
phase: "offline",
|
|
stage: null,
|
|
attempt,
|
|
generation,
|
|
lastFailure,
|
|
retryAt: null,
|
|
};
|
|
}
|
|
|
|
function connectingState(
|
|
intent: SupervisorIntent,
|
|
generation: number,
|
|
attempt: number,
|
|
lastFailure: ConnectionAttemptError | null,
|
|
stage: SupervisorConnectionState["stage"] = "preparing",
|
|
): SupervisorConnectionState {
|
|
return {
|
|
desired: true,
|
|
network: intent.network,
|
|
phase: "connecting",
|
|
stage,
|
|
attempt,
|
|
generation,
|
|
lastFailure,
|
|
retryAt: null,
|
|
};
|
|
}
|
|
|
|
function failureFromExit<A>(
|
|
target: ConnectionTarget,
|
|
exit: Exit.Exit<A, TracedAttemptFailure>,
|
|
established: boolean,
|
|
stable: boolean,
|
|
): AttemptOutcome {
|
|
if (Exit.isSuccess(exit) || Cause.hasInterruptsOnly(exit.cause)) {
|
|
return { _tag: "Interrupted", established, stable, resetRetry: false };
|
|
}
|
|
const typedFailure = exit.cause.reasons.find(Cause.isFailReason);
|
|
if (typedFailure) {
|
|
return {
|
|
_tag: "Failure",
|
|
established,
|
|
stable,
|
|
failure: typedFailure.error,
|
|
};
|
|
}
|
|
return {
|
|
_tag: "Failure",
|
|
established,
|
|
stable,
|
|
failure: {
|
|
error: new ConnectionTransientError({
|
|
reason: "transport",
|
|
detail: `${target.label} connection failed unexpectedly.`,
|
|
}),
|
|
attemptSpan: Option.none(),
|
|
},
|
|
};
|
|
}
|
|
|
|
export class EnvironmentSupervisor extends Context.Service<
|
|
EnvironmentSupervisor,
|
|
{
|
|
readonly target: ConnectionTarget;
|
|
readonly state: SubscriptionRef.SubscriptionRef<SupervisorConnectionState>;
|
|
readonly session: SubscriptionRef.SubscriptionRef<Option.Option<RpcSession.RpcSession>>;
|
|
readonly prepared: SubscriptionRef.SubscriptionRef<Option.Option<PreparedConnection>>;
|
|
readonly connect: Effect.Effect<void>;
|
|
readonly disconnect: Effect.Effect<void>;
|
|
readonly retryNow: Effect.Effect<void>;
|
|
}
|
|
>()("@t3tools/client-runtime/connection/supervisor/EnvironmentSupervisor") {}
|
|
|
|
export const make = Effect.fn("EnvironmentSupervisor.make")(function* (
|
|
entry: ConnectionCatalogEntry,
|
|
options?: EnvironmentSupervisorOptions,
|
|
): Effect.fn.Return<
|
|
EnvironmentSupervisor["Service"],
|
|
never,
|
|
| Connectivity.Connectivity
|
|
| ConnectionDriver.ConnectionDriver
|
|
| Scope.Scope
|
|
| ConnectionWakeups.ConnectionWakeups
|
|
> {
|
|
const target = entry.target;
|
|
// Relay-specific handling applies when any route is T3 Connect, since the
|
|
// attempt or the live session may be using it.
|
|
const usesRelay = connectionRoutes(entry).some(
|
|
(route) => route.target._tag === "RelayConnectionTarget",
|
|
);
|
|
const setupTimeoutDetail = `${target.label} did not respond during connection setup.${
|
|
usesRelay ? ` ${NETWORK_BLOCKING_HINT}` : ""
|
|
}`;
|
|
yield* annotateTarget(target);
|
|
|
|
const connectivity = yield* Connectivity.Connectivity;
|
|
const driver = yield* ConnectionDriver.ConnectionDriver;
|
|
const wakeups = yield* ConnectionWakeups.ConnectionWakeups;
|
|
const initialIntent: SupervisorIntent = {
|
|
desired: options?.initiallyDesired ?? false,
|
|
network: yield* connectivity.status,
|
|
};
|
|
const intent = yield* Ref.make(initialIntent);
|
|
const signals = yield* Queue.unbounded<SupervisorSignal>();
|
|
const resetRetryState = 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<SupervisorConnectionState>(
|
|
!initialIntent.desired
|
|
? availableState(initialIntent, 0)
|
|
: initialIntent.network === "offline"
|
|
? offlineState(initialIntent, 0, 0, null)
|
|
: connectingState(initialIntent, 0, 1, null),
|
|
);
|
|
const session = yield* SubscriptionRef.make<Option.Option<RpcSession.RpcSession>>(Option.none());
|
|
const prepared = yield* SubscriptionRef.make<Option.Option<PreparedConnection>>(Option.none());
|
|
// Learned routes arrive while a session is live, so the route list is read
|
|
// fresh by each attempt and check instead of being fixed at creation.
|
|
const currentEntry = yield* Ref.make(entry);
|
|
const currentRoutes = Ref.get(currentEntry).pipe(Effect.map(connectionRoutes));
|
|
// Set when a better route answered while connected over a worse one; the
|
|
// replacement attempt tries it first. Cleared once an attempt starts.
|
|
const preferredRouteId = yield* Ref.make(Option.none<string>());
|
|
// Monotonic deadline per route that answered a check but failed to connect.
|
|
const routeCooldowns = yield* Ref.make<ReadonlyMap<string, bigint>>(new Map());
|
|
|
|
const attemptEntry = Effect.gen(function* () {
|
|
const latest = yield* Ref.get(currentEntry);
|
|
const routes = connectionRoutes(latest);
|
|
const preferred = yield* Ref.getAndSet(preferredRouteId, Option.none());
|
|
if (Option.isNone(preferred)) return latest;
|
|
const route = routes.find(
|
|
(candidate) => connectionRouteId(candidate.target) === preferred.value,
|
|
);
|
|
return route === undefined
|
|
? latest
|
|
: entryWithRoutes(latest, [route, ...routes.filter((candidate) => candidate !== route)]);
|
|
});
|
|
|
|
const learnRoutesFrom = Effect.fnUntraced(function* (
|
|
learn: NonNullable<EnvironmentSupervisorOptions["learnRoutes"]>,
|
|
lease: ConnectionDriver.EnvironmentConnectionLease,
|
|
) {
|
|
const config = yield* lease.session.initialConfig.pipe(Effect.option);
|
|
const reported = Option.getOrUndefined(config)?.directEndpoints;
|
|
if (reported === undefined) return;
|
|
const activeId = connectionRouteId(lease.prepared.target);
|
|
const activeRoute = (yield* currentRoutes).find(
|
|
(route) => connectionRouteId(route.target) === activeId,
|
|
);
|
|
if (activeRoute === undefined) return;
|
|
const updated = yield* learn({ activeRoute, reported });
|
|
if (Option.isSome(updated)) {
|
|
yield* Ref.set(currentEntry, updated.value);
|
|
// A learned route may rank above the one in use.
|
|
yield* requestBetterRouteCheck(lease);
|
|
}
|
|
});
|
|
|
|
const routeIndex = (lease: ConnectionDriver.EnvironmentConnectionLease) =>
|
|
currentRoutes.pipe(
|
|
Effect.map((routes) =>
|
|
routes.findIndex(
|
|
(route) => connectionRouteId(route.target) === connectionRouteId(lease.prepared.target),
|
|
),
|
|
),
|
|
);
|
|
|
|
/**
|
|
* Preflights the routes ranked above the one in use and signals the best
|
|
* that would connect. Routes without a cheap check (T3 Connect, SSH) never
|
|
* pass, so they are fallbacks, not destinations.
|
|
*/
|
|
const checkBetterRoutes = Effect.fnUntraced(function* (
|
|
lease: ConnectionDriver.EnvironmentConnectionLease,
|
|
) {
|
|
const latest = yield* Ref.get(currentEntry);
|
|
const current = yield* routeIndex(lease);
|
|
if (current <= 0) return;
|
|
const now = yield* Clock.monotonicTimeNanos;
|
|
const cooldowns = yield* Ref.get(routeCooldowns);
|
|
const better = connectionRoutes(latest)
|
|
.slice(0, current)
|
|
.filter((route) => (cooldowns.get(connectionRouteId(route.target)) ?? 0n) <= now);
|
|
const passed = yield* Effect.forEach(
|
|
better,
|
|
(route: ConnectionRoute) => driver.preflight(latest, route),
|
|
{ concurrency: "unbounded" },
|
|
);
|
|
const index = passed.indexOf(true);
|
|
if (index === -1) return;
|
|
// The check may outlive the session it was asked about.
|
|
const live = yield* SubscriptionRef.get(session);
|
|
if (Option.isNone(live) || live.value !== lease.session) return;
|
|
yield* signal({
|
|
_tag: "BetterRouteAvailable",
|
|
routeId: connectionRouteId(better[index]!.target),
|
|
});
|
|
});
|
|
|
|
// One check at a time; a trigger during a check is dropped, not queued.
|
|
const betterRouteChecks = yield* Queue.dropping<ConnectionDriver.EnvironmentConnectionLease>(1);
|
|
const requestBetterRouteCheck = (lease: ConnectionDriver.EnvironmentConnectionLease) =>
|
|
routeIndex(lease).pipe(
|
|
Effect.flatMap((index) =>
|
|
index > 0 ? Queue.offer(betterRouteChecks, lease).pipe(Effect.asVoid) : Effect.void,
|
|
),
|
|
);
|
|
|
|
const clearLease = Effect.all(
|
|
[SubscriptionRef.set(session, Option.none()), SubscriptionRef.set(prepared, Option.none())],
|
|
{ discard: true },
|
|
);
|
|
|
|
const setState = Effect.fn("EnvironmentSupervisor.setState")(function* (
|
|
next: SupervisorConnectionState,
|
|
) {
|
|
yield* SubscriptionRef.set(state, next);
|
|
});
|
|
|
|
const signal = Effect.fn("EnvironmentSupervisor.signal")(function* (next: SupervisorSignal) {
|
|
yield* Queue.offer(signals, next);
|
|
});
|
|
|
|
const logManagedRelayAccountChange = Effect.logInfo(
|
|
"Managed relay account changed; restarting the environment connection.",
|
|
).pipe(
|
|
Effect.annotateLogs({
|
|
"environment.id": target.environmentId,
|
|
"environment.label": target.label,
|
|
}),
|
|
);
|
|
|
|
const reportProgress = Effect.fn("EnvironmentSupervisor.reportProgress")(function* (
|
|
attempt: number,
|
|
generation: number,
|
|
lastFailure: ConnectionAttemptError | null,
|
|
progress: ConnectionDriver.ConnectionDriverProgress,
|
|
) {
|
|
if ("prepared" in progress) {
|
|
yield* SubscriptionRef.set(prepared, Option.some(progress.prepared));
|
|
}
|
|
yield* setState(
|
|
connectingState(yield* Ref.get(intent), generation, attempt, lastFailure, progress.stage),
|
|
);
|
|
});
|
|
|
|
const establishConnection = Effect.fnUntraced(function* (
|
|
attempt: number,
|
|
generation: number,
|
|
lastFailure: ConnectionAttemptError | null,
|
|
) {
|
|
return yield* driver.connect(yield* attemptEntry, (progress) =>
|
|
reportProgress(attempt, generation, lastFailure, progress),
|
|
);
|
|
});
|
|
|
|
const traceRelayEstablishment = (
|
|
effect: Effect.Effect<
|
|
ConnectionDriver.EnvironmentConnectionLease,
|
|
ConnectionAttemptError,
|
|
Scope.Scope
|
|
>,
|
|
attempt: number,
|
|
generation: number,
|
|
pendingRetry: Option.Option<PendingRetryTrace>,
|
|
) => {
|
|
const traced = Effect.gen(function* () {
|
|
const attemptSpan = yield* Effect.currentSpan.pipe(Effect.orDie);
|
|
yield* annotateTarget(target);
|
|
yield* Effect.annotateCurrentSpan({
|
|
"connection.attempt": attempt,
|
|
"connection.generation": generation,
|
|
"connection.retry.failure_count": Option.match(pendingRetry, {
|
|
onNone: () => 0,
|
|
onSome: (retry) => retry.failureCount,
|
|
}),
|
|
});
|
|
const lease = yield* effect.pipe(
|
|
Effect.mapError((error): TracedAttemptFailure => ({
|
|
error,
|
|
attemptSpan: Option.some(attemptSpan),
|
|
})),
|
|
);
|
|
return { attemptSpan: Option.some(attemptSpan), lease };
|
|
}).pipe(Effect.withSpan("relay.connection.attempt", { root: true }));
|
|
|
|
return Option.match(pendingRetry, {
|
|
onNone: () => traced,
|
|
onSome: (retry) =>
|
|
traced.pipe(
|
|
Effect.linkSpans(retry.previousAttempt, {
|
|
"connection.retry.delay_ms": retry.delayMs,
|
|
"connection.retry.reason": retry.reason,
|
|
}),
|
|
),
|
|
}).pipe(withRelayClientTracing);
|
|
};
|
|
|
|
const establishTracedConnection = Effect.fnUntraced(function* (
|
|
attempt: number,
|
|
generation: number,
|
|
lastFailure: ConnectionAttemptError | null,
|
|
pendingRetry: Option.Option<PendingRetryTrace>,
|
|
) {
|
|
if (usesRelay) {
|
|
return yield* traceRelayEstablishment(
|
|
establishConnection(attempt, generation, lastFailure),
|
|
attempt,
|
|
generation,
|
|
pendingRetry,
|
|
);
|
|
}
|
|
return yield* establishConnection(attempt, generation, lastFailure).pipe(
|
|
Effect.map((lease) => ({
|
|
attemptSpan: Option.none<Tracer.Span>(),
|
|
lease,
|
|
})),
|
|
Effect.mapError((error): TracedAttemptFailure => ({
|
|
error,
|
|
attemptSpan: Option.none(),
|
|
})),
|
|
);
|
|
});
|
|
|
|
const waitForEstablishmentInterrupt = Effect.fnUntraced(function* () {
|
|
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 "ConnectRequested":
|
|
case "BetterRouteAvailable":
|
|
break;
|
|
case "Wakeup":
|
|
if (next.reason === "application-active-reconnect") {
|
|
return true;
|
|
}
|
|
if (next.reason === "credentials-changed" && usesRelay) {
|
|
yield* logManagedRelayAccountChange;
|
|
return false;
|
|
}
|
|
break;
|
|
}
|
|
}
|
|
});
|
|
|
|
// Signals that end a connected lease whatever its health: "reset" ends it
|
|
// and restarts the retry ladder, "end" ends it, undefined keeps it.
|
|
const isRelayLease = (lease: ConnectionDriver.EnvironmentConnectionLease) =>
|
|
lease.prepared.target._tag === "RelayConnectionTarget";
|
|
|
|
const connectedLeaseEnd = Effect.fnUntraced(function* (
|
|
next: SupervisorSignal,
|
|
lease: ConnectionDriver.EnvironmentConnectionLease,
|
|
) {
|
|
if (next._tag === "DisconnectRequested") {
|
|
return "end" as const;
|
|
}
|
|
if (next._tag === "BetterRouteAvailable") {
|
|
// Replaced like a long resume: the new attempt prefers the better route
|
|
// and its session takes over the durable subscriptions.
|
|
yield* Ref.set(preferredRouteId, Option.some(next.routeId));
|
|
return "reset" 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;
|
|
}
|
|
// Only a session over T3 Connect holds the old account's credential.
|
|
if (next.reason === "credentials-changed" && isRelayLease(lease)) {
|
|
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;
|
|
}
|
|
// A socket opened on the previous network may now be unroutable.
|
|
return next.reason === "application-active-probe" || next.reason === "network-changed"
|
|
? QUICK_CONNECTION_PROBE_TIMEOUT
|
|
: undefined;
|
|
case "ConnectRequested":
|
|
case "DisconnectRequested":
|
|
case "BetterRouteAvailable":
|
|
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,
|
|
),
|
|
);
|
|
// Ticks always run: the check skips the preferred route itself, and the
|
|
// route list can grow while connected (learned routes).
|
|
yield* Stream.tick(BETTER_ROUTE_CHECK_INTERVAL).pipe(
|
|
Stream.drop(1),
|
|
Stream.runForEach(() => requestBetterRouteCheck(lease)),
|
|
Effect.forkScoped,
|
|
);
|
|
for (;;) {
|
|
const next = yield* takeSignal;
|
|
const end = yield* connectedLeaseEnd(next, lease);
|
|
if (end !== undefined) {
|
|
return end === "reset";
|
|
}
|
|
// A new network or a return to the app may have brought a better route back.
|
|
if (
|
|
(next._tag === "NetworkChanged" && next.network !== "offline") ||
|
|
(next._tag === "Wakeup" && ConnectionWakeups.resetsRetryBackoff(next.reason))
|
|
) {
|
|
yield* requestBetterRouteCheck(lease);
|
|
}
|
|
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;
|
|
}
|
|
const endDuringProbe = yield* connectedLeaseEnd(probeEvent.signal, lease);
|
|
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;
|
|
}
|
|
}
|
|
}
|
|
});
|
|
|
|
const runAttempt = Effect.fnUntraced(function* (
|
|
attempt: number,
|
|
generation: number,
|
|
lastFailure: ConnectionAttemptError | null,
|
|
pendingRetry: Option.Option<PendingRetryTrace>,
|
|
ignoreOffline: boolean,
|
|
) {
|
|
const switchingTo = yield* Ref.get(preferredRouteId);
|
|
yield* SubscriptionRef.set(prepared, Option.none());
|
|
const establishment = yield* Effect.raceAllFirst([
|
|
exitUnlessInterrupted(
|
|
establishTracedConnection(attempt, generation, lastFailure, pendingRetry),
|
|
).pipe(
|
|
Effect.map((exit): EstablishmentEvent => ({
|
|
_tag: "Completed",
|
|
exit,
|
|
})),
|
|
),
|
|
waitForEstablishmentInterrupt().pipe(
|
|
Effect.map((resetRetry): EstablishmentEvent => ({
|
|
_tag: "Interrupted",
|
|
resetRetry,
|
|
})),
|
|
),
|
|
// Each route may use the full setup time before the next is tried.
|
|
Effect.sleep(
|
|
Duration.times(establishmentTimeout, connectionRoutes(yield* Ref.get(currentEntry)).length),
|
|
).pipe(Effect.as<EstablishmentEvent>({ _tag: "TimedOut" })),
|
|
]);
|
|
|
|
if (establishment._tag === "Interrupted") {
|
|
return {
|
|
_tag: "Interrupted",
|
|
established: false,
|
|
stable: false,
|
|
resetRetry: establishment.resetRetry,
|
|
} satisfies AttemptOutcome;
|
|
}
|
|
if (establishment._tag === "TimedOut") {
|
|
return {
|
|
_tag: "Failure",
|
|
established: false,
|
|
stable: false,
|
|
failure: {
|
|
error: new ConnectionTransientError({
|
|
reason: "timeout",
|
|
detail: setupTimeoutDetail,
|
|
}),
|
|
attemptSpan: Option.none(),
|
|
},
|
|
} satisfies AttemptOutcome;
|
|
}
|
|
if (Exit.isFailure(establishment.exit)) {
|
|
const isUnexpectedDefect =
|
|
!Cause.hasInterruptsOnly(establishment.exit.cause) &&
|
|
!establishment.exit.cause.reasons.some(Cause.isFailReason);
|
|
const outcome = failureFromExit(target, establishment.exit, false, false);
|
|
if (isUnexpectedDefect) {
|
|
const defect = establishment.exit.cause.reasons.find(Cause.isDieReason)?.defect;
|
|
yield* Effect.logError("Connection attempt failed with an unexpected defect.").pipe(
|
|
Effect.annotateLogs({
|
|
"environment.id": target.environmentId,
|
|
"environment.label": target.label,
|
|
"cause.reason_count": establishment.exit.cause.reasons.length,
|
|
...safeErrorLogAttributes(defect),
|
|
}),
|
|
);
|
|
}
|
|
return outcome;
|
|
}
|
|
|
|
const active = establishment.exit.value;
|
|
if (Option.isSome(switchingTo)) {
|
|
const landed = connectionRouteId(active.lease.prepared.target);
|
|
if (landed !== switchingTo.value) {
|
|
const until =
|
|
(yield* Clock.monotonicTimeNanos) + BigInt(BETTER_ROUTE_COOLDOWN_MS) * 1_000_000n;
|
|
yield* Ref.update(routeCooldowns, (current) =>
|
|
new Map(current).set(switchingTo.value, until),
|
|
);
|
|
}
|
|
}
|
|
const currentIntent = yield* Ref.get(intent);
|
|
if (!currentIntent.desired || (currentIntent.network === "offline" && !ignoreOffline)) {
|
|
return {
|
|
_tag: "Interrupted",
|
|
established: false,
|
|
stable: false,
|
|
resetRetry: false,
|
|
} satisfies AttemptOutcome;
|
|
}
|
|
|
|
const connectedAt = yield* Clock.currentTimeMillis;
|
|
yield* SubscriptionRef.set(prepared, Option.some(active.lease.prepared));
|
|
yield* SubscriptionRef.set(session, Option.some(active.lease.session));
|
|
if (options?.learnRoutes !== undefined) {
|
|
yield* learnRoutesFrom(options.learnRoutes, active.lease).pipe(Effect.forkScoped);
|
|
}
|
|
yield* setState({
|
|
desired: true,
|
|
network: currentIntent.network,
|
|
phase: "connected",
|
|
stage: null,
|
|
attempt,
|
|
generation,
|
|
lastFailure: null,
|
|
retryAt: null,
|
|
});
|
|
|
|
const connectedExit = yield* Effect.raceFirst(
|
|
active.lease.session.closed.pipe(
|
|
Effect.mapError((error): TracedAttemptFailure => ({
|
|
error,
|
|
attemptSpan: active.attemptSpan,
|
|
})),
|
|
),
|
|
monitorConnectedLease(active.lease).pipe(
|
|
Effect.mapError((error): TracedAttemptFailure => ({
|
|
error,
|
|
attemptSpan: active.attemptSpan,
|
|
})),
|
|
),
|
|
).pipe(exitUnlessInterrupted);
|
|
const connectedForMs = (yield* Clock.currentTimeMillis) - connectedAt;
|
|
if (Exit.isSuccess(connectedExit)) {
|
|
return {
|
|
_tag: "Interrupted",
|
|
established: true,
|
|
stable: connectedForMs >= BACKOFF_RESET_AFTER_MS,
|
|
resetRetry: connectedExit.value,
|
|
} satisfies AttemptOutcome;
|
|
}
|
|
const outcome = failureFromExit(
|
|
target,
|
|
connectedExit,
|
|
true,
|
|
connectedForMs >= BACKOFF_RESET_AFTER_MS,
|
|
);
|
|
if (outcome._tag === "Failure") {
|
|
// A live session ending is otherwise invisible in the client trace, so
|
|
// record why, and how long it lasted, as its own root span.
|
|
yield* Effect.void.pipe(
|
|
Effect.withSpan("EnvironmentSupervisor.connectionLost", {
|
|
root: true,
|
|
attributes: {
|
|
"environment.id": target.environmentId,
|
|
"environment.label": target.label,
|
|
"environment.target.kind": target._tag,
|
|
"connection.connected_ms": connectedForMs,
|
|
"connection.failure.reason": outcome.failure.error.reason,
|
|
"connection.failure.detail": outcome.failure.error.detail,
|
|
},
|
|
}),
|
|
);
|
|
}
|
|
return outcome;
|
|
}, Effect.ensuring(clearLease));
|
|
|
|
const waitForRetrySignal = Effect.fnUntraced(function* (delayMs: number) {
|
|
// @effect-diagnostics-next-line raceFirstWithSleepToTimeout:off - the sleep is the retry delay (false), not a timeout around the signal loop
|
|
return yield* Effect.raceFirst(
|
|
Effect.sleep(delayMs).pipe(Effect.as(false)),
|
|
Effect.gen(function* () {
|
|
for (;;) {
|
|
const next = yield* Queue.take(signals);
|
|
switch (next._tag) {
|
|
case "Wakeup":
|
|
return ConnectionWakeups.resetsRetryBackoff(next.reason);
|
|
case "ConnectRequested":
|
|
case "DisconnectRequested":
|
|
case "RetryRequested":
|
|
case "NetworkChanged":
|
|
return false;
|
|
case "BetterRouteAvailable":
|
|
break;
|
|
}
|
|
}
|
|
}),
|
|
);
|
|
});
|
|
|
|
// A better route only matters to a live session, so idle states ignore it.
|
|
const waitForSignal = Queue.take(signals).pipe(
|
|
Effect.repeat({ while: (next) => next._tag === "BetterRouteAvailable" }),
|
|
Effect.map(
|
|
(next) => next._tag === "Wakeup" && ConnectionWakeups.resetsRetryBackoff(next.reason),
|
|
),
|
|
);
|
|
|
|
const run = Effect.fnUntraced(function* () {
|
|
let failureCount = 0;
|
|
let generation = 0;
|
|
let latestFailure: ConnectionAttemptError | null = null;
|
|
let pendingRetry = Option.none<PendingRetryTrace>();
|
|
const resetRetryLadder = () => {
|
|
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)) {
|
|
failureCount = 0;
|
|
latestFailure = null;
|
|
pendingRetry = Option.none();
|
|
}
|
|
const currentIntent = yield* Ref.get(intent);
|
|
if (!currentIntent.desired) {
|
|
resetRetryLadder();
|
|
latestFailure = null;
|
|
yield* clearLease;
|
|
yield* setState(availableState(currentIntent, generation));
|
|
yield* waitForSignal;
|
|
continue;
|
|
}
|
|
if (currentIntent.network === "offline" && !replacing) {
|
|
yield* clearLease;
|
|
yield* setState(offlineState(currentIntent, generation, failureCount + 1, latestFailure));
|
|
const applicationActivated = yield* waitForSignal;
|
|
if (applicationActivated) {
|
|
resetRetryLadder();
|
|
}
|
|
continue;
|
|
}
|
|
|
|
const attempt = failureCount + 1;
|
|
const nextGeneration = generation + 1;
|
|
const outcome: AttemptOutcome = yield* Effect.scoped(
|
|
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 failedProbe = yield* Ref.getAndSet(probeUnanswered, false);
|
|
if (outcome.established) {
|
|
generation = nextGeneration;
|
|
if (outcome.stable) {
|
|
resetRetryLadder();
|
|
latestFailure = null;
|
|
}
|
|
}
|
|
if (outcome._tag === "Interrupted") {
|
|
if (outcome.resetRetry) {
|
|
resetRetryLadder();
|
|
replacing = true;
|
|
}
|
|
continue;
|
|
}
|
|
|
|
const attemptSpan: Option.Option<Tracer.Span> = outcome.failure.attemptSpan;
|
|
const error: ConnectionAttemptError = outcome.failure.error;
|
|
latestFailure = error;
|
|
if (error._tag === "ConnectionBlockedError") {
|
|
const blockedIntent = yield* Ref.get(intent);
|
|
yield* setState({
|
|
desired: blockedIntent.desired,
|
|
network: blockedIntent.network,
|
|
phase: "blocked",
|
|
stage: null,
|
|
attempt,
|
|
generation,
|
|
lastFailure: error,
|
|
retryAt: null,
|
|
});
|
|
const applicationActivated = yield* waitForSignal;
|
|
if (applicationActivated) {
|
|
resetRetryLadder();
|
|
}
|
|
continue;
|
|
}
|
|
|
|
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, yield* Random.next);
|
|
pendingRetry = Option.map(attemptSpan, (previousAttempt) => ({
|
|
previousAttempt,
|
|
failureCount,
|
|
delayMs,
|
|
reason: error.reason,
|
|
}));
|
|
const failedIntent = yield* Ref.get(intent);
|
|
yield* setState({
|
|
desired: failedIntent.desired,
|
|
network: failedIntent.network,
|
|
phase: "backoff",
|
|
stage: null,
|
|
attempt,
|
|
generation,
|
|
lastFailure: error,
|
|
retryAt: (yield* Clock.currentTimeMillis) + delayMs,
|
|
});
|
|
const applicationActivated = yield* waitForRetrySignal(delayMs);
|
|
if (applicationActivated) {
|
|
resetRetryLadder();
|
|
}
|
|
}
|
|
});
|
|
|
|
yield* connectivity.changes.pipe(
|
|
Stream.runForEach((network) =>
|
|
Ref.modify(intent, (current) =>
|
|
current.network === network ? [false, current] : ([true, { ...current, network }] as const),
|
|
).pipe(
|
|
Effect.flatMap((changed) =>
|
|
changed ? signal({ _tag: "NetworkChanged", network }) : Effect.void,
|
|
),
|
|
),
|
|
),
|
|
Effect.forkScoped,
|
|
);
|
|
yield* wakeups.changes.pipe(
|
|
Stream.runForEach((reason) => signal({ _tag: "Wakeup", reason })),
|
|
Effect.forkScoped,
|
|
);
|
|
yield* Queue.take(betterRouteChecks).pipe(
|
|
Effect.flatMap(checkBetterRoutes),
|
|
Effect.forever,
|
|
Effect.forkScoped,
|
|
);
|
|
yield* run().pipe(Effect.forkScoped);
|
|
|
|
const connect = Ref.update(intent, (current) => ({
|
|
...current,
|
|
desired: true,
|
|
})).pipe(
|
|
Effect.andThen(signal({ _tag: "ConnectRequested" })),
|
|
Effect.withSpan("EnvironmentSupervisor.connect"),
|
|
);
|
|
|
|
const disconnect = Ref.update(intent, (current) => ({
|
|
...current,
|
|
desired: false,
|
|
})).pipe(
|
|
Effect.andThen(signal({ _tag: "DisconnectRequested" })),
|
|
Effect.withSpan("EnvironmentSupervisor.disconnect"),
|
|
);
|
|
|
|
const retryNow = Ref.set(resetRetryState, true).pipe(
|
|
Effect.andThen(signal({ _tag: "RetryRequested" })),
|
|
Effect.withSpan("EnvironmentSupervisor.retryNow"),
|
|
);
|
|
|
|
yield* Effect.addFinalizer(() => Queue.shutdown(signals).pipe(Effect.andThen(clearLease)));
|
|
|
|
return EnvironmentSupervisor.of({
|
|
target,
|
|
state,
|
|
session,
|
|
prepared,
|
|
connect,
|
|
disconnect,
|
|
retryNow,
|
|
});
|
|
});
|