Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 12 additions & 12 deletions packages/opencode/src/kilocode/background-process/runner.ts
Original file line number Diff line number Diff line change
Expand Up @@ -64,23 +64,23 @@ export namespace BackgroundProcessRunner {
let size = (await file.stat()).size
let queue = Promise.resolve()
const append = (chunk: Buffer) => {
if (!process.stdout.destroyed) process.stdout.write(chunk)
queue = queue.then(async () => {
if (size + chunk.length <= MAX) {
await file.write(chunk)
size += chunk.length
return
} else {
await file.close()
const source = Bun.file(input.log)
const old = size
? Buffer.from(await source.slice(Math.max(0, size - KEEP), size).arrayBuffer())
: Buffer.alloc(0)
const next = Buffer.concat([old, chunk])
const tail = next.subarray(Math.max(0, next.length - KEEP))
await Filesystem.write(input.log, tail, MODE)
file = await open(input.log, "a", MODE)
size = tail.length
}
await file.close()
const source = Bun.file(input.log)
const old = size
? Buffer.from(await source.slice(Math.max(0, size - KEEP), size).arrayBuffer())
: Buffer.alloc(0)
const next = Buffer.concat([old, chunk])
const tail = next.subarray(Math.max(0, next.length - KEEP))
await Filesystem.write(input.log, tail, MODE)
file = await open(input.log, "a", MODE)
size = tail.length
if (!process.stdout.destroyed) process.stdout.write(chunk)
})
}
return {
Expand Down
53 changes: 36 additions & 17 deletions packages/opencode/test/session/prompt.test.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,3 @@
import { NodeFileSystem } from "@effect/platform-node"
import { FetchHttpClient } from "effect/unstable/http"
// kilocode_change start
Expand Down Expand Up @@ -160,21 +160,44 @@
const run = SessionRunState.layer.pipe(Layer.provide(status))
const infra = Layer.mergeAll(NodeFileSystem.layer, CrossSpawnSpawner.defaultLayer)

const processorCreateStarted: Array<() => void> = []
// kilocode_change start
const agent: AgentSvc.Info = {
name: "build",
mode: "primary",
native: true,
permission: Permission.fromConfig({ "*": "allow" }),
model: ref,
options: {},
}
const fastAgents = Layer.mock(AgentSvc.Service)({
get: () => Effect.succeed(agent),
list: () => Effect.succeed([agent]),
defaultInfo: () => Effect.succeed(agent),
defaultAgent: () => Effect.succeed(agent.name),
guardRequirements: () => Effect.void,
})

const processorCreateStarted: Deferred.Deferred<void>[] = []
const blockingProcessor = Layer.succeed(
SessionProcessor.Service,
SessionProcessor.Service.of({
create: () => Effect.sync(() => processorCreateStarted.shift()?.()).pipe(Effect.andThen(Effect.never)),
create: () =>
Effect.gen(function* () {
const started = processorCreateStarted.shift()
if (started) yield* Deferred.succeed(started, undefined).pipe(Effect.ignore)
return yield* Effect.never
}),
}),
)
// kilocode_change end

function makePrompt(input?: { processor?: "blocking" }) {
const deps = Layer.mergeAll(
Session.defaultLayer,
Snapshot.defaultLayer,
LLM.defaultLayer,
Env.defaultLayer,
AgentSvc.defaultLayer,
input?.processor === "blocking" ? fastAgents : AgentSvc.defaultLayer, // kilocode_change
Command.defaultLayer,
Permission.defaultLayer,
Plugin.defaultLayer,
Expand Down Expand Up @@ -357,14 +380,6 @@
},
})

function defer<T>() {
let resolve!: (value: T | PromiseLike<T>) => void
const promise = new Promise<T>((done) => {
resolve = done
})
return { promise, resolve }
}

const succeedVoid = (deferred: Deferred.Deferred<void>) => {
Effect.runSync(Deferred.succeed(deferred, void 0).pipe(Effect.ignore))
}
Expand Down Expand Up @@ -1149,10 +1164,12 @@
parts: [{ type: "text", text: "first" }],
})

const firstCreate = defer<void>()
processorCreateStarted.push(firstCreate.resolve)
// kilocode_change start
const firstCreate = yield* Deferred.make<void>()
processorCreateStarted.push(firstCreate)
const first = yield* prompt.loop({ sessionID: chat.id }).pipe(Effect.forkChild)
yield* Effect.promise(() => firstCreate.promise)
yield* awaitWithTimeout(Deferred.await(firstCreate), "processor.create did not start for first turn")
// kilocode_change end

yield* prompt.cancel(chat.id)
const firstExit = yield* Fiber.await(first)
Expand All @@ -1175,10 +1192,12 @@
parts: [{ type: "text", text: "second" }],
})

const secondCreate = defer<void>()
processorCreateStarted.push(secondCreate.resolve)
// kilocode_change start
const secondCreate = yield* Deferred.make<void>()
processorCreateStarted.push(secondCreate)
const second = yield* prompt.loop({ sessionID: chat.id }).pipe(Effect.forkChild)
yield* Effect.promise(() => secondCreate.promise)
yield* awaitWithTimeout(Deferred.await(secondCreate), "processor.create did not start for second turn")
// kilocode_change end

yield* prompt.cancel(chat.id)
const secondExit = yield* Fiber.await(second)
Expand Down
Loading