|
| 1 | +import { |
| 2 | + type AgentSession, |
| 3 | + type QueuedMessage, |
| 4 | + sendableQueuePrefixLength, |
| 5 | +} from "@posthog/shared"; |
| 6 | +import { afterEach, describe, expect, it } from "vitest"; |
| 7 | +import { sessionStore, sessionStoreSetters } from "./sessionStore"; |
| 8 | + |
| 9 | +const RUN = "run-queue"; |
| 10 | +const TASK = "task-queue"; |
| 11 | + |
| 12 | +function seedQueue(messages: QueuedMessage[]) { |
| 13 | + sessionStoreSetters.setSession({ |
| 14 | + taskRunId: RUN, |
| 15 | + taskId: TASK, |
| 16 | + events: [], |
| 17 | + messageQueue: messages, |
| 18 | + pendingPermissions: new Map(), |
| 19 | + status: "connected", |
| 20 | + } as unknown as AgentSession); |
| 21 | +} |
| 22 | + |
| 23 | +function queue(): QueuedMessage[] { |
| 24 | + return sessionStore.getState().sessions[RUN].messageQueue; |
| 25 | +} |
| 26 | + |
| 27 | +function msg(id: string, content: string): QueuedMessage { |
| 28 | + return { id, content, queuedAt: 1 }; |
| 29 | +} |
| 30 | + |
| 31 | +afterEach(() => sessionStoreSetters.removeSession(RUN)); |
| 32 | + |
| 33 | +describe("moveQueuedMessage", () => { |
| 34 | + it("moves a message to a later position, preserving the others' order", () => { |
| 35 | + seedQueue([msg("a", "A"), msg("b", "B"), msg("c", "C")]); |
| 36 | + |
| 37 | + sessionStoreSetters.moveQueuedMessage(TASK, 0, 2); |
| 38 | + |
| 39 | + expect(queue().map((m) => m.id)).toEqual(["b", "c", "a"]); |
| 40 | + }); |
| 41 | + |
| 42 | + it("moves a message to an earlier position", () => { |
| 43 | + seedQueue([msg("a", "A"), msg("b", "B"), msg("c", "C")]); |
| 44 | + |
| 45 | + sessionStoreSetters.moveQueuedMessage(TASK, 2, 0); |
| 46 | + |
| 47 | + expect(queue().map((m) => m.id)).toEqual(["c", "a", "b"]); |
| 48 | + }); |
| 49 | + |
| 50 | + it.each([ |
| 51 | + ["same index", 1, 1], |
| 52 | + ["from out of range", 5, 0], |
| 53 | + ["to out of range", 0, 9], |
| 54 | + ["negative index", -1, 0], |
| 55 | + ])("is a no-op for %s", (_label, from, to) => { |
| 56 | + seedQueue([msg("a", "A"), msg("b", "B"), msg("c", "C")]); |
| 57 | + |
| 58 | + sessionStoreSetters.moveQueuedMessage(TASK, from, to); |
| 59 | + |
| 60 | + expect(queue().map((m) => m.id)).toEqual(["a", "b", "c"]); |
| 61 | + }); |
| 62 | +}); |
| 63 | + |
| 64 | +describe("updateQueuedMessage", () => { |
| 65 | + it("replaces content and rawPrompt in place, keeping id and position", () => { |
| 66 | + seedQueue([msg("a", "A"), msg("b", "B")]); |
| 67 | + |
| 68 | + sessionStoreSetters.updateQueuedMessage(TASK, "a", { |
| 69 | + content: "edited A", |
| 70 | + rawPrompt: "edited A raw", |
| 71 | + }); |
| 72 | + |
| 73 | + expect(queue().map((m) => m.id)).toEqual(["a", "b"]); |
| 74 | + expect(queue()[0].content).toBe("edited A"); |
| 75 | + expect(queue()[0].rawPrompt).toBe("edited A raw"); |
| 76 | + expect(queue()[1].content).toBe("B"); |
| 77 | + }); |
| 78 | + |
| 79 | + it("clears a stale rawPrompt when the patch omits it (local edit)", () => { |
| 80 | + seedQueue([{ id: "a", content: "A", rawPrompt: "old raw", queuedAt: 1 }]); |
| 81 | + |
| 82 | + sessionStoreSetters.updateQueuedMessage(TASK, "a", { content: "edited" }); |
| 83 | + |
| 84 | + expect(queue()[0].content).toBe("edited"); |
| 85 | + expect(queue()[0].rawPrompt).toBeUndefined(); |
| 86 | + }); |
| 87 | + |
| 88 | + it("is a no-op when the target id is not in the queue", () => { |
| 89 | + seedQueue([msg("a", "A")]); |
| 90 | + |
| 91 | + sessionStoreSetters.updateQueuedMessage(TASK, "missing", { |
| 92 | + content: "edited", |
| 93 | + }); |
| 94 | + |
| 95 | + expect(queue()[0].content).toBe("A"); |
| 96 | + }); |
| 97 | +}); |
| 98 | + |
| 99 | +describe("sendableQueuePrefixLength", () => { |
| 100 | + const q = (ids: string[]): Pick<AgentSession, "messageQueue"> => ({ |
| 101 | + messageQueue: ids.map((id) => msg(id, id)), |
| 102 | + }); |
| 103 | + |
| 104 | + it("returns the full length when nothing is being edited", () => { |
| 105 | + expect(sendableQueuePrefixLength(q(["a", "b", "c"]))).toBe(3); |
| 106 | + }); |
| 107 | + |
| 108 | + it("stops at the edited message, so only earlier messages count", () => { |
| 109 | + expect( |
| 110 | + sendableQueuePrefixLength({ |
| 111 | + ...q(["a", "b", "c"]), |
| 112 | + editingQueuedId: "b", |
| 113 | + }), |
| 114 | + ).toBe(1); |
| 115 | + }); |
| 116 | + |
| 117 | + it("returns 0 when the head message is being edited", () => { |
| 118 | + expect( |
| 119 | + sendableQueuePrefixLength({ |
| 120 | + ...q(["a", "b", "c"]), |
| 121 | + editingQueuedId: "a", |
| 122 | + }), |
| 123 | + ).toBe(0); |
| 124 | + }); |
| 125 | + |
| 126 | + it("returns the full length when the edited id already left the queue", () => { |
| 127 | + expect( |
| 128 | + sendableQueuePrefixLength({ |
| 129 | + ...q(["a", "b", "c"]), |
| 130 | + editingQueuedId: "gone", |
| 131 | + }), |
| 132 | + ).toBe(3); |
| 133 | + }); |
| 134 | +}); |
| 135 | + |
| 136 | +describe("editing hold on the drain", () => { |
| 137 | + it("set/clear stores and releases the edit hold", () => { |
| 138 | + seedQueue([msg("a", "A"), msg("b", "B")]); |
| 139 | + |
| 140 | + sessionStoreSetters.setEditingQueuedMessage(TASK, "b"); |
| 141 | + expect(sessionStore.getState().sessions[RUN].editingQueuedId).toBe("b"); |
| 142 | + |
| 143 | + sessionStoreSetters.clearEditingQueuedMessage(TASK); |
| 144 | + expect( |
| 145 | + sessionStore.getState().sessions[RUN].editingQueuedId, |
| 146 | + ).toBeUndefined(); |
| 147 | + }); |
| 148 | + |
| 149 | + it("stopAtEdited drains only the messages before the edited one", () => { |
| 150 | + seedQueue([msg("a", "A"), msg("b", "B"), msg("c", "C")]); |
| 151 | + sessionStoreSetters.setEditingQueuedMessage(TASK, "b"); |
| 152 | + |
| 153 | + const combined = sessionStoreSetters.dequeueMessagesAsText(TASK, { |
| 154 | + stopAtEdited: true, |
| 155 | + }); |
| 156 | + |
| 157 | + expect(combined).toBe("A"); |
| 158 | + // The edited message and everything after it stay queued. |
| 159 | + expect(queue().map((m) => m.id)).toEqual(["b", "c"]); |
| 160 | + }); |
| 161 | + |
| 162 | + it("stopAtEdited sends nothing when the head message is being edited", () => { |
| 163 | + seedQueue([msg("a", "A"), msg("b", "B")]); |
| 164 | + sessionStoreSetters.setEditingQueuedMessage(TASK, "a"); |
| 165 | + |
| 166 | + expect( |
| 167 | + sessionStoreSetters.dequeueMessagesAsText(TASK, { stopAtEdited: true }), |
| 168 | + ).toBeNull(); |
| 169 | + expect( |
| 170 | + sessionStoreSetters.dequeueMessages(TASK, { stopAtEdited: true }), |
| 171 | + ).toEqual([]); |
| 172 | + expect(queue().map((m) => m.id)).toEqual(["a", "b"]); |
| 173 | + }); |
| 174 | + |
| 175 | + it("dequeueMessages with stopAtEdited returns the sendable prefix as raw items", () => { |
| 176 | + seedQueue([msg("a", "A"), msg("b", "B"), msg("c", "C")]); |
| 177 | + sessionStoreSetters.setEditingQueuedMessage(TASK, "c"); |
| 178 | + |
| 179 | + const drained = sessionStoreSetters.dequeueMessages(TASK, { |
| 180 | + stopAtEdited: true, |
| 181 | + }); |
| 182 | + |
| 183 | + expect(drained.map((m) => m.id)).toEqual(["a", "b"]); |
| 184 | + expect(queue().map((m) => m.id)).toEqual(["c"]); |
| 185 | + }); |
| 186 | + |
| 187 | + it("drains the whole queue when stopAtEdited is not set, even mid-edit", () => { |
| 188 | + seedQueue([msg("a", "A"), msg("b", "B"), msg("c", "C")]); |
| 189 | + sessionStoreSetters.setEditingQueuedMessage(TASK, "b"); |
| 190 | + |
| 191 | + // Cancel / recall paths pull everything back regardless of the edit hold. |
| 192 | + const combined = sessionStoreSetters.dequeueMessagesAsText(TASK); |
| 193 | + |
| 194 | + expect(combined).toBe("A\n\nB\n\nC"); |
| 195 | + expect(queue()).toEqual([]); |
| 196 | + }); |
| 197 | +}); |
| 198 | + |
| 199 | +describe("sequential drain (max: 1)", () => { |
| 200 | + it("dequeueMessagesAsText drains only the head message, leaving the rest", () => { |
| 201 | + seedQueue([msg("a", "A"), msg("b", "B"), msg("c", "C")]); |
| 202 | + |
| 203 | + const first = sessionStoreSetters.dequeueMessagesAsText(TASK, { max: 1 }); |
| 204 | + expect(first).toBe("A"); |
| 205 | + expect(queue().map((m) => m.id)).toEqual(["b", "c"]); |
| 206 | + |
| 207 | + // The turn-end drain fires again per turn; each call takes the next head. |
| 208 | + const second = sessionStoreSetters.dequeueMessagesAsText(TASK, { max: 1 }); |
| 209 | + expect(second).toBe("B"); |
| 210 | + expect(queue().map((m) => m.id)).toEqual(["c"]); |
| 211 | + }); |
| 212 | + |
| 213 | + it("dequeueMessages drains only the head message as a raw item", () => { |
| 214 | + seedQueue([msg("a", "A"), msg("b", "B")]); |
| 215 | + |
| 216 | + const drained = sessionStoreSetters.dequeueMessages(TASK, { max: 1 }); |
| 217 | + |
| 218 | + expect(drained.map((m) => m.id)).toEqual(["a"]); |
| 219 | + expect(queue().map((m) => m.id)).toEqual(["b"]); |
| 220 | + }); |
| 221 | + |
| 222 | + it("takes min(max, edit boundary): head sends, edited message and rest stay", () => { |
| 223 | + seedQueue([msg("a", "A"), msg("b", "B"), msg("c", "C")]); |
| 224 | + sessionStoreSetters.setEditingQueuedMessage(TASK, "b"); |
| 225 | + |
| 226 | + const first = sessionStoreSetters.dequeueMessagesAsText(TASK, { |
| 227 | + stopAtEdited: true, |
| 228 | + max: 1, |
| 229 | + }); |
| 230 | + expect(first).toBe("A"); |
| 231 | + expect(queue().map((m) => m.id)).toEqual(["b", "c"]); |
| 232 | + |
| 233 | + // Next drain sends nothing: the new head is the message being edited. |
| 234 | + expect( |
| 235 | + sessionStoreSetters.dequeueMessagesAsText(TASK, { |
| 236 | + stopAtEdited: true, |
| 237 | + max: 1, |
| 238 | + }), |
| 239 | + ).toBeNull(); |
| 240 | + expect(queue().map((m) => m.id)).toEqual(["b", "c"]); |
| 241 | + }); |
| 242 | +}); |
0 commit comments