iop/apps/node/internal/node/node.go
toki b4c6550eab feat: edge/node architecture updates and agent-task integration
- Add node store implementation for edge app
- Add adapters factory for node app
- Update edge and node transport layers
- Update domain rules for edge and node
- Add bin scripts for edge and node
- Update configs and documentation
- Add agent-task node_centralized_mgmt directory
2026-05-03 10:51:29 +09:00

203 lines
5.4 KiB
Go

// Package node is the core IOP Node service. It implements
// transport.Handler and orchestrates routing → adapter execution.
package node
import (
"context"
"encoding/json"
"fmt"
"io"
"os"
"strings"
"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 "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(
nodeID string,
router runtime.Router,
registry *adapters.Registry,
st *store.Store,
logger *zap.Logger,
) *Node {
return &Node{
nodeID: nodeID,
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(),
}
printEdgeMessage(os.Stdout, rr.Input)
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)
}
if err := n.store.InsertRun(ctx, store.RunRecord{
RunID: spec.RunID,
Adapter: spec.Adapter,
Model: spec.Model,
Status: "running",
CreatedAt: time.Now(),
}); err != nil {
n.logger.Warn("store: insert run", zap.String("run_id", spec.RunID), zap.Error(err))
}
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, out: os.Stdout}
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))
}
if err := n.store.CompleteRun(context.Background(), spec.RunID, status, errMsg); err != nil {
n.logger.Warn("store: complete run", zap.String("run_id", spec.RunID), zap.Error(err))
}
return execErr
}
// 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
out io.Writer
response strings.Builder
}
func (s *sessionSink) Emit(_ context.Context, event runtime.RuntimeEvent) error {
s.printEvent(event)
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 (s *sessionSink) printEvent(event runtime.RuntimeEvent) {
if s.out == nil {
return
}
switch event.Type {
case runtime.EventTypeStart:
fmt.Fprintf(s.out, "[node-event] start run_id=%s\n", event.RunID)
case runtime.EventTypeDelta:
s.response.WriteString(event.Delta)
case runtime.EventTypeComplete:
fmt.Fprintf(s.out, "[node-event] complete run_id=%s detail=%q\n", event.RunID, event.Message)
printTaggedMessage(s.out, "node-message", s.response.String())
case runtime.EventTypeError:
fmt.Fprintf(s.out, "[node-event] error run_id=%s detail=%q\n", event.RunID, event.Error)
printTaggedMessage(s.out, "node-message", s.response.String())
default:
fmt.Fprintf(s.out, "[node-event] %s run_id=%s detail=%q\n", event.Type, event.RunID, event.Message)
}
}
func printEdgeMessage(out io.Writer, input map[string]any) {
if out == nil {
return
}
if input == nil {
printTaggedMessage(out, "edge-message", "")
return
}
if prompt, ok := input["prompt"].(string); ok {
printTaggedMessage(out, "edge-message", prompt)
return
}
b, err := json.MarshalIndent(input, "", " ")
if err != nil {
printTaggedMessage(out, "edge-message", fmt.Sprintf("%v", input))
return
}
fmt.Fprintf(out, "[edge-message]\n%s\n", b)
}
func printTaggedMessage(out io.Writer, tag, message string) {
message = strings.TrimSpace(message)
if message == "" {
fmt.Fprintf(out, "[%s] <empty>\n", tag)
return
}
fmt.Fprintf(out, "[%s] %s\n", tag, message)
}
func structAsMap(s *structpb.Struct) map[string]any {
if s == nil {
return nil
}
return s.AsMap()
}