diff --git a/.github/workflows/build.yaml b/.github/workflows/build.yaml index ebadc1650..af70f284d 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 diff --git a/.github/workflows/test.yaml b/.github/workflows/test.yaml index e0ee98767..7544c4f57 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 @@ -63,9 +92,21 @@ 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 + go test $(go list ./... | grep 'test/e2e') --timeout 30m -v --image catalyst-e2e-test - name: Upload coverage reports diff --git a/Dockerfile b/Dockerfile index 5facef6ea..da9e8ee63 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 diff --git a/manifest.yaml b/manifest.yaml index 2f57982cf..121f609bb 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 diff --git a/scripts/livepeer-nginx b/scripts/livepeer-nginx index e141502e1..4b79a353a 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 5c2c34478..c4631c02d 100644 --- a/test/e2e/box_record_test.go +++ b/test/e2e/box_record_test.go @@ -3,17 +3,23 @@ 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" - "golang.org/x/sync/errgroup" ) func TestBoxRecording(t *testing.T) { @@ -29,22 +35,31 @@ 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 { - 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 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, @@ -54,7 +69,23 @@ func startBoxWithEnv(ctx context.Context, t *testing.T, hostname, network string ShmSize: 1000000000, WaitingFor: wait.NewLogStrategy("API server listening").WithStartupTimeout(3 * time.Minute), Env: map[string]string{ - "LP_API_FRONTEND": "false", + "LP_API_FRONTEND": "false", + "E2E_PUBLIC_URL": publicURL, + "E2E_ORCHESTRATOR_URL": fmt.Sprintf("https://%s:8936", hostname), + }, + 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\"|" \ + -e "s|\"orchAddr\": \"127.0.0.1:8936\"|\"orchAddr\": \"${E2E_ORCHESTRATOR_URL}\"|g" \ + -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 +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`, }, } container, err := testcontainers.GenericContainer(ctx, testcontainers.GenericContainerRequest{ @@ -97,6 +128,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/e2e_test.go b/test/e2e/e2e_test.go index 1e9111c7e..986716ce7 100644 --- a/test/e2e/e2e_test.go +++ b/test/e2e/e2e_test.go @@ -138,11 +138,29 @@ 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) + 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 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) @@ -152,19 +170,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 1f463ba8e..0166803b4 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,31 @@ 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 + } + } +} + +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" { + continue + } + if p.Broadcaster { + p.OrchAddr = orchestratorURL + } + if p.Orchestrator { + p.HTTPRPCAddr = "https://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 new file mode 100644 index 000000000..932a156c8 --- /dev/null +++ b/test/e2e/quick_tunnel_test.go @@ -0,0 +1,243 @@ +package e2e + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "io" + "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()) + } + } +} + +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) + + 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) { + 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", ""), + ) +} + +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 28dc166b1..d87f705dc 100644 --- a/test/e2e/vod_test.go +++ b/test/e2e/vod_test.go @@ -3,8 +3,12 @@ package e2e import ( "bytes" "context" + "encoding/json" "fmt" + "io" "net/http" + "net/url" + "strings" "testing" "time" @@ -41,14 +45,24 @@ 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) + callbacks := 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(callbacks.apiServerURL) + mistConfig.setNonLoopbackOrchestrator(h) + c := startCatalyst(ctx, t, h, network.name, mistConfig) defer c.Terminate(ctx) waitForCatalystAPI(t, c) // when - processVod(t, m, 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) @@ -115,14 +129,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 +189,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) string { + 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,23 +204,49 @@ 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))) - 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() + 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) + } } func requireOutputFiles(ctx context.Context, t *testing.T, m *minioContainer) { cli := minioClient(t, m) var files []string - timeoutAt := time.Now().Add(5 * time.Minute) + // 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", @@ -202,6 +266,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) {