iop/apps/node/internal/node/node.go
toki 0e02958684 feat(node): inference provider availability probing을 추가한다
Adapter에 ProviderProber 인터페이스를 도입하여 provider endpoint와
target model의 가용성을 확인한다.

- runtime에 ProviderStatus, ProviderProbeResult 타입을 추가한다
- ollama/vllm/mock adapter에 ProbeProvider를 구현한다
- node handleCapabilitiesCommand에서 Prober 인터페이스를 호출하여
  target-aware provider status를 응답한다
- 각 adapter 및 node 테스트를 추가한다
- 로드맵 문서를 갱신한다
2026-06-14 16:39:49 +09:00

652 lines
20 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"
"errors"
"fmt"
"io"
"os"
"sort"
"strconv"
"strings"
"time"
"sync"
"go.uber.org/zap"
"google.golang.org/protobuf/proto"
"google.golang.org/protobuf/types/known/structpb"
"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
store *store.Store
runs *runManager
globalPermit *permitManager // node-wide limit across all adapters
adapterPermitsMu sync.Mutex
adapterPermits map[string]*permitManager // per adapter-key limit
out io.Writer
logger *zap.Logger
}
// New creates a Node. It satisfies transport.Handler.
// globalConcurrency is the node-wide maximum for simultaneous executions across
// all adapters (>= 1 enforces limit; <= 0 means unlimited at the node level).
// Per-adapter MaxConcurrency from Capabilities is applied independently on top.
// out receives console debug output; pass os.Stdout for production, io.Discard in tests.
func New(
nodeID string,
router runtime.Router,
st *store.Store,
globalConcurrency int,
out io.Writer,
logger *zap.Logger,
) *Node {
if out == nil {
out = os.Stdout
}
return &Node{
nodeID: nodeID,
router: router,
store: st,
runs: newRunManager(),
globalPermit: newPermitManager(globalConcurrency),
adapterPermits: make(map[string]*permitManager),
out: out,
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("target", req.GetTarget()),
)
rr := runtime.RunRequest{
RunID: req.GetRunId(),
Adapter: req.GetAdapter(),
Target: req.GetTarget(),
SessionID: req.GetSessionId(),
SessionMode: sessionModeFromProto(req.GetSessionMode()),
Background: req.GetBackground(),
Workspace: req.GetWorkspace(),
Policy: structAsMap(req.GetPolicy()),
Input: structAsMap(req.GetInput()),
TimeoutSec: int(req.GetTimeoutSec()),
Metadata: req.GetMetadata(),
}
printEdgeMessage(n.out, rr.Input)
spec, adapter, err := n.router.ResolveAdapter(ctx, rr)
if err != nil {
return fmt.Errorf("node: resolve: %w", err)
}
// Concurrency gate (foreground and background use the same path).
// Policy: reject-on-exceed (no queue, no wait).
//
// Two independent limits are enforced:
// 1. global (node-wide): n.globalPermit — shared across all adapter keys.
// 2. adapter-key: per-adapter Capabilities.MaxConcurrency permit.
//
// Both must be acquired. If global is acquired but adapter-key fails, global
// is immediately released. The reject limit reported to the caller is the
// stricter of the two.
adapterCap := 0
if caps, capsErr := adapter.Capabilities(ctx); capsErr == nil {
adapterCap = caps.MaxConcurrency
}
adapterPM := n.adapterPermitFor(spec.Adapter, adapterCap)
if !n.globalPermit.tryAcquire() {
limit := int(n.globalPermit.limit)
n.logger.Warn("global concurrency limit exceeded",
zap.String("run_id", spec.RunID),
zap.Int("global_limit", limit),
)
n.rejectRun(ctx, sess, spec, limit)
return fmt.Errorf("node: run %s: %w", spec.RunID, ErrConcurrencyLimitExceeded)
}
if !adapterPM.tryAcquire() {
n.globalPermit.release()
n.logger.Warn("adapter concurrency limit exceeded",
zap.String("run_id", spec.RunID),
zap.String("adapter", spec.Adapter),
zap.Int("adapter_limit", adapterCap),
)
n.rejectRun(ctx, sess, spec, adapterCap)
return fmt.Errorf("node: run %s: %w", spec.RunID, ErrConcurrencyLimitExceeded)
}
if err := n.store.InsertRun(ctx, store.RunRecord{
RunID: spec.RunID,
Adapter: spec.Adapter,
Target: spec.Target,
SessionID: normalizeSessionID(spec.SessionID),
Background: spec.Background,
Status: "running",
CreatedAt: time.Now(),
}); err != nil {
n.logger.Warn("store: insert run", zap.String("run_id", spec.RunID), zap.Error(err))
}
execCtx, cancel := context.WithCancel(ctx)
if spec.TimeoutSec > 0 {
execCtx, cancel = context.WithTimeout(ctx, time.Duration(spec.TimeoutSec)*time.Second)
}
h := &runHandle{
runID: spec.RunID,
adapter: spec.Adapter,
target: spec.Target,
sessionID: normalizeSessionID(spec.SessionID),
cancel: cancel,
done: make(chan struct{}),
}
n.runs.register(h)
sink := &sessionSink{
sess: sess,
out: n.out,
nodeID: n.nodeID,
sessionID: normalizeSessionID(spec.SessionID),
background: spec.Background,
}
run := func() error {
defer n.globalPermit.release()
defer adapterPM.release()
defer cancel()
defer n.runs.deregister(spec.RunID)
defer close(h.done)
execErr := adapter.Execute(execCtx, spec, sink)
n.completeRun(spec, execErr)
return execErr
}
if spec.Background {
go func() { _ = run() }()
return nil
}
return run()
}
// adapterPermitFor returns the shared permitManager for the given adapter
// instance key, creating one with the given limit if not yet seen.
// limit <= 0 means unlimited.
func (n *Node) adapterPermitFor(adapterKey string, limit int) *permitManager {
n.adapterPermitsMu.Lock()
defer n.adapterPermitsMu.Unlock()
pm, ok := n.adapterPermits[adapterKey]
if !ok {
pm = newPermitManager(limit)
n.adapterPermits[adapterKey] = pm
}
return pm
}
// rejectRun records a rejected run in the store and sends an error RunEvent to
// the session so Edge can observe the rejection through the normal event stream.
func (n *Node) rejectRun(ctx context.Context, sess *transport.Session, spec runtime.ExecutionSpec, limit int) {
errMsg := fmt.Sprintf("concurrency limit exceeded (limit=%d)", limit)
if err := n.store.InsertRun(ctx, store.RunRecord{
RunID: spec.RunID,
Adapter: spec.Adapter,
Target: spec.Target,
SessionID: normalizeSessionID(spec.SessionID),
Background: spec.Background,
Status: "rejected",
CreatedAt: time.Now(),
}); err != nil {
n.logger.Warn("store: insert rejected run", zap.String("run_id", spec.RunID), zap.Error(err))
} else if err := n.store.CompleteRun(ctx, spec.RunID, "rejected", errMsg); err != nil {
n.logger.Warn("store: complete rejected run", zap.String("run_id", spec.RunID), zap.Error(err))
}
if sess != nil && sess.IsAlive() {
re := &iop.RunEvent{
RunId: spec.RunID,
Type: string(runtime.EventTypeError),
Error: errMsg,
Timestamp: time.Now().UnixNano(),
SessionId: normalizeSessionID(spec.SessionID),
Background: spec.Background,
NodeId: n.nodeID,
}
if err := sess.Send(re); err != nil {
n.logger.Warn("session: send reject event", zap.String("run_id", spec.RunID), zap.Error(err))
}
}
}
func (n *Node) completeRun(spec runtime.ExecutionSpec, execErr error) {
status := "completed"
errMsg := ""
if execErr != nil {
if errors.Is(execErr, runtime.ErrRunCancelled) {
status = "cancelled"
} else {
status = "failed"
errMsg = execErr.Error()
}
n.logger.Warn("run ended", zap.String("run_id", spec.RunID), zap.String("status", status), 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))
}
}
// OnCancel cancels a running execution or terminates an adapter session.
func (n *Node) OnCancel(_ context.Context, _ *transport.Session, req *iop.CancelRequest) error {
n.logger.Info("cancel request", zap.String("run_id", req.GetRunId()), zap.String("action", req.GetAction().String()))
switch cancelActionFromProto(req.GetAction()) {
case runtime.CancelActionTerminateSession:
adapter, err := n.router.LookupAdapter(req.GetAdapter())
if err != nil {
return fmt.Errorf("node: %w", err)
}
terminator, ok := adapter.(runtime.SessionTerminator)
if !ok {
return fmt.Errorf("node: adapter %q does not support session termination", req.GetAdapter())
}
return terminator.TerminateSession(context.Background(), req.GetTarget(), normalizeSessionID(req.GetSessionId()))
default:
n.runs.cancelRun(req.GetRunId())
return nil
}
}
// OnCommandRequest handles an incoming node command from a transport Session.
func (n *Node) OnCommandRequest(ctx context.Context, sess *transport.Session, req *iop.NodeCommandRequest) (*iop.NodeCommandResponse, error) {
n.logger.Info("command request",
zap.String("request_id", req.GetRequestId()),
zap.String("type", req.GetType().String()),
zap.String("adapter", req.GetAdapter()),
zap.String("target", req.GetTarget()),
)
cmdType, ok := protoCommandTypeToDomain(req.GetType())
if !ok {
return n.commandErrorResponse(req, fmt.Sprintf("node: unsupported command type %q", req.GetType().String())), nil
}
execCtx := ctx
if req.GetTimeoutSec() > 0 {
var cancel context.CancelFunc
execCtx, cancel = context.WithTimeout(ctx, time.Duration(req.GetTimeoutSec())*time.Second)
defer cancel()
}
switch cmdType {
case runtime.CommandTypeCapabilities:
return n.handleCapabilitiesCommand(execCtx, req), nil
case runtime.CommandTypeTransportStatus:
return n.handleTransportStatusCommand(sess, req), nil
default:
return n.dispatchAdapterCommand(execCtx, req), nil
}
}
func (n *Node) handleCapabilitiesCommand(ctx context.Context, req *iop.NodeCommandRequest) *iop.NodeCommandResponse {
adapter, err := n.router.LookupAdapter(req.GetAdapter())
if err != nil {
return n.commandErrorResponse(req, fmt.Sprintf("node: %s", err.Error()))
}
caps, err := adapter.Capabilities(ctx)
if err != nil {
return n.commandErrorResponse(req, err.Error())
}
targets := append([]string(nil), caps.Targets...)
providerStatus := caps.ProviderStatus
providerDetail := ""
if prober, ok := adapter.(runtime.ProviderProber); ok {
probeRes, err := prober.ProbeProvider(ctx, req.GetTarget())
if err != nil {
providerStatus = runtime.ProviderStatusUnavailable
providerDetail = err.Error()
} else {
providerStatus = probeRes.Status
providerDetail = probeRes.Detail
if len(probeRes.Targets) > 0 {
targets = append([]string(nil), probeRes.Targets...)
}
}
}
sort.Strings(targets)
n.adapterPermitsMu.Lock()
pm, ok := n.adapterPermits[req.GetAdapter()]
n.adapterPermitsMu.Unlock()
inFlight := 0
if ok {
inFlight = pm.activeCount()
}
result := map[string]string{
"adapter": caps.AdapterName,
"instance_key": caps.InstanceKey,
"targets": strings.Join(targets, ","),
"max_concurrency": strconv.Itoa(caps.MaxConcurrency),
"provider_status": string(runtime.NormalizeProviderStatus(providerStatus)),
"capacity": strconv.Itoa(caps.MaxConcurrency),
"in_flight": strconv.Itoa(inFlight),
"queued": "0",
}
if providerDetail != "" {
result["provider_detail"] = providerDetail
}
providerSnapshot := &iop.ProviderSnapshot{
Adapter: req.GetAdapter(),
Status: string(runtime.NormalizeProviderStatus(providerStatus)),
Capacity: int32(caps.MaxConcurrency),
InFlight: int32(inFlight),
Queued: 0,
}
return &iop.NodeCommandResponse{
RequestId: req.GetRequestId(),
Type: req.GetType(),
Adapter: req.GetAdapter(),
Target: req.GetTarget(),
SessionId: req.GetSessionId(),
Result: result,
ProviderSnapshots: []*iop.ProviderSnapshot{providerSnapshot},
}
}
func sessionConnected(sess *transport.Session) (alive bool) {
if sess == nil {
return false
}
defer func() {
if r := recover(); r != nil {
alive = false
}
}()
return sess.IsAlive()
}
func (n *Node) handleTransportStatusCommand(sess *transport.Session, req *iop.NodeCommandRequest) *iop.NodeCommandResponse {
connected := "false"
state := "disconnected"
if sessionConnected(sess) {
connected = "true"
state = "connected"
}
result := map[string]string{
"node_id": n.nodeID,
"connected": connected,
"state": state,
"adapter": req.GetAdapter(),
"target": req.GetTarget(),
"session_id": normalizeSessionID(req.GetSessionId()),
}
return &iop.NodeCommandResponse{
RequestId: req.GetRequestId(),
Type: req.GetType(),
Adapter: req.GetAdapter(),
Target: req.GetTarget(),
SessionId: normalizeSessionID(req.GetSessionId()),
Result: result,
}
}
func (n *Node) dispatchAdapterCommand(ctx context.Context, req *iop.NodeCommandRequest) *iop.NodeCommandResponse {
adapter, err := n.router.LookupAdapter(req.GetAdapter())
if err != nil {
return n.commandErrorResponse(req, fmt.Sprintf("node: %s", err.Error()))
}
handler, ok := adapter.(runtime.CommandHandler)
if !ok {
return n.commandErrorResponse(req, fmt.Sprintf("node: adapter %q does not support commands", req.GetAdapter()))
}
domainReq := toDomainCommandRequest(req)
domainResp, err := handler.HandleCommand(ctx, domainReq)
if err != nil {
return n.commandErrorResponse(req, err.Error())
}
return toProtoCommandResponse(domainResp)
}
func (n *Node) commandErrorResponse(req *iop.NodeCommandRequest, msg string) *iop.NodeCommandResponse {
return &iop.NodeCommandResponse{
RequestId: req.GetRequestId(),
Type: req.GetType(),
Adapter: req.GetAdapter(),
Target: req.GetTarget(),
SessionId: req.GetSessionId(),
Error: msg,
}
}
// protoCommandTypeToDomain maps a proto NodeCommandType to its runtime
// CommandType. The second return is false for UNSPECIFIED or unknown values
// so callers can reject unsupported commands before dispatching.
func protoCommandTypeToDomain(t iop.NodeCommandType) (runtime.CommandType, bool) {
switch t {
case iop.NodeCommandType_NODE_COMMAND_TYPE_USAGE_STATUS:
return runtime.CommandTypeUsageStatus, true
case iop.NodeCommandType_NODE_COMMAND_TYPE_CAPABILITIES:
return runtime.CommandTypeCapabilities, true
case iop.NodeCommandType_NODE_COMMAND_TYPE_SESSION_LIST:
return runtime.CommandTypeSessionList, true
case iop.NodeCommandType_NODE_COMMAND_TYPE_TRANSPORT_STATUS:
return runtime.CommandTypeTransportStatus, true
case iop.NodeCommandType_NODE_COMMAND_TYPE_OLLAMA_API:
return runtime.CommandTypeOllamaAPI, true
default:
return "", false
}
}
func domainCommandTypeToProto(t runtime.CommandType) iop.NodeCommandType {
switch t {
case runtime.CommandTypeUsageStatus:
return iop.NodeCommandType_NODE_COMMAND_TYPE_USAGE_STATUS
case runtime.CommandTypeCapabilities:
return iop.NodeCommandType_NODE_COMMAND_TYPE_CAPABILITIES
case runtime.CommandTypeSessionList:
return iop.NodeCommandType_NODE_COMMAND_TYPE_SESSION_LIST
case runtime.CommandTypeTransportStatus:
return iop.NodeCommandType_NODE_COMMAND_TYPE_TRANSPORT_STATUS
case runtime.CommandTypeOllamaAPI:
return iop.NodeCommandType_NODE_COMMAND_TYPE_OLLAMA_API
default:
return iop.NodeCommandType_NODE_COMMAND_TYPE_UNSPECIFIED
}
}
func toDomainCommandRequest(req *iop.NodeCommandRequest) runtime.CommandRequest {
cmdType, _ := protoCommandTypeToDomain(req.GetType())
return runtime.CommandRequest{
RequestID: req.GetRequestId(),
Type: cmdType,
Adapter: req.GetAdapter(),
Target: req.GetTarget(),
SessionID: normalizeSessionID(req.GetSessionId()),
TimeoutSec: int(req.GetTimeoutSec()),
Metadata: req.GetMetadata(),
}
}
func toProtoCommandResponse(resp runtime.CommandResponse) *iop.NodeCommandResponse {
out := &iop.NodeCommandResponse{
RequestId: resp.RequestID,
Type: domainCommandTypeToProto(resp.Type),
Adapter: resp.Adapter,
Target: resp.Target,
SessionId: resp.SessionID,
Result: resp.Result,
}
if resp.UsageStatus != nil {
out.UsageStatus = &iop.AgentUsageStatus{
RawOutput: resp.UsageStatus.RawOutput,
DailyLimit: resp.UsageStatus.DailyLimit,
DailyResetTime: resp.UsageStatus.DailyResetTime,
WeeklyLimit: resp.UsageStatus.WeeklyLimit,
WeeklyResetTime: resp.UsageStatus.WeeklyResetTime,
Metadata: resp.UsageStatus.Metadata,
}
}
return out
}
type protoSender interface {
Send(m proto.Message) error
}
// sessionSink wraps a transport.Session to implement runtime.EventSink.
type sessionSink struct {
sess protoSender
out io.Writer
nodeID string
sessionID string
background bool
streaming bool
lineEnded bool
}
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(),
SessionId: s.sessionID,
Background: s.background,
NodeId: s.nodeID,
}
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:
s.streaming = false
s.lineEnded = true
fmt.Fprintf(s.out, "[node-event] start run_id=%s\n", event.RunID)
case runtime.EventTypeDelta:
if event.Delta == "" {
return
}
if !s.streaming {
fmt.Fprint(s.out, "[node-message] ")
s.streaming = true
}
fmt.Fprint(s.out, event.Delta)
s.lineEnded = strings.HasSuffix(event.Delta, "\n")
case runtime.EventTypeComplete:
if s.streaming && !s.lineEnded {
fmt.Fprintln(s.out)
}
s.streaming = false
s.lineEnded = true
fmt.Fprintf(s.out, "[node-event] complete run_id=%s detail=%q\n", event.RunID, event.Message)
case runtime.EventTypeError:
if s.streaming && !s.lineEnded {
fmt.Fprintln(s.out)
}
s.streaming = false
s.lineEnded = true
fmt.Fprintf(s.out, "[node-event] error run_id=%s detail=%q\n", event.RunID, event.Error)
case runtime.EventTypeCancelled:
if s.streaming && !s.lineEnded {
fmt.Fprintln(s.out)
}
s.streaming = false
s.lineEnded = true
fmt.Fprintf(s.out, "[node-event] cancelled run_id=%s\n", event.RunID)
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()
}
func sessionModeFromProto(m iop.RunSessionMode) runtime.SessionMode {
if m == iop.RunSessionMode_RUN_SESSION_MODE_REQUIRE_EXISTING {
return runtime.SessionModeRequireExisting
}
return runtime.SessionModeCreateIfMissing
}
func cancelActionFromProto(a iop.CancelAction) runtime.CancelAction {
if a == iop.CancelAction_CANCEL_ACTION_TERMINATE_SESSION {
return runtime.CancelActionTerminateSession
}
return runtime.CancelActionCancelRun
}
func normalizeSessionID(id string) string {
if id == "" {
return runtime.DefaultSessionID
}
return id
}