diff --git a/README.md b/README.md index 4858144..9fbc89f 100644 --- a/README.md +++ b/README.md @@ -128,6 +128,20 @@ Node는 모델 런타임, CLI Agent, 도구 실행을 담당한다. Node는 Cont Node는 가능한 한 단순한 실행 단위로 유지한다. 정책, 전체 시스템 조정, 다중 Edge 운영 판단을 Node에 밀어 넣지 않고, 전달받은 실행을 안정적으로 수행하는 데 집중한다. +### Worker Structure + +Edge, Node, Control Plane은 별도 `iop-worker` 앱으로 분리하지 않고, 각 Go 서비스 내부의 공통 Worker 모듈을 사용한다. 공통 처리 모델은 `Job Queue`, `Worker Pool`, `Job Status`, `Retry`, `Timeout`, `Cancel`이며, Worker가 담당하는 역할은 서비스별 책임에 맞춰 분리한다. + +- **Edge Worker** + - 사용자 요청 처리 흐름에 붙는 짧은 병렬/비동기 작업을 담당한다. + - intent 분석, history refinement, routing 보조, 응답 validation, fallback 판단, stream 종료 후 usage/log/metric 기록을 처리한다. +- **Node Worker** + - 모델 런타임과 로컬 프로세스에 붙는 작업을 담당한다. + - runtime adapter 처리, model process 상태 감시, local queue 처리, streaming relay 보조, local metric/log flush를 처리한다. +- **Control Plane Worker** + - 운영/관리/스케줄 기반 작업을 담당한다. + - node health 수집, model registry 동기화, policy/config 배포, drain/reload 명령, benchmark/job 실행, 운영 리포트와 cleanup 작업을 처리한다. + ## Execution Model ### Adapter and Target @@ -195,7 +209,7 @@ NomadCode | `apps/node` | Edge에 연결되어 adapter execution을 수행하는 Node agent | | `apps/edge` | Node registry, 설정 전달, routing, stream relay를 담당하는 Edge skeleton | | `apps/control-plane` | 향후 여러 Edge를 연결하고 운영 화면을 제공할 중앙 관리 계층 | -| `apps/worker` | 향후 비동기 작업 처리 가능성을 위한 placeholder | +| `apps/worker` | 현재 placeholder이며, Worker 구조는 우선 각 Go 서비스 내부 공통 모듈 방향으로 둔다 | | `packages` | 설정, 인증, 정책, 작업, 관측성, 버전 등 공통 패키지 | | `proto` | 앱 간 메시지 계약 원본과 생성물 | | `configs` | 현재 개발용 설정 예시 | diff --git a/apps/edge/cmd/edge/console_events.go b/apps/edge/cmd/edge/console_events.go index d76f9e2..aef2622 100644 --- a/apps/edge/cmd/edge/console_events.go +++ b/apps/edge/cmd/edge/console_events.go @@ -126,16 +126,31 @@ func (r *consoleEventRouter) printNodeEventLocked(event *iop.EdgeNodeEvent) { } label := r.edgeNodeLabel(event) + detail := transportCloseDetail(event.GetMetadata()) switch event.GetType() { case eventpkg.TypeNodeConnected: - fmt.Fprintf(r.out, "[node-%s-event] connected reason=%q\n", label, event.GetReason()) + fmt.Fprintf(r.out, "[node-%s-event] connected reason=%q%s\n", label, event.GetReason(), detail) case eventpkg.TypeNodeDisconnected: - fmt.Fprintf(r.out, "[node-%s-event] disconnected reason=%q\n", label, event.GetReason()) + fmt.Fprintf(r.out, "[node-%s-event] disconnected reason=%q%s\n", label, event.GetReason(), detail) default: - fmt.Fprintf(r.out, "[node-%s-event] %s reason=%q\n", label, event.GetType(), event.GetReason()) + fmt.Fprintf(r.out, "[node-%s-event] %s reason=%q%s\n", label, event.GetType(), event.GetReason(), detail) } } +func transportCloseDetail(metadata map[string]string) string { + if metadata == nil { + return "" + } + detail := "" + if reason := metadata[eventpkg.MetadataTransportCloseReason]; reason != "" { + detail += fmt.Sprintf(" transport_close_reason=%q", reason) + } + if err := metadata[eventpkg.MetadataTransportCloseError]; err != "" { + detail += fmt.Sprintf(" transport_close_error=%q", err) + } + return detail +} + func (r *consoleEventRouter) printAsyncLocked(event *iop.RunEvent) { if r.out == nil { return diff --git a/apps/edge/cmd/edge/console_test.go b/apps/edge/cmd/edge/console_test.go index 8fc73de..dabbcbe 100644 --- a/apps/edge/cmd/edge/console_test.go +++ b/apps/edge/cmd/edge/console_test.go @@ -421,10 +421,20 @@ func TestConsoleEventRouterPrintsNodeLifecycleEvents(t *testing.T) { NodeId: "node-1", Alias: "alias-1", Reason: eventpkg.ReasonTransportClosed, + Metadata: map[string]string{ + eventpkg.MetadataTransportCloseReason: "heartbeat_timeout", + eventpkg.MetadataTransportCloseError: "no heartbeat response within 10s", + }, }) got := out.String() if !strings.Contains(got, `[node-alias-1-event] disconnected reason="transport_closed"`) { t.Errorf("expected disconnected lifecycle output, got:\n%s", got) } + if !strings.Contains(got, `transport_close_reason="heartbeat_timeout"`) { + t.Errorf("expected transport close reason, got:\n%s", got) + } + if !strings.Contains(got, `transport_close_error="no heartbeat response within 10s"`) { + t.Errorf("expected transport close error, got:\n%s", got) + } } diff --git a/apps/edge/internal/transport/integration_test.go b/apps/edge/internal/transport/integration_test.go index 209fe3c..7a98c7d 100644 --- a/apps/edge/internal/transport/integration_test.go +++ b/apps/edge/internal/transport/integration_test.go @@ -184,6 +184,12 @@ func TestEdgeServerIntegration(t *testing.T) { if event.GetNodeId() != wantNodeID || event.GetAlias() != "test-node" || event.GetReason() != eventpkg.ReasonTransportClosed { t.Fatalf("unexpected disconnected event: %+v", event) } + if event.GetMetadata()[eventpkg.MetadataTransportCloseReason] != toki.DisconnectReasonRemoteClosed { + t.Fatalf("transport close reason: got %q want %q", event.GetMetadata()[eventpkg.MetadataTransportCloseReason], toki.DisconnectReasonRemoteClosed) + } + if event.GetMetadata()[eventpkg.MetadataTransportCloseError] == "" { + t.Fatalf("expected transport close error metadata, got %+v", event.GetMetadata()) + } case <-time.After(2 * time.Second): t.Fatal("timeout waiting for node disconnected event") } diff --git a/apps/edge/internal/transport/server.go b/apps/edge/internal/transport/server.go index e9dd28a..68ae3f0 100644 --- a/apps/edge/internal/transport/server.go +++ b/apps/edge/internal/transport/server.go @@ -154,7 +154,10 @@ func (s *Server) onNodeConnected(client *toki.TcpClient) { }) client.AddDisconnectListener(func(_ *toki.TcpClient) { s.registry.Unregister(rec.ID) - s.logger.Info("node unregistered", zap.String("node_id", rec.ID)) + transportInfo := client.DisconnectInfo() + fields := []zap.Field{zap.String("node_id", rec.ID)} + fields = append(fields, transportDisconnectFields(transportInfo)...) + s.logger.Info("node unregistered", fields...) reason := events.ReasonTransportClosed if s.stopping.Load() { reason = events.ReasonEdgeShutdown @@ -165,7 +168,7 @@ func (s *Server) onNodeConnected(client *toki.TcpClient) { rec.ID, rec.Alias, reason, - nil, + transportDisconnectMetadata(transportInfo), )) }) s.registry.Register(entry) @@ -212,3 +215,28 @@ func safePrefix(s string) string { } return s } + +func transportDisconnectMetadata(info toki.DisconnectInfo) map[string]string { + metadata := make(map[string]string, 2) + if info.Reason != "" { + metadata[events.MetadataTransportCloseReason] = info.Reason + } + if info.Error != "" { + metadata[events.MetadataTransportCloseError] = info.Error + } + if len(metadata) == 0 { + return nil + } + return metadata +} + +func transportDisconnectFields(info toki.DisconnectInfo) []zap.Field { + fields := make([]zap.Field, 0, 2) + if info.Reason != "" { + fields = append(fields, zap.String("transport_close_reason", info.Reason)) + } + if info.Error != "" { + fields = append(fields, zap.String("transport_close_error", info.Error)) + } + return fields +} diff --git a/apps/node/internal/bootstrap/module.go b/apps/node/internal/bootstrap/module.go index 162dea6..96cb744 100644 --- a/apps/node/internal/bootstrap/module.go +++ b/apps/node/internal/bootstrap/module.go @@ -120,8 +120,22 @@ func printEdgeEvent(out io.Writer, event *iop.EdgeNodeEvent) { } switch event.GetType() { case events.TypeEdgeDisconnected: - fmt.Fprintf(out, "[edge-event] disconnected reason=%q\n", event.GetReason()) + fmt.Fprintf(out, "[edge-event] disconnected reason=%q%s\n", event.GetReason(), transportCloseDetail(event.GetMetadata())) default: - fmt.Fprintf(out, "[edge-event] %s reason=%q\n", event.GetType(), event.GetReason()) + fmt.Fprintf(out, "[edge-event] %s reason=%q%s\n", event.GetType(), event.GetReason(), transportCloseDetail(event.GetMetadata())) } } + +func transportCloseDetail(metadata map[string]string) string { + if metadata == nil { + return "" + } + detail := "" + if reason := metadata[events.MetadataTransportCloseReason]; reason != "" { + detail += fmt.Sprintf(" transport_close_reason=%q", reason) + } + if err := metadata[events.MetadataTransportCloseError]; err != "" { + detail += fmt.Sprintf(" transport_close_error=%q", err) + } + return detail +} diff --git a/apps/node/internal/transport/integration_test.go b/apps/node/internal/transport/integration_test.go index 58b072b..3eb8de9 100644 --- a/apps/node/internal/transport/integration_test.go +++ b/apps/node/internal/transport/integration_test.go @@ -186,6 +186,12 @@ func TestNodeClientIntegration(t *testing.T) { if event.GetNodeId() != "test-node" || event.GetAlias() != "test-alias" || event.GetReason() != eventpkg.ReasonTransportClosed { t.Fatalf("unexpected edge disconnected event: %+v", event) } + if event.GetMetadata()[eventpkg.MetadataTransportCloseReason] != toki.DisconnectReasonRemoteClosed { + t.Fatalf("transport close reason: got %q want %q", event.GetMetadata()[eventpkg.MetadataTransportCloseReason], toki.DisconnectReasonRemoteClosed) + } + if event.GetMetadata()[eventpkg.MetadataTransportCloseError] == "" { + t.Fatalf("expected transport close error metadata, got %+v", event.GetMetadata()) + } case <-time.After(2 * time.Second): t.Fatal("timeout waiting for edge disconnected event") } diff --git a/apps/node/internal/transport/session.go b/apps/node/internal/transport/session.go index 2a44202..634e550 100644 --- a/apps/node/internal/transport/session.go +++ b/apps/node/internal/transport/session.go @@ -85,14 +85,15 @@ func newSession(client *toki.TcpClient, logger *zap.Logger, nodeID, alias string }) client.AddDisconnectListener(func(_ *toki.TcpClient) { - logger.Info("disconnected from edge") + transportInfo := client.DisconnectInfo() + logger.Info("disconnected from edge", transportDisconnectFields(transportInfo)...) s.emitEvent(events.NewEdgeNodeEvent( events.SourceNode, events.TypeEdgeDisconnected, nodeID, alias, s.disconnectReason(), - nil, + transportDisconnectMetadata(transportInfo), )) }) @@ -152,3 +153,28 @@ func (s *Session) disconnectReason() string { } return reason } + +func transportDisconnectMetadata(info toki.DisconnectInfo) map[string]string { + metadata := make(map[string]string, 2) + if info.Reason != "" { + metadata[events.MetadataTransportCloseReason] = info.Reason + } + if info.Error != "" { + metadata[events.MetadataTransportCloseError] = info.Error + } + if len(metadata) == 0 { + return nil + } + return metadata +} + +func transportDisconnectFields(info toki.DisconnectInfo) []zap.Field { + fields := make([]zap.Field, 0, 2) + if info.Reason != "" { + fields = append(fields, zap.String("transport_close_reason", info.Reason)) + } + if info.Error != "" { + fields = append(fields, zap.String("transport_close_error", info.Error)) + } + return fields +} diff --git a/packages/events/events.go b/packages/events/events.go index b806a2d..3208367 100644 --- a/packages/events/events.go +++ b/packages/events/events.go @@ -20,6 +20,9 @@ const ( ReasonTransportClosed = "transport_closed" ReasonLocalShutdown = "local_shutdown" ReasonEdgeShutdown = "edge_shutdown" + + MetadataTransportCloseReason = "transport_close_reason" + MetadataTransportCloseError = "transport_close_error" ) func NewEdgeNodeEvent(source, eventType, nodeID, alias, reason string, metadata map[string]string) *iop.EdgeNodeEvent {