yandex: close the document connection and the writer on Stop

YandexDocsTransport had no Stop of its own. After BaseTransport.Stop the
reader stayed in conn.ReadMessage until the server's next message and
then returned without closing the socket, and the writer stayed blocked
on the write queue forever. Each stop therefore left a WebSocket
attached to the document (a ghost participant) and two goroutines
behind; on a mobile client that is every VPN toggle.

Stop now closes the session's connection, and the writer also waits on
Done() so it exits with the transport.
This commit is contained in:
damnurmum
2026-09-25 18:39:20 +03:00
parent ea52949d45
commit c837df704a
2 changed files with 93 additions and 3 deletions
+71
View File
@@ -0,0 +1,71 @@
package yandex
import (
"net/http"
"net/http/httptest"
"runtime"
"strings"
"testing"
"time"
"github.com/gorilla/websocket"
"openflux/transport"
)
func writerLoops() int {
buf := make([]byte, 1<<20)
return strings.Count(string(buf[:runtime.Stack(buf, true)]), "(*YandexDocsTransport).writerLoop")
}
// Stop must close the document connection and end the writer, not leave
// the socket (a participant in the document) and goroutines behind.
func TestStopClosesDocumentConnection(t *testing.T) {
closed := make(chan struct{})
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
c, err := (&websocket.Upgrader{}).Upgrade(w, r, nil)
if err != nil {
return
}
defer c.Close()
for {
if _, _, err := c.ReadMessage(); err != nil {
close(closed)
return
}
}
}))
defer srv.Close()
conn, _, err := websocket.DefaultDialer.Dial("ws"+strings.TrimPrefix(srv.URL, "http"), nil)
if err != nil {
t.Fatal(err)
}
tr := NewYandexDocsTransport("https://docs.example/d", transport.DefaultConfig())
if err := tr.BaseTransport.Start(); err != nil {
t.Fatal(err)
}
tr.session = &DocSession{Conn: conn, WriteQueue: make(chan []byte, 1)}
before := writerLoops()
go tr.writerLoop()
deadline := time.Now().Add(2 * time.Second)
for writerLoops() == before && time.Now().Before(deadline) {
time.Sleep(5 * time.Millisecond)
}
if err := tr.Stop(); err != nil {
t.Fatal(err)
}
select {
case <-closed:
case <-time.After(2 * time.Second):
t.Fatal("document connection still open after Stop")
}
deadline = time.Now().Add(2 * time.Second)
for writerLoops() > before && time.Now().Before(deadline) {
time.Sleep(5 * time.Millisecond)
}
if writerLoops() > before {
t.Fatal("writer goroutine still running after Stop")
}
}
+22 -3
View File
@@ -114,6 +114,21 @@ func (t *YandexDocsTransport) Start() error {
return nil
}
// Stop also closes the document connection. Otherwise the reader sits in
// ReadMessage until the server's next message and then leaves the socket
// open, keeping a participant attached to the document after the transport
// is gone.
func (t *YandexDocsTransport) Stop() error {
err := t.BaseTransport.Stop()
t.Mu.RLock()
session := t.session
t.Mu.RUnlock()
if session != nil && session.Conn != nil {
_ = session.Conn.Close()
}
return err
}
func (t *YandexDocsTransport) Send(data []byte) error {
if !t.IsConnected() {
return fmt.Errorf("transport not connected")
@@ -288,11 +303,15 @@ func (t *YandexDocsTransport) writerLoop() {
var pending []byte
for t.IsRunning() {
if pending == nil {
packet, ok := <-queue
if !ok {
select {
case packet, ok := <-queue:
if !ok {
return
}
pending = packet
case <-t.Done():
return
}
pending = packet
}
t.Mu.RLock()