forked from pingdotgg/t3code
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathCodexChatGptHandoff.ts
More file actions
63 lines (62 loc) · 2.67 KB
/
Copy pathCodexChatGptHandoff.ts
File metadata and controls
63 lines (62 loc) · 2.67 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
import type { ChatGptHandoffInput, ChatGptHandoffState } from "@t3tools/contracts";
import { ProviderSetupError } from "@t3tools/contracts";
import * as Effect from "effect/Effect";
import * as Option from "effect/Option";
import * as Stream from "effect/Stream";
import { FetchHttpClient } from "effect/unstable/http";
import * as ServerSecretStore from "../auth/ServerSecretStore.ts";
import * as ServerEnvironment from "../environment/ServerEnvironment.ts";
import { makeCodexChatGptAuth } from "./CodexChatGptAuth.ts";
// The primary owns OAuth for this stream; the destination owns the refresh session.
export function subscribeChatGptHandoff(
input: ChatGptHandoffInput,
owner: string,
endpoints: { discoveryUrl: string; resource: string } | undefined = undefined,
) {
return Stream.unwrap(
Effect.gen(function* () {
const bytes = new Map<string, Uint8Array>();
yield* Effect.addFinalizer(() => Effect.sync(() => bytes.clear()));
const store = ServerSecretStore.ServerSecretStore.of({
get: (key) => Effect.sync(() => Option.fromUndefinedOr(bytes.get(key))),
set: (key, value) =>
Effect.sync(() => {
bytes.set(key, value);
}),
remove: (key) =>
Effect.sync(() => {
bytes.delete(key);
}),
create: () => Effect.die("Handoff does not create persistent secrets."),
getOrCreateRandom: () => Effect.die("Handoff uses the destination environment identity."),
});
const auth = yield* makeCodexChatGptAuth({
...endpoints,
instanceId: input.instanceId,
reconnectProfile: input.profile,
telemetryFlow: "primary_handoff",
defaultReturnUrl: input.returnUrl,
}).pipe(
Effect.provideService(ServerSecretStore.ServerSecretStore, store),
Effect.provideService(
ServerEnvironment.ServerEnvironmentIdentity,
ServerEnvironment.ServerEnvironmentIdentity.of({
getEnvironmentId: Effect.succeed(input.environmentId),
}),
),
Effect.provide(FetchHttpClient.layer),
);
yield* auth.controller.start(owner, Effect.void, "chatgpt", input.returnUrl, "server");
return auth.controller.subscribe(owner).pipe(
Stream.takeUntil((state) => ["succeeded", "failed", "cancelled"].includes(state.phase)),
Stream.mapEffect((state): Effect.Effect<ChatGptHandoffState, ProviderSetupError> =>
state.phase === "succeeded"
? auth.exportProfile.pipe(
Effect.map((profile) => ({ phase: "finished" as const, profile })),
)
: Effect.succeed({ phase: "auth" as const, state }),
),
);
}),
);
}