Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 18 additions & 0 deletions integration_test/pytransformer_contract/bc_helpers_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -248,6 +248,24 @@ func makeEventWithCredentials(messageID, versionID string, credentials []types.C
return ev
}

// makeEvents creates n TransformerEvents for versionID, optionally carrying library version ids.
//
// Message ids are prefixed with versionID so a test sharing one mock config backend across
// subtests can scope its assertions to its own requests.
func makeEvents(versionID string, n int, libraryVersionIDs ...string) []types.TransformerEvent {
libraries := make([]backendconfig.LibraryT, len(libraryVersionIDs))
for i, id := range libraryVersionIDs {
libraries[i] = backendconfig.LibraryT{VersionID: id}
}

events := make([]types.TransformerEvent, n)
for i := range events {
events[i] = makeEvent(fmt.Sprintf("%s-msg-%d", versionID, i+1), versionID)
events[i].Libraries = libraries
}
return events
}

// configBackendEntry controls what the mock config backend returns for a given versionId.
//
// When statusCode is 0 (default), the entry is treated as a normal transformation:
Expand Down
715 changes: 715 additions & 0 deletions integration_test/pytransformer_contract/config_backend_auth_test.go

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -111,7 +111,7 @@ def transformEvent(event, metadata):
wg.Go(func() {
msgID := fmt.Sprintf("msg-cookie-iso-%d", i)
events := []types.TransformerEvent{makeEvent(msgID, versionID)}
status, items := sendRawTransform(t, pyURL, events)
status, _, items := sendRawTransform(t, pyURL, events)

res := result{idx: i, status: status, items: items}
if len(items) == 1 && items[0].StatusCode == http.StatusOK {
Expand Down
39 changes: 24 additions & 15 deletions integration_test/pytransformer_contract/dns_cache_contract_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -103,7 +103,7 @@ def transformEvent(event, metadata):
// Send one request per mock server, sequentially
for _, name := range names {
events := []types.TransformerEvent{makeEvent("msg-"+name, versionIDs[name])}
status, items := sendRawTransform(t, pyURL, events)
status, _, items := sendRawTransform(t, pyURL, events)
requireCorrectServer(t, name, status, items)
}

Expand All @@ -116,7 +116,7 @@ def transformEvent(event, metadata):
// Now repeat the first two — DNS cache should still resolve correctly
for _, name := range []string{"alpha", "bravo"} {
events := []types.TransformerEvent{makeEvent("msg-"+name+"-repeat", versionIDs[name])}
status, items := sendRawTransform(t, pyURL, events)
status, _, items := sendRawTransform(t, pyURL, events)
requireCorrectServer(t, name, status, items)
}

Expand Down Expand Up @@ -144,7 +144,7 @@ def transformEvent(event, metadata):
idx, n := i, name
wg.Go(func() {
events := []types.TransformerEvent{makeEvent("msg-parallel-"+n, versionIDs[n])}
status, items := sendRawTransform(t, pyURL, events)
status, _, items := sendRawTransform(t, pyURL, events)
results[idx] = result{name: n, status: status, items: items}
})
}
Expand All @@ -168,7 +168,7 @@ def transformEvent(event, metadata):
// Same transformation 3 times — DNS cache must remain correct
for i := range 3 {
events := []types.TransformerEvent{makeEvent(fmt.Sprintf("msg-repeat-%d", i), versionIDs["alpha"])}
status, items := sendRawTransform(t, pyURL, events)
status, _, items := sendRawTransform(t, pyURL, events)
requireCorrectServer(t, "alpha", status, items)
}

Expand Down Expand Up @@ -224,7 +224,7 @@ def transformEvent(event, metadata):

t.Run("OverrideResolvesToCorrectServer", func(t *testing.T) {
events := []types.TransformerEvent{makeEvent("msg-override-1", versionID)}
status, items := sendRawTransform(t, pyURL, events)
status, _, items := sendRawTransform(t, pyURL, events)

require.Equal(t, http.StatusOK, status)
require.Len(t, items, 1)
Expand All @@ -242,7 +242,7 @@ def transformEvent(event, metadata):
callsBefore := mockCalls.Load()
for i := range 3 {
events := []types.TransformerEvent{makeEvent(fmt.Sprintf("msg-override-repeat-%d", i), versionID)}
status, items := sendRawTransform(t, pyURL, events)
status, _, items := sendRawTransform(t, pyURL, events)

require.Equal(t, http.StatusOK, status)
require.Len(t, items, 1)
Expand Down Expand Up @@ -325,22 +325,22 @@ def transformEvent(event, metadata):

wg.Go(func() {
events := []types.TransformerEvent{makeEvent("msg-combined-override-1", versionOverride)}
s, items := sendRawTransform(t, pyURL, events)
s, _, items := sendRawTransform(t, pyURL, events)
results[0] = result{"override-1", "override-server", s, items}
})
wg.Go(func() {
events := []types.TransformerEvent{makeEvent("msg-combined-cached-1", versionCached)}
s, items := sendRawTransform(t, pyURL, events)
s, _, items := sendRawTransform(t, pyURL, events)
results[1] = result{"cached-1", "cached-server", s, items}
})
wg.Go(func() {
events := []types.TransformerEvent{makeEvent("msg-combined-override-2", versionOverride)}
s, items := sendRawTransform(t, pyURL, events)
s, _, items := sendRawTransform(t, pyURL, events)
results[2] = result{"override-2", "override-server", s, items}
})
wg.Go(func() {
events := []types.TransformerEvent{makeEvent("msg-combined-cached-2", versionCached)}
s, items := sendRawTransform(t, pyURL, events)
s, _, items := sendRawTransform(t, pyURL, events)
results[3] = result{"cached-2", "cached-server", s, items}
})
wg.Wait()
Expand Down Expand Up @@ -435,27 +435,36 @@ func newMockAPIServer(t *testing.T, name string) (*httptest.Server, *atomic.Int6
return srv, calls
}

// sendRawTransform sends events directly to pytransformer's /customTransform
// endpoint and returns the HTTP status code and parsed response items.
// Unlike the usertransformer.Client, this allows inspecting raw HTTP status.
// sendRawTransform sends events directly to pytransformer's /customTransform endpoint and
// returns the HTTP status code, the response headers and the parsed response items. Unlike the
// usertransformer.Client, this allows inspecting the raw HTTP status and the retry-contract
// headers (X-Rudder-Should-Retry / X-Rudder-Error-Reason); discard the headers with _ when a
// test does not assert on them.
func sendRawTransform(
t *testing.T,
baseURL string,
events []types.TransformerEvent,
) (
int,
http.Header,
[]types.TransformerResponse,
) {
t.Helper()
payload := make([]any, len(events))
for i, ev := range events {
payload[i] = map[string]any{
item := map[string]any{
"message": ev.Message,
"metadata": ev.Metadata,
"destination": map[string]any{
"Transformations": ev.Destination.Transformations,
},
}
// Only when set: pytransformer reads libraries per event, and an empty array would
// change the cache key for every test that does not use libraries.
if len(ev.Libraries) > 0 {
item["libraries"] = ev.Libraries
}
payload[i] = item
}
body, err := jsonrs.Marshal(payload)
require.NoError(t, err)
Expand All @@ -471,5 +480,5 @@ func sendRawTransform(
var items []types.TransformerResponse
require.NoError(t, jsonrs.NewDecoder(resp.Body).Decode(&items))

return resp.StatusCode, items
return resp.StatusCode, resp.Header.Clone(), items
}
Original file line number Diff line number Diff line change
Expand Up @@ -86,7 +86,7 @@ def transformEvent(event, metadata):
{
ev := makeEvent("single-evt", versionID)
ev.Message["n"] = 1
status, items := sendRawTransform(t, pyURL, []types.TransformerEvent{ev})
status, _, items := sendRawTransform(t, pyURL, []types.TransformerEvent{ev})
require.Equal(t, http.StatusOK, status, "single-event /customTransform failed")
require.Len(t, items, 1)
require.Equalf(t, http.StatusOK, items[0].StatusCode,
Expand Down Expand Up @@ -116,7 +116,7 @@ def transformEvent(event, metadata):
batch[j] = ev
inputs = append(inputs, inputEvent{messageID: messageID, n: n})
}
status, items := sendRawTransform(t, pyURL, batch)
status, _, items := sendRawTransform(t, pyURL, batch)
require.Equalf(t, http.StatusOK, status, "request %d /customTransform failed", r)
require.Lenf(t, items, eventsPerRequest, "request %d expected %d items", r, eventsPerRequest)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -191,7 +191,7 @@ def transformEvent(event, metadata):
// Our cap (1s) fires first. The transformation fails with a 400
// per-event error (non-retryable user code HTTP timeout).
events := []types.TransformerEvent{makeEvent("msg-bigger-1", versionBiggerTimeout)}
status, items := sendRawTransform(t, pyURL, events)
status, _, items := sendRawTransform(t, pyURL, events)
require.Equal(t, http.StatusOK, status,
"/customTransform HTTP response must be 200 (per-event errors are in the payload)")
require.Len(t, items, 1)
Expand All @@ -215,7 +215,7 @@ def transformEvent(event, metadata):
// The user passes timeout=0.1 s; the server replies after 2s.
// The user's cap fires first (0.5s < our 1s cap < server's 2s delay).
events := []types.TransformerEvent{makeEvent("msg-smaller-1", versionSmallerTimeout)}
status, items := sendRawTransform(t, pyURL, events)
status, _, items := sendRawTransform(t, pyURL, events)
require.Equal(t, http.StatusOK, status)
require.Len(t, items, 1)
require.Equal(t, http.StatusBadRequest, items[0].StatusCode,
Expand Down Expand Up @@ -272,13 +272,13 @@ def transformEvent(event, metadata):
newConns.Store(0)

ev1 := makeEvent("msg-pool-1", versionID)
status1, items1 := sendRawTransform(t, poolURL, []types.TransformerEvent{ev1})
status1, _, items1 := sendRawTransform(t, poolURL, []types.TransformerEvent{ev1})
require.Equal(t, http.StatusOK, status1)
require.Len(t, items1, 1)
require.Equal(t, http.StatusOK, items1[0].StatusCode, "first request must succeed")

ev2 := makeEvent("msg-pool-2", versionID)
status2, items2 := sendRawTransform(t, poolURL, []types.TransformerEvent{ev2})
status2, _, items2 := sendRawTransform(t, poolURL, []types.TransformerEvent{ev2})
require.Equal(t, http.StatusOK, status2)
require.Len(t, items2, 1)
require.Equal(t, http.StatusOK, items2[0].StatusCode, "second request must succeed")
Expand Down Expand Up @@ -352,13 +352,13 @@ def transformEvent(event, metadata):
newConns.Store(0)

ev1 := makeEvent("msg-shared-alpha", versionIDAlpha)
status1, items1 := sendRawTransform(t, sharedURL, []types.TransformerEvent{ev1})
status1, _, items1 := sendRawTransform(t, sharedURL, []types.TransformerEvent{ev1})
require.Equal(t, http.StatusOK, status1)
require.Len(t, items1, 1)
require.Equal(t, http.StatusOK, items1[0].StatusCode, "alpha request must succeed")

ev2 := makeEvent("msg-shared-beta", versionIDBeta)
status2, items2 := sendRawTransform(t, sharedURL, []types.TransformerEvent{ev2})
status2, _, items2 := sendRawTransform(t, sharedURL, []types.TransformerEvent{ev2})
require.Equal(t, http.StatusOK, status2)
require.Len(t, items2, 1)
require.Equal(t, http.StatusOK, items2[0].StatusCode, "beta request must succeed")
Expand Down Expand Up @@ -424,7 +424,7 @@ def transformEvent(event, metadata):
events := []types.TransformerEvent{makeEvent("msg-slow-drip-1", versionID)}

start := time.Now()
status, items := sendRawTransform(t, pyURL, events)
status, _, items := sendRawTransform(t, pyURL, events)
elapsed := time.Since(start)
t.Logf("slow-drip request elapsed: %s", elapsed)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -113,7 +113,7 @@ def transformEvent(event, metadata):

// 1) Prime the only L1 slot with a CPU-bound transformation.
fillerEvent := makeEvent("filler-1", fillerVersionID)
status, items := sendRawTransform(t, pyURL, []types.TransformerEvent{fillerEvent})
status, _, items := sendRawTransform(t, pyURL, []types.TransformerEvent{fillerEvent})
require.Equal(t, http.StatusOK, status, "filler request must return 200")
require.Len(t, items, 1, "filler request must produce one response item")
require.Equal(t, http.StatusOK, items[0].StatusCode,
Expand All @@ -125,7 +125,7 @@ def transformEvent(event, metadata):
for i := range eventsPerBatch {
ioBatch1[i] = makeEvent(fmt.Sprintf("io1-%d", i), ioVersionID)
}
status, items = sendRawTransform(t, pyURL, ioBatch1)
status, _, items = sendRawTransform(t, pyURL, ioBatch1)
require.Equal(t, http.StatusOK, status, "first I/O batch must return 200")
require.Len(t, items, eventsPerBatch, "first I/O batch must produce one item per event")
for i, item := range items {
Expand All @@ -148,7 +148,7 @@ def transformEvent(event, metadata):
for i := range eventsPerBatch {
ioBatch2[i] = makeEvent(fmt.Sprintf("io2-%d", i), ioVersionID)
}
status, items = sendRawTransform(t, pyURL, ioBatch2)
status, _, items = sendRawTransform(t, pyURL, ioBatch2)
require.Equal(t, http.StatusOK, status, "second I/O batch must return 200")
require.Len(t, items, eventsPerBatch, "second I/O batch must produce one item per event")
for i, item := range items {
Expand Down Expand Up @@ -255,7 +255,7 @@ def transformEvent(event, metadata):
ev.Message["do_http"] = true
ioBatch[i] = ev
}
status, items := sendRawTransform(t, pyURL, ioBatch)
status, _, items := sendRawTransform(t, pyURL, ioBatch)
require.Equal(t, http.StatusOK, status, "I/O batch must return 200")
require.Len(t, items, eventsPerBatch)
for i, item := range items {
Expand All @@ -274,7 +274,7 @@ def transformEvent(event, metadata):
ev.Message["do_http"] = false
cpuBatch[i] = ev
}
status, items = sendRawTransform(t, pyURL, cpuBatch)
status, _, items = sendRawTransform(t, pyURL, cpuBatch)
require.Equal(t, http.StatusOK, status, "CPU batch must return 200")
require.Len(t, items, eventsPerBatch)
for i, item := range items {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -98,7 +98,7 @@ def transformEvent(event, metadata):
)

events := []types.TransformerEvent{makeEvent("msg-module-retry-1", versionID)}
status, items := sendRawTransform(t, pyURL, events)
status, _, items := sendRawTransform(t, pyURL, events)
require.Equal(t, http.StatusOK, status)
require.Len(t, items, 1)
require.Equal(t, http.StatusOK, items[0].StatusCode,
Expand Down Expand Up @@ -190,7 +190,7 @@ def transformEvent(event, metadata):
)

events := []types.TransformerEvent{makeEvent("msg-call-budget-1", versionID)}
status, items := sendRawTransform(t, pyURL, events)
status, _, items := sendRawTransform(t, pyURL, events)
require.Equal(t, http.StatusOK, status)
require.Len(t, items, 1)
require.Equal(t, http.StatusOK, items[0].StatusCode,
Expand Down Expand Up @@ -278,7 +278,7 @@ def transformEvent(event, metadata):
)

events := []types.TransformerEvent{makeEvent("msg-verify-false-1", versionID)}
status, items := sendRawTransform(t, pyURL, events)
status, _, items := sendRawTransform(t, pyURL, events)
require.Equal(t, http.StatusOK, status)
require.Len(t, items, 1)
require.Equal(t, http.StatusOK, items[0].StatusCode,
Expand Down Expand Up @@ -331,7 +331,7 @@ def transformEvent(event, metadata):
)

events := []types.TransformerEvent{makeEvent("msg-req-verify-false-1", versionID)}
status, items := sendRawTransform(t, pyURL, events)
status, _, items := sendRawTransform(t, pyURL, events)
require.Equal(t, http.StatusOK, status)
require.Len(t, items, 1)
require.Equal(t, http.StatusOK, items[0].StatusCode,
Expand Down Expand Up @@ -391,7 +391,7 @@ def transformEvent(event, metadata):

hits.Store(0)
events := []types.TransformerEvent{makeEvent("msg-subclass-1", versionID)}
status, items := sendRawTransform(t, pyURL, events)
status, _, items := sendRawTransform(t, pyURL, events)
require.Equal(t, http.StatusOK, status)
require.Len(t, items, 1)
require.Equal(t, http.StatusOK, items[0].StatusCode,
Expand Down Expand Up @@ -455,7 +455,7 @@ def transformEvent(event, metadata):
hitsA.Store(0)
hitsB.Store(0)
events := []types.TransformerEvent{makeEvent("msg-mount-prefix-1", versionID)}
status, items := sendRawTransform(t, pyURL, events)
status, _, items := sendRawTransform(t, pyURL, events)
require.Equal(t, http.StatusOK, status)
require.Len(t, items, 1)
require.Equal(t, http.StatusOK, items[0].StatusCode,
Expand Down Expand Up @@ -514,7 +514,7 @@ def transformEvent(event, metadata):

observed.reset()
events := []types.TransformerEvent{makeEvent("msg-headers-auth-1", versionID)}
status, items := sendRawTransform(t, pyURL, events)
status, _, items := sendRawTransform(t, pyURL, events)
require.Equal(t, http.StatusOK, status)
require.Len(t, items, 1)
require.Equal(t, http.StatusOK, items[0].StatusCode,
Expand Down Expand Up @@ -576,7 +576,7 @@ def transformEvent(event, metadata):

newConns.Store(0)
events := []types.TransformerEvent{makeEvent("msg-close-noop-1", versionID)}
status, items := sendRawTransform(t, pyURL, events)
status, _, items := sendRawTransform(t, pyURL, events)
require.Equal(t, http.StatusOK, status)
require.Len(t, items, 1)
require.Equal(t, http.StatusOK, items[0].StatusCode,
Expand Down Expand Up @@ -639,13 +639,13 @@ def transformEvent(event, metadata):
newConns.Store(0)

evA := makeEvent("msg-mgd-alpha", versionIDAlpha)
statusA, itemsA := sendRawTransform(t, pyURL, []types.TransformerEvent{evA})
statusA, _, itemsA := sendRawTransform(t, pyURL, []types.TransformerEvent{evA})
require.Equal(t, http.StatusOK, statusA)
require.Len(t, itemsA, 1)
require.Equal(t, http.StatusOK, itemsA[0].StatusCode)

evB := makeEvent("msg-mgd-beta", versionIDBeta)
statusB, itemsB := sendRawTransform(t, pyURL, []types.TransformerEvent{evB})
statusB, _, itemsB := sendRawTransform(t, pyURL, []types.TransformerEvent{evB})
require.Equal(t, http.StatusOK, statusB)
require.Len(t, itemsB, 1)
require.Equal(t, http.StatusOK, itemsB[0].StatusCode)
Expand Down Expand Up @@ -766,7 +766,7 @@ def transformEvent(event, metadata):
for i := range evs {
evs[i] = makeEvent(fmt.Sprintf("msg-%s-%d", variant.name, i), versionID)
}
status, items := sendRawTransform(t, pyURL, evs)
status, _, items := sendRawTransform(t, pyURL, evs)
require.Equal(t, http.StatusOK, status)
require.Len(t, items, eventsPerRun)
for i, it := range items {
Expand Down Expand Up @@ -888,7 +888,7 @@ def transformEvent(event, metadata):
for i := range evs {
evs[i] = makeEvent(fmt.Sprintf("msg-pool-maxsize-%d", i), versionID)
}
status, items := sendRawTransform(t, pyURL, evs)
status, _, items := sendRawTransform(t, pyURL, evs)
require.Equal(t, http.StatusOK, status)
require.Len(t, items, eventsPerRun)
for i, it := range items {
Expand Down
Loading
Loading