Documentation
¶
Overview ¶
Package kafkaconnect provides Kafka publish and consume capabilities for workflows. It exposes CRUD routes for managing Kafka configurations, a produce host function for publishing messages, and a background consumer loop that polls for messages and logs them.
Index ¶
- func New() plugin.Plugin
- type Config
- type Plugin
- func (p *Plugin) Info() plugin.PluginInfo
- func (p *Plugin) Init(ctx context.Context, env *plugin.Environment) error
- func (p *Plugin) Migrations() []plugin.Migration
- func (p *Plugin) RegisterHostFunctions(scope plugin.FuncRegistry) error
- func (p *Plugin) RegisterRoutes(mux *http.ServeMux) error
- func (p *Plugin) Run(ctx context.Context) error
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
Types ¶
type Config ¶
type Config struct {
RestProxyURL string `json:"rest_proxy_url,omitempty"` // Confluent REST Proxy URL
}
Config holds optional configuration for the kafka-connect plugin.
type Plugin ¶
type Plugin struct {
// contains filtered or unexported fields
}
Plugin implements Kafka publish and consume integration for workflows.
func (*Plugin) Info ¶
func (p *Plugin) Info() plugin.PluginInfo
Info returns plugin metadata for discovery and documentation.
func (*Plugin) Init ¶
Init initializes the plugin with the given environment. It creates an HTTP client with a 10-second timeout for REST proxy requests.
func (*Plugin) Migrations ¶
Migrations returns the database schema for Kafka config storage. Tables are idempotent (IF NOT EXISTS) and safe to run multiple times.
func (*Plugin) RegisterHostFunctions ¶
func (p *Plugin) RegisterHostFunctions(scope plugin.FuncRegistry) error
RegisterHostFunctions registers workflow-callable functions on the scoped function registry. The plugin name is implicit -- each plugin gets its own scope, so function names need not be globally unique.