← Public packages
@kentcdodds/ai
Kody tool-using agent turns with Vercel AI SDK and Cloudflare AI Gateway.
src/runs.ts
174 lines · 4.9 KB · TypeScriptimport { packageStorage, workflows } from 'kody:runtime'
import { migrateAiFromValues } from './legacy-value.ts'
import { agentTurnStream } from './index.ts'
import type {
AgentTurnCancelOutput,
AgentTurnInput,
AgentTurnNextOutput,
AgentTurnStartOutput,
AgentTurnStreamEvent,
} from './types.ts'
type StoredRun = {
sessionId: string
runId: string
conversationId?: string
status: 'running' | 'completed' | 'failed' | 'cancelled'
events: AgentTurnStreamEvent[]
error: string | null
createdAt: string
updatedAt: string
}
type StartAgentTurnInput = AgentTurnInput & {
sessionId?: string
}
type ReadNextInput = {
sessionId?: string
runId?: string
cursor?: number
}
type CancelInput = {
sessionId?: string
runId?: string
}
type RunStoredInput = {
sessionId: string
runId: string
turn?: AgentTurnInput
}
function ensureStorage() {
return packageStorage()
}
function key(sessionId: string, runId: string) {
return 'agent-turn-run:' + sessionId + ':' + runId
}
/** Start a background agent turn stored in package storage; returns polling handles. */
export async function startAgentTurn(
input: StartAgentTurnInput = { messages: [] },
): Promise<AgentTurnStartOutput> {
const sessionId = String(input.sessionId || crypto.randomUUID())
const runId = crypto.randomUUID()
const conversationId = input.conversationId || crypto.randomUUID()
const run: StoredRun = {
sessionId,
runId,
conversationId,
status: 'running',
events: [],
error: null,
createdAt: new Date().toISOString(),
updatedAt: new Date().toISOString(),
}
await ensureStorage().set(key(sessionId, runId), run)
if (workflows?.create) {
await workflows.create({
exportName: './worker',
runAt: new Date().toISOString(),
idempotencyKey: 'agent-turn:' + sessionId + ':' + runId,
params: { sessionId, runId, turn: { ...input, sessionId, conversationId } },
})
} else {
await runStoredAgentTurn({
sessionId,
runId,
turn: { ...input, sessionId, conversationId },
})
}
return { ok: true, runId, sessionId, conversationId }
}
/** Read the next slice of buffered events for a stored agent turn run. */
export async function readNextAgentTurnEvents(
input: ReadNextInput = {},
): Promise<AgentTurnNextOutput> {
const sessionId = String(input.sessionId || '')
const runId = String(input.runId || '')
const cursor = Number(input.cursor || 0)
const run = (await ensureStorage().get(key(sessionId, runId))) as StoredRun | null
if (!run) throw new Error('Agent turn run not found: ' + runId)
const events = (run.events || []).slice(cursor)
return {
ok: true,
events,
nextCursor: cursor + events.length,
done: run.status !== 'running',
}
}
/** Cancel a stored agent turn run by marking it cancelled in package storage. */
export async function cancelAgentTurn(input: CancelInput = {}): Promise<AgentTurnCancelOutput> {
const sessionId = String(input.sessionId || '')
const runId = String(input.runId || '')
const run = (await ensureStorage().get(key(sessionId, runId))) as StoredRun | null
if (run) {
await ensureStorage().set(key(sessionId, runId), {
...run,
status: 'cancelled',
error: 'Turn cancelled.',
updatedAt: new Date().toISOString(),
})
}
return { ok: true, cancelled: true }
}
/** Execute a stored agent turn, appending buffered events until completion or failure. */
export async function runStoredAgentTurn(input: RunStoredInput) {
const sessionId = String(input.sessionId)
const runId = String(input.runId)
const events: AgentTurnStreamEvent[] = []
try {
for await (const event of agentTurnStream(input.turn || { messages: [] })) events.push(event)
const stored = (await ensureStorage().get(key(sessionId, runId))) as StoredRun | null
if (stored?.status === 'cancelled') return { ok: true, cancelled: true }
await ensureStorage().set(key(sessionId, runId), {
...stored,
sessionId,
runId,
conversationId: input.turn?.conversationId,
status: 'completed',
events,
error: null,
updatedAt: new Date().toISOString(),
})
return { ok: true }
} catch (error) {
events.push({
type: 'error',
message: error instanceof Error ? error.message : String(error),
phase: 'run',
})
const stored = (await ensureStorage().get(key(sessionId, runId))) as StoredRun | null
if (stored?.status !== 'cancelled') {
await ensureStorage().set(key(sessionId, runId), {
...stored,
sessionId,
runId,
status: 'failed',
events,
error: events.at(-1)!.message,
updatedAt: new Date().toISOString(),
})
}
throw error
}
}
/**
* Default export for `./runs`: run a stored agent turn workflow body.
* @example
* import runAgentTurnWorker from 'kody:@kentcdodds/ai/runs'
* await runAgentTurnWorker({ sessionId: 'sess-1', runId: 'run-1', turn: { messages: [{ role: 'user', content: 'Hi' }] } })
*/
export default async function runAgentTurnWorker(input: RunStoredInput = {} as RunStoredInput) {
if (!input?.sessionId || !input?.runId) {
return await migrateAiFromValues()
}
return await runStoredAgentTurn(input)
}