Documentation
¶
Overview ¶
Package subscriber implements the subscriber service to forward incoming data to remote services.
Index ¶
- Constants
- type BalanceMode
- type Config
- type HTTP
- type PointsWriter
- type Service
- func (s *Service) Close() error
- func (s *Service) Open() error
- func (s *Service) PrepareReloadTLSCertificates(conf Config) (func() error, error)
- func (s *Service) Send(request *coordinator.WritePointsRequest)
- func (s *Service) Statistics(tags map[string]string) []models.Statistic
- func (s *Service) WithLogger(log *zap.Logger)
- type Statistics
- type UDP
- type WriteRequest
Constants ¶
const ( // DefaultHTTPTimeout is the default HTTP timeout for a Config. DefaultHTTPTimeout = 30 * time.Second // DefaultWriteConcurrency is the default write concurrency for a Config. DefaultWriteConcurrency = 40 // DefaultWriteBufferSize is the default write buffer size for a Config. DefaultWriteBufferSize = 1000 )
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type BalanceMode ¶
type BalanceMode int
BalanceMode specifies what balance mode to use on a subscription.
const ( // ALL indicates to send writes to all subscriber destinations. ALL BalanceMode = iota // ANY indicates to send writes to a single subscriber destination, round robin. ANY )
type Config ¶
type Config struct {
// Whether to enable to Subscriber service
Enabled bool `toml:"enabled"`
HTTPTimeout toml.Duration `toml:"http-timeout"`
// InsecureSkipVerify gets passed to the http client, if true, it will
// skip https certificate verification. Defaults to false
InsecureSkipVerify bool `toml:"insecure-skip-verify"`
// configure the path to the PEM encoded CA certs file. If the
// empty string, the default system certs will be used
CaCerts string `toml:"ca-certs"`
// RootCA configures the CA pool used to verify subscription endpoint server
// certificates. It is combined with the legacy CaCerts (ca-certs) setting,
// which continues to work on its own for backwards compatibility.
RootCA *tlsconfig.CAConfig `toml:"root-ca"`
// Certificate and PrivateKey are the client certificate the subscriber
// presents to HTTPS endpoints for mutual TLS. Empty means no client
// certificate is presented.
Certificate string `toml:"certificate"`
PrivateKey string `toml:"private-key"`
// InsecureCertificate is true if the client certificate's file permissions
// should be ignored when it is loaded.
InsecureCertificate bool `toml:"insecure-certificate"`
// IgnoreCertSanityChecks loads the certificate even when it fails the
// checks that decide whether it can be used at all. The checks currently
// only cover server certificates, so this has no effect on the subscriber's
// client certificate today.
IgnoreCertSanityChecks bool `toml:"ignore-cert-sanity-checks"`
// The number of writer goroutines processing the write channel.
WriteConcurrency int `toml:"write-concurrency"`
// The number of in-flight writes buffered in the write channel.
WriteBufferSize int `toml:"write-buffer-size"`
// TotalBufferBytes is the total size in bytes allocated to buffering across all subscriptions.
// Each named subscription will receive an even division of the total.
TotalBufferBytes toml.SSize `toml:"total-buffer-bytes"`
// TLS is a base tls config to use for https clients.
TLS *tls.Config `toml:"-"`
}
Config represents a configuration of the subscriber service.
func (Config) Diagnostics ¶ added in v1.3.0
func (c Config) Diagnostics() (*diagnostics.Diagnostics, error)
Diagnostics returns a diagnostics representation of a subset of the Config.
func (Config) TLSManagerOpts ¶ added in v1.13.0
func (c Config) TLSManagerOpts() []tlsconfig.TLSConfigManagerOpt
TLSManagerOpts returns the list of TLS manager options specified by c.
type HTTP ¶ added in v1.0.0
type HTTP struct {
// contains filtered or unexported fields
}
HTTP supports writing points over HTTP using the line protocol.
func NewHTTPS ¶ added in v1.1.0
func NewHTTPS(addr string, timeout time.Duration, dialTLSContext func(ctx context.Context, network, addr string) (net.Conn, error)) (*HTTP, error)
NewHTTPS returns a new HTTPS points writer with default options and HTTPS configured.
dialTLSContext dials the writer's TLS connections and is the only source of its TLS configuration. It is normally tlsconfig.TLSConfigManager.DialContext, which resolves the configuration on each connection. That is what allows a reloaded configuration to reach a writer that already exists: the writer keeps dialing through the manager rather than through a configuration captured when it was built. It may be nil to use Go's defaults, as NewHTTP does.
timeout bounds the whole request, including the connection and TLS handshake, so no separate handshake timeout is needed.
func (*HTTP) WritePointsContext ¶ added in v1.9.3
func (h *HTTP) WritePointsContext(ctx context.Context, request WriteRequest) (destination string, err error)
WritePoints writes points over HTTP transport.
type PointsWriter ¶
type PointsWriter interface {
WritePointsContext(ctx context.Context, request WriteRequest) (destination string, err error)
}
PointsWriter is an interface for writing points to a subscription destination. Only WritePoints() needs to be satisfied. PointsWriter implementations must be goroutine safe.
type Service ¶
type Service struct {
MetaClient interface {
Databases() []meta.DatabaseInfo
WaitForDataChanged() chan struct{}
}
NewPointsWriter func(u url.URL) (PointsWriter, error)
Logger *zap.Logger
// contains filtered or unexported fields
}
Service manages forking the incoming data from InfluxDB to defined third party destinations. Subscriptions are defined per database and retention policy.
func NewService ¶
func NewService(c Config, certMonitor *tlsconfig.TLSCertMonitor) *Service
NewService returns a subscriber service with given settings
func (*Service) Close ¶
Close terminates the subscription service. It will return an error if Open was not called first.
func (*Service) PrepareReloadTLSCertificates ¶ added in v1.13.0
PrepareReloadTLSCertificates verifies that the client TLS settings in conf can be loaded. On success it returns a function that applies them; subsequent HTTPS connections then use the new settings, including from writers that already exist, because writers dial through the TLS manager. Connections that are already established are left alone.
A service with TLS disabled has no TLS manager, so there is nothing to reload and the returned function is a no-op.
func (*Service) Send ¶ added in v1.9.3
func (s *Service) Send(request *coordinator.WritePointsRequest)
func (*Service) Statistics ¶ added in v1.0.0
Statistics returns statistics for periodic monitoring.
func (*Service) WithLogger ¶ added in v1.2.0
WithLogger sets the logger on the service.
type Statistics ¶ added in v1.0.0
Statistics maintains the statistics for the subscriber service.
type UDP ¶
type UDP struct {
// contains filtered or unexported fields
}
UDP supports writing points over UDP using the line protocol.
func (*UDP) WritePointsContext ¶ added in v1.9.3
func (u *UDP) WritePointsContext(_ context.Context, request WriteRequest) (destination string, err error)
WritePoints writes points over UDP transport.
type WriteRequest ¶ added in v1.9.3
type WriteRequest struct {
Database string
RetentionPolicy string
// contains filtered or unexported fields
}
WriteRequest is a parsed write request.
func NewWriteRequest ¶ added in v1.9.3
func NewWriteRequest(r *coordinator.WritePointsRequest, log *zap.Logger) (wr WriteRequest, numInvalid int64)
func (*WriteRequest) Length ¶ added in v1.9.3
func (w *WriteRequest) Length() int
func (*WriteRequest) PointAt ¶ added in v1.9.3
func (w *WriteRequest) PointAt(i int) []byte
PointAt uses pointOffsets to slice the lineProtocol buffer and retrieve the i_th point in the request. It includes the trailing newline.
func (*WriteRequest) SizeOf ¶ added in v1.9.3
func (w *WriteRequest) SizeOf() int