Documentation
¶
Index ¶
- Constants
- func CreateTestLogger(t *testing.T, enabled bool) *zap.Logger
- func CreateTestToken(audience string) string
- type AuthService
- type FakeServer
- func (fs *FakeServer) DisconnectConfigStream()
- func (fs *FakeServer) GRPCAddress() string
- func (fs *FakeServer) GetConfigurationUpdates(stream pb.KubernetesInfoService_GetConfigurationUpdatesServer) error
- func (fs *FakeServer) SendConfig(resp *pb.GetConfigurationUpdatesResponse)
- func (fs *FakeServer) SendKubernetesNetworkFlows(stream pb.KubernetesInfoService_SendKubernetesNetworkFlowsServer) error
- func (fs *FakeServer) SendKubernetesResources(stream pb.KubernetesInfoService_SendKubernetesResourcesServer) error
- func (fs *FakeServer) SendLogs(stream pb.KubernetesInfoService_SendLogsServer) error
- func (fs *FakeServer) Start() error
- func (fs *FakeServer) Stop()
- type FakeServerTestHarness
- func (h *FakeServerTestHarness) DialGRPC(t *testing.T) *grpc.ClientConn
- func (h *FakeServerTestHarness) LogCurrentState()
- func (h *FakeServerTestHarness) ResetState()
- func (h *FakeServerTestHarness) SetBadInitialCommit(bad bool)
- func (h *FakeServerTestHarness) Start() error
- func (h *FakeServerTestHarness) Stop()
- func (h *FakeServerTestHarness) WaitForCondition(condition func() bool, description string) error
- func (h *FakeServerTestHarness) WaitForConnection() error
- type LogEntry
- type OnboardRequest
- type OnboardResponse
- type ProxyServer
- type ServerState
- func (s *ServerState) AllStreamsOpened() bool
- func (s *ServerState) CheckAndClearBadInitialCommit() bool
- func (s *ServerState) GetCiliumFlowsReceived() int
- func (s *ServerState) GetFiveTupleFlowsReceived() int
- func (s *ServerState) GetResourcesReceived() int
- func (s *ServerState) GetSummary() map[string]any
- func (s *ServerState) IncrementResourcesReceived()
- func (s *ServerState) IsBadInitialCommit() bool
- func (s *ServerState) IsConnectionSuccessful() bool
- func (s *ServerState) IsResourceSnapshotComplete() bool
- func (s *ServerState) MarkConfigStreamOpened()
- func (s *ServerState) MarkFlowsStreamOpened()
- func (s *ServerState) MarkLogsStreamOpened()
- func (s *ServerState) MarkResourcesStreamOpened()
- func (s *ServerState) RecordAuthRequest()
- func (s *ServerState) RecordCiliumFlow()
- func (s *ServerState) RecordFiveTupleFlow()
- func (s *ServerState) RecordKeepalive(stream string)
- func (s *ServerState) RecordOnboardRequest()
- func (s *ServerState) RecordResourceSnapshot()
- func (s *ServerState) Reset()
- func (s *ServerState) SetBadInitialCommit(bad bool)
- type StreamState
- type TestConfig
- type TokenRequest
- type TokenResponse
Constants ¶
const ( AllowedGrantType = "client_credentials" InvalidGrantError = "invalid_grant" )
const ( AuthorizationHeader = "authorization" DefaultClientID = "client_id_1" DefaultClientSecret = "client_secret_1" )
Variables ¶
This section is empty.
Functions ¶
func CreateTestLogger ¶
CreateTestLogger creates a logger for tests.
func CreateTestToken ¶
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 (fs *FakeServer) SendLogs(stream pb.KubernetesInfoService_SendLogsServer) error
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 ¶
func (LogEntry) MarshalLogObject ¶
func (l LogEntry) MarshalLogObject(enc zapcore.ObjectEncoder) error
type OnboardRequest ¶
type OnboardResponse ¶
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.