Skip to content

Commit 21a70bf

Browse files
fix(llminternal): classify a dropped live connection by type, not error text
RunLive's reader and sender goroutines both report into the same errChan and the flow acts on whichever arrives first, so the two must agree on whether a dropped connection is resumable. They did not agree on Windows. isResumable matched substrings. Its list covers the reader's websocket close text ("close 1006 ... unexpected EOF") and the POSIX sender text ("write: broken pipe"), but not the Windows sender text ("wsasend: An established connection was aborted by the software in your host machine."). When the sender won the race on Windows, the same connection loss that would have resumed was pushed to the caller as fatal and the session stopped after one connection. errors.Is against the POSIX constants does not close the gap, because Go does not map WSA error numbers onto them. The chain is *net.OpError -> *os.SyscallError -> syscall.Errno(10053). Match the transport failure by type instead. Any *net.OpError reaching this path is a failed read or write on an already-established live socket, since the dial is handled separately above, so it is resumable on every platform. The substring checks stay for the websocket-level cases (1006, 1008, GoAway), which are not *net.OpError. isResumable moves from a closure inside RunLive to a package-level function so it can be tested directly. Behaviour is otherwise unchanged. Fixes #1602
1 parent b531451 commit 21a70bf

2 files changed

Lines changed: 201 additions & 15 deletions

File tree

‎internal/llminternal/base_flow.go‎

Lines changed: 35 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@ import (
2424
"iter"
2525
"log"
2626
"maps"
27+
"net"
2728
"slices"
2829
"strings"
2930
"sync"
@@ -347,6 +348,40 @@ func waitBeforeReconnect(ctx context.Context, sess *liveSessionImpl, d time.Dura
347348
}
348349
}
349350

351+
// isResumable reports whether a live-connection error means the socket is
352+
// gone and the flow should reconnect, rather than surface the error and stop.
353+
//
354+
// Both the reader and the sender goroutine report into the same errChan and
355+
// the flow acts on whichever arrives first, so the two must classify the same
356+
// connection loss the same way. They do not produce the same text: the reader
357+
// sees the websocket close ("close 1006 ... unexpected EOF"), while the sender
358+
// sees the raw socket write failure, whose wording is platform-specific
359+
// ("write: broken pipe" on Linux, "wsasend: An established connection was
360+
// aborted by the software in your host machine." on Windows). Matching the
361+
// transport failure by type rather than by text keeps the verdict the same on
362+
// every platform.
363+
func isResumable(err error) bool {
364+
if err == nil {
365+
return false
366+
}
367+
if err == io.EOF {
368+
return true
369+
}
370+
// A failed read or write on the underlying socket of an already-established
371+
// connection. The dial lives on a different path, which handles its own
372+
// errors, so reaching here means a live connection dropped.
373+
var opErr *net.OpError
374+
if errors.As(err, &opErr) {
375+
return true
376+
}
377+
errStr := err.Error()
378+
return strings.Contains(errStr, "broken pipe") ||
379+
strings.Contains(errStr, "connection reset") ||
380+
strings.Contains(errStr, "EOF") ||
381+
strings.Contains(errStr, "1008") ||
382+
strings.Contains(errStr, "GoAway")
383+
}
384+
350385
func (f *Flow) RunLive(ctx agent.InvocationContext) (agent.LiveSession, iter.Seq2[*session.Event, error], error) {
351386
clientProvider, ok := f.Model.(interface {
352387
Client() *genai.Client
@@ -393,21 +428,6 @@ func (f *Flow) RunLive(ctx agent.InvocationContext) (agent.LiveSession, iter.Seq
393428
OutputAudioTranscription: runCfg.Live.OutputAudioTranscription,
394429
}
395430

396-
isResumable := func(err error) bool {
397-
if err == nil {
398-
return false
399-
}
400-
if err == io.EOF {
401-
return true
402-
}
403-
errStr := err.Error()
404-
return strings.Contains(errStr, "broken pipe") ||
405-
strings.Contains(errStr, "connection reset") ||
406-
strings.Contains(errStr, "EOF") ||
407-
strings.Contains(errStr, "1008") ||
408-
strings.Contains(errStr, "GoAway")
409-
}
410-
411431
iCtx, isIContext := ctx.(*icontext.InvocationContext)
412432

413433
policy := f.reconnect
Lines changed: 166 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,166 @@
1+
// Copyright 2026 Google LLC
2+
//
3+
// Licensed under the Apache License, Version 2.0 (the "License");
4+
// you may not use this file except in compliance with the License.
5+
// You may obtain a copy of the License at
6+
//
7+
// http://www.apache.org/licenses/LICENSE-2.0
8+
//
9+
// Unless required by applicable law or agreed to in writing, software
10+
// distributed under the License is distributed on an "AS IS" BASIS,
11+
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
// See the License for the specific language governing permissions and
13+
// limitations under the License.
14+
15+
package llminternal
16+
17+
import (
18+
"errors"
19+
"fmt"
20+
"io"
21+
"net"
22+
"os"
23+
"syscall"
24+
"testing"
25+
)
26+
27+
// socketErr rebuilds the error the net package returns for a failed read or
28+
// write on a live socket: *net.OpError wrapping *os.SyscallError wrapping the
29+
// platform errno. The errno value is what differs across platforms, so each
30+
// case passes one explicitly rather than relying on the host's own spelling.
31+
func socketErr(op, syscallName string, errno syscall.Errno) error {
32+
return &net.OpError{
33+
Op: op,
34+
Net: "tcp",
35+
Source: &net.TCPAddr{IP: net.IPv4(127, 0, 0, 1), Port: 1000},
36+
Addr: &net.TCPAddr{IP: net.IPv4(127, 0, 0, 1), Port: 1001},
37+
Err: os.NewSyscallError(syscallName, errno),
38+
}
39+
}
40+
41+
// Windows sockets report a dropped connection with WSA error numbers, which Go
42+
// does not map onto the POSIX ECONN* constants. Spelled numerically so the test
43+
// compiles and means the same thing on every platform.
44+
const (
45+
wsaeConnAborted = syscall.Errno(10053)
46+
wsaeConnReset = syscall.Errno(10054)
47+
)
48+
49+
func TestIsResumable(t *testing.T) {
50+
tests := []struct {
51+
name string
52+
err error
53+
want bool
54+
}{
55+
{"nil is not resumable", nil, false},
56+
{"io.EOF", io.EOF, true},
57+
58+
// The reader path. The websocket layer turns a dropped connection into
59+
// a close error whose text carries the code, on every platform.
60+
{
61+
name: "reader: abnormal closure 1006",
62+
err: errors.New("failed to receive message: websocket: close 1006 (abnormal closure): unexpected EOF"),
63+
want: true,
64+
},
65+
{
66+
name: "reader: policy violation 1008",
67+
err: errors.New("websocket: close 1008 (policy violation)"),
68+
want: true,
69+
},
70+
{"reader: GoAway", errors.New("GoAway received"), true},
71+
72+
// The sender path. A write straight to the socket surfaces the raw
73+
// platform error, and its wording is not the reader's. All four of
74+
// these mean the same thing: the connection this flow was using is
75+
// gone, so the flow must reconnect rather than give up.
76+
{
77+
name: "sender: linux EPIPE",
78+
err: socketErr("write", "write", syscall.EPIPE),
79+
want: true,
80+
},
81+
{
82+
name: "sender: linux ECONNRESET",
83+
err: socketErr("write", "write", syscall.ECONNRESET),
84+
want: true,
85+
},
86+
{
87+
name: "sender: windows WSAECONNABORTED",
88+
err: socketErr("write", "wsasend", wsaeConnAborted),
89+
want: true,
90+
},
91+
{
92+
name: "sender: windows WSAECONNRESET",
93+
err: socketErr("write", "wsasend", wsaeConnReset),
94+
want: true,
95+
},
96+
{
97+
name: "reader: windows WSAECONNRESET on recv",
98+
err: socketErr("read", "wsarecv", wsaeConnReset),
99+
want: true,
100+
},
101+
{
102+
name: "sender: wrapped socket error still resumable",
103+
err: fmt.Errorf("sending realtime input: %w", socketErr("write", "wsasend", wsaeConnAborted)),
104+
want: true,
105+
},
106+
107+
// Errors that are about the exchange rather than the transport must
108+
// still terminate the flow.
109+
{
110+
name: "protocol error 1002 is fatal",
111+
err: errors.New("websocket: close 1002 (protocol error): fatal"),
112+
want: false,
113+
},
114+
{
115+
name: "application error is fatal",
116+
err: errors.New("model refused the request"),
117+
want: false,
118+
},
119+
}
120+
121+
for _, tc := range tests {
122+
t.Run(tc.name, func(t *testing.T) {
123+
if got := isResumable(tc.err); got != tc.want {
124+
t.Errorf("isResumable(%v) = %v, want %v", tc.err, got, tc.want)
125+
}
126+
})
127+
}
128+
}
129+
130+
// The reader and the sender both report into the same errChan and RunLive acts
131+
// on whichever lands first, so a single dropped connection must get the same
132+
// verdict whichever goroutine saw it. Before the *net.OpError check, the two
133+
// disagreed on Windows: the reader's close text matched "EOF" and resumed,
134+
// while the sender's wsasend error matched nothing and terminated the session.
135+
func TestIsResumableAgreesAcrossReaderAndSender(t *testing.T) {
136+
platforms := []struct {
137+
name string
138+
readerErr error
139+
senderErr error
140+
}{
141+
{
142+
name: "linux",
143+
readerErr: errors.New("failed to receive message: websocket: close 1006 (abnormal closure): unexpected EOF"),
144+
senderErr: socketErr("write", "write", syscall.EPIPE),
145+
},
146+
{
147+
name: "windows",
148+
readerErr: errors.New("failed to receive message: websocket: close 1006 (abnormal closure): unexpected EOF"),
149+
senderErr: socketErr("write", "wsasend", wsaeConnAborted),
150+
},
151+
}
152+
153+
for _, p := range platforms {
154+
t.Run(p.name, func(t *testing.T) {
155+
reader, sender := isResumable(p.readerErr), isResumable(p.senderErr)
156+
if reader != sender {
157+
t.Errorf("same connection loss classified differently: reader=%v sender=%v; "+
158+
"whether the session resumes would depend on which goroutine reported first",
159+
reader, sender)
160+
}
161+
if !reader {
162+
t.Errorf("a dropped connection should be resumable, got reader=%v sender=%v", reader, sender)
163+
}
164+
})
165+
}
166+
}

0 commit comments

Comments
 (0)