task: implement sdk runner

This commit is contained in:
YeonGyu-Kim
2026-03-16 14:43:47 +09:00
parent 8c4fa47e5e
commit 96f4b3b56c
20 changed files with 1350 additions and 137 deletions
+3
View File
@@ -0,0 +1,3 @@
import type { CreateOmoRunnerOptions, OmoRunner } from "./types"
export declare function createOmoRunner(options: CreateOmoRunnerOptions): OmoRunner
+106
View File
@@ -0,0 +1,106 @@
import { beforeEach, describe, expect, it, mock } from "bun:test"
import type { ServerConnection } from "../../../src/cli/run/types"
const mockCreateServerConnection = mock(
async (): Promise<ServerConnection> => ({
client: {} as never,
cleanup: mock(() => {}),
}),
)
const mockExecuteRunSession = mock(async (_options: unknown) => ({
exitCode: 0,
sessionId: "ses_runner",
result: {
sessionId: "ses_runner",
success: true,
durationMs: 10,
messageCount: 1,
summary: "done",
},
}))
mock.module("../../../src/cli/run/server-connection", () => ({
createServerConnection: mockCreateServerConnection,
}))
mock.module("../../../src/cli/run/run-engine", () => ({
executeRunSession: mockExecuteRunSession,
}))
const { createOmoRunner } = await import("./create-omo-runner")
describe("createOmoRunner", () => {
beforeEach(() => {
mockCreateServerConnection.mockClear()
mockExecuteRunSession.mockClear()
})
it("reuses the same connection and enables question-aware execution", async () => {
const runner = createOmoRunner({
directory: "/repo",
agent: "atlas",
})
const first = await runner.run("first")
const second = await runner.run("second", { agent: "prometheus" })
expect(first.summary).toBe("done")
expect(second.summary).toBe("done")
expect(mockCreateServerConnection).toHaveBeenCalledTimes(1)
expect(mockExecuteRunSession).toHaveBeenNthCalledWith(1, expect.objectContaining({
directory: "/repo",
agent: "atlas",
questionPermission: "allow",
questionToolEnabled: true,
renderOutput: false,
}))
expect(mockExecuteRunSession).toHaveBeenNthCalledWith(2, expect.objectContaining({
agent: "prometheus",
}))
await runner.close()
})
it("streams normalized events", async () => {
mockExecuteRunSession.mockImplementationOnce(async (options: { eventObserver?: { onEvent?: (event: unknown) => Promise<void> } }) => {
await options.eventObserver?.onEvent?.({
type: "session.started",
sessionId: "ses_runner",
agent: "Atlas (Plan Executor)",
resumed: false,
})
await options.eventObserver?.onEvent?.({
type: "session.completed",
sessionId: "ses_runner",
result: {
sessionId: "ses_runner",
success: true,
durationMs: 10,
messageCount: 1,
summary: "done",
},
})
return {
exitCode: 0,
sessionId: "ses_runner",
result: {
sessionId: "ses_runner",
success: true,
durationMs: 10,
messageCount: 1,
summary: "done",
},
}
})
const runner = createOmoRunner({ directory: "/repo" })
const seenTypes: string[] = []
for await (const event of runner.stream("stream")) {
seenTypes.push(event.type)
}
expect(seenTypes).toEqual(["session.started", "session.completed"])
await runner.close()
})
})
+186
View File
@@ -0,0 +1,186 @@
import { createServerConnection } from "../../../src/cli/run/server-connection"
import { executeRunSession } from "../../../src/cli/run/run-engine"
import type { RunEventObserver, ServerConnection } from "../../../src/cli/run/types"
import type {
CreateOmoRunnerOptions,
OmoRunInvocationOptions,
OmoRunner,
RunResult,
StreamEvent,
} from "./types"
class AsyncEventQueue<T> implements AsyncIterable<T> {
private readonly values: T[] = []
private readonly waiters: Array<(value: IteratorResult<T>) => void> = []
private closed = false
push(value: T): void {
if (this.closed) return
const waiter = this.waiters.shift()
if (waiter) {
waiter({ value, done: false })
return
}
this.values.push(value)
}
close(): void {
if (this.closed) return
this.closed = true
while (this.waiters.length > 0) {
const waiter = this.waiters.shift()
waiter?.({ value: undefined, done: true })
}
}
[Symbol.asyncIterator](): AsyncIterator<T> {
return {
next: () => {
const value = this.values.shift()
if (value !== undefined) {
return Promise.resolve({ value, done: false })
}
if (this.closed) {
return Promise.resolve({ value: undefined, done: true })
}
return new Promise<IteratorResult<T>>((resolve) => {
this.waiters.push(resolve)
})
},
}
}
}
export function createOmoRunner(options: CreateOmoRunnerOptions): OmoRunner {
const {
directory,
agent,
port,
model,
attach,
includeRawEvents = false,
onIdle,
onQuestion,
onComplete,
onError,
} = options
let connectionPromise: Promise<ServerConnection> | null = null
let closed = false
let activeRun: Promise<unknown> | null = null
const silentLogger = {
log: () => {},
error: () => {},
}
const ensureConnection = async (): Promise<ServerConnection> => {
if (closed) {
throw new Error("Runner is closed")
}
if (connectionPromise === null) {
const controller = new AbortController()
connectionPromise = createServerConnection({
port,
attach,
signal: controller.signal,
logger: silentLogger,
})
}
return await connectionPromise
}
const createObserver = (
queue?: AsyncEventQueue<StreamEvent>,
): RunEventObserver => ({
includeRawEvents,
onEvent: async (event) => {
queue?.push(event as StreamEvent)
},
onIdle,
onQuestion,
onComplete,
onError,
})
const runOnce = async (
prompt: string,
invocationOptions: OmoRunInvocationOptions | undefined,
observer: RunEventObserver,
): Promise<RunResult> => {
if (activeRun !== null) {
throw new Error("Runner already has an active operation")
}
const connection = await ensureConnection()
const execution = executeRunSession({
client: connection.client,
message: prompt,
directory,
agent: invocationOptions?.agent ?? agent,
model: invocationOptions?.model ?? model,
sessionId: invocationOptions?.sessionId,
questionPermission: "allow",
questionToolEnabled: true,
renderOutput: false,
logger: silentLogger,
eventObserver: observer,
signal: invocationOptions?.signal,
})
activeRun = execution
const abortHandler = () => {
void observer.onError?.({
type: "session.error",
sessionId: invocationOptions?.sessionId ?? "",
error: "Aborted by caller",
})
}
invocationOptions?.signal?.addEventListener("abort", abortHandler, { once: true })
try {
const { result } = await execution
return result
} finally {
invocationOptions?.signal?.removeEventListener("abort", abortHandler)
activeRun = null
}
}
return {
async run(prompt, invocationOptions) {
return await runOnce(prompt, invocationOptions, createObserver())
},
stream(prompt, invocationOptions) {
const queue = new AsyncEventQueue<StreamEvent>()
const execution = runOnce(prompt, invocationOptions, createObserver(queue))
.catch((error) => {
queue.push({
type: "session.error",
sessionId: invocationOptions?.sessionId ?? "",
error: error instanceof Error ? error.message : String(error),
})
})
.finally(() => {
queue.close()
})
return {
async *[Symbol.asyncIterator]() {
try {
for await (const event of queue) {
yield event
}
} finally {
await execution
}
},
}
},
async close() {
closed = true
const connection = await connectionPromise
connection?.cleanup()
connectionPromise = null
},
}
}
+8
View File
@@ -0,0 +1,8 @@
export { createOmoRunner } from "./create-omo-runner"
export type {
CreateOmoRunnerOptions,
OmoRunInvocationOptions,
OmoRunner,
RunResult,
StreamEvent,
} from "./types"
+8
View File
@@ -0,0 +1,8 @@
export { createOmoRunner } from "./create-omo-runner"
export type {
CreateOmoRunnerOptions,
OmoRunInvocationOptions,
OmoRunner,
RunResult,
StreamEvent,
} from "./types"
+95
View File
@@ -0,0 +1,95 @@
export interface RunResult {
sessionId: string
success: boolean
durationMs: number
messageCount: number
summary: string
}
export type StreamEvent =
| {
type: "session.started"
sessionId: string
agent: string
resumed: boolean
model?: { providerID: string; modelID: string }
}
| {
type: "message.delta"
sessionId: string
messageId?: string
partId?: string
delta: string
}
| {
type: "message.completed"
sessionId: string
messageId?: string
partId?: string
text: string
}
| {
type: "tool.started"
sessionId: string
toolName: string
input?: unknown
}
| {
type: "tool.completed"
sessionId: string
toolName: string
output?: string
status: "completed" | "error"
}
| {
type: "session.idle"
sessionId: string
}
| {
type: "session.question"
sessionId: string
toolName: string
input?: unknown
question?: string
}
| {
type: "session.completed"
sessionId: string
result: RunResult
}
| {
type: "session.error"
sessionId: string
error: string
}
| {
type: "raw"
sessionId: string
payload: unknown
}
export interface OmoRunInvocationOptions {
sessionId?: string
signal?: AbortSignal
agent?: string
model?: string
}
export interface CreateOmoRunnerOptions {
directory: string
agent?: string
port?: number
model?: string
attach?: string
includeRawEvents?: boolean
onIdle?: (event: Extract<StreamEvent, { type: "session.idle" }>) => void | Promise<void>
onQuestion?: (event: Extract<StreamEvent, { type: "session.question" }>) => void | Promise<void>
onComplete?: (event: Extract<StreamEvent, { type: "session.completed" }>) => void | Promise<void>
onError?: (event: Extract<StreamEvent, { type: "session.error" }>) => void | Promise<void>
}
export interface OmoRunner {
run(prompt: string, options?: OmoRunInvocationOptions): Promise<RunResult>
stream(prompt: string, options?: OmoRunInvocationOptions): AsyncIterable<StreamEvent>
close(): Promise<void>
}
+95
View File
@@ -0,0 +1,95 @@
export interface RunResult {
sessionId: string
success: boolean
durationMs: number
messageCount: number
summary: string
}
export type StreamEvent =
| {
type: "session.started"
sessionId: string
agent: string
resumed: boolean
model?: { providerID: string; modelID: string }
}
| {
type: "message.delta"
sessionId: string
messageId?: string
partId?: string
delta: string
}
| {
type: "message.completed"
sessionId: string
messageId?: string
partId?: string
text: string
}
| {
type: "tool.started"
sessionId: string
toolName: string
input?: unknown
}
| {
type: "tool.completed"
sessionId: string
toolName: string
output?: string
status: "completed" | "error"
}
| {
type: "session.idle"
sessionId: string
}
| {
type: "session.question"
sessionId: string
toolName: string
input?: unknown
question?: string
}
| {
type: "session.completed"
sessionId: string
result: RunResult
}
| {
type: "session.error"
sessionId: string
error: string
}
| {
type: "raw"
sessionId: string
payload: unknown
}
export interface OmoRunInvocationOptions {
sessionId?: string
signal?: AbortSignal
agent?: string
model?: string
}
export interface CreateOmoRunnerOptions {
directory: string
agent?: string
port?: number
model?: string
attach?: string
includeRawEvents?: boolean
onIdle?: (event: Extract<StreamEvent, { type: "session.idle" }>) => void | Promise<void>
onQuestion?: (event: Extract<StreamEvent, { type: "session.question" }>) => void | Promise<void>
onComplete?: (event: Extract<StreamEvent, { type: "session.completed" }>) => void | Promise<void>
onError?: (event: Extract<StreamEvent, { type: "session.error" }>) => void | Promise<void>
}
export interface OmoRunner {
run(prompt: string, options?: OmoRunInvocationOptions): Promise<RunResult>
stream(prompt: string, options?: OmoRunInvocationOptions): AsyncIterable<StreamEvent>
close(): Promise<void>
}