Skip to content
File

Blob: src/store.ts

typescript291 lines
1import { EventEmitter } from "node:events";
2import crypto from "node:crypto";
3import type { WorkflowId } from "./workflows.js";
4 
5// ── Types ──────────────────────────────────────────────────────────
6 
7export type AgentId = "claude" | "codex";
8 
9export interface Message {
10 id: number;
11 sessionId: string;
12 from: AgentId;
13 to: AgentId;
14 body: string;
15 status: "pending" | "delivered";
16 createdAt: number;
17}
18 
19export interface Session {
20 id: string;
21 goal: string;
22 initiator: AgentId;
23 responder: AgentId;
24 workflow: WorkflowId;
25 currentPhaseIndex: number;
26 status: "active" | "closed";
27 autoForward: boolean;
28 lastSeenByAgent: Record<AgentId, number>;
29 messages: Message[];
30 nextMessageId: number;
31 createdAt: number;
32 closedAt?: number;
33 closedBy?: AgentId | "user";
34 closedReason?: string;
35}
36 
37export type UIEvent =
38 | { type: "session-created"; session: Session }
39 | { type: "session-updated"; session: Session }
40 | { type: "session-deleted"; sessionId: string }
41 | { type: "message-added"; sessionId: string; message: Message }
42 | { type: "message-updated"; sessionId: string; message: Message }
43 | { type: "message-removed"; sessionId: string; messageId: number };
44 
45export type WaitForMessageResult =
46 | { status: "received"; message: Message; skippedMessageIds: number[] }
47 | { status: "timeout" }
48 | { status: "closed" };
49 
50// ── Store ──────────────────────────────────────────────────────────
51 
52export class Store extends EventEmitter {
53 private sessions = new Map<string, Session>();
54 
55 createSession(goal: string, initiator: AgentId, workflow: WorkflowId): Session {
56 const id = "collab-" + crypto.randomBytes(4).toString("hex");
57 const responder: AgentId = initiator === "claude" ? "codex" : "claude";
58 const session: Session = {
59 id,
60 goal,
61 initiator,
62 responder,
63 workflow,
64 currentPhaseIndex: 0,
65 status: "active",
66 autoForward: true,
67 lastSeenByAgent: { claude: 0, codex: 0 },
68 messages: [],
69 nextMessageId: 1,
70 createdAt: Date.now(),
71 };
72 this.sessions.set(id, session);
73 this.emit("ui-update", { type: "session-created", session } satisfies UIEvent);
74 return session;
75 }
76 
77 getSession(id: string): Session | undefined {
78 return this.sessions.get(id);
79 }
80 
81 listSessions(): Session[] {
82 return [...this.sessions.values()];
83 }
84 
85 updateSession(
86 id: string,
87 updates: { autoForward?: boolean; goal?: string; currentPhaseIndex?: number },
88 ): Session | undefined {
89 const session = this.sessions.get(id);
90 if (!session) return undefined;
91 if (updates.autoForward !== undefined) session.autoForward = updates.autoForward;
92 if (updates.goal !== undefined) session.goal = updates.goal;
93 if (updates.currentPhaseIndex !== undefined) session.currentPhaseIndex = updates.currentPhaseIndex;
94 this.emit("ui-update", { type: "session-updated", session } satisfies UIEvent);
95 return session;
96 }
97 
98 deleteSession(id: string): boolean {
99 const existed = this.sessions.delete(id);
100 if (existed) {
101 this.emit("ui-update", { type: "session-deleted", sessionId: id } satisfies UIEvent);
102 }
103 return existed;
104 }
105 
106 closeSession(sessionId: string, closedBy: AgentId | "user", reason?: string): Session | undefined {
107 const session = this.sessions.get(sessionId);
108 if (!session) return undefined;
109 if (session.status === "closed") return session;
110 
111 session.status = "closed";
112 session.closedAt = Date.now();
113 session.closedBy = closedBy;
114 if (reason && reason.trim()) session.closedReason = reason.trim();
115 
116 this.emit("ui-update", { type: "session-updated", session } satisfies UIEvent);
117 this.emit(`message:${sessionId}`);
118 return session;
119 }
120 
121 addMessage(sessionId: string, from: AgentId, body: string): Message | undefined {
122 const session = this.sessions.get(sessionId);
123 if (!session) return undefined;
124 if (session.status === "closed") return undefined;
125 
126 const to: AgentId = from === "claude" ? "codex" : "claude";
127 const msg: Message = {
128 id: session.nextMessageId++,
129 sessionId,
130 from,
131 to,
132 body,
133 status: session.autoForward ? "delivered" : "pending",
134 createdAt: Date.now(),
135 };
136 session.messages.push(msg);
137 
138 this.emit("ui-update", { type: "message-added", sessionId, message: msg } satisfies UIEvent);
139 
140 if (msg.status === "delivered") {
141 this.emit(`message:${sessionId}`, msg);
142 }
143 
144 return msg;
145 }
146 
147 getLatestUnseenOutbound(sessionId: string, from: AgentId): Message | undefined {
148 const session = this.sessions.get(sessionId);
149 if (!session) return undefined;
150 
151 const to: AgentId = from === "claude" ? "codex" : "claude";
152 const seenId = session.lastSeenByAgent[to] ?? 0;
153 for (let i = session.messages.length - 1; i >= 0; i--) {
154 const m = session.messages[i];
155 if (m.from === from && m.to === to && m.id > seenId) {
156 return m;
157 }
158 }
159 return undefined;
160 }
161 
162 deliverMessage(sessionId: string, messageId: number): Message | undefined {
163 const session = this.sessions.get(sessionId);
164 if (!session) return undefined;
165 
166 const msg = session.messages.find((m) => m.id === messageId);
167 if (!msg || msg.status !== "pending") return undefined;
168 
169 msg.status = "delivered";
170 this.emit("ui-update", { type: "message-updated", sessionId, message: msg } satisfies UIEvent);
171 this.emit(`message:${sessionId}`, msg);
172 return msg;
173 }
174 
175 editMessage(sessionId: string, messageId: number, newBody: string): Message | undefined {
176 const session = this.sessions.get(sessionId);
177 if (!session) return undefined;
178 
179 const msg = session.messages.find((m) => m.id === messageId);
180 if (!msg || msg.status !== "pending") return undefined;
181 
182 msg.body = newBody;
183 this.emit("ui-update", { type: "message-updated", sessionId, message: msg } satisfies UIEvent);
184 return msg;
185 }
186 
187 rejectMessage(sessionId: string, messageId: number): boolean {
188 const session = this.sessions.get(sessionId);
189 if (!session) return false;
190 
191 const idx = session.messages.findIndex((m) => m.id === messageId && m.status === "pending");
192 if (idx === -1) return false;
193 
194 session.messages.splice(idx, 1);
195 this.emit("ui-update", { type: "message-removed", sessionId, messageId } satisfies UIEvent);
196 return true;
197 }
198 
199 getLatestDeliveredForAgent(sessionId: string, agent: AgentId, afterId: number): Message | undefined {
200 const session = this.sessions.get(sessionId);
201 if (!session) return undefined;
202 
203 for (let i = session.messages.length - 1; i >= 0; i--) {
204 const m = session.messages[i];
205 if (m.to === agent && m.status === "delivered" && m.id > afterId) {
206 return m;
207 }
208 }
209 return undefined;
210 }
211 
212 getSkippedDeliveredIdsForAgent(sessionId: string, agent: AgentId, afterId: number, latestId: number): number[] {
213 const session = this.sessions.get(sessionId);
214 if (!session) return [];
215 
216 return session.messages
217 .filter((m) => m.to === agent && m.status === "delivered" && m.id > afterId && m.id < latestId)
218 .map((m) => m.id);
219 }
220 
221 markSeen(sessionId: string, agent: AgentId, messageId: number): void {
222 const session = this.sessions.get(sessionId);
223 if (!session) return;
224 session.lastSeenByAgent[agent] = Math.max(session.lastSeenByAgent[agent] ?? 0, messageId);
225 }
226 
227 waitForMessage(
228 sessionId: string,
229 agent: AgentId,
230 afterId: number,
231 timeoutMs: number,
232 ): Promise<WaitForMessageResult> {
233 const buildReceived = (msg: Message): WaitForMessageResult => {
234 const skippedMessageIds = this.getSkippedDeliveredIdsForAgent(sessionId, agent, afterId, msg.id);
235 this.markSeen(sessionId, agent, msg.id);
236 return { status: "received", message: msg, skippedMessageIds };
237 };
238 
239 const session = this.sessions.get(sessionId);
240 if (session?.status === "closed") return Promise.resolve({ status: "closed" });
241 
242 const existing = this.getLatestDeliveredForAgent(sessionId, agent, afterId);
243 if (existing) return Promise.resolve(buildReceived(existing));
244 
245 return new Promise((resolve) => {
246 let settled = false;
247 const eventName = `message:${sessionId}`;
248 
249 const onMessage = () => {
250 const current = this.sessions.get(sessionId);
251 if (current?.status === "closed" && !settled) {
252 settled = true;
253 clearTimeout(timer);
254 this.removeListener(eventName, onMessage);
255 resolve({ status: "closed" });
256 return;
257 }
258 
259 const msg = this.getLatestDeliveredForAgent(sessionId, agent, afterId);
260 if (msg && !settled) {
261 settled = true;
262 clearTimeout(timer);
263 this.removeListener(eventName, onMessage);
264 resolve(buildReceived(msg));
265 }
266 };
267 
268 const timer = setTimeout(() => {
269 if (!settled) {
270 settled = true;
271 this.removeListener(eventName, onMessage);
272 resolve({ status: "timeout" });
273 }
274 }, timeoutMs);
275 
276 // Register listener BEFORE re-checking to close the race window
277 this.on(eventName, onMessage);
278 
279 // Re-check: a message may have arrived between the fast-path check and
280 // the listener registration above.
281 const arrived = this.getLatestDeliveredForAgent(sessionId, agent, afterId);
282 if (arrived && !settled) {
283 settled = true;
284 clearTimeout(timer);
285 this.removeListener(eventName, onMessage);
286 resolve(buildReceived(arrived));
287 }
288 });
289 }
290}