summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorKyren223 <Kyren223@proton.me>2025-07-09 13:47:35 +0300
committerKyren223 <Kyren223@proton.me>2025-07-19 18:32:12 +0300
commit402e2b41084fd6558a73121734a1bae66e4ae5ad (patch)
tree2cc699d7985e6f9c16ec29e090e523590b35a50c
parenta81cd075d7807574746ba08b9d2b919565fa6aaf (diff)
Fixed bug where server couldn't shutdown unless all clients closed their
connections willingly
-rw-r--r--internal/server/server.go25
1 files changed, 15 insertions, 10 deletions
diff --git a/internal/server/server.go b/internal/server/server.go
index f8ce1a3..a9f1852 100644
--- a/internal/server/server.go
+++ b/internal/server/server.go
@@ -32,6 +32,8 @@ var nodeId int64 = 0
const CertFile = "EKO_SERVER_CERT_FILE"
+const ReadCheckCancelledInterval = 1 * time.Second
+
func getTlsConfig() *tls.Config {
path, ok := os.LookupEnv(CertFile)
if !ok {
@@ -181,9 +183,8 @@ func (s *server) Run() {
listener, err := tls.Listen("tcp4", ":"+strconv.Itoa(int(s.Port)), getTlsConfig())
if err != nil {
- // TODO: we need certs even for dev
slog.Error("error starting server", "error", err)
- os.Exit(1)
+ assert.Abort("see logs")
}
assert.AddFlush(listener)
@@ -260,23 +261,19 @@ func (server *server) handleConnection(conn net.Conn) {
for payload := range writeQueue {
packet := packet.NewPacket(packet.NewMsgPackEncoder(payload))
if _, err := packet.Into(conn); err != nil {
- // TODO: probably should add this to prevent the
- // "use of closed connection" error, as it's intended to happen
- // and once it happens we can just return
- // if !errors.Is(err, net.ErrClosed) {
- // log.Println(addr, err)
- // }
slog.ErrorContext(ctx, "error sending packet", "error", err, "packet", packet.LogValue(), "payload", payload)
return
}
slog.InfoContext(ctx, "packet sent", "packet", packet.LogValue(), "payload", payload)
}
+ slog.InfoContext(ctx, "writer done")
}()
// Writer closer
go func() {
writerWg.Wait()
sess.CloseWriteQueue() // causes writer to return
+ slog.InfoContext(ctx, "closing writer...")
}()
// Processor
@@ -290,6 +287,7 @@ func (server *server) handleConnection(conn net.Conn) {
for request := range framer.Out {
processPacket(localCtx, sess, request)
}
+ slog.InfoContext(ctx, "processor done")
}()
// NOTE: IMPROTANT LEGAL STUFF
@@ -299,14 +297,20 @@ func (server *server) handleConnection(conn net.Conn) {
// Reader
buffer := make([]byte, 512)
for {
+ conn.SetReadDeadline(time.Now().Add(ReadCheckCancelledInterval))
n, err := conn.Read(buffer)
if err != nil {
if errors.Is(err, io.EOF) {
slog.InfoContext(ctx, "closing gracefully")
+ break
+ } else if ne, ok := err.(net.Error); ok && ne.Timeout() {
+ if ctx.Err() != nil {
+ slog.InfoContext(ctx, "reader context done", "error", ctx.Err())
+ } // else continue normally
} else {
- slog.ErrorContext(ctx, "failed reading from buffer", "error", err)
+ slog.ErrorContext(ctx, "failed reading from conn", "error", err)
+ break
}
- break
}
err = framer.Push(ctx, buffer[:n])
@@ -323,6 +327,7 @@ func (server *server) handleConnection(conn net.Conn) {
}
}
close(framer.Out) // stop processing
+ slog.InfoContext(ctx, "reader done, closed framer")
<-done
}