Compare commits

..

1 Commits

Author SHA1 Message Date
Kit Langton e0608e74ac refactor(core): use Latch for plugin supervisor ready gate 2026-08-19 23:44:58 -04:00
5 changed files with 70 additions and 143 deletions
+5 -5
View File
@@ -3,7 +3,7 @@ export { Service, type Interface } from "./supervisor-service.js"
import type { Plugin as PluginDefinition } from "@opencode-ai/plugin/effect/plugin"
import { Event } from "@opencode-ai/schema/config"
import { Cause, Deferred, Effect, Layer, Schema, Stream } from "effect"
import { Cause, Effect, Latch, Layer, Schema, Stream } from "effect"
import path from "path"
import { pathToFileURL } from "url"
import { ConfigPluginSource } from "../config/plugin/source.js"
@@ -137,7 +137,7 @@ export const layer = Layer.effect(
const sdk = yield* SdkPlugins.Service
const sources = yield* ConfigPluginSource.Service
const bus = yield* Bus.Service
const ready = { current: yield* Deferred.make<void>() }
const ready = yield* Latch.make()
let observed = 0
const activate = Effect.fn("PluginSupervisor.activate")(function* () {
@@ -164,7 +164,7 @@ export const layer = Layer.effect(
Stream.mapEffect(() =>
Effect.gen(function* () {
observed++
if (yield* Deferred.isDone(ready.current)) ready.current = yield* Deferred.make<void>()
yield* ready.close
return observed
}),
),
@@ -176,12 +176,12 @@ export const layer = Layer.effect(
Stream.runForEach((target) =>
Effect.gen(function* () {
yield* activate()
if (observed === target) yield* Deferred.succeed(ready.current, undefined)
if (observed === target) yield* ready.open
}).pipe(Effect.catchCause((cause) => Effect.logError("failed to reload plugins", { cause }))),
),
Effect.forkScoped({ startImmediately: true }),
)
return Service.of({ flush: Effect.suspend(() => Deferred.await(ready.current)) })
return Service.of({ flush: ready.await })
}),
)
+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()
}
})