Merge branch 'main' into agent-network-validate-proxy-cluster

This commit is contained in:
Maycon Santos
2026-09-12 10:45:11 +02:00
committed by GitHub
506 changed files with 31909 additions and 11042 deletions
+143
View File
@@ -0,0 +1,143 @@
//go:build e2e
package agentnetwork
import (
"context"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/netbirdio/netbird/shared/management/http/api"
)
// joinGroup places the PAT's own user into the group so caller-scoped answers
// (GET /api/agent-network/agent-config) see the policies sourced from it, and
// restores the previous auto-groups on cleanup. Self-service updates of one's
// own auto_groups are permitted for every role, so this needs no second user.
func joinGroup(t *testing.T, ctx context.Context, groupID string) {
t.Helper()
me, err := srv.API().Users.Current(ctx)
require.NoError(t, err, "read current user")
before := append([]string(nil), me.AutoGroups...)
_, err = srv.API().Users.Update(ctx, me.Id, api.PutApiUsersUserIdJSONRequestBody{
Role: me.Role,
IsBlocked: me.IsBlocked,
AutoGroups: append(append([]string(nil), before...), groupID),
})
require.NoError(t, err, "add the caller to the policy source group")
t.Cleanup(func() {
_, _ = srv.API().Users.Update(context.Background(), me.Id, api.PutApiUsersUserIdJSONRequestBody{
Role: me.Role,
IsBlocked: me.IsBlocked,
AutoGroups: before,
})
})
}
// configProvider returns the agent-config entry for the named provider, nil
// when the answer does not offer it. The suite shares one account, so other
// tests' fixtures may add unrelated providers to the caller's answer.
func configProvider(cfg api.AgentNetworkAgentConfig, name string) *api.AgentNetworkAgentConfigProvider {
for i := range cfg.Providers {
if cfg.Providers[i].Name == name {
return &cfg.Providers[i]
}
}
return nil
}
// TestAgentConfigAllowlistOfDeclaredModels reproduces the post-#7221 field
// report: a provider carrying a declared model set plus a policy guardrail
// whose allowlist holds those same declared ids must advertise the models on
// GET /api/agent-network/agent-config — the guardrail was built FROM the
// provider's model list (the dashboard's allowlist picker persists the
// declared ids verbatim), so nothing about the setup excludes them.
//
// The plain case passes today. The path-style case (Bedrock; Vertex has the
// same shape) fails: the declared id is compared through the proxy parser's
// canonical form (region prefix and version suffix stripped) while the
// allowlist entry is not, so the raw-vs-raw pair never intersects and the
// caller sees an empty model list. The same one-sided normalization sits in
// policyPermitsModel, so the proxy also denies the model at request time —
// the guardrail meant to allow exactly this model turns it off end to end.
func TestAgentConfigAllowlistOfDeclaredModels(t *testing.T) {
ctx := context.Background()
cases := []struct {
name string
catalogID string
upstream string
declared string
}{
{
name: "plain-declared-id",
catalogID: "openai_api",
upstream: "https://api.openai.com",
declared: "gpt-4o-mini",
},
{
// The operator declares the id AWS issues — region-prefixed
// inference profile with a version suffix — and the allowlist
// picker copies it as-is.
name: "bedrock-declared-id",
catalogID: "bedrock_api",
upstream: "https://bedrock-runtime.eu-central-1.amazonaws.com",
declared: "eu.anthropic.claude-sonnet-4-5-20250929-v1:0",
},
}
for _, tc := range cases {
tc := tc
t.Run(tc.name, func(t *testing.T) {
grp, err := srv.API().Groups.Create(ctx, api.PostApiGroupsJSONRequestBody{Name: "e2e-agentcfg-" + tc.name})
require.NoError(t, err, "create source group")
t.Cleanup(func() { _ = srv.API().Groups.Delete(context.Background(), grp.Id) })
joinGroup(t, ctx, grp.Id)
providerName := "e2e-agentcfg-" + tc.name
prov, err := srv.CreateProvider(ctx, api.AgentNetworkProviderRequest{
Name: providerName,
ProviderId: tc.catalogID,
UpstreamUrl: tc.upstream,
ApiKey: ptr("sk-dummy-e2e-key"),
Enabled: ptr(true),
Models: &[]api.AgentNetworkProviderModel{{Id: tc.declared, InputPer1k: 0.001, OutputPer1k: 0.002}},
})
require.NoError(t, err, "create provider")
t.Cleanup(func() { _ = srv.DeleteProvider(context.Background(), prov.Id) })
// Allowlist exactly the declared model, the way the dashboard
// builds a guardrail from the provider's model list.
var gr api.AgentNetworkGuardrailRequest
gr.Name = "e2e-agentcfg-" + tc.name
gr.Checks.ModelAllowlist.Enabled = true
gr.Checks.ModelAllowlist.Models = []string{tc.declared}
guard, err := srv.CreateGuardrail(ctx, gr)
require.NoError(t, err, "create guardrail")
t.Cleanup(func() { _ = srv.DeleteGuardrail(context.Background(), guard.Id) })
pol, err := srv.CreatePolicy(ctx, api.AgentNetworkPolicyRequest{
Name: "e2e-agentcfg-" + tc.name,
Enabled: ptr(true),
SourceGroups: []string{grp.Id},
DestinationProviderIds: []string{prov.Id},
GuardrailIds: &[]string{guard.Id},
})
require.NoError(t, err, "create policy")
t.Cleanup(func() { _ = srv.DeletePolicy(context.Background(), pol.Id) })
cfg, err := srv.GetAgentConfig(ctx)
require.NoError(t, err, "read the caller-scoped agent config")
require.True(t, cfg.Configured, "the account endpoint is bootstrapped by TestMain")
entry := configProvider(cfg, providerName)
require.NotNil(t, entry, "the policy authorizes the caller for the provider, so it must be offered")
assert.False(t, entry.AllModelsAllowed, "an allowlist guardrail restricts the provider")
assert.Equal(t, []string{tc.declared}, entry.Models,
"the allowlist holds the provider's own declared id, so that model must be advertised")
})
}
}
@@ -0,0 +1,235 @@
//go:build e2e
package agentnetwork
import (
"context"
"net/http"
"os"
"strings"
"testing"
"time"
"github.com/stretchr/testify/require"
"github.com/netbirdio/netbird/shared/management/client/rest"
"github.com/netbirdio/netbird/shared/management/http/api"
)
// credentialCase is one vendor to try the save-time check against. The key is
// the real one the suite already sources; corrupting it is what produces the
// refusal, so the pair of cases differ only in the credential.
type credentialCase struct {
name string
catalogID string
upstream string
apiKey string
}
// liveCredentialCases mirrors the discovery matrix's env gating so a partial
// key set still yields partial coverage. Vertex is left out: its credential is
// a service-account keyfile, and mangling one produces a client-side parse
// failure rather than the vendor refusal this is about.
func liveCredentialCases() []credentialCase {
var cases []credentialCase
if k := os.Getenv("OPENAI_TOKEN"); k != "" {
cases = append(cases, credentialCase{
name: "openai", catalogID: "openai_api",
upstream: "https://api.openai.com", apiKey: k,
})
}
if k := os.Getenv("ANTHROPIC_TOKEN"); k != "" {
cases = append(cases, credentialCase{
name: "anthropic", catalogID: "anthropic_api",
upstream: "https://api.anthropic.com", apiKey: k,
})
}
if k := os.Getenv("AWS_BEARER_TOKEN_BEDROCK"); k != "" {
region := os.Getenv("AWS_REGION")
if region == "" {
region = "eu-central-1"
}
cases = append(cases, credentialCase{
name: "bedrock", catalogID: "bedrock_api",
upstream: "https://bedrock-runtime." + region + ".amazonaws.com", apiKey: k,
})
}
return cases
}
// TestLiveProviderCredentialCheck drives the save-time check against the real
// vendors. A unit test can only assert that a mocked refusal is classified;
// what it cannot show is that these vendors refuse a bad key on their listing
// endpoint at all, which is the assumption the whole feature rests on.
//
// The good-key case matters just as much as the bad one: a check that refused
// everything would pass a test asserting only the refusal, and would make the
// product unusable.
//
// The suite asserts on the vendors themselves, so it inherits their
// availability: the check blocks on 5xx and 429 by design, and a vendor outage
// or a rate limit during a run fails "a good credential saves" with a
// perfectly valid key. There is no retry here on purpose — a retry loop would
// also mask the outage classification these tests exist to prove. Re-run the
// job.
func TestLiveProviderCredentialCheck(t *testing.T) {
cases := liveCredentialCases()
if len(cases) == 0 {
t.Skip("no live provider credentials in the environment; source ~/.llm-keys to run")
}
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Minute)
defer cancel()
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
t.Run("a good credential saves", func(t *testing.T) {
prov, err := srv.CreateProvider(ctx, credentialProviderRequest(tc, "e2e-cred-ok-"+tc.name, tc.apiKey))
require.NoError(t, err, "the suite's own credential must pass its check")
t.Cleanup(func() { _ = srv.DeleteProvider(context.Background(), prov.Id) })
require.NotEmpty(t, prov.Id)
})
t.Run("a rejected credential is refused", func(t *testing.T) {
_, err := srv.CreateProvider(ctx, credentialProviderRequest(tc, "e2e-cred-bad-"+tc.name, corrupt(tc.apiKey)))
require.Error(t, err, "a key the vendor rejects must not save")
var apiErr *rest.APIError
require.ErrorAs(t, err, &apiErr)
require.Equal(t, http.StatusUnprocessableEntity, apiErr.StatusCode,
"a refused credential is the caller's problem to fix, not a server fault")
require.Contains(t, strings.ToLower(apiErr.Message), "rejected the credential",
"the message must name the credential rather than the url")
// The record must be absent, not merely unusable: a provider
// saved despite its check is the state this prevents.
all, listErr := srv.ListProviders(ctx)
require.NoError(t, listErr)
for _, p := range all {
require.NotEqual(t, "e2e-cred-bad-"+tc.name, p.Name, "a refused provider must not be stored")
}
})
})
}
}
// TestLiveProviderUrlCheck points a real credential at a host that is not the
// vendor's API. It is the half of the split a wrong key cannot exercise: the
// operator has to be told the URL is at fault while their key is fine.
//
// One vendor, deliberately. The transport classification under test happens
// before any vendor is reached, so running it per configured vendor would
// repeat the same code path and multiply the wall-clock of a suite that
// already creates real records. cases[0] is whichever vendor the environment
// supplies first.
func TestLiveProviderUrlCheck(t *testing.T) {
cases := liveCredentialCases()
if len(cases) == 0 {
t.Skip("no live provider credentials in the environment; source ~/.llm-keys to run")
}
tc := cases[0]
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute)
defer cancel()
// A name that resolves nowhere. The check has to reach a verdict without
// the vendor's help, which is the transport half of the classification.
req := credentialProviderRequest(tc, "e2e-cred-badurl", tc.apiKey)
req.UpstreamUrl = "https://not-a-real-vendor-host.netbird-e2e.invalid"
_, err := srv.CreateProvider(ctx, req)
require.Error(t, err, "an upstream that does not resolve must not save")
var apiErr *rest.APIError
require.ErrorAs(t, err, &apiErr)
require.Equal(t, http.StatusUnprocessableEntity, apiErr.StatusCode)
require.Contains(t, strings.ToLower(apiErr.Message), "could not be reached",
"the message must name the url rather than the credential")
// An error is not the same fact as an absent record: a handler that saved
// first and reported afterwards would satisfy everything above.
all, listErr := srv.ListProviders(ctx)
require.NoError(t, listErr)
for _, p := range all {
require.NotEqual(t, "e2e-cred-badurl", p.Name, "a refused provider must not be stored")
}
}
// TestLiveProviderUpdateKeepsTheWorkingKey is the state the check exists to
// prevent on the update path: a rejected rotation that has already replaced
// the credential would take a working provider down.
//
// Also one vendor: the behaviour is in the manager's merge, not in any
// vendor's response, and each run creates and mutates a real provider record.
func TestLiveProviderUpdateKeepsTheWorkingKey(t *testing.T) {
cases := liveCredentialCases()
if len(cases) == 0 {
t.Skip("no live provider credentials in the environment; source ~/.llm-keys to run")
}
tc := cases[0]
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Minute)
defer cancel()
prov, err := srv.CreateProvider(ctx, credentialProviderRequest(tc, "e2e-cred-rotate", tc.apiKey))
require.NoError(t, err)
t.Cleanup(func() { _ = srv.DeleteProvider(context.Background(), prov.Id) })
rotation := credentialProviderRequest(tc, "e2e-cred-rotate", corrupt(tc.apiKey))
_, err = srv.UpdateProvider(ctx, prov.Id, rotation)
require.Error(t, err, "a rotation the vendor rejects must not be stored")
var apiErr *rest.APIError
require.ErrorAs(t, err, &apiErr)
require.Equal(t, http.StatusUnprocessableEntity, apiErr.StatusCode)
// The stored key is never returned by the API, so the proof that it
// survived is that an edit which reuses it still passes its check. A
// replaced key would fail here exactly as the rotation just did.
//
// The trailing slash is what makes that an actual check: an edit touching
// neither the url, the key nor the catalog entry is stored without asking
// the vendor anything, so a rename alone would pass whatever is on the
// record. Only the host is read out of the upstream, so the same vendor is
// reached — but the string differs, and the check runs.
recheck := credentialProviderRequest(tc, "e2e-cred-rotate-renamed", "")
recheck.UpstreamUrl = tc.upstream + "/"
updated, err := srv.UpdateProvider(ctx, prov.Id, recheck)
require.NoError(t, err, "the working key must still be the stored one")
require.Equal(t, "e2e-cred-rotate-renamed", updated.Name)
}
// credentialProviderRequest builds a create/update body for a case. An empty apiKey is
// omitted rather than sent blank, which is how the form asks to keep whatever
// is already stored.
func credentialProviderRequest(tc credentialCase, name, apiKey string) api.AgentNetworkProviderRequest {
req := api.AgentNetworkProviderRequest{
Name: name,
ProviderId: tc.catalogID,
UpstreamUrl: tc.upstream,
Enabled: ptr(true),
}
if apiKey != "" {
req.ApiKey = &apiKey
}
return req
}
// corrupt returns a key the vendor will reject while keeping the shape of the
// original. Replacing the last character rather than appending keeps any
// length or prefix validation satisfied, so the refusal comes from the vendor
// checking the secret rather than from it rejecting an obviously malformed
// one.
func corrupt(key string) string {
if key == "" {
return key
}
last := key[len(key)-1]
replacement := byte('A')
if last == 'A' {
replacement = 'B'
}
return key[:len(key)-1] + string(replacement)
}
@@ -0,0 +1,132 @@
//go:build e2e
package agentnetwork
import (
"context"
"testing"
"time"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/netbirdio/netbird/e2e/harness"
"github.com/netbirdio/netbird/shared/management/http/api"
)
// TestModelAllowlistOfDeclaredIDsServed drives the setup an operator actually
// builds for a path-routed provider: the models are declared in the form the
// vendor issues (Bedrock's region-prefixed, versioned inference-profile id;
// Vertex's model@version), and the guardrail allowlist is built from that
// declared list — the dashboard's allowlist picker persists the declared ids
// verbatim. A request for the declared model must be served end to end, and a
// model outside the allowlist must still be denied.
//
// TestModelAllowlistEnforced never caught this because it registers and
// allowlists the pre-normalized catalog form (see the catalogModel comment
// there and the one in providerRequest: "register the normalized form here or
// routing fails as model_not_routable") — the harness encoded the
// canonicalization workaround instead of the shape operators configure.
func TestModelAllowlistOfDeclaredIDsServed(t *testing.T) {
var providers []providerCase
for _, pc := range availableProviders() {
if pc.kind == harness.WireBedrock || pc.kind == harness.WireVertex {
providers = append(providers, pc)
}
}
if len(providers) == 0 {
t.Skip("no path-routed provider keys set (AWS_BEARER_TOKEN_BEDROCK / GOOGLE_VERTEX_*); source ~/.llm-keys")
}
ctx, cancel := context.WithTimeout(context.Background(), 20*time.Minute)
defer cancel()
grp, err := srv.API().Groups.Create(ctx, api.PostApiGroupsJSONRequestBody{Name: "e2e-declared-allowlist"})
require.NoError(t, err, "create group")
t.Cleanup(func() { _ = srv.API().Groups.Delete(context.Background(), grp.Id) })
ephemeral := false
sk, err := srv.API().SetupKeys.Create(ctx, api.PostApiSetupKeysJSONRequestBody{
Name: "e2e-declared-allowlist-client",
Type: "reusable",
ExpiresIn: 86400,
UsageLimit: 0,
AutoGroups: []string{grp.Id},
Ephemeral: &ephemeral,
})
require.NoError(t, err, "mint setup key")
t.Cleanup(func() { _ = srv.API().SetupKeys.Delete(context.Background(), sk.Id) })
// Providers declaring the raw vendor-issued model id — NOT the normalized
// catalog form providerRequest would register.
ids := make([]string, 0, len(providers))
declared := make([]string, 0, len(providers))
for _, pc := range providers {
req := providerRequest(pc)
req.Models = &[]api.AgentNetworkProviderModel{{Id: pc.model, InputPer1k: 0.001, OutputPer1k: 0.002}}
prov, perr := srv.CreateProvider(ctx, req)
require.NoError(t, perr, "create provider %s", pc.name)
id := prov.Id
ids = append(ids, id)
declared = append(declared, pc.model)
t.Cleanup(func() { _ = srv.DeleteProvider(context.Background(), id) })
}
// Guardrail allowlisting the declared ids verbatim, the way the dashboard
// builds an allowlist from the providers' model lists.
var gr api.AgentNetworkGuardrailRequest
gr.Name = "e2e-declared-allowlist"
gr.Checks.ModelAllowlist.Enabled = true
gr.Checks.ModelAllowlist.Models = declared
guard, err := srv.CreateGuardrail(ctx, gr)
require.NoError(t, err, "create guardrail")
t.Cleanup(func() { _ = srv.DeleteGuardrail(context.Background(), guard.Id) })
enabled := true
pol, err := srv.CreatePolicy(ctx, api.AgentNetworkPolicyRequest{
Name: "e2e-declared-allowlist",
Enabled: &enabled,
SourceGroups: []string{grp.Id},
DestinationProviderIds: ids,
GuardrailIds: &[]string{guard.Id},
})
require.NoError(t, err, "create policy")
t.Cleanup(func() { _ = srv.DeletePolicy(context.Background(), pol.Id) })
settings, err := srv.GetSettings(ctx)
require.NoError(t, err, "read settings for endpoint")
require.NotEmpty(t, settings.Endpoint, "agent-network endpoint must be assigned")
proxyToken, err := srv.CreateProxyTokenCLI(ctx, "e2e-proxy-declared-allowlist")
require.NoError(t, err, "mint proxy token via CLI")
px, err := harness.StartProxy(ctx, srv, proxyToken)
require.NoError(t, err, "start proxy")
t.Cleanup(func() { _ = px.Terminate(context.Background()) })
cl, err := harness.StartClient(ctx, srv, sk.Key)
require.NoError(t, err, "start client")
t.Cleanup(func() { _ = cl.Terminate(context.Background()) })
require.NoError(t, cl.WaitConnected(ctx, 90*time.Second), "client must connect to management")
// Probe first: the GET resolves the endpoint (DNS error fails) and its first packet wakes the lazy proxy peer, so WaitProxyPeer sees it connected; any HTTP status counts.
proxyIP, err := cl.ResolveProxyIP(ctx, settings.Endpoint)
require.NoError(t, err, "resolve agent-network endpoint to proxy IP")
if err := cl.WaitProxyPeer(ctx, 180*time.Second); err != nil {
t.Fatalf("client did not see the proxy peer: %v\n=== proxy logs ===\n%s", err, px.Logs(context.Background()))
}
for _, pc := range providers {
pc := pc
t.Run(pc.name, func(t *testing.T) {
// The model the operator declared and allowlisted is served end to
// end: the route must claim it and the guardrail must permit it,
// both through the canonicalization the parser applies at request
// time — whatever id form the operator configured.
assert.Equal(t, 200, sendModel(ctx, t, cl, settings.Endpoint, proxyIP, pc, pc.model),
"the declared and allowlisted model must be served for %s", pc.name)
// A model outside the allowlist stays denied.
assert.Equal(t, 403, sendModel(ctx, t, cl, settings.Endpoint, proxyIP, pc, disallowedModel(pc)),
"model outside the allowlist must be denied for %s", pc.name)
})
}
}
+9 -3
View File
@@ -16,14 +16,20 @@ import (
func ptr[T any](v T) *T { return &v }
// newProvider creates an OpenAI-catalog provider with a dummy key (these tests
// never call the upstream) and registers cleanup.
// newProvider creates an OpenAI-catalog provider these tests can hang a policy
// off, and registers cleanup. Nothing here calls the upstream.
func newProvider(t *testing.T, ctx context.Context, name string) api.AgentNetworkProvider {
t.Helper()
// A provider save is credential-checked against the vendor, and every
// caller here wants a provider row to hang a policy off rather than a
// working upstream. A private address is left unchecked — the proxy would
// reach it through the tunnel, management cannot reach it at all — which
// keeps this fixture independent of whether the run has vendor keys, and
// covers the unchecked-provider-still-saves path while it is at it.
prov, err := srv.CreateProvider(ctx, api.AgentNetworkProviderRequest{
Name: name,
ProviderId: "openai_api",
UpstreamUrl: "https://api.openai.com",
UpstreamUrl: "https://10.255.255.1",
ApiKey: ptr("sk-dummy-e2e-key"),
})
require.NoError(t, err, "create provider %q", name)
+1 -1
View File
@@ -3,7 +3,7 @@
# artifact), so this mirrors its alpine runtime + entrypoint while compiling the
# CGO-free client inline. BuildKit cache mounts keep rebuilds incremental.
FROM golang:1.25-bookworm AS builder
FROM golang:1.26.7-bookworm AS builder
WORKDIR /src
COPY go.mod go.sum ./
RUN --mount=type=cache,target=/go/pkg/mod go mod download
+9
View File
@@ -100,6 +100,15 @@ func (c *Combined) SetProviderEnabled(ctx context.Context, id string, enabled bo
return err
}
// GetAgentConfig returns the caller-scoped self-service connection config —
// the answer the dashboard's "Connect Agent" view renders for the PAT's user.
// Providers appear only when the caller's own groups intersect an enabled
// policy's source groups, so tests must place the PAT user into the policy's
// source group first (via the Users API auto-groups).
func (c *Combined) GetAgentConfig(ctx context.Context) (api.AgentNetworkAgentConfig, error) {
return anRequest[api.AgentNetworkAgentConfig](ctx, c, http.MethodGet, "/api/agent-network/agent-config", nil)
}
// CreatePolicy creates an agent-network policy.
func (c *Combined) CreatePolicy(ctx context.Context, req api.AgentNetworkPolicyRequest) (api.AgentNetworkPolicy, error) {
return anRequest[api.AgentNetworkPolicy](ctx, c, http.MethodPost, "/api/agent-network/policies", req)
+20 -3
View File
@@ -31,6 +31,9 @@ const (
// Client is a running NetBird client container joined to the combined server.
type Client struct {
container testcontainers.Container
// name is the container hostname the agent reports to management at
// registration — the name the peer appears under in the peers API.
name string
}
// clientOptions is what the ClientOption values assemble.
@@ -99,24 +102,38 @@ func StartClient(ctx context.Context, c *Combined, setupKey string, opts ...Clie
if err != nil {
return nil, fmt.Errorf("start client container: %w", err)
}
return &Client{container: ctr}, nil
return &Client{container: ctr, name: o.name}, nil
}
// Hostname returns the container hostname the agent reports to management —
// the name the registered peer appears under in the peers API.
func (cl *Client) Hostname() string {
return cl.name
}
// Restart bounces the client connection (netbird down/up) so it pulls a fresh
// network map — the documented workaround for a freshly-joined client not yet
// seeing a synthesized agent-network service.
func (cl *Client) Restart(ctx context.Context) error {
return cl.Up(ctx)
}
// Up re-runs `netbird up` inside the client with the given extra flags (e.g.
// "--allow-remote-jobs"), bouncing the connection first so the new config is
// picked up and re-synced to management. Used to toggle peer options that ride
// on the login/sync request without recreating the container.
func (cl *Client) Up(ctx context.Context, extraArgs ...string) error {
if _, _, err := cl.container.Exec(ctx, []string{"netbird", "down"}, tcexec.Multiplexed()); err != nil {
return fmt.Errorf("netbird down: %w", err)
}
time.Sleep(2 * time.Second)
code, reader, err := cl.container.Exec(ctx, []string{"netbird", "up"}, tcexec.Multiplexed())
code, reader, err := cl.container.Exec(ctx, append([]string{"netbird", "up"}, extraArgs...), tcexec.Multiplexed())
if err != nil {
return fmt.Errorf("netbird up: %w", err)
}
if code != 0 {
out, _ := io.ReadAll(reader)
return fmt.Errorf("netbird up exited %d: %s", code, string(out))
return fmt.Errorf("netbird up %v exited %d: %s", extraArgs, code, string(out))
}
return nil
}
+47
View File
@@ -0,0 +1,47 @@
//go:build e2e
// Package remotejobs holds the container-based e2e suite for the remote-jobs
// opt-in (PR #7153) and the debug-bundle job parameters anonymize_level /
// upload_url (PR #7147). A combined server is built and bootstrapped once per
// package run (TestMain) and shared via srv; each test registers its own client
// and cleans it up.
package remotejobs
import (
"context"
"fmt"
"os"
"testing"
"time"
"github.com/netbirdio/netbird/e2e/harness"
)
// srv is the shared combined server for the package, PAT-authenticated by the
// time any Test runs.
var srv *harness.Combined
func TestMain(m *testing.M) {
os.Exit(run(m))
}
func run(m *testing.M) int {
// Generous timeout to cover a cold image build on first run.
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Minute)
defer cancel()
var err error
srv, err = harness.StartCombined(ctx)
if err != nil {
fmt.Fprintf(os.Stderr, "e2e: start combined server: %v\n", err)
return 1
}
defer func() { _ = srv.Terminate(context.Background()) }()
if _, err := srv.Bootstrap(ctx); err != nil {
fmt.Fprintf(os.Stderr, "e2e: bootstrap admin PAT: %v\n", err)
return 1
}
return m.Run()
}
+197
View File
@@ -0,0 +1,197 @@
//go:build e2e
package remotejobs
import (
"context"
"strings"
"testing"
"time"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/netbirdio/netbird/e2e/harness"
"github.com/netbirdio/netbird/shared/management/http/api"
)
const (
refusedReason = "remote jobs are not enabled on this peer"
// testUploadURL is the debug-bundle upload URL the job subtests pass; it only
// needs to be a well-formed https URL with a host (see ValidateBundleUploadURL).
testUploadURL = "https://uploads.example.com/bundle"
)
// TestRemoteJobsOptInAndBundleParams exercises the two PRs end-to-end against a
// live management server and a real client:
//
// - #7153: the peer's remote-jobs opt-in defaults off, is reported to
// management (visible via the peers API as remote_jobs_allowed), and gates
// job execution on the client — a streamed job is refused until the peer
// opts in with `netbird up --allow-remote-jobs`, after which it runs.
// - #7147: the debug-bundle job's anonymize_level is validated (an unknown
// value is rejected at creation) and normalized (trimmed + lowercased) in
// the stored job the API returns.
func TestRemoteJobsOptInAndBundleParams(t *testing.T) {
ctx := context.Background()
// A group for the setup key to auto-assign; peers must land in some group.
grp, err := srv.API().Groups.Create(ctx, api.PostApiGroupsJSONRequestBody{Name: "e2e-remotejobs"})
require.NoError(t, err, "create group")
t.Cleanup(func() { _ = srv.API().Groups.Delete(context.Background(), grp.Id) })
sk, err := srv.API().SetupKeys.Create(ctx, api.PostApiSetupKeysJSONRequestBody{
Name: "e2e-remotejobs",
Type: "reusable",
ExpiresIn: 86400,
UsageLimit: 0,
AutoGroups: []string{grp.Id},
})
require.NoError(t, err, "mint setup key")
require.NotEmpty(t, sk.Key, "setup key plaintext")
t.Cleanup(func() { _ = srv.API().SetupKeys.Delete(context.Background(), sk.Id) })
// Start the client with a plain `netbird up` (remote jobs NOT enabled).
cl, err := harness.StartClient(ctx, srv, sk.Key)
require.NoError(t, err, "start client")
t.Cleanup(func() { _ = cl.Terminate(context.Background()) })
require.NoError(t, cl.WaitConnected(ctx, 90*time.Second), "client must connect to management")
peerID := waitForPeer(ctx, t, cl.Hostname())
t.Run("opt-in flag defaults to false and is reported to management (#7153)", func(t *testing.T) {
p, err := srv.API().Peers.Get(ctx, peerID)
require.NoError(t, err)
allowed := remoteJobsAllowed(p)
require.NotNil(t, allowed, "remote_jobs_allowed must be present on the peer API")
assert.False(t, *allowed, "a peer that ran plain `netbird up` must default to opt-out")
})
t.Run("anonymize_level is validated and normalized (#7147)", func(t *testing.T) {
// Unknown level is rejected at job creation.
_, err := srv.API().Peers.Jobs(peerID).Create(ctx, bundleJob("bogus", testUploadURL))
require.Error(t, err, "an unknown anonymize_level must be rejected")
assert.Contains(t, strings.ToLower(err.Error()), "anonymize_level",
"the rejection must name the offending field")
// A messy but valid level is normalized (trimmed + lowercased) in the
// stored job the API echoes back.
job, err := srv.API().Peers.Jobs(peerID).Create(ctx, bundleJob(" Strict ", testUploadURL))
require.NoError(t, err, "a valid anonymize_level must be accepted")
bw, err := job.Workload.AsBundleWorkloadResponse()
require.NoError(t, err, "job workload must be a bundle")
require.NotNil(t, bw.Parameters.AnonymizeLevel)
assert.Equal(t, "strict", *bw.Parameters.AnonymizeLevel,
"anonymize_level must be normalized to trimmed lowercase")
waitForJobTerminal(ctx, t, peerID, job.Id) // let it settle before the next create
})
t.Run("a job is refused while the peer has not opted in (#7153 enforcement)", func(t *testing.T) {
job, err := srv.API().Peers.Jobs(peerID).Create(ctx, bundleJob("default", testUploadURL))
require.NoError(t, err, "job creation itself is allowed; enforcement is on the client")
final := waitForJobTerminal(ctx, t, peerID, job.Id)
assert.Equal(t, api.JobResponseStatusFailed, final.Status, "the client must refuse the job")
require.NotNil(t, final.FailedReason)
assert.Contains(t, *final.FailedReason, refusedReason,
"the failure must be the opt-out refusal, not some other error")
})
t.Run("opting in flips the flag and lets the job run (#7153)", func(t *testing.T) {
require.NoError(t, cl.Up(ctx, "--allow-remote-jobs"), "re-run up with --allow-remote-jobs")
require.NoError(t, cl.WaitConnected(ctx, 90*time.Second), "client must reconnect")
// The new opt-in must round-trip to management and surface on the API.
require.Eventually(t, func() bool {
p, err := srv.API().Peers.Get(ctx, peerID)
if err != nil {
return false
}
allowed := remoteJobsAllowed(p)
return allowed != nil && *allowed
}, 60*time.Second, 2*time.Second, "remote_jobs_allowed must become true after opt-in")
// The same job that was refused before must now be accepted for
// execution: whatever its outcome, it must NOT be the opt-out refusal.
job, err := srv.API().Peers.Jobs(peerID).Create(ctx, bundleJob("default", testUploadURL))
require.NoError(t, err)
final := waitForJobTerminal(ctx, t, peerID, job.Id)
if final.Status == api.JobResponseStatusFailed && final.FailedReason != nil {
assert.NotContains(t, *final.FailedReason, refusedReason,
"once opted in, the job must not be refused for opt-out; any failure must be for another reason (e.g. upload)")
}
})
}
// remoteJobsAllowed returns the peer's remote-jobs opt-in flag from the API
// response (nil if the peer or its local flags are absent).
func remoteJobsAllowed(p *api.Peer) *bool {
if p == nil || p.LocalFlags == nil {
return nil
}
return p.LocalFlags.RemoteJobsAllowed
}
// bundleJob builds a debug-bundle job request with the given anonymize_level
// (omitted when empty) and upload_url (omitted when empty).
func bundleJob(anonymizeLevel, uploadURL string) api.JobRequest {
params := api.BundleParameters{
Anonymize: true,
LogFileCount: 1,
}
if anonymizeLevel != "" {
params.AnonymizeLevel = &anonymizeLevel
}
if uploadURL != "" {
params.UploadUrl = &uploadURL
}
var wl api.WorkloadRequest
// FromBundleWorkloadRequest cannot fail for a well-formed value.
_ = wl.FromBundleWorkloadRequest(api.BundleWorkloadRequest{
Type: api.WorkloadTypeBundle,
Parameters: params,
})
return api.JobRequest{Workload: wl}
}
// waitForPeer polls the peers API until the client that registered under the
// given hostname appears and returns its ID. Matching by hostname rather than
// taking the first list entry keeps the test correct if the account ever holds
// more than one peer (a shared bootstrap account, or a second client added to
// the package).
func waitForPeer(ctx context.Context, t *testing.T, hostname string) string {
t.Helper()
var peerID string
require.Eventually(t, func() bool {
peers, err := srv.API().Peers.List(ctx)
if err != nil {
return false
}
for _, p := range peers {
if p.Hostname == hostname {
peerID = p.Id
return true
}
}
return false
}, 60*time.Second, 2*time.Second, "the client peer must register with management")
return peerID
}
// waitForJobTerminal polls a job until it leaves the pending state, then returns
// the final response.
func waitForJobTerminal(ctx context.Context, t *testing.T, peerID, jobID string) *api.JobResponse {
t.Helper()
var final *api.JobResponse
require.Eventually(t, func() bool {
j, err := srv.API().Peers.Jobs(peerID).Get(ctx, jobID)
if err != nil || j == nil {
return false
}
if j.Status == api.JobResponseStatusPending {
return false
}
final = j
return true
}, 120*time.Second, 2*time.Second, "job must reach a terminal state")
return final
}