Compare commits

...
16 changed files with 577 additions and 83 deletions
+11 -9
View File
@@ -91,8 +91,7 @@ const layer = Layer.effect(
})
const selectable = (agent: Info | undefined) =>
agent && agent.mode !== "subagent" && !agent.hidden ? agent : undefined
const selectedDefault = () => {
const data = state.get()
const selectedDefault = (data: Data) => {
const configured = data.default ? selectable(data.agents.get(data.default)) : undefined
if (configured) return configured
const build = selectable(data.agents.get(ID.make("build")))
@@ -107,23 +106,26 @@ const layer = Layer.effect(
transform: state.transform,
reload: state.reload,
get: Effect.fn("Agent.get")(function* (id) {
return state.get().agents.get(id)
return (yield* state.read()).agents.get(id)
}),
resolve: Effect.fnUntraced(function* (id) {
if (id !== undefined) return state.get().agents.get(ID.make(id))
return selectedDefault()
const data = yield* state.read()
if (id !== undefined) return data.agents.get(ID.make(id))
return selectedDefault(data)
}),
select: Effect.fn("Agent.select")(function* (id) {
const data = yield* state.read()
if (id !== undefined) {
const selected = ID.make(id)
return { id: selected, info: state.get().agents.get(selected) }
return { id: selected, info: data.agents.get(selected) }
}
const info = selectedDefault()
const info = selectedDefault(data)
return { id: info?.id ?? defaultID, info }
}),
list: Effect.fn("Agent.list")(function* () {
const agents = Array.fromIterable(state.get().agents.values())
const selected = selectedDefault()
const data = yield* state.read()
const agents = Array.fromIterable(data.agents.values())
const selected = selectedDefault(data)
if (!selected) return agents
return [selected, ...agents.filter((agent) => agent.id !== selected.id)]
}),
+26 -13
View File
@@ -144,11 +144,11 @@ const layer = Layer.effect(
provider: {
get: Effect.fn("Catalog.provider.get")(function* (providerID) {
return state.get().providers.get(providerID)?.provider
return (yield* state.read()).providers.get(providerID)?.provider
}),
all: Effect.fn("Catalog.provider.all")(function* () {
return Array.fromIterable(state.get().providers.values()).map((record) => record.provider)
return Array.fromIterable((yield* state.read()).providers.values()).map((record) => record.provider)
}),
available: Effect.fn("Catalog.provider.available")(function* () {
@@ -161,15 +161,16 @@ const layer = Layer.effect(
model: {
get: Effect.fn("Catalog.model.get")(function* (providerID, modelID) {
const record = state.get().providers.get(providerID)
const record = (yield* state.read()).providers.get(providerID)
if (!record) return
const model = record.models.get(modelID)
return model && projectModel(model, record.provider)
}),
all: Effect.fn("Catalog.model.all")(function* () {
const data = yield* state.read()
return pipe(
Array.fromIterable(state.get().providers.values()),
Array.fromIterable(data.providers.values()),
Array.flatMap((record) => {
return Array.fromIterable(record.models.values()).map((model) => projectModel(model, record.provider))
}),
@@ -178,10 +179,18 @@ const layer = Layer.effect(
}),
available: Effect.fn("Catalog.model.available")(function* () {
const providers = new Set((yield* result.provider.available()).map((provider) => provider.id))
const active = new Map((yield* integrations.list()).map((integration) => [integration.id, integration]))
const data = yield* state.read()
const models: Model.Info[] = []
for (const record of state.get().providers.values()) {
if (!providers.has(record.provider.id)) continue
for (const record of data.providers.values()) {
if (
!available(
record.provider,
active.get(record.provider.integrationID ?? Integration.ID.make(record.provider.id)),
)
) {
continue
}
for (const model of record.models.values()) {
if (!model.enabled) continue
models.push(projectModel(model, record.provider))
@@ -194,12 +203,16 @@ const layer = Layer.effect(
}),
default: Effect.fn("Catalog.model.default")(function* () {
const defaultModel = state.get().defaultModel
const data = yield* state.read()
const defaultModel = data.defaultModel
if (defaultModel) {
const provider = yield* result.provider.get(defaultModel.providerID)
if (provider && (yield* result.provider.available()).some((item) => item.id === provider.id)) {
const model = yield* result.model.get(defaultModel.providerID, defaultModel.modelID)
if (model?.enabled) return model
const record = data.providers.get(defaultModel.providerID)
const model = record?.models.get(defaultModel.modelID)
if (record && model?.enabled) {
const integration = yield* integrations.get(
record.provider.integrationID ?? Integration.ID.make(record.provider.id),
)
if (available(record.provider, integration)) return projectModel(model, record.provider)
}
}
@@ -207,7 +220,7 @@ const layer = Layer.effect(
}),
small: Effect.fn("Catalog.model.small")(function* (providerID) {
const record = state.get().providers.get(providerID)
const record = (yield* state.read()).providers.get(providerID)
if (!record) return
const models = pipe(
Array.fromIterable(record.models.values()),
+8 -8
View File
@@ -71,15 +71,15 @@ export const layer = Layer.effect(
return Service.of({
reload: state.reload,
transform: state.transform,
get: Effect.fn("Command.get")((name) =>
Effect.sync(() => {
const definition = state.get().get(name)
return definition ? info(definition) : undefined
}),
),
list: Effect.fn("Command.list")(() => Effect.sync(() => Array.from(state.get().values(), info))),
get: Effect.fn("Command.get")(function* (name) {
const definition = (yield* state.read()).get(name)
return definition ? info(definition) : undefined
}),
list: Effect.fn("Command.list")(function* () {
return Array.from((yield* state.read()).values(), info)
}),
execute: Effect.fn("Command.execute")(function* (input) {
const definition = state.get().get(input.name)
const definition = (yield* state.read()).get(input.name)
if (!definition)
return yield* new NotFoundError({ command: input.name, message: `Command not found: ${input.name}` })
return yield* definition.execute(input.invocation).pipe(
+1 -1
View File
@@ -88,7 +88,7 @@ export const layer = (options?: Options) =>
})
const list = Effect.fn("InstructionDiscovery.list")(function* () {
const current = state.get()
const current = yield* state.read()
if (!current.available) return Instructions.unavailable
return Array.from(current.files.values())
})
+12 -16
View File
@@ -402,9 +402,8 @@ const layer = Layer.effect(
}
yield* Effect.gen(function* () {
const implementation = state
.get()
.integrations.get(attempt.integrationID)
const implementation = (yield* state.read()).integrations
.get(attempt.integrationID)
?.implementations.get(attempt.methodID)
const persistence = yield* Effect.sync(() => attempt.label ?? implementation?.label?.(exit.value)).pipe(
Effect.flatMap((label) =>
@@ -545,7 +544,7 @@ const layer = Layer.effect(
readonly answer?: Form.Answer
readonly label?: string
}) {
const method = state.get().integrations.get(input.integrationID)?.implementations.get(input.methodID)
const method = (yield* state.read()).integrations.get(input.integrationID)?.implementations.get(input.methodID)
if (!method) {
return yield* Effect.die(new Error(`OAuth method not found: ${input.integrationID}/${input.methodID}`))
}
@@ -596,8 +595,7 @@ const layer = Layer.effect(
readonly methodID: MethodID
readonly label?: string
}) {
const method = state
.get()
const method = (yield* state.read())
.integrations.get(input.integrationID)
?.methods.find((method) => method.type === "command" && method.id === input.methodID)
if (!method || method.type !== "command" || !method.command[0]) {
@@ -661,19 +659,19 @@ const layer = Layer.effect(
transform: state.transform,
reload: state.reload,
get: Effect.fn("Integration.get")(function* (id) {
const entry = state.get().integrations.get(id)
const entry = (yield* state.read()).integrations.get(id)
if (!entry) return undefined
return project(entry, resolveConnections(entry, yield* credentials.list(id)))
}),
list: Effect.fn("Integration.list")(function* () {
const saved = Map.groupBy(yield* credentials.all(), (credential) => credential.integrationID)
return Array.from(state.get().integrations.values(), (entry) =>
return Array.from((yield* state.read()).integrations.values(), (entry) =>
project(entry, resolveConnections(entry, saved.get(entry.ref.id) ?? [])),
).toSorted((a, b) => a.name.localeCompare(b.name))
}),
connection: {
active: Effect.fn("Integration.connection.active")(function* (id) {
const entry = state.get().integrations.get(id)
const entry = (yield* state.read()).integrations.get(id)
return resolveConnections(entry, yield* credentials.list(id))[0]
}),
resolve: Effect.fn("Integration.connection.resolve")(function* (connection) {
@@ -684,20 +682,18 @@ const layer = Layer.effect(
const credential = yield* credentials.get(connection.id)
if (!credential) return undefined
if (credential.value.type === "key") return credential.value
const implementation = state
.get()
.integrations.get(credential.integrationID)
?.implementations.get(credential.value.methodID)
if (!implementation?.refresh) return credential.value
const now = yield* Clock.currentTimeMillis
if (credential.value.expires > now + Duration.toMillis(Duration.minutes(5))) return credential.value
const implementation = (yield* state.read()).integrations
.get(credential.integrationID)
?.implementations.get(credential.value.methodID)
if (!implementation?.refresh) return credential.value
const value = yield* authorize(implementation.refresh(credential.value))
yield* credentials.update(credential.id, { value })
return value
}),
key: Effect.fn("Integration.connection.key")(function* (input) {
const method = state
.get()
const method = (yield* state.read())
.integrations.get(input.integrationID)
?.methods.find((method) => method.type === "key")
if (!method) return yield* Effect.die(new Error(`Key method not found: ${input.integrationID}`))
+13 -5
View File
@@ -706,17 +706,18 @@ export const layer = (options?: Options) =>
finalize: reconcile,
})
// Suspend so each await sees current entries; a bare Map iterator is exhausted after one run.
const whenAllReady = Effect.suspend(() =>
Effect.forEach(Array.from(entries.values()), (entry) => entry.startup.await, {
const whenAllReady = Effect.gen(function* () {
yield* state.read()
yield* Effect.forEach(Array.from(entries.values()), (entry) => entry.startup.await, {
concurrency: "unbounded",
discard: true,
}),
)
})
})
return Service.of({
transform: state.transform,
reload: state.reload,
servers: Effect.fn("MCP.servers")(function* () {
yield* state.read()
return Array.from(entries)
.toSorted(([a], [b]) => a.localeCompare(b))
.map(([name, entry]) => new ServerInfo({ name, status: entry.status, integrationID: entry.integrationID }))
@@ -727,6 +728,7 @@ export const layer = (options?: Options) =>
yield* state.reload()
}),
connect: Effect.fn("MCP.connect")(function* (server) {
yield* state.read()
const name = ServerName.make(server)
yield* Effect.gen(function* () {
const target = yield* requireServer(name)
@@ -735,6 +737,7 @@ export const layer = (options?: Options) =>
}).pipe(locks.withLock(name))
}),
disconnect: Effect.fn("MCP.disconnect")(function* (server) {
yield* state.read()
const name = ServerName.make(server)
yield* Effect.gen(function* () {
const target = yield* requireServer(name)
@@ -744,6 +747,7 @@ export const layer = (options?: Options) =>
}).pipe(locks.withLock(name))
}),
remove: Effect.fn("MCP.remove")(function* (server) {
yield* state.read()
const name = ServerName.make(server)
yield* requireServer(name)
overrides.set(name, false)
@@ -756,6 +760,7 @@ export const layer = (options?: Options) =>
.toSorted((a, b) => a.server.localeCompare(b.server) || a.name.localeCompare(b.name))
}),
callTool: Effect.fn("MCP.callTool")(function* (input) {
yield* state.read()
const target = yield* requireServer(input.server)
yield* target.entry.startup.await
if (!target.entry.client)
@@ -790,11 +795,13 @@ export const layer = (options?: Options) =>
.toSorted((a, b) => a.server.localeCompare(b.server))
}),
prompts: Effect.fn("MCP.prompts")(function* () {
yield* state.read()
return Array.from(entries.values())
.flatMap((entry) => entry.prompts ?? [])
.toSorted((a, b) => a.server.localeCompare(b.server) || a.name.localeCompare(b.name))
}),
prompt: Effect.fn("MCP.prompt")(function* (input) {
yield* state.read()
const target = yield* requireServer(input.server)
yield* target.entry.startup.await
if (!target.entry.client) return undefined
@@ -849,6 +856,7 @@ export const layer = (options?: Options) =>
})
}),
readResource: Effect.fn("MCP.readResource")(function* (input) {
yield* state.read()
const target = yield* requireServer(input.server)
yield* target.entry.startup.await
if (!target.entry.client) return undefined
+1
View File
@@ -111,6 +111,7 @@ const layer = Layer.effect(
transform: state.transform,
reload: state.reload,
list: Effect.fn("Reference.list")(function* () {
yield* state.read()
return Array.from(materialized.values())
}),
})
+2 -2
View File
@@ -116,10 +116,10 @@ const layer = Layer.effect(
transform: state.transform,
reload: state.reload,
get: Effect.fn("Skill.get")(function* (id) {
return state.get().skills.get(id)
return (yield* state.read()).skills.get(id)
}),
list: Effect.fn("Skill.list")(function* () {
return Array.from(state.get().skills.values())
return Array.from((yield* state.read()).skills.values())
}),
})
}),
+113 -22
View File
@@ -1,6 +1,6 @@
export * as State from "./state.js"
import { Clock, Context, Deferred, Effect, Scope, Semaphore } from "effect"
import { Clock, Context, Deferred, Effect, Exit, Scope, Semaphore } from "effect"
/**
* A replayable transform applied to a draft during reload.
@@ -34,24 +34,46 @@ type Batch = {
active: boolean
readonly flush: boolean
readonly reloads: Set<Reload>
readonly done: Deferred.Deferred<void>
}
type Commit = {
active: boolean
}
const CurrentBatch = Context.Reference<Batch | undefined>("@opencode/State/CurrentBatch", {
defaultValue: () => undefined,
})
const CurrentCommit = Context.Reference<Commit | undefined>("@opencode/State/CurrentCommit", {
defaultValue: () => undefined,
})
const reloadDebounce = 500
/** flush: false is terminal teardown: states whose transforms are removed stop rebuilding, including pending reloads. */
export function batch<A, E, R>(effect: Effect.Effect<A, E, R>, options: { readonly flush?: boolean } = {}) {
return Effect.gen(function* () {
const current = yield* CurrentBatch
if (current?.active && options.flush !== false) return yield* effect
const batch: Batch = { active: true, flush: options.flush !== false, reloads: new Set() }
const exit = yield* effect.pipe(Effect.provideService(CurrentBatch, batch), Effect.exit)
batch.active = false
if (batch.flush) yield* Effect.forEach(batch.reloads, (reload) => reload(), { discard: true })
return yield* exit
})
return Effect.uninterruptibleMask((restore) =>
Effect.gen(function* () {
const current = yield* CurrentBatch
if (current?.active && options.flush !== false) return yield* restore(effect)
const batch: Batch = {
active: true,
flush: options.flush !== false,
reloads: new Set(),
done: Deferred.makeUnsafe<void>(),
}
const exit = yield* restore(effect.pipe(Effect.provideService(CurrentBatch, batch))).pipe(Effect.exit)
batch.active = false
const reloaded = batch.flush
? yield* Effect.forEach(batch.reloads, (reload) => reload().pipe(Effect.exit)).pipe(
Effect.provideService(CurrentBatch, batch),
)
: []
yield* Deferred.succeed(batch.done, undefined)
const failure = reloaded.find(Exit.isFailure)
if (failure) return yield* failure
return yield* exit
}),
)
}
export const inherit = Effect.fnUntraced(function* () {
@@ -75,21 +97,38 @@ export interface Options<State, DraftApi> {
export interface Interface<State, DraftApi> extends Transformable<DraftApi> {
readonly get: () => State
readonly read: () => Effect.Effect<State>
}
export function create<State, DraftApi>(options: Options<State, DraftApi>): Interface<State, DraftApi> {
let state = options.initial()
let transforms: { run: TransformCallback<DraftApi> }[] = []
let generation = 0
let reloadedGeneration = 0
let reloading: Deferred.Deferred<void> | undefined
let requestedAt = 0
let running = false
let closed = false
let committing = false
let waiters: { generation: number; done: Deferred.Deferred<void> }[] = []
const semaphore = Semaphore.makeUnsafe(1)
const batches = new Set<Batch>()
const commit = Effect.fn("State.commit")(function* (next: State) {
state = next
if (options.finalize) yield* options.finalize(options.draft(next))
if (options.finalize) {
const current: Commit = { active: true }
committing = true
yield* options.finalize(options.draft(next)).pipe(
Effect.provideService(CurrentCommit, current),
Effect.ensuring(
Effect.sync(() => {
current.active = false
committing = false
}),
),
)
}
})
const materialize = Effect.fnUntraced(function* () {
@@ -104,7 +143,50 @@ export function create<State, DraftApi>(options: Options<State, DraftApi>): Inte
yield* commit(next)
})
const materializeReload = () => semaphore.withPermit(materialize())
const materializeCurrent = Effect.fnUntraced(function* () {
const target = generation
const exit = yield* materialize().pipe(Effect.exit)
if (Exit.isSuccess(exit)) reloadedGeneration = target
for (const batch of batches) {
if (!batch.active) batches.delete(batch)
}
if (waiters.length) {
const completed = waiters.filter((waiter) => waiter.generation <= target)
waiters = waiters.filter((waiter) => waiter.generation > target)
for (const waiter of completed) Deferred.doneUnsafe(waiter.done, exit)
}
return yield* exit
})
const materializeReload = () => semaphore.withPermit(materializeCurrent())
const materializePending = (): Effect.Effect<void> =>
Effect.suspend(() => {
if (reloadedGeneration >= generation) return Effect.void
if (reloading) return Deferred.await(reloading)
const done = Deferred.makeUnsafe<void>()
reloading = done
return Effect.gen(function* () {
yield* Effect.forEach(batches, (batch) => Deferred.await(batch.done), { discard: true }).pipe(
Effect.andThen(
semaphore.withPermit(
Effect.suspend(() => {
if (reloadedGeneration >= generation) return Effect.void
return materializeCurrent()
}),
),
),
Effect.exit,
Effect.flatMap((exit) =>
Effect.sync(() => {
reloading = undefined
Deferred.doneUnsafe(done, exit)
}),
),
Effect.forkDetach,
)
return yield* Deferred.await(done)
})
})
const rebuild = (): Effect.Effect<void> =>
Effect.gen(function* () {
@@ -112,15 +194,13 @@ export function create<State, DraftApi>(options: Options<State, DraftApi>): Inte
const remaining = requestedAt + reloadDebounce - clock.currentTimeMillisUnsafe()
if (remaining > 0) yield* Effect.sleep(remaining)
if (clock.currentTimeMillisUnsafe() < requestedAt + reloadDebounce) return yield* rebuild()
if (reloadedGeneration >= generation || waiters.length === 0) {
running = false
return
}
const target = generation
const exit = yield* materializeReload().pipe(Effect.exit)
const completed = waiters.filter((waiter) => waiter.generation <= target)
waiters = waiters.filter((waiter) => waiter.generation > target)
yield* Effect.forEach(completed, (waiter) => Deferred.done(waiter.done, exit), {
concurrency: "unbounded",
discard: true,
})
yield* materializePending().pipe(Effect.exit)
if (generation > target) return yield* rebuild()
running = false
})
@@ -136,11 +216,20 @@ export function create<State, DraftApi>(options: Options<State, DraftApi>): Inte
running = true
yield* rebuild().pipe(Effect.forkDetach)
}
yield* Deferred.await(done)
return yield* Deferred.await(done)
})
return {
get: () => state,
read: Effect.fnUntraced(function* () {
while (reloadedGeneration < generation) {
if ((yield* CurrentCommit)?.active && committing) return state
const batch = yield* CurrentBatch
if (batch?.active || (batch && batches.has(batch))) return state
yield* materializePending()
}
return state
}),
transform: Effect.fn("State.transform")(function* (update) {
yield* Effect.annotateCurrentSpan("state", options.name ?? "anonymous")
const scope = yield* Scope.Scope
@@ -162,21 +251,23 @@ export function create<State, DraftApi>(options: Options<State, DraftApi>): Inte
closed = true
return
}
batches.add(batch)
batch.reloads.add(materializeReload)
return
}
yield* materialize()
yield* materializeCurrent()
})
}),
),
)
const batch = yield* CurrentBatch
if (batch?.active) batches.add(batch)
yield* semaphore.withPermit(
Effect.sync(() => {
transforms = [...transforms, transform]
}),
)
yield* Scope.addFinalizer(scope, dispose)
const batch = yield* CurrentBatch
if (batch?.active) batch.reloads.add(materializeReload)
else yield* materializeReload()
return { dispose }
+5 -4
View File
@@ -97,7 +97,7 @@ const layer = Layer.effect(
}
const defaultProvider = Effect.fn("WebSearch.default")(function* () {
const data = state.get()
const data = yield* state.read()
const stored = data.selection === undefined ? yield* kv.get(ProviderKey) : undefined
const decoded = Schema.decodeUnknownOption(Selection)(stored)
if (stored !== undefined && Option.isNone(decoded)) yield* kv.remove(ProviderKey)
@@ -111,8 +111,9 @@ const layer = Layer.effect(
})
const resolve = Effect.fn("WebSearch.resolve")(function* (input: Input) {
const providers = state.get().providers
if (input.providerID) return yield* requireProvider(providers, input.providerID)
if (input.providerID) {
return yield* requireProvider((yield* state.read()).providers, input.providerID)
}
const provider = yield* defaultProvider()
if (!provider) return yield* new ProviderRequiredError()
return provider
@@ -122,7 +123,7 @@ const layer = Layer.effect(
transform: state.transform,
reload: state.reload,
providers: Effect.fn("WebSearch.providers")(function* () {
return Array.from(state.get().providers.values(), (provider) => ({
return Array.from((yield* state.read()).providers.values(), (provider) => ({
id: provider.id,
name: provider.name,
})).toSorted((a, b) => a.name.localeCompare(b.name))
+40
View File
@@ -54,6 +54,46 @@ describe("Skill", () => {
}),
)
it.effect("exposes reloaded skills to synchronous update listeners", () =>
Effect.gen(function* () {
const skill = yield* Skill.Service
const bus = yield* Bus.Service
let description = "Initial"
yield* skill.transform((draft) => draft.add(info("review", description)))
const observed: Skill.Info[][] = []
const unsubscribe = yield* bus.listen((event) =>
event.type === Skill.Event.Updated.type
? skill.list().pipe(
Effect.map((skills) => observed.push(skills)),
Effect.asVoid,
)
: Effect.void,
)
yield* Effect.addFinalizer(() => unsubscribe)
description = "Updated"
const reload = yield* skill.reload().pipe(Effect.forkChild({ startImmediately: true }))
expect(yield* skill.list()).toEqual([info("review", "Updated")])
yield* Fiber.join(reload)
expect(observed).toEqual([[info("review", "Updated")]])
}),
)
it.effect("reloads skills before a direct lookup", () =>
Effect.gen(function* () {
const skill = yield* Skill.Service
let description = "Initial"
yield* skill.transform((draft) => draft.add(info("review", description)))
description = "Updated"
const reload = yield* skill.reload().pipe(Effect.forkChild({ startImmediately: true }))
expect(yield* skill.get(Skill.ID.make("review"))).toEqual(info("review", "Updated"))
yield* Fiber.join(reload)
}),
)
it.effect("restores earlier values when an updating transform is disposed", () =>
Effect.gen(function* () {
const skill = yield* Skill.Service
+331
View File
@@ -76,6 +76,249 @@ describe("State", () => {
}),
)
it.effect("reloads immediately on read", () =>
Effect.gen(function* () {
let value = "first"
let finalized = 0
const state = State.create({
initial: () => ({ values: [] as string[] }),
draft: (draft) => ({ add: (item: string) => draft.values.push(item) }),
finalize: () => Effect.sync(() => finalized++),
})
yield* state.transform((editor) => {
editor.add(value)
})
finalized = 0
value = "second"
const reload = yield* state.reload().pipe(Effect.forkChild({ startImmediately: true }))
expect(state.get().values).toEqual(["first"])
expect(finalized).toBe(0)
expect((yield* state.read()).values).toEqual(["second"])
yield* Fiber.join(reload)
expect(finalized).toBe(1)
yield* TestClock.adjust("500 millis")
expect(finalized).toBe(1)
}),
)
it.effect("single-flights concurrent reads", () =>
Effect.gen(function* () {
const rebuilding = yield* Deferred.make<void>()
const release = yield* Deferred.make<void>()
let block = false
let finalized = 0
const state = State.create({
initial: () => ({ values: [] as string[] }),
draft: (draft) => ({ add: (item: string) => draft.values.push(item) }),
finalize: () =>
Effect.gen(function* () {
finalized++
if (!block) return
yield* Deferred.succeed(rebuilding, undefined)
yield* Deferred.await(release)
}),
})
yield* state.transform((editor) => editor.add("value"))
finalized = 0
block = true
const reload = yield* state.reload().pipe(Effect.forkChild({ startImmediately: true }))
const first = yield* state.read().pipe(Effect.forkChild({ startImmediately: true }))
yield* Deferred.await(rebuilding)
const second = yield* state.read().pipe(Effect.forkChild({ startImmediately: true }))
yield* Deferred.succeed(release, undefined)
yield* Fiber.join(first)
yield* Fiber.join(second)
yield* Fiber.join(reload)
expect(finalized).toBe(1)
}),
)
it.effect("includes reloads requested during a read", () =>
Effect.gen(function* () {
const rebuilding = yield* Deferred.make<void>()
const release = yield* Deferred.make<void>()
let value = "first"
let block = false
let finalized = 0
const state = State.create({
initial: () => ({ values: [] as string[] }),
draft: (draft) => ({ add: (item: string) => draft.values.push(item) }),
finalize: () =>
Effect.gen(function* () {
finalized++
if (!block) return
block = false
yield* Deferred.succeed(rebuilding, undefined)
yield* Deferred.await(release)
}),
})
yield* state.transform((editor) => editor.add(value))
finalized = 0
block = true
value = "second"
const first = yield* state.reload().pipe(Effect.forkChild({ startImmediately: true }))
const read = yield* state.read().pipe(Effect.forkChild({ startImmediately: true }))
yield* Deferred.await(rebuilding)
value = "third"
const second = yield* state.reload().pipe(Effect.forkChild({ startImmediately: true }))
yield* Deferred.succeed(release, undefined)
expect((yield* Fiber.join(read)).values).toEqual(["third"])
yield* Fiber.join(first)
yield* Fiber.join(second)
expect(finalized).toBe(2)
}),
)
it.effect("shares a failed reload between concurrent reads", () =>
Effect.gen(function* () {
const rebuilding = yield* Deferred.make<void>()
const release = yield* Deferred.make<void>()
let fail = false
let finalized = 0
const state = State.create({
initial: () => ({ values: [] as string[] }),
draft: (draft) => ({ add: (item: string) => draft.values.push(item) }),
finalize: () =>
Effect.gen(function* () {
finalized++
if (!fail) return
yield* Deferred.succeed(rebuilding, undefined)
yield* Deferred.await(release)
return yield* Effect.die("failed")
}),
})
yield* state.transform((editor) => editor.add("value"))
finalized = 0
fail = true
const reload = yield* state.reload().pipe(Effect.exit, Effect.forkChild({ startImmediately: true }))
const first = yield* state.read().pipe(Effect.exit, Effect.forkChild({ startImmediately: true }))
yield* Deferred.await(rebuilding)
const second = yield* state.read().pipe(Effect.exit, Effect.forkChild({ startImmediately: true }))
yield* Deferred.succeed(release, undefined)
expect(Exit.isFailure(yield* Fiber.join(first))).toBeTrue()
expect(Exit.isFailure(yield* Fiber.join(second))).toBeTrue()
expect(Exit.isFailure(yield* Fiber.join(reload))).toBeTrue()
expect(finalized).toBe(1)
fail = false
yield* state.read()
expect(finalized).toBe(2)
}),
)
it.effect("allows nested finalizers to read their committed states", () =>
Effect.gen(function* () {
let observe = false
const first: State.Interface<{ values: string[] }, { add: (item: string) => void }> = State.create({
initial: () => ({ values: [] as string[] }),
draft: (draft) => ({ add: (item: string) => draft.values.push(item) }),
finalize: () => (observe ? second.read().pipe(Effect.asVoid) : Effect.void),
})
const second: State.Interface<{ values: string[] }, { add: (item: string) => void }> = State.create({
initial: () => ({ values: [] as string[] }),
draft: (draft) => ({ add: (item: string) => draft.values.push(item) }),
finalize: () => (observe ? first.read().pipe(Effect.asVoid) : Effect.void),
})
yield* first.transform((draft) => draft.add("first"))
yield* second.transform((draft) => draft.add("second"))
observe = true
const firstReload = yield* first.reload().pipe(Effect.forkChild({ startImmediately: true }))
const secondReload = yield* second.reload().pipe(Effect.forkChild({ startImmediately: true }))
expect((yield* first.read()).values).toEqual(["first"])
yield* Fiber.join(firstReload)
yield* Fiber.join(secondReload)
}),
)
it.effect("allows concurrent finalizers to read committed states", () =>
Effect.gen(function* () {
const firstEntered = yield* Deferred.make<void>()
const secondEntered = yield* Deferred.make<void>()
let observe = false
const first: State.Interface<{ values: string[] }, { add: (item: string) => void }> = State.create({
initial: () => ({ values: [] as string[] }),
draft: (draft) => ({ add: (item: string) => draft.values.push(item) }),
finalize: () =>
observe
? Deferred.succeed(firstEntered, undefined).pipe(
Effect.andThen(Deferred.await(secondEntered)),
Effect.andThen(second.read()),
Effect.asVoid,
)
: Effect.void,
})
const second: State.Interface<{ values: string[] }, { add: (item: string) => void }> = State.create({
initial: () => ({ values: [] as string[] }),
draft: (draft) => ({ add: (item: string) => draft.values.push(item) }),
finalize: () =>
observe
? Deferred.succeed(secondEntered, undefined).pipe(
Effect.andThen(Deferred.await(firstEntered)),
Effect.andThen(first.read()),
Effect.asVoid,
)
: Effect.void,
})
yield* first.transform((draft) => draft.add("first"))
yield* second.transform((draft) => draft.add("second"))
observe = true
const firstReload = yield* first.reload().pipe(Effect.forkChild({ startImmediately: true }))
const secondReload = yield* second.reload().pipe(Effect.forkChild({ startImmediately: true }))
const firstRead = yield* first.read().pipe(Effect.forkChild({ startImmediately: true }))
const secondRead = yield* second.read().pipe(Effect.forkChild({ startImmediately: true }))
expect((yield* Fiber.join(firstRead)).values).toEqual(["first"])
expect((yield* Fiber.join(secondRead)).values).toEqual(["second"])
yield* Fiber.join(firstReload)
yield* Fiber.join(secondReload)
}),
)
it.effect("expires finalizer read access before later reloads", () =>
Effect.gen(function* () {
const release = yield* Deferred.make<void>()
const observed = yield* Deferred.make<string>()
let value = "first"
let observe = false
const state: State.Interface<{ values: string[] }, { add: (item: string) => void }> = State.create({
initial: () => ({ values: [] as string[] }),
draft: (draft) => ({ add: (item: string) => draft.values.push(item) }),
finalize: () =>
observe
? Deferred.await(release).pipe(
Effect.andThen(state.read()),
Effect.flatMap((current) => Deferred.succeed(observed, current.values[0] ?? "")),
Effect.forkDetach,
Effect.asVoid,
)
: Effect.void,
})
yield* state.transform((draft) => draft.add(value))
observe = true
const first = yield* state.reload().pipe(Effect.forkChild({ startImmediately: true }))
yield* state.read()
yield* Fiber.join(first)
observe = false
value = "second"
const second = yield* state.reload().pipe(Effect.forkChild({ startImmediately: true }))
yield* Deferred.succeed(release, undefined)
expect(yield* Deferred.await(observed)).toBe("second")
yield* Fiber.join(second)
}),
)
it.effect("disposes a transform once and rebuilds remaining state", () =>
Effect.gen(function* () {
const state = State.create({
@@ -191,6 +434,94 @@ describe("State", () => {
}),
)
it.effect("does not publish pending reloads inside an active batch", () =>
Effect.gen(function* () {
let finalized = 0
const state = State.create({
initial: () => ({ values: [] as string[] }),
draft: (draft) => ({ add: (item: string) => draft.values.push(item) }),
finalize: () => Effect.sync(() => finalized++),
})
const reload = yield* state.reload().pipe(Effect.forkChild({ startImmediately: true }))
const observe = yield* Deferred.make<void>()
const reader = yield* Deferred.await(observe).pipe(
Effect.andThen(state.read()),
Effect.forkChild({ startImmediately: true }),
)
yield* State.batch(
Effect.gen(function* () {
yield* state.transform((draft) => draft.add("first"))
expect((yield* state.read()).values).toEqual([])
yield* Deferred.succeed(observe, undefined)
yield* Effect.yieldNow
expect(finalized).toBe(0)
yield* state.transform((draft) => draft.add("second"))
}),
)
yield* Fiber.join(reload)
expect((yield* Fiber.join(reader)).values).toEqual(["first", "second"])
expect(finalized).toBe(1)
}),
)
it.effect("releases pending reads when a batch is interrupted", () =>
Effect.gen(function* () {
const started = yield* Deferred.make<void>()
const release = yield* Deferred.make<void>()
const state = State.create({
initial: () => ({ values: [] as string[] }),
draft: (draft) => ({ add: (item: string) => draft.values.push(item) }),
})
const reload = yield* state.reload().pipe(Effect.forkChild({ startImmediately: true }))
const batch = yield* State.batch(
Effect.gen(function* () {
yield* state.transform((draft) => draft.add("value"))
yield* Deferred.succeed(started, undefined)
yield* Deferred.await(release)
}),
).pipe(Effect.forkChild({ startImmediately: true }))
yield* Deferred.await(started)
const reader = yield* state.read().pipe(Effect.forkChild({ startImmediately: true }))
yield* Fiber.interrupt(batch)
expect((yield* Fiber.join(reader)).values).toEqual(["value"])
yield* Fiber.join(reload)
}),
)
it.effect("flushes remaining domains after a batched reload fails", () =>
Effect.gen(function* () {
let failing = true
const first = State.create({
initial: () => ({ values: [] as string[] }),
draft: (draft) => ({ add: (item: string) => draft.values.push(item) }),
finalize: () => (failing ? Effect.die("failed") : Effect.void),
})
const second = State.create({
initial: () => ({ values: [] as string[] }),
draft: (draft) => ({ add: (item: string) => draft.values.push(item) }),
})
const firstReload = yield* first.reload().pipe(Effect.exit, Effect.forkChild({ startImmediately: true }))
const secondReload = yield* second.reload().pipe(Effect.forkChild({ startImmediately: true }))
const result = yield* State.batch(
Effect.gen(function* () {
yield* first.transform((draft) => draft.add("first"))
yield* second.transform((draft) => draft.add("second"))
}),
).pipe(Effect.exit)
expect(Exit.isFailure(result)).toBeTrue()
expect(Exit.isFailure(yield* Fiber.join(firstReload))).toBeTrue()
yield* Fiber.join(secondReload)
expect((yield* second.read()).values).toEqual(["second"])
failing = false
}),
)
it.effect("debounces reload bursts", () =>
Effect.gen(function* () {
let finalized = 0
+4
View File
@@ -125,6 +125,10 @@ data = await loadCatalog()
await ctx.catalog.reload()
```
`reload()` immediately marks the domain stale and schedules a debounced reload that replays every active transform in
order. A concurrent domain read immediately performs the pending reload; otherwise, it runs in the background after the
debounce window. In either case, `reload()` resolves after the actual reload completes.
Available reload operations are:
```ts
+3 -1
View File
@@ -108,7 +108,9 @@ data = yield * loadCatalog()
yield * ctx.catalog.reload()
```
Reload belongs to the domain, not an individual registration. `ctx.catalog.reload()` reruns every active catalog transform and publishes the rebuilt catalog.
`reload()` immediately marks the domain stale and schedules a debounced reload that replays every active transform in
order. A concurrent domain read immediately performs the pending reload; otherwise, it runs in the background after the
debounce window. In either case, `reload()` resolves after the actual reload completes.
Available reload operations are:
@@ -209,7 +209,9 @@ effect: (ctx) =>
}),
```
Call `reload` when external state used by a transform changes. Reload replays every transform in order.
Call `reload` when external state used by a transform changes. It immediately marks the domain stale and schedules a
debounced reload that replays every transform in order. A concurrent domain read immediately performs the pending reload;
otherwise, it runs in the background after the debounce window. In either case, `reload` resolves after the actual reload.
```ts title="plugins/models-effect.ts"
effect: (ctx) =>
@@ -177,7 +177,10 @@ export default Plugin.define({
})
```
`reload` replays every catalog transform in order, so the output-price policy still filters the refreshed models.
`reload` immediately marks the catalog stale and schedules a debounced reload that replays every transform in order,
so the output-price policy still filters the refreshed models. A concurrent catalog read immediately performs the pending
reload; otherwise, it runs in the background after the debounce window. In either case, `reload` resolves after the
actual reload.
## API