- 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
1281 lines
36 KiB
Text
1281 lines
36 KiB
Text
<!-- task=node_centralized_mgmt plan=0 tag=REFACTOR -->
|
|
|
|
# Jenkins 스타일 노드 중앙 관리 구조 리팩토링
|
|
|
|
## 이 파일을 읽는 구현 에이전트에게
|
|
|
|
각 항목의 체크리스트를 완료하면 즉시 체크 표시한다. 각 중간 검증 명령을 실제로 실행하고 출력을 `CODE_REVIEW.md`의 검증 결과 섹션에 붙여 넣는다. 모든 항목 완료 후 최종 검증을 실행하고 결과를 기록한다.
|
|
|
|
## 배경
|
|
|
|
현재 구조는 각 node가 자신의 `node.yaml`에서 ID, 어댑터 설정, 런타임 설정을 자체 선언한다. 이는 노드가 많아질수록 분산 관리의 복잡도를 높이고, edge가 어떤 노드가 연결해야 하는지 사전에 알 수 없어 위조·중복 접속을 막을 방법이 없다. Jenkins master/agent 모델과 동일하게, edge가 노드를 사전 등록(hash 토큰 발급)하고 모든 설정을 중앙 관리하며, node 바이너리는 `edge_addr + token` 두 값만 받아 기동하는 구조로 전환한다.
|
|
|
|
## 의존 관계 및 구현 순서
|
|
|
|
```
|
|
REFACTOR-1 (proto)
|
|
└─ REFACTOR-2 (config)
|
|
├─ REFACTOR-3 (edge: NodeStore + NodeEntry)
|
|
│ └─ REFACTOR-4 (edge: transport + bootstrap)
|
|
└─ REFACTOR-5 (node: transport 핸드셰이크)
|
|
└─ REFACTOR-6 (node: adapter factory + bootstrap)
|
|
└─ REFACTOR-7 (테스트 갱신)
|
|
```
|
|
|
|
---
|
|
|
|
### [REFACTOR-1] proto: RegisterRequest / RegisterResponse / NodeConfigPayload 추가
|
|
|
|
#### 문제
|
|
|
|
`proto/iop/runtime.proto`에 node → edge 등록 메시지가 없다. 현재 handshake는 edge가 먼저 `CapabilityRequest`를 보내는 방식이지만, 새 설계에서는 node가 먼저 `RegisterRequest{token}`을 보내고 edge가 `RegisterResponse{config}`로 응답한다.
|
|
|
|
#### 해결 방법
|
|
|
|
`proto/iop/runtime.proto`에 아래 메시지를 추가한다. `CapabilityRequest`/`CapabilityResponse`는 wire 호환성 유지를 위해 삭제하지 않고 그대로 남긴다.
|
|
|
|
```protobuf
|
|
// 추가할 메시지 (runtime.proto 기존 내용 뒤에 append)
|
|
|
|
// RegisterRequest is sent by node to edge immediately on connect.
|
|
message RegisterRequest {
|
|
string token = 1;
|
|
}
|
|
|
|
// RegisterResponse is sent by edge to node in response to RegisterRequest.
|
|
message RegisterResponse {
|
|
bool accepted = 1;
|
|
string node_id = 2;
|
|
string alias = 3;
|
|
string reason = 4; // rejection reason
|
|
NodeConfigPayload config = 5;
|
|
}
|
|
|
|
// NodeConfigPayload carries all configuration edge pushes to the node.
|
|
message NodeConfigPayload {
|
|
repeated AdapterConfig adapters = 1;
|
|
NodeRuntimeConfig runtime = 2;
|
|
}
|
|
|
|
// AdapterConfig describes one adapter to enable on the node.
|
|
message AdapterConfig {
|
|
string type = 1; // "mock" | "ollama" | "vllm" | "cli"
|
|
bool enabled = 2;
|
|
google.protobuf.Struct settings = 3;
|
|
}
|
|
|
|
// NodeRuntimeConfig is the runtime tuning pushed to the node.
|
|
message NodeRuntimeConfig {
|
|
int32 concurrency = 1;
|
|
string workspace_root = 2;
|
|
}
|
|
```
|
|
|
|
`make proto`를 실행해 `proto/gen/iop/runtime.pb.go`를 재생성한다. 생성 파일은 직접 수정하지 않는다.
|
|
|
|
#### 수정 파일 및 체크리스트
|
|
|
|
- [x] `proto/iop/runtime.proto` — 위 5개 메시지 추가
|
|
- [x] `proto/gen/iop/runtime.pb.go` — `make proto`로 재생성 (직접 수정 금지)
|
|
|
|
#### 테스트 작성
|
|
|
|
skip — proto 파일 자체는 컴파일 검증으로 충분하다.
|
|
|
|
#### 중간 검증
|
|
|
|
```bash
|
|
make proto
|
|
go build ./...
|
|
```
|
|
|
|
컴파일 오류 없이 통과해야 한다.
|
|
|
|
---
|
|
|
|
### [REFACTOR-2] packages/config: NodeConfig 단순화 + EdgeConfig에 노드 정의 추가
|
|
|
|
#### 문제
|
|
|
|
`packages/config/config.go`의 `NodeConfig`(L7-16)에 `NodeInfo`, `RuntimeConf`, `SQLiteConf`, `AdaptersConf`가 있다. 이 값들은 이제 edge가 관리하므로 node config에서 제거한다. `EdgeConfig`(L18-23)에는 어떤 노드가 연결해야 하는지에 대한 정의가 전혀 없으므로 `Nodes []NodeDefinition`을 추가한다.
|
|
|
|
#### 해결 방법
|
|
|
|
**`packages/config/config.go` 변경:**
|
|
|
|
Before (`NodeConfig`, L7-16):
|
|
```go
|
|
type NodeConfig struct {
|
|
Node NodeInfo `mapstructure:"node" yaml:"node"`
|
|
Transport TransportConf `mapstructure:"transport" yaml:"transport"`
|
|
TLS TLSConf `mapstructure:"tls" yaml:"tls"`
|
|
Runtime RuntimeConf `mapstructure:"runtime" yaml:"runtime"`
|
|
SQLite SQLiteConf `mapstructure:"sqlite" yaml:"sqlite"`
|
|
Logging LoggingConf `mapstructure:"logging" yaml:"logging"`
|
|
Metrics MetricsConf `mapstructure:"metrics" yaml:"metrics"`
|
|
Adapters AdaptersConf `mapstructure:"adapters" yaml:"adapters"`
|
|
}
|
|
```
|
|
|
|
After:
|
|
```go
|
|
type NodeConfig struct {
|
|
Transport TransportConf `mapstructure:"transport" yaml:"transport"`
|
|
Logging LoggingConf `mapstructure:"logging" yaml:"logging"`
|
|
Metrics MetricsConf `mapstructure:"metrics" yaml:"metrics"`
|
|
}
|
|
```
|
|
|
|
Before (`TransportConf`, L34-36):
|
|
```go
|
|
type TransportConf struct {
|
|
EdgeAddr string `mapstructure:"edge_addr" yaml:"edge_addr"`
|
|
}
|
|
```
|
|
|
|
After:
|
|
```go
|
|
type TransportConf struct {
|
|
EdgeAddr string `mapstructure:"edge_addr" yaml:"edge_addr"`
|
|
Token string `mapstructure:"token" yaml:"token"`
|
|
}
|
|
```
|
|
|
|
Before (`EdgeConfig`, L18-23):
|
|
```go
|
|
type EdgeConfig struct {
|
|
Server EdgeServerConf `mapstructure:"server" yaml:"server"`
|
|
TLS TLSConf `mapstructure:"tls" yaml:"tls"`
|
|
Logging LoggingConf `mapstructure:"logging" yaml:"logging"`
|
|
Metrics MetricsConf `mapstructure:"metrics" yaml:"metrics"`
|
|
}
|
|
```
|
|
|
|
After:
|
|
```go
|
|
type EdgeConfig struct {
|
|
Server EdgeServerConf `mapstructure:"server" yaml:"server"`
|
|
TLS TLSConf `mapstructure:"tls" yaml:"tls"`
|
|
Logging LoggingConf `mapstructure:"logging" yaml:"logging"`
|
|
Metrics MetricsConf `mapstructure:"metrics" yaml:"metrics"`
|
|
Nodes []NodeDefinition `mapstructure:"nodes" yaml:"nodes"`
|
|
}
|
|
```
|
|
|
|
새 타입 `NodeDefinition` 추가 (EdgeConfig 아래):
|
|
```go
|
|
// NodeDefinition is the edge-side record for a pre-registered node.
|
|
type NodeDefinition struct {
|
|
Alias string `mapstructure:"alias" yaml:"alias"`
|
|
Token string `mapstructure:"token" yaml:"token"`
|
|
Adapters AdaptersConf `mapstructure:"adapters" yaml:"adapters"`
|
|
Runtime RuntimeConf `mapstructure:"runtime" yaml:"runtime"`
|
|
}
|
|
```
|
|
|
|
`setDefaults`(L117-125) 변경: node 관련 기본값 제거, token 기본값 없음 (필수값).
|
|
|
|
Before:
|
|
```go
|
|
func setDefaults(v *viper.Viper) {
|
|
v.SetDefault("transport.edge_addr", "localhost:9090")
|
|
v.SetDefault("runtime.concurrency", 4)
|
|
v.SetDefault("runtime.workspace_root", "/tmp/iop/workspace")
|
|
v.SetDefault("sqlite.dsn", "file:iop.db?cache=shared&mode=rwc")
|
|
v.SetDefault("logging.level", "info")
|
|
v.SetDefault("metrics.port", 9091)
|
|
v.SetDefault("tls.enabled", false)
|
|
}
|
|
```
|
|
|
|
After:
|
|
```go
|
|
func setDefaults(v *viper.Viper) {
|
|
v.SetDefault("transport.edge_addr", "localhost:9090")
|
|
v.SetDefault("logging.level", "info")
|
|
v.SetDefault("metrics.port", 9091)
|
|
}
|
|
```
|
|
|
|
**`configs/node.yaml` 변경:**
|
|
|
|
```yaml
|
|
transport:
|
|
edge_addr: "localhost:9090"
|
|
token: "changeme"
|
|
|
|
logging:
|
|
level: "info"
|
|
|
|
metrics:
|
|
port: 9091
|
|
```
|
|
|
|
**`configs/edge.yaml` 변경:**
|
|
|
|
```yaml
|
|
server:
|
|
listen: "0.0.0.0:9090"
|
|
|
|
tls:
|
|
enabled: false
|
|
|
|
logging:
|
|
level: "info"
|
|
|
|
metrics:
|
|
port: 9092
|
|
|
|
nodes:
|
|
- alias: "local-node"
|
|
token: "changeme"
|
|
adapters:
|
|
ollama:
|
|
enabled: false
|
|
base_url: "http://localhost:11434"
|
|
vllm:
|
|
enabled: false
|
|
endpoint: "http://localhost:8000"
|
|
cli:
|
|
enabled: false
|
|
profiles:
|
|
claude:
|
|
command: "claude"
|
|
args: []
|
|
env: []
|
|
runtime:
|
|
concurrency: 4
|
|
workspace_root: "/tmp/iop/workspace"
|
|
```
|
|
|
|
#### 수정 파일 및 체크리스트
|
|
|
|
- [x] `packages/config/config.go` — NodeConfig 단순화, TransportConf에 Token 추가, NodeDefinition 추가, EdgeConfig에 Nodes 추가, setDefaults 정리
|
|
- [x] `configs/node.yaml` — transport + token + logging + metrics만 남김
|
|
- [x] `configs/edge.yaml` — nodes 섹션 추가
|
|
- [x] `NodeInfo`, `RuntimeConf`, `SQLiteConf` 타입은 node가 더 이상 직접 사용하지 않지만 타입 자체는 config.go에 남겨 edge의 `NodeDefinition`에서 재사용한다
|
|
|
|
#### 테스트 작성
|
|
|
|
skip — config는 이후 bootstrap 변경 후 통합 검증한다.
|
|
|
|
#### 중간 검증
|
|
|
|
```bash
|
|
go build ./packages/config/...
|
|
```
|
|
|
|
---
|
|
|
|
### [REFACTOR-3] edge: NodeStore 신설 + NodeEntry 갱신
|
|
|
|
#### 문제
|
|
|
|
edge가 노드를 사전에 알 수 있는 저장소가 없다. `NodeEntry`(registry.go L13-17)에 alias, config 필드가 없다.
|
|
|
|
#### 해결 방법
|
|
|
|
**신규 파일 `apps/edge/internal/node/store.go` 작성:**
|
|
|
|
```go
|
|
package node
|
|
|
|
import (
|
|
"fmt"
|
|
"sync"
|
|
|
|
"iop/packages/config"
|
|
)
|
|
|
|
// NodeRecord is the pre-registered node definition stored in edge.
|
|
type NodeRecord struct {
|
|
ID string
|
|
Alias string
|
|
Token string
|
|
Adapters config.AdaptersConf
|
|
Runtime config.RuntimeConf
|
|
}
|
|
|
|
// NodeStore holds pre-registered node definitions, keyed by token.
|
|
type NodeStore struct {
|
|
mu sync.RWMutex
|
|
byToken map[string]*NodeRecord
|
|
byID map[string]*NodeRecord
|
|
}
|
|
|
|
func NewNodeStore() *NodeStore {
|
|
return &NodeStore{
|
|
byToken: make(map[string]*NodeRecord),
|
|
byID: make(map[string]*NodeRecord),
|
|
}
|
|
}
|
|
|
|
func (s *NodeStore) Add(rec *NodeRecord) {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
s.byToken[rec.Token] = rec
|
|
s.byID[rec.ID] = rec
|
|
}
|
|
|
|
func (s *NodeStore) FindByToken(token string) (*NodeRecord, bool) {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
r, ok := s.byToken[token]
|
|
return r, ok
|
|
}
|
|
|
|
func (s *NodeStore) FindByID(id string) (*NodeRecord, bool) {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
r, ok := s.byID[id]
|
|
return r, ok
|
|
}
|
|
|
|
// LoadFromConfig seeds the store from EdgeConfig.Nodes.
|
|
// ID is derived as "node-{alias}" if not provided — caller should use UUIDs in production.
|
|
func LoadFromConfig(defs []config.NodeDefinition) (*NodeStore, error) {
|
|
s := NewNodeStore()
|
|
seen := make(map[string]bool)
|
|
for i, d := range defs {
|
|
if d.Token == "" {
|
|
return nil, fmt.Errorf("node[%d] alias=%q: token must not be empty", i, d.Alias)
|
|
}
|
|
if seen[d.Token] {
|
|
return nil, fmt.Errorf("node[%d] alias=%q: duplicate token", i, d.Alias)
|
|
}
|
|
seen[d.Token] = true
|
|
s.Add(&NodeRecord{
|
|
ID: "node-" + d.Alias,
|
|
Alias: d.Alias,
|
|
Token: d.Token,
|
|
Adapters: d.Adapters,
|
|
Runtime: d.Runtime,
|
|
})
|
|
}
|
|
return s, nil
|
|
}
|
|
```
|
|
|
|
**`apps/edge/internal/node/registry.go` 변경:**
|
|
|
|
`NodeEntry`(L13-17) 변경:
|
|
|
|
Before:
|
|
```go
|
|
type NodeEntry struct {
|
|
NodeID string
|
|
Client *toki.TcpClient
|
|
Adapters []*iop.AdapterInfo
|
|
}
|
|
```
|
|
|
|
After:
|
|
```go
|
|
type NodeEntry struct {
|
|
NodeID string
|
|
Alias string
|
|
Client *toki.TcpClient
|
|
}
|
|
```
|
|
|
|
`Adapters []*iop.AdapterInfo` 제거 — adapter 정보는 edge의 NodeStore에서 관리하므로 런타임 registry에 중복 보관할 필요가 없다.
|
|
|
|
#### 수정 파일 및 체크리스트
|
|
|
|
- [x] `apps/edge/internal/node/store.go` — 신규 작성
|
|
- [x] `apps/edge/internal/node/registry.go` — NodeEntry에서 Adapters 제거, Alias 추가
|
|
|
|
#### 테스트 작성
|
|
|
|
신규 파일 `apps/edge/internal/node/store_test.go` 작성:
|
|
|
|
- `TestLoadFromConfig_Success`: alias 2개 정의 → `FindByToken` 조회 성공
|
|
- `TestLoadFromConfig_DuplicateToken`: 중복 token → 오류 반환
|
|
- `TestLoadFromConfig_EmptyToken`: token 빈 값 → 오류 반환
|
|
|
|
#### 중간 검증
|
|
|
|
```bash
|
|
go test ./apps/edge/internal/node/...
|
|
```
|
|
|
|
---
|
|
|
|
### [REFACTOR-4] edge: transport 핸드셰이크 변경 + bootstrap 재배선
|
|
|
|
#### 문제
|
|
|
|
`apps/edge/internal/transport/server.go`의 `onNodeConnected`(L74-112)이 연결 직후 `CapabilityRequest`를 보낸다. 새 설계에서는 node가 먼저 `RegisterRequest{token}`을 보내고 edge가 검증 후 `RegisterResponse{config}`를 응답한다.
|
|
|
|
`apps/edge/internal/bootstrap/module.go`에 `NodeStore`가 없다.
|
|
|
|
#### 해결 방법
|
|
|
|
**`apps/edge/internal/transport/server.go` 전면 변경:**
|
|
|
|
`Server` 구조체에 `nodeStore` 추가:
|
|
|
|
Before (L36-42):
|
|
```go
|
|
type Server struct {
|
|
tcp *toki.TcpServer
|
|
listen string
|
|
registry *edgenode.Registry
|
|
logger *zap.Logger
|
|
}
|
|
```
|
|
|
|
After:
|
|
```go
|
|
type Server struct {
|
|
tcp *toki.TcpServer
|
|
listen string
|
|
registry *edgenode.Registry
|
|
nodeStore *edgenode.NodeStore
|
|
logger *zap.Logger
|
|
}
|
|
```
|
|
|
|
`NewServer` 시그니처 변경:
|
|
|
|
Before (L44):
|
|
```go
|
|
func NewServer(listen string, registry *edgenode.Registry, logger *zap.Logger) (*Server, error) {
|
|
```
|
|
|
|
After:
|
|
```go
|
|
func NewServer(listen string, registry *edgenode.Registry, nodeStore *edgenode.NodeStore, logger *zap.Logger) (*Server, error) {
|
|
```
|
|
|
|
`edgeParserMap`(L23-34) 변경: `CapabilityResponse` 제거, `RegisterRequest` 추가:
|
|
|
|
Before:
|
|
```go
|
|
func edgeParserMap() toki.ParserMap {
|
|
return toki.ParserMap{
|
|
toki.TypeNameOf(&iop.RunEvent{}): func(b []byte) (proto.Message, error) {
|
|
m := &iop.RunEvent{}
|
|
return m, proto.Unmarshal(b, m)
|
|
},
|
|
toki.TypeNameOf(&iop.CapabilityResponse{}): func(b []byte) (proto.Message, error) {
|
|
m := &iop.CapabilityResponse{}
|
|
return m, proto.Unmarshal(b, m)
|
|
},
|
|
}
|
|
}
|
|
```
|
|
|
|
After:
|
|
```go
|
|
func edgeParserMap() toki.ParserMap {
|
|
return toki.ParserMap{
|
|
toki.TypeNameOf(&iop.RunEvent{}): func(b []byte) (proto.Message, error) {
|
|
m := &iop.RunEvent{}
|
|
return m, proto.Unmarshal(b, m)
|
|
},
|
|
toki.TypeNameOf(&iop.RegisterRequest{}): func(b []byte) (proto.Message, error) {
|
|
m := &iop.RegisterRequest{}
|
|
return m, proto.Unmarshal(b, m)
|
|
},
|
|
}
|
|
}
|
|
```
|
|
|
|
`onNodeConnected`(L74-112) 전면 재작성:
|
|
|
|
Before (핵심 로직):
|
|
```go
|
|
func (s *Server) onNodeConnected(client *toki.TcpClient) {
|
|
// ... RunEvent listener 등록 ...
|
|
go func() {
|
|
resp, err := toki.SendRequestTyped[*iop.CapabilityRequest, *iop.CapabilityResponse](...)
|
|
// ... registry.Register(entry) ...
|
|
}()
|
|
}
|
|
```
|
|
|
|
After:
|
|
```go
|
|
func (s *Server) onNodeConnected(client *toki.TcpClient) {
|
|
s.logger.Info("node connection established")
|
|
|
|
toki.AddListenerTyped[*iop.RunEvent](&client.Communicator, func(e *iop.RunEvent) {
|
|
s.logger.Debug("run event received",
|
|
zap.String("run_id", e.GetRunId()),
|
|
zap.String("type", e.GetType()),
|
|
)
|
|
})
|
|
|
|
toki.AddRequestListenerTyped[*iop.RegisterRequest, *iop.RegisterResponse](
|
|
&client.Communicator,
|
|
func(req *iop.RegisterRequest) (*iop.RegisterResponse, error) {
|
|
rec, ok := s.nodeStore.FindByToken(req.GetToken())
|
|
if !ok {
|
|
s.logger.Warn("unknown token", zap.String("token_prefix", safePrefix(req.GetToken())))
|
|
return &iop.RegisterResponse{Accepted: false, Reason: "unknown token"}, nil
|
|
}
|
|
|
|
entry := &edgenode.NodeEntry{
|
|
NodeID: rec.ID,
|
|
Alias: rec.Alias,
|
|
Client: client,
|
|
}
|
|
client.AddDisconnectListener(func(_ *toki.TcpClient) {
|
|
s.registry.Unregister(rec.ID)
|
|
s.logger.Info("node unregistered", zap.String("node_id", rec.ID))
|
|
})
|
|
s.registry.Register(entry)
|
|
s.logger.Info("node registered",
|
|
zap.String("node_id", rec.ID),
|
|
zap.String("alias", rec.Alias),
|
|
)
|
|
|
|
return &iop.RegisterResponse{
|
|
Accepted: true,
|
|
NodeId: rec.ID,
|
|
Alias: rec.Alias,
|
|
Config: buildConfigPayload(rec),
|
|
}, nil
|
|
},
|
|
)
|
|
}
|
|
|
|
func safePrefix(s string) string {
|
|
if len(s) > 8 {
|
|
return s[:8] + "..."
|
|
}
|
|
return s
|
|
}
|
|
|
|
func buildConfigPayload(rec *edgenode.NodeRecord) *iop.NodeConfigPayload {
|
|
payload := &iop.NodeConfigPayload{
|
|
Runtime: &iop.NodeRuntimeConfig{
|
|
Concurrency: int32(rec.Runtime.Concurrency),
|
|
WorkspaceRoot: rec.Runtime.WorkspaceRoot,
|
|
},
|
|
}
|
|
addAdapter := func(typ string, enabled bool, settings map[string]any) {
|
|
if !enabled {
|
|
return
|
|
}
|
|
st, _ := structpb.NewStruct(settings)
|
|
payload.Adapters = append(payload.Adapters, &iop.AdapterConfig{
|
|
Type: typ,
|
|
Enabled: true,
|
|
Settings: st,
|
|
})
|
|
}
|
|
addAdapter("mock", true, nil) // mock is always enabled
|
|
addAdapter("ollama", rec.Adapters.Ollama.Enabled, map[string]any{
|
|
"base_url": rec.Adapters.Ollama.BaseURL,
|
|
})
|
|
addAdapter("vllm", rec.Adapters.Vllm.Enabled, map[string]any{
|
|
"endpoint": rec.Adapters.Vllm.Endpoint,
|
|
})
|
|
if rec.Adapters.CLI.Enabled {
|
|
profiles := make(map[string]any)
|
|
for name, p := range rec.Adapters.CLI.Profiles {
|
|
profiles[name] = map[string]any{
|
|
"command": p.Command,
|
|
"args": p.Args,
|
|
"env": p.Env,
|
|
}
|
|
}
|
|
addAdapter("cli", true, map[string]any{"profiles": profiles})
|
|
}
|
|
return payload
|
|
}
|
|
```
|
|
|
|
필요한 import 추가: `"google.golang.org/protobuf/types/known/structpb"`
|
|
|
|
**`apps/edge/internal/bootstrap/module.go` 변경:**
|
|
|
|
`NewServer` 호출부에 `nodeStore` 추가, `NodeStore` 생성 Provider 추가:
|
|
|
|
Before (L28-29):
|
|
```go
|
|
func(cfg *config.EdgeConfig, reg *edgenode.Registry, logger *zap.Logger) (*transport.Server, error) {
|
|
return transport.NewServer(cfg.Server.Listen, reg, logger)
|
|
},
|
|
```
|
|
|
|
After:
|
|
```go
|
|
func() *edgenode.NodeStore { return edgenode.NewNodeStore() }, // placeholder; seeded in Invoke
|
|
|
|
func(cfg *config.EdgeConfig, reg *edgenode.Registry, ns *edgenode.NodeStore, logger *zap.Logger) (*transport.Server, error) {
|
|
return transport.NewServer(cfg.Server.Listen, reg, ns, logger)
|
|
},
|
|
```
|
|
|
|
`fx.Invoke`에서 `NodeStore`를 config으로 시드:
|
|
|
|
```go
|
|
fx.Invoke(func(lc fx.Lifecycle, srv *transport.Server, ns *edgenode.NodeStore, cfg *config.EdgeConfig, logger *zap.Logger) {
|
|
lc.Append(fx.Hook{
|
|
OnStart: func(ctx context.Context) error {
|
|
seeded, err := edgenode.LoadFromConfig(cfg.Nodes)
|
|
if err != nil {
|
|
return fmt.Errorf("edge: seed node store: %w", err)
|
|
}
|
|
for _, rec := range seeded.All() {
|
|
ns.Add(rec)
|
|
}
|
|
// ... start server + metrics ...
|
|
},
|
|
})
|
|
})
|
|
```
|
|
|
|
`NodeStore.All()` 메서드도 `store.go`에 추가한다:
|
|
```go
|
|
func (s *NodeStore) All() []*NodeRecord {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
out := make([]*NodeRecord, 0, len(s.byID))
|
|
for _, r := range s.byID {
|
|
out = append(out, r)
|
|
}
|
|
return out
|
|
}
|
|
```
|
|
|
|
#### 수정 파일 및 체크리스트
|
|
|
|
- [x] `apps/edge/internal/transport/server.go` — edgeParserMap 변경, Server 구조체 변경, NewServer 시그니처 변경, onNodeConnected 재작성, buildConfigPayload 추가, safePrefix 추가
|
|
- [x] `apps/edge/internal/bootstrap/module.go` — NodeStore provider 추가, NewServer 호출 갱신, OnStart에서 NodeStore 시드
|
|
- [x] `apps/edge/internal/node/store.go` — All() 메서드 추가 (REFACTOR-3에서 작성한 파일에 추가)
|
|
|
|
#### 테스트 작성
|
|
|
|
skip — integration_test.go는 REFACTOR-7에서 갱신한다.
|
|
|
|
#### 중간 검증
|
|
|
|
```bash
|
|
go build ./apps/edge/...
|
|
```
|
|
|
|
---
|
|
|
|
### [REFACTOR-5] node: transport 핸드셰이크 변경 + DialEdge 시그니처 변경
|
|
|
|
#### 문제
|
|
|
|
`apps/node/internal/transport/session.go`의 `Handler` 인터페이스(L16-19)에 `OnCapabilityRequest`가 있고, `newSession`(L29-64)에 `CapabilityRequest` 리스너가 등록되어 있다. `DialEdge`(client.go L19-37)가 handler와 nodeID를 파라미터로 받는다.
|
|
|
|
새 설계에서 node는 연결 직후 `RegisterRequest{token}`을 보내고 `RegisterResponse`를 받은 뒤에야 `node.Node`를 생성한다. `DialEdge` 시점에는 아직 handler가 없으므로 Session이 handler를 나중에 주입받을 수 있어야 한다.
|
|
|
|
#### 해결 방법
|
|
|
|
**`apps/node/internal/transport/session.go` 변경:**
|
|
|
|
`Handler` 인터페이스(L16-19)에서 `OnCapabilityRequest` 제거:
|
|
|
|
Before:
|
|
```go
|
|
type Handler interface {
|
|
OnRunRequest(ctx context.Context, sess *Session, req *iop.RunRequest) error
|
|
OnCapabilityRequest(ctx context.Context, sess *Session) (*iop.CapabilityResponse, error)
|
|
OnCancel(ctx context.Context, sess *Session, req *iop.CancelRequest) error
|
|
}
|
|
```
|
|
|
|
After:
|
|
```go
|
|
type Handler interface {
|
|
OnRunRequest(ctx context.Context, sess *Session, req *iop.RunRequest) error
|
|
OnCancel(ctx context.Context, sess *Session, req *iop.CancelRequest) error
|
|
}
|
|
```
|
|
|
|
`Session` 구조체에 mutable handler 추가 (nodeID 제거):
|
|
|
|
Before (L22-27):
|
|
```go
|
|
type Session struct {
|
|
client *toki.TcpClient
|
|
nodeID string
|
|
logger *zap.Logger
|
|
cancelFns sync.Map
|
|
}
|
|
```
|
|
|
|
After:
|
|
```go
|
|
type Session struct {
|
|
client *toki.TcpClient
|
|
logger *zap.Logger
|
|
mu sync.RWMutex
|
|
handler Handler
|
|
cancelFns sync.Map
|
|
}
|
|
```
|
|
|
|
`SetHandler` 메서드 추가:
|
|
```go
|
|
// SetHandler attaches the message handler. Called after registration completes.
|
|
func (s *Session) SetHandler(h Handler) {
|
|
s.mu.Lock()
|
|
s.handler = h
|
|
s.mu.Unlock()
|
|
}
|
|
```
|
|
|
|
`newSession` 시그니처에서 `handler Handler`, `nodeID string` 제거, CapabilityRequest 리스너 제거, handler 호출 시 nil 체크:
|
|
|
|
Before (L29):
|
|
```go
|
|
func newSession(client *toki.TcpClient, handler Handler, nodeID string, logger *zap.Logger) *Session {
|
|
```
|
|
|
|
After:
|
|
```go
|
|
func newSession(client *toki.TcpClient, logger *zap.Logger) *Session {
|
|
s := &Session{client: client, logger: logger}
|
|
|
|
toki.AddListenerTyped[*iop.RunRequest](&client.Communicator, func(req *iop.RunRequest) {
|
|
go func() {
|
|
s.mu.RLock()
|
|
h := s.handler
|
|
s.mu.RUnlock()
|
|
if h == nil {
|
|
return
|
|
}
|
|
if err := h.OnRunRequest(context.Background(), s, req); err != nil {
|
|
logger.Warn("run request error",
|
|
zap.String("run_id", req.GetRunId()),
|
|
zap.Error(err),
|
|
)
|
|
}
|
|
}()
|
|
})
|
|
|
|
toki.AddListenerTyped[*iop.CancelRequest](&client.Communicator, func(req *iop.CancelRequest) {
|
|
s.mu.RLock()
|
|
h := s.handler
|
|
s.mu.RUnlock()
|
|
if h == nil {
|
|
return
|
|
}
|
|
if err := h.OnCancel(context.Background(), s, req); err != nil {
|
|
logger.Warn("cancel error",
|
|
zap.String("run_id", req.GetRunId()),
|
|
zap.Error(err),
|
|
)
|
|
}
|
|
})
|
|
|
|
client.AddDisconnectListener(func(_ *toki.TcpClient) {
|
|
logger.Info("disconnected from edge")
|
|
})
|
|
|
|
return s
|
|
}
|
|
```
|
|
|
|
**`apps/node/internal/transport/parser.go` 변경:**
|
|
|
|
`nodeParserMap`에서 `CapabilityRequest` 제거, `RegisterResponse` 추가:
|
|
|
|
Before:
|
|
```go
|
|
func nodeParserMap() toki.ParserMap {
|
|
return toki.ParserMap{
|
|
toki.TypeNameOf(&iop.RunRequest{}): ...,
|
|
toki.TypeNameOf(&iop.CancelRequest{}): ...,
|
|
toki.TypeNameOf(&iop.CapabilityRequest{}): ...,
|
|
}
|
|
}
|
|
```
|
|
|
|
After:
|
|
```go
|
|
func nodeParserMap() toki.ParserMap {
|
|
return toki.ParserMap{
|
|
toki.TypeNameOf(&iop.RunRequest{}): func(b []byte) (proto.Message, error) {
|
|
m := &iop.RunRequest{}
|
|
return m, proto.Unmarshal(b, m)
|
|
},
|
|
toki.TypeNameOf(&iop.CancelRequest{}): func(b []byte) (proto.Message, error) {
|
|
m := &iop.CancelRequest{}
|
|
return m, proto.Unmarshal(b, m)
|
|
},
|
|
toki.TypeNameOf(&iop.RegisterResponse{}): func(b []byte) (proto.Message, error) {
|
|
m := &iop.RegisterResponse{}
|
|
return m, proto.Unmarshal(b, m)
|
|
},
|
|
}
|
|
}
|
|
```
|
|
|
|
**`apps/node/internal/transport/client.go` 전면 재작성:**
|
|
|
|
Before (전체):
|
|
```go
|
|
const (
|
|
heartbeatIntervalSec = 30
|
|
heartbeatWaitSec = 10
|
|
)
|
|
|
|
func DialEdge(ctx context.Context, addr string, handler Handler, nodeID string, logger *zap.Logger) (*Session, error) {
|
|
// ...
|
|
sess := newSession(client, handler, nodeID, logger)
|
|
// ...
|
|
return sess, nil
|
|
}
|
|
```
|
|
|
|
After:
|
|
```go
|
|
package transport
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"net"
|
|
"strconv"
|
|
"time"
|
|
|
|
toki "git.toki-labs.com/toki/common-proto-socket/go"
|
|
"go.uber.org/zap"
|
|
|
|
iop "iop/proto/gen/iop"
|
|
)
|
|
|
|
const (
|
|
heartbeatIntervalSec = 30
|
|
heartbeatWaitSec = 10
|
|
registerTimeout = 10 * time.Second
|
|
)
|
|
|
|
// RegisterResult is returned by DialEdge after successful registration.
|
|
type RegisterResult struct {
|
|
Session *Session
|
|
NodeID string
|
|
Alias string
|
|
Config *iop.NodeConfigPayload
|
|
}
|
|
|
|
// DialEdge connects to edge, performs the registration handshake, and returns
|
|
// a RegisterResult. Call result.Session.SetHandler after creating node.Node.
|
|
func DialEdge(ctx context.Context, addr, token string, logger *zap.Logger) (*RegisterResult, error) {
|
|
host, portStr, err := net.SplitHostPort(addr)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("transport: invalid addr %q: %w", addr, err)
|
|
}
|
|
port, err := strconv.Atoi(portStr)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("transport: invalid port %q: %w", portStr, err)
|
|
}
|
|
|
|
client, err := toki.DialTcp(ctx, host, port, heartbeatIntervalSec, heartbeatWaitSec, nodeParserMap())
|
|
if err != nil {
|
|
return nil, fmt.Errorf("transport: dial edge %s: %w", addr, err)
|
|
}
|
|
|
|
resp, err := toki.SendRequestTyped[*iop.RegisterRequest, *iop.RegisterResponse](
|
|
&client.Communicator,
|
|
&iop.RegisterRequest{Token: token},
|
|
registerTimeout,
|
|
)
|
|
if err != nil {
|
|
_ = client.Close()
|
|
return nil, fmt.Errorf("transport: register: %w", err)
|
|
}
|
|
if !resp.GetAccepted() {
|
|
_ = client.Close()
|
|
return nil, fmt.Errorf("transport: register rejected: %s", resp.GetReason())
|
|
}
|
|
|
|
sess := newSession(client, logger)
|
|
logger.Info("registered with edge",
|
|
zap.String("node_id", resp.GetNodeId()),
|
|
zap.String("alias", resp.GetAlias()),
|
|
)
|
|
return &RegisterResult{
|
|
Session: sess,
|
|
NodeID: resp.GetNodeId(),
|
|
Alias: resp.GetAlias(),
|
|
Config: resp.GetConfig(),
|
|
}, nil
|
|
}
|
|
```
|
|
|
|
#### 수정 파일 및 체크리스트
|
|
|
|
- [x] `apps/node/internal/transport/session.go` — Handler 인터페이스에서 OnCapabilityRequest 제거, Session에 mu + handler 추가 (nodeID 제거), SetHandler 추가, newSession 재작성
|
|
- [x] `apps/node/internal/transport/parser.go` — CapabilityRequest 제거, RegisterResponse 추가
|
|
- [x] `apps/node/internal/transport/client.go` — RegisterResult 추가, DialEdge 전면 재작성
|
|
|
|
#### 테스트 작성
|
|
|
|
skip — integration_test.go는 REFACTOR-7에서 갱신한다.
|
|
|
|
#### 중간 검증
|
|
|
|
```bash
|
|
go build ./apps/node/internal/transport/...
|
|
```
|
|
|
|
---
|
|
|
|
### [REFACTOR-6] node: adapter factory + bootstrap 재구성
|
|
|
|
#### 문제
|
|
|
|
`apps/node/internal/bootstrap/module.go`가 `config.NodeConfig`의 어댑터 설정으로 어댑터를 정적 초기화한다(L38-54). 새 설계에서는 `DialEdge` 이후 수신된 `NodeConfigPayload`를 기반으로 동적 초기화해야 한다. 또한 `node.New`(node.go L31-45)가 `*config.NodeConfig`를 받아 `nodeID`를 추출하는데, ID는 이제 registration 결과에서 온다.
|
|
|
|
어댑터(ollama, vllm, cli)의 `New()` 함수가 `config.XxxConf`를 받으므로 proto payload로부터 어댑터를 생성하는 factory가 필요하다.
|
|
|
|
#### 해결 방법
|
|
|
|
**`apps/node/internal/node/node.go` 변경:**
|
|
|
|
`node.New` 시그니처에서 `*config.NodeConfig` 제거, `nodeID string` 직접 파라미터로 변경:
|
|
|
|
Before (L31-45):
|
|
```go
|
|
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,
|
|
}
|
|
}
|
|
```
|
|
|
|
After:
|
|
```go
|
|
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,
|
|
}
|
|
}
|
|
```
|
|
|
|
`OnCapabilityRequest` 메서드(L112-127) 삭제 — Handler 인터페이스에서 제거되었으므로 node.Node도 구현할 필요 없다.
|
|
|
|
`sessionSink`와 `structAsMap`은 그대로 유지한다.
|
|
|
|
**신규 파일 `apps/node/internal/adapters/factory.go` 작성:**
|
|
|
|
```go
|
|
package adapters
|
|
|
|
import (
|
|
"fmt"
|
|
|
|
"go.uber.org/zap"
|
|
"google.golang.org/protobuf/types/known/structpb"
|
|
|
|
"iop/apps/node/internal/adapters/cli"
|
|
"iop/apps/node/internal/adapters/mock"
|
|
"iop/apps/node/internal/adapters/ollama"
|
|
"iop/apps/node/internal/adapters/vllm"
|
|
"iop/packages/config"
|
|
iop "iop/proto/gen/iop"
|
|
)
|
|
|
|
// BuildFromPayload creates a Registry from a NodeConfigPayload received from edge.
|
|
func BuildFromPayload(payload *iop.NodeConfigPayload, logger *zap.Logger) (*Registry, error) {
|
|
reg := NewRegistry()
|
|
reg.Register(mock.New(logger)) // mock is always available
|
|
|
|
for _, ac := range payload.GetAdapters() {
|
|
if !ac.GetEnabled() {
|
|
continue
|
|
}
|
|
switch ac.GetType() {
|
|
case "mock":
|
|
// already registered above
|
|
case "ollama":
|
|
cfg := ollamaConfFromStruct(ac.GetSettings())
|
|
reg.Register(ollama.New(cfg, logger))
|
|
case "vllm":
|
|
cfg := vllmConfFromStruct(ac.GetSettings())
|
|
reg.Register(vllm.New(cfg, logger))
|
|
case "cli":
|
|
cfg := cliConfFromStruct(ac.GetSettings())
|
|
reg.Register(cli.New(cfg, logger))
|
|
default:
|
|
return nil, fmt.Errorf("adapters: unknown adapter type %q", ac.GetType())
|
|
}
|
|
}
|
|
return reg, nil
|
|
}
|
|
|
|
func ollamaConfFromStruct(s *structpb.Struct) config.OllamaConf {
|
|
if s == nil {
|
|
return config.OllamaConf{}
|
|
}
|
|
m := s.AsMap()
|
|
cfg := config.OllamaConf{Enabled: true}
|
|
if v, ok := m["base_url"].(string); ok {
|
|
cfg.BaseURL = v
|
|
}
|
|
return cfg
|
|
}
|
|
|
|
func vllmConfFromStruct(s *structpb.Struct) config.VllmConf {
|
|
if s == nil {
|
|
return config.VllmConf{}
|
|
}
|
|
m := s.AsMap()
|
|
cfg := config.VllmConf{Enabled: true}
|
|
if v, ok := m["endpoint"].(string); ok {
|
|
cfg.Endpoint = v
|
|
}
|
|
return cfg
|
|
}
|
|
|
|
func cliConfFromStruct(s *structpb.Struct) config.CLIConf {
|
|
if s == nil {
|
|
return config.CLIConf{}
|
|
}
|
|
m := s.AsMap()
|
|
cfg := config.CLIConf{Enabled: true, Profiles: make(map[string]config.CLIProfileConf)}
|
|
if profiles, ok := m["profiles"].(map[string]any); ok {
|
|
for name, pAny := range profiles {
|
|
p, ok := pAny.(map[string]any)
|
|
if !ok {
|
|
continue
|
|
}
|
|
prof := config.CLIProfileConf{}
|
|
if cmd, ok := p["command"].(string); ok {
|
|
prof.Command = cmd
|
|
}
|
|
if args, ok := p["args"].([]any); ok {
|
|
for _, a := range args {
|
|
if s, ok := a.(string); ok {
|
|
prof.Args = append(prof.Args, s)
|
|
}
|
|
}
|
|
}
|
|
if envs, ok := p["env"].([]any); ok {
|
|
for _, e := range envs {
|
|
if s, ok := e.(string); ok {
|
|
prof.Env = append(prof.Env, s)
|
|
}
|
|
}
|
|
}
|
|
cfg.Profiles[name] = prof
|
|
}
|
|
}
|
|
return cfg
|
|
}
|
|
```
|
|
|
|
**`apps/node/internal/bootstrap/module.go` 전면 재구성:**
|
|
|
|
fx provider 방식 대신 OnStart 훅 내부에서 순차 초기화:
|
|
|
|
After (전체 재작성):
|
|
```go
|
|
package bootstrap
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
|
|
"go.uber.org/fx"
|
|
"go.uber.org/zap"
|
|
|
|
"iop/apps/node/internal/adapters"
|
|
"iop/apps/node/internal/node"
|
|
"iop/apps/node/internal/router"
|
|
"iop/apps/node/internal/store"
|
|
"iop/apps/node/internal/transport"
|
|
"iop/packages/config"
|
|
"iop/packages/observability"
|
|
)
|
|
|
|
func Module(cfg *config.NodeConfig) fx.Option {
|
|
return fx.Options(
|
|
fx.Provide(
|
|
func() *config.NodeConfig { return cfg },
|
|
|
|
func(cfg *config.NodeConfig) (*zap.Logger, error) {
|
|
return observability.NewLogger(cfg.Logging.Level)
|
|
},
|
|
),
|
|
|
|
fx.Invoke(func(lc fx.Lifecycle, cfg *config.NodeConfig, logger *zap.Logger) {
|
|
var sess *transport.Session
|
|
lc.Append(fx.Hook{
|
|
OnStart: func(ctx context.Context) error {
|
|
result, err := transport.DialEdge(ctx, cfg.Transport.EdgeAddr, cfg.Transport.Token, logger)
|
|
if err != nil {
|
|
return fmt.Errorf("bootstrap: dial edge: %w", err)
|
|
}
|
|
|
|
reg, err := adapters.BuildFromPayload(result.Config, logger)
|
|
if err != nil {
|
|
_ = result.Session.Close()
|
|
return fmt.Errorf("bootstrap: build adapters: %w", err)
|
|
}
|
|
|
|
dsn := "file:iop.db?cache=shared&mode=rwc"
|
|
if result.Config.GetRuntime().GetWorkspaceRoot() != "" {
|
|
dsn = "file:" + result.Config.GetRuntime().GetWorkspaceRoot() + "/iop.db?cache=shared&mode=rwc"
|
|
}
|
|
st, err := store.New(dsn, logger)
|
|
if err != nil {
|
|
_ = result.Session.Close()
|
|
return fmt.Errorf("bootstrap: store: %w", err)
|
|
}
|
|
|
|
rtr := router.New(reg, logger)
|
|
n := node.New(result.NodeID, rtr, reg, st, logger)
|
|
result.Session.SetHandler(n)
|
|
sess = result.Session
|
|
|
|
go func() {
|
|
if err := observability.ServeMetrics(cfg.Metrics.Port); err != nil {
|
|
logger.Warn("metrics server exited", zap.Error(err))
|
|
}
|
|
}()
|
|
return nil
|
|
},
|
|
OnStop: func(_ context.Context) error {
|
|
if sess != nil {
|
|
return sess.Close()
|
|
}
|
|
return nil
|
|
},
|
|
})
|
|
}),
|
|
)
|
|
}
|
|
```
|
|
|
|
#### 수정 파일 및 체크리스트
|
|
|
|
- [x] `apps/node/internal/node/node.go` — New() 시그니처에서 *config.NodeConfig 제거 → nodeID string 직접 파라미터, OnCapabilityRequest 메서드 삭제, config import 제거
|
|
- [x] `apps/node/internal/adapters/factory.go` — 신규 작성
|
|
- [x] `apps/node/internal/bootstrap/module.go` — 전면 재구성 (OnStart 내부 순차 초기화)
|
|
- [x] `apps/node/cmd/node/main.go` — config print 커맨드에서 더 이상 존재하지 않는 필드 참조 제거 여부 확인 (NodeConfig 구조체가 단순해졌으므로 코드 자체는 동작하지만, 출력 내용이 달라짐)
|
|
|
|
#### 테스트 작성
|
|
|
|
skip — node 단위 테스트는 REFACTOR-7에서 갱신한다.
|
|
|
|
#### 중간 검증
|
|
|
|
```bash
|
|
go build ./apps/node/...
|
|
```
|
|
|
|
---
|
|
|
|
### [REFACTOR-7] 테스트 전체 갱신
|
|
|
|
#### 문제
|
|
|
|
모든 테스트가 현재 handshake(CapabilityRequest 기반)와 이전 구조(NodeConfig, OnCapabilityRequest)를 기준으로 작성되어 있다.
|
|
|
|
#### 해결 방법
|
|
|
|
**`apps/edge/internal/transport/integration_test.go` 재작성:**
|
|
|
|
- `TestEdgeServerIntegration`: CapabilityRequest/Response 흐름 제거 → mock node가 연결 후 `RegisterRequest{token}`을 보내고 `RegisterResponse` 수신 확인. `waitForRegistryEntry` 조건을 RegisterResponse 수신 이후로 변경.
|
|
- mock nodeParser에 `RegisterResponse` 파서 추가, `CapabilityRequest` 파서 제거.
|
|
- edge 서버 생성 시 `NodeStore`를 주입 (`LoadFromConfig`로 test token 등록).
|
|
|
|
**`apps/node/internal/transport/integration_test.go` 재작성:**
|
|
|
|
- `mockHandler`에서 `OnCapabilityRequest` 제거.
|
|
- `TestNodeClientIntegration`: mock edge가 `RegisterRequest`를 받아 `RegisterResponse`로 응답. CapabilityRequest 전송 부분(L116-127) 제거.
|
|
- mock edgeParser에 `RegisterRequest` 추가, `CapabilityResponse` 제거.
|
|
- `transport.DialEdge` 새 시그니처(`addr, token string`, handler 없음) 반영.
|
|
|
|
**`apps/node/internal/node/node_test.go` 갱신:**
|
|
|
|
- `TestOnCapabilityRequest`(L82-107) 삭제 — Handler 인터페이스에서 제거됨.
|
|
- `makeNode`(L65-80): `cfg := &config.NodeConfig{...}` → `node.New("test-node", ...)` 직접 호출로 변경.
|
|
- `TestOnRunRequest_RouterError`, `TestOnRunRequest_AdapterNotFound`, `TestOnRunRequest_Success`, `TestOnCancel_CallsCancelFn` — node.New 시그니처 변경만 반영, 로직은 그대로 유지.
|
|
|
|
**`apps/edge/internal/node/registry_test.go` 갱신:**
|
|
|
|
- `TestRegistry_RegisterAndCount`(L12-22): `NodeEntry`에서 `Adapters` 필드 제거됨 → 초기화 코드에서 `Adapters` 제거.
|
|
|
|
#### 수정 파일 및 체크리스트
|
|
|
|
- [x] `apps/edge/internal/transport/integration_test.go` — RegisterRequest/Response 기반으로 재작성
|
|
- [x] `apps/node/internal/transport/integration_test.go` — RegisterRequest/Response 기반으로 재작성, DialEdge 시그니처 반영
|
|
- [x] `apps/node/internal/node/node_test.go` — TestOnCapabilityRequest 삭제, makeNode 수정, 나머지 테스트 node.New 시그니처 반영
|
|
- [x] `apps/edge/internal/node/registry_test.go` — NodeEntry Adapters 제거 반영
|
|
|
|
#### 테스트 작성
|
|
|
|
위 갱신 외 신규 테스트:
|
|
- `apps/node/internal/adapters/factory_test.go`:
|
|
- `TestBuildFromPayload_MockAlwaysPresent`: 빈 payload → mock adapter만 존재
|
|
- `TestBuildFromPayload_OllamaEnabled`: ollama enabled 포함 payload → ollama + mock
|
|
- `TestBuildFromPayload_UnknownType`: 알 수 없는 type → 오류
|
|
|
|
#### 중간 검증
|
|
|
|
```bash
|
|
go test ./apps/edge/internal/node/...
|
|
go test ./apps/node/internal/node/...
|
|
go test ./apps/node/internal/adapters/...
|
|
```
|
|
|
|
---
|
|
|
|
## 수정 파일 요약
|
|
|
|
| 파일 | 항목 |
|
|
|------|------|
|
|
| `proto/iop/runtime.proto` | REFACTOR-1 |
|
|
| `proto/gen/iop/runtime.pb.go` | REFACTOR-1 (생성) |
|
|
| `packages/config/config.go` | REFACTOR-2 |
|
|
| `configs/node.yaml` | REFACTOR-2 |
|
|
| `configs/edge.yaml` | REFACTOR-2 |
|
|
| `apps/edge/internal/node/store.go` (신규) | REFACTOR-3 |
|
|
| `apps/edge/internal/node/registry.go` | REFACTOR-3 |
|
|
| `apps/edge/internal/transport/server.go` | REFACTOR-4 |
|
|
| `apps/edge/internal/bootstrap/module.go` | REFACTOR-4 |
|
|
| `apps/node/internal/transport/session.go` | REFACTOR-5 |
|
|
| `apps/node/internal/transport/parser.go` | REFACTOR-5 |
|
|
| `apps/node/internal/transport/client.go` | REFACTOR-5 |
|
|
| `apps/node/internal/node/node.go` | REFACTOR-6 |
|
|
| `apps/node/internal/adapters/factory.go` (신규) | REFACTOR-6 |
|
|
| `apps/node/internal/bootstrap/module.go` | REFACTOR-6 |
|
|
| `apps/edge/internal/transport/integration_test.go` | REFACTOR-7 |
|
|
| `apps/node/internal/transport/integration_test.go` | REFACTOR-7 |
|
|
| `apps/node/internal/node/node_test.go` | REFACTOR-7 |
|
|
| `apps/edge/internal/node/registry_test.go` | REFACTOR-7 |
|
|
| `apps/node/internal/adapters/factory_test.go` (신규) | REFACTOR-7 |
|
|
| `apps/edge/internal/node/store_test.go` (신규) | REFACTOR-3 |
|
|
|
|
## 최종 검증
|
|
|
|
```bash
|
|
make proto # proto 재생성 확인
|
|
go build ./... # 전체 컴파일
|
|
go test ./... # 전체 테스트
|
|
```
|
|
|
|
모든 테스트 PASS, 컴파일 오류 없음이 기대 결과이다.
|