From f16f29cb1ec6952942b0033327690ef711147cb8 Mon Sep 17 00:00:00 2001 From: pirosiki197 Date: Wed, 24 Jun 2026 01:32:11 +0900 Subject: [PATCH 1/2] fix(loki-stream): reconnect when websocket is unexpectedly closed --- go.mod | 2 +- go.sum | 4 +- pkg/infrastructure/log/loki/streamer.go | 93 +++++++++++++------------ 3 files changed, 50 insertions(+), 49 deletions(-) diff --git a/go.mod b/go.mod index 74524da84..eae55ee7b 100644 --- a/go.mod +++ b/go.mod @@ -16,6 +16,7 @@ require ( github.com/aws/aws-sdk-go-v2/service/s3 v1.102.2 github.com/bep/debounce v1.2.1 github.com/cert-manager/cert-manager v1.20.2 + github.com/coder/websocket v1.8.14 github.com/docker/cli v29.5.2+incompatible github.com/friendsofgo/errors v0.9.2 github.com/gliderlabs/ssh v0.3.8 @@ -37,7 +38,6 @@ require ( github.com/prometheus/common v0.68.0 github.com/regclient/regclient v0.11.5 github.com/samber/lo v1.53.0 - github.com/shiguredo/websocket v1.6.1 github.com/sourcegraph/conc v0.3.1-0.20240121214520-5f936abd7ae8 github.com/spf13/cobra v1.10.2 github.com/spf13/pflag v1.0.10 diff --git a/go.sum b/go.sum index e97bb98ce..284858860 100644 --- a/go.sum +++ b/go.sum @@ -97,6 +97,8 @@ github.com/cloudflare/circl v1.6.3 h1:9GPOhQGF9MCYUeXyMYlqTR6a5gTrgR/fBLXvUgtVcg github.com/cloudflare/circl v1.6.3/go.mod h1:2eXP6Qfat4O/Yhh8BznvKnJ+uzEoTQ6jVKJRn81BiS4= github.com/codahale/rfc6979 v0.0.0-20141003034818-6a90f24967eb h1:EDmT6Q9Zs+SbUoc7Ik9EfrFqcylYqgPZ9ANSbTAntnE= github.com/codahale/rfc6979 v0.0.0-20141003034818-6a90f24967eb/go.mod h1:ZjrT6AXHbDs86ZSdt/osfBi5qfexBrKUdONk989Wnk4= +github.com/coder/websocket v1.8.14 h1:9L0p0iKiNOibykf283eHkKUHHrpG7f65OE3BhhO7v9g= +github.com/coder/websocket v1.8.14/go.mod h1:NX3SzP+inril6yawo5CQXx8+fk145lPDC6pumgx0mVg= github.com/containerd/cgroups/v3 v3.1.3 h1:eUNflyMddm18+yrDmZPn3jI7C5hJ9ahABE5q6dyLYXQ= github.com/containerd/cgroups/v3 v3.1.3/go.mod h1:PKZ2AcWmSBsY/tJUVhtS/rluX0b1uq1GmPO1ElCmbOw= github.com/containerd/console v1.0.5 h1:R0ymNeydRqH2DmakFNdmjR2k0t7UPuiOV/N/27/qqsc= @@ -443,8 +445,6 @@ github.com/sergi/go-diff v1.3.2-0.20230802210424-5b0b94c5c0d3 h1:n661drycOFuPLCN github.com/sergi/go-diff v1.3.2-0.20230802210424-5b0b94c5c0d3/go.mod h1:A0bzQcvG0E7Rwjx0REVgAGH58e96+X0MeOfepqsbeW4= github.com/shibumi/go-pathspec v1.3.0 h1:QUyMZhFo0Md5B8zV8x2tesohbb5kfbpTi9rBnKh5dkI= github.com/shibumi/go-pathspec v1.3.0/go.mod h1:Xutfslp817l2I1cZvgcfeMQJG5QnU2lh5tVaaMCl3jE= -github.com/shiguredo/websocket v1.6.1 h1:xGZ5LmGjQLfGaCcxZrI2/z0en24eJ3VunnNK1RmJdzg= -github.com/shiguredo/websocket v1.6.1/go.mod h1:qUnxxJOWcK8Q7Q+o31UucO9HJQuhcygUYMuLRa88XTw= github.com/shirou/gopsutil/v4 v4.26.3 h1:2ESdQt90yU3oXF/CdOlRCJxrP+Am1aBYubTMTfxJ1qc= github.com/shirou/gopsutil/v4 v4.26.3/go.mod h1:LZ6ewCSkBqUpvSOf+LsTGnRinC6iaNUNMGBtDkJBaLQ= github.com/sigstore/sigstore v1.10.5 h1:KqrOjDhNOVY+uOzQFat2FrGLClPPCb3uz8pK3wuI+ow= diff --git a/pkg/infrastructure/log/loki/streamer.go b/pkg/infrastructure/log/loki/streamer.go index 3252868a1..666f88c1e 100644 --- a/pkg/infrastructure/log/loki/streamer.go +++ b/pkg/infrastructure/log/loki/streamer.go @@ -11,8 +11,8 @@ import ( "text/template" "time" + "github.com/coder/websocket" "github.com/friendsofgo/errors" - "github.com/shiguredo/websocket" "github.com/traPtitech/neoshowcase/pkg/domain" ) @@ -122,65 +122,66 @@ func (l *lokiStreamer) Stream(ctx context.Context, app *domain.Application, begi if err != nil { return nil, errors.Wrap(err, "templating logQL") } - q := make(url.Values) - q.Set("query", logQL) - q.Set("limit", "100") - // ensure start time is not in the future to prevent Loki API errors - start := min(time.Now().UnixNano(), begin.UnixNano()) - q.Set("start", fmt.Sprintf("%d", start)) - - conn, _, err := websocket.DefaultDialer.DialContext(ctx, l.streamEndpoint()+"?"+q.Encode(), nil) - if err != nil { - return nil, errors.Wrap(err, "failed to dial to stream ws endpoint") - } ch := make(chan *domain.ContainerLog, 100) ctx, cancel := context.WithCancel(ctx) go func() { - <-ctx.Done() - _ = conn.Close() defer close(ch) - }() - go func() { defer cancel() - defer slog.InfoContext(ctx, "closing loki websocket stream") - slog.InfoContext(ctx, "new loki websocket stream") - for { - typ, b, err := conn.ReadMessage() - select { // check if context was cancelled - case <-ctx.Done(): - return - default: - } + lastSeenTime := begin + backoffCount := 0 + maxBackoff := 5 + + for backoffCount < maxBackoff { + start := lastSeenTime.Add(time.Nanosecond) + q := make(url.Values) + q.Set("query", logQL) + q.Set("limit", "100") + q.Set("start", fmt.Sprintf("%d", start.UnixNano())) + + conn, _, err := websocket.Dial(ctx, l.streamEndpoint()+"?"+q.Encode(), nil) if err != nil { - slog.ErrorContext(ctx, "failed to read ws message", "error", err) + slog.ErrorContext(ctx, "failed to dial to stream ws endpoint", "error", err) return } - switch typ { - case websocket.TextMessage: - var res streamResponse - err = json.NewDecoder(bytes.NewReader(b)).Decode(&res) - if err != nil { - slog.ErrorContext(ctx, "failed to decode ws message", "error", err) - continue // fail-safe - } - logs, err := res.Streams.toSortedResponse(true) - if err != nil { - slog.ErrorContext(ctx, "failed to decode ws message", "error", err) - continue // fail-safe + + for { + typ, b, err := conn.Read(ctx) + if errors.Is(err, context.Canceled) { + return + } else if err != nil { + slog.WarnContext(ctx, "failed to read ws message", "error", err) + backoffCount += 1 + break // retry } - for _, l := range logs { - select { - case ch <- l: - default: + + switch typ { + case websocket.MessageText: + var res streamResponse + err := json.Unmarshal(b, &res) + if err != nil { + slog.ErrorContext(ctx, "failed to decode ws message", "error", err) + continue // fail-safe + } + logs, err := res.Streams.toSortedResponse(true) + if err != nil { + slog.ErrorContext(ctx, "failed to decode ws message", "error", err) + continue // fail-safe } + for _, l := range logs { + switch { + case l.Time.Before(lastSeenTime): + continue + case l.Time.After(lastSeenTime): + lastSeenTime = l.Time + } + ch <- l + } + case websocket.MessageBinary: + // ignore } - case websocket.BinaryMessage: - // ignore - case websocket.CloseMessage: - return } } }() From ed89b389b97b3cb9955a8af1f6ffe0efc818349d Mon Sep 17 00:00:00 2001 From: pirosiki197 Date: Wed, 24 Jun 2026 10:40:26 +0900 Subject: [PATCH 2/2] fix --- pkg/infrastructure/log/loki/streamer.go | 102 ++++++++++++++---------- 1 file changed, 62 insertions(+), 40 deletions(-) diff --git a/pkg/infrastructure/log/loki/streamer.go b/pkg/infrastructure/log/loki/streamer.go index 666f88c1e..0cf8039b1 100644 --- a/pkg/infrastructure/log/loki/streamer.go +++ b/pkg/infrastructure/log/loki/streamer.go @@ -5,6 +5,7 @@ import ( "context" "encoding/json" "fmt" + "iter" "log/slog" "net/http" "net/url" @@ -131,56 +132,31 @@ func (l *lokiStreamer) Stream(ctx context.Context, app *domain.Application, begi defer cancel() lastSeenTime := begin - backoffCount := 0 - maxBackoff := 5 - for backoffCount < maxBackoff { - start := lastSeenTime.Add(time.Nanosecond) - q := make(url.Values) - q.Set("query", logQL) - q.Set("limit", "100") - q.Set("start", fmt.Sprintf("%d", start.UnixNano())) - - conn, _, err := websocket.Dial(ctx, l.streamEndpoint()+"?"+q.Encode(), nil) + for { + logSeq, err := l.readStream(ctx, logQL, lastSeenTime) if err != nil { - slog.ErrorContext(ctx, "failed to dial to stream ws endpoint", "error", err) + slog.ErrorContext(ctx, "failed to read log stream", "error", err) return } - - for { - typ, b, err := conn.Read(ctx) + for l, err := range logSeq { if errors.Is(err, context.Canceled) { return } else if err != nil { - slog.WarnContext(ctx, "failed to read ws message", "error", err) - backoffCount += 1 + slog.WarnContext(ctx, "unexpected error occurred while reading log stream", "error", err) break // retry } - switch typ { - case websocket.MessageText: - var res streamResponse - err := json.Unmarshal(b, &res) - if err != nil { - slog.ErrorContext(ctx, "failed to decode ws message", "error", err) - continue // fail-safe - } - logs, err := res.Streams.toSortedResponse(true) - if err != nil { - slog.ErrorContext(ctx, "failed to decode ws message", "error", err) - continue // fail-safe - } - for _, l := range logs { - switch { - case l.Time.Before(lastSeenTime): - continue - case l.Time.After(lastSeenTime): - lastSeenTime = l.Time - } - ch <- l - } - case websocket.MessageBinary: - // ignore + switch { + case l.Time.After(lastSeenTime): + lastSeenTime = l.Time + case l.Time.Before(lastSeenTime): + continue + } + select { + case ch <- l: + case <-ctx.Done(): + return } } } @@ -188,3 +164,49 @@ func (l *lokiStreamer) Stream(ctx context.Context, app *domain.Application, begi return ch, nil } + +func (l *lokiStreamer) readStream(ctx context.Context, query string, start time.Time) (iter.Seq2[*domain.ContainerLog, error], error) { + q := make(url.Values) + q.Set("query", query) + q.Set("limit", "100") + q.Set("start", fmt.Sprintf("%d", start.UnixNano())) + + conn, _, err := websocket.Dial(ctx, l.streamEndpoint()+"?"+q.Encode(), nil) + if err != nil { + return nil, err + } + + return func(yield func(*domain.ContainerLog, error) bool) { + defer conn.CloseNow() + + for { + typ, b, err := conn.Read(ctx) + if err != nil { + yield(nil, err) + return + } + + switch typ { + case websocket.MessageText: + var res streamResponse + err := json.Unmarshal(b, &res) + if err != nil { + slog.ErrorContext(ctx, "failed to decode ws message", "error", err) + continue // fail-safe + } + logs, err := res.Streams.toSortedResponse(true) + if err != nil { + slog.ErrorContext(ctx, "failed to decode ws message", "error", err) + continue // fail-safe + } + for _, l := range logs { + if !yield(l, nil) { + return + } + } + case websocket.MessageBinary: + // ignore + } + } + }, nil +}