iop/apps/node/internal/node/node_concurrency_integration_test.go

668 lines
21 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/store"
"iop/apps/node/internal/transport"
runtime "iop/packages/go/agentruntime"
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
}
}
// TestOverDispatchSafety_RejectEventObservedByEdge is an integration test that
// verifies the RunEvent{type:"error"} sent by rejectRun() carries the concurrency
// unavailable reason and actually arrives at the mock Edge server when the adapter's
// concurrency capacity limit is exceeded (capacity=1).
//
// Setup:
// 1. Start mock Edge TCP server.
// 2. Dial Node-side session via transport.DialEdge.
// 3. Create a real node.Node with an adapter whose MaxConcurrency=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 concurrency unavailable reason
// arrives at the Edge.
func TestOverDispatchSafety_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 adapter capacity=1.
// A second concurrent run overflows the adapter concurrency limit.
sa := newQueuedSlowAdapter("slow", 1, 0, 0)
rtr := &fixedRouter{
adapterName: "slow",
adapters: map[string]runtime.Provider{"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, nil)
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.
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(), "concurrency unavailable") {
t.Fatalf("expected concurrency unavailable 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")
}
}
}
// TestOnRunRequest_SynthesizedTerminalObservedByEdge verifies that when an
// adapter returns without emitting a terminal event, Node synthesizes the
// appropriate terminal RunEvent and it arrives at the mock Edge server.
func TestOnRunRequest_SynthesizedTerminalObservedByEdge(t *testing.T) {
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
logger := zap.NewNop()
listenAddr := getFreeListenAddr(t)
edgeClientCh := make(chan *toki.TcpClient, 1)
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()
dialResult, err := transport.DialEdge(ctx, listenAddr, "test-token", logger)
if err != nil {
t.Fatalf("DialEdge: %v", err)
}
defer dialResult.Session.Close()
var edgeClient *toki.TcpClient
select {
case edgeClient = <-edgeClientCh:
case <-time.After(3 * time.Second):
t.Fatal("edge server did not accept connection")
}
runEventCh := make(chan *iop.RunEvent, 8)
toki.AddListenerTyped[*iop.RunEvent](&edgeClient.Communicator, func(event *iop.RunEvent) {
runEventCh <- event
})
// synthesizeAdapter returns without emitting any terminal event.
synAdapter := &synthesizeNoTerminalAdapter{}
rtr := &fixedRouter{
adapterName: "synth",
adapters: map[string]runtime.Provider{"synth": synAdapter},
}
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, nil)
dialResult.Session.SetHandler(n)
req := &iop.RunRequest{RunId: "run-synth-ok", Adapter: "synth", Target: "v1"}
if err := edgeClient.Send(req); err != nil {
t.Fatalf("send run request: %v", err)
}
deadline := time.After(5 * time.Second)
for {
select {
case ev := <-runEventCh:
if ev.GetRunId() == "run-synth-ok" && ev.GetType() == string(runtime.EventTypeComplete) {
// Success: synthesized complete event observed.
return
}
case <-deadline:
t.Fatal("timeout: did not receive synthesized complete event")
}
}
}
// TestOnRunRequest_SynthesizedErrorObservedByEdge verifies that when an
// adapter returns a non-cancel error, Node synthesizes an error terminal event.
func TestOnRunRequest_SynthesizedErrorObservedByEdge(t *testing.T) {
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
logger := zap.NewNop()
listenAddr := getFreeListenAddr(t)
edgeClientCh := make(chan *toki.TcpClient, 1)
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()
dialResult, err := transport.DialEdge(ctx, listenAddr, "test-token", logger)
if err != nil {
t.Fatalf("DialEdge: %v", err)
}
defer dialResult.Session.Close()
var edgeClient *toki.TcpClient
select {
case edgeClient = <-edgeClientCh:
case <-time.After(3 * time.Second):
t.Fatal("edge server did not accept connection")
}
runEventCh := make(chan *iop.RunEvent, 8)
toki.AddListenerTyped[*iop.RunEvent](&edgeClient.Communicator, func(event *iop.RunEvent) {
runEventCh <- event
})
// errorAdapter returns a failure error without emitting terminal events.
errAdapter := &synthesizeErrorAdapter{err: fmt.Errorf("provider timeout")}
rtr := &fixedRouter{
adapterName: "synth-err",
adapters: map[string]runtime.Provider{"synth-err": errAdapter},
}
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, nil)
dialResult.Session.SetHandler(n)
req := &iop.RunRequest{RunId: "run-synth-err", Adapter: "synth-err", Target: "v1"}
if err := edgeClient.Send(req); err != nil {
t.Fatalf("send run request: %v", err)
}
deadline := time.After(5 * time.Second)
for {
select {
case ev := <-runEventCh:
if ev.GetRunId() == "run-synth-err" && ev.GetType() == string(runtime.EventTypeError) {
if !strings.Contains(ev.GetError(), "provider timeout") {
t.Fatalf("expected error message in synthesized event, got %q", ev.GetError())
}
return
}
case <-deadline:
t.Fatal("timeout: did not receive synthesized error event")
}
}
}
// synthesizeNoTerminalAdapter returns nil (success) without emitting any terminal event.
type synthesizeNoTerminalAdapter struct{}
func (a *synthesizeNoTerminalAdapter) Name() string { return "synth" }
func (a *synthesizeNoTerminalAdapter) Capabilities(_ context.Context) (runtime.Capabilities, error) {
return runtime.Capabilities{AdapterName: "synth", Targets: []string{"v1"}, MaxConcurrency: 1}, nil
}
func (a *synthesizeNoTerminalAdapter) Execute(_ context.Context, _ runtime.ExecutionSpec, _ runtime.EventSink) error {
// Returns success without emitting any terminal event.
return nil
}
// synthesizeErrorAdapter returns an error without emitting any terminal event.
type synthesizeErrorAdapter struct {
err error
}
func (a *synthesizeErrorAdapter) Name() string { return "synth-err" }
func (a *synthesizeErrorAdapter) Capabilities(_ context.Context) (runtime.Capabilities, error) {
return runtime.Capabilities{AdapterName: "synth-err", Targets: []string{"v1"}, MaxConcurrency: 1}, nil
}
func (a *synthesizeErrorAdapter) Execute(_ context.Context, _ runtime.ExecutionSpec, _ runtime.EventSink) error {
return a.err
}
// cancelAdapter returns runtime.ErrRunCancelled without emitting any terminal event.
// This simulates a real adapter that cancels without emitting a terminal event,
// relying on Node's synthAndEmitTerminal to synthesize the cancelled event.
type cancelAdapter struct {
started chan struct{}
done chan struct{}
}
func newCancelAdapter(name string) *cancelAdapter {
return &cancelAdapter{
started: make(chan struct{}),
done: make(chan struct{}),
}
}
func (a *cancelAdapter) Name() string { return "cancel" }
func (a *cancelAdapter) Capabilities(_ context.Context) (runtime.Capabilities, error) {
return runtime.Capabilities{AdapterName: "cancel", Targets: []string{"v1"}, MaxConcurrency: 1}, nil
}
func (a *cancelAdapter) Execute(_ context.Context, _ runtime.ExecutionSpec, _ runtime.EventSink) error {
close(a.started)
<-a.done
return runtime.ErrRunCancelled
}
func (a *cancelAdapter) releaseRun() { close(a.done) }
// TestOnRunRequest_SynthesizedCancelledObservedByEdge verifies that when an
// adapter returns runtime.ErrRunCancelled without emitting any terminal event,
// Node synthesizes a RunEvent{type:"cancelled"} and it arrives at the mock Edge server.
// This is a regression test for the inflight accounting leak scenario where
// ErrRunCancelled was not being sent as a cancelled terminal event.
func TestOnRunRequest_SynthesizedCancelledObservedByEdge(t *testing.T) {
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
logger := zap.NewNop()
listenAddr := getFreeListenAddr(t)
edgeClientCh := make(chan *toki.TcpClient, 1)
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()
dialResult, err := transport.DialEdge(ctx, listenAddr, "test-token", logger)
if err != nil {
t.Fatalf("DialEdge: %v", err)
}
defer dialResult.Session.Close()
var edgeClient *toki.TcpClient
select {
case edgeClient = <-edgeClientCh:
case <-time.After(3 * time.Second):
t.Fatal("edge server did not accept connection")
}
runEventCh := make(chan *iop.RunEvent, 8)
toki.AddListenerTyped[*iop.RunEvent](&edgeClient.Communicator, func(event *iop.RunEvent) {
runEventCh <- event
})
// cancelAdapter returns runtime.ErrRunCancelled without emitting any terminal event.
cancelAdapt := newCancelAdapter("cancel")
rtr := &fixedRouter{
adapterName: "cancel",
adapters: map[string]runtime.Provider{"cancel": cancelAdapt},
}
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, nil)
dialResult.Session.SetHandler(n)
req := &iop.RunRequest{RunId: "run-cancel-test", Adapter: "cancel", Target: "v1"}
if err := edgeClient.Send(req); err != nil {
t.Fatalf("send run request: %v", err)
}
// Wait until the adapter's Execute has started so the run is in-flight.
select {
case <-cancelAdapt.started:
// Adapter is now blocking on cancelAdapt.done.
case <-time.After(5 * time.Second):
t.Fatal("timeout: adapter did not start")
}
// Now release the adapter so it returns ErrRunCancelled.
cancelAdapt.releaseRun()
// Collect RunEvents until we see the cancelled event for this run, then verify
// no additional terminal event (complete/error) arrives for the same run id.
deadline := time.After(5 * time.Second)
for {
select {
case ev := <-runEventCh:
if ev.GetRunId() == "run-cancel-test" && ev.GetType() == string(runtime.EventTypeCancelled) {
// We saw the expected cancelled event; now settle and verify no
// additional complete/error terminal event follows for the same run.
settleDeadline := time.After(2 * time.Second)
for {
select {
case extraEv := <-runEventCh:
if extraEv.GetRunId() == "run-cancel-test" {
switch extraEv.GetType() {
case string(runtime.EventTypeComplete), string(runtime.EventTypeError):
t.Fatalf("unexpected additional terminal event %q after cancelled for run-cancel-test", extraEv.GetType())
}
// Non-terminal event (e.g. delta); ignore and keep settling.
continue
}
default:
// No more events within 2s; settle passed.
goto done
case <-settleDeadline:
// No additional events; settle passed.
goto done
}
}
done:
// Success: synthesized cancelled event observed by Edge with no follow-up terminal.
return
}
case <-deadline:
t.Fatal("timeout: did not receive RunEvent{type:cancelled} for run-cancel-test")
}
}
}
// TestIntegration_ResolveAdapterErrorObservedByEdge verifies that when ResolveAdapter
// fails, Node sends a RunEvent{type:"error"} to the session so Edge can observe the
// failure and avoid inflight slot leaks.
func TestIntegration_ResolveAdapterErrorObservedByEdge(t *testing.T) {
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
logger := zap.NewNop()
listenAddr := getFreeListenAddr(t)
edgeClientCh := make(chan *toki.TcpClient, 1)
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()
dialResult, err := transport.DialEdge(ctx, listenAddr, "test-token", logger)
if err != nil {
t.Fatalf("DialEdge: %v", err)
}
defer dialResult.Session.Close()
var edgeClient *toki.TcpClient
select {
case edgeClient = <-edgeClientCh:
case <-time.After(3 * time.Second):
t.Fatal("edge server did not accept connection")
}
runEventCh := make(chan *iop.RunEvent, 4)
toki.AddListenerTyped[*iop.RunEvent](&edgeClient.Communicator, func(event *iop.RunEvent) {
runEventCh <- event
})
// Use an errorRouter that always fails ResolveAdapter.
rtr := &errorRouter{err: fmt.Errorf("adapter not found")}
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, nil)
dialResult.Session.SetHandler(n)
// Send RunRequest from Edge — Node will fail at ResolveAdapter.
req := &iop.RunRequest{RunId: "run-resolve-fail", Adapter: "nonexistent", Target: "v1"}
if err := edgeClient.Send(req); err != nil {
t.Fatalf("send run request: %v", err)
}
deadline := time.After(5 * time.Second)
for {
select {
case ev := <-runEventCh:
if ev.GetRunId() == "run-resolve-fail" && ev.GetType() == string(runtime.EventTypeError) {
if ev.GetError() == "" {
t.Fatal("ResolveAdapter error RunEvent has empty error field")
}
if !strings.Contains(ev.GetError(), "adapter not found") {
t.Fatalf("expected 'adapter not found' in error, got %q", ev.GetError())
}
// Success: error event for failed ResolveAdapter observed by Edge.
return
}
case <-deadline:
t.Fatal("timeout: did not receive error event for failed ResolveAdapter")
}
}
}