Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions src/app/SessionWindowApp.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,7 @@ import {
} from "@/features/chat/stores/sessionWindowStore";
import { useBerdctlQueuedMessageDrain } from "@/features/berdctl/bridge/useBerdctlQueuedMessageDrain";
import { ChatView } from "@/features/chat/ui/ChatView";
import { ReleasedQueuedMessageDrain } from "@/features/chat/ui/ReleasedQueuedMessageDrain";
import { BackgroundQueuedMessageDrain } from "@/features/chat/ui/BackgroundQueuedMessageDrain";
import { useWorkspaceNameRequestQueue } from "@/features/chat/hooks/useWorkspaceNameRequestQueue";
import { ProjectWorkspaceStartupNameDialog } from "@/features/projects/ui/ProjectWorkspaceStartupNameDialog";
import { Button } from "@/shared/ui/button";
Expand Down Expand Up @@ -437,7 +437,7 @@ export function SessionWindowApp({

return (
<>
<ReleasedQueuedMessageDrain
<BackgroundQueuedMessageDrain
sessionId={sessionId}
ownerReady={phase === "ready"}
/>
Expand Down
8 changes: 4 additions & 4 deletions src/app/main.berdctl.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -101,15 +101,15 @@ describe("main entrypoint berdctl bridge loading", () => {
const sessionBranch = mainSource.slice(sessionBranchStart, mainBranchStart);
const mainBranch = mainSource.slice(mainBranchStart);
expect(sessionBranch).not.toContain("<OptionalBerdctlBridge />");
expect(mainBranch).toContain("<ReleasedQueuedMessageDrain />");
expect(mainBranch).toContain("<BackgroundQueuedMessageDrain />");
expect(mainBranch).toContain("<OptionalBerdctlBridge />");
});

it("mounts the released queued-message drain unconditionally", () => {
it("mounts the background queued-message drain unconditionally", () => {
expect(mainSource).toContain(
'import { ReleasedQueuedMessageDrain } from "@/features/chat/ui/ReleasedQueuedMessageDrain"',
'import { BackgroundQueuedMessageDrain } from "@/features/chat/ui/BackgroundQueuedMessageDrain"',
);
expect(mainSource).not.toContain("OptionalReleasedQueuedMessageDrain");
expect(mainSource).not.toContain("OptionalBackgroundQueuedMessageDrain");
});

it("reports a failed dynamic bridge import without boot-failing the app", () => {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { SessionDispatchContentionError } from "@/features/chat/lib/sessionDispatchAcquisition";
import type { SessionDispatchReleaseWaiter } from "@/features/chat/lib/sessionTargetCoordinator";
import { QueuedMessageOwnershipLostError } from "@/features/chat/lib/preCommitSendRejection";
import { resetReclaimedQueueReconciliationForTesting } from "@/features/chat/lib/reclaimedQueueReconciliation";
import {
type QueuedMessageRecord,
useChatStore,
Expand Down Expand Up @@ -80,6 +81,7 @@ function resetChatStore(): void {
describe("useBerdctlQueuedMessageDrain", () => {
beforeEach(() => {
vi.clearAllMocks();
resetReclaimedQueueReconciliationForTesting();
mocks.sendPromptToExistingSessionInBackground.mockResolvedValue(undefined);
mocks.sendQueuedPromptToExistingSessionInBackground.mockResolvedValue(
undefined,
Expand Down
43 changes: 16 additions & 27 deletions src/features/berdctl/bridge/useBerdctlQueuedMessageDrain.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,13 +15,17 @@ import {
useChatStore,
} from "@/features/chat/stores/chatStore";
import { useChatSessionStore } from "@/features/chat/stores/chatSessionStore";
import { loadPersistedMessageQueues } from "@/features/chat/stores/queuePersistence";
import { useSessionWindowStore } from "@/features/chat/stores/sessionWindowStore";
import {
isBerdctlCrossSessionQueuedMessage,
sendPromptToExistingSessionInBackground,
} from "@/features/berdctl/commands/runtime/sessionSend";
import { SessionDispatchContentionError } from "@/features/chat/lib/sessionDispatchAcquisition";
import {
isReclaimedQueueReconciliationPending,
requestReclaimedQueueReconciliation,
subscribeReclaimedQueueReconciliation,
} from "@/features/chat/lib/reclaimedQueueReconciliation";

const drainingSessionIds = new Set<string>();
const activeOwners = new Set<string>();
Expand All @@ -35,7 +39,6 @@ type ContentionWaiter = {
};

const contentionWaiters = new Map<string, ContentionWaiter>();
let ownershipRefreshSequence = 0;

function ownerIdFor(scopedSessionId?: string): string {
return scopedSessionId ? `session:${scopedSessionId}` : "global";
Expand Down Expand Up @@ -63,28 +66,9 @@ function scheduleContentionResume(
queueMicrotask(() => drainQueuedMessage(sessionId, waiter.ownerId));
}

async function refreshReclaimedQueues(
previousOpenSessions: Record<string, string>,
openSessions: Record<string, string>,
): Promise<void> {
const reclaimedSessionIds = Object.keys(previousOpenSessions).filter(
(sessionId) => !(sessionId in openSessions),
);
if (reclaimedSessionIds.length === 0) {
drainReadyQueuedMessages();
return;
}
const sequence = ++ownershipRefreshSequence;
const persistedQueues = await loadPersistedMessageQueues();
if (sequence !== ownershipRefreshSequence) return;
useChatStore
.getState()
.reconcileQueuedMessages(persistedQueues, reclaimedSessionIds);
drainReadyQueuedMessages();
}

function drainQueuedMessage(queuedSessionId: string, ownerId: string): void {
if (!activeOwners.has(ownerId)) return;
if (isReclaimedQueueReconciliationPending(queuedSessionId)) return;
const sessionExists = Boolean(
useChatSessionStore.getState().getSession(queuedSessionId),
);
Expand Down Expand Up @@ -251,16 +235,20 @@ export function useBerdctlQueuedMessageDrain(
(!previousState.hasLoadedSnapshot ||
state.openSessions !== previousState.openSessions)
) {
if (!previousState.hasLoadedSnapshot) {
drainReadyQueuedMessages();
} else {
void refreshReclaimedQueues(
if (
!previousState.hasLoadedSnapshot ||
!requestReclaimedQueueReconciliation(
previousState.openSessions,
state.openSessions,
);
)
) {
drainReadyQueuedMessages();
}
}
});
const unsubscribeReclaimedQueues = subscribeReclaimedQueueReconciliation(
() => drainReadyQueuedMessages(queuedSessionId),
);
const unsubscribeSessionStore = useChatSessionStore.subscribe(
(state, previousState) => {
reconcileContentionWaiters(queuedSessionId);
Expand Down Expand Up @@ -343,6 +331,7 @@ export function useBerdctlQueuedMessageDrain(
);
return () => {
unsubscribeWindowStore?.();
unsubscribeReclaimedQueues();
unsubscribeSessionStore();
unsubscribeChatStore();
activeOwners.delete(ownerId);
Expand Down
12 changes: 1 addition & 11 deletions src/features/berdctl/commands/runtime/sessionSend.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,9 +14,9 @@ export {
SessionDispatchUnresolvedError,
} from "@/features/chat/lib/queuedSessionSend";
import { formatIncludedWorkspacesPrompt } from "@/features/chat/lib/workspaceAttachments";
import type { QueuedMessageRecord } from "@/features/chat/stores/chatStore";
import type { MessageMetadata } from "@/shared/types/messages";
import type { ChatSendOptions } from "@/features/chat/types";
export { isBerdctlCrossSessionQueuedMessage } from "@/features/chat/lib/queuedMessageOrigin";

export const BERDCTL_CROSS_SESSION_ORIGIN =
"berdctl_cross_session" satisfies NonNullable<MessageMetadata["origin"]>;
Expand All @@ -32,16 +32,6 @@ export function berdctlCrossSessionSendOptions(): ChatSendOptions {
};
}

export function isBerdctlCrossSessionQueuedMessage(
message: QueuedMessageRecord | undefined,
): boolean {
return (
message?.kind === "transport-ready" &&
message.payload.sendOptions?.userMessageMetadata?.origin ===
BERDCTL_CROSS_SESSION_ORIGIN
);
}

export async function sendPromptToExistingSessionInBackground(
sessionId: string,
prompt: string,
Expand Down
Loading