- edge_config_mapper: 엣지 설정 매핑 기능 추가 - edge_node_id: 엣지 노드 ID 생성/관리 기능 추가 - node_router_registry: 노드 라우터 레지스트리 기능 추가 - node_writer_injection: 노드 writer 주입 기능 추가 - edge console 및 server 업데이트 - node bootstrap module 업데이트
389 lines
11 KiB
Go
389 lines
11 KiB
Go
package main
|
|
|
|
import (
|
|
"bufio"
|
|
"context"
|
|
"fmt"
|
|
"io"
|
|
"strings"
|
|
"time"
|
|
|
|
toki "git.toki-labs.com/toki/common-proto-socket/go"
|
|
"go.uber.org/zap"
|
|
"google.golang.org/protobuf/types/known/structpb"
|
|
|
|
edgenode "iop/apps/edge/internal/node"
|
|
"iop/apps/edge/internal/transport"
|
|
"iop/packages/config"
|
|
"iop/packages/observability"
|
|
iop "iop/proto/gen/iop"
|
|
)
|
|
|
|
func runConsole(ctx context.Context, cfg *config.EdgeConfig, in io.Reader, out io.Writer) error {
|
|
logger, err := observability.NewLogger(cfg.Logging.Level, cfg.Logging.Pretty)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer func() { _ = logger.Sync() }()
|
|
|
|
registry := edgenode.NewRegistry()
|
|
nodeStore, err := edgenode.LoadFromConfig(cfg.Nodes)
|
|
if err != nil {
|
|
return fmt.Errorf("edge: seed node store: %w", err)
|
|
}
|
|
|
|
server, err := transport.NewServer(cfg.Server.Listen, registry, nodeStore, logger)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
events := newConsoleEventRouter(out, registry, logger)
|
|
server.SetRunEventHandler(events.Handle)
|
|
|
|
if err := server.Start(ctx); err != nil {
|
|
return err
|
|
}
|
|
defer func() { _ = server.Stop() }()
|
|
|
|
go func() {
|
|
if err := observability.ServeMetrics(cfg.Metrics.Port); err != nil {
|
|
logger.Warn("metrics server exited", zap.Error(err))
|
|
}
|
|
}()
|
|
|
|
// Console state: target node, adapter, agent, session, background mode.
|
|
target := &consoleTarget{
|
|
Adapter: cfg.Console.Adapter,
|
|
Agent: cfg.Console.ResolveAgent(),
|
|
SessionID: normalizeConsoleSessionID(cfg.Console.SessionID),
|
|
Background: cfg.Console.Background,
|
|
TimeoutSec: cfg.Console.TimeoutSec,
|
|
}
|
|
|
|
fmt.Fprintf(out, "IOP Edge console listening on %s\n", cfg.Server.Listen)
|
|
fmt.Fprintf(out, "Console target node=%s adapter=%s agent=%s session=%s background=%v\n",
|
|
target.NodeRef, target.Adapter, target.Agent, target.SessionID, target.Background)
|
|
fmt.Fprintln(out, "Start node.sh on another host, then type a message here.")
|
|
fmt.Fprintln(out, "Commands: /nodes, /node <id|alias>, /session <id>, /background on|off, /terminate-session, /status, /exit")
|
|
|
|
scanner := bufio.NewScanner(in)
|
|
for {
|
|
fmt.Fprint(out, "edge> ")
|
|
if !scanner.Scan() {
|
|
fmt.Fprintln(out)
|
|
return scanner.Err()
|
|
}
|
|
|
|
message := strings.TrimSpace(scanner.Text())
|
|
lower := strings.ToLower(message)
|
|
switch {
|
|
case message == "":
|
|
continue
|
|
case lower == "/exit" || lower == "/quit" || lower == "exit" || lower == "quit":
|
|
fmt.Fprintln(out, "bye")
|
|
return nil
|
|
case lower == "/nodes":
|
|
printNodes(out, registry, target.NodeRef)
|
|
case strings.HasPrefix(lower, "/node "):
|
|
parts := strings.Fields(message)
|
|
if len(parts) < 2 || parts[1] == "" {
|
|
fmt.Fprintln(out, "usage: /node <id|alias>")
|
|
continue
|
|
}
|
|
if _, err := resolveConsoleNode(registry, parts[1]); err != nil {
|
|
fmt.Fprintf(out, "error: %v\n", err)
|
|
continue
|
|
}
|
|
target.NodeRef = parts[1]
|
|
fmt.Fprintf(out, "node → %s\n", target.NodeRef)
|
|
case strings.HasPrefix(lower, "/session "):
|
|
parts := strings.Fields(message)
|
|
if len(parts) < 2 || parts[1] == "" {
|
|
fmt.Fprintln(out, "usage: /session <id>")
|
|
continue
|
|
}
|
|
target.SessionID = parts[1]
|
|
fmt.Fprintf(out, "session → %s\n", target.SessionID)
|
|
case strings.HasPrefix(lower, "/background "):
|
|
parts := strings.Fields(lower)
|
|
if len(parts) < 2 {
|
|
fmt.Fprintln(out, "usage: /background on|off")
|
|
continue
|
|
}
|
|
switch parts[1] {
|
|
case "on":
|
|
target.Background = true
|
|
fmt.Fprintln(out, "background → on")
|
|
case "off":
|
|
target.Background = false
|
|
fmt.Fprintln(out, "background → off")
|
|
default:
|
|
fmt.Fprintln(out, "usage: /background on|off")
|
|
}
|
|
case lower == "/terminate-session":
|
|
handleTerminateSession(ctx, registry, out, target)
|
|
case lower == "/status":
|
|
if err := sendConsoleStatus(ctx, registry, out, target); err != nil {
|
|
fmt.Fprintf(out, "error: %v\n", err)
|
|
}
|
|
default:
|
|
if err := sendConsoleRun(ctx, registry, events, out, target, message); err != nil {
|
|
fmt.Fprintf(out, "error: %v\n", err)
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
type consoleTarget struct {
|
|
NodeRef string
|
|
Adapter string
|
|
Agent string
|
|
SessionID string
|
|
Background bool
|
|
TimeoutSec int
|
|
}
|
|
|
|
func resolveConsoleNode(registry *edgenode.Registry, nodeRef string) (*edgenode.NodeEntry, error) {
|
|
return registry.Resolve(nodeRef)
|
|
}
|
|
|
|
func printNodes(out io.Writer, registry *edgenode.Registry, selectedRef string) {
|
|
nodes := registry.All()
|
|
if len(nodes) == 0 {
|
|
fmt.Fprintln(out, "no nodes connected")
|
|
return
|
|
}
|
|
for _, node := range nodes {
|
|
marker := " "
|
|
if selectedRef != "" && (node.NodeID == selectedRef || node.Alias == selectedRef) {
|
|
marker = "* "
|
|
}
|
|
fmt.Fprintf(out, "%s%s (%s)\n", marker, node.NodeID, node.Alias)
|
|
}
|
|
}
|
|
|
|
// buildRunRequest constructs the RunRequest for a console send.
|
|
// The wire schema still uses RunRequest.model; for the CLI adapter this value
|
|
// semantically carries the selected agent/profile name.
|
|
func buildRunRequest(adapter, agent, sessionID string, background bool, timeoutSec int, message string) (*iop.RunRequest, string, error) {
|
|
input, err := structpb.NewStruct(map[string]any{"prompt": message})
|
|
if err != nil {
|
|
return nil, "", err
|
|
}
|
|
if timeoutSec <= 0 {
|
|
timeoutSec = 30
|
|
}
|
|
runID := fmt.Sprintf("manual-%d", time.Now().UnixNano())
|
|
req := &iop.RunRequest{
|
|
RunId: runID,
|
|
Adapter: adapter,
|
|
Model: agent,
|
|
SessionId: normalizeConsoleSessionID(sessionID),
|
|
SessionMode: iop.RunSessionMode_RUN_SESSION_MODE_CREATE_IF_MISSING,
|
|
Background: background,
|
|
Input: input,
|
|
TimeoutSec: int32(timeoutSec),
|
|
Metadata: map[string]string{
|
|
"source": "edge-console",
|
|
},
|
|
}
|
|
return req, runID, nil
|
|
}
|
|
|
|
func sendConsoleRun(ctx context.Context, registry *edgenode.Registry, events *consoleEventRouter, out io.Writer, target *consoleTarget, message string) error {
|
|
entry, err := resolveConsoleNode(registry, target.NodeRef)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
req, runID, err := buildRunRequest(target.Adapter, target.Agent, target.SessionID, target.Background, target.TimeoutSec, message)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
var runEvents <-chan *iop.RunEvent
|
|
var unregister func()
|
|
if !target.Background {
|
|
runEvents, unregister = events.Register(runID)
|
|
defer unregister()
|
|
}
|
|
|
|
nodeAlias := entry.Alias
|
|
if nodeAlias == "" {
|
|
nodeAlias = entry.NodeID
|
|
}
|
|
fmt.Fprintf(out, "[edge] sent run_id=%s node=%s adapter=%s agent=%s session=%s background=%v\n",
|
|
runID, nodeAlias, target.Adapter, target.Agent, req.GetSessionId(), target.Background)
|
|
if err := entry.Client.Send(req); err != nil {
|
|
return err
|
|
}
|
|
|
|
if target.Background {
|
|
fmt.Fprintf(out, "[edge] background run dispatched, events will arrive asynchronously\n")
|
|
return nil
|
|
}
|
|
|
|
timer := time.NewTimer(time.Duration(target.TimeoutSec+5) * time.Second)
|
|
defer timer.Stop()
|
|
|
|
var response *consoleResponseStream
|
|
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return ctx.Err()
|
|
case <-timer.C:
|
|
return fmt.Errorf("timed out waiting for node response")
|
|
case event := <-runEvents:
|
|
if event.GetRunId() != runID {
|
|
continue
|
|
}
|
|
label := events.nodeLabel(event)
|
|
if response == nil {
|
|
response = newConsoleResponseStream(out, fmt.Sprintf("[node-%s-message] ", label))
|
|
}
|
|
|
|
switch event.GetType() {
|
|
case "start":
|
|
fmt.Fprintf(out, "[node-%s-event] start run_id=%s\n", label, runID)
|
|
case "delta":
|
|
response.Write(event.GetDelta())
|
|
case "complete":
|
|
response.Finish()
|
|
fmt.Fprintf(out, "[node-%s-event] complete run_id=%s detail=%q\n", label, runID, event.GetMessage())
|
|
return nil
|
|
case "cancelled":
|
|
response.FinishIfStarted()
|
|
fmt.Fprintf(out, "[node-%s-event] cancelled run_id=%s\n", label, runID)
|
|
return nil
|
|
case "error":
|
|
response.FinishIfStarted()
|
|
fmt.Fprintf(out, "[node-%s-event] error run_id=%s detail=%q\n", label, runID, event.GetError())
|
|
return fmt.Errorf("node reported error")
|
|
default:
|
|
fmt.Fprintf(out, "[node-%s-event] %s run_id=%s detail=%q\n", label, event.GetType(), runID, event.GetMessage())
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
var sendTerminateSessionFunc = sendTerminateSession
|
|
|
|
func handleTerminateSession(ctx context.Context, registry *edgenode.Registry, out io.Writer, target *consoleTarget) {
|
|
if label, err := sendTerminateSessionFunc(ctx, registry, target); err != nil {
|
|
fmt.Fprintf(out, "error: %v\n", err)
|
|
} else {
|
|
fmt.Fprintf(out, "terminated session %s node=%s\n", target.SessionID, label)
|
|
}
|
|
}
|
|
|
|
func sendTerminateSession(ctx context.Context, registry *edgenode.Registry, target *consoleTarget) (string, error) {
|
|
entry, err := resolveConsoleNode(registry, target.NodeRef)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
req := &iop.CancelRequest{
|
|
Adapter: target.Adapter,
|
|
Model: target.Agent,
|
|
SessionId: normalizeConsoleSessionID(target.SessionID),
|
|
Action: iop.CancelAction_CANCEL_ACTION_TERMINATE_SESSION,
|
|
}
|
|
label := entry.Alias
|
|
if label == "" {
|
|
label = entry.NodeID
|
|
}
|
|
return label, entry.Client.Send(req)
|
|
}
|
|
|
|
func normalizeConsoleSessionID(id string) string {
|
|
if id == "" {
|
|
return "default"
|
|
}
|
|
return id
|
|
}
|
|
|
|
func buildNodeCommandRequest(adapter, agent, sessionID string, timeoutSec int) (*iop.NodeCommandRequest, string) {
|
|
if timeoutSec <= 0 {
|
|
timeoutSec = 30
|
|
}
|
|
reqID := fmt.Sprintf("status-%d", time.Now().UnixNano())
|
|
req := &iop.NodeCommandRequest{
|
|
RequestId: reqID,
|
|
Type: iop.NodeCommandType_NODE_COMMAND_TYPE_USAGE_STATUS,
|
|
Adapter: adapter,
|
|
Model: agent,
|
|
SessionId: normalizeConsoleSessionID(sessionID),
|
|
TimeoutSec: int32(timeoutSec),
|
|
}
|
|
return req, reqID
|
|
}
|
|
|
|
func statusWaitTimeout(req *iop.NodeCommandRequest) time.Duration {
|
|
return time.Duration(req.GetTimeoutSec()+5) * time.Second
|
|
}
|
|
|
|
func sendConsoleStatus(ctx context.Context, registry *edgenode.Registry, out io.Writer, target *consoleTarget) error {
|
|
entry, err := resolveConsoleNode(registry, target.NodeRef)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
req, _ := buildNodeCommandRequest(target.Adapter, target.Agent, target.SessionID, target.TimeoutSec)
|
|
|
|
nodeAlias := entry.Alias
|
|
if nodeAlias == "" {
|
|
nodeAlias = entry.NodeID
|
|
}
|
|
fmt.Fprintf(out, "[edge] sent command=status node=%s adapter=%s agent=%s session=%s\n", nodeAlias, target.Adapter, target.Agent, req.GetSessionId())
|
|
|
|
timeout := statusWaitTimeout(req)
|
|
resp, err := toki.SendRequestTyped[*iop.NodeCommandRequest, *iop.NodeCommandResponse](
|
|
&entry.Client.Communicator,
|
|
req,
|
|
timeout,
|
|
)
|
|
if err != nil {
|
|
return fmt.Errorf("transport error: %w", err)
|
|
}
|
|
|
|
if resp.GetError() != "" {
|
|
return fmt.Errorf("node reported error: %s", resp.GetError())
|
|
}
|
|
|
|
formatUsageStatus(out, nodeAlias, target.Agent, req.GetSessionId(), resp.GetUsageStatus())
|
|
return nil
|
|
}
|
|
|
|
func formatUsageStatus(out io.Writer, nodeAlias, agent, sessionID string, status *iop.AgentUsageStatus) {
|
|
fmt.Fprintf(out, "[node-%s-status] agent=%s session=%s\n", nodeAlias, agent, sessionID)
|
|
|
|
if status == nil {
|
|
fmt.Fprintln(out, "no usage status provided")
|
|
return
|
|
}
|
|
|
|
hasParsedLimits := false
|
|
if status.GetDailyLimit() != "" {
|
|
fmt.Fprintf(out, "Daily limit: %s (resets %s)\n", status.GetDailyLimit(), status.GetDailyResetTime())
|
|
hasParsedLimits = true
|
|
}
|
|
if status.GetWeeklyLimit() != "" {
|
|
fmt.Fprintf(out, "Weekly limit: %s (resets %s)\n", status.GetWeeklyLimit(), status.GetWeeklyResetTime())
|
|
hasParsedLimits = true
|
|
}
|
|
|
|
if !hasParsedLimits {
|
|
if status.GetRawOutput() == "" {
|
|
fmt.Fprintln(out, "raw output did not include parsed limits and was empty")
|
|
} else {
|
|
fmt.Fprintln(out, "raw output did not include parsed limits:")
|
|
lines := strings.Split(status.GetRawOutput(), "\n")
|
|
for i, line := range lines {
|
|
if i >= 5 {
|
|
fmt.Fprintln(out, "...")
|
|
break
|
|
}
|
|
fmt.Fprintln(out, line)
|
|
}
|
|
}
|
|
}
|
|
}
|