feat(core): reload commands from config change feed (#39160)
This commit is contained in:
parent
713658c07b
commit
92807d0bb9
2 changed files with 254 additions and 8 deletions
|
|
@ -27,7 +27,19 @@ export const Plugin = define({
|
|||
)
|
||||
}).pipe(Effect.map((documents) => documents.flat()))
|
||||
})
|
||||
const loaded = { documents: yield* load() }
|
||||
const loaded: { documents: { commands: Config.Info["commands"] }[] } = { documents: [] }
|
||||
const reload = load().pipe(
|
||||
Effect.tap((documents) => Effect.sync(() => (loaded.documents = documents))),
|
||||
Effect.andThen(ctx.command.reload()),
|
||||
)
|
||||
// Subscribe before the initial scan so updates racing the scan still trigger a rebuild.
|
||||
yield* config.changes().pipe(
|
||||
Stream.filterEffect((update) => Effect.map(config.entries(), (entries) => isCommandSource(entries, update.path))),
|
||||
Stream.debounce("100 millis"),
|
||||
Stream.runForEach(() => reload),
|
||||
Effect.forkScoped({ startImmediately: true }),
|
||||
)
|
||||
loaded.documents = yield* load()
|
||||
yield* ctx.command.transform((draft) => {
|
||||
for (const document of loaded.documents) {
|
||||
for (const [name, command] of Object.entries(document.commands ?? {})) {
|
||||
|
|
@ -48,17 +60,25 @@ export const Plugin = define({
|
|||
})
|
||||
yield* ctx.event.subscribe().pipe(
|
||||
Stream.filter((event) => event.type === "config.updated"),
|
||||
Stream.runForEach(() =>
|
||||
load().pipe(
|
||||
Effect.tap((documents) => Effect.sync(() => (loaded.documents = documents))),
|
||||
Effect.andThen(ctx.command.reload()),
|
||||
),
|
||||
),
|
||||
Stream.runForEach(() => reload),
|
||||
Effect.forkScoped({ startImmediately: true }),
|
||||
)
|
||||
}),
|
||||
})
|
||||
|
||||
// Matches anything at or under <root>/{command,commands}. No file-suffix check:
|
||||
// directory-level events such as renames carry no per-file paths.
|
||||
function isCommandSource(entries: Config.Entry[], file: string) {
|
||||
return entries.some(
|
||||
(entry) =>
|
||||
entry.type === "directory" &&
|
||||
["command", "commands"].some((name) => {
|
||||
const directory = path.join(entry.path, name)
|
||||
return file === directory || file.startsWith(directory + path.sep)
|
||||
}),
|
||||
)
|
||||
}
|
||||
|
||||
function loadDirectory(fs: FSUtil.Interface, directory: string) {
|
||||
return Effect.gen(function* () {
|
||||
const files = yield* fs
|
||||
|
|
|
|||
|
|
@ -1,7 +1,8 @@
|
|||
import fs from "fs/promises"
|
||||
import path from "path"
|
||||
import { describe, expect } from "bun:test"
|
||||
import { Effect, PubSub, Schema, Stream } from "effect"
|
||||
import { Effect, Fiber, PubSub, Schema, Stream } from "effect"
|
||||
import * as TestClock from "effect/testing/TestClock"
|
||||
import { Config as ConfigSchema } from "@opencode-ai/schema/config"
|
||||
import { Command } from "@opencode-ai/core/command"
|
||||
import { Agent } from "@opencode-ai/core/agent"
|
||||
|
|
@ -108,4 +109,229 @@ Review files`,
|
|||
),
|
||||
),
|
||||
)
|
||||
|
||||
for (const testCase of sourceCases()) {
|
||||
it.effect(`rebuilds commands when a source file is ${testCase.name}`, () =>
|
||||
Effect.acquireRelease(
|
||||
Effect.promise(() => tmpdir()),
|
||||
(tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]()),
|
||||
).pipe(
|
||||
Effect.flatMap((tmp) =>
|
||||
Effect.gen(function* () {
|
||||
const directory = path.join(tmp.path, "commands")
|
||||
yield* Effect.promise(() => fs.mkdir(directory, { recursive: true }))
|
||||
yield* testCase.prepare(directory)
|
||||
|
||||
const command = yield* Command.Service
|
||||
const bus = yield* Bus.Service
|
||||
const test = yield* Config.Test
|
||||
yield* ConfigCommandPlugin.Plugin.effect(
|
||||
host({
|
||||
command: {
|
||||
list: () => Effect.die("unused command.list"),
|
||||
transform: command.transform,
|
||||
reload: command.reload,
|
||||
},
|
||||
}),
|
||||
)
|
||||
|
||||
// Verify inside the subscription so the update event is a read barrier:
|
||||
// committed state must be visible at event delivery time.
|
||||
let received = 0
|
||||
const changed = yield* bus.subscribe(Command.Event.Updated).pipe(
|
||||
Stream.take(1),
|
||||
Stream.tap(() => Effect.sync(() => received++)),
|
||||
Stream.mapEffect(() => testCase.verify(command)),
|
||||
Stream.runDrain,
|
||||
Effect.forkScoped({ startImmediately: true }),
|
||||
)
|
||||
yield* Effect.yieldNow
|
||||
|
||||
const updates = yield* testCase.mutate(directory)
|
||||
yield* Effect.forEach(updates, (update) => test.emitChange(update), { discard: true })
|
||||
yield* advance(Effect.sync(() => received === 1))
|
||||
yield* Fiber.join(changed)
|
||||
}).pipe(Effect.provide(Config.testLayer([directoryEntry(tmp.path)]))),
|
||||
),
|
||||
),
|
||||
)
|
||||
}
|
||||
|
||||
it.effect("coalesces updates inside the debounce window into one rebuild", () =>
|
||||
Effect.acquireRelease(
|
||||
Effect.promise(() => tmpdir()),
|
||||
(tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]()),
|
||||
).pipe(
|
||||
Effect.flatMap((tmp) =>
|
||||
Effect.gen(function* () {
|
||||
const directory = path.join(tmp.path, "commands")
|
||||
yield* Effect.promise(() => fs.mkdir(directory, { recursive: true }))
|
||||
|
||||
const command = yield* Command.Service
|
||||
const test = yield* Config.Test
|
||||
let reloads = 0
|
||||
yield* ConfigCommandPlugin.Plugin.effect(
|
||||
host({
|
||||
command: {
|
||||
list: () => Effect.die("unused command.list"),
|
||||
transform: command.transform,
|
||||
reload: () => command.reload().pipe(Effect.tap(() => Effect.sync(() => reloads++))),
|
||||
},
|
||||
}),
|
||||
)
|
||||
yield* Effect.promise(() => fs.writeFile(path.join(directory, "review.md"), "Review once"))
|
||||
yield* test.emitChange({ type: "create", path: path.join(directory, "review.md") })
|
||||
yield* test.emitChange({ type: "update", path: path.join(directory, "review.md") })
|
||||
yield* test.emitChange({ type: "update", path: path.join(directory, "review.md") })
|
||||
yield* advance(Effect.sync(() => reloads >= 1))
|
||||
expect(reloads).toBe(1)
|
||||
|
||||
yield* Effect.promise(() => fs.writeFile(path.join(directory, "review.md"), "Review twice"))
|
||||
yield* test.emitChange({ type: "update", path: path.join(directory, "review.md") })
|
||||
yield* advance(Effect.sync(() => reloads >= 2))
|
||||
expect(reloads).toBe(2)
|
||||
expect((yield* command.get("review"))?.template).toBe("Review twice")
|
||||
}).pipe(Effect.provide(Config.testLayer([directoryEntry(tmp.path)]))),
|
||||
),
|
||||
),
|
||||
)
|
||||
|
||||
it.effect("ignores updates outside command source directories", () =>
|
||||
Effect.acquireRelease(
|
||||
Effect.promise(() => tmpdir()),
|
||||
(tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]()),
|
||||
).pipe(
|
||||
Effect.flatMap((tmp) =>
|
||||
Effect.gen(function* () {
|
||||
const directory = path.join(tmp.path, "commands")
|
||||
yield* Effect.promise(() => fs.mkdir(directory, { recursive: true }))
|
||||
|
||||
const command = yield* Command.Service
|
||||
const test = yield* Config.Test
|
||||
let reloads = 0
|
||||
yield* ConfigCommandPlugin.Plugin.effect(
|
||||
host({
|
||||
command: {
|
||||
list: () => Effect.die("unused command.list"),
|
||||
transform: command.transform,
|
||||
reload: () => command.reload().pipe(Effect.tap(() => Effect.sync(() => reloads++))),
|
||||
},
|
||||
}),
|
||||
)
|
||||
|
||||
yield* test.emitChange({ type: "create", path: path.join(tmp.path, "notes", "todo.md") })
|
||||
yield* test.emitChange({ type: "update", path: path.join(tmp.path, "opencode.json") })
|
||||
yield* drain
|
||||
expect(reloads).toBe(0)
|
||||
|
||||
// The feed stays live after unrelated updates.
|
||||
yield* Effect.promise(() => fs.writeFile(path.join(directory, "review.md"), "Review related"))
|
||||
yield* test.emitChange({ type: "create", path: path.join(directory, "review.md") })
|
||||
yield* advance(Effect.sync(() => reloads >= 1))
|
||||
expect((yield* command.get("review"))?.template).toBe("Review related")
|
||||
}).pipe(Effect.provide(Config.testLayer([directoryEntry(tmp.path)]))),
|
||||
),
|
||||
),
|
||||
)
|
||||
})
|
||||
|
||||
function directoryEntry(directory: string) {
|
||||
return new Config.Directory({ type: "directory", path: AbsolutePath.make(directory) })
|
||||
}
|
||||
|
||||
// Defers on a real macrotask so pubsub delivery, fiber hops, and filesystem
|
||||
// work can complete between TestClock advances.
|
||||
const settle = Effect.promise(() => new Promise((resolve) => setTimeout(resolve, 1)))
|
||||
|
||||
// Drives every pending TestClock timer to completion: the plugin's 100ms
|
||||
// stream debounce and State's 500ms reload debounce both register their sleeps
|
||||
// from separate fiber hops that a single adjust can miss, so the loop
|
||||
// alternates real-macrotask settles with adjusts until the condition holds.
|
||||
// Extra adjusts are harmless when nothing is pending.
|
||||
const advance = Effect.fnUntraced(function* (condition: Effect.Effect<boolean>) {
|
||||
for (let attempt = 0; attempt < 100; attempt++) {
|
||||
if (yield* condition) return
|
||||
yield* settle
|
||||
yield* TestClock.adjust("500 millis")
|
||||
}
|
||||
throw new Error("condition never became true")
|
||||
})
|
||||
|
||||
// Advances far enough that any matching update would have rebuilt, so a
|
||||
// zero-rebuild assertion afterwards is meaningful.
|
||||
const drain = Effect.gen(function* () {
|
||||
for (let round = 0; round < 4; round++) {
|
||||
yield* settle
|
||||
yield* TestClock.adjust("500 millis")
|
||||
}
|
||||
yield* settle
|
||||
})
|
||||
|
||||
function sourceCases() {
|
||||
return [
|
||||
{
|
||||
name: "created",
|
||||
prepare: () => Effect.void,
|
||||
mutate: (directory: string) =>
|
||||
Effect.promise(async () => {
|
||||
const file = path.join(directory, "review.md")
|
||||
await fs.writeFile(file, "Review created")
|
||||
return [{ type: "create" as const, path: file }]
|
||||
}),
|
||||
verify: (command: Command.Interface) =>
|
||||
Effect.gen(function* () {
|
||||
expect((yield* command.get("review"))?.template).toBe("Review created")
|
||||
}),
|
||||
},
|
||||
{
|
||||
name: "updated",
|
||||
prepare: (directory: string) =>
|
||||
Effect.promise(() => fs.writeFile(path.join(directory, "review.md"), "Review first")),
|
||||
mutate: (directory: string) =>
|
||||
Effect.promise(async () => {
|
||||
const file = path.join(directory, "review.md")
|
||||
await fs.writeFile(file, "Review updated")
|
||||
return [{ type: "update" as const, path: file }]
|
||||
}),
|
||||
verify: (command: Command.Interface) =>
|
||||
Effect.gen(function* () {
|
||||
expect((yield* command.get("review"))?.template).toBe("Review updated")
|
||||
}),
|
||||
},
|
||||
{
|
||||
name: "renamed",
|
||||
prepare: (directory: string) =>
|
||||
Effect.promise(() => fs.writeFile(path.join(directory, "review.md"), "Review renamed")),
|
||||
mutate: (directory: string) =>
|
||||
Effect.promise(async () => {
|
||||
const previous = path.join(directory, "review.md")
|
||||
const next = path.join(directory, "release.md")
|
||||
await fs.rename(previous, next)
|
||||
return [
|
||||
{ type: "delete" as const, path: previous },
|
||||
{ type: "create" as const, path: next },
|
||||
]
|
||||
}),
|
||||
verify: (command: Command.Interface) =>
|
||||
Effect.gen(function* () {
|
||||
expect(yield* command.get("review")).toBeUndefined()
|
||||
expect((yield* command.get("release"))?.template).toBe("Review renamed")
|
||||
}),
|
||||
},
|
||||
{
|
||||
name: "deleted",
|
||||
prepare: (directory: string) =>
|
||||
Effect.promise(() => fs.writeFile(path.join(directory, "review.md"), "Review deleted")),
|
||||
mutate: (directory: string) =>
|
||||
Effect.promise(async () => {
|
||||
const file = path.join(directory, "review.md")
|
||||
await fs.unlink(file)
|
||||
return [{ type: "delete" as const, path: file }]
|
||||
}),
|
||||
verify: (command: Command.Interface) =>
|
||||
Effect.gen(function* () {
|
||||
expect(yield* command.get("review")).toBeUndefined()
|
||||
}),
|
||||
},
|
||||
] as const
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue