mirror of
https://github.com/anomalyco/opencode.git
synced 2026-09-07 23:24:35 +00:00
fix(desktop): resend unacknowledged browser state and correlate network hops
This commit is contained in:
parent
d84be2d7ef
commit
89a6eb4e18
4 changed files with 112 additions and 36 deletions
|
|
@ -4,7 +4,7 @@ import { Browser } from "@opencode-ai/plugin-browser/rpc"
|
|||
import { OpenCode } from "@opencode-ai/client/effect"
|
||||
import { SessionID } from "@opencode-ai/schema/session-id"
|
||||
import type { BrowserWindow } from "electron"
|
||||
import { Deferred, Effect, ManagedRuntime, Queue, Schema, Stream } from "effect"
|
||||
import { Deferred, Effect, ManagedRuntime, Queue, Schedule, Schema, Stream } from "effect"
|
||||
import { HttpClient, HttpClientRequest } from "effect/unstable/http"
|
||||
import { BrowserPaneEvent } from "../shared/ipc-rpc/events"
|
||||
import { createBrowserPage, type BrowserPage } from "./browser-chromium"
|
||||
|
|
@ -103,13 +103,17 @@ export function createBrowserPane() {
|
|||
outbound,
|
||||
effect.pipe(Effect.catchCause((cause) => Effect.logWarning("Browser send failed", cause))),
|
||||
)
|
||||
// Report state before publishing it locally or completing a command.
|
||||
// Report state before publishing it locally or completing a command. The server's copy of
|
||||
// the inventory resolves every tab ID, so a state that never arrived must not count as published.
|
||||
entry.report = (event) => {
|
||||
const local = Effect.sync(() => publish(entry, event))
|
||||
if (event.type !== "state") return send(local)
|
||||
send(
|
||||
(event.type === "state"
|
||||
? rpc.state({ ...attachment, state: event.state ?? { tabs: [], focusedTabID: null } }, options)
|
||||
: Effect.void
|
||||
).pipe(Effect.ensuring(Effect.sync(() => publish(entry, event)))),
|
||||
rpc.state({ ...attachment, state: event.state ?? { tabs: [], focusedTabID: null } }, options).pipe(
|
||||
Effect.retry({ times: 4, schedule: Schedule.exponential("250 millis") }),
|
||||
Effect.tapError(() => Effect.sync(() => (entry.lastState = undefined))),
|
||||
Effect.ensuring(local),
|
||||
),
|
||||
)
|
||||
}
|
||||
const receive = client.event.subscribe().pipe(
|
||||
|
|
@ -159,12 +163,11 @@ export function createBrowserPane() {
|
|||
}),
|
||||
),
|
||||
Effect.ensuring(Effect.sync(() => entry.requests.delete(message.requestID))),
|
||||
// Only a server cancel or tab close aborts the signal; any other failure means
|
||||
// the command could not be retrieved, and the attachment is no longer trustworthy.
|
||||
// A request that vanished (cancelled before retrieval) or an operation this desktop
|
||||
// cannot decode fails only that request; the server times it out. Transport loss
|
||||
// surfaces through the event stream and attach call instead.
|
||||
Effect.catchCause((cause) =>
|
||||
abort.signal.aborted
|
||||
? Effect.void
|
||||
: Effect.logError("Browser command failed", cause).pipe(Effect.andThen(Effect.sync(stop))),
|
||||
abort.signal.aborted ? Effect.void : Effect.logWarning("Browser command failed", cause),
|
||||
),
|
||||
Effect.forkScoped,
|
||||
)
|
||||
|
|
|
|||
|
|
@ -6,20 +6,27 @@ type Request = {
|
|||
info: Browser.NetworkRequest
|
||||
nativeID: string
|
||||
sessionID?: string
|
||||
started: number
|
||||
/** Monotonic start; unknown for a WebSocket until its handshake is sent. */
|
||||
started?: number
|
||||
request: Pick<Protocol.Network.Request, "headers" | "postData" | "hasPostData">
|
||||
response?: Pick<Protocol.Network.Response, "headers" | "mimeType">
|
||||
headersTruncated: boolean
|
||||
postDataTruncated: boolean
|
||||
redirected?: boolean
|
||||
/** ExtraInfo arrives once per redirect hop, in order; these mark which hops consumed theirs. */
|
||||
wire: { request: boolean; response: boolean }
|
||||
}
|
||||
type Wire = { headers: Protocol.Network.Headers; statusCode: number }
|
||||
const levels = ["debug", "info", "warning", "error"] as const
|
||||
// Values the model must not read; the header name still shows it was sent.
|
||||
const redacted = new Set(["cookie", "set-cookie", "authorization", "proxy-authorization"])
|
||||
|
||||
export function createDiagnostics(cdp: Cdp) {
|
||||
const messages: Browser.ConsoleEntry[] = []
|
||||
const requests = new Map<string, Request>()
|
||||
const current = new Map<string, Request>()
|
||||
// Chromium reuses one request ID across a redirect chain; every hop is retained in order.
|
||||
const hops = new Map<string, Request[]>()
|
||||
const latest = (key: string) => hops.get(key)?.at(-1)
|
||||
// Network-stack headers and wire status arrive as ExtraInfo events, in either order relative to
|
||||
// the renderer-side events; stash whichever comes first.
|
||||
const extra = new Map<string, { request?: Protocol.Network.Headers; response?: Wire }>()
|
||||
|
|
@ -84,7 +91,7 @@ export function createDiagnostics(cdp: Cdp) {
|
|||
key: string,
|
||||
sessionID: string | undefined,
|
||||
nativeID: string,
|
||||
started: number,
|
||||
started: number | undefined,
|
||||
wallTime: number,
|
||||
info: Pick<Browser.NetworkRequest, "url" | "method" | "resourceType">,
|
||||
request: Request["request"],
|
||||
|
|
@ -97,6 +104,7 @@ export function createDiagnostics(cdp: Cdp) {
|
|||
request: { ...request, headers: headers.headers, postData: request.postData?.slice(0, 20_000) },
|
||||
headersTruncated: headers.truncated,
|
||||
postDataTruncated: (request.postData?.length ?? 0) > 20_000,
|
||||
wire: { request: false, response: false },
|
||||
info: {
|
||||
id: `${scope}:${++sequence}`,
|
||||
...info,
|
||||
|
|
@ -106,7 +114,8 @@ export function createDiagnostics(cdp: Cdp) {
|
|||
},
|
||||
}
|
||||
requests.set(entry.info.id, entry)
|
||||
current.set(key, entry)
|
||||
hops.set(key, [...(hops.get(key) ?? []), entry])
|
||||
// A stash can only belong to this hop: earlier hops would have consumed it on arrival.
|
||||
const stashed = extra.get(key)
|
||||
if (stashed?.request) requestHeaders(entry, stashed.request)
|
||||
if (stashed?.response) responseInfo(entry, stashed.response)
|
||||
|
|
@ -115,8 +124,10 @@ export function createDiagnostics(cdp: Cdp) {
|
|||
const first = requests.values().next().value
|
||||
if (first) {
|
||||
requests.delete(first.info.id)
|
||||
if (current.get(`${first.sessionID ?? ""}:${first.nativeID}`) === first)
|
||||
current.delete(`${first.sessionID ?? ""}:${first.nativeID}`)
|
||||
const firstKey = `${first.sessionID ?? ""}:${first.nativeID}`
|
||||
const rest = hops.get(firstKey)?.filter((hop) => hop !== first) ?? []
|
||||
if (rest.length) hops.set(firstKey, rest)
|
||||
if (!rest.length) hops.delete(firstKey)
|
||||
}
|
||||
droppedRequests++
|
||||
}
|
||||
|
|
@ -126,17 +137,19 @@ export function createDiagnostics(cdp: Cdp) {
|
|||
const trimmed = trimHeaders({ ...request.request.headers, ...headers })
|
||||
request.request = { ...request.request, headers: trimmed.headers }
|
||||
request.headersTruncated ||= trimmed.truncated
|
||||
request.wire.request = true
|
||||
}
|
||||
const responseInfo = (request: Request, wire: Wire) => {
|
||||
const trimmed = trimHeaders({ ...request.response?.headers, ...wire.headers })
|
||||
request.response = { mimeType: request.response?.mimeType ?? "", headers: trimmed.headers }
|
||||
request.headersTruncated ||= trimmed.truncated
|
||||
request.info = { ...request.info, statusCode: wire.statusCode }
|
||||
request.wire.response = true
|
||||
}
|
||||
const finish = (key: string, timestamp: number, failure?: string) => {
|
||||
const request = current.get(key)
|
||||
const request = latest(key)
|
||||
if (!request) return
|
||||
const durationMs = Math.max(0, (timestamp - request.started) * 1000)
|
||||
const durationMs = request.started === undefined ? 0 : Math.max(0, (timestamp - request.started) * 1000)
|
||||
request.info =
|
||||
failure === undefined
|
||||
? { ...request.info, state: "completed", durationMs }
|
||||
|
|
@ -144,17 +157,18 @@ export function createDiagnostics(cdp: Cdp) {
|
|||
}
|
||||
cdp.on("Network.requestWillBeSent", (event, sessionID) => {
|
||||
const key = `${sessionID ?? ""}:${event.requestId}`
|
||||
const previous = current.get(key)
|
||||
const previous = latest(key)
|
||||
if (previous && event.redirectResponse) {
|
||||
const response = trimHeaders(event.redirectResponse.headers)
|
||||
// Wire headers for this hop may already be merged in; the renderer copy lacks set-cookie.
|
||||
const response = trimHeaders({ ...event.redirectResponse.headers, ...previous.response?.headers })
|
||||
previous.response = { mimeType: event.redirectResponse.mimeType, headers: response.headers }
|
||||
previous.headersTruncated ||= response.truncated
|
||||
previous.redirected = true
|
||||
previous.info = {
|
||||
...previous.info,
|
||||
state: "completed",
|
||||
statusCode: event.redirectResponse.status,
|
||||
durationMs: Math.max(0, (event.timestamp - previous.started) * 1000),
|
||||
statusCode: previous.info.statusCode ?? event.redirectResponse.status,
|
||||
durationMs: Math.max(0, (event.timestamp - (previous.started ?? event.timestamp)) * 1000),
|
||||
}
|
||||
}
|
||||
begin(
|
||||
|
|
@ -171,14 +185,16 @@ export function createDiagnostics(cdp: Cdp) {
|
|||
{ headers: event.request.headers, hasPostData: event.request.hasPostData, postData: event.request.postData },
|
||||
)
|
||||
})
|
||||
// ExtraInfo for a redirect hop can land after the renderer already started the next hop, so it
|
||||
// goes to the earliest hop that has not consumed its own rather than to the latest.
|
||||
cdp.on("Network.requestWillBeSentExtraInfo", (event, sessionID) => {
|
||||
const key = `${sessionID ?? ""}:${event.requestId}`
|
||||
const request = current.get(key)
|
||||
const request = hops.get(key)?.find((hop) => !hop.wire.request)
|
||||
if (request) return requestHeaders(request, event.headers)
|
||||
extra.set(key, { ...extra.get(key), request: event.headers })
|
||||
})
|
||||
cdp.on("Network.responseReceived", (event, sessionID) => {
|
||||
const request = current.get(`${sessionID ?? ""}:${event.requestId}`)
|
||||
const request = latest(`${sessionID ?? ""}:${event.requestId}`)
|
||||
if (!request) return
|
||||
const headers = trimHeaders({ ...event.response.headers, ...request.response?.headers })
|
||||
request.response = { mimeType: event.response.mimeType, headers: headers.headers }
|
||||
|
|
@ -188,7 +204,7 @@ export function createDiagnostics(cdp: Cdp) {
|
|||
})
|
||||
cdp.on("Network.responseReceivedExtraInfo", (event, sessionID) => {
|
||||
const key = `${sessionID ?? ""}:${event.requestId}`
|
||||
const request = current.get(key)
|
||||
const request = hops.get(key)?.find((hop) => !hop.wire.response)
|
||||
if (request) return responseInfo(request, event)
|
||||
extra.set(key, { ...extra.get(key), response: event })
|
||||
})
|
||||
|
|
@ -198,14 +214,14 @@ export function createDiagnostics(cdp: Cdp) {
|
|||
`${sessionID ?? ""}:${event.requestId}`,
|
||||
sessionID,
|
||||
event.requestId,
|
||||
0,
|
||||
undefined,
|
||||
Date.now() / 1000,
|
||||
{ url: event.url, method: "GET", resourceType: "websocket" },
|
||||
{ headers: {}, hasPostData: false },
|
||||
)
|
||||
})
|
||||
cdp.on("Network.webSocketWillSendHandshakeRequest", (event, sessionID) => {
|
||||
const request = current.get(`${sessionID ?? ""}:${event.requestId}`)
|
||||
const request = latest(`${sessionID ?? ""}:${event.requestId}`)
|
||||
if (!request) return
|
||||
request.started = event.timestamp
|
||||
request.info = { ...request.info, timestampMs: event.wallTime * 1000 }
|
||||
|
|
@ -213,7 +229,7 @@ export function createDiagnostics(cdp: Cdp) {
|
|||
})
|
||||
cdp.on("Network.webSocketHandshakeResponseReceived", (event, sessionID) => {
|
||||
const key = `${sessionID ?? ""}:${event.requestId}`
|
||||
const request = current.get(key)
|
||||
const request = latest(key)
|
||||
if (!request) return
|
||||
responseInfo(request, { headers: event.response.headers, statusCode: event.response.status })
|
||||
finish(key, event.timestamp)
|
||||
|
|
@ -221,6 +237,10 @@ export function createDiagnostics(cdp: Cdp) {
|
|||
cdp.on("Network.webSocketFrameError", (event, sessionID) =>
|
||||
finish(`${sessionID ?? ""}:${event.requestId}`, event.timestamp, event.errorMessage),
|
||||
)
|
||||
cdp.on("Network.webSocketClosed", (event, sessionID) => {
|
||||
const key = `${sessionID ?? ""}:${event.requestId}`
|
||||
if (latest(key)?.info.state === "pending") finish(key, event.timestamp, "Closed before the handshake completed")
|
||||
})
|
||||
cdp.on("Network.loadingFinished", (event, sessionID) =>
|
||||
finish(`${sessionID ?? ""}:${event.requestId}`, event.timestamp),
|
||||
)
|
||||
|
|
@ -231,7 +251,7 @@ export function createDiagnostics(cdp: Cdp) {
|
|||
clear() {
|
||||
messages.length = 0
|
||||
requests.clear()
|
||||
current.clear()
|
||||
hops.clear()
|
||||
extra.clear()
|
||||
droppedMessages = 0
|
||||
droppedRequests = 0
|
||||
|
|
@ -346,8 +366,9 @@ function trimHeaders(headers: Protocol.Network.Headers) {
|
|||
let size = 0
|
||||
let truncated = entries.length > 100
|
||||
for (const [key, value] of entries.slice(0, 100)) {
|
||||
const name = key.slice(0, 2_048)
|
||||
const text = String(value)
|
||||
// Lower-case names so renderer and wire copies of one header merge instead of duplicating.
|
||||
const name = key.toLowerCase().slice(0, 2_048)
|
||||
const text = redacted.has(name) ? "<redacted>" : String(value)
|
||||
const remaining = Math.max(0, 16_000 - size - name.length)
|
||||
if (!remaining) {
|
||||
truncated = true
|
||||
|
|
|
|||
|
|
@ -47,6 +47,7 @@ test.skipIf(!!process.env.CI)(
|
|||
let native: ReturnType<typeof Bun.spawn> | undefined
|
||||
let proxy: Bun.Server<undefined> | undefined
|
||||
let tunnels = 0
|
||||
let rejectedStates = 0
|
||||
try {
|
||||
const url = await Promise.race([
|
||||
ready.promise,
|
||||
|
|
@ -56,6 +57,8 @@ test.skipIf(!!process.env.CI)(
|
|||
])
|
||||
// Exercise the real HTTP boundary with delayed state acknowledgments, as on
|
||||
// a remote server. Neither endpoint can rely on synchronous UI/state updates.
|
||||
// The first inventory that carries a tab is rejected twice: the desktop must
|
||||
// keep it pending and resend, or the server never learns the tab exists.
|
||||
proxy = Bun.serve({
|
||||
hostname: "127.0.0.1",
|
||||
port: 0,
|
||||
|
|
@ -63,8 +66,13 @@ test.skipIf(!!process.env.CI)(
|
|||
async fetch(request) {
|
||||
const incoming = new URL(request.url)
|
||||
if (incoming.pathname.endsWith("/tunnel.open")) tunnels++
|
||||
if (incoming.pathname.endsWith("/experimental.browser/state"))
|
||||
if (incoming.pathname.endsWith("/experimental.browser/state")) {
|
||||
if (rejectedStates < 2 && (await request.clone().text()).includes('"tab_')) {
|
||||
rejectedStates++
|
||||
return new Response(null, { status: 503 })
|
||||
}
|
||||
await new Promise((resolve) => setTimeout(resolve, 75))
|
||||
}
|
||||
return fetch(new Request(new URL(incoming.pathname + incoming.search, url), request), { decompress: false })
|
||||
},
|
||||
})
|
||||
|
|
@ -90,6 +98,7 @@ test.skipIf(!!process.env.CI)(
|
|||
]),
|
||||
).toBe(0)
|
||||
expect(tunnels).toBeGreaterThan(0)
|
||||
expect(rejectedStates).toBe(2)
|
||||
} finally {
|
||||
if (native?.exitCode === null) {
|
||||
native.kill()
|
||||
|
|
|
|||
|
|
@ -63,6 +63,17 @@ async function main() {
|
|||
response.end('<iframe name="inner" title="Inner frame" src="/frame"></iframe>')
|
||||
return
|
||||
}
|
||||
// One Chromium request ID spans both hops; each hop must keep its own wire headers.
|
||||
if (request.url === "/redirect") {
|
||||
response.writeHead(302, { location: "/redirected", "set-cookie": "hop=1; Path=/" })
|
||||
response.end()
|
||||
return
|
||||
}
|
||||
if (request.url === "/redirected") {
|
||||
response.writeHead(200, { "content-type": "text/plain" })
|
||||
response.end("landed")
|
||||
return
|
||||
}
|
||||
// Revalidation: the wire answers 304 while the renderer reports the cached 200.
|
||||
if (request.url === "/etag") {
|
||||
if (request.headers["if-none-match"] === '"v1"') {
|
||||
|
|
@ -182,6 +193,8 @@ async function main() {
|
|||
websocket: await new Promise((resolve, reject) => { const socket = new WebSocket('ws://localhost:${address.port}/hmr'); socket.onmessage = event => { resolve(event.data); socket.close(); }; socket.onerror = () => reject(new Error('WebSocket failed')); }),
|
||||
cookie: (document.cookie = 'wire=1', await fetch('/api/test?cookie').then(response => response.ok)),
|
||||
etag: [await fetch('/etag').then(response => response.status), await fetch('/etag').then(response => response.status)],
|
||||
redirect: await fetch('/redirect').then(response => response.text()),
|
||||
refused: await new Promise((resolve) => { const socket = new WebSocket('ws://127.0.0.1:1/refused'); socket.onerror = () => resolve('refused'); socket.onopen = () => resolve('opened'); }),
|
||||
}))()`,
|
||||
})
|
||||
assert.deepEqual(networkProof.value, {
|
||||
|
|
@ -191,10 +204,38 @@ async function main() {
|
|||
websocket: "rpc websocket proof",
|
||||
cookie: true,
|
||||
etag: [200, 200],
|
||||
redirect: "landed",
|
||||
refused: "refused",
|
||||
})
|
||||
const sockets = await call("network.list", { tabID, resourceType: "websocket" })
|
||||
assert.equal(sockets.requests[0]?.statusCode, 101, JSON.stringify(sockets))
|
||||
assert.equal(sockets.requests[0]?.state, "completed")
|
||||
const opened = sockets.requests.find((request) => request.url.includes("/hmr"))
|
||||
const refused = sockets.requests.find((request) => request.url.includes("/refused"))
|
||||
assert.equal(opened?.statusCode, 101, JSON.stringify(sockets))
|
||||
assert.equal(opened?.state, "completed")
|
||||
assert.equal(refused?.state, "failed", JSON.stringify(refused))
|
||||
assert(
|
||||
refused?.state === "failed" && refused.durationMs >= 0 && refused.durationMs < 60_000,
|
||||
JSON.stringify(refused),
|
||||
)
|
||||
const chain = await call("network.list", { tabID, urlContains: "/redirect" })
|
||||
assert.deepEqual(
|
||||
chain.requests.map((request) => [request.url.endsWith("/redirect"), request.statusCode]),
|
||||
[
|
||||
[true, 302],
|
||||
[false, 200],
|
||||
],
|
||||
JSON.stringify(chain),
|
||||
)
|
||||
const hop = await call("network.get", { tabID, id: chain.requests[0].id })
|
||||
const landed = await call("network.get", { tabID, id: chain.requests[1].id })
|
||||
const names = (headers: { name: string }[]) => headers.map((header) => header.name)
|
||||
assert(names(hop.responseHeaders).includes("set-cookie"), JSON.stringify(hop.responseHeaders))
|
||||
assert(names(hop.responseHeaders).includes("location"), JSON.stringify(hop.responseHeaders))
|
||||
assert(!names(landed.responseHeaders).includes("location"), JSON.stringify(landed.responseHeaders))
|
||||
assert(
|
||||
hop.responseHeaders.every((header) => header.name !== "set-cookie" || header.value === "<redacted>"),
|
||||
JSON.stringify(hop.responseHeaders),
|
||||
)
|
||||
const revalidated = await call("network.list", { tabID, urlContains: "/etag" })
|
||||
assert.deepEqual(
|
||||
revalidated.requests.map((request) => request.statusCode),
|
||||
|
|
@ -203,10 +244,12 @@ async function main() {
|
|||
)
|
||||
const withCookie = await call("network.list", { tabID, urlContains: "/api/test?cookie" })
|
||||
const cookieDetail = await call("network.get", { tabID, id: withCookie.requests[0].id })
|
||||
// The wire header proves the cookie was sent; its value never reaches the model.
|
||||
assert(
|
||||
cookieDetail.requestHeaders.some((header) => header.name.toLowerCase() === "cookie" && header.value === "wire=1"),
|
||||
cookieDetail.requestHeaders.some((header) => header.name === "cookie" && header.value === "<redacted>"),
|
||||
JSON.stringify(cookieDetail.requestHeaders),
|
||||
)
|
||||
assert(!cookieDetail.requestHeaders.some((header) => header.name !== header.name.toLowerCase()))
|
||||
const upload = await rpc.write({ text: "server upload bytes" }, { location })
|
||||
await fails("trace.stop", { tabID }, /browser\.trace\.start/)
|
||||
await fails("cpu.stop", { tabID }, /browser\.cpu\.start/)
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue