flowstream

package
v0.8.0 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: 17 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type FlowFilterDirection

type FlowFilterDirection int32

FlowFilterDirection controls which endpoint of a flow the directional filters are matched against.

const (
	FlowFilterDirectionBoth FlowFilterDirection = 0
	FlowFilterDirectionFrom FlowFilterDirection = 1
	FlowFilterDirectionTo   FlowFilterDirection = 2
)

type FlowStreamFilter

type FlowStreamFilter struct {
	Namespaces       []string
	PodNames         []string
	PodLabelSelector string
	ServiceNames     []string
	FlowTypes        []apisv1.FlowType
	IPs              []string
	Direction        FlowFilterDirection
}

FlowStreamFilter represents the parsed query parameters for the flow stream endpoint. All specified filters are AND-ed. Within each filter, values are OR-ed.

type FlowStreamSubscriber

type FlowStreamSubscriber interface {
	// Subscribe starts streaming flows matching the given filter.
	// It returns a channel of FlowStreamEvent and a channel of errors.
	// The caller should read from both channels until they are closed.
	// Cancel the context to stop the stream.
	Subscribe(ctx context.Context, filter *FlowStreamFilter) (<-chan apisv1.FlowStreamEvent, <-chan error)
}

FlowStreamSubscriber provides a channel-based interface for streaming flow data. Implementations connect to the FlowAggregator's gRPC FlowStreamService (see grpc.go) and relay flow events to the caller.

type GRPCConfig

type GRPCConfig struct {
	Address string
	// TLSConfig is the TLS configuration used for the gRPC connection.
	TLSConfig *tls.Config
}

GRPCConfig holds the connection parameters for the FlowAggregator gRPC server. The FlowStreamService uses server-side TLS only (no client authentication).

type GRPCFlowStreamSubscriber

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

GRPCFlowStreamSubscriber connects to the FlowAggregator's FlowStreamService over gRPC and implements the FlowStreamSubscriber interface.

func NewGRPCFlowStreamSubscriber

func NewGRPCFlowStreamSubscriber(logger logr.Logger, cfg GRPCConfig) (*GRPCFlowStreamSubscriber, error)

func (*GRPCFlowStreamSubscriber) Close

func (h *GRPCFlowStreamSubscriber) Close() error

func (*GRPCFlowStreamSubscriber) Subscribe

func (h *GRPCFlowStreamSubscriber) Subscribe(ctx context.Context, filter *FlowStreamFilter) (<-chan apisv1.FlowStreamEvent, <-chan error)

type SSEHandler

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

SSEHandler handles the SSE endpoint for flow streaming.

Known gap, deliberate for now: this endpoint is authenticated but not authorized per user. The subscriber reaches the Flow Aggregator over antrea-ui's own mTLS gRPC connection, so unlike every other API route, the caller's Kubernetes RBAC has no say in what they see. As an interim measure the route is restricted to the built-in admin and to Kubernetes cluster admins (requireFlowVisibility in pkg/server/api/flowstream.go), which narrows who is exposed but does not close the gap: within that set, every caller still sees every exported flow. Authorization is being implemented upstream in antrea-io/antrea#8221; see the "Flow data is not yet per-user" section of docs/authentication.md.

func NewSSEHandler

func NewSSEHandler(logger logr.Logger, handler FlowStreamSubscriber) *SSEHandler

func (*SSEHandler) StreamFlows

func (h *SSEHandler) StreamFlows(c *gin.Context)

StreamFlows handles GET /api/v1/flows/stream as an SSE endpoint.

Directories

Path Synopsis
Package testing is a generated GoMock package.
Package testing is a generated GoMock package.

Jump to

Keyboard shortcuts

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