diff --git a/apps/server/test/threads/thread-send-dispatch.test.ts b/apps/server/test/threads/thread-send-dispatch.test.ts index acab5bdc3..98a8ac746 100644 --- a/apps/server/test/threads/thread-send-dispatch.test.ts +++ b/apps/server/test/threads/thread-send-dispatch.test.ts @@ -1,5 +1,6 @@ import { archiveThread, + updateHost, getQueuedThreadMessage, getThread, listEvents, @@ -19,6 +20,7 @@ import { import { groupHostDaemonEvents } from "@bb/host-daemon-contract"; import { afterEach, describe, expect, it, vi } from "vitest"; import type { TelemetryService } from "../../src/services/system/telemetry.js"; +import * as queuedDispatch from "../../src/services/threads/queued-message-dispatch.js"; import * as threadEvents from "../../src/services/threads/thread-events.js"; import { runQueuedMessageDispatch } from "../../src/services/threads/queued-message-dispatch.js"; import { @@ -344,6 +346,81 @@ describe("user message telemetry", () => { }); }); +describe("active thread admission before host readiness", () => { + it.each(["stale", "fresh"] as const)( + "rejects a start with a %s snapshot without waking a suspended host", + async (snapshot) => { + await withTestHarness(async (harness) => { + const { environment, thread } = seedProviderThreadFixture({ + harness, + value: 3987, + }); + updateHost(harness.db, harness.hub, environment.hostId, { + phase: "suspended", + suspendedAt: Date.now(), + }); + applyLoggedThreadLifecycleEvent(harness.deps, { + event: { type: "run.started" }, + threadId: thread.id, + }); + const current = getThread(harness.db, thread.id); + expect(current?.status).toBe("active"); + if (current === null) throw new Error("Missing seeded thread"); + const readiness = vi + .spyOn(queuedDispatch, "requestQueuedMachineReadiness") + .mockImplementation(() => {}); + + await expect( + acceptThreadSendRequest(harness.deps, { + thread: snapshot === "stale" ? thread : current, + payload: { input: textInput("begin another turn"), mode: "start" }, + }), + ).rejects.toMatchObject({ + status: 409, + body: { + code: "thread_not_writable", + details: { reason: "already_active" }, + }, + }); + expect(listQueuedThreadMessages(harness.db, thread.id)).toEqual([]); + expect(readiness).not.toHaveBeenCalled(); + }); + }, + ); + + it("keeps the host wait for queue-if-active after a stale idle snapshot", async () => { + await withTestHarness(async (harness) => { + const { environment, thread } = seedProviderThreadFixture({ + harness, + value: 3988, + }); + updateHost(harness.db, harness.hub, environment.hostId, { + phase: "suspended", + suspendedAt: Date.now(), + }); + applyLoggedThreadLifecycleEvent(harness.deps, { + event: { type: "run.started" }, + threadId: thread.id, + }); + const readiness = vi + .spyOn(queuedDispatch, "requestQueuedMachineReadiness") + .mockImplementation(() => {}); + + await expect( + acceptThreadSendRequest(harness.deps, { + thread, + payload: { input: textInput("follow up later"), mode: "queue-if-active" }, + }), + ).resolves.toMatchObject({ + delivery: "queued", + queuedMessage: { waitingOn: { kind: "host-offline" } }, + }); + expect(listQueuedThreadMessages(harness.db, thread.id)).toHaveLength(1); + expect(readiness).toHaveBeenCalledWith(harness.deps, environment.hostId); + }); + }); +}); + describe("startup queue waits", () => { it("steers provisioning input into the first turn in queue order", async () => { await withTestHarness(async (harness) => {