refactor(core): replace background job service (#34559)
This commit is contained in:
@@ -48,7 +48,7 @@ import { ShareNext } from "@/share/share-next"
|
||||
import { SessionShare } from "@/share/session"
|
||||
import { Npm } from "@opencode-ai/core/npm"
|
||||
import { memoMap } from "@opencode-ai/core/effect/memo-map"
|
||||
import { BackgroundJob } from "@/background/job"
|
||||
import { Job } from "@/job"
|
||||
import { RuntimeFlags } from "@/effect/runtime-flags"
|
||||
import { EventV2Bridge } from "@/event-v2-bridge"
|
||||
|
||||
@@ -74,7 +74,7 @@ export const AppLayer = Layer.mergeAll(
|
||||
Todo.defaultLayer,
|
||||
Session.defaultLayer,
|
||||
SessionStatus.defaultLayer,
|
||||
BackgroundJob.defaultLayer,
|
||||
Job.defaultLayer,
|
||||
RuntimeFlags.defaultLayer,
|
||||
EventV2Bridge.defaultLayer,
|
||||
SessionRunState.defaultLayer,
|
||||
|
||||
@@ -1,32 +1,34 @@
|
||||
import { LayerNode } from "@opencode-ai/core/effect/layer-node"
|
||||
import { BackgroundJob as CoreBackgroundJob } from "@opencode-ai/core/background-job"
|
||||
import { Service, make } from "@opencode-ai/core/job"
|
||||
import { InstanceState } from "@/effect/instance-state"
|
||||
import { Effect, Layer } from "effect"
|
||||
|
||||
export {
|
||||
Service,
|
||||
type ExtendInput,
|
||||
type BackgroundAllInput,
|
||||
type BlockInput,
|
||||
type BlockResult,
|
||||
type Info,
|
||||
type Interface,
|
||||
type StartInput,
|
||||
type Status,
|
||||
type WaitInput,
|
||||
type WaitResult,
|
||||
} from "@opencode-ai/core/background-job"
|
||||
} from "@opencode-ai/core/job"
|
||||
|
||||
/** Keeps the legacy service instance-scoped while sharing the core registry engine. */
|
||||
export const layer = Layer.effect(
|
||||
CoreBackgroundJob.Service,
|
||||
Service,
|
||||
Effect.gen(function* () {
|
||||
const state = yield* InstanceState.make(() => CoreBackgroundJob.make)
|
||||
return CoreBackgroundJob.Service.of({
|
||||
const state = yield* InstanceState.make(() => make)
|
||||
return Service.of({
|
||||
list: () => InstanceState.useEffect(state, (jobs) => jobs.list()),
|
||||
get: (id) => InstanceState.useEffect(state, (jobs) => jobs.get(id)),
|
||||
start: (input) => InstanceState.useEffect(state, (jobs) => jobs.start(input)),
|
||||
extend: (input) => InstanceState.useEffect(state, (jobs) => jobs.extend(input)),
|
||||
wait: (input) => InstanceState.useEffect(state, (jobs) => jobs.wait(input)),
|
||||
waitForPromotion: (id) => InstanceState.useEffect(state, (jobs) => jobs.waitForPromotion(id)),
|
||||
promote: (id) => InstanceState.useEffect(state, (jobs) => jobs.promote(id)),
|
||||
block: (input) => InstanceState.useEffect(state, (jobs) => jobs.block(input)),
|
||||
background: (id) => InstanceState.useEffect(state, (jobs) => jobs.background(id)),
|
||||
backgroundAll: (input) => InstanceState.useEffect(state, (jobs) => jobs.backgroundAll(input)),
|
||||
cancel: (id) => InstanceState.useEffect(state, (jobs) => jobs.cancel(id)),
|
||||
})
|
||||
}),
|
||||
@@ -34,6 +36,6 @@ export const layer = Layer.effect(
|
||||
|
||||
export const defaultLayer = layer
|
||||
|
||||
export const node = LayerNode.make({ service: CoreBackgroundJob.Service, layer, deps: [] })
|
||||
export const node = LayerNode.make({ service: Service, layer, deps: [] })
|
||||
|
||||
export * as BackgroundJob from "./job"
|
||||
export * as Job from "./job"
|
||||
@@ -1,6 +1,6 @@
|
||||
import { Account } from "@/account/account"
|
||||
import { Agent } from "@/agent/agent"
|
||||
import { BackgroundJob } from "@/background/job"
|
||||
import { Job } from "@/job"
|
||||
import { Config } from "@/config/config"
|
||||
import { InstanceState } from "@/effect/instance-state"
|
||||
import { RuntimeFlags } from "@/effect/runtime-flags"
|
||||
@@ -33,7 +33,7 @@ export const experimentalHandlers = HttpApiBuilder.group(InstanceHttpApi, "exper
|
||||
const registry = yield* ToolRegistry.Service
|
||||
const worktreeSvc = yield* Worktree.Service
|
||||
const sessions = yield* Session.Service
|
||||
const background = yield* BackgroundJob.Service
|
||||
const jobs = yield* Job.Service
|
||||
const flags = yield* RuntimeFlags.Service
|
||||
|
||||
const capabilities = Effect.fn("ExperimentalHttpApi.capabilities")(function* () {
|
||||
@@ -159,15 +159,7 @@ export const experimentalHandlers = HttpApiBuilder.group(InstanceHttpApi, "exper
|
||||
params: { sessionID: SessionID }
|
||||
}) {
|
||||
if (!flags.experimentalBackgroundSubagents) return false
|
||||
const jobs = (yield* background.list()).filter(
|
||||
(job) =>
|
||||
job.type === "task" &&
|
||||
job.status === "running" &&
|
||||
job.metadata?.parentSessionId === ctx.params.sessionID &&
|
||||
job.metadata.background !== true,
|
||||
)
|
||||
const promoted = yield* Effect.forEach(jobs, (job) => background.promote(job.id), { concurrency: "unbounded" })
|
||||
return promoted.some((job) => job !== undefined)
|
||||
return (yield* jobs.backgroundAll({ sessionID: ctx.params.sessionID, type: "task" })).length > 0
|
||||
})
|
||||
|
||||
const resource = Effect.fn("ExperimentalHttpApi.resource")(function* () {
|
||||
|
||||
@@ -7,7 +7,7 @@ import * as Observability from "@opencode-ai/core/observability"
|
||||
import { Account } from "@/account/account"
|
||||
import { Agent } from "@/agent/agent"
|
||||
import { Auth } from "@/auth"
|
||||
import { BackgroundJob } from "@/background/job"
|
||||
import { Job } from "@/job"
|
||||
import { Command } from "@/command"
|
||||
import { Config } from "@/config/config"
|
||||
import { Workspace } from "@/control-plane/workspace"
|
||||
@@ -233,7 +233,7 @@ const app = LayerNode.group([
|
||||
Session.node,
|
||||
SessionProjector.node,
|
||||
SessionStatus.node,
|
||||
BackgroundJob.node,
|
||||
Job.node,
|
||||
RuntimeFlags.node,
|
||||
EventV2Bridge.node,
|
||||
SessionRunState.node,
|
||||
|
||||
@@ -2,7 +2,7 @@ import { LayerNode } from "@opencode-ai/core/effect/layer-node"
|
||||
import { InstanceState } from "@/effect/instance-state"
|
||||
import { SessionV1 } from "@opencode-ai/core/v1/session"
|
||||
import { Runner } from "@/effect/runner"
|
||||
import { BackgroundJob } from "@/background/job"
|
||||
import { Job } from "@/job"
|
||||
import { Effect, Latch, Layer, Scope, Context } from "effect"
|
||||
import { Session } from "./session"
|
||||
import { SessionID } from "./schema"
|
||||
@@ -29,7 +29,7 @@ export class Service extends Context.Service<Service, Interface>()("@opencode/Se
|
||||
export const layer = Layer.effect(
|
||||
Service,
|
||||
Effect.gen(function* () {
|
||||
const background = yield* BackgroundJob.Service
|
||||
const jobs = yield* Job.Service
|
||||
const status = yield* SessionStatus.Service
|
||||
|
||||
const state = yield* InstanceState.make(
|
||||
@@ -75,7 +75,7 @@ export const layer = Layer.effect(
|
||||
})
|
||||
|
||||
const cancel = Effect.fn("SessionRunState.cancel")(function* (sessionID: SessionID) {
|
||||
yield* cancelBackgroundJobs(background, sessionID)
|
||||
yield* cancelJobs(jobs, sessionID)
|
||||
const data = yield* InstanceState.get(state)
|
||||
const existing = data.runners.get(sessionID)
|
||||
if (!existing) {
|
||||
@@ -108,31 +108,25 @@ export const layer = Layer.effect(
|
||||
}),
|
||||
)
|
||||
|
||||
export const defaultLayer = layer.pipe(
|
||||
Layer.provide(BackgroundJob.defaultLayer),
|
||||
Layer.provide(SessionStatus.defaultLayer),
|
||||
)
|
||||
export const defaultLayer = layer.pipe(Layer.provide(Job.defaultLayer), Layer.provide(SessionStatus.defaultLayer))
|
||||
|
||||
const cancelBackgroundJobs = Effect.fn("SessionRunState.cancelBackgroundJobs")(function* (
|
||||
background: BackgroundJob.Interface,
|
||||
sessionID: SessionID,
|
||||
) {
|
||||
const jobs = yield* background.list()
|
||||
const cancelJobs = Effect.fn("SessionRunState.cancelJobs")(function* (jobs: Job.Interface, sessionID: SessionID) {
|
||||
const running = yield* jobs.list()
|
||||
const pending = new Set<string>([sessionID])
|
||||
const cancelled = new Set<string>()
|
||||
const matches = (job: BackgroundJob.Info) => {
|
||||
const matches = (job: Job.Info) => {
|
||||
if (job.status !== "running") return false
|
||||
if (cancelled.has(job.id)) return false
|
||||
if (pending.has(job.id)) return true
|
||||
if (typeof job.metadata?.sessionId === "string" && pending.has(job.metadata.sessionId)) return true
|
||||
return typeof job.metadata?.parentSessionId === "string" && pending.has(job.metadata.parentSessionId)
|
||||
}
|
||||
let batch = jobs.filter(matches)
|
||||
let batch = running.filter(matches)
|
||||
while (batch.length > 0) {
|
||||
yield* Effect.forEach(
|
||||
batch,
|
||||
(job) =>
|
||||
background.cancel(job.id).pipe(
|
||||
jobs.cancel(job.id).pipe(
|
||||
Effect.tap(() =>
|
||||
Effect.sync(() => {
|
||||
cancelled.add(job.id)
|
||||
@@ -143,7 +137,7 @@ const cancelBackgroundJobs = Effect.fn("SessionRunState.cancelBackgroundJobs")(f
|
||||
),
|
||||
{ concurrency: "unbounded", discard: true },
|
||||
)
|
||||
batch = jobs.filter(matches)
|
||||
batch = running.filter(matches)
|
||||
}
|
||||
})
|
||||
|
||||
@@ -151,6 +145,6 @@ function busyError(sessionID: SessionID) {
|
||||
return new Session.BusyError({ sessionID })
|
||||
}
|
||||
|
||||
export const node = LayerNode.make({ service: Service, layer: layer, deps: [BackgroundJob.node, SessionStatus.node] })
|
||||
export const node = LayerNode.make({ service: Service, layer: layer, deps: [Job.node, SessionStatus.node] })
|
||||
|
||||
export * as SessionRunState from "./run-state"
|
||||
|
||||
@@ -4,7 +4,7 @@ import { Slug } from "@opencode-ai/core/util/slug"
|
||||
import { SessionV1 } from "@opencode-ai/core/v1/session"
|
||||
import { serviceUse } from "@opencode-ai/core/effect/service-use"
|
||||
import path from "path"
|
||||
import { BackgroundJob } from "@/background/job"
|
||||
import { Job } from "@/job"
|
||||
import { Decimal } from "decimal.js"
|
||||
import type { ProviderMetadata, Usage } from "@opencode-ai/llm"
|
||||
import { InstallationVersion } from "@opencode-ai/core/installation/version"
|
||||
@@ -491,13 +491,13 @@ export type Patch = Omit<Partial<Info>, "time" | "share" | "summary" | "revert"
|
||||
export const layer: Layer.Layer<
|
||||
Service,
|
||||
never,
|
||||
BackgroundJob.Service | RuntimeFlags.Service | Database.Service | EventV2Bridge.Service
|
||||
Job.Service | RuntimeFlags.Service | Database.Service | EventV2Bridge.Service
|
||||
> = Layer.effect(
|
||||
Service,
|
||||
Effect.gen(function* () {
|
||||
const { db } = yield* Database.Service
|
||||
const database = yield* Database.Service
|
||||
const background = yield* BackgroundJob.Service
|
||||
const jobs = yield* Job.Service
|
||||
const events = yield* EventV2Bridge.Service
|
||||
const flags = yield* RuntimeFlags.Service
|
||||
|
||||
@@ -618,7 +618,7 @@ export const layer: Layer.Layer<
|
||||
Effect.catchCause(() => Effect.succeed(false)),
|
||||
)
|
||||
|
||||
if (hasInstance) yield* cancelBackgroundJobs(background, sessionID)
|
||||
if (hasInstance) yield* cancelJobs(jobs, sessionID)
|
||||
const kids = yield* children(sessionID)
|
||||
for (const child of kids) {
|
||||
yield* remove(child.id)
|
||||
@@ -941,7 +941,7 @@ export const layer: Layer.Layer<
|
||||
)
|
||||
|
||||
export const defaultLayer = layer.pipe(
|
||||
Layer.provide(BackgroundJob.defaultLayer),
|
||||
Layer.provide(Job.defaultLayer),
|
||||
Layer.provide(Database.defaultLayer),
|
||||
Layer.provide(EventV2Bridge.defaultLayer),
|
||||
Layer.provide(
|
||||
@@ -953,19 +953,16 @@ export const defaultLayer = layer.pipe(
|
||||
Layer.provide(RuntimeFlags.defaultLayer),
|
||||
)
|
||||
|
||||
const cancelBackgroundJobs = Effect.fn("Session.cancelBackgroundJobs")(function* (
|
||||
background: BackgroundJob.Interface,
|
||||
sessionID: SessionID,
|
||||
) {
|
||||
const jobs = yield* background.list()
|
||||
const cancelJobs = Effect.fn("Session.cancelJobs")(function* (jobs: Job.Interface, sessionID: SessionID) {
|
||||
const running = yield* jobs.list()
|
||||
yield* Effect.forEach(
|
||||
jobs.filter((job) => {
|
||||
running.filter((job) => {
|
||||
if (job.status !== "running") return false
|
||||
if (job.id === sessionID) return true
|
||||
if (job.metadata?.sessionId === sessionID) return true
|
||||
return job.metadata?.parentSessionId === sessionID
|
||||
}),
|
||||
(job) => background.cancel(job.id),
|
||||
(job) => jobs.cancel(job.id),
|
||||
{ concurrency: "unbounded", discard: true },
|
||||
)
|
||||
})
|
||||
@@ -1098,7 +1095,7 @@ export function* listGlobal(input?: {
|
||||
export const node = LayerNode.make({
|
||||
service: Service,
|
||||
layer: layer,
|
||||
deps: [BackgroundJob.node, RuntimeFlags.node, Database.node, EventV2Bridge.node],
|
||||
deps: [Job.node, RuntimeFlags.node, Database.node, EventV2Bridge.node],
|
||||
})
|
||||
|
||||
export * as Session from "./session"
|
||||
|
||||
@@ -48,7 +48,7 @@ import { EventV2Bridge } from "@/event-v2-bridge"
|
||||
import { Agent } from "../agent/agent"
|
||||
import { Skill } from "../skill"
|
||||
import { Permission } from "@/permission"
|
||||
import { BackgroundJob } from "@/background/job"
|
||||
import { Job } from "@/job"
|
||||
import { RuntimeFlags } from "@/effect/runtime-flags"
|
||||
import { ProviderV2 } from "@opencode-ai/core/provider"
|
||||
import { ModelV2 } from "@opencode-ai/core/model"
|
||||
@@ -325,7 +325,7 @@ export const defaultLayer = Layer.suspend(() =>
|
||||
Layer.provide(Skill.defaultLayer),
|
||||
Layer.provide(Agent.defaultLayer),
|
||||
Layer.provide(Session.defaultLayer),
|
||||
Layer.provide(BackgroundJob.defaultLayer),
|
||||
Layer.provide(Job.defaultLayer),
|
||||
Layer.provide(Provider.defaultLayer),
|
||||
Layer.provide(LSP.defaultLayer),
|
||||
Layer.provide(Instruction.defaultLayer),
|
||||
@@ -426,7 +426,7 @@ export const node = LayerNode.make({
|
||||
Agent.node,
|
||||
Skill.node,
|
||||
Session.node,
|
||||
BackgroundJob.node,
|
||||
Job.node,
|
||||
Provider.node,
|
||||
LSP.node,
|
||||
Instruction.node,
|
||||
|
||||
@@ -2,7 +2,7 @@ import * as Tool from "./tool"
|
||||
import DESCRIPTION from "./task.txt"
|
||||
import { ToolJsonSchema } from "./json-schema"
|
||||
import { SessionV1 } from "@opencode-ai/core/v1/session"
|
||||
import { BackgroundJob } from "@/background/job"
|
||||
import { Job } from "@/job"
|
||||
import { Session } from "@/session/session"
|
||||
import { SessionID, MessageID } from "../session/schema"
|
||||
import { MessageV2 } from "../session/message-v2"
|
||||
@@ -33,11 +33,11 @@ const BACKGROUND_STARTED = [
|
||||
"DO NOT sleep, poll for progress, ask the task for status, or duplicate this task's work — avoid working with the same files or topics it is using.",
|
||||
"Work on non-overlapping tasks, or briefly tell the user what you launched and end your response.",
|
||||
].join("\n")
|
||||
const BACKGROUND_UPDATED = [
|
||||
"Additional context sent to the running background task.",
|
||||
const BACKGROUND_ALREADY_RUNNING = [
|
||||
"The task is already working in the background.",
|
||||
"The task is still working in the background. You will be notified automatically when it finishes.",
|
||||
"DO NOT sleep, poll for progress, ask the task for status, or duplicate this task's work — avoid working with the same files or topics it is using.",
|
||||
"Work on non-overlapping tasks, or briefly tell the user what you sent and end your response.",
|
||||
"Work on non-overlapping tasks, or briefly tell the user it is still running and end your response.",
|
||||
].join("\n")
|
||||
|
||||
const BaseParameterFields = {
|
||||
@@ -82,7 +82,7 @@ export const TaskTool = Tool.define(
|
||||
id,
|
||||
Effect.gen(function* () {
|
||||
const agent = yield* Agent.Service
|
||||
const background = yield* BackgroundJob.Service
|
||||
const jobs = yield* Job.Service
|
||||
const config = yield* Config.Service
|
||||
const sessions = yield* Session.Service
|
||||
const scope = yield* Scope.Scope
|
||||
@@ -229,7 +229,7 @@ export const TaskTool = Tool.define(
|
||||
})
|
||||
|
||||
const notify = Effect.fn("TaskTool.notifyBackgroundResult")(function* (jobID: string) {
|
||||
yield* background.wait({ id: jobID }).pipe(
|
||||
yield* jobs.wait({ id: jobID }).pipe(
|
||||
Effect.flatMap((result) => {
|
||||
if (result.info?.status === "completed") return inject("completed", result.info.output ?? "")
|
||||
if (result.info?.status === "error") return inject("error", result.info.error ?? "")
|
||||
@@ -239,7 +239,8 @@ export const TaskTool = Tool.define(
|
||||
)
|
||||
})
|
||||
|
||||
if (yield* background.extend({ id: nextSession.id, run: runTask() })) {
|
||||
const existing = yield* jobs.get(nextSession.id)
|
||||
if (existing?.status === "running") {
|
||||
return {
|
||||
title: params.description,
|
||||
metadata: {
|
||||
@@ -250,24 +251,17 @@ export const TaskTool = Tool.define(
|
||||
output: renderOutput({
|
||||
sessionID: nextSession.id,
|
||||
state: "running",
|
||||
summary: "Background task updated",
|
||||
text: BACKGROUND_UPDATED,
|
||||
summary: "Background task already running",
|
||||
text: BACKGROUND_ALREADY_RUNNING,
|
||||
}),
|
||||
}
|
||||
}
|
||||
|
||||
const info = yield* background.start({
|
||||
const info = yield* jobs.start({
|
||||
id: nextSession.id,
|
||||
type: id,
|
||||
title: params.description,
|
||||
metadata,
|
||||
onPromote: Effect.all([
|
||||
ctx.metadata({
|
||||
title: params.description,
|
||||
metadata: { ...metadata, background: true, jobId: nextSession.id },
|
||||
}),
|
||||
notify(nextSession.id),
|
||||
]),
|
||||
run: runTask().pipe(Effect.onInterrupt(() => ops.cancel(nextSession.id))),
|
||||
})
|
||||
|
||||
@@ -289,6 +283,7 @@ export const TaskTool = Tool.define(
|
||||
}
|
||||
|
||||
if (runInBackground) {
|
||||
yield* jobs.background(info.id)
|
||||
yield* notify(info.id)
|
||||
return backgroundResult()
|
||||
}
|
||||
@@ -306,23 +301,27 @@ export const TaskTool = Tool.define(
|
||||
}),
|
||||
() =>
|
||||
Effect.gen(function* () {
|
||||
const result = yield* Effect.raceFirst(
|
||||
background.wait({ id: nextSession.id }).pipe(Effect.map((waited) => waited.info)),
|
||||
background.waitForPromotion(nextSession.id),
|
||||
)
|
||||
if (result?.metadata?.background === true) return backgroundResult()
|
||||
if (result?.status === "error") return yield* Effect.fail(new Error(result.error ?? "Task failed"))
|
||||
if (result?.status === "cancelled") return yield* Effect.fail(new Error("Task cancelled"))
|
||||
const result = yield* jobs.block({ id: nextSession.id, sessionID: ctx.sessionID })
|
||||
if (result?.type === "backgrounded") {
|
||||
yield* ctx.metadata({
|
||||
title: params.description,
|
||||
metadata: { ...metadata, background: true, jobId: nextSession.id },
|
||||
})
|
||||
yield* notify(nextSession.id)
|
||||
return backgroundResult()
|
||||
}
|
||||
if (result?.info.status === "error")
|
||||
return yield* Effect.fail(new Error(result.info.error ?? "Task failed"))
|
||||
if (result?.info.status === "cancelled") return yield* Effect.fail(new Error("Task cancelled"))
|
||||
return {
|
||||
title: params.description,
|
||||
metadata,
|
||||
output: renderOutput({ sessionID: nextSession.id, state: "completed", text: result?.output ?? "" }),
|
||||
output: renderOutput({ sessionID: nextSession.id, state: "completed", text: result?.info.output ?? "" }),
|
||||
}
|
||||
}),
|
||||
(_, exit) =>
|
||||
Effect.gen(function* () {
|
||||
if (Exit.hasInterrupts(exit))
|
||||
yield* Effect.all([cancel, background.cancel(nextSession.id)], { discard: true })
|
||||
if (Exit.hasInterrupts(exit)) yield* Effect.all([cancel, jobs.cancel(nextSession.id)], { discard: true })
|
||||
}).pipe(
|
||||
Effect.ensuring(
|
||||
Effect.sync(() => {
|
||||
|
||||
Reference in New Issue
Block a user