From 3fd61f00186e238ac9f02880edcbb7ea1dc4be2a Mon Sep 17 00:00:00 2001 From: prabhaks Date: Wed, 26 Aug 2026 07:55:41 -0700 Subject: [PATCH] Update stream deletion assertions for the new async DELETE contract parseablehq/parseable#1770 makes DELETE /logstream/{stream} return 202 Accepted instead of 200 OK, since deletion now runs in the background rather than blocking the response. It also makes recreating a stream immediately after deleting it return 409 while the old stream's deletion is still in flight, instead of succeeding right away. Updates DeleteStream to expect 202, and adds a bounded retry-on-409 to the stream creation helpers so tests that delete and immediately recreate the same stream name (a common setup/teardown pattern here) keep working without needing changes at every call site. --- test_utils.go | 85 +++++++++++++++++++++++++++++++++++++-------------- 1 file changed, 62 insertions(+), 23 deletions(-) diff --git a/test_utils.go b/test_utils.go index caabfc9..ba07aed 100644 --- a/test_utils.go +++ b/test_utils.go @@ -21,6 +21,7 @@ import ( "encoding/json" "fmt" "io" + "net/http" "os/exec" "strings" "testing" @@ -63,40 +64,74 @@ func Sleep() { time.Sleep(sleepDuration) } +// doWithRetryOn409 retries a request while the server reports 409 Conflict, +// which now happens when a stream with the same name was deleted moments +// earlier and its deletion is still running in the background (deletion +// became asynchronous in parseablehq/parseable#1770 -- a stream isn't +// actually gone the instant DeleteStream returns). buildRequest is called +// again on every attempt since a request (and any body reader it holds) +// can't be reused once sent. +func doWithRetryOn409(t *testing.T, client HTTPClient, buildRequest func() (*http.Request, error)) *http.Response { + deadline := time.Now().Add(15 * time.Second) + for { + req, err := buildRequest() + require.NoErrorf(t, err, "failed to build request: %s", err) + response, err := client.Do(req) + require.NoErrorf(t, err, "Request failed: %s", err) + if response.StatusCode != 409 || time.Now().After(deadline) { + return response + } + response.Body.Close() + time.Sleep(200 * time.Millisecond) + } +} + func CreateStream(t *testing.T, client HTTPClient, stream string) { - req, _ := client.NewRequest("PUT", "logstream/"+stream, nil) - response, err := client.Do(req) - require.NoErrorf(t, err, "Request failed: %s", err) + response := doWithRetryOn409(t, client, func() (*http.Request, error) { + return client.NewRequest("PUT", "logstream/"+stream, nil) + }) require.Equalf(t, 200, response.StatusCode, "Server returned http code: %s", response.Status) } func CreateStreamWithHeader(t *testing.T, client HTTPClient, stream string, header map[string]string) { - req, _ := client.NewRequest("PUT", "logstream/"+stream, nil) - for k, v := range header { - req.Header.Add(k, v) - } - response, err := client.Do(req) - require.NoErrorf(t, err, "Request failed: %s", err) + response := doWithRetryOn409(t, client, func() (*http.Request, error) { + req, err := client.NewRequest("PUT", "logstream/"+stream, nil) + if err != nil { + return req, err + } + for k, v := range header { + req.Header.Add(k, v) + } + return req, nil + }) require.Equalf(t, 200, response.StatusCode, "Server returned http code: %s", response.Status) } func CreateStreamWithCustompartitionError(t *testing.T, client HTTPClient, stream string, header map[string]string) { - req, _ := client.NewRequest("PUT", "logstream/"+stream, nil) - for k, v := range header { - req.Header.Add(k, v) - } - response, _ := client.Do(req) + response := doWithRetryOn409(t, client, func() (*http.Request, error) { + req, err := client.NewRequest("PUT", "logstream/"+stream, nil) + if err != nil { + return req, err + } + for k, v := range header { + req.Header.Add(k, v) + } + return req, nil + }) require.Equalf(t, 500, response.StatusCode, "Server returned http code: %s", response.Status) } func CreateStreamWithSchemaBody(t *testing.T, client HTTPClient, stream string, header map[string]string, schema_payload string) { - - req, _ := client.NewRequest("PUT", "logstream/"+stream, bytes.NewBufferString(schema_payload)) - for k, v := range header { - req.Header.Add(k, v) - } - response, err := client.Do(req) - require.NoErrorf(t, err, "Request failed: %s", err) + response := doWithRetryOn409(t, client, func() (*http.Request, error) { + req, err := client.NewRequest("PUT", "logstream/"+stream, bytes.NewBufferString(schema_payload)) + if err != nil { + return req, err + } + for k, v := range header { + req.Header.Add(k, v) + } + return req, nil + }) require.Equalf(t, 200, response.StatusCode, "Server returned http code: %s", response.Status) } @@ -113,7 +148,11 @@ func DeleteStream(t *testing.T, client HTTPClient, stream string) { req, _ := client.NewRequest("DELETE", "logstream/"+stream, nil) response, err := client.Do(req) require.NoErrorf(t, err, "Request failed: %s", err) - require.Equalf(t, 200, response.StatusCode, "Server returned http code: %s", response.Status) + // Deletion runs in the background (parseablehq/parseable#1770); 202 means + // accepted, not finished. Callers that immediately recreate the same + // stream name rely on CreateStream et al. retrying past a transient 409 + // while the old stream's deletion is still in flight. + require.Equalf(t, 202, response.StatusCode, "Server returned http code: %s", response.Status) } func DeleteAlert(t *testing.T, client HTTPClient, alert_id string) { @@ -455,7 +494,7 @@ func checkAPIAccess(t *testing.T, queryClient HTTPClient, ingestClient HTTPClien req, _ = queryClient.NewRequest("DELETE", "logstream/"+stream, nil) response, err = queryClient.Do(req) require.NoErrorf(t, err, "Request failed: %s", err) - require.Equalf(t, 200, response.StatusCode, "Server returned http code: %s and response: %s", response.Status, readAsString(response.Body)) + require.Equalf(t, 202, response.StatusCode, "Server returned http code: %s and response: %s", response.Status, readAsString(response.Body)) case "writer": // Check access to non-protected API