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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
53 changes: 53 additions & 0 deletions .github/workflows/reverse-proxy-e2e.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,53 @@
name: Reverse Proxy E2E

on:
pull_request:
paths:
- '.github/workflows/reverse-proxy-e2e.yml'
- '**/*.go'
- 'e2e/**'
- 'shared/management/**'
- 'combined/Dockerfile*'
- 'proxy/Dockerfile*'
- '.dockerignore'
- 'go.mod'
- 'go.sum'
workflow_dispatch:

permissions:
contents: read

concurrency:
group: ${{ github.workflow }}-${{ github.ref }}
cancel-in-progress: true

jobs:
reverse-proxy:
runs-on: ubuntu-latest
timeout-minutes: 30
steps:
- name: Checkout code
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
persist-credentials: false

- name: Install Go
uses: actions/setup-go@4a3601121dd01d1626a1e23e37211e3254c1c06c # v6.4.0
with:
go-version-file: go.mod

- name: Set up Buildx
uses: docker/setup-buildx-action@d7f5e7f509e45cec5c76c4d5afdd7de93d0b3df5 # v4.1.0

- name: Cache Docker layers
uses: actions/cache@27d5ce7f107fe9357f9df03efb73ab90386fccae # v5.0.5
with:
path: /tmp/.reverse-proxy-buildx-cache
key: ${{ runner.os }}-reverse-proxy-e2e-${{ hashFiles('go.sum', 'combined/Dockerfile.multistage', 'proxy/Dockerfile.multistage') }}
restore-keys: |
${{ runner.os }}-reverse-proxy-e2e-

- name: Run reverse-proxy E2E
env:
NB_E2E_BUILDX_CACHE: /tmp/.reverse-proxy-buildx-cache
run: go test -tags e2e -timeout 25m -v ./e2e/reverseproxy/...
4 changes: 2 additions & 2 deletions e2e/harness/cert.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,8 +18,8 @@ import (

// writeSelfSignedCert generates a self-signed TLS cert/key pair covering the
// given DNS names and writes them as tls.crt / tls.key in dir. The proxy serves
// this for the agent-network endpoint; the client curls with -k, so validity
// chains don't matter — the proxy just needs a usable cert to present.
// this for the agent-network endpoint; test clients either trust the generated
// certificate explicitly or opt out of verification inside the test network.
func writeSelfSignedCert(dir string, dnsNames []string) error {
priv, err := ecdsa.GenerateKey(elliptic.P256(), rand.Reader)
if err != nil {
Expand Down
228 changes: 228 additions & 0 deletions e2e/harness/http_upstream.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,228 @@
//go:build e2e

package harness

import (
"bufio"
"context"
"encoding/json"
"fmt"
"io"
"net"
"net/http"
"net/url"
"os"
"path/filepath"
"strings"
"sync/atomic"
"time"

"github.com/docker/docker/api/types/container"
"github.com/docker/go-connections/nat"
"github.com/testcontainers/testcontainers-go"
"github.com/testcontainers/testcontainers-go/wait"
)

const (
httpUpstreamImage = "nginx:alpine"
httpUpstreamAlias = "httpupstream"
httpUpstreamPort = "8080/tcp"
httpUpstreamLogPath = "/tmp/e2e-access.log"
)

// httpUpstreamNginxConf returns request details from a plain HTTP upstream and
// writes every request URI to an unbuffered file. Tests use unique path markers
// to distinguish requests that reached the upstream from proxy-generated
// responses with the same status.
const httpUpstreamNginxConf = `pid /tmp/nginx.pid;
worker_processes 1;
events {}
http {
log_format e2e escape=json '{"uri":"$request_uri"}';
access_log /tmp/e2e-access.log e2e;

server {
listen 8080;
location / {
default_type application/json;
return 200 '{"source":"netbird-e2e-http-upstream","path":"$uri","request_uri":"$request_uri","configured_auth":"$http_x_e2e_proxy_key","authorization":"$http_authorization","cookie":"$http_cookie","netbird_user":"$http_x_netbird_user","netbird_groups":"$http_x_netbird_groups"}';
}
}
}
`

// HTTPUpstream is a plain HTTP server on the combined server's network. It
// reflects selected request properties and records request URIs for tests that
// need to distinguish proxy responses from upstream responses.
type HTTPUpstream struct {
container testcontainers.Container
workDir string
directAddress string
syncSequence atomic.Uint64
URL string
}

// StartHTTPUpstream runs a request-observing nginx server on the shared network.
func StartHTTPUpstream(ctx context.Context, c *Combined) (*HTTPUpstream, error) {
workDir, err := os.MkdirTemp("/tmp", "nb-e2e-http-upstream-*")
if err != nil {
return nil, fmt.Errorf("create HTTP upstream work dir: %w", err)
}
if err := os.Chmod(workDir, 0o755); err != nil { //nolint:gosec // throwaway e2e config dir, must be traversable by the container uid
_ = os.RemoveAll(workDir)
return nil, fmt.Errorf("chmod HTTP upstream dir: %w", err)
}
if err := os.WriteFile(filepath.Join(workDir, "nginx.conf"), []byte(httpUpstreamNginxConf), 0o644); err != nil { //nolint:gosec // non-secret e2e config
_ = os.RemoveAll(workDir)
return nil, fmt.Errorf("write HTTP upstream config: %w", err)
}

req := testcontainers.ContainerRequest{
Image: httpUpstreamImage,
ExposedPorts: []string{httpUpstreamPort},
Networks: []string{c.network.Name},
NetworkAliases: map[string][]string{c.network.Name: {httpUpstreamAlias}},
Cmd: []string{"nginx", "-c", "/conf/nginx.conf", "-g", "daemon off;"},
HostConfigModifier: func(hc *container.HostConfig) {
hc.Binds = append(hc.Binds, workDir+":/conf:ro")
},
WaitingFor: wait.ForListeningPort(httpUpstreamPort).WithStartupTimeout(60 * time.Second),
}

ctr, err := testcontainers.GenericContainer(ctx, testcontainers.GenericContainerRequest{
ContainerRequest: req,
Started: true,
})
if err != nil {
cleanupHTTPUpstreamStart(ctr, workDir)
return nil, fmt.Errorf("start HTTP upstream container: %w", err)
}
host, err := ctr.Host(ctx)
if err != nil {
cleanupHTTPUpstreamStart(ctr, workDir)
return nil, fmt.Errorf("HTTP upstream container host: %w", err)
}
mapped, err := ctr.MappedPort(ctx, nat.Port(httpUpstreamPort))
if err != nil {
cleanupHTTPUpstreamStart(ctr, workDir)
return nil, fmt.Errorf("HTTP upstream mapped port: %w", err)
}

return &HTTPUpstream{
container: ctr,
workDir: workDir,
directAddress: net.JoinHostPort(host, mapped.Port()),
URL: "http://" + httpUpstreamAlias + ":8080",
}, nil
}

func cleanupHTTPUpstreamStart(ctr testcontainers.Container, workDir string) {
if ctr != nil {
cleanupCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
_ = ctr.Terminate(cleanupCtx)
cancel()
}
_ = os.RemoveAll(workDir)
}

// Synchronize waits until the upstream has finalized all requests received
// before a direct barrier request. The fixture runs one nginx worker and writes
// access logs without buffering, so observing the barrier in the file proves
// any earlier proxied request has also been recorded.
func (u *HTTPUpstream) Synchronize(ctx context.Context) error {
marker := fmt.Sprintf("__e2e_upstream_barrier_%d__", u.syncSequence.Add(1))
requestURL := (&url.URL{
Scheme: "http",
Host: u.directAddress,
Path: "/" + marker,
}).String()
req, err := http.NewRequestWithContext(ctx, http.MethodGet, requestURL, nil)
if err != nil {
return fmt.Errorf("create HTTP upstream barrier request: %w", err)
}
transport := &http.Transport{}
defer transport.CloseIdleConnections()
client := &http.Client{Transport: transport, Timeout: 5 * time.Second}
resp, err := client.Do(req)
if err != nil {
return fmt.Errorf("send HTTP upstream barrier request: %w", err)
}
defer resp.Body.Close()
if _, err := io.Copy(io.Discard, resp.Body); err != nil {
return fmt.Errorf("read HTTP upstream barrier response: %w", err)
}
if resp.StatusCode != http.StatusOK {
return fmt.Errorf("HTTP upstream barrier returned status %d", resp.StatusCode)
}

waitCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
defer cancel()
ticker := time.NewTicker(25 * time.Millisecond)
defer ticker.Stop()
for {
count, err := u.RequestCount(waitCtx, marker)
if err != nil {
return err
}
switch {
case count == 1:
return nil
case count > 1:
return fmt.Errorf("HTTP upstream recorded barrier %q %d times", marker, count)
}

select {
case <-waitCtx.Done():
return fmt.Errorf("wait for HTTP upstream barrier: %w", waitCtx.Err())
case <-ticker.C:
}
}
}

// RequestCount returns the number of recorded request URIs containing marker.
func (u *HTTPUpstream) RequestCount(ctx context.Context, marker string) (int, error) {
if marker == "" {
return 0, fmt.Errorf("request marker is empty")
}

reader, err := u.container.CopyFileFromContainer(ctx, httpUpstreamLogPath)
if err != nil {
return 0, fmt.Errorf("read HTTP upstream access log: %w", err)
}
defer reader.Close()

count := 0
scanner := bufio.NewScanner(reader)
for scanner.Scan() {
var entry struct {
URI string `json:"uri"`
}
if err := json.Unmarshal(scanner.Bytes(), &entry); err != nil {
return 0, fmt.Errorf("decode HTTP upstream access log: %w", err)
}
if strings.Contains(entry.URI, marker) {
count++
}
}
if err := scanner.Err(); err != nil {
return 0, fmt.Errorf("scan HTTP upstream access log: %w", err)
}
return count, nil
}

// Logs returns the HTTP upstream container logs for failure diagnostics.
func (u *HTTPUpstream) Logs(ctx context.Context) string {
return containerLogs(ctx, u.container)
}

// Terminate stops the HTTP upstream container and removes its work directory.
func (u *HTTPUpstream) Terminate(ctx context.Context) error {
var err error
if u.container != nil {
err = u.container.Terminate(ctx)
}
if u.workDir != "" {
_ = os.RemoveAll(u.workDir)
}
return err
}
Loading
Loading