yao/trace/node.go
Max 47a35b7e63 Refactor trace manager for improved state management and concurrency
- Replaced the previous node and space management with a channel-based state management system in `manager.go`, enhancing concurrency handling.
- Implemented methods for loading and saving trace updates to disk in `local/driver.go` and `store/driver.go`, allowing for persistent state across sessions.
- Updated subscription handling in `subscription.go` to streamline the process of broadcasting updates to active subscribers.
- Enhanced node and trace status management, including cancellation and completion states, to provide better control over trace execution.
- Adjusted related tests to ensure compatibility with the new state management approach.
2025-11-18 18:08:20 +08:00

266 lines
7.1 KiB
Go

package trace
import (
"fmt"
"time"
"github.com/yaoapp/yao/trace/types"
)
// node implements the Node interface for custom node operations
type node struct {
manager *manager
data *types.TraceNode
}
// Info logs info message (public method, broadcasts event)
func (n *node) Info(format string, args ...any) types.Node {
n.logWithBroadcast("info", format, args...)
return n
}
// Debug logs debug message (public method, broadcasts event)
func (n *node) Debug(format string, args ...any) types.Node {
n.logWithBroadcast("debug", format, args...)
return n
}
// Error logs error message (public method, broadcasts event)
func (n *node) Error(format string, args ...any) types.Node {
n.logWithBroadcast("error", format, args...)
return n
}
// Warn logs warning message (public method, broadcasts event)
func (n *node) Warn(format string, args ...any) types.Node {
n.logWithBroadcast("warn", format, args...)
return n
}
// logWithBroadcast logs and broadcasts event (for external calls)
func (n *node) logWithBroadcast(level string, format string, args ...any) {
log := n.log(level, format, args...)
// Broadcast event
n.manager.addUpdateAndBroadcast(&types.TraceUpdate{
Type: types.UpdateTypeLogAdded,
TraceID: n.manager.traceID,
NodeID: n.data.ID,
Timestamp: log.Timestamp,
Data: log,
})
}
// log logs without broadcasting (for internal Manager calls)
func (n *node) log(level string, format string, args ...any) *types.TraceLog {
message := fmt.Sprintf(format, args...)
log := &types.TraceLog{
Timestamp: time.Now().Unix(),
Level: level,
Message: message,
NodeID: n.data.ID,
}
// Save log (ignore errors for non-critical logging)
_ = n.manager.driver.SaveLog(n.manager.ctx, n.manager.traceID, log)
return log
}
// Add creates next sequential node
func (n *node) Add(input types.TraceInput, option types.TraceNodeOption) (types.Node, error) {
now := time.Now().Unix()
// Create child node data
childNodeData := &types.TraceNode{
ID: genNodeID(),
ParentID: n.data.ID,
Children: []*types.TraceNode{},
TraceNodeOption: option,
Status: types.StatusRunning,
Input: input,
CreatedAt: now,
StartTime: now,
UpdatedAt: now,
}
// Add to parent's children
n.data.Children = append(n.data.Children, childNodeData)
// Save both nodes
if err := n.manager.driver.SaveNode(n.manager.ctx, n.manager.traceID, childNodeData); err != nil {
return nil, err
}
if err := n.manager.driver.SaveNode(n.manager.ctx, n.manager.traceID, n.data); err != nil {
return nil, err
}
// Return Node interface
return &node{
manager: n.manager,
data: childNodeData,
}, nil
}
// Parallel creates multiple concurrent child nodes
func (n *node) Parallel(parallelInputs []types.TraceParallelInput) ([]types.Node, error) {
now := time.Now().Unix()
nodeInterfaces := make([]types.Node, 0, len(parallelInputs))
// Create multiple child nodes
for _, input := range parallelInputs {
childNodeData := &types.TraceNode{
ID: genNodeID(),
ParentID: n.data.ID,
Children: []*types.TraceNode{},
TraceNodeOption: input.Option,
Status: types.StatusRunning,
Input: input.Input,
CreatedAt: now,
StartTime: now,
UpdatedAt: now,
}
n.data.Children = append(n.data.Children, childNodeData)
// Save node
if err := n.manager.driver.SaveNode(n.manager.ctx, n.manager.traceID, childNodeData); err != nil {
return nil, err
}
// Create Node interface wrapper
nodeInterfaces = append(nodeInterfaces, &node{
manager: n.manager,
data: childNodeData,
})
}
// Save parent node
if err := n.manager.driver.SaveNode(n.manager.ctx, n.manager.traceID, n.data); err != nil {
return nil, err
}
return nodeInterfaces, nil
}
// Join joins multiple nodes into one
func (n *node) Join(nodes []*types.TraceNode, input types.TraceInput, option types.TraceNodeOption) (types.Node, error) {
now := time.Now().Unix()
// Create join node data
joinNodeData := &types.TraceNode{
ID: genNodeID(),
ParentID: n.data.ID,
Children: []*types.TraceNode{},
TraceNodeOption: option,
Status: types.StatusRunning,
Input: input,
CreatedAt: now,
StartTime: now,
UpdatedAt: now,
}
// Save join node
if err := n.manager.driver.SaveNode(n.manager.ctx, n.manager.traceID, joinNodeData); err != nil {
return nil, err
}
// Return Node interface
return &node{
manager: n.manager,
data: joinNodeData,
}, nil
}
// ID returns the node ID
func (n *node) ID() string {
return n.data.ID
}
// SetOutput sets the node output
func (n *node) SetOutput(output types.TraceOutput) error {
n.data.Output = output
n.data.UpdatedAt = time.Now().Unix()
return n.manager.driver.SaveNode(n.manager.ctx, n.manager.traceID, n.data)
}
// SetMetadata sets node metadata
func (n *node) SetMetadata(key string, value any) error {
if n.data.Metadata == nil {
n.data.Metadata = make(map[string]any)
}
n.data.Metadata[key] = value
n.data.UpdatedAt = time.Now().Unix()
return n.manager.driver.SaveNode(n.manager.ctx, n.manager.traceID, n.data)
}
// SetStatus sets the node status
func (n *node) SetStatus(status string) error {
n.data.Status = types.NodeStatus(status)
n.data.UpdatedAt = time.Now().Unix()
return n.manager.driver.SaveNode(n.manager.ctx, n.manager.traceID, n.data)
}
// Complete marks the node as completed (public method, broadcasts event)
// Optional output parameter: if provided, sets the output before completing
func (n *node) Complete(output ...types.TraceOutput) error {
if err := n.complete(output...); err != nil {
return err
}
// Broadcast event
n.manager.addUpdateAndBroadcast(&types.TraceUpdate{
Type: types.UpdateTypeNodeComplete,
TraceID: n.manager.traceID,
NodeID: n.data.ID,
Timestamp: n.data.EndTime,
Data: n.data.ToCompleteData(),
})
return nil
}
// complete marks as completed without broadcasting (for Manager calls)
func (n *node) complete(output ...types.TraceOutput) error {
now := time.Now().Unix()
// Set output if provided
if len(output) > 0 {
n.data.Output = output[0]
}
n.data.Status = types.StatusCompleted
n.data.EndTime = now
n.data.UpdatedAt = now
return n.manager.driver.SaveNode(n.manager.ctx, n.manager.traceID, n.data)
}
// Fail marks the node as failed (public method, broadcasts event)
func (n *node) Fail(err error) error {
// Log error first
n.Error("Node failed: %v", err)
if saveErr := n.fail(err); saveErr != nil {
return saveErr
}
// Broadcast event
n.manager.addUpdateAndBroadcast(&types.TraceUpdate{
Type: types.UpdateTypeNodeFailed,
TraceID: n.manager.traceID,
NodeID: n.data.ID,
Timestamp: n.data.EndTime,
Data: n.data.ToFailedData(err),
})
return nil
}
// fail marks as failed without broadcasting (for Manager calls)
func (n *node) fail(err error) error {
now := time.Now().Unix()
// Update status
n.data.Status = types.StatusFailed
n.data.EndTime = now
n.data.UpdatedAt = now
return n.manager.driver.SaveNode(n.manager.ctx, n.manager.traceID, n.data)
}