Writing a connector
Connectors, the blocks and sources they carry, and the editor schema.
A connector owns a resource (a server, a connection pool, a client) and everything a flow reaches it through. Adding one is how you teach Octo to talk to something new.
runtime/connectors/logger is the smallest complete example, and the one this page follows.
The shape
A connector package registers three or four things from init, and a binary picks it up by blank import. Nothing scans, nothing is configured by path.
| Register | For |
|---|---|
core.MustRegisterConnector | the connector itself |
core.MustRegisterBlock | each block it contributes |
core.RegisterConnectorMeta / RegisterBlockMeta | the editor schema |
core.RegisterExtension | package-level palette defaults (group, icon) |
Sources are the exception: they are not a registry. A connector implements the optional core.SourceProvider interface and builds its own, so a source closes over the connector's live connections instead of reaching for a global.
The connector
core.Connector is two methods. Start acquires resources and must not block; Stop releases them.
func init() {
core.MustRegisterConnector("logger", func() core.Connector {
return &Connector{}
})
// Package-level editor defaults: everything in this package lands in the Data
// palette group with the ScrollText icon unless it says otherwise.
core.RegisterExtension(core.ExtensionMeta{Group: "Data", Icon: "ScrollText"})
core.RegisterConnectorMeta(core.ConnectorMeta{
Type: "logger",
Label: "Logger",
Settings: reflect.TypeFor[connectorSettings](),
})
}
type Connector struct {
logger *slog.Logger
file *os.File
}
func (c *Connector) Start(ctx context.Context, config types.ConnectorConfig) error {
var set connectorSettings
if err := config.Settings.Decode(&set); err != nil {
return err
}
// ... open the file, build the logger
return nil
}
func (c *Connector) Stop(context.Context) error {
// ... close what Start opened
return nil
}Settings are a typed struct decoded from the YAML, not a map. Give every field a default so the connector can be declared bare: - { name: audit, type: logger } is a valid declaration because every setting on logger has one.
Some connectors are resolved implicitly by source type and are never declared in YAML at all. queue and cron work this way. If your connector has no settings worth configuring, do not make users declare it; let sources resolve it on demand.
Blocks
A block is a core.MessageProcessor: it receives the current message and returns the next one. Returning nil drops the message and skips the rest of the chain; returning an error aborts it.
func init() {
core.MustRegisterBlock("log", newLog)
core.RegisterBlockMeta(core.BlockMeta{
Type: "log",
Label: "Log",
Category: core.CategoryProcessor,
Description: "Wire-tap that logs the message and passes it through unchanged.",
Config: reflect.TypeFor[logSettings](),
})
}Blocks live in the connector's package, so importing the connector registers its blocks too. A block binds to a named connector instance by looking it up, which is what lets the log block write through a specific logger connector's output.
Blocks with sub-flows
A block may run nested chains the flow author writes inline, the way if runs its then or an ai-agent runs a tool. Nothing about that is reserved for the runtime's own packages. The sub-flows are fields of the block's settings struct, and the block builds them through the seam every factory is handed.
type gateSettings struct {
// CEL boolean expression.
Condition string `json:"condition" octo:"label=Condition,type=cel,required"`
// Chain run when the condition holds.
Then *types.FlowConfig `json:"then" octo:"label=Then,type=flow,required"`
// Chain run when it does not; optional.
Else *types.FlowConfig `json:"else" octo:"label=Else,type=flow"`
}
func newGate(raw types.Settings, deps core.BlockDeps) (core.MessageProcessor, error) {
var cfg gateSettings
if err := raw.DecodeStrict(&cfg); err != nil { // a misspelled slot is a build error
return nil, err
}
flows, err := core.SubFlowsOf(deps) // nil outside the engine; fail here, not later
if err != nil {
return nil, err
}
then, err := flows.Branch(core.BranchThen, *cfg.Then)
if err != nil {
return nil, err
}
// …compile the condition, build else the same way, return the processor.
}Four rules keep it honest. The branch name must be the field's json name: flows.Branch("then", …) addresses the chain as <block>[then], and the debug address resolver reads that same name off your schema, so --break-at gate.check[then].charge works with nothing else told. A list of chains uses flows.Member(name, index, …), addressed by each member's own name or its index, and a type=flow-list / case-list / route-list / tool-list field in the schema. flows.Root() says whether the block sits in a flow's own chain; a block that continues the flow itself (one implementing core.ContinuationAware) needs a root chain and should refuse to build elsewhere. deps.Scheduler is the flow's worker pool, for a block that scatters work the way fork does.
In the flow file the slots are written as top-level keys beside type, or under settings; the runtime folds the two together, so both spellings decode the same way.
Sources
A source is a flow's entry point. Implement core.SourceProvider on the connector:
func (c *Connector) NewSource(
cfg types.SourceConfig, out chan<- *types.Message, deps core.SourceDeps,
) (core.MessageSource, error) {
// cfg.Type selects which source to build when a connector offers several.
}The runtime owns the channel: Start must not send after Stop returns, and Start must not block. Acquire what you need and do the work on your own goroutine.
Declare the sources in the connector's meta so the editor knows they exist:
core.RegisterConnectorMeta(core.ConnectorMeta{
Type: "http",
Label: "HTTP",
Sources: []core.SourceMeta{{ /* ... */ }},
})The editor schema comes from your Go
octo schema generates the whole editor palette from the registered metadata and the octo: struct tags on your settings types. A settings struct with no tags is a block the editor cannot render.
type logSettings struct {
// CEL expression to log. Defaults to the JSON body.
Message string `json:"message" octo:"label=Message,type=cel"`
// Log level.
Level string `json:"level" octo:"label=Level,type=enum,enum=debug|info|warn|error,default=info"`
// Optional named logger connector to write to.
Logger string `json:"logger" octo:"label=Logger,ref=connector:logger"`
// Attach the whole message as structured log attributes.
Full bool `json:"full" octo:"label=Full message,default=false"`
}ref=connector:logger is what turns a field into a dropdown of the config's logger connectors.
Ship it
- Blank-import the package in
runtime/octo/main.go. - Document every type you registered. Add
apps/docs/content/docs/reference/connectors/<name>.mdxwith anocto_typesfrontmatter list naming your connector, its blocks and its sources. - Run
task docs:check.
Step 2 is not optional. scripts/check-docs-drift.mjs compares octo schema against the frontmatter and fails CI on any type that is registered but undocumented, or documented but not registered.
See also
- Runtime services, for capability that is not referenced from YAML.
- Connectors and blocks, the same model from a flow author's side.
- Processing pipeline, where a block sits in a message's life.