Repository navigation
Expand file tree
/
Copy pathdriver.ts
More file actions
205 lines (192 loc) · 8 KB
/
Copy pathdriver.ts
File metadata and controls
205 lines (192 loc) · 8 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
import * as Context from "effect/Context";
import * as Effect from "effect/Effect";
import * as Exit from "effect/Exit";
import * as Fiber from "effect/Fiber";
import * as Layer from "effect/Layer";
import * as Scope from "effect/Scope";
import * as HttpClient from "effect/http/HttpClient";
import type { ConnectionCatalogEntry, ConnectionRoute } from "./catalog.ts";
import type {
ConnectionAttemptError,
ConnectionAttemptStage,
PreparedConnection,
} from "./model.ts";
import { ConnectionTransientError } from "./model.ts";
import * as ConnectionResolver from "./resolver.ts";
import { connectionRoutes, routeEntry, routeHttpBaseUrl } from "./routes.ts";
import * as RpcSession from "../rpc/session.ts";
import { fetchRemoteEnvironmentDescriptor } from "../environment/descriptor.ts";
export type ConnectionDriverProgress =
| {
readonly stage: "preparing";
}
| {
readonly stage: Exclude<ConnectionAttemptStage, "preparing">;
readonly prepared: PreparedConnection;
};
export interface EnvironmentConnectionLease {
readonly prepared: PreparedConnection;
readonly session: RpcSession.RpcSession;
}
/**
* The result of an unauthenticated reachability check. T3 Connect and SSH
* routes have no cheap check, so they are "unchecked".
*/
export type RouteCheck = "answered" | "silent" | "unchecked";
/** How long a direct route has to answer before it counts as unreachable from here. */
const ROUTE_CHECK_TIMEOUT_MS = 2_500;
export class ConnectionDriver extends Context.Service<
ConnectionDriver,
{
readonly connect: (
entry: ConnectionCatalogEntry,
reportProgress: (progress: ConnectionDriverProgress) => Effect.Effect<void>,
) => Effect.Effect<EnvironmentConnectionLease, ConnectionAttemptError, Scope.Scope>;
/** Whether a direct route answers as the entry's environment, without credentials. */
readonly checkRoute: (
entry: ConnectionCatalogEntry,
route: ConnectionRoute,
) => Effect.Effect<RouteCheck>;
/**
* Whether a direct route answers and accepts this client's credential,
* without opening a socket. Switching to a route that fails this would
* drop a working connection for nothing.
*/
readonly preflight: (
entry: ConnectionCatalogEntry,
route: ConnectionRoute,
) => Effect.Effect<boolean>;
}
>()("@t3tools/client-runtime/connection/driver/ConnectionDriver") {}
/**
* Connects over the first route, in preference order, that is worth trying.
* Every route is checked at once, but a route only waits for its own check,
* so a reachable LAN address connects without waiting on a silent tailnet
* one. A silent route is skipped on the first pass so a LAN address from
* another network costs one short check, not a connection timeout. A route
* that fails to connect moves on to the next: a signed-out T3 Connect must
* not hide a working LAN. Silent routes are tried last, since a check is not
* proof.
*
* The reported error is a transient one when any route failed transiently,
* so the supervisor keeps retrying a route that may come back; a blocked
* error is reported only when every attempted route was blocked.
*/
export const connectOverRoutes = Effect.fn("ConnectionDriver.connectOverRoutes")(function* <R>(
entry: ConnectionCatalogEntry,
checkRoute: (route: ConnectionRoute) => Effect.Effect<RouteCheck>,
connectRoute: (
route: ConnectionRoute,
) => Effect.Effect<EnvironmentConnectionLease, ConnectionAttemptError, R | Scope.Scope>,
) {
const routes = connectionRoutes(entry);
const checks =
routes.length === 1
? []
: yield* Effect.forEach(routes, (route) => Effect.forkChild(checkRoute(route)));
const attemptScope = yield* Scope.Scope;
let transient: ConnectionAttemptError | undefined;
let blocked: ConnectionAttemptError | undefined;
// Each route gets its own scope so a half-open session closes before the next try.
const attempt = Effect.fnUntraced(function* (route: ConnectionRoute) {
const routeScope = yield* Scope.fork(attemptScope);
const result = yield* connectRoute(route).pipe(
Scope.provide(routeScope),
Effect.onExit((exit) => (Exit.isSuccess(exit) ? Effect.void : Scope.close(routeScope, exit))),
Effect.result,
);
if (result._tag === "Failure") {
if (result.failure._tag === "ConnectionTransientError") transient ??= result.failure;
else blocked ??= result.failure;
}
return result;
});
const silent: Array<ConnectionRoute> = [];
for (const [index, route] of routes.entries()) {
const check = checks.length === 0 ? "unchecked" : yield* Fiber.join(checks[index]!);
if (check === "silent") {
silent.push(route);
continue;
}
const result = yield* attempt(route);
if (result._tag === "Success") return result.success;
// An incompatible server is the same server on every route.
if (result.failure.reason === "unsupported") return yield* result.failure;
}
for (const route of silent) {
const result = yield* attempt(route);
if (result._tag === "Success") return result.success;
if (result.failure.reason === "unsupported") return yield* result.failure;
}
return yield* (
transient ??
blocked ??
new ConnectionTransientError({
reason: "endpoint-unavailable",
detail: `${entry.target.label} did not answer on any saved route.`,
})
);
});
/** @public Service construction is part of the canonical Effect module API. */
export const make = Effect.gen(function* () {
const resolver = yield* ConnectionResolver.ConnectionResolver;
const sessions = yield* RpcSession.RpcSessionFactory;
const httpClient = yield* HttpClient.HttpClient;
const checkRoute = (entry: ConnectionCatalogEntry, route: ConnectionRoute) => {
const httpBaseUrl = routeHttpBaseUrl(route);
if (httpBaseUrl === null) return Effect.succeed<RouteCheck>("unchecked");
// The descriptor is public, so this sends no credential to whatever
// answers at a saved LAN address on a different network.
return fetchRemoteEnvironmentDescriptor({
httpBaseUrl,
timeoutMs: ROUTE_CHECK_TIMEOUT_MS,
}).pipe(
Effect.map((descriptor): RouteCheck =>
descriptor.environmentId === entry.target.environmentId ? "answered" : "silent",
),
Effect.orElseSucceed((): RouteCheck => "silent"),
Effect.provideService(HttpClient.HttpClient, httpClient),
Effect.withSpan("ConnectionDriver.checkRoute", {
attributes: { "connection.target.kind": route.target._tag },
}),
);
};
const connect = Effect.fn("ConnectionDriver.connect")(function* (
entry: ConnectionCatalogEntry,
reportProgress: (progress: ConnectionDriverProgress) => Effect.Effect<void>,
) {
const target = entry.target;
yield* Effect.annotateCurrentSpan({
"connection.environment.id": target.environmentId,
"connection.target.kind": target._tag,
"connection.route.count": connectionRoutes(entry).length,
});
yield* reportProgress({ stage: "preparing" });
return yield* connectOverRoutes(
entry,
(route) => checkRoute(entry, route),
Effect.fnUntraced(function* (route) {
const prepared = yield* resolver.prepare(routeEntry(entry, route));
yield* reportProgress({ stage: "opening", prepared });
const session = yield* sessions.connect(prepared);
yield* reportProgress({ stage: "synchronizing", prepared });
yield* session.ready;
return { prepared, session } satisfies EnvironmentConnectionLease;
}),
);
});
const preflight = (entry: ConnectionCatalogEntry, route: ConnectionRoute) =>
checkRoute(entry, route).pipe(
Effect.flatMap((check) =>
check === "answered"
? resolver.prepare(routeEntry(entry, route)).pipe(
Effect.as(true),
Effect.orElseSucceed(() => false),
)
: Effect.succeed(false),
),
Effect.withSpan("ConnectionDriver.preflight"),
);
return ConnectionDriver.of({ connect, checkRoute, preflight });
});
export const layer = Layer.effect(ConnectionDriver, make);