Skip to content
Open
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
85 changes: 62 additions & 23 deletions test_utils.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ import (
"encoding/json"
"fmt"
"io"
"net/http"
"os/exec"
"strings"
"testing"
Expand Down Expand Up @@ -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)
}

Expand All @@ -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) {
Expand Down Expand Up @@ -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
Expand Down
Loading