import { EnvironmentId } from "@t3tools/contracts";
import * as Context from "effect/Context";
import * as Effect from "effect/Effect";
import * as Equal from "effect/Equal";
import * as Exit from "effect/Exit";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as Ref from "effect/Ref";
import * as Schema from "effect/Schema";
import * as Scope from "effect/Scope";
import * as Semaphore from "effect/Semaphore";
import * as Stream from "effect/Stream";
import * as SubscriptionRef from "effect/SubscriptionRef";
import * as ClientCapabilities from "../platform/capabilities.ts";
import {
type ConnectionCatalogEntry,
type ConnectionRegistration,
type PlatformConnectionRegistration,
type PrimaryConnectionRegistration,
SshConnectionProfile,
connectionRegistrationCatalogEntry,
} from "./catalog.ts";
import * as ConnectionCredentialStore from "./credentialStore.ts";
import * as ConnectionProfileStore from "./profileStore.ts";
import * as Connectivity from "./connectivity.ts";
import type {
ConnectionAttemptError,
ConnectionTarget,
NetworkStatus,
SupervisorConnectionState,
} from "./model.ts";
import { ConnectionBlockedError } from "./model.ts";
import * as Persistence from "../platform/persistence.ts";
import * as EnvironmentSupervisor from "./supervisor.ts";
import * as ConnectionDriver from "./driver.ts";
import * as ConnectionWakeups from "./wakeups.ts";
import {
GitHubRoutingPermissions,
gitHubRoutingConnectionKey,
} from "./githubRoutingPermissions.ts";
const isSshConnectionProfile = Schema.is(SshConnectionProfile);
function unsupportedState(
entry: ConnectionCatalogEntry,
): Pick<ConnectionCatalogEntry, "unsupportedReason" | "serverUpdateRequired"> {
return {
...(entry.unsupportedReason === undefined
? {}
: { unsupportedReason: entry.unsupportedReason }),
...(entry.serverUpdateRequired === true ? { serverUpdateRequired: true } : {}),
};
}
export class EnvironmentNotRegisteredError extends Schema.TaggedError<EnvironmentNotRegisteredError>()(
"EnvironmentNotRegisteredError",
{
environmentId: EnvironmentId,
},
) {
override get message(): string {
return `Environment ${this.environmentId} is not registered.`;
}
}
export class PlatformEnvironmentRemovalError extends Schema.TaggedError<PlatformEnvironmentRemovalError>()(
"PlatformEnvironmentRemovalError",
{
environmentId: EnvironmentId,
},
) {
override get message(): string {
return `Platform-managed environment ${this.environmentId} cannot be removed.`;
}
}
export class EnvironmentRegistry extends Context.Service<
EnvironmentRegistry,
{
readonly entries: SubscriptionRef.SubscriptionRef<
ReadonlyMap<EnvironmentId, ConnectionCatalogEntry>
>;
readonly networkStatus: SubscriptionRef.SubscriptionRef<NetworkStatus>;
readonly start: Effect.Effect<void>;
readonly register: (
registration: ConnectionRegistration,
) => Effect.Effect<void, Persistence.ConnectionPersistenceError>;
readonly registerPlatform: (registration: PrimaryConnectionRegistration) => Effect.Effect<void>;
readonly reconcilePlatform: (
registrations: ReadonlyArray<PlatformConnectionRegistration>,
) => Effect.Effect<void>;
readonly remove: (
environmentId: EnvironmentId,
) => Effect.Effect<
void,
| Persistence.ConnectionPersistenceError
| ConnectionAttemptError
| EnvironmentNotRegisteredError
| PlatformEnvironmentRemovalError
>;
readonly removeRelayEnvironments: () => Effect.Effect<
void,
| Persistence.ConnectionPersistenceError
| ConnectionAttemptError
| PlatformEnvironmentRemovalError
>;
readonly retryNow: (environmentId: EnvironmentId) => Effect.Effect<void>;
/**
* Switches a saved environment on or off. Off drops the socket, stops the
* retry ladder, and persists so the next launch stays off. Registration,
* credentials, and cache are untouched.
*/
readonly setEnabled: (
environmentId: EnvironmentId,
enabled: boolean,
) => Effect.Effect<
void,
| EnvironmentNotRegisteredError
| Persistence.ConnectionPersistenceError
| ConnectionBlockedError
>;
readonly setCompatibility: (
environmentId: EnvironmentId,
error: ConnectionBlockedError | null,
) => Effect.Effect<void, Persistence.ConnectionPersistenceError>;
readonly state: (
environmentId: EnvironmentId,
) => Effect.Effect<SupervisorConnectionState, EnvironmentNotRegisteredError>;
readonly stateChanges: (
environmentId: EnvironmentId,
) => Stream.Stream<SupervisorConnectionState, EnvironmentNotRegisteredError>;
readonly run: <A, E, R>(
environmentId: EnvironmentId,
effect: Effect.Effect<A, E, R>,
) => Effect.Effect<
A,
E | EnvironmentNotRegisteredError,
Exclude<R, EnvironmentSupervisor.EnvironmentSupervisor>
>;
readonly runStream: <A, E, R>(
environmentId: EnvironmentId,
stream: Stream.Stream<A, E, R>,
) => Stream.Stream<
A,
E | EnvironmentNotRegisteredError,
Exclude<R, EnvironmentSupervisor.EnvironmentSupervisor>
>;
readonly followStream: <A, E, R>(
environmentId: EnvironmentId,
stream: Stream.Stream<A, E, R>,
) => Stream.Stream<A, E, Exclude<R, EnvironmentSupervisor.EnvironmentSupervisor>>;
}
>()("@t3tools/client-runtime/connection/registry/EnvironmentRegistry") {}
interface EnvironmentServiceScope {
readonly entry: ConnectionCatalogEntry;
readonly supervisor: EnvironmentSupervisor.EnvironmentSupervisor["Service"];
readonly scope: Scope.Closeable;
}
/** @public Service construction is part of the canonical Effect module API. */
export const make = Effect.gen(function* () {
const registryScope = yield* Scope.Scope;
const storage = yield* Persistence.ConnectionTargetStore;
const registrations = yield* Persistence.ConnectionRegistrationStore;
const cache = yield* Persistence.EnvironmentCacheStore;
const ownedDataCleanup = yield* Persistence.EnvironmentOwnedDataCleanup;
const profiles = yield* ConnectionProfileStore.ConnectionProfileStore;
const credentials = yield* ConnectionCredentialStore.ConnectionCredentialStore;
const githubRoutingPermissions = yield* GitHubRoutingPermissions;
const connectivity = yield* Connectivity.Connectivity;
const driver = yield* ConnectionDriver.ConnectionDriver;
const wakeups = yield* ConnectionWakeups.ConnectionWakeups;
const ssh = yield* ClientCapabilities.SshEnvironmentGateway;
const persistedTargets = yield* storage.list;
const disabledEnvironmentIds = new Set(yield* storage.listDisabled);
const initialEntries = new Map(
yield* Effect.forEach(
persistedTargets,
Effect.fn("EnvironmentRegistry.loadCatalogEntry")(function* (target) {
const profile =
target._tag === "BearerConnectionTarget" || target._tag === "SshConnectionTarget"
? yield* profiles.get(target.connectionId)
: Option.none();
return [
target.environmentId,
{
target,
profile,
enabled: !disabledEnvironmentIds.has(target.environmentId),
} satisfies ConnectionCatalogEntry,
] as const;
}),
{ concurrency: "unbounded" },
),
);
const entries =
yield* SubscriptionRef.make<ReadonlyMap<EnvironmentId, ConnectionCatalogEntry>>(initialEntries);
const networkStatus = yield* SubscriptionRef.make(yield* connectivity.status);
const serviceScopes = yield* SubscriptionRef.make<
ReadonlyMap<EnvironmentId, EnvironmentServiceScope>
>(new Map());
const platformEnvironmentIds = yield* Ref.make<ReadonlySet<EnvironmentId>>(new Set());
const persistedTargetsByEnvironment = yield* Ref.make<
ReadonlyMap<EnvironmentId, ConnectionTarget>
>(new Map(persistedTargets.map((target) => [target.environmentId, target])));
interface LeaseLock {
readonly semaphore: Semaphore.Semaphore;
readonly users: number;
}
const leaseLocks = yield* Ref.make<ReadonlyMap<EnvironmentId, LeaseLock>>(new Map());
const leaseLocksGuard = yield* Semaphore.make(1);
const started = yield* Ref.make(false);
const withLeaseLock = <A, E, R>(
environmentId: EnvironmentId,
effect: Effect.Effect<A, E, R>,
): Effect.Effect<A, E, R> =>
Effect.acquireUseRelease(
leaseLocksGuard.withPermits(1)(
Effect.gen(function* () {
const current = yield* Ref.get(leaseLocks);
const existing = current.get(environmentId);
if (existing !== undefined) {
yield* Ref.set(
leaseLocks,
new Map(current).set(environmentId, {
semaphore: existing.semaphore,
users: existing.users + 1,
}),
);
return existing.semaphore;
}
const semaphore = yield* Semaphore.make(1);
yield* Ref.set(leaseLocks, new Map(current).set(environmentId, { semaphore, users: 1 }));
return semaphore;
}),
),
(semaphore) => semaphore.withPermits(1)(effect),
(semaphore) =>
leaseLocksGuard.withPermits(1)(
Ref.update(leaseLocks, (current) => {
const existing = current.get(environmentId);
if (existing === undefined || existing.semaphore !== semaphore) {
return current;
}
const next = new Map(current);
if (existing.users === 1) {
next.delete(environmentId);
} else {
next.set(environmentId, {
semaphore,
users: existing.users - 1,
});
}
return next;
}),
),
).pipe(Effect.withSpan("EnvironmentRegistry.withLeaseLock"));
const getEntry = Effect.fn("EnvironmentRegistry.getEntry")(function* (
environmentId: EnvironmentId,
) {
const entry = (yield* SubscriptionRef.get(entries)).get(environmentId);
if (entry === undefined) {
return yield* new EnvironmentNotRegisteredError({
environmentId,
});
}
return entry;
});
const closeServiceScope = Effect.fn("EnvironmentRegistry.closeServiceScope")(function* (
environmentId: EnvironmentId,
) {
const current = yield* SubscriptionRef.get(serviceScopes);
const lease = current.get(environmentId);
if (lease === undefined) {
return;
}
const next = new Map(current);
next.delete(environmentId);
yield* SubscriptionRef.set(serviceScopes, next);
yield* Scope.close(lease.scope, Exit.void);
});
const createServiceScope = Effect.fn("EnvironmentRegistry.createServiceScope")(
(entry: ConnectionCatalogEntry) =>
Effect.uninterruptible(
Effect.gen(function* () {
const environmentId = entry.target.environmentId;
const scope = yield* Scope.fork(registryScope);
const supervisor = yield* EnvironmentSupervisor.make(entry, {
initiallyDesired: false,
}).pipe(
Effect.provideService(Connectivity.Connectivity, connectivity),
Effect.provideService(ConnectionDriver.ConnectionDriver, driver),
Effect.provideService(ConnectionWakeups.ConnectionWakeups, wakeups),
Scope.provide(scope),
Effect.onError(() => Scope.close(scope, Exit.void)),
);
if (entry.enabled) {
yield* supervisor.connect;
}
yield* SubscriptionRef.update(serviceScopes, (current) => {
const next = new Map(current);
next.set(environmentId, { entry, supervisor, scope });
return next;
});
yield* SubscriptionRef.changes(supervisor.state).pipe(
Stream.runForEach((state) =>
state.phase === "blocked" && state.lastFailure?.reason === "unsupported"
? setCompatibility(environmentId, state.lastFailure).pipe(
Effect.catch((error) =>
Effect.logWarning("Could not disable an unsupported environment.", {
environmentId,
error,
}),
),
)
: Effect.void,
),
Effect.forkIn(scope),
);
return supervisor;
}),
),
);
const acquireSupervisor = Effect.fn("EnvironmentRegistry.acquireSupervisor")(function* (
environmentId: EnvironmentId,
) {
return yield* withLeaseLock(
environmentId,
Effect.gen(function* () {
const entry = yield* getEntry(environmentId);
const existing = (yield* SubscriptionRef.get(serviceScopes)).get(environmentId);
if (existing !== undefined) {
if (Equal.equals(existing.entry, entry)) {
return existing.supervisor;
}
yield* closeServiceScope(environmentId);
}
return yield* createServiceScope(entry);
}),
);
});
const run: EnvironmentRegistry["Service"]["run"] = Effect.fn("EnvironmentRegistry.run")(
function* <A, E, R>(environmentId: EnvironmentId, effect: Effect.Effect<A, E, R>) {
const supervisor = yield* acquireSupervisor(environmentId);
return yield* Effect.provideService(
effect,
EnvironmentSupervisor.EnvironmentSupervisor,
supervisor,
);
},
);
const runStream: EnvironmentRegistry["Service"]["runStream"] = <A, E, R>(
environmentId: EnvironmentId,
stream: Stream.Stream<A, E, R>,
) =>
Stream.unwrap(
acquireSupervisor(environmentId).pipe(
Effect.map((supervisor) =>
Stream.provideService(stream, EnvironmentSupervisor.EnvironmentSupervisor, supervisor),
),
),
);
const followStream: EnvironmentRegistry["Service"]["followStream"] = <A, E, R>(
environmentId: EnvironmentId,
stream: Stream.Stream<A, E, R>,
) =>
Stream.concat(
Stream.fromEffect(SubscriptionRef.get(entries)),
SubscriptionRef.changes(entries),
).pipe(
Stream.map((current) => Option.fromUndefinedOr(current.get(environmentId))),
Stream.changes,
Stream.switchMap(
Option.match({
onNone: () => Stream.empty,
onSome: () =>
Stream.unwrap(
acquireSupervisor(environmentId).pipe(
Effect.match({
onFailure: () => Stream.empty,
onSuccess: (supervisor) =>
Stream.provideService(
stream,
EnvironmentSupervisor.EnvironmentSupervisor,
supervisor,
),
}),
),
),
}),
),
);
const start = Effect.gen(function* () {
if (yield* Ref.getAndSet(started, true)) {
return;
}
yield* Effect.forEach(
persistedTargets,
(target) =>
acquireSupervisor(target.environmentId).pipe(
Effect.catchTag("EnvironmentNotRegisteredError", () => Effect.void),
),
{
concurrency: "unbounded",
discard: true,
},
);
}).pipe(Effect.withSpan("EnvironmentRegistry.start"));
const installEntryLocked = Effect.fn("EnvironmentRegistry.installEntryLocked")(function* (
entry: ConnectionCatalogEntry,
options?: { readonly retainEquivalentRuntime?: boolean },
) {
const target = entry.target;
const previous = (yield* SubscriptionRef.get(entries)).get(target.environmentId);
const existingScope = (yield* SubscriptionRef.get(serviceScopes)).get(target.environmentId);
if (
options?.retainEquivalentRuntime === true &&
previous !== undefined &&
Equal.equals(previous, entry) &&
existingScope !== undefined &&
Equal.equals(existingScope.entry, entry)
) {
return;
}
yield* closeServiceScope(target.environmentId);
yield* SubscriptionRef.update(entries, (current) => {
const next = new Map(current);
next.set(target.environmentId, entry);
return next;
});
yield* createServiceScope(entry);
});
const register = Effect.fn("EnvironmentRegistry.register")(function* (
registration: ConnectionRegistration,
) {
const registered = connectionRegistrationCatalogEntry(registration);
const environmentId = registered.target.environmentId;
yield* withLeaseLock(
environmentId,
Effect.gen(function* () {
if ((yield* Ref.get(platformEnvironmentIds)).has(environmentId)) {
return;
}
// Editing a saved environment must preserve its disabled state.
const previous = (yield* SubscriptionRef.get(entries)).get(environmentId);
const entry: ConnectionCatalogEntry =
previous === undefined
? registered
: {
...registered,
enabled: previous.enabled,
...(previous.unsupportedReason !== undefined &&
gitHubRoutingConnectionKey(previous) === gitHubRoutingConnectionKey(registered)
? unsupportedState(previous)
: {}),
};
if (
previous !== undefined &&
gitHubRoutingConnectionKey(previous) !== gitHubRoutingConnectionKey(entry)
) {
yield* githubRoutingPermissions.forget(environmentId).pipe(
Effect.mapError(
(error) =>
new Persistence.ConnectionPersistenceError({
operation: "register-connection",
message: error.message,
}),
),
);
}
yield* registrations.register(registration);
yield* Ref.update(persistedTargetsByEnvironment, (current) => {
const next = new Map(current);
next.set(environmentId, registration.target);
return next;
});
yield* installEntryLocked(entry);
}),
);
});
const installPlatformRegistration = Effect.fn("EnvironmentRegistry.installPlatformRegistration")(
function* (registration: PlatformConnectionRegistration) {
const registered = connectionRegistrationCatalogEntry(registration);
const target = registered.target;
yield* withLeaseLock(
target.environmentId,
Effect.gen(function* () {
const previous = (yield* SubscriptionRef.get(entries)).get(target.environmentId);
const entry: ConnectionCatalogEntry =
previous?.unsupportedReason !== undefined &&
gitHubRoutingConnectionKey(previous) === gitHubRoutingConnectionKey(registered)
? { ...registered, enabled: false, ...unsupportedState(previous) }
: registered;
const persistedTarget = (yield* Ref.get(persistedTargetsByEnvironment)).get(
target.environmentId,
);
if (
persistedTarget !== undefined ||
(previous !== undefined &&
gitHubRoutingConnectionKey(previous) !== gitHubRoutingConnectionKey(entry))
) {
const revoked = yield* githubRoutingPermissions.forget(target.environmentId).pipe(
Effect.tapError((error) =>
Effect.logWarning(
"Could not clear GitHub routing permission for a platform environment.",
{
environmentId: target.environmentId,
error,
},
),
),
Effect.exit,
);
if (Exit.isFailure(revoked)) return;
}
yield* Ref.update(platformEnvironmentIds, (current) => {
const next = new Set(current);
next.add(target.environmentId);
return next;
});
// Secondary desktop-local backends (e.g. a parallel WSL backend) live
// on their own loopback origin, so they authenticate with a bearer
// token instead of the primary's same-origin cookie. Stash it where
// the resolver's bearer broker looks it up.
if (registration._tag === "BearerConnectionRegistration") {
yield* credentials.put(registration.target.connectionId, registration.credential).pipe(
Effect.catch((error) =>
Effect.logWarning("Could not store the platform bearer credential.", {
environmentId: target.environmentId,
error,
}),
),
);
}
if (persistedTarget !== undefined) {
yield* registrations.remove(persistedTarget).pipe(
Effect.tap(() =>
Ref.update(persistedTargetsByEnvironment, (current) => {
const next = new Map(current);
next.delete(target.environmentId);
return next;
}),
),
Effect.catch((error) =>
Effect.logWarning(
"Could not remove a persisted registration shadowed by a platform environment.",
{
environmentId: target.environmentId,
error,
},
),
),
);
}
yield* installEntryLocked(entry, { retainEquivalentRuntime: true });
}),
);
},
);
// Tear down a platform-managed environment that the host no longer reports
// (e.g. the user turned the parallel WSL backend off). Platform environments
// bypass the user-facing `remove` guard since they are reconciled from the
// bootstrap rather than removed by hand.
const removePlatformEnvironment = Effect.fn("EnvironmentRegistry.removePlatformEnvironment")(
function* (environmentId: EnvironmentId) {
yield* withLeaseLock(
environmentId,
Effect.gen(function* () {
const entry = (yield* SubscriptionRef.get(entries)).get(environmentId);
const revoked = yield* githubRoutingPermissions.forget(environmentId).pipe(
Effect.tapError((error) =>
Effect.logWarning(
"Could not clear GitHub routing permission after platform removal.",
{
environmentId,
error,
},
),
),
Effect.exit,
);
if (Exit.isFailure(revoked)) return;
yield* Ref.update(platformEnvironmentIds, (current) => {
const next = new Set(current);
next.delete(environmentId);
return next;
});
yield* closeServiceScope(environmentId);
yield* SubscriptionRef.update(entries, (current) => {
const next = new Map(current);
next.delete(environmentId);
return next;
});
if (entry !== undefined && entry.target._tag === "BearerConnectionTarget") {
yield* credentials.remove(entry.target.connectionId).pipe(
Effect.catch((error) =>
Effect.logWarning("Could not clear the platform bearer credential.", {
environmentId,
error,
}),
),
);
}
yield* Effect.all(
[
cache.clear(environmentId).pipe(
Effect.catch((error) =>
Effect.logWarning("Could not clear cached environment data after removal.", {
environmentId,
error,
}),
),
),
ownedDataCleanup.clear(environmentId),
],
{ concurrency: "unbounded", discard: true },
);
}),
);
},
);
const registerPlatform = Effect.fn("EnvironmentRegistry.registerPlatform")(function* (
registration: PrimaryConnectionRegistration,
) {
yield* installPlatformRegistration(registration);
});
// Reconcile the full set of platform-managed environments against what the
// host currently reports: add/refresh the desired ones and tear down any
// platform environment that disappeared (WSL toggled off, distro switched).
const reconcilePlatform = Effect.fn("EnvironmentRegistry.reconcilePlatform")(function* (
platformRegistrations: ReadonlyArray<PlatformConnectionRegistration>,
) {
const desiredIds = new Set(
platformRegistrations.map((registration) => registration.target.environmentId),
);
const currentPlatformIds = yield* Ref.get(platformEnvironmentIds);
yield* Effect.forEach(
currentPlatformIds,
(environmentId) =>
desiredIds.has(environmentId) ? Effect.void : removePlatformEnvironment(environmentId),
{ discard: true },
);
yield* Effect.forEach(platformRegistrations, installPlatformRegistration, { discard: true });
});
const remove = Effect.fn("EnvironmentRegistry.remove")(function* (environmentId: EnvironmentId) {
return yield* withLeaseLock(
environmentId,
Effect.gen(function* () {
if ((yield* Ref.get(platformEnvironmentIds)).has(environmentId)) {
return yield* new PlatformEnvironmentRemovalError({
environmentId,
});
}
const target = (yield* getEntry(environmentId)).target;
const profile =
target._tag === "BearerConnectionTarget" || target._tag === "SshConnectionTarget"
? yield* profiles.get(target.connectionId)
: Option.none();
yield* githubRoutingPermissions.forget(environmentId);
yield* registrations.remove(target);
yield* Ref.update(persistedTargetsByEnvironment, (current) => {
const next = new Map(current);
next.delete(environmentId);
return next;
});
yield* closeServiceScope(environmentId);
yield* SubscriptionRef.update(entries, (current) => {
const next = new Map(current);
next.delete(environmentId);
return next;
});
yield* Effect.all(
[
cache.clear(environmentId).pipe(
Effect.catch((error) =>
Effect.logWarning("Could not clear cached environment data after removal.", {
environmentId,
error,
}),
),
),
ownedDataCleanup.clear(environmentId),
],
{ concurrency: "unbounded", discard: true },
);
if (
target._tag === "SshConnectionTarget" &&
Option.isSome(profile) &&
isSshConnectionProfile(profile.value)
) {
yield* ssh.disconnect(profile.value.target).pipe(
Effect.tapError((error) =>
Effect.logWarning("Could not disconnect the managed SSH environment.", {
environmentId,
error,
}),
),
Effect.ignore,
);
}
}),
);
});
const removeRelayEnvironments = Effect.fn("EnvironmentRegistry.removeRelayEnvironments")(
function* () {
const relayEnvironmentIds = [...(yield* SubscriptionRef.get(entries)).values()]
.filter((entry) => entry.target._tag === "RelayConnectionTarget")
.map((entry) => entry.target.environmentId);
yield* Effect.forEach(
relayEnvironmentIds,
(environmentId) =>
remove(environmentId).pipe(
Effect.catchTag("EnvironmentNotRegisteredError", () => Effect.void),
),
{
concurrency: "unbounded",
discard: true,
},
);
},
);
const retryNow = (environmentId: EnvironmentId) =>
acquireSupervisor(environmentId).pipe(
Effect.flatMap((supervisor) => supervisor.retryNow),
Effect.catchTag("EnvironmentNotRegisteredError", () => Effect.void),
Effect.withSpan("EnvironmentRegistry.retryNow"),
);
const setEnabled = Effect.fn("EnvironmentRegistry.setEnabled")(function* (
environmentId: EnvironmentId,
enabled: boolean,
) {
yield* withLeaseLock(
environmentId,
Effect.gen(function* () {
const entry = yield* getEntry(environmentId);
if (enabled && entry.unsupportedReason !== undefined) {
return yield* new ConnectionBlockedError({
reason: "unsupported",
detail: entry.unsupportedReason,
});
}
if (entry.enabled === enabled) {
return;
}
// Platform-managed environments are reconciled from the host and are
// never persisted, so only user-saved ones write the flag.
if (!(yield* Ref.get(platformEnvironmentIds)).has(environmentId)) {
yield* registrations.setEnabled(environmentId, enabled);
}
const next: ConnectionCatalogEntry = { ...entry, enabled };
// Update the lease in place so the supervisor keeps its generation and
// durable streams; `installEntryLocked` would tear it down instead.
const lease = (yield* SubscriptionRef.get(serviceScopes)).get(environmentId);
if (lease !== undefined) {
yield* SubscriptionRef.update(serviceScopes, (current) => {
const nextScopes = new Map(current);
nextScopes.set(environmentId, { ...lease, entry: next });
return nextScopes;
});
}
yield* SubscriptionRef.update(entries, (current) => {
const nextEntries = new Map(current);
nextEntries.set(environmentId, next);
return nextEntries;
});
if (lease !== undefined) {
yield* enabled ? lease.supervisor.connect : lease.supervisor.disconnect;
} else if (enabled) {
yield* createServiceScope(next);
}
// The supervisor only owns the RPC session. A managed SSH backend and
// its tunnel outlive it, so switching off tears those down as well.
if (
!enabled &&
entry.target._tag === "SshConnectionTarget" &&
Option.isSome(entry.profile) &&
isSshConnectionProfile(entry.profile.value)
) {
yield* ssh.disconnect(entry.profile.value.target).pipe(
Effect.tapError((error) =>
Effect.logWarning("Could not disconnect the switched-off SSH environment.", {
environmentId,
error,
}),
),
Effect.ignore,
);
}
}),
);
});
const state = Effect.fn("EnvironmentRegistry.state")(function* (environmentId: EnvironmentId) {
const supervisor = yield* acquireSupervisor(environmentId);
return yield* SubscriptionRef.get(supervisor.state);
});
const stateChanges = (environmentId: EnvironmentId) =>
followStream(
environmentId,
Stream.unwrap(
EnvironmentSupervisor.EnvironmentSupervisor.pipe(
Effect.map((supervisor) => SubscriptionRef.changes(supervisor.state)),
),
),
);
yield* Effect.addFinalizer(() =>
SubscriptionRef.get(serviceScopes).pipe(
Effect.flatMap((current) =>
Effect.forEach(current.values(), (lease) => Scope.close(lease.scope, Exit.void), {
concurrency: "unbounded",
discard: true,
}),
),
),
);
yield* connectivity.changes.pipe(
Stream.runForEach((status) => SubscriptionRef.set(networkStatus, status)),
Effect.forkScoped,
);
const setCompatibility = Effect.fn("EnvironmentRegistry.setCompatibility")(function* (
environmentId: EnvironmentId,
error: ConnectionBlockedError | null,
) {
yield* withLeaseLock(
environmentId,
Effect.gen(function* () {
const entry = (yield* SubscriptionRef.get(entries)).get(environmentId);
if (
entry === undefined ||
(entry.unsupportedReason === (error?.message ?? undefined) &&
entry.serverUpdateRequired === (error?.serverUpdateRequired ?? undefined))
)
return;
const {
unsupportedReason: _previousReason,
serverUpdateRequired: _previousUpdateRequired,
...rest
} = entry;
const next: ConnectionCatalogEntry =
error === null
? rest
: {
...rest,
enabled: false,
unsupportedReason: error.message,
...(error.serverUpdateRequired === true ? { serverUpdateRequired: true } : {}),
};
if (
error !== null &&
entry.enabled &&
!(yield* Ref.get(platformEnvironmentIds)).has(environmentId)
) {
yield* registrations.setEnabled(environmentId, false);
}
const lease = (yield* SubscriptionRef.get(serviceScopes)).get(environmentId);
if (lease !== undefined) {
yield* SubscriptionRef.update(serviceScopes, (current) =>
new Map(current).set(environmentId, { ...lease, entry: next }),
);
if (error !== null) yield* lease.supervisor.disconnect;
}
yield* SubscriptionRef.update(entries, (current) =>
new Map(current).set(environmentId, next),
);
}),
);
});
return EnvironmentRegistry.of({
entries,
networkStatus,
start,
register,
registerPlatform,
reconcilePlatform,
remove,
removeRelayEnvironments,
retryNow,
setEnabled,
setCompatibility,
state,
stateChanges,
run,
runStream,
followStream,
});
});
export const layer = Layer.effect(EnvironmentRegistry, make);