diff --git a/pkg/listen/proxy/binary_body_test.go b/pkg/listen/proxy/binary_body_test.go new file mode 100644 index 00000000..ecfe4bb7 --- /dev/null +++ b/pkg/listen/proxy/binary_body_test.go @@ -0,0 +1,133 @@ +package proxy + +import ( + "bytes" + "encoding/base64" + "encoding/json" + "io" + "net/http" + "net/http/httptest" + "net/url" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/hookdeck/hookdeck-cli/pkg/websocket" +) + +type receivedRequest struct { + body []byte + contentType string + contentLength int64 +} + +// forwardAttempt runs one attempt through processAttempt against a local +// server and returns what that server received. +func forwardAttempt(t *testing.T, request websocket.AttemptRequest) receivedRequest { + t.Helper() + + received := make(chan receivedRequest, 1) + local := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + body, _ := io.ReadAll(r.Body) + received <- receivedRequest{body: body, contentType: r.Header.Get("Content-Type"), contentLength: r.ContentLength} + w.WriteHeader(http.StatusOK) + })) + t.Cleanup(local.Close) + + target, err := url.Parse(local.URL) + require.NoError(t, err) + p := New(&Config{URL: target, NoHealthcheck: true}, nil, newRecordingRenderer()) + + p.processAttempt(websocket.IncomingMessage{Attempt: &websocket.Attempt{ + Body: websocket.AttemptBody{Path: "/webhooks", EventID: "evt_test", AttemptId: "evt_test", Request: request}, + }}) + + select { + case r := <-received: + return r + case <-time.After(5 * time.Second): + t.Fatal("local server never received the request") + return receivedRequest{} + } +} + +func headersJSON(t *testing.T, headers map[string]string) json.RawMessage { + t.Helper() + raw, err := json.Marshal(headers) + require.NoError(t, err) + return raw +} + +// Every byte value, so any UTF-8 round trip along the way would show. +func allBytes() []byte { + b := make([]byte, 256) + for i := range b { + b[i] = byte(i) + } + return b +} + +func TestProcessAttemptForwardsBinaryBodyByteExact(t *testing.T) { + body := allBytes() + + got := forwardAttempt(t, websocket.AttemptRequest{ + Method: http.MethodPost, + BodyFormat: websocket.BodyFormatBinary, + DataBase64: base64.StdEncoding.EncodeToString(body), + // A stale Content-Length must not win over the decoded body's length. + Headers: headersJSON(t, map[string]string{"content-type": "application/octet-stream", "content-length": "1"}), + }) + + assert.Equal(t, body, got.body) + assert.Equal(t, int64(len(body)), got.contentLength) + assert.Equal(t, "application/octet-stream", got.contentType) +} + +func TestProcessAttemptForwardsBinaryMultipartUnparsed(t *testing.T) { + contentType := "multipart/form-data; boundary=hookdeck-boundary" + var body bytes.Buffer + body.WriteString("--hookdeck-boundary\r\nContent-Disposition: form-data; name=\"field\"\r\n\r\nvalue\r\n") + body.WriteString("--hookdeck-boundary\r\nContent-Disposition: form-data; name=\"file\"; filename=\"f.bin\"\r\nContent-Type: application/octet-stream\r\n\r\n") + body.Write(allBytes()) + body.WriteString("\r\n--hookdeck-boundary--\r\n") + + got := forwardAttempt(t, websocket.AttemptRequest{ + Method: http.MethodPost, + BodyFormat: websocket.BodyFormatBinary, + DataBase64: base64.StdEncoding.EncodeToString(body.Bytes()), + Headers: headersJSON(t, map[string]string{"content-type": contentType}), + }) + + assert.Equal(t, body.Bytes(), got.body, "boundary and file part bytes must match the original") + assert.Equal(t, contentType, got.contentType, "the boundary travels in the original Content-Type") +} + +func TestProcessAttemptForwardsTextBodyFromDataString(t *testing.T) { + // The shape every server sends today, and the only one older servers send. + got := forwardAttempt(t, websocket.AttemptRequest{ + Method: http.MethodPost, + DataString: `{"hello":"wörld"}`, + Headers: headersJSON(t, map[string]string{"content-type": "application/json"}), + }) + + assert.Equal(t, `{"hello":"wörld"}`, string(got.body)) + assert.Equal(t, int64(len(`{"hello":"wörld"}`)), got.contentLength) +} + +func TestAttemptRequestDecodesFromTheServerWireFormat(t *testing.T) { + var msg websocket.IncomingMessage + require.NoError(t, json.Unmarshal([]byte(`{"event":"attempt","body":{"cli_path":"/","request":{"method":"POST","headers":{},"body_format":"binary","data_base64":"AP+A"}}}`), &msg)) + require.NotNil(t, msg.Attempt) + + body, err := msg.Attempt.Body.Request.Body() + require.NoError(t, err) + assert.Equal(t, []byte{0x00, 0xff, 0x80}, body) + assert.True(t, msg.Attempt.Body.Request.IsBinary()) +} + +func TestAttemptRequestRejectsInvalidBase64(t *testing.T) { + _, err := websocket.AttemptRequest{BodyFormat: websocket.BodyFormatBinary, DataBase64: "not base64!"}.Body() + assert.Error(t, err) +} diff --git a/pkg/listen/proxy/proxy.go b/pkg/listen/proxy/proxy.go index 3edc7e6a..d47b5756 100644 --- a/pkg/listen/proxy/proxy.go +++ b/pkg/listen/proxy/proxy.go @@ -1,6 +1,7 @@ package proxy import ( + "bytes" "context" "crypto/tls" "encoding/json" @@ -13,7 +14,6 @@ import ( "os" "os/signal" "strconv" - "strings" "sync" "sync/atomic" "syscall" @@ -440,8 +440,27 @@ func (p *Proxy) processAttempt(msg websocket.IncomingMessage) { req.Header.Set(key, unquoted_value) } - req.Body = ioutil.NopCloser(strings.NewReader(webhookEvent.Body.Request.DataString)) - req.ContentLength = int64(len(webhookEvent.Body.Request.DataString)) + // Binary bodies (data_base64) are forwarded as the original bytes, with the + // original Content-Type (multipart boundary included) already set above. + // net/http ignores Content-Length in req.Header and uses ContentLength. + body, err := webhookEvent.Body.Request.Body() + if err != nil { + p.renderer.OnEventError(eventID, webhookEvent, fmt.Errorf("decoding binary request body: %w", err), time.Now()) + // Fail the attempt now rather than leave Hookdeck waiting for its timeout. + if wsClient := p.currentWebSocketClient(); wsClient != nil { + wsClient.SendMessage(&websocket.OutgoingMessage{ + ErrorAttemptResponse: &websocket.ErrorAttemptResponse{ + Event: "attempt_response", + Body: websocket.ErrorAttemptBody{ + AttemptId: webhookEvent.Body.AttemptId, + Error: true, + }, + }}) + } + return + } + req.Body = ioutil.NopCloser(bytes.NewReader(body)) + req.ContentLength = int64(len(body)) // For interactive mode: start 100ms timer and HTTP request concurrently requestStartTime := time.Now() diff --git a/pkg/listen/tui/model.go b/pkg/listen/tui/model.go index 8f9ebc52..4334d7c3 100644 --- a/pkg/listen/tui/model.go +++ b/pkg/listen/tui/model.go @@ -366,8 +366,11 @@ func (m *Model) buildEventDetailsContent(event *EventInfo) (string, requestCopyC requestCopy.headers = strings.TrimSuffix(requestCopy.headers, "\n") content.WriteString("\n") - // Request body - if event.Data.Body.Request.DataString != "" { + // Request body. Binary bodies are summarised rather than printed: they + // are not text, and the copied request would not reproduce them. + if event.Data.Body.Request.IsBinary() { + content.WriteString(faintStyle.Render(binaryBodySummary(event.Data.Body.Request)) + "\n") + } else if event.Data.Body.Request.DataString != "" { // Try to pretty print JSON requestCopy.body = m.prettyPrintJSON(event.Data.Body.Request.DataString) content.WriteString(requestCopy.body + "\n") @@ -438,6 +441,16 @@ func (c *requestCopyContent) buildRequest() { c.request = strings.Join(parts, "\n\n") } +// binaryBodySummary describes a binary request body by size, since its bytes +// cannot be shown as text. +func binaryBodySummary(req websocket.AttemptRequest) string { + body, err := req.Body() + if err != nil { + return "(binary body could not be decoded)" + } + return fmt.Sprintf("(binary body, %d bytes)", len(body)) +} + // prettyPrintJSON attempts to pretty print JSON, returns original if not valid JSON. // It uses json.Indent so object key order is preserved exactly as received rather // than being sorted (which json.Marshal would do). diff --git a/pkg/websocket/attempt_messages.go b/pkg/websocket/attempt_messages.go index 1df32f74..c13f50ca 100644 --- a/pkg/websocket/attempt_messages.go +++ b/pkg/websocket/attempt_messages.go @@ -1,16 +1,47 @@ package websocket import ( + "encoding/base64" "encoding/json" ) +// CapabilitiesHeader advertises, on every websocket connect, which attempt +// formats this CLI understands. The server only sends a binary body +// (data_base64) to a session that advertised CapabilityBinaryBody; without it, +// binary events fail on the server instead of reaching an older CLI that would +// forward an empty body. +const CapabilitiesHeader = "X-Hookdeck-CLI-Capabilities" + +// CapabilityBinaryBody means the CLI forwards request.data_base64 as raw bytes. +const CapabilityBinaryBody = "binary" + +// BodyFormatBinary marks an attempt whose body is carried in DataBase64. +const BodyFormatBinary = "binary" + type AttemptRequest struct { Method string `json:"method"` Timeout int64 `json:"timeout"` DataString string `json:"data_string"` + BodyFormat string `json:"body_format,omitempty"` + DataBase64 string `json:"data_base64,omitempty"` Headers json.RawMessage `json:"headers"` } +// IsBinary reports whether the body travels as base64 rather than as DataString. +func (r AttemptRequest) IsBinary() bool { + return r.BodyFormat == BodyFormatBinary || r.DataBase64 != "" +} + +// Body returns the exact bytes to forward to the local server. Binary bodies +// are decoded from DataBase64; text bodies are DataString as sent, which keeps +// working against servers that predate data_base64. +func (r AttemptRequest) Body() ([]byte, error) { + if r.IsBinary() { + return base64.StdEncoding.DecodeString(r.DataBase64) + } + return []byte(r.DataString), nil +} + type AttemptBody struct { Path string `json:"cli_path"` EventID string `json:"event_id"` diff --git a/pkg/websocket/client.go b/pkg/websocket/client.go index 9d84429a..6d7b6157 100644 --- a/pkg/websocket/client.go +++ b/pkg/websocket/client.go @@ -303,6 +303,7 @@ func (c *Client) connect(ctx context.Context) error { header.Set("Accept-Encoding", "identity") header.Set("User-Agent", useragent.GetEncodedUserAgent()) header.Set("X-Hookdeck-Client-User-Agent", useragent.GetEncodedHookdeckUserAgent()) + header.Set(CapabilitiesHeader, CapabilityBinaryBody) header.Set("Websocket-Id", c.WebSocketID) header.Set("X-Team-Id", c.TeamID) header.Set("Authorization", "Basic "+basicAuth(c.CLIKey, "")) diff --git a/pkg/websocket/client_test.go b/pkg/websocket/client_test.go index 6cb87c20..10b6445b 100644 --- a/pkg/websocket/client_test.go +++ b/pkg/websocket/client_test.go @@ -72,6 +72,9 @@ func TestConnectSendsSessionRecreationHeaders(t *testing.T) { if got := captured.Get("Websocket-Id"); got != "cses_test" { t.Errorf("Websocket-Id = %q, want %q", got, "cses_test") } + if got := captured.Get(CapabilitiesHeader); got != CapabilityBinaryBody { + t.Errorf("%s = %q, want %q", CapabilitiesHeader, got, CapabilityBinaryBody) + } if got := captured.Get("X-Webhook-Ids"); got != "web_abc,web_def" { t.Errorf("X-Webhook-Ids = %q, want %q", got, "web_abc,web_def") } diff --git a/test/acceptance/listen_binary_test.go b/test/acceptance/listen_binary_test.go new file mode 100644 index 00000000..07d6f6f9 --- /dev/null +++ b/test/acceptance/listen_binary_test.go @@ -0,0 +1,600 @@ +//go:build listen + +package acceptance + +import ( + "archive/tar" + "bytes" + "compress/gzip" + "crypto/sha256" + "encoding/binary" + "encoding/json" + "fmt" + "image" + "image/color" + "image/jpeg" + "image/png" + "io" + "math" + "mime" + "mime/multipart" + "net/http" + "net/http/httptest" + "net/textproto" + "net/url" + "os" + "path/filepath" + "runtime" + "strconv" + "sync" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// binaryDeliveryEnvVar gates the binary delivery tests. They need server +// support that ships separately from the CLI: Core must send data_base64 to +// sessions advertising the binary capability, and ingestion must store these +// bodies as binary (BINARY_PAYLOADS_ENABLED). Against a server without that, +// binary events fail with CLI_BINARY_UNSUPPORTED and nothing reaches the local +// server. Remove the gate once the server side is deployed. +const binaryDeliveryEnvVar = "HOOKDECK_CLI_TESTING_BINARY_DELIVERY" + +// multipartBinaryDeliveryEnvVar additionally enables the multipart cases, which +// also need ingestion to store multipart/form-data as binary +// (hookdeck/http-ingestion#547, hookdeck/core#5708). Until then multipart is +// stored as UTF-8 text and non-UTF-8 file parts cannot survive. +const multipartBinaryDeliveryEnvVar = "HOOKDECK_CLI_TESTING_MULTIPART_BINARY_DELIVERY" + +// legacyCLIVersion is the last release before binary delivery: it sends no +// X-Hookdeck-CLI-Capabilities header and only reads data_string. +const legacyCLIVersion = "3.0.3" + +// legacyCLIChecksums pins the SHA-256 of each release tarball, from the +// release's published checksum files. +var legacyCLIChecksums = map[string]string{ + "darwin_amd64": "a2ccc50db7211cadb20d23f8025b69b1b8ba180d92cfb4c92d7fd74e907f73e8", + "darwin_arm64": "b26f0e4a077c80cdde5cef8f71f189e27bcdf4cbc5792ce89f765f310f147764", + "linux_amd64": "73ecce58128efb8dce09657085f642bb470b618857ff4d53d400b585cf65d764", + "linux_arm64": "a98f6e46ca3d147d37af12cdf8a08651818097c20f521cf6ebc1131b9ef65bf7", +} + +// downloadLegacyCLI fetches the released v3.0.3 binary for this platform, +// verifies it against the pinned checksum, and returns its path. +func downloadLegacyCLI(t *testing.T) string { + t.Helper() + platform := runtime.GOOS + "_" + runtime.GOARCH + want, ok := legacyCLIChecksums[platform] + require.True(t, ok, "no pinned v%s release for %s", legacyCLIVersion, platform) + + asset := fmt.Sprintf("hookdeck_%s_%s.tar.gz", legacyCLIVersion, platform) + url := fmt.Sprintf("https://github.com/hookdeck/hookdeck-cli/releases/download/v%s/%s", legacyCLIVersion, asset) + client := &http.Client{Timeout: 2 * time.Minute} + resp, err := client.Get(url) + require.NoError(t, err, "download %s", url) + defer resp.Body.Close() + require.Equal(t, http.StatusOK, resp.StatusCode, "download %s", url) + archive, err := io.ReadAll(resp.Body) + require.NoError(t, err) + require.Equal(t, want, fmt.Sprintf("%x", sha256.Sum256(archive)), "checksum mismatch for %s", asset) + + gz, err := gzip.NewReader(bytes.NewReader(archive)) + require.NoError(t, err) + tr := tar.NewReader(gz) + for { + hdr, err := tr.Next() + require.NoError(t, err, "hookdeck binary not found in %s", asset) + if hdr.Name != "hookdeck" { + continue + } + binary := filepath.Join(t.TempDir(), "hookdeck-"+legacyCLIVersion) + f, err := os.OpenFile(binary, os.O_CREATE|os.O_WRONLY, 0o755) + require.NoError(t, err) + _, err = io.Copy(f, tr) + require.NoError(t, err) + require.NoError(t, f.Close()) + return binary + } +} + +func requireBinaryDelivery(t *testing.T) { + t.Helper() + if testing.Short() { + t.Skip("Skipping acceptance test in short mode") + } + if os.Getenv(binaryDeliveryEnvVar) == "" { + t.Skipf("set %s=1 once the server delivers binary bodies to the CLI", binaryDeliveryEnvVar) + } +} + +func requireMultipartBinaryDelivery(t *testing.T) { + t.Helper() + if os.Getenv(multipartBinaryDeliveryEnvVar) == "" { + t.Skipf("set %s=1 once ingestion stores multipart bodies as binary", multipartBinaryDeliveryEnvVar) + } +} + +// --------------------------------------------------------------------------- +// Local app +// --------------------------------------------------------------------------- + +// receivedRequest is one request as the local app saw it. Files holds what the +// app wrote to disk: the raw body for single-file requests, or every file part +// of a multipart form, keyed by form field name. +type receivedRequest struct { + contentType string + eventID string + attemptCount int + body []byte + files map[string]string +} + +// localApp stands in for the user's application: it saves incoming bodies to +// disk the way an upload endpoint would, using only the standard library. +type localApp struct { + server *httptest.Server + port string + received chan receivedRequest + + mu sync.Mutex + // status returns the HTTP status for the nth request (1-based). + status func(n int) int + count int +} + +func startLocalApp(t *testing.T) *localApp { + t.Helper() + app := &localApp{received: make(chan receivedRequest, 16), status: func(int) int { return http.StatusOK }} + dir := t.TempDir() + + app.server = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + app.mu.Lock() + app.count++ + n := app.count + status := app.status(n) + app.mu.Unlock() + + req := receivedRequest{ + contentType: r.Header.Get("Content-Type"), + eventID: r.Header.Get("X-Hookdeck-EventID"), + files: map[string]string{}, + } + req.attemptCount, _ = strconv.Atoi(r.Header.Get("X-Hookdeck-Attempt-Count")) + + mediaType, _, _ := mime.ParseMediaType(req.contentType) + if mediaType == "multipart/form-data" { + if err := r.ParseMultipartForm(32 << 20); err != nil { + t.Errorf("local app: ParseMultipartForm: %v", err) + } else { + for field, headers := range r.MultipartForm.File { + src, err := headers[0].Open() + require.NoError(t, err) + path := filepath.Join(dir, fmt.Sprintf("%d-%s-%s", n, field, filepath.Base(headers[0].Filename))) + dst, err := os.Create(path) + require.NoError(t, err) + _, err = io.Copy(dst, src) + require.NoError(t, err) + require.NoError(t, dst.Close()) + require.NoError(t, src.Close()) + req.files[field] = path + } + } + } else { + body, _ := io.ReadAll(r.Body) + req.body = body + exts, _ := mime.ExtensionsByType(mediaType) + ext := ".bin" + if len(exts) > 0 { + ext = exts[0] + } + path := filepath.Join(dir, fmt.Sprintf("%d-body%s", n, ext)) + require.NoError(t, os.WriteFile(path, body, 0o600)) + req.files["body"] = path + } + + select { + case app.received <- req: + default: + } + w.WriteHeader(status) + })) + t.Cleanup(app.server.Close) + + u, err := url.Parse(app.server.URL) + require.NoError(t, err) + app.port = u.Port() + require.NotEmpty(t, app.port) + return app +} + +func (a *localApp) setStatus(status func(n int) int) { + a.mu.Lock() + defer a.mu.Unlock() + a.status = status +} + +func (a *localApp) next(t *testing.T, within time.Duration, outputs ...*syncBuffer) receivedRequest { + t.Helper() + select { + case r := <-a.received: + return r + case <-time.After(within): + for _, o := range outputs { + t.Logf("listen output: %s", o.String()) + } + t.Fatalf("no request reached the local app within %s; if the event failed with CLI_BINARY_UNSUPPORTED, the server does not deliver binary bodies to this CLI", within) + return receivedRequest{} + } +} + +func (a *localApp) expectNothing(t *testing.T, within time.Duration) { + t.Helper() + select { + case r := <-a.received: + t.Fatalf("expected no delivery, got %d bytes of %q", len(r.body), r.contentType) + case <-time.After(within): + } +} + +// --------------------------------------------------------------------------- +// Fixtures: real media, generated deterministically +// --------------------------------------------------------------------------- + +type fixture struct { + name string + contentType string + data []byte + // validate decodes the file the app saved, proving it is usable, not just + // the same length. + validate func(t *testing.T, data []byte) +} + +func testImage() image.Image { + img := image.NewRGBA(image.Rect(0, 0, 32, 32)) + for y := 0; y < 32; y++ { + for x := 0; x < 32; x++ { + img.Set(x, y, color.RGBA{R: uint8(x * 8), G: uint8(y * 8), B: uint8((x ^ y) * 8), A: 255}) + } + } + return img +} + +func pngFixture(t *testing.T) fixture { + var buf bytes.Buffer + require.NoError(t, png.Encode(&buf, testImage())) + return fixture{name: "picture.png", contentType: "image/png", data: buf.Bytes(), validate: func(t *testing.T, data []byte) { + img, err := png.Decode(bytes.NewReader(data)) + require.NoError(t, err, "saved PNG must decode") + assert.Equal(t, testImage().Bounds(), img.Bounds()) + }} +} + +func jpegFixture(t *testing.T) fixture { + var buf bytes.Buffer + require.NoError(t, jpeg.Encode(&buf, testImage(), &jpeg.Options{Quality: 90})) + return fixture{name: "photo.jpg", contentType: "image/jpeg", data: buf.Bytes(), validate: func(t *testing.T, data []byte) { + img, err := jpeg.Decode(bytes.NewReader(data)) + require.NoError(t, err, "saved JPEG must decode") + assert.Equal(t, testImage().Bounds(), img.Bounds()) + }} +} + +// wavFixture is 0.1s of a 440Hz tone as 16-bit PCM, which covers every byte +// value in its samples. +func wavFixture() fixture { + const sampleRate, samples = 8000, 800 + var pcm bytes.Buffer + for i := 0; i < samples; i++ { + v := int16(math.Sin(2*math.Pi*440*float64(i)/sampleRate) * 32000) + _ = binary.Write(&pcm, binary.LittleEndian, v) + } + var buf bytes.Buffer + buf.WriteString("RIFF") + _ = binary.Write(&buf, binary.LittleEndian, uint32(36+pcm.Len())) + buf.WriteString("WAVEfmt ") + for _, v := range []any{uint32(16), uint16(1), uint16(1), uint32(sampleRate), uint32(sampleRate * 2), uint16(2), uint16(16)} { + _ = binary.Write(&buf, binary.LittleEndian, v) + } + buf.WriteString("data") + _ = binary.Write(&buf, binary.LittleEndian, uint32(pcm.Len())) + buf.Write(pcm.Bytes()) + return fixture{name: "tone.wav", contentType: "audio/wav", data: buf.Bytes(), validate: func(t *testing.T, data []byte) { + require.GreaterOrEqual(t, len(data), 44) + assert.Equal(t, "RIFF", string(data[0:4])) + assert.Equal(t, "WAVE", string(data[8:12])) + assert.Equal(t, uint32(samples*2), binary.LittleEndian.Uint32(data[40:44]), "WAV data chunk length") + }} +} + +// pdfFixture is a minimal one-page PDF with the customary binary comment line. +func pdfFixture() fixture { + objects := []string{ + "<< /Type /Catalog /Pages 2 0 R >>", + "<< /Type /Pages /Kids [3 0 R] /Count 1 >>", + "<< /Type /Page /Parent 2 0 R /MediaBox [0 0 72 72] >>", + } + var buf bytes.Buffer + buf.WriteString("%PDF-1.4\n%\xe2\xe3\xcf\xd3\n") + offsets := make([]int, len(objects)) + for i, obj := range objects { + offsets[i] = buf.Len() + fmt.Fprintf(&buf, "%d 0 obj\n%s\nendobj\n", i+1, obj) + } + xref := buf.Len() + fmt.Fprintf(&buf, "xref\n0 %d\n0000000000 65535 f \n", len(objects)+1) + for _, off := range offsets { + fmt.Fprintf(&buf, "%010d 00000 n \n", off) + } + fmt.Fprintf(&buf, "trailer\n<< /Size %d /Root 1 0 R >>\nstartxref\n%d\n%%%%EOF\n", len(objects)+1, xref) + return fixture{name: "doc.pdf", contentType: "application/pdf", data: buf.Bytes(), validate: func(t *testing.T, data []byte) { + assert.True(t, bytes.HasPrefix(data, []byte("%PDF-1.4\n%\xe2\xe3\xcf\xd3")), "PDF header and binary marker") + assert.True(t, bytes.HasSuffix(data, []byte("%%EOF\n")), "PDF trailer") + }} +} + +func allBytesFixture() fixture { + data := make([]byte, 256) + for i := range data { + data[i] = byte(i) + } + return fixture{name: "all-bytes.bin", contentType: "application/octet-stream", data: data, validate: func(*testing.T, []byte) {}} +} + +// multipartUpload builds a form with a text field and one file part per +// fixture, the way a browser or SDK upload would. +func multipartUpload(t *testing.T, files ...fixture) (contentType string, body []byte) { + var buf bytes.Buffer + w := multipart.NewWriter(&buf) + require.NoError(t, w.WriteField("caption", "binary delivery acceptance")) + for i, f := range files { + h := make(textproto.MIMEHeader) + h.Set("Content-Disposition", fmt.Sprintf(`form-data; name="file%d"; filename=%q`, i, f.name)) + h.Set("Content-Type", f.contentType) + part, err := w.CreatePart(h) + require.NoError(t, err) + _, err = part.Write(f.data) + require.NoError(t, err) + } + require.NoError(t, w.Close()) + return w.FormDataContentType(), buf.Bytes() +} + +// assertSavedFile checks the file the app wrote is the file that was sent, +// byte for byte, and still decodes as its format. +func assertSavedFile(t *testing.T, path string, want fixture) { + t.Helper() + got, err := os.ReadFile(path) + require.NoError(t, err) + assert.Equal(t, sha256.Sum256(want.data), sha256.Sum256(got), + "%s: saved file differs from the original (%d bytes sent, %d saved)", want.name, len(want.data), len(got)) + want.validate(t, got) +} + +// --------------------------------------------------------------------------- +// Harness +// --------------------------------------------------------------------------- + +type binaryListenSetup struct { + cli *CLIRunner + sourceName string + sourceURL string + connID string +} + +func newBinaryListenSetup(t *testing.T, prefix string) binaryListenSetup { + t.Helper() + cli := NewCLIRunner(t) + timestamp := generateTimestamp() + sourceName := prefix + "-" + timestamp + + var conn Connection + require.NoError(t, cli.RunJSON(&conn, + "gateway", "connection", "create", + "--name", prefix+"-conn-"+timestamp, + "--source-name", sourceName, + "--source-type", "WEBHOOK", + "--destination-name", prefix+"-dst-"+timestamp, + "--destination-type", "CLI", + "--destination-cli-path", "/", + )) + require.NotEmpty(t, conn.ID) + t.Cleanup(func() { deleteConnection(t, cli, conn.ID) }) + + var src Source + require.NoError(t, cli.RunJSON(&src, "gateway", "source", "get", conn.Source.ID)) + require.NotEmpty(t, src.URL) + return binaryListenSetup{cli: cli, sourceName: sourceName, sourceURL: src.URL, connID: conn.ID} +} + +// waitForTunnel gives listen time to connect and fails fast if it exited. +func waitForTunnel(t *testing.T, done chan error, outputs ...*syncBuffer) { + t.Helper() + t.Log("Waiting 12 seconds for the tunnel to connect...") + time.Sleep(12 * time.Second) + select { + case err := <-done: + for _, o := range outputs { + t.Logf("listen output: %s", o.String()) + } + t.Fatalf("listen exited before forwarding anything: %v", err) + default: + } +} + +// postRaw sends body to a source URL with the given Content-Type, unmodified. +func postRaw(t *testing.T, sourceURL, contentType string, body []byte) { + t.Helper() + client := &http.Client{Timeout: 10 * time.Second} + resp, err := client.Post(sourceURL, contentType, bytes.NewReader(body)) + require.NoError(t, err, "POST to source URL failed") + defer resp.Body.Close() + require.True(t, resp.StatusCode >= 200 && resp.StatusCode < 300, + "POST to source URL returned %d", resp.StatusCode) +} + +// --------------------------------------------------------------------------- +// Tests +// --------------------------------------------------------------------------- + +// TestListenForwardsBinaryFilesThatTheAppCanSave sends real files (images, a +// PDF, audio inside a multipart upload, and every byte value) through a source +// and `hookdeck listen`, and checks the local app can save each one back to a +// file identical to the original. +func TestListenForwardsBinaryFilesThatTheAppCanSave(t *testing.T) { + requireBinaryDelivery(t) + + s := newBinaryListenSetup(t, "test-bin") + app := startLocalApp(t) + _, stdout, stderr, done := startListenCapturingOutput(t, s.cli, "listen", app.port, s.sourceName, "--output", "compact") + waitForTunnel(t, done, stdout, stderr) + + // Raw bodies: the content types ingestion stores as binary. + for _, f := range []fixture{allBytesFixture(), pngFixture(t), jpegFixture(t), pdfFixture()} { + f := f + t.Run("raw "+f.contentType, func(t *testing.T) { + postRaw(t, s.sourceURL, f.contentType, f.data) + got := app.next(t, 45*time.Second, stdout, stderr) + assert.Equal(t, f.contentType, got.contentType, "original Content-Type must be forwarded") + assertSavedFile(t, got.files["body"], f) + }) + } + + t.Run("multipart upload with picture and audio", func(t *testing.T) { + requireMultipartBinaryDelivery(t) + files := []fixture{pngFixture(t), wavFixture(), allBytesFixture()} + contentType, body := multipartUpload(t, files...) + + postRaw(t, s.sourceURL, contentType, body) + got := app.next(t, 45*time.Second, stdout, stderr) + + assert.Equal(t, contentType, got.contentType, "boundary travels in the original Content-Type") + require.Len(t, got.files, len(files), "the app should save every file part") + for i, f := range files { + assertSavedFile(t, got.files[fmt.Sprintf("file%d", i)], f) + } + }) +} + +// TestListenRetriesBinaryBodyByteExact fails the first delivery, retries the +// event, and checks the retried attempt carries the same bytes. Retries resolve +// the CLI session through a different path from first attempts. +func TestListenRetriesBinaryBodyByteExact(t *testing.T) { + requireBinaryDelivery(t) + + s := newBinaryListenSetup(t, "test-bin-retry") + app := startLocalApp(t) + app.setStatus(func(n int) int { + if n == 1 { + return http.StatusInternalServerError + } + return http.StatusOK + }) + _, stdout, stderr, done := startListenCapturingOutput(t, s.cli, "listen", app.port, s.sourceName, "--output", "compact") + waitForTunnel(t, done, stdout, stderr) + + f := pngFixture(t) + postRaw(t, s.sourceURL, f.contentType, f.data) + + first := app.next(t, 45*time.Second, stdout, stderr) + require.NotEmpty(t, first.eventID, "delivery should carry X-Hookdeck-EventID") + assert.Equal(t, 1, first.attemptCount) + assertSavedFile(t, first.files["body"], f) + + s.cli.RunExpectSuccess("gateway", "event", "retry", first.eventID) + + second := app.next(t, 45*time.Second, stdout, stderr) + assert.Equal(t, first.eventID, second.eventID, "the retry should deliver the same event") + assert.Equal(t, 2, second.attemptCount) + assert.Equal(t, f.contentType, second.contentType) + assertSavedFile(t, second.files["body"], f) +} + +// TestListenBinaryWithOldAndNewCLIListening runs the released v3.0.3 CLI, which +// predates binary delivery, and this CLI on the same source at the same time, +// as separate CLI clients. Each session gets its own event: this CLI receives +// the exact bytes, while the old CLI's event fails closed (or, for multipart, +// is delivered as the lossy text it always got). +func TestListenBinaryWithOldAndNewCLIListening(t *testing.T) { + requireBinaryDelivery(t) + legacyBinary := downloadLegacyCLI(t) + + s := newBinaryListenSetup(t, "test-bin-mixed") + legacyCLI := newSeparateCLIClient(t, s.cli) + newApp := startLocalApp(t) + oldApp := startLocalApp(t) + _, newOut, newErr, newDone := startListenCapturingOutput(t, s.cli, "listen", newApp.port, s.sourceName, "--output", "compact") + _, oldOut, oldErr, oldDone := startListenBinaryCapturingOutput(t, legacyCLI, legacyBinary, "listen", oldApp.port, s.sourceName, "--output", "compact") + waitForTunnel(t, newDone, newOut, newErr) + waitForTunnel(t, oldDone, oldOut, oldErr) + + t.Run("raw image", func(t *testing.T) { + f := pngFixture(t) + postRaw(t, s.sourceURL, f.contentType, f.data) + + got := newApp.next(t, 45*time.Second, newOut, newErr) + assertSavedFile(t, got.files["body"], f) + oldApp.expectNothing(t, 20*time.Second) + + assertEventFailedWithCode(t, s, "CLI_BINARY_UNSUPPORTED") + }) + + t.Run("multipart upload", func(t *testing.T) { + requireMultipartBinaryDelivery(t) + files := []fixture{pngFixture(t), wavFixture()} + contentType, body := multipartUpload(t, files...) + postRaw(t, s.sourceURL, contentType, body) + + got := newApp.next(t, 45*time.Second, newOut, newErr) + for i, f := range files { + assertSavedFile(t, got.files[fmt.Sprintf("file%d", i)], f) + } + + // The old CLI still gets the upload, decoded as UTF-8 text as before + // multipart moved to binary, so its file parts are not byte-exact. + legacy := oldApp.next(t, 45*time.Second, oldOut, oldErr) + assert.Equal(t, contentType, legacy.contentType) + legacyPNG, err := os.ReadFile(legacy.files["file0"]) + require.NoError(t, err) + assert.NotEqual(t, files[0].data, legacyPNG, "the old CLI cannot receive non-UTF-8 bytes intact") + }) +} + +// newSeparateCLIClient returns a runner with its own config file and its own +// `hookdeck ci` login, so it is a separate CLI client. Event ids hash the CLI +// client id, so two listeners sharing one client collapse into a single event +// delivered to only one of them. The new config starts as a copy of base's so +// settings such as api_base carry over. +func newSeparateCLIClient(t *testing.T, base *CLIRunner) *CLIRunner { + t.Helper() + configPath := filepath.Join(t.TempDir(), "legacy-cli-config.toml") + if base.configPath != "" { + if existing, err := os.ReadFile(base.configPath); err == nil { + require.NoError(t, os.WriteFile(configPath, existing, 0o600)) + } + } + return NewCLIRunnerWithConfigPath(t, configPath) +} + +// assertEventFailedWithCode waits for an event on the connection to fail with +// the given error code. +func assertEventFailedWithCode(t *testing.T, s binaryListenSetup, code string) { + t.Helper() + var last []map[string]any + for i := 0; i < propagationAttempts; i++ { + var resp struct { + Models []map[string]any `json:"models"` + } + require.NoError(t, s.cli.RunJSON(&resp, "gateway", "event", "list", "--connection-id", s.connID, "--limit", "20")) + last = resp.Models + for _, e := range resp.Models { + if e["status"] == "FAILED" && e["error_code"] == code { + return + } + } + time.Sleep(propagationInterval) + } + raw, _ := json.Marshal(last) + t.Fatalf("no event on %s failed with %s; events: %s", s.connID, code, raw) +} diff --git a/test/acceptance/listen_test.go b/test/acceptance/listen_test.go index 10a285f7..8206f6c4 100644 --- a/test/acceptance/listen_test.go +++ b/test/acceptance/listen_test.go @@ -63,6 +63,17 @@ func startListenCapturingOutput(t *testing.T, cli *CLIRunner, extraArgs ...strin require.NoError(t, buildCmd.Run(), "failed to build CLI binary") t.Cleanup(func() { _ = os.Remove(binary) }) + return startListenBinaryCapturingOutput(t, cli, binary, extraArgs...) +} + +// startListenBinaryCapturingOutput runs `listen` from an already-built CLI +// binary, for example an older release, with the runner's config file. +func startListenBinaryCapturingOutput(t *testing.T, cli *CLIRunner, binary string, extraArgs ...string) (*exec.Cmd, *syncBuffer, *syncBuffer, chan error) { + t.Helper() + + projectRoot, err := filepath.Abs("../..") + require.NoError(t, err, "Failed to get project root") + cmd := exec.Command(binary, extraArgs...) cmd.Dir = projectRoot