commit 5acac28
delthas
·
2023-12-11 16:37:55 +0000 UTC
parent 686d365
Send pings, retry connecting after sleep/long pause
2 files changed,
+40,
-10
M
app.go
M
app.go
+0,
-5
1@@ -412,11 +412,6 @@ func (app *App) tryConnect() (conn net.Conn, err error) {
2 return nil, fmt.Errorf("connect: %v", err)
3 }
4
5- if tcpConn, ok := conn.(*net.TCPConn); ok {
6- tcpConn.SetKeepAlive(true)
7- tcpConn.SetKeepAlivePeriod(15 * time.Second)
8- }
9-
10 if app.cfg.TLS {
11 host, _, _ := net.SplitHostPort(addr) // should succeed since net.Dial did.
12 conn = tls.Client(conn, &tls.Config{
+40,
-5
1@@ -5,6 +5,8 @@ import (
2 "fmt"
3 "net"
4 "strings"
5+ "sync/atomic"
6+ "time"
7 "unicode"
8 )
9
10@@ -14,6 +16,11 @@ func ChanInOut(conn net.Conn) (in <-chan Message, out chan<- Message) {
11 in_ := make(chan Message, chanCapacity)
12 out_ := make(chan Message, chanCapacity)
13
14+ const keepAlive = 30 * time.Second
15+ const maxRTT = 10 * time.Second
16+ var last atomic.Value
17+ last.Store(time.Now())
18+
19 go func() {
20 r := bufio.NewScanner(conn)
21 for r.Scan() {
22@@ -23,17 +30,45 @@ func ChanInOut(conn net.Conn) (in <-chan Message, out chan<- Message) {
23 if err != nil {
24 continue
25 }
26+ now := time.Now()
27+ last.Store(now)
28+ conn.SetReadDeadline(now.Add(keepAlive + maxRTT))
29 in_ <- msg
30 }
31 close(in_)
32 }()
33
34 go func() {
35- for msg := range out_ {
36- // TODO send messages by batches
37- _, err := fmt.Fprintf(conn, "%s\r\n", msg.String())
38- if err != nil {
39- break
40+ t := time.NewTicker(time.Second)
41+ defer t.Stop()
42+ outer:
43+ for {
44+ select {
45+ case msg, ok := <-out_:
46+ if !ok {
47+ break outer
48+ }
49+ last.Store(time.Now())
50+ // TODO send messages by batches
51+ _, err := fmt.Fprintf(conn, "%s\r\n", msg.String())
52+ if err != nil {
53+ break outer
54+ }
55+ case <-t.C:
56+ now := time.Now()
57+ if last.Load().(time.Time).Add(keepAlive).After(now) {
58+ continue
59+ }
60+ if last.Load().(time.Time).Add(keepAlive + maxRTT).Before(now) {
61+ // probably out of sleep, reset connection
62+ conn.Close()
63+ continue
64+ }
65+ last.Store(now)
66+ _, err := fmt.Fprint(conn, "PING _\r\n")
67+ if err != nil {
68+ break outer
69+ }
70 }
71 }
72 _ = conn.Close()