mirror of
https://github.com/Sea-Haven-Industries/open-swe.git
synced 2026-09-30 11:33:14 +00:00
fix: stream proxy preflight + optimistic thread staleTime (#1502)
* fix: preflight stream proxy before SSE starts, let optimistic thread survive refetch * fix: seed sidebar thread details as stale so mark-viewed fetch still fires --------- Co-authored-by: open-swe[bot] <open-swe@users.noreply.github.com>
This commit is contained in:
parent
8fa98398fa
commit
b8035a4758
4 changed files with 29 additions and 14 deletions
|
|
@ -996,19 +996,15 @@ async def api_thread_stream_events(
|
|||
session: dict[str, Any] = _SESSION_DEP,
|
||||
) -> StreamingResponse:
|
||||
body = await request.body()
|
||||
|
||||
async def event_generator():
|
||||
async for chunk in proxy_dashboard_thread_stream_events(
|
||||
thread_id,
|
||||
session["sub"],
|
||||
body,
|
||||
email=session.get("email"),
|
||||
content_type=request.headers.get("content-type", "application/json"),
|
||||
):
|
||||
yield chunk
|
||||
|
||||
stream = await proxy_dashboard_thread_stream_events(
|
||||
thread_id,
|
||||
session["sub"],
|
||||
body,
|
||||
email=session.get("email"),
|
||||
content_type=request.headers.get("content-type", "application/json"),
|
||||
)
|
||||
return StreamingResponse(
|
||||
event_generator(),
|
||||
stream,
|
||||
media_type="text/event-stream",
|
||||
headers={"Cache-Control": "no-cache", "Connection": "keep-alive"},
|
||||
)
|
||||
|
|
|
|||
|
|
@ -1032,8 +1032,18 @@ async def proxy_dashboard_thread_stream_events(
|
|||
email: str | None = None,
|
||||
content_type: str = "application/json",
|
||||
) -> AsyncIterator[bytes]:
|
||||
# Preflight here (not in the generator) so auth/content-type failures
|
||||
# surface as real HTTP errors before the SSE response starts streaming.
|
||||
_require_json_content_type(content_type)
|
||||
await _authorized_thread_metadata(thread_id, login, email=email)
|
||||
return _stream_thread_events(thread_id, body, content_type)
|
||||
|
||||
|
||||
async def _stream_thread_events(
|
||||
thread_id: str,
|
||||
body: bytes,
|
||||
content_type: str,
|
||||
) -> AsyncIterator[bytes]:
|
||||
url = f"{langgraph_url().rstrip('/')}/threads/{thread_id}/stream/events"
|
||||
headers = _langgraph_proxy_headers(content_type=content_type, accept="text/event-stream")
|
||||
|
||||
|
|
|
|||
|
|
@ -424,7 +424,7 @@ async def test_proxy_endpoints_enforce_thread_ownership(monkeypatch) -> None:
|
|||
assert exc_info.value.status_code == 404
|
||||
|
||||
with pytest.raises(HTTPException) as exc_info:
|
||||
await anext(thread_api.proxy_dashboard_thread_stream_events("tid", "intruder", b"{}"))
|
||||
await thread_api.proxy_dashboard_thread_stream_events("tid", "intruder", b"{}")
|
||||
assert exc_info.value.status_code == 404
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -31,7 +31,12 @@ export function useSeedAgentThreadDetails(
|
|||
useEffect(() => {
|
||||
for (const thread of threads) {
|
||||
if (thread.id === activeThreadId) continue
|
||||
queryClient.setQueryData(agentThreadKeys.detail(thread.id), thread)
|
||||
// Seed as already-stale: the detail GET is what marks a thread viewed
|
||||
// server-side, so opening a seeded entry must still refetch despite the
|
||||
// detail query's `staleTime` (which exists for the optimistic seed).
|
||||
queryClient.setQueryData(agentThreadKeys.detail(thread.id), thread, {
|
||||
updatedAt: 0,
|
||||
})
|
||||
}
|
||||
}, [activeThreadId, queryClient, threads])
|
||||
}
|
||||
|
|
@ -51,6 +56,10 @@ export function useAgentThread(threadId: string) {
|
|||
return useQuery({
|
||||
queryKey: agentThreadKeys.detail(threadId),
|
||||
queryFn: () => agentsApi.getThread(threadId),
|
||||
// Lets the optimistic detail seeded by `AgentsHome` survive until the
|
||||
// proxied run.start stamps the server-side thread; an immediate refetch
|
||||
// would 404 and bounce the route back to /agents.
|
||||
staleTime: 30_000,
|
||||
})
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue