mirror of
https://github.com/p1neappleXpress/OpenFlux.git
synced 2026-10-02 05:04:39 +08:00
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:
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user