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 } 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";
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";
// 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 {
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 };
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;
}
/**
* 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;
const setupTimeoutDetail = `${target.label} did not respond during connection setup.${
target._tag === "RelayConnectionTarget" ? ` ${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());
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(entry, (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 (target._tag === "RelayConnectionTarget") {
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":
break;
case "Wakeup":
if (next.reason === "application-active-reconnect") {
return true;
}
if (next.reason === "credentials-changed" && target._tag === "RelayConnectionTarget") {
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 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* 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;
}
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;
}
}
}
});
const runAttempt = Effect.fnUntraced(function* (
attempt: number,
generation: number,
lastFailure: ConnectionAttemptError | null,
pendingRetry: Option.Option<PendingRetryTrace>,
ignoreOffline: boolean,
) {
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,
})),
),
Effect.sleep(CONNECTION_ESTABLISHMENT_TIMEOUT).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;
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));
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;
}
return failureFromExit(target, connectedExit, true, connectedForMs >= BACKOFF_RESET_AFTER_MS);
}, 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.isApplicationActiveWakeup(next.reason);
case "ConnectRequested":
case "DisconnectRequested":
case "RetryRequested":
case "NetworkChanged":
return false;
}
}
}),
);
});
const waitForSignal = Queue.take(signals).pipe(
Effect.map(
(next) => next._tag === "Wakeup" && ConnectionWakeups.isApplicationActiveWakeup(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* 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,
});
});