iop/apps/node/internal/node/node.go
toki c46874055a feat: edge node unit tests and related updates
- Add edge node unit tests and transport package
- Add node test client and related artifacts
- Update bootstrap, node, and config modules
- Add proto generated files
- Update Makefile and configuration files
2026-05-02 20:09:55 +09:00

161 lines
4.2 KiB
Go

// Package node is the core IOP Node service. It implements
// transport.Handler and orchestrates routing → adapter execution.
package node
import (
"context"
"fmt"
"time"
"go.uber.org/zap"
"google.golang.org/protobuf/types/known/structpb"
"iop/apps/node/internal/adapters"
"iop/apps/node/internal/runtime"
"iop/apps/node/internal/store"
"iop/apps/node/internal/transport"
"iop/packages/config"
iop "iop/proto/gen/iop"
)
// Node implements transport.Handler and coordinates the full execution pipeline.
type Node struct {
nodeID string
router runtime.Router
registry *adapters.Registry
store *store.Store
logger *zap.Logger
}
// New creates a Node. It satisfies transport.Handler.
func New(
cfg *config.NodeConfig,
router runtime.Router,
registry *adapters.Registry,
st *store.Store,
logger *zap.Logger,
) *Node {
return &Node{
nodeID: cfg.Node.ID,
router: router,
registry: registry,
store: st,
logger: logger,
}
}
// OnRunRequest handles an incoming RunRequest from a transport Session.
func (n *Node) OnRunRequest(ctx context.Context, sess *transport.Session, req *iop.RunRequest) error {
n.logger.Info("run request received",
zap.String("run_id", req.GetRunId()),
zap.String("adapter", req.GetAdapter()),
zap.String("model", req.GetModel()),
)
rr := runtime.RunRequest{
RunID: req.GetRunId(),
Adapter: req.GetAdapter(),
Model: req.GetModel(),
Workspace: req.GetWorkspace(),
Policy: structAsMap(req.GetPolicy()),
Input: structAsMap(req.GetInput()),
TimeoutSec: int(req.GetTimeoutSec()),
Metadata: req.GetMetadata(),
}
spec, err := n.router.Resolve(ctx, rr)
if err != nil {
return fmt.Errorf("node: resolve: %w", err)
}
adapter, ok := n.registry.Get(spec.Adapter)
if !ok {
return fmt.Errorf("node: adapter %q not found after routing", spec.Adapter)
}
_ = n.store.InsertRun(ctx, store.RunRecord{
RunID: spec.RunID,
Adapter: spec.Adapter,
Model: spec.Model,
Status: "running",
CreatedAt: time.Now(),
})
execCtx := ctx
if spec.TimeoutSec > 0 {
var cancel context.CancelFunc
execCtx, cancel = context.WithTimeout(ctx, time.Duration(spec.TimeoutSec)*time.Second)
defer cancel()
sess.RegisterCancel(spec.RunID, cancel)
defer sess.DeregisterCancel(spec.RunID)
}
sink := &sessionSink{sess: sess}
execErr := adapter.Execute(execCtx, spec, sink)
status := "completed"
errMsg := ""
if execErr != nil {
status = "failed"
errMsg = execErr.Error()
n.logger.Warn("run failed", zap.String("run_id", spec.RunID), zap.Error(execErr))
}
_ = n.store.CompleteRun(context.Background(), spec.RunID, status, errMsg)
return execErr
}
// OnCapabilityRequest returns the node's available adapters.
func (n *Node) OnCapabilityRequest(_ context.Context, _ *transport.Session) (*iop.CapabilityResponse, error) {
resp := &iop.CapabilityResponse{NodeId: n.nodeID}
for _, a := range n.registry.All() {
caps, err := a.Capabilities(context.Background())
if err != nil {
n.logger.Warn("capabilities query failed", zap.String("adapter", a.Name()), zap.Error(err))
continue
}
resp.Adapters = append(resp.Adapters, &iop.AdapterInfo{
Name: caps.AdapterName,
Models: caps.Models,
MaxConcurrency: int32(caps.MaxConcurrency),
})
}
return resp, nil
}
// OnCancel cancels a running execution.
func (n *Node) OnCancel(_ context.Context, sess *transport.Session, req *iop.CancelRequest) error {
n.logger.Info("cancel request", zap.String("run_id", req.GetRunId()))
sess.CancelRun(req.GetRunId())
return nil
}
// sessionSink wraps a transport.Session to implement runtime.EventSink.
type sessionSink struct {
sess *transport.Session
}
func (s *sessionSink) Emit(_ context.Context, event runtime.RuntimeEvent) error {
re := &iop.RunEvent{
RunId: event.RunID,
Type: string(event.Type),
Delta: event.Delta,
Message: event.Message,
Error: event.Error,
Metadata: event.Metadata,
Timestamp: event.Timestamp.UnixNano(),
}
if event.Usage != nil {
re.Usage = &iop.Usage{
InputTokens: int32(event.Usage.InputTokens),
OutputTokens: int32(event.Usage.OutputTokens),
}
}
return s.sess.Send(re)
}
func structAsMap(s *structpb.Struct) map[string]any {
if s == nil {
return nil
}
return s.AsMap()
}