Store and show an event's query string, and pass it on to an HTTP target when set (closes #312)
check / check (push) Successful in 6m32s
check / check (push) Successful in 6m32s
The receiver keeps the query string of the request it received on the event, in a new `raw_query` column of the per-webhook `events` table; a resubmitted copy carries its original's. The event log and the event's page show it with the request, the event log leaving out one over 32 KiB as it does request headers. The archive and log targets carry it. An HTTP target gets a `forwardQuery` setting, off by default, on both target forms and in the target list: on, each delivery, replays and resubmits included, appends the query string to the target URL, joined with `&` to one it already has. Model: opus-5-5
This commit is contained in:
@@ -110,6 +110,7 @@ type Task struct {
|
||||
MaxRetries int
|
||||
|
||||
Method string
|
||||
RawQuery string
|
||||
Headers string
|
||||
ContentType string
|
||||
Body *string
|
||||
@@ -1752,6 +1753,7 @@ func buildEventFromTask(task *Task) database.Event {
|
||||
event := database.Event{
|
||||
EntrypointID: task.EntrypointID,
|
||||
Method: task.Method,
|
||||
RawQuery: task.RawQuery,
|
||||
Headers: task.Headers,
|
||||
ContentType: task.ContentType,
|
||||
}
|
||||
@@ -2102,6 +2104,7 @@ func buildRecoveryTask(
|
||||
TargetConfig: target.Config,
|
||||
MaxRetries: target.MaxRetries,
|
||||
Method: event.Method,
|
||||
RawQuery: event.RawQuery,
|
||||
Headers: event.Headers,
|
||||
ContentType: event.ContentType,
|
||||
Body: bodyPtr,
|
||||
|
||||
@@ -673,6 +673,11 @@ func TestRecoverPendingDeliveries(t *testing.T) {
|
||||
t, s.WebhookDB, s.WebhookID, targetID, 3,
|
||||
)
|
||||
|
||||
// A recovered delivery still carries its event's query string.
|
||||
require.NoError(t, s.WebhookDB.Model(&database.Event{}).
|
||||
Where("webhook_id = ?", s.WebhookID).
|
||||
Update("raw_query", eventQuery).Error)
|
||||
|
||||
s.Engine.ExportRecoverPendingDeliveries(
|
||||
context.Background(), s.WebhookDB,
|
||||
s.WebhookID,
|
||||
@@ -687,6 +692,8 @@ func TestRecoverPendingDeliveries(t *testing.T) {
|
||||
database.TargetTypeLog,
|
||||
task.TargetType,
|
||||
)
|
||||
|
||||
assert.Equal(t, eventQuery, task.RawQuery)
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatalf("expected task %d", i)
|
||||
}
|
||||
|
||||
@@ -11,6 +11,7 @@ import (
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
@@ -1990,6 +1991,10 @@ func assertLogLineComplete(
|
||||
"log line must contain the full request headers",
|
||||
)
|
||||
|
||||
assert.Contains(t, out, "raw_query="+strconv.Quote(event.RawQuery),
|
||||
"log line must contain the query string",
|
||||
)
|
||||
|
||||
assert.Contains(t, out, event.EntrypointID,
|
||||
"log line must contain the entrypoint id",
|
||||
)
|
||||
@@ -2012,6 +2017,7 @@ func TestDeliverLog_LogsFullContent(t *testing.T) {
|
||||
event := seedEvent(
|
||||
t, db, `{"log-body-marker":"abc123"}`,
|
||||
)
|
||||
event.RawQuery = eventQuery
|
||||
|
||||
dlv := seedDelivery(
|
||||
t, db, event.ID, uuid.New().String(),
|
||||
|
||||
@@ -39,6 +39,9 @@ type TargetConfigForm struct {
|
||||
// Timeout is the HTTP target's per-request timeout in seconds,
|
||||
// empty when unset.
|
||||
Timeout string
|
||||
// ForwardQuery is the HTTP target's setting that passes each
|
||||
// event's query string on to it.
|
||||
ForwardQuery bool
|
||||
// Expiry is the database (archive) target's row expiry.
|
||||
Expiry string
|
||||
// Rotation is the database (archive) target's rotation.
|
||||
@@ -64,9 +67,10 @@ func NewTargetConfigForm(
|
||||
}
|
||||
|
||||
return TargetConfigForm{
|
||||
URL: cfg.URL,
|
||||
Headers: FormatTargetHeaders(cfg.Headers),
|
||||
Timeout: FormatTargetTimeout(cfg.Timeout),
|
||||
URL: cfg.URL,
|
||||
Headers: FormatTargetHeaders(cfg.Headers),
|
||||
Timeout: FormatTargetTimeout(cfg.Timeout),
|
||||
ForwardQuery: cfg.ForwardQuery,
|
||||
}, nil
|
||||
case database.TargetTypeSlack:
|
||||
cfg, err := parseSlackConfig(t.Config)
|
||||
|
||||
@@ -171,6 +171,13 @@ func httpConfigFields(t *database.Target) []ConfigField {
|
||||
})
|
||||
}
|
||||
|
||||
if cfg.ForwardQuery {
|
||||
fields = append(fields, ConfigField{
|
||||
Label: "Query string",
|
||||
Value: "passed on to this target",
|
||||
})
|
||||
}
|
||||
|
||||
fields = append(fields, maxRetriesField(t))
|
||||
|
||||
return fields
|
||||
|
||||
@@ -223,7 +223,8 @@ func TestNewTargetViews_HTTP(t *testing.T) {
|
||||
Type: database.TargetTypeHTTP,
|
||||
Config: `{"url":"` + viewExampleHook + `",` +
|
||||
`"timeout":30,` +
|
||||
`"headers":{"Authorization":"Bearer sekrit"}}`,
|
||||
`"headers":{"Authorization":"Bearer sekrit"},` +
|
||||
`"forwardQuery":true}`,
|
||||
MaxRetries: 5,
|
||||
})
|
||||
|
||||
@@ -235,6 +236,7 @@ func TestNewTargetViews_HTTP(t *testing.T) {
|
||||
"Destination URL": viewMaskedOrigin,
|
||||
"Timeout": "30s",
|
||||
"Headers": "1 configured",
|
||||
"Query string": "passed on to this target",
|
||||
viewMaxRetries: "5",
|
||||
},
|
||||
fields,
|
||||
|
||||
@@ -184,6 +184,7 @@ func (t *databaseTarget) archive(d *database.Delivery) error {
|
||||
WebhookID: webhookID,
|
||||
EntrypointID: d.Event.EntrypointID,
|
||||
Method: d.Event.Method,
|
||||
RawQuery: d.Event.RawQuery,
|
||||
Headers: d.Event.Headers,
|
||||
Body: d.Event.Body,
|
||||
ContentType: d.Event.ContentType,
|
||||
|
||||
@@ -101,6 +101,7 @@ type archivedEvent struct {
|
||||
WebhookID string
|
||||
EntrypointID string
|
||||
Method string
|
||||
RawQuery string
|
||||
Headers string
|
||||
Body string
|
||||
ContentType string
|
||||
|
||||
@@ -360,6 +360,7 @@ func writeRow(w io.Writer, ev *archivedEvent, period string) error {
|
||||
"webhook_id": ev.WebhookID,
|
||||
"entrypoint_id": ev.EntrypointID,
|
||||
"method": ev.Method,
|
||||
"raw_query": ev.RawQuery,
|
||||
"headers": ev.Headers,
|
||||
"body": ev.Body,
|
||||
"content_type": ev.ContentType,
|
||||
|
||||
@@ -166,6 +166,7 @@ func TestArchiveExport_MatchesStoredRows(t *testing.T) {
|
||||
WebhookID: exportWebhookID,
|
||||
EntrypointID: "ep-1",
|
||||
Method: "POST",
|
||||
RawQuery: eventQuery,
|
||||
Headers: `{"X-Test":["yes"]}`,
|
||||
Body: body,
|
||||
ContentType: testContentType,
|
||||
@@ -215,12 +216,13 @@ func assertExportedRow(
|
||||
assert.Equal(t, row.WebhookID, ev["webhook_id"])
|
||||
assert.Equal(t, row.EntrypointID, ev["entrypoint_id"])
|
||||
assert.Equal(t, row.Method, ev["method"])
|
||||
assert.Equal(t, row.RawQuery, ev["raw_query"])
|
||||
assert.Equal(t, row.Headers, ev["headers"])
|
||||
assert.Equal(t, row.ContentType, ev["content_type"])
|
||||
|
||||
if row.Body != binaryBody {
|
||||
assert.Equal(t, row.Body, ev["body"])
|
||||
assert.Len(t, ev, 9, "the nine columns and nothing else: %v", ev)
|
||||
assert.Len(t, ev, 10, "the ten columns and nothing else: %v", ev)
|
||||
|
||||
return
|
||||
}
|
||||
@@ -229,7 +231,7 @@ func assertExportedRow(
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, binaryBody, string(body))
|
||||
assert.Equal(t, "base64", ev["body_encoding"])
|
||||
assert.Len(t, ev, 10, "the nine columns and body_encoding: %v", ev)
|
||||
assert.Len(t, ev, 11, "the ten columns and body_encoding: %v", ev)
|
||||
}
|
||||
|
||||
// TestArchiveExport_Empty proves an archive with nothing in it exports
|
||||
|
||||
@@ -85,6 +85,7 @@ func TestDeliverDatabase_ArchivesEvent(t *testing.T) {
|
||||
|
||||
webhookDB := testWebhookDB(t)
|
||||
event := seedEvent(t, webhookDB, `{"archived":true}`)
|
||||
event.RawQuery = eventQuery
|
||||
d := seedDatabaseTargetDelivery(t, webhookDB, event, tgt)
|
||||
|
||||
env.eng.ExportDeliverDatabase(webhookDB, d)
|
||||
@@ -113,6 +114,7 @@ func TestDeliverDatabase_ArchivesEvent(t *testing.T) {
|
||||
assert.Equal(t, event.ID, rows[0].EventID)
|
||||
assert.Equal(t, event.WebhookID, rows[0].WebhookID)
|
||||
assert.Equal(t, event.Method, rows[0].Method)
|
||||
assert.Equal(t, eventQuery, rows[0].RawQuery)
|
||||
assert.JSONEq(t, `{"archived":true}`, rows[0].Body)
|
||||
}
|
||||
|
||||
|
||||
@@ -8,6 +8,7 @@ import (
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/url"
|
||||
"sort"
|
||||
"sync"
|
||||
"time"
|
||||
@@ -32,6 +33,11 @@ type HTTPTargetConfig struct {
|
||||
URL string `json:"url"`
|
||||
Headers map[string]string `json:"headers,omitempty"`
|
||||
Timeout int `json:"timeout,omitempty"`
|
||||
|
||||
// ForwardQuery passes each event's query string on to the target,
|
||||
// appended to URL. Off, the target URL is sent exactly as
|
||||
// configured.
|
||||
ForwardQuery bool `json:"forwardQuery,omitempty"`
|
||||
}
|
||||
|
||||
// httpCore holds the retry, backoff, and circuit-breaker
|
||||
@@ -444,6 +450,10 @@ func (t *httpTarget) doHTTPRequest(
|
||||
)
|
||||
}
|
||||
|
||||
if cfg.ForwardQuery {
|
||||
appendQuery(req.URL, event.RawQuery)
|
||||
}
|
||||
|
||||
originScoped := applyRequestHeaders(
|
||||
req, event, cfg, t.eng.userAgent(),
|
||||
)
|
||||
@@ -474,6 +484,19 @@ func (t *httpTarget) doHTTPRequest(
|
||||
return resp.StatusCode, string(body), dur, nil
|
||||
}
|
||||
|
||||
// appendQuery adds an event's query string to a delivery's URL, joined
|
||||
// with "&" to any query string the target URL already has.
|
||||
func appendQuery(u *url.URL, rawQuery string) {
|
||||
switch {
|
||||
case rawQuery == "":
|
||||
return
|
||||
case u.RawQuery == "":
|
||||
u.RawQuery = rawQuery
|
||||
default:
|
||||
u.RawQuery += "&" + rawQuery
|
||||
}
|
||||
}
|
||||
|
||||
// clientForRequest returns the client for one delivery attempt.
|
||||
// originScoped is the header set applyRequestHeaders built for that
|
||||
// attempt; a request with neither a per-target timeout nor an
|
||||
|
||||
@@ -0,0 +1,172 @@
|
||||
package delivery_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"testing"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
"sneak.berlin/go/webhooker/internal/database"
|
||||
"sneak.berlin/go/webhooker/internal/delivery"
|
||||
)
|
||||
|
||||
// eventQuery is the query string the events in these tests arrived
|
||||
// with.
|
||||
const eventQuery = "a=1&b=2"
|
||||
|
||||
// httpTargetConfig is the stored configuration of an HTTP target at
|
||||
// targetURL.
|
||||
func httpTargetConfig(
|
||||
t *testing.T, targetURL string, forwardQuery bool,
|
||||
) string {
|
||||
t.Helper()
|
||||
|
||||
cfg, err := json.Marshal(delivery.HTTPTargetConfig{
|
||||
URL: targetURL, ForwardQuery: forwardQuery,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
|
||||
return string(cfg)
|
||||
}
|
||||
|
||||
// deliverWithQuery sends one event that arrived with eventQuery to an
|
||||
// HTTP target configured with cfg, through the path a received event's
|
||||
// delivery takes, and returns the attempt it recorded.
|
||||
func deliverWithQuery(t *testing.T, cfg string) database.DeliveryResult {
|
||||
t.Helper()
|
||||
|
||||
s := newISetup(t)
|
||||
|
||||
event := iSeedEvent(t, s.WebhookDB, s.WebhookID, "{}")
|
||||
d := iSeedDelivery(
|
||||
t, s.WebhookDB, event.ID, uuid.NewString(),
|
||||
database.DeliveryStatusPending,
|
||||
)
|
||||
task := iTask(
|
||||
d, event, s.WebhookID, d.TargetID, "query", cfg, 0, 1, &event.Body,
|
||||
)
|
||||
task.RawQuery = eventQuery
|
||||
|
||||
s.Engine.ExportProcessNewTask(context.TODO(), &task)
|
||||
|
||||
var result database.DeliveryResult
|
||||
|
||||
require.NoError(t, s.WebhookDB.Where(
|
||||
"delivery_id = ?", d.ID,
|
||||
).First(&result).Error)
|
||||
|
||||
return result
|
||||
}
|
||||
|
||||
// TestDeliverHTTP_ForwardQuery proves the URL a delivery is sent to:
|
||||
// with the target's setting off, the target URL exactly as configured;
|
||||
// with it on, the event's query string appended, joined with "&" to a
|
||||
// query string the target URL already has.
|
||||
func TestDeliverHTTP_ForwardQuery(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
// The target URL's path, without and with a query string of its
|
||||
// own.
|
||||
const (
|
||||
plain = "/in"
|
||||
withQuery = "/in?key=k"
|
||||
)
|
||||
|
||||
tests := map[string]struct {
|
||||
path string
|
||||
forwardQuery bool
|
||||
want string
|
||||
}{
|
||||
"off": {
|
||||
path: plain, want: plain,
|
||||
},
|
||||
"off, the target URL has a query string": {
|
||||
path: withQuery, want: withQuery,
|
||||
},
|
||||
"on": {
|
||||
path: plain, forwardQuery: true, want: plain + "?" + eventQuery,
|
||||
},
|
||||
"on, the target URL has a query string": {
|
||||
path: withQuery, forwardQuery: true,
|
||||
want: withQuery + "&" + eventQuery,
|
||||
},
|
||||
}
|
||||
|
||||
for name, tc := range tests {
|
||||
t.Run(name, func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
received := make(chan string, 1)
|
||||
|
||||
ts := httptest.NewServer(http.HandlerFunc(
|
||||
func(w http.ResponseWriter, r *http.Request) {
|
||||
received <- r.RequestURI
|
||||
|
||||
w.WriteHeader(http.StatusOK)
|
||||
},
|
||||
))
|
||||
t.Cleanup(ts.Close)
|
||||
|
||||
result := deliverWithQuery(t, httpTargetConfig(
|
||||
t, ts.URL+tc.path, tc.forwardQuery,
|
||||
))
|
||||
|
||||
assert.True(t, result.Success)
|
||||
require.Len(t, received, 1)
|
||||
assert.Equal(t, tc.want, <-received)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// TestDeliverHTTP_ForwardedQueryKeepsTheTargetURLMasked proves the
|
||||
// credential in a target URL's own query string stays masked once the
|
||||
// event's query string is appended to it: in a response or error that
|
||||
// echoes the URL the target was sent, as the event log's Redactor shows
|
||||
// it, and in the error a failed connection stores.
|
||||
func TestDeliverHTTP_ForwardedQueryKeepsTheTargetURLMasked(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
const secret = "s3cr3t"
|
||||
|
||||
received := make(chan string, 1)
|
||||
|
||||
ts := httptest.NewServer(http.HandlerFunc(
|
||||
func(w http.ResponseWriter, r *http.Request) {
|
||||
received <- r.RequestURI
|
||||
|
||||
w.WriteHeader(http.StatusBadRequest)
|
||||
},
|
||||
))
|
||||
t.Cleanup(ts.Close)
|
||||
|
||||
target := &database.Target{
|
||||
Type: database.TargetTypeHTTP,
|
||||
Config: httpTargetConfig(t, ts.URL+"/in?token="+secret, true),
|
||||
}
|
||||
|
||||
deliverWithQuery(t, target.Config)
|
||||
require.Len(t, received, 1)
|
||||
|
||||
sent := <-received
|
||||
require.Equal(t, "/in?token="+secret+"&"+eventQuery, sent)
|
||||
|
||||
redactor := delivery.NewRedactor(target)
|
||||
|
||||
for _, echoed := range []string{sent, ts.URL + sent} {
|
||||
shown := redactor.Redact("rejected " + echoed)
|
||||
assert.NotContains(t, shown, secret, echoed)
|
||||
assert.Contains(t, shown, delivery.RedactionMarker, echoed)
|
||||
}
|
||||
|
||||
// Nothing listens on port 1.
|
||||
failed := deliverWithQuery(t, httpTargetConfig(
|
||||
t, "http://127.0.0.1:1/in?token="+secret, true,
|
||||
))
|
||||
require.NotEmpty(t, failed.Error)
|
||||
assert.NotContains(t, failed.Error, secret)
|
||||
assert.NotContains(t, failed.Error, eventQuery)
|
||||
}
|
||||
@@ -9,9 +9,9 @@ import (
|
||||
)
|
||||
|
||||
// logTarget is a fire-and-forget target that logs the entire
|
||||
// inbound webhook — the full request body and headers, plus
|
||||
// the method, content type, and the webhook and entrypoint
|
||||
// ids — then records a single successful attempt.
|
||||
// inbound webhook — the full request body, query string and
|
||||
// headers, plus the method, content type, and the webhook and
|
||||
// entrypoint ids — then records a single successful attempt.
|
||||
//
|
||||
// This is the one log call in the service that deliberately writes
|
||||
// unbounded client-chosen bytes, so it is the one exception to the
|
||||
@@ -46,6 +46,7 @@ func (t *logTarget) Deliver(
|
||||
"webhook_id", d.Event.WebhookID,
|
||||
"entrypoint_id", d.Event.EntrypointID,
|
||||
"method", d.Event.Method,
|
||||
"raw_query", d.Event.RawQuery,
|
||||
"content_type", d.Event.ContentType,
|
||||
"headers", d.Event.Headers,
|
||||
"body", d.Event.Body,
|
||||
|
||||
Reference in New Issue
Block a user