manager

package
v3.65.6 Latest Latest
Warning

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

Go to latest
Published: Aug 10, 2022 License: MIT Imports: 28 Imported by: 0

Documentation

Overview

Package manager implements the types.Manager interface used for creating and sharing resources across a Benthos service.

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func AddExamples

func AddExamples(c *Config)

AddExamples inserts example caches and conditions if none exist in the config.

func DocumentPlugin

func DocumentPlugin(
	typeString, description string,
	configSanitiser PluginConfigSanitiser,
)

DocumentPlugin adds a description and an optional configuration sanitiser function to the definition of a registered plugin. This improves the documentation generated by PluginDescriptions.

func PluginDescriptions

func PluginDescriptions() string

PluginDescriptions generates and returns a markdown formatted document listing each registered plugin and an example configuration for it.

func RegisterPlugin

func RegisterPlugin(
	typeString string,
	configConstructor PluginConfigConstructor,
	constructor PluginConstructor,
)

RegisterPlugin registers a plugin by a unique name so that it can be constucted similar to regular resources. A constructor for both the plugin itself as well as its configuration struct must be provided.

A constructed resource plugin can be any type and is wrapped as an interface{} type.

func SanitiseConfig

func SanitiseConfig(conf Config) (interface{}, error)

SanitiseConfig creates a sanitised version of a manager config.

func Spec

func Spec() docs.FieldSpecs

Spec returns a field spec for the manager configuration.

func SwapMetrics

func SwapMetrics(mgr types.Manager, stats metrics.Type) types.Manager

SwapMetrics attempts to swap the underlying metrics implementation of a manager. This function does nothing if the manager type is not a *Type.

Types

type APIReg

type APIReg interface {
	RegisterEndpoint(path, desc string, h http.HandlerFunc)
}

APIReg is an interface representing an API builder.

type Config

type Config struct {
	Inputs     map[string]input.Config     `json:"inputs,omitempty" yaml:"inputs,omitempty"`
	Conditions map[string]condition.Config `json:"conditions,omitempty" yaml:"conditions,omitempty"`
	Processors map[string]processor.Config `json:"processors,omitempty" yaml:"processors,omitempty"`
	Outputs    map[string]output.Config    `json:"outputs,omitempty" yaml:"outputs,omitempty"`
	Caches     map[string]cache.Config     `json:"caches,omitempty" yaml:"caches,omitempty"`
	RateLimits map[string]ratelimit.Config `json:"rate_limits,omitempty" yaml:"rate_limits,omitempty"`
	Plugins    map[string]PluginConfig     `json:"plugins,omitempty" yaml:"plugins,omitempty"`
}

Config contains all configuration fields for a Benthos service manager.

func NewConfig

func NewConfig() Config

NewConfig returns a Config with default values.

func (*Config) AddFrom

func (c *Config) AddFrom(extra *Config) error

AddFrom takes another Config and adds all of its resources to itself. If there are any resource name collisions an error is returned.

func (Config) Sanitised

func (c Config) Sanitised(removeDeprecated bool) (interface{}, error)

Sanitised returns a sanitised version of the config, meaning sections that aren't relevant to behaviour are removed. Also optionally removes deprecated fields.

type ErrResourceNotFound

type ErrResourceNotFound string

ErrResourceNotFound represents an error where a named resource could not be accessed because it was not found by the manager.

func (ErrResourceNotFound) Error

func (e ErrResourceNotFound) Error() string

Error implements the standard error interface.

type OptFunc

type OptFunc func(*Type)

OptFunc is an opt setting for a manager type.

func OptSetBloblangEnvironment

func OptSetBloblangEnvironment(env *bloblang.Environment) OptFunc

OptSetBloblangEnvironment determines the environment from which the manager parses bloblang functions and methods. This option is for internal use only.

func OptSetEnvironment

func OptSetEnvironment(e *bundle.Environment) OptFunc

OptSetEnvironment determines the environment from which the manager initializes components and resources. This option is for internal use only.

type PluginConfig

type PluginConfig struct {
	Type   string      `json:"type" yaml:"type"`
	Plugin interface{} `json:"plugin" yaml:"plugin"`
}

PluginConfig is a config struct representing a resource plugin.

func (*PluginConfig) UnmarshalJSON

func (p *PluginConfig) UnmarshalJSON(bytes []byte) error

UnmarshalJSON ensures that when parsing configs that are in a map or slice the default values are still applied.

func (*PluginConfig) UnmarshalYAML

func (p *PluginConfig) UnmarshalYAML(unmarshal func(interface{}) error) error

UnmarshalYAML ensures that when parsing configs that are in a map or slice the default values are still applied.

type PluginConfigConstructor

type PluginConfigConstructor func() interface{}

PluginConfigConstructor is a func that returns a pointer to a new and fully populated configuration struct for a plugin type. It is valid to return a pointer to an empty struct (&struct{}{}) if no configuration fields are needed.

type PluginConfigSanitiser

type PluginConfigSanitiser func(conf interface{}) interface{}

PluginConfigSanitiser is a function that takes a configuration object for a plugin and returns a sanitised (minimal) version of it for printing in examples and plugin documentation.

This function is useful for when a plugins configuration struct is very large and complex, but can sometimes be expressed in a more concise way without losing the original intent.

type PluginConstructor

type PluginConstructor func(
	config interface{},
	manager types.Manager,
	logger log.Modular,
	metrics metrics.Type,
) (interface{}, error)

PluginConstructor is a func that constructs a Benthos resource plugin. These are shareable resources that can be accessed by any other type within Benthos.

The configuration object will be the result of the PluginConfigConstructor after overlaying the user configuration.

type ResourceConfig

type ResourceConfig struct {
	// Called manager for backwards compatibility.
	Manager            Config             `json:"resources,omitempty" yaml:"resources,omitempty"`
	ResourceInputs     []input.Config     `json:"input_resources,omitempty" yaml:"input_resources,omitempty"`
	ResourceProcessors []processor.Config `json:"processor_resources,omitempty" yaml:"processor_resources,omitempty"`
	ResourceOutputs    []output.Config    `json:"output_resources,omitempty" yaml:"output_resources,omitempty"`
	ResourceCaches     []cache.Config     `json:"cache_resources,omitempty" yaml:"cache_resources,omitempty"`
	ResourceRateLimits []ratelimit.Config `json:"rate_limit_resources,omitempty" yaml:"rate_limit_resources,omitempty"`
}

ResourceConfig contains fields for specifying resource components at the root of a Benthos config.

func NewResourceConfig

func NewResourceConfig() ResourceConfig

NewResourceConfig creates a ResourceConfig with default values.

func (*ResourceConfig) AddFrom

func (r *ResourceConfig) AddFrom(extra *ResourceConfig) error

AddFrom takes another Config and adds all of its resources to itself. If there are any resource name collisions an error is returned.

type Type

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

Type is an implementation of types.Manager, which is expected by Benthos components that need to register service wide behaviours such as HTTP endpoints and event listeners, and obtain service wide shared resources such as caches and other resources.

func New

func New(conf Config, apiReg APIReg, log log.Modular, stats metrics.Type) (*Type, error)

New returns an instance of manager.Type, which can be shared amongst components and logical threads of a Benthos service.

func NewV2

func NewV2(conf ResourceConfig, apiReg APIReg, log log.Modular, stats metrics.Type, opts ...OptFunc) (*Type, error)

NewV2 returns an instance of manager.Type, which can be shared amongst components and logical threads of a Benthos service.

func (*Type) AccessCache

func (t *Type) AccessCache(ctx context.Context, name string, fn func(types.Cache)) error

AccessCache attempts to access a cache resource by a unique identifier and executes a closure function with the cache as an argument. Returns an error if the cache does not exist (or is otherwise inaccessible).

During the execution of the provided closure it is guaranteed that the resource will not be closed or removed. However, it is possible for the resource to be accessed by any number of components in parallel.

func (*Type) AccessInput

func (t *Type) AccessInput(ctx context.Context, name string, fn func(types.Input)) error

AccessInput attempts to access an input resource by a unique identifier and executes a closure function with the input as an argument. Returns an error if the input does not exist (or is otherwise inaccessible).

During the execution of the provided closure it is guaranteed that the resource will not be closed or removed. However, it is possible for the resource to be accessed by any number of components in parallel.

func (*Type) AccessOutput

func (t *Type) AccessOutput(ctx context.Context, name string, fn func(types.OutputWriter)) error

AccessOutput attempts to access an output resource by a unique identifier and executes a closure function with the output as an argument. Returns an error if the output does not exist (or is otherwise inaccessible).

During the execution of the provided closure it is guaranteed that the resource will not be closed or removed. However, it is possible for the resource to be accessed by any number of components in parallel.

func (*Type) AccessProcessor

func (t *Type) AccessProcessor(ctx context.Context, name string, fn func(types.Processor)) error

AccessProcessor attempts to access a processor resource by a unique identifier and executes a closure function with the processor as an argument. Returns an error if the processor does not exist (or is otherwise inaccessible).

During the execution of the provided closure it is guaranteed that the resource will not be closed or removed. However, it is possible for the resource to be accessed by any number of components in parallel.

func (*Type) AccessRateLimit

func (t *Type) AccessRateLimit(ctx context.Context, name string, fn func(types.RateLimit)) error

AccessRateLimit attempts to access a rate limit resource by a unique identifier and executes a closure function with the rate limit as an argument. Returns an error if the rate limit does not exist (or is otherwise inaccessible).

During the execution of the provided closure it is guaranteed that the resource will not be closed or removed. However, it is possible for the resource to be accessed by any number of components in parallel.

func (*Type) BloblEnvironment

func (t *Type) BloblEnvironment() *bloblang.Environment

BloblEnvironment returns a Bloblang environment used by the manager. This is for internal use only.

func (*Type) CloseAsync

func (t *Type) CloseAsync()

CloseAsync triggers the shut down of all resource types that implement the lifetime interface types.Closable.

func (*Type) Environment

func (t *Type) Environment() *bundle.Environment

Environment returns a bundle environment used by the manager. This is for internal use only.

func (*Type) ForChildComponent

func (t *Type) ForChildComponent(id string) types.Manager

ForChildComponent returns a variant of this manager to be used by a particular component identifer, which is a child of the current component, where observability components will be automatically tagged with the label.

func (*Type) ForComponent

func (t *Type) ForComponent(id string) types.Manager

ForComponent returns a variant of this manager to be used by a particular component identifer, where observability components will be automatically tagged with the label.

func (*Type) ForStream

func (t *Type) ForStream(id string) types.Manager

ForStream returns a variant of this manager to be used by a particular stream identifer, where APIs registered will be namespaced by that id.

func (*Type) GetCache

func (t *Type) GetCache(name string) (types.Cache, error)

GetCache attempts to find a service wide cache by its name.

func (*Type) GetCondition

func (t *Type) GetCondition(name string) (types.Condition, error)

GetCondition attempts to find a service wide condition by its name.

func (*Type) GetDocs

func (t *Type) GetDocs(name string, ctype docs.Type) (docs.ComponentSpec, bool)

GetDocs returns a documentation spec for an implementation of a component.

func (*Type) GetInput

func (t *Type) GetInput(name string) (types.Input, error)

GetInput attempts to find a service wide input by its name.

func (*Type) GetOutput

func (t *Type) GetOutput(name string) (types.OutputWriter, error)

GetOutput attempts to find a service wide output by its name.

func (*Type) GetPipe

func (t *Type) GetPipe(name string) (<-chan types.Transaction, error)

GetPipe attempts to obtain and return a named output Pipe

func (*Type) GetPlugin

func (t *Type) GetPlugin(name string) (interface{}, error)

GetPlugin attempts to find a service wide resource plugin by its name.

func (*Type) GetProcessor

func (t *Type) GetProcessor(name string) (types.Processor, error)

GetProcessor attempts to find a service wide processor by its name.

func (*Type) GetRateLimit

func (t *Type) GetRateLimit(name string) (types.RateLimit, error)

GetRateLimit attempts to find a service wide rate limit by its name.

func (*Type) Label

func (t *Type) Label() string

Label returns the current component label held by a manager.

func (*Type) Logger

func (t *Type) Logger() log.Modular

Logger returns a logger preset with the current component context.

func (*Type) Metrics

func (t *Type) Metrics() metrics.Type

Metrics returns an aggregator preset with the current component context.

func (*Type) NewBuffer

func (t *Type) NewBuffer(conf buffer.Config) (buffer.Type, error)

NewBuffer attempts to create a new buffer component from a config.

func (*Type) NewCache

func (t *Type) NewCache(conf cache.Config) (types.Cache, error)

NewCache attempts to create a new cache component from a config.

func (*Type) NewInput

func (t *Type) NewInput(conf input.Config, hasBatchProc bool, pipelines ...types.PipelineConstructorFunc) (types.Input, error)

NewInput attempts to create a new input component from a config.

TODO: V4 Remove the dumb batch field.

func (*Type) NewOutput

func (t *Type) NewOutput(conf output.Config, pipelines ...types.PipelineConstructorFunc) (types.Output, error)

NewOutput attempts to create a new output component from a config.

func (*Type) NewProcessor

func (t *Type) NewProcessor(conf processor.Config) (types.Processor, error)

NewProcessor attempts to create a new processor component from a config.

func (*Type) NewRateLimit

func (t *Type) NewRateLimit(conf ratelimit.Config) (types.RateLimit, error)

NewRateLimit attempts to create a new rate limit component from a config.

func (*Type) RegisterEndpoint

func (t *Type) RegisterEndpoint(apiPath, desc string, h http.HandlerFunc)

RegisterEndpoint registers a server wide HTTP endpoint.

func (*Type) SetPipe

func (t *Type) SetPipe(name string, tran <-chan types.Transaction)

SetPipe registers a new transaction chan to a named pipe.

func (*Type) StoreCache

func (t *Type) StoreCache(ctx context.Context, name string, conf cache.Config) error

StoreCache attempts to store a new cache resource. If an existing resource has the same name it is closed and removed _before_ the new one is initialized in order to avoid duplicate connections.

func (*Type) StoreInput

func (t *Type) StoreInput(ctx context.Context, name string, conf input.Config) error

StoreInput attempts to store a new input resource. If an existing resource has the same name it is closed and removed _before_ the new one is initialized in order to avoid duplicate connections.

func (*Type) StoreOutput

func (t *Type) StoreOutput(ctx context.Context, name string, conf output.Config) error

StoreOutput attempts to store a new output resource. If an existing resource has the same name it is closed and removed _before_ the new one is initialized in order to avoid duplicate connections.

func (*Type) StoreProcessor

func (t *Type) StoreProcessor(ctx context.Context, name string, conf processor.Config) error

StoreProcessor attempts to store a new processor resource. If an existing resource has the same name it is closed and removed _before_ the new one is initialized in order to avoid duplicate connections.

func (*Type) StoreRateLimit

func (t *Type) StoreRateLimit(ctx context.Context, name string, conf ratelimit.Config) error

StoreRateLimit attempts to store a new rate limit resource. If an existing resource has the same name it is closed and removed _before_ the new one is initialized in order to avoid duplicate connections.

func (*Type) UnsetPipe

func (t *Type) UnsetPipe(name string, tran <-chan types.Transaction)

UnsetPipe removes a named pipe transaction chan.

func (*Type) WaitForClose

func (t *Type) WaitForClose(timeout time.Duration) error

WaitForClose blocks until either all closable resource types are shut down or a timeout occurs.

func (*Type) WithMetricsMapping

func (t *Type) WithMetricsMapping(m *imetrics.Mapping) *Type

WithMetricsMapping returns a manager with the stored metrics exporter wrapped with a mapping.

Jump to

Keyboard shortcuts

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