refactor: update-plane milestone archive and edge node registry refactor
- Archive runtime-reconnect-config-refresh milestone and SDD documents - Refactor edge node registry and transport server - Update PHASE.md for update-plane-self-update-foundation
This commit is contained in:
parent
7d69b9f2a1
commit
0210e20eab
6 changed files with 131 additions and 39 deletions
|
|
@ -12,7 +12,7 @@ Edge 재시작이나 일시 단절 뒤에도 Node가 지정된 retry 정책으
|
|||
|
||||
## 상태
|
||||
|
||||
[진행중]
|
||||
[완료]
|
||||
|
||||
## 승격 조건
|
||||
|
||||
|
|
@ -75,20 +75,21 @@ Edge가 refresh된 NodeConfigPayload를 연결 Node에 전달하고, Node는 변
|
|||
|
||||
현재 dev field 환경의 bootstrap 테스트를 refresh/reconnect 완료 증거로 유지한다.
|
||||
|
||||
- [ ] [field-reconnect] `toki-labs.com` dev-runtime에서 Edge restart 후 GX10 Linux Node와 OneXPlayer Windows Node가 bootstrap 재실행 없이 재접속한다. 검증: Control Plane status와 Edge node TCP connection 수가 3개 Node 기준으로 회복된다.
|
||||
- [ ] [field-refresh] `gx10-vllm=4`, `onexplayer-lemonade=3` capacity 변경을 refresh로 반영하고 Edge process restart 없이 `qwen3.6:35b` 동시 4개 요청이 성공한다. 검증: aggregate 기대 배정은 capacity 4/3 기준 `gx10-vllm` 2개, `onexplayer-lemonade` 2개이며, 요청별 node id 확정은 dispatch trace가 구현된 경우에만 판정한다.
|
||||
- [ ] [docs] bootstrap/reconnect/refresh 운영 방법과 제한 사항이 `apps/edge/README.md`, `docs/edge-local-dev-guide.md`, agent-test field 기준에 반영된다.
|
||||
- [x] [field-reconnect] `toki-labs.com` dev-runtime에서 Edge restart 후 GX10 Linux Node와 OneXPlayer Windows Node가 bootstrap 재실행 없이 재접속한다. 검증: Control Plane status와 Edge node TCP connection 수가 3개 Node 기준으로 회복된다.
|
||||
- [x] [field-refresh] `gx10-vllm=4`, `onexplayer-lemonade=3` capacity 변경을 refresh로 반영하고 Edge process restart 없이 `qwen3.6:35b` 동시 4개 요청이 성공한다. 검증: aggregate 기대 배정은 capacity 4/3 기준 `gx10-vllm` 2개, `onexplayer-lemonade` 2개이며, 요청별 node id 확정은 dispatch trace가 구현된 경우에만 판정한다.
|
||||
- [x] [docs] bootstrap/reconnect/refresh 운영 방법과 제한 사항이 `apps/edge/README.md`, `docs/edge-local-dev-guide.md`, agent-test field 기준에 반영된다.
|
||||
|
||||
## 완료 리뷰
|
||||
|
||||
- 상태: 없음
|
||||
- 요청일: 없음
|
||||
- 완료 근거: `agent-task/archive/2026/06/m-runtime-reconnect-config-refresh/01_reconnect_supervisor/complete.log`의 `Roadmap Completion`에 따라 `retry-policy`, `supervisor`, `exit-limit`가 PASS로 확인되었다. `dedupe`는 `TestEdgeServerDuplicateRegistrationReason`와 `TestEdgeServerReconnectAfterUnregisterAccepted` 존재 및 `go test -count=1 ./apps/edge/internal/transport` PASS로 확인되었다. `02_refresh_contract/complete.log`와 `03+02_edge_refresh_apply/complete.log`의 `Roadmap Completion`에 따라 `entrypoint`, `diff-classify`, `atomic-apply`, `ops-report`가 PASS로 확인되었다. `04+02_config_push_contract/complete.log`, `05+04_node_adapter_diff/complete.log`, `06+05_node_drain_runtime_config/complete.log`의 `Roadmap Completion`에 따라 `push-contract`, `adapter-diff`, `drain`, `runtime-config`가 PASS로 확인되었다.
|
||||
- 상태: 통과
|
||||
- 요청일: 2026-06-24
|
||||
- 완료 근거: `agent-task/archive/2026/06/m-runtime-reconnect-config-refresh/01_reconnect_supervisor/complete.log`의 `Roadmap Completion`에 따라 `retry-policy`, `supervisor`, `exit-limit`가 PASS로 확인되었다. `dedupe`는 `TestEdgeServerDuplicateRegistrationReason`, `TestEdgeServerReconnectAfterUnregisterAccepted`, `TestRegistryRegisterIfAbsentRejectsConcurrentDuplicate`, `TestRegistryUnregisterIfClientIgnoresStaleConnection` 존재와 `go test -count=1 ./apps/edge/internal/node ./apps/edge/internal/transport` PASS로 확인되었다. `02_refresh_contract/complete.log`와 `03+02_edge_refresh_apply/complete.log`의 `Roadmap Completion`에 따라 `entrypoint`, `diff-classify`, `atomic-apply`, `ops-report`가 PASS로 확인되었다. `04+02_config_push_contract/complete.log`, `05+04_node_adapter_diff/complete.log`, `06+05_node_drain_runtime_config/complete.log`의 `Roadmap Completion`에 따라 `push-contract`, `adapter-diff`, `drain`, `runtime-config`가 PASS로 확인되었다. `07+01,03,04,05,06_field_docs_smoke/complete.log`의 `Roadmap Completion`에 따라 `field-reconnect`, `field-refresh`, `docs`가 PASS로 확인되었다. 최종 종결 감사에서 `go test -count=1 ./...`, `./scripts/e2e-smoke.sh`, `git diff --check`가 PASS였다.
|
||||
- 검토 항목:
|
||||
- [x] SDD gate가 승인되어 구현 잠금이 해제되었다
|
||||
- [ ] 모든 기능 Task와 Task 안의 검증이 충족되었다
|
||||
- [ ] 사용자가 field dev 환경 결과를 확인했다
|
||||
- 리뷰 코멘트: reconnect, refresh, node-config Epic Task는 완료 근거가 확인되어 완료 처리했다. 남은 미완료 Task는 `field-reconnect`, `field-refresh`, `docs`이며 field dev 환경 확인과 문서/agent-test 반영 완료 근거가 필요하다.
|
||||
- [x] 모든 기능 Task와 Task 안의 검증이 충족되었다
|
||||
- [x] field dev 환경 결과와 문서/agent-test 반영 근거가 완료 로그와 리뷰 PASS로 확인되었다
|
||||
- 남은 차단 항목: 없음
|
||||
- 리뷰 코멘트: 코드 레벨 종결 감사 중 같은 node id의 동시 등록 경쟁 구간을 확인해 `RegisterIfAbsent`와 client identity 기반 unregister로 보강했다. 모든 Milestone 기능 Task와 검증, SDD gate, field evidence, 최종 회귀가 충족되어 완료 archive 대상으로 판단한다.
|
||||
|
||||
## 범위 제외
|
||||
|
||||
|
|
@ -16,8 +16,8 @@ Control Plane은 release manifest, desired version, rollout/audit view를 제공
|
|||
완료, 검토중, 진행중, 계획, 스케치 순서로 두어 아래로 갈수록 미래 작업에 가까워지게 정렬한다.
|
||||
스케치 Milestone은 아직 구현 가능한 계획이 아니므로 사용자 검토와 구체화 후 `[계획]`으로 승격한다.
|
||||
|
||||
- [진행중] Edge/Node Runtime Reconnect와 Config Refresh
|
||||
- 경로: `agent-roadmap/phase/update-plane-self-update-foundation/milestones/runtime-reconnect-config-refresh.md`
|
||||
- [완료] Edge/Node Runtime Reconnect와 Config Refresh
|
||||
- 경로: `agent-roadmap/archive/phase/update-plane-self-update-foundation/milestones/runtime-reconnect-config-refresh.md`
|
||||
- 요약: Edge 단절 후 Node 10초 간격 10회 재접속과 종료 정책, 운영 중 config refresh/diff/apply/Node 전파의 MVP 경계를 구현 가능한 계획으로 정리한다.
|
||||
|
||||
- [스케치] Update Plane 안정 프로토콜
|
||||
|
|
|
|||
|
|
@ -52,6 +52,23 @@ func NewRegistry() *Registry {
|
|||
func (r *Registry) Register(entry *NodeEntry) {
|
||||
r.mu.Lock()
|
||||
defer r.mu.Unlock()
|
||||
r.registerLocked(entry)
|
||||
}
|
||||
|
||||
// RegisterIfAbsent registers entry only when the node id is not already
|
||||
// connected. The check and insert happen under one lock so concurrent duplicate
|
||||
// registration attempts cannot both be accepted by the transport server.
|
||||
func (r *Registry) RegisterIfAbsent(entry *NodeEntry) bool {
|
||||
r.mu.Lock()
|
||||
defer r.mu.Unlock()
|
||||
if _, exists := r.byID[entry.NodeID]; exists {
|
||||
return false
|
||||
}
|
||||
r.registerLocked(entry)
|
||||
return true
|
||||
}
|
||||
|
||||
func (r *Registry) registerLocked(entry *NodeEntry) {
|
||||
if entry.AgentKind == "" {
|
||||
entry.AgentKind = config.AgentKindGenericNode
|
||||
}
|
||||
|
|
@ -109,14 +126,32 @@ func (r *Registry) Unregister(nodeID string) {
|
|||
r.mu.Lock()
|
||||
defer r.mu.Unlock()
|
||||
if entry, ok := r.byID[nodeID]; ok {
|
||||
if entry.Alias != "" {
|
||||
delete(r.byAlias, entry.Alias)
|
||||
}
|
||||
delete(r.byIndex, entry.Index)
|
||||
delete(r.byID, nodeID)
|
||||
r.unregisterLocked(nodeID, entry)
|
||||
}
|
||||
}
|
||||
|
||||
// UnregisterIfClient removes nodeID only when the currently registered entry
|
||||
// belongs to client. Late disconnect callbacks from rejected or superseded
|
||||
// connections must not clear the live registry entry.
|
||||
func (r *Registry) UnregisterIfClient(nodeID string, client *toki.TcpClient) bool {
|
||||
r.mu.Lock()
|
||||
defer r.mu.Unlock()
|
||||
entry, ok := r.byID[nodeID]
|
||||
if !ok || entry.Client != client {
|
||||
return false
|
||||
}
|
||||
r.unregisterLocked(nodeID, entry)
|
||||
return true
|
||||
}
|
||||
|
||||
func (r *Registry) unregisterLocked(nodeID string, entry *NodeEntry) {
|
||||
if entry.Alias != "" {
|
||||
delete(r.byAlias, entry.Alias)
|
||||
}
|
||||
delete(r.byIndex, entry.Index)
|
||||
delete(r.byID, nodeID)
|
||||
}
|
||||
|
||||
func (r *Registry) Get(nodeID string) (*NodeEntry, bool) {
|
||||
r.mu.RLock()
|
||||
defer r.mu.RUnlock()
|
||||
|
|
|
|||
|
|
@ -1,6 +1,8 @@
|
|||
package node_test
|
||||
|
||||
import (
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
|
||||
toki "git.toki-labs.com/toki/proto-socket/go"
|
||||
|
|
@ -22,6 +24,62 @@ func TestRegistry_RegisterAndCount(t *testing.T) {
|
|||
}
|
||||
}
|
||||
|
||||
func TestRegistryRegisterIfAbsentRejectsConcurrentDuplicate(t *testing.T) {
|
||||
reg := edgenode.NewRegistry()
|
||||
|
||||
var accepted int32
|
||||
var wg sync.WaitGroup
|
||||
start := make(chan struct{})
|
||||
for i := 0; i < 32; i++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
<-start
|
||||
if reg.RegisterIfAbsent(&edgenode.NodeEntry{
|
||||
NodeID: "node-dup",
|
||||
Alias: "dup",
|
||||
Client: &toki.TcpClient{},
|
||||
}) {
|
||||
atomic.AddInt32(&accepted, 1)
|
||||
}
|
||||
}()
|
||||
}
|
||||
close(start)
|
||||
wg.Wait()
|
||||
|
||||
if accepted != 1 {
|
||||
t.Fatalf("accepted registrations: got %d want 1", accepted)
|
||||
}
|
||||
if reg.Count() != 1 {
|
||||
t.Fatalf("registry count: got %d want 1", reg.Count())
|
||||
}
|
||||
}
|
||||
|
||||
func TestRegistryUnregisterIfClientIgnoresStaleConnection(t *testing.T) {
|
||||
reg := edgenode.NewRegistry()
|
||||
live := &toki.TcpClient{}
|
||||
stale := &toki.TcpClient{}
|
||||
reg.Register(&edgenode.NodeEntry{NodeID: "node-1", Alias: "alias-1", Client: live})
|
||||
|
||||
if reg.UnregisterIfClient("node-1", stale) {
|
||||
t.Fatal("stale client should not unregister live entry")
|
||||
}
|
||||
entry, ok := reg.Get("node-1")
|
||||
if !ok {
|
||||
t.Fatal("live entry was removed by stale client")
|
||||
}
|
||||
if entry.Client != live {
|
||||
t.Fatal("live entry client changed unexpectedly")
|
||||
}
|
||||
|
||||
if !reg.UnregisterIfClient("node-1", live) {
|
||||
t.Fatal("live client should unregister entry")
|
||||
}
|
||||
if reg.Count() != 0 {
|
||||
t.Fatalf("registry count after live unregister: got %d want 0", reg.Count())
|
||||
}
|
||||
}
|
||||
|
||||
func TestRegistry_Resolve_ByAliasOrID(t *testing.T) {
|
||||
reg := edgenode.NewRegistry()
|
||||
entry := &edgenode.NodeEntry{
|
||||
|
|
|
|||
|
|
@ -184,26 +184,6 @@ func (s *Server) onNodeConnected(client *toki.TcpClient) {
|
|||
return &iop.RegisterResponse{Accepted: false, Reason: "unknown token"}, nil
|
||||
}
|
||||
|
||||
if _, ok := s.registry.Get(rec.ID); ok {
|
||||
reason := "node already connected"
|
||||
s.logger.Warn("duplicate registration rejected",
|
||||
zap.String("node_id", rec.ID),
|
||||
zap.String("agent_kind", rec.AgentKind),
|
||||
)
|
||||
s.emitNodeEvent(events.NewEdgeNodeEvent(
|
||||
events.SourceEdge,
|
||||
events.TypeNodeRegistrationFailed,
|
||||
rec.ID,
|
||||
rec.Alias,
|
||||
events.ReasonDuplicateConnection,
|
||||
map[string]string{
|
||||
events.MetadataFailureReason: events.ReasonDuplicateConnection,
|
||||
events.MetadataAgentKind: rec.AgentKind,
|
||||
},
|
||||
))
|
||||
return &iop.RegisterResponse{Accepted: false, Reason: reason}, nil
|
||||
}
|
||||
|
||||
cfg, err := edgenode.BuildConfigPayload(rec)
|
||||
if err != nil {
|
||||
s.logger.Error("build config payload failed",
|
||||
|
|
@ -264,9 +244,27 @@ func (s *Server) onNodeConnected(client *toki.TcpClient) {
|
|||
reason,
|
||||
meta,
|
||||
))
|
||||
s.registry.Unregister(rec.ID)
|
||||
s.registry.UnregisterIfClient(rec.ID, client)
|
||||
})
|
||||
s.registry.Register(entry)
|
||||
if !s.registry.RegisterIfAbsent(entry) {
|
||||
reason := "node already connected"
|
||||
s.logger.Warn("duplicate registration rejected",
|
||||
zap.String("node_id", rec.ID),
|
||||
zap.String("agent_kind", rec.AgentKind),
|
||||
)
|
||||
s.emitNodeEvent(events.NewEdgeNodeEvent(
|
||||
events.SourceEdge,
|
||||
events.TypeNodeRegistrationFailed,
|
||||
rec.ID,
|
||||
rec.Alias,
|
||||
events.ReasonDuplicateConnection,
|
||||
map[string]string{
|
||||
events.MetadataFailureReason: events.ReasonDuplicateConnection,
|
||||
events.MetadataAgentKind: rec.AgentKind,
|
||||
},
|
||||
))
|
||||
return &iop.RegisterResponse{Accepted: false, Reason: reason}, nil
|
||||
}
|
||||
s.logger.Info("node registered",
|
||||
zap.String("node_id", rec.ID),
|
||||
zap.String("alias", rec.Alias),
|
||||
|
|
|
|||
Loading…
Reference in a new issue