Compare commits

..

2 Commits

Author SHA1 Message Date
Kit Langton 7c1075b586 refactor(server): use Latch for shutdown gate 2026-08-19 23:45:29 -04:00
Kit Langton 20ff543ff2 fix(tui): honest reconnect overlay copy without a managed service (#43561) 2026-08-20 03:33:57 +00:00
7 changed files with 78 additions and 147 deletions
+5 -5
View File
@@ -3,7 +3,7 @@ export * as ServerProcess from "./process"
import { NodeHttpServer } from "@effect/platform-node"
import { SessionRestart } from "@opencode-ai/core/session/execution/restart"
import { hasPtyConnectTicketURL } from "@opencode-ai/protocol/groups/pty"
import { Cause, Context, Deferred, Effect, Exit, Layer, Option, Ref, Scope } from "effect"
import { Cause, Context, Effect, Exit, Latch, Layer, Option, Ref, Scope } from "effect"
import { HttpMiddleware, HttpRouter, HttpServer, HttpServerRequest, HttpServerResponse } from "effect/unstable/http"
import { createServer } from "node:http"
import { ServerAuth } from "./auth"
@@ -47,7 +47,7 @@ export const start = Effect.fn("ServerProcess.start")(function* <E, R>(
if (!password) return yield* Effect.fail(new Error("Missing server password"))
const hostname = options.hostname ?? "127.0.0.1"
const port = Option.fromNullishOr(options.port)
const shutdown = yield* Deferred.make<void>()
const shutdown = yield* Latch.make()
const status = yield* Status.make()
const bound = yield* listen({ hostname, port })
const application = yield* Ref.make(Option.none<App>())
@@ -61,7 +61,7 @@ export const start = Effect.fn("ServerProcess.start")(function* <E, R>(
)
.pipe(withoutParentSpan)
if (lifecycle)
yield* lifecycle.onListen(bound.http.address, Deferred.succeed(shutdown, undefined).pipe(Effect.asVoid)).pipe(
yield* lifecycle.onListen(bound.http.address, shutdown.open.pipe(Effect.asVoid)).pipe(
Effect.flatMap((cleanup) =>
Effect.addFinalizer(() => Scope.close(bound.scope, Exit.void).pipe(Effect.andThen(cleanup))),
),
@@ -101,7 +101,7 @@ export const start = Effect.fn("ServerProcess.start")(function* <E, R>(
const app = Context.get(context, HttpRouter.HttpRouter).asHttpEffect()
yield* Ref.set(application, Option.some(transform ? transform(app) : app))
yield* status.ready
return { address: bound.http.address, shutdown: Deferred.await(shutdown) }
return { address: bound.http.address, shutdown: shutdown.await }
}).pipe(
Effect.catchCause((cause) => {
if (!lifecycle || Cause.hasInterruptsOnly(cause)) return Effect.failCause(cause)
@@ -119,7 +119,7 @@ export const start = Effect.fn("ServerProcess.start")(function* <E, R>(
}),
)
if (!lifecycle) return yield* boot
return yield* Effect.raceFirst(boot, Deferred.await(shutdown).pipe(Effect.andThen(Effect.interrupt)))
return yield* Effect.raceFirst(boot, shutdown.await.pipe(Effect.andThen(Effect.interrupt)))
})
function listen(options: { readonly hostname: string; readonly port: Option.Option<number> }) {
+1 -1
View File
@@ -1376,7 +1376,7 @@ function App(props: { pair?: DialogPairCredentials }) {
<StartupLoading ready={plugins.ready} />
</Show>
<Show when={showReconnecting()}>
<Reconnecting />
<Reconnecting managed={client.restart !== undefined} />
</Show>
<MigrationOverlay />
<Toast />
+7 -3
View File
@@ -2,7 +2,7 @@ import { RGBA } from "@opentui/core"
import { useTheme } from "../context/theme"
import { Spinner } from "./spinner"
export function Reconnecting() {
export function Reconnecting(props: { managed?: boolean }) {
const theme = useTheme("elevated")
return (
@@ -28,8 +28,12 @@ export function Reconnecting() {
paddingRight={2}
gap={1}
>
<Spinner color={theme.text.default}>Restarting service...</Spinner>
<text fg={theme.text.subdued}>Your session will resume automatically.</text>
<Spinner color={theme.text.default}>{props.managed ? "Restarting service..." : "Connection lost..."}</Spinner>
<text fg={theme.text.subdued}>
{props.managed
? "Your session will resume automatically."
: "Reconnecting to the server automatically."}
</text>
</box>
</box>
)
+51 -13
View File
@@ -12,10 +12,8 @@ import { readJson, writeJsonAtomic } from "../util/persistence"
import {
createModelPreferenceRepository,
cycleModelVariant,
favoriteModels,
modelPreferenceKey,
normalizeModelVariant,
recentModels,
type ModelPreference,
type ModelPreferenceModel,
} from "../model-preference"
@@ -34,6 +32,19 @@ export function parseModel(model: string) {
}
}
export function recentModels(model: ModelPreferenceModel, recent: ModelPreferenceModel[]) {
const seen = new Set<string>()
return [model, ...recent]
.filter((item) => {
const key = modelPreferenceKey(item)
if (seen.has(key)) return false
seen.add(key)
return true
})
.slice(0, 10)
.map((item) => ({ providerID: item.providerID, modelID: item.modelID }))
}
export const { use: useLocal, provider: LocalProvider } = createSimpleContext({
name: "Local",
init: () => {
@@ -140,16 +151,37 @@ export const { use: useLocal, provider: LocalProvider } = createSimpleContext({
const pendingSelectionCommits = new Map<string, string>()
const selectionKey = (value: ModelSelection) =>
`${modelPreferenceKey(value)}:${normalizeModelVariant(value.variant) ?? "default"}`
function applyPreferences(value: ModelPreference) {
batch(() => {
const saveState = {
pending: false,
}
function savePreferences() {
if (!preferences.ready) {
saveState.pending = true
return
}
saveState.pending = false
void repository
.patch({
recent: preferences.recent,
favorite: preferences.favorite,
variant: preferences.variant,
})
.catch(() => undefined)
}
repository
.load()
.then((value) => {
setPreferences("recent", value.recent)
setPreferences("favorite", value.favorite)
setPreferences("variant", value.variant)
setPreferences("ready", true)
})
}
onCleanup(repository.subscribe(applyPreferences))
.catch(() => {})
.finally(() => {
setPreferences("ready", true)
if (saveState.pending) savePreferences()
})
const fallbackModel = createMemo(() => {
if (args.model) {
@@ -357,7 +389,7 @@ export const { use: useLocal, provider: LocalProvider } = createSimpleContext({
if (!next) return
if (!selectModel({ ...next })) return
setPreferences("recent", recentModels(next, preferences.recent))
void repository.addRecent(next).catch(() => undefined)
savePreferences()
},
set(model: { providerID: string; modelID: string }, options?: { recent?: boolean }) {
batch(() => {
@@ -365,7 +397,7 @@ export const { use: useLocal, provider: LocalProvider } = createSimpleContext({
if (!selectModel(model)) return
if (options?.recent) {
setPreferences("recent", recentModels(model, preferences.recent))
void repository.addRecent(model).catch(() => undefined)
savePreferences()
}
})
},
@@ -375,8 +407,14 @@ export const { use: useLocal, provider: LocalProvider } = createSimpleContext({
const exists = preferences.favorite.some(
(x) => x.providerID === model.providerID && x.modelID === model.modelID,
)
setPreferences("favorite", favoriteModels(model, preferences.favorite, !exists))
void repository.setFavorite(model, !exists).catch(() => undefined)
const next = exists
? preferences.favorite.filter((x) => x.providerID !== model.providerID || x.modelID !== model.modelID)
: [model, ...preferences.favorite]
setPreferences(
"favorite",
next.map((x) => ({ providerID: x.providerID, modelID: x.modelID })),
)
savePreferences()
})
},
variant: {
@@ -401,7 +439,7 @@ export const { use: useLocal, provider: LocalProvider } = createSimpleContext({
setSessionDraft(route.data.sessionID, { ...m, variant: normalizeModelVariant(value) })
}
setPreferences("variant", modelPreferenceKey(m), normalizeModelVariant(value))
void repository.saveVariant(m, value).catch(() => undefined)
savePreferences()
},
cycle() {
const variants = this.list()
+11 -77
View File
@@ -1,7 +1,5 @@
import { readJson, writeJsonAtomic } from "./util/persistence"
import { isRecord } from "./util/record"
import { watch } from "node:fs"
import path from "node:path"
export type ModelPreferenceModel = {
providerID: string
@@ -45,24 +43,6 @@ export function modelPreferenceKey(model: ModelPreferenceModel) {
return `${model.providerID}/${model.modelID}`
}
export function recentModels(model: ModelPreferenceModel, recent: ModelPreferenceModel[]) {
const seen = new Set<string>()
return [model, ...recent]
.filter((item) => {
const key = modelPreferenceKey(item)
if (seen.has(key)) return false
seen.add(key)
return true
})
.slice(0, 10)
.map((item) => ({ providerID: item.providerID, modelID: item.modelID }))
}
export function favoriteModels(model: ModelPreferenceModel, favorite: ModelPreferenceModel[], enabled: boolean) {
const current = favorite.filter((item) => modelPreferenceKey(item) !== modelPreferenceKey(model))
return enabled ? [model, ...current] : current
}
export function cycleModelVariant(current: string | undefined, variants: string[]) {
const named = variants.filter((variant) => variant !== "default")
if (named.length === 0) return undefined
@@ -100,78 +80,32 @@ function patch(value: Partial<ModelPreference>) {
}
export function createModelPreferenceRepository(filePath: string) {
let pending = Promise.resolve()
let revision = 0
let watcher: ReturnType<typeof watch> | undefined
let reload: ReturnType<typeof setTimeout> | undefined
const listeners = new Set<(value: ModelPreference) => void>()
const state = {
pending: Promise.resolve(),
}
const read = () =>
readJson<unknown>(filePath)
.then(decodeModelPreference)
.catch(() => decodeModelPreference(undefined))
function update(change: (current: ModelPreference) => Partial<ModelPreference>) {
const result = pending.then(async () => {
const { Flock } = await import("@opencode-ai/util/flock")
return Flock.withLock(
filePath,
async () => {
const current = await read()
const next = { ...current, ...patch(change(preference(current))) }
await writeJsonAtomic(filePath, next)
},
{ dir: path.join(path.dirname(filePath), "locks") },
)
const result = state.pending.then(async () => {
const current = await read()
const next = { ...current, ...patch(change(preference(current))) }
await writeJsonAtomic(filePath, next)
})
pending = result.then(
() => undefined,
() => undefined,
)
state.pending = result.catch(() => undefined)
return result
}
function load() {
return pending.then(read).then(preference)
}
async function refresh() {
const current = ++revision
const value = await load()
if (current !== revision) return
listeners.forEach((listener) => listener(value))
return state.pending.then(read).then(preference)
}
return {
load,
addRecent(model: ModelPreferenceModel) {
return update((current) => ({ recent: recentModels(model, current.recent) }))
},
setFavorite(model: ModelPreferenceModel, enabled: boolean) {
return update((current) => ({ favorite: favoriteModels(model, current.favorite, enabled) }))
},
subscribe(listener: (value: ModelPreference) => void) {
listeners.add(listener)
void refresh()
if (!watcher) {
watcher = watch(path.dirname(filePath), (_event, filename) => {
const changed = filename?.toString()
const name = path.basename(filePath)
if (changed !== undefined && changed !== name && !changed.startsWith(name + ".")) return
clearTimeout(reload)
reload = setTimeout(() => void refresh(), 50)
})
watcher.on("error", () => {
watcher?.close()
watcher = undefined
})
}
return () => {
listeners.delete(listener)
if (listeners.size > 0) return
clearTimeout(reload)
watcher?.close()
watcher = undefined
}
patch(value: Partial<ModelPreference>) {
return update(() => value)
},
async resolveVariant(model: ModelPreferenceModel) {
return (await load()).variant[modelPreferenceKey(model)]
+1 -2
View File
@@ -1,6 +1,5 @@
import { expect, test } from "bun:test"
import { parseModel } from "../../src/context/local"
import { recentModels } from "../../src/model-preference"
import { parseModel, recentModels } from "../../src/context/local"
test("parses model IDs containing slashes", () => {
expect(parseModel("provider/family/model")).toEqual({
+2 -46
View File
@@ -19,7 +19,7 @@ test("repairs known model preferences and preserves unrelated fields", () => {
})
})
test("atomically serializes model preference updates", async () => {
test("atomically serializes patches and variant updates", async () => {
await using tmp = await tmpdir()
const file = path.join(tmp.path, "model.json")
await Bun.write(file, JSON.stringify({ unrelated: "keep", favorite: [], variant: {} }))
@@ -28,7 +28,7 @@ test("atomically serializes model preference updates", async () => {
const anthropic = { providerID: "anthropic", modelID: "claude/sonnet" }
await Promise.all([
repository.addRecent(openai),
repository.patch({ recent: [openai] }),
repository.saveVariant(openai, "high"),
repository.saveVariant(anthropic, "low"),
])
@@ -43,47 +43,3 @@ test("atomically serializes model preference updates", async () => {
expect(await repository.resolveVariant(openai)).toBeUndefined()
expect((await Bun.file(file).json()).variant).toEqual({ "anthropic/claude/sonnet": "low" })
})
test("serializes updates across repositories", async () => {
await using tmp = await tmpdir()
const file = path.join(tmp.path, "model.json")
await Bun.write(file, JSON.stringify({ recent: [], favorite: [], variant: {} }))
const first = createModelPreferenceRepository(file)
const second = createModelPreferenceRepository(file)
const openai = { providerID: "openai", modelID: "gpt-5" }
const anthropic = { providerID: "anthropic", modelID: "claude-sonnet" }
await Promise.all([first.setFavorite(openai, true), second.addRecent(anthropic), second.saveVariant(openai, "high")])
expect(await first.load()).toEqual({
recent: [anthropic],
favorite: [openai],
variant: { "openai/gpt-5": "high" },
})
})
test("subscribes to updates from another repository", async () => {
await using tmp = await tmpdir()
const file = path.join(tmp.path, "model.json")
await Bun.write(file, JSON.stringify({ recent: [], favorite: [], variant: {} }))
const first = createModelPreferenceRepository(file)
const second = createModelPreferenceRepository(file)
const openai = { providerID: "openai", modelID: "gpt-5" }
const changed = Promise.withResolvers<void>()
const unsubscribe = first.subscribe((value) => {
if (value.favorite.some((item) => item.providerID === openai.providerID && item.modelID === openai.modelID))
changed.resolve()
})
try {
await second.setFavorite(openai, true)
await Promise.race([
changed.promise,
Bun.sleep(2_000).then(() => {
throw new Error("timed out waiting for model preference update")
}),
])
} finally {
unsubscribe()
}
})