Skip to content

Commit 0d23f6f

Browse files
refactor(sse): add sse-frame utility and refine SSE hardening
1 parent 3c55b15 commit 0d23f6f

12 files changed

Lines changed: 89 additions & 71 deletions

File tree

backend/src/routes/sse-writer.ts

Lines changed: 2 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,5 @@
1+
import { encodeSSEFrame } from '../utils/sse-frame'
2+
13
export interface QueuedSSEWriterInput {
24
write: (chunk: Uint8Array) => Promise<unknown> | void
35
onError: (error: unknown) => void
@@ -9,14 +11,8 @@ export interface QueuedSSEWriter {
911
close: () => void
1012
}
1113

12-
const sharedEncoder = new TextEncoder()
1314
const MAX_QUEUED_FRAMES = 1024
1415

15-
export function encodeSSEFrame(event: string, data: string): Uint8Array {
16-
const head = event ? `event: ${event}\n` : ''
17-
return sharedEncoder.encode(`${head}data: ${data}\n\n`)
18-
}
19-
2016
export function createQueuedSSEWriter(input: QueuedSSEWriterInput): QueuedSSEWriter {
2117
const queue: Uint8Array[] = []
2218
let draining = false

backend/src/services/sse-aggregator.ts

Lines changed: 5 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -2,8 +2,9 @@ import { EventSource } from 'eventsource'
22
import { logger } from '../utils/logger'
33
import { ENV } from '@opencode-manager/shared/config/env'
44
import { DEFAULTS } from '@opencode-manager/shared/config'
5+
import type { SSEEventEnvelope, SSEEventPayload } from '@opencode-manager/shared'
56
import { getOpenCodeBasicAuthHeader, type OpenCodePasswordResolver } from './opencode/auth'
6-
import { encodeSSEFrame } from '../routes/sse-writer'
7+
import { encodeSSEFrame } from '../utils/sse-frame'
78

89
type SSEClientCallback = (event: string, data: string) => void
910
type SSEClientFrameWriter = (frame: Uint8Array) => void
@@ -18,17 +19,7 @@ interface SSEClient {
1819
activeSessionId: string | null
1920
}
2021

21-
export interface SSEEvent {
22-
type: string
23-
properties: Record<string, unknown>
24-
}
25-
26-
interface GlobalEventEnvelope {
27-
directory?: string
28-
project?: string
29-
workspace?: string
30-
payload: SSEEvent
31-
}
22+
export type SSEEvent = SSEEventPayload
3223

3324
export interface PendingActionsFetcher {
3425
getJson<T>(path: string, opts?: { directory?: string; signal?: AbortSignal }): Promise<T>
@@ -333,9 +324,9 @@ class SSEAggregator {
333324
}
334325

335326
private handleUpstreamMessage(data: string): void {
336-
let envelope: GlobalEventEnvelope
327+
let envelope: SSEEventEnvelope
337328
try {
338-
envelope = JSON.parse(data) as GlobalEventEnvelope
329+
envelope = JSON.parse(data) as SSEEventEnvelope
339330
} catch {
340331
return
341332
}

backend/src/utils/sse-frame.ts

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
const sharedEncoder = new TextEncoder()
2+
3+
export function encodeSSEFrame(event: string, data: string): Uint8Array {
4+
const head = event ? `event: ${event}\n` : ''
5+
return sharedEncoder.encode(`${head}data: ${data}\n\n`)
6+
}

backend/test/routes/sse-writer.test.ts

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
import { describe, it, expect, vi } from 'vitest'
2-
import { createQueuedSSEWriter, encodeSSEFrame } from '../../src/routes/sse-writer'
2+
import { createQueuedSSEWriter } from '../../src/routes/sse-writer'
3+
import { encodeSSEFrame } from '../../src/utils/sse-frame'
34

45
describe('encodeSSEFrame', () => {
56
const decoder = new TextDecoder()

frontend/src/hooks/useOpenCode.ts

Lines changed: 25 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,7 @@ import { parseNetworkError } from "../lib/opencode-errors";
1313
import { showToast } from "../lib/toast";
1414
import { useSessionStatus } from "../stores/sessionStatusStore";
1515
import { useSendErrorStore } from "../stores/sendErrorStore";
16-
import { invalidateSessionListCaches } from "../lib/queryInvalidation";
16+
import { invalidateSessionListCaches, messagesQueryKey } from "../lib/queryInvalidation";
1717

1818
type AssistantMessage = components["schemas"]["AssistantMessage"];
1919

@@ -140,7 +140,7 @@ export const useMessages = (opcodeUrl: string | null | undefined, sessionID: str
140140
const client = useOpenCodeClient(opcodeUrl, directory);
141141

142142
return useQuery({
143-
queryKey: ["opencode", "messages", opcodeUrl, sessionID, directory],
143+
queryKey: messagesQueryKey(opcodeUrl, sessionID, directory),
144144
queryFn: async () => {
145145
const response = await client!.listMessages(sessionID!)
146146
return response as MessageWithParts[]
@@ -404,15 +404,15 @@ export const useSendPrompt = (opcodeUrl: string | null | undefined, directory?:
404404
);
405405
const userMessageInfo = createOptimisticUserMessageInfo(sessionID, optimisticUserID, model, agent, variant);
406406

407-
const messagesQueryKey = ["opencode", "messages", opcodeUrl, sessionID, directory];
408-
await queryClient.cancelQueries({ queryKey: messagesQueryKey });
407+
const queryKey = messagesQueryKey(opcodeUrl, sessionID, directory);
408+
await queryClient.cancelQueries({ queryKey });
409409

410410
const optimisticMessageWithParts: MessageWithParts = {
411411
info: userMessageInfo,
412412
parts: userMessageParts,
413413
}
414414
queryClient.setQueryData<MessageWithParts[]>(
415-
messagesQueryKey,
415+
queryKey,
416416
(old) => [...(old || []), optimisticMessageWithParts],
417417
);
418418

@@ -483,11 +483,11 @@ export const useSendPrompt = (opcodeUrl: string | null | undefined, directory?:
483483
},
484484
onError: (error, variables) => {
485485
const { sessionID } = variables;
486-
const messagesQueryKey = ["opencode", "messages", opcodeUrl, sessionID, directory];
486+
const queryKey = messagesQueryKey(opcodeUrl, sessionID, directory);
487487

488488
setSessionStatus(sessionID, { type: "idle" });
489489
queryClient.setQueryData<MessageWithParts[]>(
490-
messagesQueryKey,
490+
queryKey,
491491
(old) => old?.filter((msgWithParts) => !msgWithParts.info.id.startsWith("optimistic_")),
492492
);
493493

@@ -509,17 +509,17 @@ export const useSendPrompt = (opcodeUrl: string | null | undefined, directory?:
509509
onSuccess: async (data, variables) => {
510510
const { sessionID } = variables;
511511
const { response } = data;
512-
const messagesQueryKey = ["opencode", "messages", opcodeUrl, sessionID, directory];
512+
const queryKey = messagesQueryKey(opcodeUrl, sessionID, directory);
513513

514514
useSendErrorStore.getState().clearError(sessionID);
515515

516516
if (data.queued || !response) {
517-
queryClient.invalidateQueries({ queryKey: messagesQueryKey });
517+
queryClient.invalidateQueries({ queryKey });
518518
return;
519519
}
520520

521521
queryClient.setQueryData<MessageWithParts[]>(
522-
messagesQueryKey,
522+
queryKey,
523523
(old) => {
524524
if (!old) return old;
525525

@@ -553,7 +553,7 @@ export const useAbortSession = (
553553
const retryCountRef = useRef(0);
554554

555555
const forceCompleteMessages = useCallback((targetSessionID: string) => {
556-
const queryKey = ["opencode", "messages", opcodeUrl, targetSessionID, directory];
556+
const queryKey = messagesQueryKey(opcodeUrl, targetSessionID, directory);
557557
const now = Date.now();
558558

559559
queryClient.setQueryData<MessageWithParts[]>(queryKey, (old) => {
@@ -616,7 +616,7 @@ export const useAbortSession = (
616616
}, []);
617617

618618
const isSessionComplete = useCallback((targetSessionID: string) => {
619-
const queryKey = ["opencode", "messages", opcodeUrl, targetSessionID, directory];
619+
const queryKey = messagesQueryKey(opcodeUrl, targetSessionID, directory);
620620
const messages = queryClient.getQueryData<MessageWithParts[]>(queryKey);
621621

622622
const hasIncompleteMessages = messages?.some(msgWithParts => {
@@ -632,7 +632,7 @@ export const useAbortSession = (
632632
if (!sessionID) return;
633633

634634
const unsubscribe = queryClient.getQueryCache().subscribe((event) => {
635-
const queryKey = ["opencode", "messages", opcodeUrl, sessionID, directory];
635+
const queryKey = messagesQueryKey(opcodeUrl, sessionID, directory);
636636
if (event.query.queryKey.join(",") === queryKey.join(",")) {
637637
if (isSessionComplete(sessionID) && retryIntervalRef.current) {
638638
stopRetrying();
@@ -716,15 +716,15 @@ export const useSendShell = (opcodeUrl: string | null | undefined, directory?: s
716716
);
717717
const userMessageInfo = createOptimisticUserMessageInfo(sessionID, optimisticUserID);
718718

719-
const messagesQueryKey = ["opencode", "messages", opcodeUrl, sessionID, directory];
720-
await queryClient.cancelQueries({ queryKey: messagesQueryKey });
719+
const queryKey = messagesQueryKey(opcodeUrl, sessionID, directory);
720+
await queryClient.cancelQueries({ queryKey });
721721

722722
const optimisticMessageWithParts: MessageWithParts = {
723723
info: userMessageInfo,
724724
parts: userMessageParts,
725725
}
726726
queryClient.setQueryData<MessageWithParts[]>(
727-
messagesQueryKey,
727+
queryKey,
728728
(old) => [...(old || []), optimisticMessageWithParts],
729729
);
730730

@@ -739,7 +739,7 @@ export const useSendShell = (opcodeUrl: string | null | undefined, directory?: s
739739
const { sessionID } = variables;
740740
setSessionStatus(sessionID, { type: "idle" });
741741
queryClient.setQueryData<MessageWithParts[]>(
742-
["opencode", "messages", opcodeUrl, sessionID, directory],
742+
messagesQueryKey(opcodeUrl, sessionID, directory),
743743
(old) => {
744744
if (!old) return old;
745745
return old.filter((msgWithParts) => !msgWithParts.info.id.startsWith("optimistic_"));
@@ -753,7 +753,7 @@ export const useSendShell = (opcodeUrl: string | null | undefined, directory?: s
753753
const { optimisticUserID } = data;
754754

755755
queryClient.setQueryData<MessageWithParts[]>(
756-
["opencode", "messages", opcodeUrl, sessionID, directory],
756+
messagesQueryKey(opcodeUrl, sessionID, directory),
757757
(old) => {
758758
if (!old) return old;
759759
return old.filter((msgWithParts) => msgWithParts.info.id !== optimisticUserID);
@@ -810,7 +810,7 @@ export const useLoadSkill = (
810810
setSessionStatus(sessionID, { type: "busy" });
811811

812812
const optimisticUserID = `optimistic_user_${Date.now()}_${Math.random()}`;
813-
const messagesQueryKey = ["opencode", "messages", opcodeUrl, sessionID, directory];
813+
const queryKey = messagesQueryKey(opcodeUrl, sessionID, directory);
814814

815815
const userMessageParts = createOptimisticUserMessageParts(
816816
sessionID,
@@ -823,9 +823,9 @@ export const useLoadSkill = (
823823
parts: userMessageParts,
824824
};
825825

826-
await queryClient.cancelQueries({ queryKey: messagesQueryKey });
826+
await queryClient.cancelQueries({ queryKey });
827827
queryClient.setQueryData<MessageWithParts[]>(
828-
messagesQueryKey,
828+
queryKey,
829829
(old) => [...(old || []), optimisticMessageWithParts],
830830
);
831831

@@ -834,21 +834,21 @@ export const useLoadSkill = (
834834
},
835835
onError: (error) => {
836836
if (sessionID) {
837-
const messagesQueryKey = ["opencode", "messages", opcodeUrl, sessionID, directory];
837+
const queryKey = messagesQueryKey(opcodeUrl, sessionID, directory);
838838
setSessionStatus(sessionID!, { type: "idle" });
839839
queryClient.setQueryData<MessageWithParts[]>(
840-
messagesQueryKey,
840+
queryKey,
841841
(old) => old?.filter((m) => !m.info.id.startsWith("optimistic_")),
842842
);
843843
}
844844
showToast.error(error instanceof Error ? error.message : "Failed to load skill");
845845
},
846846
onSuccess: (data) => {
847847
const { response } = data;
848-
const messagesQueryKey = ["opencode", "messages", opcodeUrl, sessionID, directory];
848+
const queryKey = messagesQueryKey(opcodeUrl, sessionID, directory);
849849

850850
queryClient.setQueryData<MessageWithParts[]>(
851-
messagesQueryKey,
851+
queryKey,
852852
(old) => {
853853
if (!old) return old;
854854

frontend/src/hooks/useRemoveMessage.ts

Lines changed: 7 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import { useMutation, useQueryClient } from '@tanstack/react-query'
22
import { createOpenCodeClient } from '@/api/opencode'
33
import { showToast } from '@/lib/toast'
4+
import { messagesQueryKey } from '@/lib/queryInvalidation'
45
import type { Message, Part, MessageWithParts } from '@/api/types'
56
import { useSessionStatus } from '@/stores/sessionStatusStore'
67

@@ -25,7 +26,7 @@ export function useRemoveMessage({ opcodeUrl, sessionId, directory }: UseRemoveM
2526
return client.revertMessage(sessionId, { messageID, partID })
2627
},
2728
onMutate: async ({ messageID }) => {
28-
const queryKey = ['opencode', 'messages', opcodeUrl, sessionId, directory]
29+
const queryKey = messagesQueryKey(opcodeUrl, sessionId, directory)
2930

3031
await queryClient.cancelQueries({ queryKey })
3132

@@ -44,7 +45,7 @@ export function useRemoveMessage({ opcodeUrl, sessionId, directory }: UseRemoveM
4445
onError: (_error, _variables, _context: RemoveMessageContext | undefined) => {
4546
if (_context?.previousMessages) {
4647
queryClient.setQueryData(
47-
['opencode', 'messages', opcodeUrl, sessionId, directory],
48+
messagesQueryKey(opcodeUrl, sessionId, directory),
4849
_context.previousMessages
4950
)
5051
}
@@ -53,7 +54,7 @@ export function useRemoveMessage({ opcodeUrl, sessionId, directory }: UseRemoveM
5354
},
5455
onSuccess: () => {
5556
queryClient.invalidateQueries({
56-
queryKey: ['opencode', 'messages', opcodeUrl, sessionId, directory]
57+
queryKey: messagesQueryKey(opcodeUrl, sessionId, directory)
5758
})
5859
queryClient.invalidateQueries({
5960
queryKey: ['opencode', 'session', opcodeUrl, sessionId, directory]
@@ -115,7 +116,7 @@ export function useRefreshMessage({ opcodeUrl, sessionId, directory }: UseRefres
115116
}
116117

117118
queryClient.setQueryData<MessageWithParts[]>(
118-
['opencode', 'messages', opcodeUrl, sessionId, directory],
119+
messagesQueryKey(opcodeUrl, sessionId, directory),
119120
(old) => [...(old || []), optimisticMessageWithParts]
120121
)
121122

@@ -146,7 +147,7 @@ export function useRefreshMessage({ opcodeUrl, sessionId, directory }: UseRefres
146147
},
147148
onSuccess: () => {
148149
queryClient.invalidateQueries({
149-
queryKey: ['opencode', 'messages', opcodeUrl, sessionId, directory]
150+
queryKey: messagesQueryKey(opcodeUrl, sessionId, directory)
150151
})
151152
queryClient.invalidateQueries({
152153
queryKey: ['opencode', 'session', opcodeUrl, sessionId, directory]
@@ -156,7 +157,7 @@ export function useRefreshMessage({ opcodeUrl, sessionId, directory }: UseRefres
156157
void variables
157158
setSessionStatus(sessionId, { type: 'idle' })
158159
queryClient.setQueryData<MessageWithParts[]>(
159-
['opencode', 'messages', opcodeUrl, sessionId, directory],
160+
messagesQueryKey(opcodeUrl, sessionId, directory),
160161
(old) => {
161162
const messages = old || []
162163
const optimisticIndex = messages.findIndex((m) => m.info.id.startsWith('optimistic_user_'))

0 commit comments

Comments
 (0)