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 }));
    }
  });