Documentation
¶
Overview ¶
Package Flow provides a revolutionary workflow orchestration library for Go.
Flow uses a single adaptive node that automatically changes behavior based on parameters, eliminating boilerplate while enabling unprecedented composability for building AI agents, complex workflows, and data processing pipelines.
Key Features:
- Single Adaptive Node: One node type that automatically adapts behavior based on parameters
- Parameter-Driven: Configure behavior through parameters, not inheritance
- Auto-Parallel: Add `parallel: true` to any batch operation for instant concurrency
- Auto-Retry: Set `retries > 0` to automatically enable retry logic with exponential backoff
- Auto-Batch: Set `batch: true` with `data` to automatically process collections
- Composable: Mix retry + batch + parallel in a single node declaration
- Thread-Safe: SharedState management for safe concurrent data sharing
Basic Usage:
state := flow.NewSharedState()
node := flow.NewNode()
node.SetParams(map[string]interface{}{
"name": "World",
})
node.SetExecFunc(func(prep interface{}) (interface{}, error) {
name := node.GetParam("name").(string)
fmt.Printf("Hello, %s!\n", name)
return "greeted", nil
})
result := node.Run(state)
Index ¶
- Constants
- type Flow
- type Node
- func (n *Node) GetParam(key string) interface{}
- func (n *Node) GetSuccessors() map[string]*Node
- func (n *Node) Next(node *Node, action string) *Node
- func (n *Node) Run(shared *SharedState) string
- func (n *Node) SetExecFunc(fn func(interface{}) (interface{}, error))
- func (n *Node) SetParams(params map[string]interface{})
- func (n *Node) SetPostFunc(fn func(*SharedState, interface{}, interface{}) string)
- func (n *Node) SetPrepFunc(fn func(*SharedState) interface{})
- type SharedState
Constants ¶
const (
// BatchCompleteAction represents the action returned when batch processing is complete
BatchCompleteAction = "batch_complete"
)
const (
// DefaultAction represents the default action when no specific action is provided
DefaultAction = "default"
)
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type Flow ¶
type Flow struct {
*Node
// contains filtered or unexported fields
}
Flow orchestrates the execution of connected nodes in a workflow. It provides sequential traversal and action-based routing between nodes.
A Flow embeds a Node and maintains a reference to the starting node. It executes nodes in sequence, following the connections defined by Next() calls based on the action strings returned by each node's execution.
func NewFlow ¶
func NewFlow() *Flow
NewFlow creates a new Flow instance. The returned Flow can be used to orchestrate node execution by setting a start node and calling Run().
Example:
flow := NewFlow().Start(firstNode) result := flow.Run(sharedState)
func (*Flow) Run ¶
func (f *Flow) Run(shared *SharedState) string
Run executes the flow starting from the start node (like PocketFlow's _orch)
type Node ¶
type Node struct {
// contains filtered or unexported fields
}
Node is the core adaptive node that automatically changes behavior based on parameters. This single node type eliminates the need for multiple specialized node types by detecting patterns in its parameters and adapting its execution accordingly.
Supported Adaptive Behaviors:
- Batch Processing: Set "batch": true with "data" to process collections
- Parallel Execution: Set "parallel": true to enable concurrent processing
- Retry Logic: Set "retries" > 0 to enable automatic retry with exponential backoff
- Composability: All patterns can be combined in a single node
Parameter Detection Priority:
- Batch Processing: "batch": true → process each item in "data"
- Retry Logic: "retries" > 0 → wrap execution with exponential backoff retry
- Single Execution: Default behavior
The Node maintains a map of parameters, successor nodes for workflow chaining, and optional user-provided functions for custom prep, exec, and post processing.
func NewNode ¶
func NewNode() *Node
NewNode creates a new adaptive Node with empty parameters and successors. The returned Node can be configured with parameters and functions to define its behavior and then executed with Run().
Example:
node := NewNode()
node.SetParams(map[string]interface{}{"retries": 3})
node.SetExecFunc(func(prep interface{}) (interface{}, error) {
// Your business logic here
return "success", nil
})
result := node.Run(sharedState)
func (*Node) GetParam ¶
GetParam retrieves a parameter value by key. Returns nil if the parameter doesn't exist.
Example:
retries := node.GetParam("retries")
if retries != nil {
retriesInt := retries.(int)
}
func (*Node) GetSuccessors ¶
GetSuccessors returns a map of all successor nodes keyed by their action strings. This is primarily used internally by Flow for traversal.
func (*Node) Next ¶
Next establishes a connection to another node for workflow chaining. The connection is triggered when this node's execution returns the specified action string. If action is empty, "default" is used.
Parameters:
- node: The target Node to connect to
- action: The action string that triggers this connection
Returns:
- *Node: The target node (for method chaining)
Example:
processor.Next(validator, "processed") validator.Next(success, "valid") validator.Next(failure, "invalid")
func (*Node) Run ¶
func (n *Node) Run(shared *SharedState) string
Run executes the node with adaptive behavior based on parameters
func (*Node) SetExecFunc ¶
SetExecFunc sets the user's business logic function
func (*Node) SetParams ¶
SetParams configures the node's parameters that control its adaptive behavior. Parameters determine which execution patterns the node will use:
- "batch": true - enables batch processing of "data" parameter
- "parallel": true - enables parallel execution (requires "batch": true)
- "parallel_limit": int - limits concurrent goroutines (default: 10)
- "retries": int - enables retry logic with exponential backoff
- "retry_delay": time.Duration - base delay for retry backoff
- "data": []interface{} - data to process in batch mode
Example:
node.SetParams(map[string]interface{}{
"data": []int{1, 2, 3, 4, 5},
"batch": true,
"parallel": true,
"retries": 3,
})
func (*Node) SetPostFunc ¶
func (n *Node) SetPostFunc(fn func(*SharedState, interface{}, interface{}) string)
SetPostFunc sets optional post-processing function
func (*Node) SetPrepFunc ¶
func (n *Node) SetPrepFunc(fn func(*SharedState) interface{})
SetPrepFunc sets optional preparation function
type SharedState ¶
type SharedState struct {
// contains filtered or unexported fields
}
SharedState provides thread-safe data sharing between nodes in a workflow. It acts as a central data store that nodes can read from and write to during execution. All operations are protected by a read-write mutex for safe concurrent access.
SharedState is typically created once per workflow execution and passed to all nodes. It supports storing any type of data and provides typed getter methods for convenience.
Example:
state := NewSharedState()
state.Set("user_id", 12345)
state.Set("results", []string{"item1", "item2"})
userID := state.GetInt("user_id")
results := state.GetSlice("results")
func NewSharedState ¶
func NewSharedState() *SharedState
NewSharedState creates a new SharedState instance with an empty data map. The returned SharedState is ready for use and thread-safe.
Example:
state := NewSharedState()
state.Set("key", "value")
func (*SharedState) Append ¶
func (s *SharedState) Append(key string, value interface{})
Append adds an item to a slice in shared state
func (*SharedState) Get ¶
func (s *SharedState) Get(key string) interface{}
Get retrieves a value from the shared state by key. Returns nil if the key doesn't exist. This operation is thread-safe for concurrent reads.
Parameters:
- key: The string key to retrieve
Returns:
- interface{}: The stored value, or nil if key doesn't exist
Example:
value := state.Get("counter")
if value != nil {
counter := value.(int)
}
func (*SharedState) GetInt ¶
func (s *SharedState) GetInt(key string) int
GetInt retrieves an int value, returning 0 if not found or not an int
func (*SharedState) GetSlice ¶
func (s *SharedState) GetSlice(key string) []interface{}
GetSlice retrieves a slice value, returning empty slice if not found
func (*SharedState) Set ¶
func (s *SharedState) Set(key string, value interface{})
Set stores a value in the shared state under the specified key. This operation is thread-safe and will overwrite any existing value for the key.
Parameters:
- key: The string key to store the value under
- value: The value to store (can be any type)
Example:
state.Set("counter", 42)
state.Set("results", []string{"a", "b", "c"})
Directories
¶
| Path | Synopsis |
|---|---|
|
examples
|
|
|
basic-greeting
command
|
|
|
batch-pattern
command
|
|
|
chatbot
command
|
|
|
composed-pattern
command
|
|
|
retry-pattern
command
|
|
|
workflow-pattern
command
|