fakeserver

package
v1.5.1 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 9, 2026 License: Apache-2.0 Imports: 31 Imported by: 0

README

FakeServer

Illumio's cloud-operator streams Kubernetes metadata, logs, and network flows out to a server that you control. The FakeServer provides a convenient way to test the cloud-operator locally without needing a full backend server to collect these streams. It simulates the necessary OAuth endpoints and provides dummy gRPC endpoints to receive the streamed data.

Running FakeServer

Standard Mode

From the repository root, run the following command:

go run ./fakeserver

This starts the FakeServer with its gRPC service listening on port 50051 and its HTTP/OAuth service listening on port 50053.

Proxy Mode

To test the cloud-operator's ability to connect through an HTTP/HTTPS proxy, you can run FakeServer bundled with a simple CONNECT-only proxy:

go run ./fakeserver --proxy

In this mode:

  • The FakeServer itself runs as above (gRPC on 50051, HTTP/OAuth on 50053).
  • A Proxy Server also starts, listening on port 8888 by default. This proxy handles CONNECT requests and tunnels them to the appropriate FakeServer port (50051 or 50053) running on the same host.

Pointing Cloud-Operator to FakeServer

1. Credentials

FakeServer accepts a default set of credentials for the initial onboarding step. You can find these constants in the FakeServer code (or configure FakeServer to expect different ones if needed). The defaults are typically:

  • Client ID: client_id_1
  • Client Secret: client_secret_1

You need to configure the cloud-operator Helm chart (usually via a values.yaml file) with these credentials under the onboardingSecret section:

# Example values snippet for cloud-operator Helm chart
onboardingSecret:
  clientId: "client_id_1"
  clientSecret: "client_secret_1"
2. Endpoints

The cloud-operator needs to know where to reach the FakeServer's HTTP/OAuth endpoints. Since cloud-operator runs inside a Kubernetes pod (often within a Docker Desktop VM or similar environment), you typically cannot use localhost. Instead, use host.docker.internal, which Docker Desktop provides as a DNS name resolving to the host machine.

Configure the following environment variables for the cloud-operator deployment:

# Example values snippet for cloud-operator Helm chart
# Adjust 'env' section based on your chart's structure
env:
  # Required: Allows connecting to the FakeServer's HTTPS endpoint
  # using its self-signed certificate without validation errors.
  tlsSkipVerify: true

  # Endpoint for the initial onboarding request
  onboardingEndpoint: "https://host.docker.internal:50053/api/v1/k8s_cluster/onboard"

  # Endpoint for exchanging credentials for an access token
  tokenEndpoint: "https://host.docker.internal:50053/api/v1/k8s_cluster/authenticate"

  # --- Configuration for Proxy Mode ---
  # If running FakeServer with the --proxy flag (or using any other proxy),
  # set the HTTPS_PROXY environment variable.
  # The proxy listens on port 8888 by default when run via 'fakeserver --proxy'.
  # httpsProxy: "http://host.docker.internal:8888" # Uncomment this line when using the proxy

  # --- Configuration for gRPC Target (Handled Internally) ---
  # The gRPC target address (e.g., host.docker.internal:50051) is usually
  # configured internally by the operator based on successful authentication,
  # or potentially via other specific environment variables if needed.
  # Ensure the operator is configured to eventually target host.docker.internal:50051
  # when connecting from the container to the host.
Important Note on Proxy Usage

When you set httpsProxy for the cloud-operator, its gRPC connections (to port 50051) and its HTTPS calls (to the OAuth endpoints on port 50053) will both be directed through the proxy specified. The proxy (either the built-in one started with fakeserver --proxy or another like Tinyproxy) must be configured to allow CONNECT requests to both host.docker.internal:50051 and host.docker.internal:50053.

3. Putting It All Together

An example values.yaml file (./fakeserver/cloud-operator.fakeserver.yaml) might look like this:

# ./fakeserver/cloud-operator.fakeserver.yaml
# Example values for running cloud-operator against FakeServer

onboardingSecret:
  clientId: "client_id_1"
  clientSecret: "client_secret_1"

deployment:
  env:
    - name: ILLUMIO_TLS_SKIP_VERIFY # Use the actual env var name expected by the operator
      value: "true"
    - name: ILLUMIO_ONBOARDING_ENDPOINT
      value: "https://host.docker.internal:50053/api/v1/k8s_cluster/onboard"
    - name: ILLUMIO_TOKEN_ENDPOINT
      value: "https://host.docker.internal:50053/api/v1/k8s_cluster/authenticate"

    # --- UNCOMMENT FOR PROXY MODE ---
    # - name: HTTPS_PROXY # Or the specific env var the operator uses for proxy
    #   value: "http://host.docker.internal:8888"
    # - name: HTTP_PROXY # Might be needed if operator makes plain HTTP calls
    #   value: "http://host.docker.internal:8888"
    # - name: NO_PROXY # Ensure Kubernetes internal communication isn't proxied
    #   value: "kubernetes.default.svc,.svc,.cluster.local"

# Add other necessary chart values (image repository, tags, resource limits, etc.)

You can then install the Helm chart using these values:

# Ensure FakeServer is running (either standard or --proxy mode)

# Install the chart (adjust chart name/path and release name as needed)
helm install illumio-operator ./charts/cloud-operator \
  --namespace illumio-cloud \
  --create-namespace \
  --values ./fakeserver/cloud-operator.fakeserver.yaml

Monitor the logs of both FakeServer and the cloud-operator pod to verify successful onboarding, authentication, and stream connections.

Documentation

Index

Constants

View Source
const (
	AllowedGrantType   = "client_credentials"
	InvalidGrantError  = "invalid_grant"
	UnauthorizedClient = "unauthorized_client"
)
View Source
const (
	AuthorizationHeader = "authorization"
	DefaultClientID     = "client_id_1"
	DefaultClientSecret = "client_secret_1"
)

Variables

This section is empty.

Functions

func CreateTestLogger

func CreateTestLogger(t *testing.T, enabled bool) *zap.Logger

CreateTestLogger creates a logger for tests.

func CreateTestToken

func CreateTestToken(audience string) string

CreateTestToken creates a signed JWT token for testing. Uses a fixed expiration time to ensure all tests generate identical tokens, allowing the operator to reconnect across test restarts.

Types

type AuthService

type AuthService struct {
	// contains filtered or unexported fields
}

AuthService provides authentication services using client credentials and a token.

type FakeServer

type FakeServer struct {
	pb.UnimplementedKubernetesInfoServiceServer

	Address     string
	HTTPAddress string

	State    *ServerState
	StopChan chan struct{}
	Token    string
	Logger   *zap.Logger

	// ConfigResponses carries config responses to the active stream. Never closed:
	// closing a channel senders write to would panic. Shutdown is signalled via configDone.
	ConfigResponses chan *pb.GetConfigurationUpdatesResponse
	// contains filtered or unexported fields
}

func (*FakeServer) DisconnectConfigStream

func (fs *FakeServer) DisconnectConfigStream()

DisconnectConfigStream ends the active stream (client sees EOF) by closing configDone, then re-arms it for the next stream. The data channel is never closed, so an in-flight SendConfig can't panic. Locked so it can't race the snapshot in SendConfig/handler.

func (*FakeServer) GRPCAddress

func (fs *FakeServer) GRPCAddress() string

GRPCAddress returns the actual address the gRPC server is listening on. Useful when the server is started with ":0" to get the OS-assigned port.

func (*FakeServer) GetConfigurationUpdates

func (fs *FakeServer) GetConfigurationUpdates(stream pb.KubernetesInfoService_GetConfigurationUpdatesServer) error

func (*FakeServer) SendConfig

func (fs *FakeServer) SendConfig(resp *pb.GetConfigurationUpdatesResponse)

SendConfig sends a response to the active config stream.

func (*FakeServer) SendKubernetesNetworkFlows

func (fs *FakeServer) SendKubernetesNetworkFlows(stream pb.KubernetesInfoService_SendKubernetesNetworkFlowsServer) error

func (*FakeServer) SendKubernetesResources

func (fs *FakeServer) SendKubernetesResources(stream pb.KubernetesInfoService_SendKubernetesResourcesServer) error

func (*FakeServer) SendLogs

func (*FakeServer) Start

func (fs *FakeServer) Start() error

func (*FakeServer) Stop

func (fs *FakeServer) Stop()

type FakeServerTestHarness

type FakeServerTestHarness struct {
	Server        *FakeServer
	EnhancedState *ServerState
	Config        TestConfig
	T             *testing.T
}

FakeServerTestHarness wraps FakeServer with test utilities.

func NewTestHarness

func NewTestHarness(t *testing.T, config TestConfig) *FakeServerTestHarness

NewTestHarness creates a new test harness.

func (*FakeServerTestHarness) DialGRPC

func (h *FakeServerTestHarness) DialGRPC(t *testing.T) *grpc.ClientConn

DialGRPC creates a gRPC client connection to the fake server with TLS and token auth.

func (*FakeServerTestHarness) LogCurrentState

func (h *FakeServerTestHarness) LogCurrentState()

LogCurrentState logs the current server state for debugging.

func (*FakeServerTestHarness) ResetState

func (h *FakeServerTestHarness) ResetState()

ResetState resets the server state for a new test phase.

func (*FakeServerTestHarness) SetBadInitialCommit

func (h *FakeServerTestHarness) SetBadInitialCommit(bad bool)

SetBadInitialCommit configures the server to fail the initial commit.

func (*FakeServerTestHarness) Start

func (h *FakeServerTestHarness) Start() error

Start starts the fake server. If AutoInitialConfigSnapshot is true (the default), it also sends the default initial config snapshot messages (UpdateConfiguration + empty ResourceSnapshotComplete) so connected clients complete the initial snapshot.

func (*FakeServerTestHarness) Stop

func (h *FakeServerTestHarness) Stop()

Stop stops the fake server.

func (*FakeServerTestHarness) WaitForCondition

func (h *FakeServerTestHarness) WaitForCondition(condition func() bool, description string) error

WaitForCondition waits for a condition to become true.

func (*FakeServerTestHarness) WaitForConnection

func (h *FakeServerTestHarness) WaitForConnection() error

WaitForConnection waits for the operator to connect and complete resource snapshot.

type LogEntry

type LogEntry map[string]any

func (LogEntry) Level

func (l LogEntry) Level() (zapcore.Level, error)

func (LogEntry) MarshalLogObject

func (l LogEntry) MarshalLogObject(enc zapcore.ObjectEncoder) error

type OnboardRequest

type OnboardRequest struct {
	OnboardingClientId     string `json:"onboardingClientId"`
	OnboardingClientSecret string `json:"onboardingClientSecret"`
}

type OnboardResponse

type OnboardResponse struct {
	ClusterClientId     string `json:"cluster_client_id"`
	ClusterClientSecret string `json:"cluster_client_secret"`
}

type ProxyServer

type ProxyServer struct {
	// contains filtered or unexported fields
}

ProxyServer represents a proxy server focused on HTTP CONNECT.

func NewProxyServer

func NewProxyServer(httpAddress string, logger *zap.Logger) *ProxyServer

NewProxyServer creates and initializes a new ProxyServer.

func (*ProxyServer) ServeHTTP

func (p *ProxyServer) ServeHTTP(w http.ResponseWriter, r *http.Request)

ServeHTTP is the entry point for all HTTP requests made to the proxy server.

func (*ProxyServer) Start

func (p *ProxyServer) Start()

Start launches the ProxyServer's HTTP listener.

func (*ProxyServer) Stop

func (p *ProxyServer) Stop() error

Stop gracefully shuts down the ProxyServer.

type ServerState

type ServerState struct {

	// Legacy fields for backward compatibility
	ConnectionSuccessful bool
	IncorrectCredentials bool
	BadIntialCommit      bool

	// Stream-specific state
	ConfigStream    StreamState
	LogsStream      StreamState
	ResourcesStream StreamState
	FlowsStream     StreamState

	// Auth state
	AuthRequests    int
	OnboardRequests int
	LastAuthTime    time.Time
	LastOnboardTime time.Time

	// Resource tracking
	ResourceSnapshotComplete bool
	ResourcesReceived        int
	MutationsReceived        int

	// Flow tracking
	CiliumFlowsReceived    int
	FiveTupleFlowsReceived int
	// contains filtered or unexported fields
}

ServerState provides detailed tracking of all server activity.

func NewServerState

func NewServerState() *ServerState

NewServerState creates a new state tracker.

func (*ServerState) AllStreamsOpened

func (s *ServerState) AllStreamsOpened() bool

AllStreamsOpened returns true if all streams have been opened.

func (*ServerState) CheckAndClearBadInitialCommit

func (s *ServerState) CheckAndClearBadInitialCommit() bool

CheckAndClearBadInitialCommit returns whether BadIntialCommit was set and, if so, clears it. Returns true when the initial commit should be treated as bad.

func (*ServerState) GetCiliumFlowsReceived

func (s *ServerState) GetCiliumFlowsReceived() int

GetCiliumFlowsReceived returns the count of Cilium flows received.

func (*ServerState) GetFiveTupleFlowsReceived

func (s *ServerState) GetFiveTupleFlowsReceived() int

GetFiveTupleFlowsReceived returns the count of FiveTuple flows received.

func (*ServerState) GetResourcesReceived

func (s *ServerState) GetResourcesReceived() int

GetResourcesReceived returns the count of resources received.

func (*ServerState) GetSummary

func (s *ServerState) GetSummary() map[string]any

GetSummary returns a summary of the current state.

func (*ServerState) IncrementResourcesReceived

func (s *ServerState) IncrementResourcesReceived()

IncrementResourcesReceived increments the count of resources received.

func (*ServerState) IsBadInitialCommit

func (s *ServerState) IsBadInitialCommit() bool

IsBadInitialCommit returns whether BadIntialCommit is set.

func (*ServerState) IsConnectionSuccessful

func (s *ServerState) IsConnectionSuccessful() bool

IsConnectionSuccessful returns whether the connection was successful.

func (*ServerState) IsResourceSnapshotComplete

func (s *ServerState) IsResourceSnapshotComplete() bool

IsResourceSnapshotComplete returns whether the resource snapshot has completed.

func (*ServerState) MarkConfigStreamOpened

func (s *ServerState) MarkConfigStreamOpened()

MarkConfigStreamOpened marks the config stream as opened.

func (*ServerState) MarkFlowsStreamOpened

func (s *ServerState) MarkFlowsStreamOpened()

MarkFlowsStreamOpened marks the flows stream as opened.

func (*ServerState) MarkLogsStreamOpened

func (s *ServerState) MarkLogsStreamOpened()

MarkLogsStreamOpened marks the logs stream as opened.

func (*ServerState) MarkResourcesStreamOpened

func (s *ServerState) MarkResourcesStreamOpened()

MarkResourcesStreamOpened marks the resources stream as opened.

func (*ServerState) RecordAuthRequest

func (s *ServerState) RecordAuthRequest()

RecordAuthRequest records an authentication request.

func (*ServerState) RecordCiliumFlow

func (s *ServerState) RecordCiliumFlow()

RecordCiliumFlow increments the Cilium flow counter.

func (*ServerState) RecordFiveTupleFlow

func (s *ServerState) RecordFiveTupleFlow()

RecordFiveTupleFlow increments the FiveTuple flow counter.

func (*ServerState) RecordKeepalive

func (s *ServerState) RecordKeepalive(stream string)

RecordKeepalive records a keepalive for the specified stream.

func (*ServerState) RecordOnboardRequest

func (s *ServerState) RecordOnboardRequest()

RecordOnboardRequest records an onboard request.

func (*ServerState) RecordResourceSnapshot

func (s *ServerState) RecordResourceSnapshot()

RecordResourceSnapshot marks resource snapshot as complete.

func (*ServerState) Reset

func (s *ServerState) Reset()

Reset resets the legacy connection/commit tracking fields for a new test phase.

func (*ServerState) SetBadInitialCommit

func (s *ServerState) SetBadInitialCommit(bad bool)

SetBadInitialCommit sets the BadIntialCommit flag.

type StreamState

type StreamState struct {
	Opened           bool
	LastActivity     time.Time
	MessagesReceived int
	KeepalivesRecv   int
}

StreamState tracks the state of individual streams.

type TestConfig

type TestConfig struct {
	GRPCAddress   string
	HTTPAddress   string
	Timeout       time.Duration
	PollInterval  time.Duration
	EnableLogging bool
	// AutoInitialConfigSnapshot controls whether Start() sends the default
	// initial config snapshot (UpdateConfiguration + empty
	// ResourceSnapshotComplete). Set to false when tests need to control the
	// initial config snapshot sequence themselves.
	AutoInitialConfigSnapshot bool
}

TestConfig holds configuration for integration tests.

func DefaultTestConfig

func DefaultTestConfig() TestConfig

DefaultTestConfig returns sensible defaults for testing. Uses fixed ports for tests that start the full operator binary (connectivity, flows).

type TokenRequest

type TokenRequest struct {
	GrantType    string
	ClientID     string // Client ID for authentication
	ClientSecret string // Client secret for authentication
}

TokenRequest is a struct to hold the request parameters for the authenticateHandler following the OAuth2.0 specification.

type TokenResponse

type TokenResponse struct {
	AccessToken string `json:"access_token,omitempty"` //nolint:tagliatelle
}

TokenResponse is a struct to hold the response parameters for the authenticateHandler following the OAuth2.0 specification.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL