File size: 2,947 Bytes
68d7816 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 | /* oxlint-disable typescript-eslint/no-unsafe-declaration-merging, eslint-plugin-import/namespace -- Event2 class+payload-interface declaration merging is the sanctioned event-declaration idiom. */
import { z } from 'zod';
import { AgentEvent2 } from '#/app/event/event2';
import { defineState } from '#/state/state';
import type { AgentTaskNotificationContext } from './task';
import type { AgentTaskInfo } from './types';
export type TaskModelState = Map<string, AgentTaskInfo>;
const taskStartedSchema = z.object({
agentId: z.string(),
info: z.custom<AgentTaskInfo>(),
});
export class TaskStarted extends AgentEvent2<z.infer<typeof taskStartedSchema>> {
static override readonly type = 'task.started';
static override readonly durable = true;
static override readonly observable = true;
static override readonly schema = taskStartedSchema;
}
export interface TaskStarted {
readonly agentId: string;
readonly info: AgentTaskInfo;
}
const taskTerminatedSchema = z.object({
agentId: z.string(),
info: z.custom<AgentTaskInfo>(),
outputTail: z.string().optional(),
});
export class TaskTerminated extends AgentEvent2<z.infer<typeof taskTerminatedSchema>> {
static override readonly type = 'task.terminated';
static override readonly durable = true;
static override readonly schema = taskTerminatedSchema;
}
export interface TaskTerminated {
readonly agentId: string;
readonly info: AgentTaskInfo;
readonly outputTail?: string;
}
export interface TaskTerminatedNoticePayload {
readonly agentId: string;
readonly info: AgentTaskInfo;
}
export class TaskTerminatedNotice extends AgentEvent2<TaskTerminatedNoticePayload> {
static override readonly type = 'task.terminated';
static override readonly observable = true;
}
export interface TaskTerminatedNotice extends TaskTerminatedNoticePayload {}
export class TaskNotified extends AgentEvent2<AgentTaskNotificationContext> {
static override readonly type = 'task.notified';
static override readonly observable = true;
}
export interface TaskNotified extends AgentTaskNotificationContext {}
const taskWaitDeliveredSchema = z.object({
agentId: z.string(),
keys: z.array(z.string()),
});
export class TaskWaitDelivered extends AgentEvent2<z.infer<typeof taskWaitDeliveredSchema>> {
static override readonly type = 'task.waitDelivered';
static override readonly durable = true;
static override readonly schema = taskWaitDeliveredSchema;
}
export interface TaskWaitDelivered {
readonly agentId: string;
readonly keys: string[];
}
export const taskKey = defineState('task', (): TaskModelState => new Map()).replayable({
schema: z.custom<TaskModelState>(),
})
.on(TaskStarted, (s, e) => {
s.set(e.info.taskId, e.info);
})
.on(TaskTerminated, (s, e, ctx) => {
s.set(e.info.taskId, e.info);
if (e instanceof TaskTerminated) {
ctx.emit(new TaskTerminatedNotice({ agentId: e.agentId, info: e.info }));
}
});
|