From 9d35c4cad3c0676d20f387e26a14a0f00f6ea152 Mon Sep 17 00:00:00 2001 From: Josh Allmann Date: Fri, 4 Sep 2026 11:42:40 -0700 Subject: [PATCH 01/13] Update catalyst-api and go-livepeer --- manifest.yaml | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/manifest.yaml b/manifest.yaml index 2f57982c..121f609b 100644 --- a/manifest.yaml +++ b/manifest.yaml @@ -26,7 +26,7 @@ box: strategy: download: bucket project: catalyst-api - commit: 7c549ec156e81fbad359738e088e1d3fca559b80 + commit: 56aff2d85fe064e610dc15d5793c23199267fb5b release: main srcFilenames: darwin-amd64: livepeer-catalyst-api-darwin-amd64.tar.gz @@ -49,7 +49,7 @@ box: strategy: download: bucket project: go-livepeer - commit: 48c5d860fec714677f465a38c8a6d52a9f384154 + commit: 38eb47d12ab1d2d874fc4c7c061aa1900b7c0bad binary: livepeer release: master archivePath: livepeer From 00332e4076978c217924d646e92385974a0613f0 Mon Sep 17 00:00:00 2001 From: Josh Allmann Date: Fri, 4 Sep 2026 12:17:36 -0700 Subject: [PATCH 02/13] Fix MinIO downloads --- Dockerfile | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/Dockerfile b/Dockerfile index 5facef6e..da9e8ee6 100644 --- a/Dockerfile +++ b/Dockerfile @@ -120,8 +120,8 @@ RUN curl -L -O https://binaries.cockroachdb.com/cockroach-v23.1.5.linux-$TARGETA && rm -rf cockroach-v23.1.5.linux-$TARGETARCH.tgz cockroach-v23.1.5.linux-$TARGETARCH \ && cockroach --version -RUN curl -o /usr/bin/minio https://dl.min.io/server/minio/release/linux-$TARGETARCH/minio \ - && curl -o /usr/bin/mc https://dl.min.io/client/mc/release/linux-$TARGETARCH/mc \ +RUN curl -fL -o /usr/bin/minio https://dl.min.io/server/minio/release/linux-$TARGETARCH/minio \ + && curl -fL -o /usr/bin/mc https://dl.min.io/client/mc/release/linux-$TARGETARCH/mc \ && chmod +x /usr/bin/minio /usr/bin/mc \ && minio --version \ && mc --version From 0a64b1ac1c9af6be95380a57585f01dbe5ff287b Mon Sep 17 00:00:00 2001 From: Josh Allmann Date: Fri, 4 Sep 2026 17:48:42 -0700 Subject: [PATCH 03/13] Use quick tunnels in E2E tests --- .github/workflows/test.yaml | 12 +++ scripts/livepeer-nginx | 16 +++- test/e2e/box_record_test.go | 137 +++++++++++++++++++++++++++- test/e2e/mist_config.go | 9 ++ test/e2e/quick_tunnel_test.go | 167 ++++++++++++++++++++++++++++++++++ test/e2e/vod_test.go | 58 +++++++++--- 6 files changed, 381 insertions(+), 18 deletions(-) create mode 100644 test/e2e/quick_tunnel_test.go diff --git a/.github/workflows/test.yaml b/.github/workflows/test.yaml index e0ee9876..6f1037ce 100644 --- a/.github/workflows/test.yaml +++ b/.github/workflows/test.yaml @@ -63,6 +63,18 @@ jobs: - name: Build Docker Box image run: make box + - name: Install cloudflared + env: + CLOUDFLARED_VERSION: 2026.7.3 + CLOUDFLARED_SHA256: 9d71c677db00134c1bd4144b7783486b654ad281b1ea62b4972098d19f770f17 + run: | + curl -fL --retry 3 \ + "https://github.com/cloudflare/cloudflared/releases/download/${CLOUDFLARED_VERSION}/cloudflared-linux-amd64" \ + -o "${RUNNER_TEMP}/cloudflared" + echo "${CLOUDFLARED_SHA256} ${RUNNER_TEMP}/cloudflared" | sha256sum --check + chmod +x "${RUNNER_TEMP}/cloudflared" + echo "${RUNNER_TEMP}" >> "${GITHUB_PATH}" + - name: Run E2E tests run: go test $(go list ./... | grep 'test/e2e') --timeout 15m --image diff --git a/scripts/livepeer-nginx b/scripts/livepeer-nginx index e141502e..4b79a353 100755 --- a/scripts/livepeer-nginx +++ b/scripts/livepeer-nginx @@ -88,20 +88,28 @@ http { proxy_pass http://127.0.0.1:3080; } - location /os-vod/ { + location /os-vod { proxy_pass http://127.0.0.1:9000; + proxy_set_header Host \$http_host; } - location /os-catalyst-vod/ { + location /os-catalyst-vod { proxy_pass http://127.0.0.1:9000; + proxy_set_header Host \$http_host; } - location /os-private/ { + location /os-private { proxy_pass http://127.0.0.1:9000; + proxy_set_header Host \$http_host; } - location /os-recordings/ { + location /os-recordings { proxy_pass http://127.0.0.1:9000; + proxy_set_header Host \$http_host; + } + + location /task-runner { + proxy_pass http://127.0.0.1:3060; } location / { diff --git a/test/e2e/box_record_test.go b/test/e2e/box_record_test.go index 5c2c3447..271ddf53 100644 --- a/test/e2e/box_record_test.go +++ b/test/e2e/box_record_test.go @@ -3,13 +3,20 @@ package e2e import ( "bytes" "context" + "encoding/json" "fmt" + "io" + "net/http" + "net/url" "os" "os/exec" + "strings" "sync" "testing" "time" + "github.com/minio/minio-go/v7" + "github.com/minio/minio-go/v7/pkg/credentials" "github.com/stretchr/testify/require" "github.com/testcontainers/testcontainers-go" "github.com/testcontainers/testcontainers-go/wait" @@ -29,10 +36,13 @@ func TestBoxRecording(t *testing.T) { defer network.Remove(ctx) boxName := randomString("box-") + publicURL := startQuickTunnel(t, "http://127.0.0.1:8888") // when - box := startBoxWithEnv(ctx, t, boxName, network.name) + box := startBoxWithEnv(ctx, t, boxName, network.name, publicURL) defer box.Terminate(ctx) + waitForBoxMinio(t, publicURL) + configureBoxObjectStores(t, publicURL) eg, ctx := errgroup.WithContext(ctx) eg.Go(func() error { @@ -44,7 +54,7 @@ func TestBoxRecording(t *testing.T) { require.NoError(t, eg.Wait()) } -func startBoxWithEnv(ctx context.Context, t *testing.T, hostname, network string) *catalystContainer { +func startBoxWithEnv(ctx context.Context, t *testing.T, hostname, network, publicURL string) *catalystContainer { req := testcontainers.ContainerRequest{ Image: "livepeer/in-a-box", Hostname: hostname, @@ -55,6 +65,18 @@ func startBoxWithEnv(ctx context.Context, t *testing.T, hostname, network string WaitingFor: wait.NewLogStrategy("API server listening").WithStartupTimeout(3 * time.Minute), Env: map[string]string{ "LP_API_FRONTEND": "false", + "E2E_PUBLIC_URL": publicURL, + }, + Cmd: []string{ + "bash", + "-ceu", + `sed -i \ + -e "s|\"api-server\": \"http://127.0.0.1:3004\"|\"api-server\": \"${E2E_PUBLIC_URL}\"|" \ + -e "s|\"own-base-url\": \"http://127.0.0.1:3060/task-runner\"|\"own-base-url\": \"${E2E_PUBLIC_URL}/task-runner\"|" \ + /etc/livepeer/full-stack.json +grep -Fq "\"api-server\": \"${E2E_PUBLIC_URL}\"" /etc/livepeer/full-stack.json +grep -Fq "\"own-base-url\": \"${E2E_PUBLIC_URL}/task-runner\"" /etc/livepeer/full-stack.json +exec /usr/local/bin/catalyst -- /usr/local/bin/MistController -c /etc/livepeer/full-stack.json`, }, } container, err := testcontainers.GenericContainer(ctx, testcontainers.GenericContainerRequest{ @@ -97,6 +119,117 @@ func startBoxWithEnv(ctx context.Context, t *testing.T, hostname, network string return catalyst } +const boxAPIToken = "f61b3cdb-d173-4a7a-a0d3-547b871a56f9" + +var boxObjectStores = map[string]string{ + "917a2f18-f7a8-4ae3-a849-6efd4aac8e59": "os-vod", + "517873a4-487c-40ad-872f-027f4bc6bd98": "os-catalyst-vod", + "cab9266f-5583-4532-9630-7be10d92affe": "os-private", + "0926e4ba-b726-4386-92ee-5c4583f62f0a": "os-recordings", +} + +func waitForBoxMinio(t *testing.T, publicURL string) { + t.Helper() + + u, err := url.Parse(publicURL) + require.NoError(t, err) + cli, err := minio.New(u.Host, &minio.Options{ + Creds: credentials.NewStaticV4("admin", "password", ""), + Secure: true, + Region: region, + BucketLookup: minio.BucketLookupPath, + }) + require.NoError(t, err) + + deadline := time.Now().Add(time.Minute) + for { + requestCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + exists, err := cli.BucketExists(requestCtx, "os-recordings") + cancel() + if err == nil && exists { + return + } + if err == nil { + err = fmt.Errorf("recordings bucket does not exist") + } + if time.Now().After(deadline) { + t.Fatalf("in-a-box MinIO tunnel did not become ready: %v", err) + } + time.Sleep(time.Second) + } +} + +func configureBoxObjectStores(t *testing.T, publicURL string) { + t.Helper() + + client := &http.Client{Timeout: 10 * time.Second} + for id, bucket := range boxObjectStores { + deadline := time.Now().Add(time.Minute) + for { + err := patchBoxObjectStore(client, publicURL, id, bucket) + if err == nil { + break + } + if time.Now().After(deadline) { + t.Fatalf("could not configure object store %s: %v", bucket, err) + } + time.Sleep(time.Second) + } + } +} + +func patchBoxObjectStore(client *http.Client, publicURL, id, bucket string) error { + storeURL, err := boxObjectStoreURL(publicURL, bucket) + if err != nil { + return err + } + payload, err := json.Marshal(map[string]string{ + "url": storeURL, + "publicUrl": strings.TrimRight(publicURL, "/") + "/" + bucket, + }) + if err != nil { + return err + } + + req, err := http.NewRequest(http.MethodPatch, strings.TrimRight(publicURL, "/")+"/api/object-store/"+id, bytes.NewReader(payload)) + if err != nil { + return err + } + req.Header.Set("Authorization", "Bearer "+boxAPIToken) + req.Header.Set("Content-Type", "application/json") + + resp, err := client.Do(req) + if err != nil { + return err + } + defer resp.Body.Close() + if resp.StatusCode != http.StatusNoContent { + body, _ := io.ReadAll(io.LimitReader(resp.Body, 4096)) + return fmt.Errorf("unexpected status %s: %s", resp.Status, strings.TrimSpace(string(body))) + } + return nil +} + +func boxObjectStoreURL(publicURL, bucket string) (string, error) { + u, err := url.Parse(publicURL) + if err != nil { + return "", err + } + if u.Scheme != "https" || u.Hostname() == "" { + return "", fmt.Errorf("invalid tunnel URL %q", publicURL) + } + u.Scheme = "s3+https" + u.User = url.UserPassword("admin", "password") + u.Path = "/" + bucket + return u.String(), nil +} + +func TestBoxObjectStoreURL(t *testing.T) { + got, err := boxObjectStoreURL("https://example.trycloudflare.com", "os-recordings") + require.NoError(t, err) + require.Equal(t, "s3+https://admin:password@example.trycloudflare.com/os-recordings", got) +} + func startRecordTester(ctx context.Context, recordingCopyOnly bool) error { startTime := time.Now() fmt.Printf("starting record tester copyOnly=%v\n", recordingCopyOnly) diff --git a/test/e2e/mist_config.go b/test/e2e/mist_config.go index 1f463ba8..de3a4d5d 100644 --- a/test/e2e/mist_config.go +++ b/test/e2e/mist_config.go @@ -25,6 +25,7 @@ type bandwidth struct { type protocol struct { Connector string `json:"connector"` + APIServer string `json:"api-server,omitempty"` RetryJoin string `json:"retry-join,omitempty"` Advertise string `json:"advertise,omitempty"` RPCAddr string `json:"rpc-addr,omitempty"` @@ -44,6 +45,14 @@ type protocol struct { Catabalancer string `json:"catabalancer,omitempty"` } +func (m *mistConfig) setAPIServer(apiServer string) { + for i := range m.Config.Protocols { + if m.Config.Protocols[i].Connector == "livepeer-catalyst-api" { + m.Config.Protocols[i].APIServer = apiServer + } + } +} + type config struct { Accesslog string `json:"accesslog"` Controller struct { diff --git a/test/e2e/quick_tunnel_test.go b/test/e2e/quick_tunnel_test.go new file mode 100644 index 00000000..5302efa5 --- /dev/null +++ b/test/e2e/quick_tunnel_test.go @@ -0,0 +1,167 @@ +package e2e + +import ( + "bytes" + "context" + "fmt" + "net/http" + "net/http/httptest" + "net/url" + "os/exec" + "regexp" + "strings" + "sync" + "testing" + "time" + + "github.com/stretchr/testify/require" +) + +var quickTunnelURLPattern = regexp.MustCompile(`https://[a-z0-9-]+\.trycloudflare\.com`) + +type synchronizedBuffer struct { + mu sync.Mutex + bytes.Buffer +} + +func (b *synchronizedBuffer) Write(p []byte) (int, error) { + b.mu.Lock() + defer b.mu.Unlock() + return b.Buffer.Write(p) +} + +func (b *synchronizedBuffer) String() string { + b.mu.Lock() + defer b.mu.Unlock() + return b.Buffer.String() +} + +// startQuickTunnel exposes origin through a temporary trycloudflare.com URL. +// The cloudflared process is stopped automatically when the test completes. +func startQuickTunnel(t *testing.T, origin string) string { + t.Helper() + + originURL, err := url.Parse(origin) + require.NoError(t, err) + require.Equal(t, "http", originURL.Scheme) + require.NotEmpty(t, originURL.Host) + + cloudflared, err := exec.LookPath("cloudflared") + require.NoError(t, err, "cloudflared is required to run the E2E tests") + + ctx, cancel := context.WithCancel(context.Background()) + output := &synchronizedBuffer{} + cmd := exec.CommandContext( + ctx, + cloudflared, + "tunnel", + "--no-autoupdate", + "--protocol", "http2", + "--url", origin, + ) + cmd.Stdout = output + cmd.Stderr = output + require.NoError(t, cmd.Start()) + + done := make(chan error, 1) + go func() { + done <- cmd.Wait() + }() + + stop := func() { + cancel() + select { + case <-done: + case <-time.After(5 * time.Second): + _ = cmd.Process.Kill() + <-done + } + } + + ticker := time.NewTicker(100 * time.Millisecond) + defer ticker.Stop() + timer := time.NewTimer(time.Minute) + defer timer.Stop() + + for { + select { + case err := <-done: + cancel() + t.Fatalf("cloudflared exited before creating a tunnel: %v\n%s", err, output.String()) + case <-ticker.C: + if publicURL := quickTunnelURLPattern.FindString(output.String()); publicURL != "" { + t.Cleanup(stop) + t.Logf("quick tunnel %s -> %s", publicURL, origin) + return publicURL + } + case <-timer.C: + stop() + t.Fatalf("timed out creating a quick tunnel for %s\n%s", origin, output.String()) + } + } +} + +func startCallbackTunnel(t *testing.T) (apiServerURL, callbackURL string) { + t.Helper() + + callbackServer := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if !strings.HasPrefix(r.URL.Path, "/task-runner") { + http.NotFound(w, r) + return + } + w.WriteHeader(http.StatusNoContent) + })) + t.Cleanup(callbackServer.Close) + + apiServerURL = startQuickTunnel(t, callbackServer.URL) + waitForPublicHTTP(t, apiServerURL+"/task-runner/ready") + return apiServerURL, apiServerURL + "/task-runner/vod-test" +} + +func waitForPublicHTTP(t *testing.T, endpoint string) { + t.Helper() + + client := &http.Client{Timeout: 5 * time.Second} + deadline := time.Now().Add(time.Minute) + for { + resp, err := client.Get(endpoint) + if err == nil { + resp.Body.Close() + if resp.StatusCode >= http.StatusOK && resp.StatusCode < http.StatusMultipleChoices { + return + } + err = fmt.Errorf("unexpected status %s", resp.Status) + } + if time.Now().After(deadline) { + t.Fatalf("tunnel did not become ready: %v", err) + } + time.Sleep(time.Second) + } +} + +func tunneledObjectStoreURL(t *testing.T, publicURL, username, password, bucket, key string) string { + t.Helper() + + u, err := url.Parse(publicURL) + require.NoError(t, err) + require.Equal(t, "https", u.Scheme) + require.NotEmpty(t, u.Hostname()) + + u.Scheme = "s3+https" + u.User = url.UserPassword(username, password) + u.Path = "/" + strings.Trim(strings.Join([]string{bucket, key}, "/"), "/") + return u.String() +} + +func TestTunneledObjectStoreURL(t *testing.T) { + require.Equal( + t, + "s3+https://access:secret@example.trycloudflare.com/bucket/path/source.mp4", + tunneledObjectStoreURL(t, "https://example.trycloudflare.com", "access", "secret", "bucket", "path/source.mp4"), + ) + require.Equal( + t, + "s3+https://access:secret@example.trycloudflare.com/bucket", + tunneledObjectStoreURL(t, "https://example.trycloudflare.com", "access", "secret", "bucket", ""), + ) +} diff --git a/test/e2e/vod_test.go b/test/e2e/vod_test.go index 28dc166b..91f3afa2 100644 --- a/test/e2e/vod_test.go +++ b/test/e2e/vod_test.go @@ -4,7 +4,10 @@ import ( "bytes" "context" "fmt" + "io" "net/http" + "net/url" + "strings" "testing" "time" @@ -41,14 +44,19 @@ func TestVod(t *testing.T) { createSourceBucket(t, m) uploadSourceVideo(ctx, t, m) createDestBucket(t, m) + storageURL := startQuickTunnel(t, fmt.Sprintf("http://127.0.0.1:%s", m.port)) + waitForTunneledMinio(ctx, t, storageURL) + apiServerURL, callbackURL := startCallbackTunnel(t) h := randomString("catalyst-") - c := startCatalyst(ctx, t, h, network.name, defaultMistConfigWithLivepeerProcess(h, sourceOutput(m))) + mistConfig := defaultMistConfigWithLivepeerProcess(h, tunneledObjectStoreURL(t, storageURL, username, password, inBucket, "")) + mistConfig.setAPIServer(apiServerURL) + c := startCatalyst(ctx, t, h, network.name, mistConfig) defer c.Terminate(ctx) waitForCatalystAPI(t, c) // when - processVod(t, m, c) + processVod(t, storageURL, callbackURL, c) // then requireOutputFiles(ctx, t, m) @@ -115,14 +123,38 @@ func createSourceBucket(t *testing.T, m *minioContainer) { createBucket(t, m, inBucket) } -func sourceOutput(m *minioContainer) string { - return fmt.Sprintf("s3+http://%s:%s@%s:9000/%s", username, password, m.hostname, inBucket) -} - func createDestBucket(t *testing.T, m *minioContainer) { createBucket(t, m, outBucket) } +func waitForTunneledMinio(ctx context.Context, t *testing.T, storageURL string) { + t.Helper() + + u, err := url.Parse(storageURL) + require.NoError(t, err) + cli, err := minio.New(u.Host, &minio.Options{ + Creds: credentials.NewStaticV4(username, password, ""), + Secure: true, + Region: region, + BucketLookup: minio.BucketLookupPath, + }) + require.NoError(t, err) + + deadline := time.Now().Add(time.Minute) + for { + requestCtx, cancel := context.WithTimeout(ctx, 5*time.Second) + _, err = cli.StatObject(requestCtx, inBucket, source, minio.StatObjectOptions{}) + cancel() + if err == nil { + return + } + if time.Now().After(deadline) { + t.Fatalf("MinIO tunnel did not become ready: %v", err) + } + time.Sleep(time.Second) + } +} + func createBucket(t *testing.T, m *minioContainer, bucket string) { err := minioClient(t, m).MakeBucket(context.Background(), bucket, minio.MakeBucketOptions{Region: region, ObjectLocking: true}) require.NoError(t, err) @@ -151,12 +183,12 @@ func waitForCatalystAPI(t *testing.T, c *catalystContainer) { require.Eventually(t, catalystAPIStarted, 5*time.Minute, time.Second) } -func processVod(t *testing.T, m *minioContainer, c *catalystContainer) { - sourceVideoURL := fmt.Sprintf("s3+http://%s:%s@%s:9000/%s/%s", username, password, m.hostname, inBucket, source) - destURL := fmt.Sprintf("s3+http://%s:%s@%s:9000/%s/", username, password, m.hostname, outBucket) +func processVod(t *testing.T, storageURL, callbackURL string, c *catalystContainer) { + sourceVideoURL := tunneledObjectStoreURL(t, storageURL, username, password, inBucket, source) + destURL := tunneledObjectStoreURL(t, storageURL, username, password, outBucket, "") var jsonData = fmt.Sprintf(`{ - "url": "%s", - "callback_url": "https://todo-callback.com", + "url": "%s", + "callback_url": "%s", "output_locations": [ { "type": "object_store", @@ -166,7 +198,7 @@ func processVod(t *testing.T, m *minioContainer, c *catalystContainer) { } } ] - }`, sourceVideoURL, destURL) + }`, sourceVideoURL, callbackURL, destURL) url := fmt.Sprintf("http://127.0.0.1:%s/api/vod", c.catalystAPIInternal) req, err := http.NewRequest("POST", url, bytes.NewBuffer([]byte(jsonData))) @@ -177,6 +209,8 @@ func processVod(t *testing.T, m *minioContainer, c *catalystContainer) { resp, err := client.Do(req) require.NoError(t, err) defer resp.Body.Close() + body, _ := io.ReadAll(io.LimitReader(resp.Body, 4096)) + require.Equal(t, http.StatusOK, resp.StatusCode, "unexpected response: %s", strings.TrimSpace(string(body))) } func requireOutputFiles(ctx context.Context, t *testing.T, m *minioContainer) { From 51f2c527ddbb5be5bdb2f010d531ae0c637fe0b2 Mon Sep 17 00:00:00 2001 From: Josh Allmann Date: Fri, 4 Sep 2026 17:56:57 -0700 Subject: [PATCH 04/13] Avoid Docker build for misspell CI --- .github/workflows/test.yaml | 31 ++++++++++++++++++++++++++++++- 1 file changed, 30 insertions(+), 1 deletion(-) diff --git a/.github/workflows/test.yaml b/.github/workflows/test.yaml index 6f1037ce..12f15fc8 100644 --- a/.github/workflows/test.yaml +++ b/.github/workflows/test.yaml @@ -46,8 +46,37 @@ jobs: with: config: config.toml + - name: Set up reviewdog + uses: reviewdog/action-setup@d8a7baabd7f3e8544ee4dbde3ee41d0011c3a93f # v1.5.0 + with: + reviewdog_version: v0.21.0 + + - name: Install misspell + env: + MISSPELL_VERSION: 0.8.0 + MISSPELL_SHA256: 0a889624797be3eed241eeef4dd2b982adadedd6556723fd00b9c3f721417b51 + run: | + archive="${RUNNER_TEMP}/misspell.tar.gz" + curl -fL --retry 3 \ + "https://github.com/golangci/misspell/releases/download/v${MISSPELL_VERSION}/misspell_${MISSPELL_VERSION}_linux_amd64.tar.gz" \ + -o "${archive}" + echo "${MISSPELL_SHA256} ${archive}" | sha256sum --check + tar -xzf "${archive}" --strip-components=1 -C "${RUNNER_TEMP}" \ + "misspell_${MISSPELL_VERSION}_linux_amd64/misspell" + echo "${RUNNER_TEMP}" >> "${GITHUB_PATH}" + - name: misspell - uses: reviewdog/action-misspell@v1 + env: + REVIEWDOG_GITHUB_API_TOKEN: ${{ secrets.GITHUB_TOKEN }} + run: | + find . -type f -print0 \ + | xargs -0 misspell \ + | reviewdog -efm="%f:%l:%c: %m" \ + -filter-mode=added \ + -name=misspell \ + -reporter=github-pr-check \ + -level=error \ + -fail-level=none - name: Install FFMPEG uses: FedericoCarboni/setup-ffmpeg@v2 From c83cf4e6e65de0e2a31493fc113bfc4784237135 Mon Sep 17 00:00:00 2001 From: Josh Allmann Date: Fri, 4 Sep 2026 18:28:47 -0700 Subject: [PATCH 05/13] Wait for E2E completion callbacks --- .github/workflows/test.yaml | 2 +- test/e2e/box_record_test.go | 23 ++++++---- test/e2e/e2e_test.go | 18 ++++++++ test/e2e/quick_tunnel_test.go | 84 +++++++++++++++++++++++++++++++++-- test/e2e/vod_test.go | 21 ++++++--- 5 files changed, 129 insertions(+), 19 deletions(-) diff --git a/.github/workflows/test.yaml b/.github/workflows/test.yaml index 12f15fc8..7544c4f5 100644 --- a/.github/workflows/test.yaml +++ b/.github/workflows/test.yaml @@ -106,7 +106,7 @@ jobs: - name: Run E2E tests run: - go test $(go list ./... | grep 'test/e2e') --timeout 15m --image + go test $(go list ./... | grep 'test/e2e') --timeout 30m -v --image catalyst-e2e-test - name: Upload coverage reports diff --git a/test/e2e/box_record_test.go b/test/e2e/box_record_test.go index 271ddf53..fe82ca6c 100644 --- a/test/e2e/box_record_test.go +++ b/test/e2e/box_record_test.go @@ -20,7 +20,6 @@ import ( "github.com/stretchr/testify/require" "github.com/testcontainers/testcontainers-go" "github.com/testcontainers/testcontainers-go/wait" - "golang.org/x/sync/errgroup" ) func TestBoxRecording(t *testing.T) { @@ -44,14 +43,20 @@ func TestBoxRecording(t *testing.T) { waitForBoxMinio(t, publicURL) configureBoxObjectStores(t, publicURL) - eg, ctx := errgroup.WithContext(ctx) - eg.Go(func() error { - return startRecordTester(ctx, false) - }) - eg.Go(func() error { - return startRecordTester(ctx, true) - }) - require.NoError(t, eg.Wait()) + for _, mode := range []struct { + name string + copyOnly bool + }{ + {name: "copy-only", copyOnly: true}, + {name: "transcoded", copyOnly: false}, + } { + t.Run(mode.name, func(t *testing.T) { + if err := startRecordTester(ctx, mode.copyOnly); err != nil { + dumpContainerLogs(ctx, t, box.Container) + require.NoError(t, err) + } + }) + } } func startBoxWithEnv(ctx context.Context, t *testing.T, hostname, network, publicURL string) *catalystContainer { diff --git a/test/e2e/e2e_test.go b/test/e2e/e2e_test.go index 1e9111c7..3133ec2f 100644 --- a/test/e2e/e2e_test.go +++ b/test/e2e/e2e_test.go @@ -138,6 +138,24 @@ func (lc *logConsumer) Accept(l testcontainers.Log) { glog.Infof("[%s] %s", lc.name, string(l.Content)) } +func dumpContainerLogs(ctx context.Context, t *testing.T, container testcontainers.Container) { + t.Helper() + + logs, err := container.Logs(ctx) + if err != nil { + t.Logf("could not retrieve container logs: %v", err) + return + } + defer logs.Close() + + output, err := io.ReadAll(io.LimitReader(logs, 4<<20)) + if err != nil { + t.Logf("could not read container logs: %v", err) + return + } + t.Logf("container logs:\n%s", output) +} + func startCatalyst(ctx context.Context, t *testing.T, hostname, network string, mc mistConfig) *catalystContainer { return startCatalystWithEnv(ctx, t, hostname, network, mc, nil) } diff --git a/test/e2e/quick_tunnel_test.go b/test/e2e/quick_tunnel_test.go index 5302efa5..932a156c 100644 --- a/test/e2e/quick_tunnel_test.go +++ b/test/e2e/quick_tunnel_test.go @@ -3,7 +3,9 @@ package e2e import ( "bytes" "context" + "encoding/json" "fmt" + "io" "net/http" "net/http/httptest" "net/url" @@ -101,21 +103,81 @@ func startQuickTunnel(t *testing.T, origin string) string { } } -func startCallbackTunnel(t *testing.T) (apiServerURL, callbackURL string) { +type callbackStatus struct { + RequestID string `json:"request_id"` + Status string `json:"status"` + Error string `json:"error,omitempty"` +} + +type callbackTunnel struct { + apiServerURL string + callbackURL string + terminal chan callbackStatus + mu sync.Mutex + last callbackStatus +} + +func startCallbackTunnel(t *testing.T) *callbackTunnel { t.Helper() + callbacks := &callbackTunnel{terminal: make(chan callbackStatus, 1)} callbackServer := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { if !strings.HasPrefix(r.URL.Path, "/task-runner") { http.NotFound(w, r) return } + if r.Method == http.MethodPost { + body, err := io.ReadAll(io.LimitReader(r.Body, 1<<20)) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + var status callbackStatus + if err := json.Unmarshal(body, &status); err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + callbacks.mu.Lock() + callbacks.last = status + callbacks.mu.Unlock() + if status.Status == "success" || status.Status == "error" { + select { + case callbacks.terminal <- status: + default: + } + } + } w.WriteHeader(http.StatusNoContent) })) t.Cleanup(callbackServer.Close) - apiServerURL = startQuickTunnel(t, callbackServer.URL) - waitForPublicHTTP(t, apiServerURL+"/task-runner/ready") - return apiServerURL, apiServerURL + "/task-runner/vod-test" + callbacks.apiServerURL = startQuickTunnel(t, callbackServer.URL) + callbacks.callbackURL = callbacks.apiServerURL + "/task-runner/vod-test" + waitForPublicHTTP(t, callbacks.apiServerURL+"/task-runner/ready") + return callbacks +} + +func (c *callbackTunnel) waitForCompletion(requestID string, timeout time.Duration) error { + timer := time.NewTimer(timeout) + defer timer.Stop() + + for { + select { + case status := <-c.terminal: + if status.RequestID != requestID { + continue + } + if status.Status == "error" { + return fmt.Errorf("VOD job %s failed: %s", requestID, status.Error) + } + return nil + case <-timer.C: + c.mu.Lock() + last := c.last + c.mu.Unlock() + return fmt.Errorf("timed out waiting for VOD job %s; last callback: %+v", requestID, last) + } + } } func waitForPublicHTTP(t *testing.T, endpoint string) { @@ -165,3 +227,17 @@ func TestTunneledObjectStoreURL(t *testing.T) { tunneledObjectStoreURL(t, "https://example.trycloudflare.com", "access", "secret", "bucket", ""), ) } + +func TestWaitForCompletion(t *testing.T) { + t.Run("success", func(t *testing.T) { + callbacks := &callbackTunnel{terminal: make(chan callbackStatus, 1)} + callbacks.terminal <- callbackStatus{RequestID: "request-1", Status: "success"} + require.NoError(t, callbacks.waitForCompletion("request-1", time.Second)) + }) + + t.Run("error", func(t *testing.T) { + callbacks := &callbackTunnel{terminal: make(chan callbackStatus, 1)} + callbacks.terminal <- callbackStatus{RequestID: "request-1", Status: "error", Error: "transcode failed"} + require.EqualError(t, callbacks.waitForCompletion("request-1", time.Second), "VOD job request-1 failed: transcode failed") + }) +} diff --git a/test/e2e/vod_test.go b/test/e2e/vod_test.go index 91f3afa2..7b54c460 100644 --- a/test/e2e/vod_test.go +++ b/test/e2e/vod_test.go @@ -3,6 +3,7 @@ package e2e import ( "bytes" "context" + "encoding/json" "fmt" "io" "net/http" @@ -46,17 +47,21 @@ func TestVod(t *testing.T) { createDestBucket(t, m) storageURL := startQuickTunnel(t, fmt.Sprintf("http://127.0.0.1:%s", m.port)) waitForTunneledMinio(ctx, t, storageURL) - apiServerURL, callbackURL := startCallbackTunnel(t) + callbacks := startCallbackTunnel(t) h := randomString("catalyst-") mistConfig := defaultMistConfigWithLivepeerProcess(h, tunneledObjectStoreURL(t, storageURL, username, password, inBucket, "")) - mistConfig.setAPIServer(apiServerURL) + mistConfig.setAPIServer(callbacks.apiServerURL) c := startCatalyst(ctx, t, h, network.name, mistConfig) defer c.Terminate(ctx) waitForCatalystAPI(t, c) // when - processVod(t, storageURL, callbackURL, c) + requestID := processVod(t, storageURL, callbacks.callbackURL, c) + if err := callbacks.waitForCompletion(requestID, 10*time.Minute); err != nil { + dumpContainerLogs(ctx, t, c.Container) + require.NoError(t, err) + } // then requireOutputFiles(ctx, t, m) @@ -183,7 +188,7 @@ func waitForCatalystAPI(t *testing.T, c *catalystContainer) { require.Eventually(t, catalystAPIStarted, 5*time.Minute, time.Second) } -func processVod(t *testing.T, storageURL, callbackURL string, c *catalystContainer) { +func processVod(t *testing.T, storageURL, callbackURL string, c *catalystContainer) string { sourceVideoURL := tunneledObjectStoreURL(t, storageURL, username, password, inBucket, source) destURL := tunneledObjectStoreURL(t, storageURL, username, password, outBucket, "") var jsonData = fmt.Sprintf(`{ @@ -211,12 +216,18 @@ func processVod(t *testing.T, storageURL, callbackURL string, c *catalystContain defer resp.Body.Close() body, _ := io.ReadAll(io.LimitReader(resp.Body, 4096)) require.Equal(t, http.StatusOK, resp.StatusCode, "unexpected response: %s", strings.TrimSpace(string(body))) + var result struct { + RequestID string `json:"request_id"` + } + require.NoError(t, json.Unmarshal(body, &result)) + require.NotEmpty(t, result.RequestID) + return result.RequestID } func requireOutputFiles(ctx context.Context, t *testing.T, m *minioContainer) { cli := minioClient(t, m) var files []string - timeoutAt := time.Now().Add(5 * time.Minute) + timeoutAt := time.Now().Add(30 * time.Second) expectedFiles := []string{ "index.m3u8", From 6b47d34041f09f6bd68a45be4191498c29d553b4 Mon Sep 17 00:00:00 2001 From: Josh Allmann Date: Fri, 4 Sep 2026 18:30:29 -0700 Subject: [PATCH 06/13] Avoid Docker build for artifact signing --- .github/workflows/build.yaml | 32 ++++++++++++++++++++++++++------ 1 file changed, 26 insertions(+), 6 deletions(-) diff --git a/.github/workflows/build.yaml b/.github/workflows/build.yaml index ebadc165..af70f284 100644 --- a/.github/workflows/build.yaml +++ b/.github/workflows/build.yaml @@ -109,12 +109,32 @@ jobs: path: releases/ - name: Generate sha256 checksum and gpg signatures for release artifacts - uses: livepeer/action-gh-checksum-and-gpg-sign@latest - with: - artifacts-dir: releases - release-name: ${{ (github.ref_type == 'tag' && github.ref_name) || github.event.pull_request.head.sha || github.sha }} - gpg-key: ${{ secrets.CI_GPG_SIGNING_KEY }} - gpg-key-passphrase: ${{ secrets.CI_GPG_SIGNING_PASSPHRASE }} + env: + RELEASE_NAME: ${{ (github.ref_type == 'tag' && github.ref_name) || github.event.pull_request.head.sha || github.sha }} + CI_GPG_SIGNING_KEY: ${{ secrets.CI_GPG_SIGNING_KEY }} + CI_GPG_SIGNING_KEY_PASSPHRASE: ${{ secrets.CI_GPG_SIGNING_PASSPHRASE }} + run: | + cd releases + sha256sum -- * > "${RELEASE_NAME}_checksums.txt" + + if [[ -z "${CI_GPG_SIGNING_KEY}" ]] || ! command -v gpg >/dev/null; then + exit 0 + fi + + printf '%s\n' "${CI_GPG_SIGNING_KEY}" | gpg --batch --import + for file in *; do + if [[ "${file}" == *.txt ]]; then + continue + fi + gpg \ + --batch \ + --no-tty \ + --passphrase "${CI_GPG_SIGNING_KEY_PASSPHRASE}" \ + --pinentry-mode loopback \ + --output "${file}.sig" \ + --detach-sign \ + "${file}" + done - name: Generate branch manifest id: branch-manifest From a638e8e3d99cfab5cec5717cda0af11a1ae2cff1 Mon Sep 17 00:00:00 2001 From: Josh Allmann Date: Fri, 4 Sep 2026 20:04:06 -0700 Subject: [PATCH 07/13] Tunnel E2E orchestrator endpoint --- test/e2e/box_record_test.go | 15 ++++++++++----- test/e2e/e2e_test.go | 34 ++++++++++++++++++++-------------- test/e2e/mist_config.go | 16 ++++++++++++++++ test/e2e/quick_tunnel_test.go | 17 +++++++++++++++++ test/e2e/vod_test.go | 4 +++- 5 files changed, 66 insertions(+), 20 deletions(-) diff --git a/test/e2e/box_record_test.go b/test/e2e/box_record_test.go index fe82ca6c..c7346087 100644 --- a/test/e2e/box_record_test.go +++ b/test/e2e/box_record_test.go @@ -36,9 +36,10 @@ func TestBoxRecording(t *testing.T) { boxName := randomString("box-") publicURL := startQuickTunnel(t, "http://127.0.0.1:8888") + orchestratorURL, orchestratorHostPort := startOrchestratorTunnel(t) // when - box := startBoxWithEnv(ctx, t, boxName, network.name, publicURL) + box := startBoxWithEnv(ctx, t, boxName, network.name, publicURL, orchestratorURL, orchestratorHostPort) defer box.Terminate(ctx) waitForBoxMinio(t, publicURL) configureBoxObjectStores(t, publicURL) @@ -59,18 +60,19 @@ func TestBoxRecording(t *testing.T) { } } -func startBoxWithEnv(ctx context.Context, t *testing.T, hostname, network, publicURL string) *catalystContainer { +func startBoxWithEnv(ctx context.Context, t *testing.T, hostname, network, publicURL, orchestratorURL, orchestratorHostPort string) *catalystContainer { req := testcontainers.ContainerRequest{ Image: "livepeer/in-a-box", Hostname: hostname, Name: hostname, Networks: []string{network}, - ExposedPorts: []string{"1935:1935/tcp", "8888:8888/tcp"}, + ExposedPorts: []string{"1935:1935/tcp", "8888:8888/tcp", fmt.Sprintf("%s:8936/tcp", orchestratorHostPort)}, ShmSize: 1000000000, WaitingFor: wait.NewLogStrategy("API server listening").WithStartupTimeout(3 * time.Minute), Env: map[string]string{ - "LP_API_FRONTEND": "false", - "E2E_PUBLIC_URL": publicURL, + "LP_API_FRONTEND": "false", + "E2E_PUBLIC_URL": publicURL, + "E2E_ORCHESTRATOR_URL": orchestratorURL, }, Cmd: []string{ "bash", @@ -78,9 +80,12 @@ func startBoxWithEnv(ctx context.Context, t *testing.T, hostname, network, publi `sed -i \ -e "s|\"api-server\": \"http://127.0.0.1:3004\"|\"api-server\": \"${E2E_PUBLIC_URL}\"|" \ -e "s|\"own-base-url\": \"http://127.0.0.1:3060/task-runner\"|\"own-base-url\": \"${E2E_PUBLIC_URL}/task-runner\"|" \ + -e "s|\"orchAddr\": \"127.0.0.1:8936\"|\"orchAddr\": \"${E2E_ORCHESTRATOR_URL}\"|g" \ + -e "s|\"serviceAddr\": \"127.0.0.1:8936\"|\"httpAddr\": \"http://0.0.0.0:8936\", \"serviceAddr\": \"${E2E_ORCHESTRATOR_URL}\"|" \ /etc/livepeer/full-stack.json grep -Fq "\"api-server\": \"${E2E_PUBLIC_URL}\"" /etc/livepeer/full-stack.json grep -Fq "\"own-base-url\": \"${E2E_PUBLIC_URL}/task-runner\"" /etc/livepeer/full-stack.json +grep -Fq "\"serviceAddr\": \"${E2E_ORCHESTRATOR_URL}\"" /etc/livepeer/full-stack.json exec /usr/local/bin/catalyst -- /usr/local/bin/MistController -c /etc/livepeer/full-stack.json`, }, } diff --git a/test/e2e/e2e_test.go b/test/e2e/e2e_test.go index 3133ec2f..579c3593 100644 --- a/test/e2e/e2e_test.go +++ b/test/e2e/e2e_test.go @@ -157,10 +157,14 @@ func dumpContainerLogs(ctx context.Context, t *testing.T, container testcontaine } func startCatalyst(ctx context.Context, t *testing.T, hostname, network string, mc mistConfig) *catalystContainer { - return startCatalystWithEnv(ctx, t, hostname, network, mc, nil) + return startCatalystWithEnv(ctx, t, hostname, network, mc, nil, nil) } -func startCatalystWithEnv(ctx context.Context, t *testing.T, hostname, network string, mc mistConfig, env map[string]string) *catalystContainer { +func startCatalystWithPublicOrchestrator(ctx context.Context, t *testing.T, hostname, network string, mc mistConfig, hostPort string) *catalystContainer { + return startCatalystWithEnv(ctx, t, hostname, network, mc, nil, []string{fmt.Sprintf("%s:8936/tcp", hostPort)}) +} + +func startCatalystWithEnv(ctx context.Context, t *testing.T, hostname, network string, mc mistConfig, env map[string]string, extraPorts []string) *catalystContainer { mcPath, err := mc.toTmpFile(t.TempDir()) require.NoError(t, err) configAbsPath := filepath.Dir(mcPath) @@ -170,19 +174,21 @@ func startCatalystWithEnv(ctx context.Context, t *testing.T, hostname, network s for k, v := range env { envVars[k] = v } + exposedPorts := []string{ + tcp(webConsolePort), + tcp(httpPort), + tcp(catalystAPIPort), + tcp(catalystAPIInternalPort), + tcp(rtmpPort), + } + exposedPorts = append(exposedPorts, extraPorts...) req := testcontainers.ContainerRequest{ - Image: params.ImageName, - ExposedPorts: []string{ - tcp(webConsolePort), - tcp(httpPort), - tcp(catalystAPIPort), - tcp(catalystAPIInternalPort), - tcp(rtmpPort), - }, - Hostname: hostname, - Name: hostname, - Networks: []string{network}, - Env: envVars, + Image: params.ImageName, + ExposedPorts: exposedPorts, + Hostname: hostname, + Name: hostname, + Networks: []string{network}, + Env: envVars, Mounts: []testcontainers.ContainerMount{{ Source: testcontainers.GenericBindMountSource{ HostPath: configAbsPath, diff --git a/test/e2e/mist_config.go b/test/e2e/mist_config.go index de3a4d5d..5a2af594 100644 --- a/test/e2e/mist_config.go +++ b/test/e2e/mist_config.go @@ -53,6 +53,22 @@ func (m *mistConfig) setAPIServer(apiServer string) { } } +func (m *mistConfig) setPublicOrchestrator(orchestratorURL string) { + for i := range m.Config.Protocols { + p := &m.Config.Protocols[i] + if p.Connector != "livepeer" { + continue + } + if p.Broadcaster { + p.OrchAddr = orchestratorURL + } + if p.Orchestrator { + p.HTTPRPCAddr = "http://0.0.0.0:8936" + p.ServiceAddr = orchestratorURL + } + } +} + type config struct { Accesslog string `json:"accesslog"` Controller struct { diff --git a/test/e2e/quick_tunnel_test.go b/test/e2e/quick_tunnel_test.go index 932a156c..6a79913f 100644 --- a/test/e2e/quick_tunnel_test.go +++ b/test/e2e/quick_tunnel_test.go @@ -6,6 +6,7 @@ import ( "encoding/json" "fmt" "io" + "net" "net/http" "net/http/httptest" "net/url" @@ -103,6 +104,22 @@ func startQuickTunnel(t *testing.T, origin string) string { } } +// startOrchestratorTunnel reserves a host port for a containerized +// orchestrator, then exposes it through a temporary public URL. The port is +// released before Docker binds it, so cloudflared can start before the +// container is configured with the URL it must advertise. +func startOrchestratorTunnel(t *testing.T) (publicURL, hostPort string) { + t.Helper() + + listener, err := net.Listen("tcp", "127.0.0.1:0") + require.NoError(t, err) + _, hostPort, err = net.SplitHostPort(listener.Addr().String()) + require.NoError(t, err) + require.NoError(t, listener.Close()) + + return startQuickTunnel(t, "http://127.0.0.1:"+hostPort), hostPort +} + type callbackStatus struct { RequestID string `json:"request_id"` Status string `json:"status"` diff --git a/test/e2e/vod_test.go b/test/e2e/vod_test.go index 7b54c460..e2b5f4e0 100644 --- a/test/e2e/vod_test.go +++ b/test/e2e/vod_test.go @@ -48,11 +48,13 @@ func TestVod(t *testing.T) { storageURL := startQuickTunnel(t, fmt.Sprintf("http://127.0.0.1:%s", m.port)) waitForTunneledMinio(ctx, t, storageURL) callbacks := startCallbackTunnel(t) + orchestratorURL, orchestratorHostPort := startOrchestratorTunnel(t) h := randomString("catalyst-") mistConfig := defaultMistConfigWithLivepeerProcess(h, tunneledObjectStoreURL(t, storageURL, username, password, inBucket, "")) mistConfig.setAPIServer(callbacks.apiServerURL) - c := startCatalyst(ctx, t, h, network.name, mistConfig) + mistConfig.setPublicOrchestrator(orchestratorURL) + c := startCatalystWithPublicOrchestrator(ctx, t, h, network.name, mistConfig, orchestratorHostPort) defer c.Terminate(ctx) waitForCatalystAPI(t, c) From e39ec4b3ea9aaa6ce83e41521e9d8c2714c1410c Mon Sep 17 00:00:00 2001 From: Josh Allmann Date: Fri, 4 Sep 2026 20:23:23 -0700 Subject: [PATCH 08/13] Preserve gRPC through orchestrator tunnel --- test/e2e/box_record_test.go | 2 +- test/e2e/mist_config.go | 2 +- test/e2e/quick_tunnel_test.go | 24 ++++++++++++------------ 3 files changed, 14 insertions(+), 14 deletions(-) diff --git a/test/e2e/box_record_test.go b/test/e2e/box_record_test.go index c7346087..0cbdef6e 100644 --- a/test/e2e/box_record_test.go +++ b/test/e2e/box_record_test.go @@ -81,7 +81,7 @@ func startBoxWithEnv(ctx context.Context, t *testing.T, hostname, network, publi -e "s|\"api-server\": \"http://127.0.0.1:3004\"|\"api-server\": \"${E2E_PUBLIC_URL}\"|" \ -e "s|\"own-base-url\": \"http://127.0.0.1:3060/task-runner\"|\"own-base-url\": \"${E2E_PUBLIC_URL}/task-runner\"|" \ -e "s|\"orchAddr\": \"127.0.0.1:8936\"|\"orchAddr\": \"${E2E_ORCHESTRATOR_URL}\"|g" \ - -e "s|\"serviceAddr\": \"127.0.0.1:8936\"|\"httpAddr\": \"http://0.0.0.0:8936\", \"serviceAddr\": \"${E2E_ORCHESTRATOR_URL}\"|" \ + -e "s|\"serviceAddr\": \"127.0.0.1:8936\"|\"httpAddr\": \"https://0.0.0.0:8936\", \"serviceAddr\": \"${E2E_ORCHESTRATOR_URL}\"|" \ /etc/livepeer/full-stack.json grep -Fq "\"api-server\": \"${E2E_PUBLIC_URL}\"" /etc/livepeer/full-stack.json grep -Fq "\"own-base-url\": \"${E2E_PUBLIC_URL}/task-runner\"" /etc/livepeer/full-stack.json diff --git a/test/e2e/mist_config.go b/test/e2e/mist_config.go index 5a2af594..1a5c0786 100644 --- a/test/e2e/mist_config.go +++ b/test/e2e/mist_config.go @@ -63,7 +63,7 @@ func (m *mistConfig) setPublicOrchestrator(orchestratorURL string) { p.OrchAddr = orchestratorURL } if p.Orchestrator { - p.HTTPRPCAddr = "http://0.0.0.0:8936" + p.HTTPRPCAddr = "https://0.0.0.0:8936" p.ServiceAddr = orchestratorURL } } diff --git a/test/e2e/quick_tunnel_test.go b/test/e2e/quick_tunnel_test.go index 6a79913f..e195b315 100644 --- a/test/e2e/quick_tunnel_test.go +++ b/test/e2e/quick_tunnel_test.go @@ -42,11 +42,15 @@ func (b *synchronizedBuffer) String() string { // startQuickTunnel exposes origin through a temporary trycloudflare.com URL. // The cloudflared process is stopped automatically when the test completes. func startQuickTunnel(t *testing.T, origin string) string { + return startQuickTunnelWithArgs(t, origin) +} + +func startQuickTunnelWithArgs(t *testing.T, origin string, args ...string) string { t.Helper() originURL, err := url.Parse(origin) require.NoError(t, err) - require.Equal(t, "http", originURL.Scheme) + require.Contains(t, []string{"http", "https"}, originURL.Scheme) require.NotEmpty(t, originURL.Host) cloudflared, err := exec.LookPath("cloudflared") @@ -54,14 +58,9 @@ func startQuickTunnel(t *testing.T, origin string) string { ctx, cancel := context.WithCancel(context.Background()) output := &synchronizedBuffer{} - cmd := exec.CommandContext( - ctx, - cloudflared, - "tunnel", - "--no-autoupdate", - "--protocol", "http2", - "--url", origin, - ) + commandArgs := []string{"tunnel", "--no-autoupdate", "--protocol", "http2", "--url", origin} + commandArgs = append(commandArgs, args...) + cmd := exec.CommandContext(ctx, cloudflared, commandArgs...) cmd.Stdout = output cmd.Stderr = output require.NoError(t, cmd.Start()) @@ -105,8 +104,9 @@ func startQuickTunnel(t *testing.T, origin string) string { } // startOrchestratorTunnel reserves a host port for a containerized -// orchestrator, then exposes it through a temporary public URL. The port is -// released before Docker binds it, so cloudflared can start before the +// orchestrator, then exposes it through a temporary public URL. The origin +// remains HTTPS/HTTP2 because Go-livepeer's control plane uses gRPC. The port +// is released before Docker binds it, so cloudflared can start before the // container is configured with the URL it must advertise. func startOrchestratorTunnel(t *testing.T) (publicURL, hostPort string) { t.Helper() @@ -117,7 +117,7 @@ func startOrchestratorTunnel(t *testing.T) (publicURL, hostPort string) { require.NoError(t, err) require.NoError(t, listener.Close()) - return startQuickTunnel(t, "http://127.0.0.1:"+hostPort), hostPort + return startQuickTunnelWithArgs(t, "https://127.0.0.1:"+hostPort, "--no-tls-verify", "--http2-origin"), hostPort } type callbackStatus struct { From f9851c4964fe552e272d61a81e419b71366d2873 Mon Sep 17 00:00:00 2001 From: Josh Allmann Date: Fri, 4 Sep 2026 20:57:02 -0700 Subject: [PATCH 09/13] Disable E2E orchestrator self-check --- test/e2e/mist_config.go | 41 ++++++++++++++++++++++------------------- 1 file changed, 22 insertions(+), 19 deletions(-) diff --git a/test/e2e/mist_config.go b/test/e2e/mist_config.go index 1a5c0786..94e11ae1 100644 --- a/test/e2e/mist_config.go +++ b/test/e2e/mist_config.go @@ -24,25 +24,26 @@ type bandwidth struct { } type protocol struct { - Connector string `json:"connector"` - APIServer string `json:"api-server,omitempty"` - RetryJoin string `json:"retry-join,omitempty"` - Advertise string `json:"advertise,omitempty"` - RPCAddr string `json:"rpc-addr,omitempty"` - RedirectPrefixes string `json:"redirect-prefixes,omitempty"` - Debug string `json:"debug,omitempty"` - HTTPAddr string `json:"http-addr,omitempty"` - HTTPAddrInternal string `json:"http-internal-addr,omitempty"` - Broadcaster bool `json:"broadcaster,omitempty"` - Orchestrator bool `json:"orchestrator,omitempty"` - Transcoder bool `json:"transcoder,omitempty"` - HTTPRPCAddr string `json:"httpAddr,omitempty"` - OrchAddr string `json:"orchAddr,omitempty"` - ServiceAddr string `json:"serviceAddr,omitempty"` - CliAddr string `json:"cliAddr,omitempty"` - RtmpAddr string `json:"rtmpAddr,omitempty"` - SourceOutput string `json:"source-output,omitempty"` - Catabalancer string `json:"catabalancer,omitempty"` + Connector string `json:"connector"` + APIServer string `json:"api-server,omitempty"` + RetryJoin string `json:"retry-join,omitempty"` + Advertise string `json:"advertise,omitempty"` + RPCAddr string `json:"rpc-addr,omitempty"` + RedirectPrefixes string `json:"redirect-prefixes,omitempty"` + Debug string `json:"debug,omitempty"` + HTTPAddr string `json:"http-addr,omitempty"` + HTTPAddrInternal string `json:"http-internal-addr,omitempty"` + Broadcaster bool `json:"broadcaster,omitempty"` + Orchestrator bool `json:"orchestrator,omitempty"` + Transcoder bool `json:"transcoder,omitempty"` + HTTPRPCAddr string `json:"httpAddr,omitempty"` + OrchAddr string `json:"orchAddr,omitempty"` + ServiceAddr string `json:"serviceAddr,omitempty"` + StartupAvailabilityCheck *bool `json:"startupAvailabilityCheck,omitempty"` + CliAddr string `json:"cliAddr,omitempty"` + RtmpAddr string `json:"rtmpAddr,omitempty"` + SourceOutput string `json:"source-output,omitempty"` + Catabalancer string `json:"catabalancer,omitempty"` } func (m *mistConfig) setAPIServer(apiServer string) { @@ -63,8 +64,10 @@ func (m *mistConfig) setPublicOrchestrator(orchestratorURL string) { p.OrchAddr = orchestratorURL } if p.Orchestrator { + startupAvailabilityCheck := false p.HTTPRPCAddr = "https://0.0.0.0:8936" p.ServiceAddr = orchestratorURL + p.StartupAvailabilityCheck = &startupAvailabilityCheck } } } From 16660d5b3145b35640226fdf3b4f603e66ffae5c Mon Sep 17 00:00:00 2001 From: Josh Allmann Date: Fri, 4 Sep 2026 21:14:26 -0700 Subject: [PATCH 10/13] Pass E2E startup probe flag value --- test/e2e/mist_config.go | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/test/e2e/mist_config.go b/test/e2e/mist_config.go index 94e11ae1..fa654ffc 100644 --- a/test/e2e/mist_config.go +++ b/test/e2e/mist_config.go @@ -39,7 +39,7 @@ type protocol struct { HTTPRPCAddr string `json:"httpAddr,omitempty"` OrchAddr string `json:"orchAddr,omitempty"` ServiceAddr string `json:"serviceAddr,omitempty"` - StartupAvailabilityCheck *bool `json:"startupAvailabilityCheck,omitempty"` + StartupAvailabilityCheck string `json:"startupAvailabilityCheck,omitempty"` CliAddr string `json:"cliAddr,omitempty"` RtmpAddr string `json:"rtmpAddr,omitempty"` SourceOutput string `json:"source-output,omitempty"` @@ -64,10 +64,9 @@ func (m *mistConfig) setPublicOrchestrator(orchestratorURL string) { p.OrchAddr = orchestratorURL } if p.Orchestrator { - startupAvailabilityCheck := false p.HTTPRPCAddr = "https://0.0.0.0:8936" p.ServiceAddr = orchestratorURL - p.StartupAvailabilityCheck = &startupAvailabilityCheck + p.StartupAvailabilityCheck = "false" } } } From 0d924d0aa2bed0f2196957ea379148368ce02315 Mon Sep 17 00:00:00 2001 From: Josh Allmann Date: Fri, 4 Sep 2026 21:38:33 -0700 Subject: [PATCH 11/13] Use non-loopback E2E orchestrator address --- test/e2e/box_record_test.go | 9 ++++---- test/e2e/e2e_test.go | 4 ---- test/e2e/mist_config.go | 43 +++++++++++++++++------------------ test/e2e/quick_tunnel_test.go | 35 ++++++++-------------------- test/e2e/vod_test.go | 5 ++-- 5 files changed, 36 insertions(+), 60 deletions(-) diff --git a/test/e2e/box_record_test.go b/test/e2e/box_record_test.go index 0cbdef6e..c4631c02 100644 --- a/test/e2e/box_record_test.go +++ b/test/e2e/box_record_test.go @@ -36,10 +36,9 @@ func TestBoxRecording(t *testing.T) { boxName := randomString("box-") publicURL := startQuickTunnel(t, "http://127.0.0.1:8888") - orchestratorURL, orchestratorHostPort := startOrchestratorTunnel(t) // when - box := startBoxWithEnv(ctx, t, boxName, network.name, publicURL, orchestratorURL, orchestratorHostPort) + box := startBoxWithEnv(ctx, t, boxName, network.name, publicURL) defer box.Terminate(ctx) waitForBoxMinio(t, publicURL) configureBoxObjectStores(t, publicURL) @@ -60,19 +59,19 @@ func TestBoxRecording(t *testing.T) { } } -func startBoxWithEnv(ctx context.Context, t *testing.T, hostname, network, publicURL, orchestratorURL, orchestratorHostPort string) *catalystContainer { +func startBoxWithEnv(ctx context.Context, t *testing.T, hostname, network, publicURL string) *catalystContainer { req := testcontainers.ContainerRequest{ Image: "livepeer/in-a-box", Hostname: hostname, Name: hostname, Networks: []string{network}, - ExposedPorts: []string{"1935:1935/tcp", "8888:8888/tcp", fmt.Sprintf("%s:8936/tcp", orchestratorHostPort)}, + ExposedPorts: []string{"1935:1935/tcp", "8888:8888/tcp"}, ShmSize: 1000000000, WaitingFor: wait.NewLogStrategy("API server listening").WithStartupTimeout(3 * time.Minute), Env: map[string]string{ "LP_API_FRONTEND": "false", "E2E_PUBLIC_URL": publicURL, - "E2E_ORCHESTRATOR_URL": orchestratorURL, + "E2E_ORCHESTRATOR_URL": fmt.Sprintf("https://%s:8936", hostname), }, Cmd: []string{ "bash", diff --git a/test/e2e/e2e_test.go b/test/e2e/e2e_test.go index 579c3593..986716ce 100644 --- a/test/e2e/e2e_test.go +++ b/test/e2e/e2e_test.go @@ -160,10 +160,6 @@ func startCatalyst(ctx context.Context, t *testing.T, hostname, network string, return startCatalystWithEnv(ctx, t, hostname, network, mc, nil, nil) } -func startCatalystWithPublicOrchestrator(ctx context.Context, t *testing.T, hostname, network string, mc mistConfig, hostPort string) *catalystContainer { - return startCatalystWithEnv(ctx, t, hostname, network, mc, nil, []string{fmt.Sprintf("%s:8936/tcp", hostPort)}) -} - func startCatalystWithEnv(ctx context.Context, t *testing.T, hostname, network string, mc mistConfig, env map[string]string, extraPorts []string) *catalystContainer { mcPath, err := mc.toTmpFile(t.TempDir()) require.NoError(t, err) diff --git a/test/e2e/mist_config.go b/test/e2e/mist_config.go index fa654ffc..0166803b 100644 --- a/test/e2e/mist_config.go +++ b/test/e2e/mist_config.go @@ -24,26 +24,25 @@ type bandwidth struct { } type protocol struct { - Connector string `json:"connector"` - APIServer string `json:"api-server,omitempty"` - RetryJoin string `json:"retry-join,omitempty"` - Advertise string `json:"advertise,omitempty"` - RPCAddr string `json:"rpc-addr,omitempty"` - RedirectPrefixes string `json:"redirect-prefixes,omitempty"` - Debug string `json:"debug,omitempty"` - HTTPAddr string `json:"http-addr,omitempty"` - HTTPAddrInternal string `json:"http-internal-addr,omitempty"` - Broadcaster bool `json:"broadcaster,omitempty"` - Orchestrator bool `json:"orchestrator,omitempty"` - Transcoder bool `json:"transcoder,omitempty"` - HTTPRPCAddr string `json:"httpAddr,omitempty"` - OrchAddr string `json:"orchAddr,omitempty"` - ServiceAddr string `json:"serviceAddr,omitempty"` - StartupAvailabilityCheck string `json:"startupAvailabilityCheck,omitempty"` - CliAddr string `json:"cliAddr,omitempty"` - RtmpAddr string `json:"rtmpAddr,omitempty"` - SourceOutput string `json:"source-output,omitempty"` - Catabalancer string `json:"catabalancer,omitempty"` + Connector string `json:"connector"` + APIServer string `json:"api-server,omitempty"` + RetryJoin string `json:"retry-join,omitempty"` + Advertise string `json:"advertise,omitempty"` + RPCAddr string `json:"rpc-addr,omitempty"` + RedirectPrefixes string `json:"redirect-prefixes,omitempty"` + Debug string `json:"debug,omitempty"` + HTTPAddr string `json:"http-addr,omitempty"` + HTTPAddrInternal string `json:"http-internal-addr,omitempty"` + Broadcaster bool `json:"broadcaster,omitempty"` + Orchestrator bool `json:"orchestrator,omitempty"` + Transcoder bool `json:"transcoder,omitempty"` + HTTPRPCAddr string `json:"httpAddr,omitempty"` + OrchAddr string `json:"orchAddr,omitempty"` + ServiceAddr string `json:"serviceAddr,omitempty"` + CliAddr string `json:"cliAddr,omitempty"` + RtmpAddr string `json:"rtmpAddr,omitempty"` + SourceOutput string `json:"source-output,omitempty"` + Catabalancer string `json:"catabalancer,omitempty"` } func (m *mistConfig) setAPIServer(apiServer string) { @@ -54,7 +53,8 @@ func (m *mistConfig) setAPIServer(apiServer string) { } } -func (m *mistConfig) setPublicOrchestrator(orchestratorURL string) { +func (m *mistConfig) setNonLoopbackOrchestrator(hostname string) { + orchestratorURL := "https://" + hostname + ":8936" for i := range m.Config.Protocols { p := &m.Config.Protocols[i] if p.Connector != "livepeer" { @@ -66,7 +66,6 @@ func (m *mistConfig) setPublicOrchestrator(orchestratorURL string) { if p.Orchestrator { p.HTTPRPCAddr = "https://0.0.0.0:8936" p.ServiceAddr = orchestratorURL - p.StartupAvailabilityCheck = "false" } } } diff --git a/test/e2e/quick_tunnel_test.go b/test/e2e/quick_tunnel_test.go index e195b315..932a156c 100644 --- a/test/e2e/quick_tunnel_test.go +++ b/test/e2e/quick_tunnel_test.go @@ -6,7 +6,6 @@ import ( "encoding/json" "fmt" "io" - "net" "net/http" "net/http/httptest" "net/url" @@ -42,15 +41,11 @@ func (b *synchronizedBuffer) String() string { // startQuickTunnel exposes origin through a temporary trycloudflare.com URL. // The cloudflared process is stopped automatically when the test completes. func startQuickTunnel(t *testing.T, origin string) string { - return startQuickTunnelWithArgs(t, origin) -} - -func startQuickTunnelWithArgs(t *testing.T, origin string, args ...string) string { t.Helper() originURL, err := url.Parse(origin) require.NoError(t, err) - require.Contains(t, []string{"http", "https"}, originURL.Scheme) + require.Equal(t, "http", originURL.Scheme) require.NotEmpty(t, originURL.Host) cloudflared, err := exec.LookPath("cloudflared") @@ -58,9 +53,14 @@ func startQuickTunnelWithArgs(t *testing.T, origin string, args ...string) strin ctx, cancel := context.WithCancel(context.Background()) output := &synchronizedBuffer{} - commandArgs := []string{"tunnel", "--no-autoupdate", "--protocol", "http2", "--url", origin} - commandArgs = append(commandArgs, args...) - cmd := exec.CommandContext(ctx, cloudflared, commandArgs...) + cmd := exec.CommandContext( + ctx, + cloudflared, + "tunnel", + "--no-autoupdate", + "--protocol", "http2", + "--url", origin, + ) cmd.Stdout = output cmd.Stderr = output require.NoError(t, cmd.Start()) @@ -103,23 +103,6 @@ func startQuickTunnelWithArgs(t *testing.T, origin string, args ...string) strin } } -// startOrchestratorTunnel reserves a host port for a containerized -// orchestrator, then exposes it through a temporary public URL. The origin -// remains HTTPS/HTTP2 because Go-livepeer's control plane uses gRPC. The port -// is released before Docker binds it, so cloudflared can start before the -// container is configured with the URL it must advertise. -func startOrchestratorTunnel(t *testing.T) (publicURL, hostPort string) { - t.Helper() - - listener, err := net.Listen("tcp", "127.0.0.1:0") - require.NoError(t, err) - _, hostPort, err = net.SplitHostPort(listener.Addr().String()) - require.NoError(t, err) - require.NoError(t, listener.Close()) - - return startQuickTunnelWithArgs(t, "https://127.0.0.1:"+hostPort, "--no-tls-verify", "--http2-origin"), hostPort -} - type callbackStatus struct { RequestID string `json:"request_id"` Status string `json:"status"` diff --git a/test/e2e/vod_test.go b/test/e2e/vod_test.go index e2b5f4e0..1ed724ef 100644 --- a/test/e2e/vod_test.go +++ b/test/e2e/vod_test.go @@ -48,13 +48,12 @@ func TestVod(t *testing.T) { storageURL := startQuickTunnel(t, fmt.Sprintf("http://127.0.0.1:%s", m.port)) waitForTunneledMinio(ctx, t, storageURL) callbacks := startCallbackTunnel(t) - orchestratorURL, orchestratorHostPort := startOrchestratorTunnel(t) h := randomString("catalyst-") mistConfig := defaultMistConfigWithLivepeerProcess(h, tunneledObjectStoreURL(t, storageURL, username, password, inBucket, "")) mistConfig.setAPIServer(callbacks.apiServerURL) - mistConfig.setPublicOrchestrator(orchestratorURL) - c := startCatalystWithPublicOrchestrator(ctx, t, h, network.name, mistConfig, orchestratorHostPort) + mistConfig.setNonLoopbackOrchestrator(h) + c := startCatalyst(ctx, t, h, network.name, mistConfig) defer c.Terminate(ctx) waitForCatalystAPI(t, c) From 499afe1122b1956fd012819434a23159b1b5f6d2 Mon Sep 17 00:00:00 2001 From: Josh Allmann Date: Fri, 4 Sep 2026 21:57:26 -0700 Subject: [PATCH 12/13] Wait for VOD output uploads --- test/e2e/vod_test.go | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/test/e2e/vod_test.go b/test/e2e/vod_test.go index 1ed724ef..7fb08137 100644 --- a/test/e2e/vod_test.go +++ b/test/e2e/vod_test.go @@ -228,7 +228,10 @@ func processVod(t *testing.T, storageURL, callbackURL string, c *catalystContain func requireOutputFiles(ctx context.Context, t *testing.T, m *minioContainer) { cli := minioClient(t, m) var files []string - timeoutAt := time.Now().Add(30 * time.Second) + // The VOD completion callback is sent after the manifests have been written, + // but segment uploads can still be in flight through the object-store tunnel. + // Keep the test container alive until those uploads are observable in MinIO. + timeoutAt := time.Now().Add(2 * time.Minute) expectedFiles := []string{ "index.m3u8", @@ -248,6 +251,7 @@ func requireOutputFiles(ctx context.Context, t *testing.T, m *minioContainer) { for timeoutAt.After(time.Now()) { files = []string{} for o := range cli.ListObjects(ctx, outBucket, minio.ListObjectsOptions{Recursive: true}) { + require.NoError(t, o.Err) files = append(files, o.Key) } if len(files) < len(expectedFiles) { From e7c75cc017c4baadbc4bbf375645a93d4453b1fb Mon Sep 17 00:00:00 2001 From: Josh Allmann Date: Fri, 4 Sep 2026 22:15:12 -0700 Subject: [PATCH 13/13] Wait for VOD startup readiness --- test/e2e/vod_test.go | 45 +++++++++++++++++++++++++++++--------------- 1 file changed, 30 insertions(+), 15 deletions(-) diff --git a/test/e2e/vod_test.go b/test/e2e/vod_test.go index 7fb08137..d87f705d 100644 --- a/test/e2e/vod_test.go +++ b/test/e2e/vod_test.go @@ -207,22 +207,37 @@ func processVod(t *testing.T, storageURL, callbackURL string, c *catalystContain }`, sourceVideoURL, callbackURL, destURL) url := fmt.Sprintf("http://127.0.0.1:%s/api/vod", c.catalystAPIInternal) - req, err := http.NewRequest("POST", url, bytes.NewBuffer([]byte(jsonData))) - require.NoError(t, err) - req.Header.Set("Content-Type", "application/json") - req.Header.Set("Authorization", "Bearer IAmAuthorized") - client := &http.Client{} - resp, err := client.Do(req) - require.NoError(t, err) - defer resp.Body.Close() - body, _ := io.ReadAll(io.LimitReader(resp.Body, 4096)) - require.Equal(t, http.StatusOK, resp.StatusCode, "unexpected response: %s", strings.TrimSpace(string(body))) - var result struct { - RequestID string `json:"request_id"` + deadline := time.Now().Add(time.Minute) + for { + req, err := http.NewRequest("POST", url, bytes.NewBuffer([]byte(jsonData))) + require.NoError(t, err) + req.Header.Set("Content-Type", "application/json") + req.Header.Set("Authorization", "Bearer IAmAuthorized") + + resp, err := (&http.Client{}).Do(req) + require.NoError(t, err) + body, err := io.ReadAll(io.LimitReader(resp.Body, 4096)) + resp.Body.Close() + require.NoError(t, err) + if resp.StatusCode == http.StatusOK { + var result struct { + RequestID string `json:"request_id"` + } + require.NoError(t, json.Unmarshal(body, &result)) + require.NotEmpty(t, result.RequestID) + return result.RequestID + } + + // The API becomes healthy before the embedded orchestrator is ready to + // accept a VOD job. Its initial 500 response is transient; retrying here + // avoids treating that startup race as an E2E failure. + if resp.StatusCode != http.StatusInternalServerError || time.Now().After(deadline) { + dumpContainerLogs(context.Background(), t, c.Container) + require.Equal(t, http.StatusOK, resp.StatusCode, "unexpected response: %s", strings.TrimSpace(string(body))) + return "" + } + time.Sleep(time.Second) } - require.NoError(t, json.Unmarshal(body, &result)) - require.NotEmpty(t, result.RequestID) - return result.RequestID } func requireOutputFiles(ctx context.Context, t *testing.T, m *minioContainer) {