From 30ff5abf75b737694181d505ebde7bba6cc26d34 Mon Sep 17 00:00:00 2001 From: ppzxc Date: Sat, 21 Mar 2026 12:04:10 +0900 Subject: [PATCH 1/5] refactor(webhook): shared HTTP transport for connection reuse MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Sender now holds two *http.Transport fields (secure/insecure) initialised in NewSender(). Send() selects the appropriate transport and wraps it in a per-call http.Client — connection pool is reused across calls. Also drains resp.Body before Close() to guarantee keep-alive reuse. --- internal/adapter/output/webhook/sender.go | 21 +++++++++++++++++---- 1 file changed, 17 insertions(+), 4 deletions(-) diff --git a/internal/adapter/output/webhook/sender.go b/internal/adapter/output/webhook/sender.go index 22331dd..c4ff192 100644 --- a/internal/adapter/output/webhook/sender.go +++ b/internal/adapter/output/webhook/sender.go @@ -5,6 +5,7 @@ import ( "context" "crypto/tls" "fmt" + "io" "net/http" "time" @@ -13,19 +14,30 @@ import ( const defaultTimeoutSec = 10 -type Sender struct{} +type Sender struct { + transport *http.Transport + insecureTransport *http.Transport +} -func NewSender() *Sender { return &Sender{} } +func NewSender() *Sender { + return &Sender{ + transport: &http.Transport{}, + insecureTransport: &http.Transport{ + TLSClientConfig: &tls.Config{InsecureSkipVerify: true}, //nolint:gosec + }, + } +} func (s *Sender) Send(ctx context.Context, out domain.Output, payload []byte) error { timeoutSec := out.TimeoutSec if timeoutSec <= 0 { timeoutSec = defaultTimeoutSec } - client := &http.Client{Timeout: time.Duration(timeoutSec) * time.Second} + t := s.transport if out.SkipTLSVerify { - client.Transport = &http.Transport{TLSClientConfig: &tls.Config{InsecureSkipVerify: true}} //nolint:gosec + t = s.insecureTransport } + client := &http.Client{Transport: t, Timeout: time.Duration(timeoutSec) * time.Second} req, err := http.NewRequestWithContext(ctx, http.MethodPost, out.URL, bytes.NewReader(payload)) if err != nil { return fmt.Errorf("create request: %w", err) @@ -39,6 +51,7 @@ func (s *Sender) Send(ctx context.Context, out domain.Output, payload []byte) er return fmt.Errorf("send: %w", err) } defer resp.Body.Close() + _, _ = io.Copy(io.Discard, resp.Body) if resp.StatusCode >= 400 { return fmt.Errorf("webhook returned %d", resp.StatusCode) } From 42d25087b083432a91ecb9a0918bddc534a43799 Mon Sep 17 00:00:00 2001 From: ppzxc Date: Sat, 21 Mar 2026 12:04:51 +0900 Subject: [PATCH 2/5] refactor(domain): remove dead InputType.IsValid() Zero callers in production or test code. Hardcoding three fixed types (BESZEL/DOZZLE/GENERIC) contradicts the generic-relay design; removing the method avoids confusion and reduces maintenance surface. --- internal/domain/input_type.go | 7 ------- 1 file changed, 7 deletions(-) diff --git a/internal/domain/input_type.go b/internal/domain/input_type.go index f9a26a5..6892cba 100644 --- a/internal/domain/input_type.go +++ b/internal/domain/input_type.go @@ -8,10 +8,3 @@ const ( InputTypeGeneric InputType = "GENERIC" ) -func (s InputType) IsValid() bool { - switch s { - case InputTypeBeszel, InputTypeDozzle, InputTypeGeneric: - return true - } - return false -} From d84cf2f35dac9151d406253e3cb6ddbd7eda3771 Mon Sep 17 00:00:00 2001 From: ppzxc Date: Sat, 21 Mar 2026 12:05:25 +0900 Subject: [PATCH 3/5] feat(http): placeholder GET /messages/{messageId} returns 501 MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Route previously delegated to h.Healthz (200 OK), misleading callers into thinking the endpoint was functional. Now returns 501 Not Implemented with a problem+json body. TDD: RED (200) → GREEN (501). --- internal/adapter/input/http/handler_test.go | 13 +++++++++++++ internal/adapter/input/http/router.go | 5 ++++- 2 files changed, 17 insertions(+), 1 deletion(-) diff --git a/internal/adapter/input/http/handler_test.go b/internal/adapter/input/http/handler_test.go index d7cde43..13253de 100644 --- a/internal/adapter/input/http/handler_test.go +++ b/internal/adapter/input/http/handler_test.go @@ -243,6 +243,19 @@ func TestDocs_AsyncAPI(t *testing.T) { } } +func TestGetMessageByID_Returns501(t *testing.T) { + router := newTestRouter(func(_ context.Context, _ domain.InputType, _ string, _ []byte) (string, error) { + return "", nil + }) + req := httptest.NewRequest(http.MethodGet, "/inputs/beszel/messages/some-id", nil) + req.Header.Set("Authorization", "Bearer test-token") + w := httptest.NewRecorder() + router.ServeHTTP(w, req) + if w.Code != http.StatusNotImplemented { + t.Errorf("status = %d, want 501", w.Code) + } +} + func TestDocs_HTML(t *testing.T) { router := newTestRouter(func(_ context.Context, _ domain.InputType, _ string, _ []byte) (string, error) { return "", nil diff --git a/internal/adapter/input/http/router.go b/internal/adapter/input/http/router.go index cb64f4c..7fbc340 100644 --- a/internal/adapter/input/http/router.go +++ b/internal/adapter/input/http/router.go @@ -44,7 +44,10 @@ func NewRouter(uc input.ReceiveMessageUseCase, resolver input.InputResolver, ws ws.ServeWS(w, req, inputTypeFromContext(req.Context())) }) r.Post("/messages", h.PostMessage) - r.Get("/messages/{messageId}", h.Healthz) // placeholder + r.Get("/messages/{messageId}", func(w http.ResponseWriter, r *http.Request) { + writeError(w, r, http.StatusNotImplemented, "Not Implemented", + "get message by ID is not yet implemented") + }) }) return r From 3d46f7580afd49046b5052d19bdf2f4e223c6801 Mon Sep 17 00:00:00 2001 From: ppzxc Date: Sat, 21 Mar 2026 12:06:28 +0900 Subject: [PATCH 4/5] test(tcp): add 1 MiB boundary-value tests for scanner limit Accepted: (maxMessageBytes-1)-byte payload is delivered normally. Exceeded: (maxMessageBytes+1)-byte payload is silently dropped by bufio.Scanner (ErrTooLong) and Receive is never called. Documents existing behaviour; no production code changed. --- internal/adapter/input/tcp/listener_test.go | 61 +++++++++++++++++++++ 1 file changed, 61 insertions(+) diff --git a/internal/adapter/input/tcp/listener_test.go b/internal/adapter/input/tcp/listener_test.go index e79aa80..77da171 100644 --- a/internal/adapter/input/tcp/listener_test.go +++ b/internal/adapter/input/tcp/listener_test.go @@ -206,6 +206,67 @@ func TestListener_GracefulShutdown(t *testing.T) { } } +func TestListener_MaxMessageSize_Accepted(t *testing.T) { + // A message body of exactly (maxMessageBytes - 1) bytes (+ delimiter) must be delivered. + mock := &mockReceiveUseCase{returnID: "msg-1"} + addr, _ := startTestListener(t, '\n', "application/json", mock) + + conn, err := net.Dial("tcp", addr) + if err != nil { + t.Fatalf("dial: %v", err) + } + defer conn.Close() + + // payload = maxMessageBytes - 1 bytes, then delimiter + payload := make([]byte, maxMessageBytes-1) + for i := range payload { + payload[i] = 'x' + } + payload = append(payload, '\n') + + if _, err := conn.Write(payload); err != nil { + t.Fatalf("write: %v", err) + } + + calls, err := waitForCalls(mock, 1, 3*time.Second) + if err != nil { + t.Fatalf("message not received: %v", err) + } + if len(calls[0].body) != maxMessageBytes-1 { + t.Errorf("body len = %d, want %d", len(calls[0].body), maxMessageBytes-1) + } +} + +func TestListener_MaxMessageSize_Exceeded(t *testing.T) { + // A message body exceeding maxMessageBytes must be silently dropped by the scanner. + mock := &mockReceiveUseCase{returnID: "msg-1"} + addr, _ := startTestListener(t, '\n', "application/json", mock) + + conn, err := net.Dial("tcp", addr) + if err != nil { + t.Fatalf("dial: %v", err) + } + defer conn.Close() + + // payload = maxMessageBytes + 1 bytes, then delimiter — scanner will reject this + payload := make([]byte, maxMessageBytes+1) + for i := range payload { + payload[i] = 'x' + } + payload = append(payload, '\n') + + if _, err := conn.Write(payload); err != nil { + t.Fatalf("write: %v", err) + } + + // Wait and verify no calls were recorded + time.Sleep(500 * time.Millisecond) + calls := mock.getCalls() + if len(calls) != 0 { + t.Errorf("expected 0 calls for oversized message, got %d", len(calls)) + } +} + func TestListener_EmptyMessageSkip(t *testing.T) { mock := &mockReceiveUseCase{returnID: "msg-1"} addr, _ := startTestListener(t, '\n', "application/json", mock) From fbba4058122f22df2845b17fbe40238d507eb1ca Mon Sep 17 00:00:00 2001 From: ppzxc Date: Sat, 21 Mar 2026 12:07:49 +0900 Subject: [PATCH 5/5] fix(service): protect builtin eval keys from ParsedData collision MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit buildEvalData() was unconditionally merging msg.ParsedData into the eval map, allowing a ParsedData field named "input" (or "id", "payload", etc.) to silently overwrite the authoritative builtin value. This could cause filter/routing expressions to evaluate against attacker-controlled data. Guard with a package-level builtinEvalKeys set; ParsedData keys that collide with builtins are silently ignored. TDD: RED (filter failed when builtin was overwritten) → GREEN. --- internal/application/service/relay_worker.go | 10 ++++-- .../application/service/relay_worker_test.go | 35 +++++++++++++++++++ 2 files changed, 43 insertions(+), 2 deletions(-) diff --git a/internal/application/service/relay_worker.go b/internal/application/service/relay_worker.go index c7513ad..6273dbf 100644 --- a/internal/application/service/relay_worker.go +++ b/internal/application/service/relay_worker.go @@ -223,6 +223,10 @@ func (w *RelayWorker) deliver(ctx context.Context, out domain.Output, payload [] return fmt.Errorf("retries exhausted: %w", lastErr) } +var builtinEvalKeys = map[string]struct{}{ + "id": {}, "input": {}, "payload": {}, "createdAt": {}, "status": {}, +} + func buildEvalData(msg domain.Message) map[string]any { data := map[string]any{ "id": msg.ID, @@ -231,9 +235,11 @@ func buildEvalData(msg domain.Message) map[string]any { "createdAt": msg.CreatedAt.Format(time.RFC3339), "status": string(msg.Status), } - // Merge ParsedData fields + // Merge ParsedData fields, skipping any key that would overwrite a builtin. for k, v := range msg.ParsedData { - data[k] = v + if _, reserved := builtinEvalKeys[k]; !reserved { + data[k] = v + } } return data } diff --git a/internal/application/service/relay_worker_test.go b/internal/application/service/relay_worker_test.go index 6dca121..f6f26a8 100644 --- a/internal/application/service/relay_worker_test.go +++ b/internal/application/service/relay_worker_test.go @@ -656,3 +656,38 @@ func TestRelayWorker_InvalidTransition_SkipsUpdate(t *testing.T) { t.Error("expected ack to be called regardless of invalid transition") } } + +func TestRelayWorker_ParsedDataDoesNotOverrideBuiltinKeys(t *testing.T) { + // ParsedData contains "input": "HACKED" — if this overwrites the builtin + // "input" key, the filter `data.input == "BESZEL"` will fail and the sender + // will never be called. The fix must protect builtin keys from ParsedData. + msg := domain.Message{ + ID: "key-collision", + Input: domain.InputTypeBeszel, + Payload: domain.RawPayload(`{}`), + Status: domain.MessageStatusPending, + Version: 1, + ParsedData: map[string]any{ + "input": "HACKED", // must NOT overwrite builtin "input" = "BESZEL" + }, + } + queue := &mockMessageQueue{messages: []domain.Message{msg}} + repo := &mockRepo{saveFn: func(_ context.Context, _ domain.Message) error { return nil }} + sender := &mockSender{} + ruleReader := &mockRuleReader{ + rule: domain.Rule{InputID: "beszel", Filter: `data.input == "BESZEL"`}, + outputs: []domain.Output{{ID: "c1", Type: domain.OutputTypeWebhook}}, + } + registry := &mockRegistry{sender: sender} + + ctx, cancel := context.WithTimeout(context.Background(), 300*time.Millisecond) + defer cancel() + + worker := service.NewRelayWorker(queue, repo, ruleReader, registry, newExprRegistry(), service.DefaultRelayWorkerConfig()) + worker.Start(ctx, 1) + time.Sleep(150 * time.Millisecond) + + if sender.count.Load() == 0 { + t.Error("builtin key 'input' was overwritten by ParsedData — filter failed when it should have passed") + } +}