iop/apps/node/internal/node/node_concurrency_integration_test.go
toki 39d1d08f33 feat(node): FIFO admission queue 구현 및 concurrency 처리 개선
- Node runner manager에서 FIFO 큐 기반 admission queue 구현
- Concurrency limit 적용으로 최대 동시 실행 job 수 제한
- Store 계층에 queue 상태 영구화 support 추가
- 관련 테스트 및 integration 테스트 업데이트
- Cloud/G06 문서 분류 조정 (archive로 이동)
2026-06-14 21:37:34 +09:00

215 lines
6.8 KiB
Go

package node_test
import (
"context"
"fmt"
"net"
"strings"
"testing"
"time"
toki "git.toki-labs.com/toki/proto-socket/go"
"go.uber.org/zap"
"google.golang.org/protobuf/proto"
"iop/apps/node/internal/node"
"iop/apps/node/internal/runtime"
"iop/apps/node/internal/store"
"iop/apps/node/internal/transport"
iop "iop/proto/gen/iop"
)
// edgeServerParserMap returns the proto parser map for a mock Edge server
// that expects to receive RunEvent messages sent by Node.
func edgeServerParserMap() 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)
},
}
}
func getFreeListenAddr(t *testing.T) string {
t.Helper()
l, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatalf("listen: %v", err)
}
addr := l.Addr().String()
l.Close()
return addr
}
// startMockEdge starts a mock Edge TCP server that accepts exactly one Node
// connection, performs the registration handshake, and returns the accepted
// toki.TcpClient so tests can attach listeners and send RunRequests.
func startMockEdge(t *testing.T, ctx context.Context, listenAddr string) *toki.TcpClient {
t.Helper()
host, portStr, _ := net.SplitHostPort(listenAddr)
port := 0
fmt.Sscanf(portStr, "%d", &port)
acceptedCh := make(chan *toki.TcpClient, 1)
server := toki.NewTcpServer(host, port, func(conn net.Conn) *toki.TcpClient {
client := toki.NewTcpClient(conn, 30, 10, edgeServerParserMap())
toki.AddRequestListenerTyped[*iop.RegisterRequest, *iop.RegisterResponse](
&client.Communicator,
func(req *iop.RegisterRequest) (*iop.RegisterResponse, error) {
return &iop.RegisterResponse{
Accepted: true,
NodeId: "test-node",
Alias: "test-alias",
Config: &iop.NodeConfigPayload{
Runtime: &iop.NodeRuntimeConfig{Concurrency: 1},
},
}, nil
},
)
acceptedCh <- client
return client
})
if err := server.Start(ctx); err != nil {
t.Fatalf("mock edge server start: %v", err)
}
t.Cleanup(func() { server.Stop() })
select {
case client := <-acceptedCh:
return client
case <-time.After(3 * time.Second):
t.Fatal("mock edge server did not accept connection within 3s")
return nil
}
}
// TestQueueOverflow_RejectEventObservedByEdge is an integration test that
// verifies the RunEvent{type:"error"} sent by rejectRun() carries the queue
// overflow reason and actually arrives at the mock Edge server when the Node's
// FIFO admission queue is full (capacity=1, max_queue=0).
//
// Setup:
// 1. Start mock Edge TCP server.
// 2. Dial Node-side session via transport.DialEdge.
// 3. Create a real node.Node with globalConcurrency=1.
// 4. Attach the session to the node via SetHandler.
// 5. From the Edge side, send a first RunRequest whose adapter blocks.
// 6. After the first run starts, send a second RunRequest from the Edge side.
// 7. Assert a RunEvent{type:"error", run_id:<second>} with a queue_full reason
// arrives at the Edge.
func TestQueueOverflow_RejectEventObservedByEdge(t *testing.T) {
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
logger := zap.NewNop()
listenAddr := getFreeListenAddr(t)
// Channel that the mock edge client will be delivered into once Node connects.
edgeClientCh := make(chan *toki.TcpClient, 1)
// Start mock Edge server (async — we need to dial first to trigger accept).
host, portStr, _ := net.SplitHostPort(listenAddr)
port := 0
fmt.Sscanf(portStr, "%d", &port)
server := toki.NewTcpServer(host, port, func(conn net.Conn) *toki.TcpClient {
client := toki.NewTcpClient(conn, 30, 10, edgeServerParserMap())
toki.AddRequestListenerTyped[*iop.RegisterRequest, *iop.RegisterResponse](
&client.Communicator,
func(req *iop.RegisterRequest) (*iop.RegisterResponse, error) {
return &iop.RegisterResponse{
Accepted: true,
NodeId: "test-node",
Alias: "test-alias",
Config: &iop.NodeConfigPayload{
Runtime: &iop.NodeRuntimeConfig{Concurrency: 1},
},
}, nil
},
)
edgeClientCh <- client
return client
})
if err := server.Start(ctx); err != nil {
t.Fatalf("mock edge server start: %v", err)
}
defer server.Stop()
// Dial Node side to Edge.
dialResult, err := transport.DialEdge(ctx, listenAddr, "test-token", logger)
if err != nil {
t.Fatalf("DialEdge: %v", err)
}
defer dialResult.Session.Close()
// Wait for Edge side to accept the connection.
var edgeClient *toki.TcpClient
select {
case edgeClient = <-edgeClientCh:
case <-time.After(3 * time.Second):
t.Fatal("edge server did not accept connection")
}
// Attach RunEvent listener on the Edge side before sending any requests.
runEventCh := make(chan *iop.RunEvent, 4)
toki.AddListenerTyped[*iop.RunEvent](&edgeClient.Communicator, func(event *iop.RunEvent) {
runEventCh <- event
})
// Build real node.Node with global concurrency=1.
// capacity=1, max_queue=0 → a second concurrent run overflows the queue.
sa := newQueuedSlowAdapter("slow", 1, 0, 0)
rtr := &fixedRouter{
adapterName: "slow",
adapters: map[string]runtime.Adapter{"slow": sa},
}
st, err := store.New(":memory:", logger)
if err != nil {
t.Fatalf("store: %v", err)
}
t.Cleanup(func() { _ = st.Close() })
n := node.New("test-node", rtr, st, 1, nil, logger)
dialResult.Session.SetHandler(n)
// Send first RunRequest from Edge → Node (adapter will block).
firstReq := &iop.RunRequest{RunId: "run-hold-1", Adapter: "slow", Target: "v1"}
if err := edgeClient.Send(firstReq); err != nil {
t.Fatalf("send first run request: %v", err)
}
// Wait until the first run's adapter has started executing so the slot is held.
waitStarted(t, sa, "run-hold-1")
// Send second RunRequest; should be rejected because capacity=1 is taken and
// the queue (max_queue=0) cannot hold it.
secondReq := &iop.RunRequest{RunId: "run-reject-2", Adapter: "slow", Target: "v1"}
if err := edgeClient.Send(secondReq); err != nil {
t.Fatalf("send second run request: %v", err)
}
// Collect RunEvents until we see the error event for the second run.
deadline := time.After(5 * time.Second)
for {
select {
case ev := <-runEventCh:
if ev.GetRunId() == "run-reject-2" && ev.GetType() == string(runtime.EventTypeError) {
if ev.GetError() == "" {
t.Fatal("reject RunEvent has empty error field")
}
if !strings.Contains(ev.GetError(), "queue full") {
t.Fatalf("expected queue_full reason in reject event, got %q", ev.GetError())
}
// Success: RunEvent{type:"error", run_id:"run-reject-2"} observed.
sa.releaseRun("run-hold-1")
return
}
case <-deadline:
t.Fatal("timeout: did not receive RunEvent{type:error} for run-reject-2")
}
}
}