Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 5 additions & 8 deletions tests/integration/a2a_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,6 @@ import (
"iter"
"net/http"
"net/http/httptest"
"path/filepath"
"strings"
"sync/atomic"
"testing"
Expand Down Expand Up @@ -103,25 +102,23 @@ func TestA2ACUJ(t *testing.T) {
defer agent.Close()

t.Log("Starting Node A (provider, region=eu)...")
_ = startBackgroundNode(t, nodeBin, hubAddr, homeA,
nodeA := startBackgroundNode(t, nodeBin, hubAddr, homeA,
"--listen", "/ip4/127.0.0.1/udp/0/quic-v1",
"--listen", "/ip4/127.0.0.1/tcp/0",
"--discovery-interval", "100ms",
"--config", writeNodeConfig(t, homeA, map[string]string{"region": "eu"}, svcDecl{Type: "a2a", Name: "echo-agent", TargetURL: agent.URL}),
)
t.Log("Starting Node B (consumer)...")
_ = startBackgroundNode(t, nodeBin, hubAddr, homeB,
nodeB := startBackgroundNode(t, nodeBin, hubAddr, homeB,
"--listen", "/ip4/127.0.0.1/udp/0/quic-v1",
"--listen", "/ip4/127.0.0.1/tcp/0",
"--discovery-interval", "100ms",
)

apiAddrA := waitForMCPAddr(t, filepath.Join(homeA, "node.log"))
apiAddrB := waitForMCPAddr(t, filepath.Join(homeB, "node.log"))
waitForAPI(t, apiAddrA)
waitForAPI(t, apiAddrB)
apiAddrA := nodeA.waitForAPI(t)
apiAddrB := nodeB.waitForAPI(t)

addrA := waitForPeerInfoInLog(t, filepath.Join(homeA, "node.log"))
addrA := nodeA.p2pAddr
connectPeer(t, apiAddrB, addrA)
waitForDHTPeers(t, apiAddrA)

Expand Down
12 changes: 5 additions & 7 deletions tests/integration/agent_ingress_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -104,26 +104,24 @@ func TestAgentIngressCUJ(t *testing.T) {
t.Fatalf("writing node config: %v", err)
}

_ = startBackgroundNode(t, nodeBin, hubAddr, homeA,
nodeA := startBackgroundNode(t, nodeBin, hubAddr, homeA,
"--listen", "/ip4/127.0.0.1/udp/0/quic-v1",
"--listen", "/ip4/127.0.0.1/tcp/0",
"--discovery-interval", "100ms",
)
// Node B hosts the agent, and is the one that advertises its service.
_ = startBackgroundNode(t, nodeBin, hubAddr, homeB,
nodeB := startBackgroundNode(t, nodeBin, hubAddr, homeB,
"--listen", "/ip4/127.0.0.1/udp/0/quic-v1",
"--listen", "/ip4/127.0.0.1/tcp/0",
"--discovery-interval", "100ms",
"--socket-path", nodeSocket,
"--config", cfgPath,
)

apiAddrA := waitForMCPAddr(t, filepath.Join(homeA, "node.log"))
apiAddrB := waitForMCPAddr(t, filepath.Join(homeB, "node.log"))
waitForAPI(t, apiAddrA)
waitForAPI(t, apiAddrB)
apiAddrA := nodeA.waitForAPI(t)
apiAddrB := nodeB.waitForAPI(t)

addrB := waitForPeerInfoInLog(t, filepath.Join(homeB, "node.log"))
addrB := nodeB.p2pAddr
peerB := extractPeerID(addrB)
connectPeer(t, apiAddrA, addrB)
waitForDHTPeers(t, apiAddrB)
Expand Down
12 changes: 5 additions & 7 deletions tests/integration/agent_policy_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -83,25 +83,23 @@ services:
t.Cleanup(func() { _ = os.RemoveAll(sockDir) })
nodeSocket := filepath.Join(sockDir, "node.sock")

_ = startBackgroundNode(t, nodeBin, hubAddr, homeA,
nodeA := startBackgroundNode(t, nodeBin, hubAddr, homeA,
"--listen", "/ip4/127.0.0.1/udp/0/quic-v1",
"--listen", "/ip4/127.0.0.1/tcp/0",
"--discovery-interval", "100ms",
"--config", configA,
)
_ = startBackgroundNode(t, nodeBin, hubAddr, homeB,
nodeB := startBackgroundNode(t, nodeBin, hubAddr, homeB,
"--listen", "/ip4/127.0.0.1/udp/0/quic-v1",
"--listen", "/ip4/127.0.0.1/tcp/0",
"--discovery-interval", "100ms",
"--socket-path", nodeSocket,
)

apiAddrA := waitForMCPAddr(t, filepath.Join(homeA, "node.log"))
apiAddrB := waitForMCPAddr(t, filepath.Join(homeB, "node.log"))
waitForAPI(t, apiAddrA)
waitForAPI(t, apiAddrB)
apiAddrA := nodeA.waitForAPI(t)
apiAddrB := nodeB.waitForAPI(t)

addrA := waitForPeerInfoInLog(t, filepath.Join(homeA, "node.log"))
addrA := nodeA.p2pAddr
connectPeer(t, apiAddrB, addrA)
waitForDHTPeers(t, apiAddrA)

Expand Down
107 changes: 9 additions & 98 deletions tests/integration/catalog_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,6 @@ import (
"fmt"
"net/http"
"os"
"os/exec"
"path/filepath"
"sort"
"strings"
Expand Down Expand Up @@ -87,62 +86,6 @@ func writeNodeConfig(t *testing.T, dir string, labels map[string]string, service
return p
}

func startBackgroundNode(t *testing.T, nodeBin string, routerAddr string, homeDir string, args ...string) *exec.Cmd {
t.Helper()
env := append(os.Environ(),
"HOME="+homeDir,
"XDG_CONFIG_HOME="+filepath.Join(homeDir, ".config"),
"SAM_API_TOKEN=test-token", // per-test overrides use --api-token-path, which wins
)
allArgs := append([]string{"run", "--control-plane", routerAddr, "--jwt", "test-jwt", "--bind-addr", "127.0.0.1:0", "--allow-loopback"}, args...)
cmd := exec.Command(nodeBin, allArgs...)
cmd.Env = env

logFile, err := os.Create(filepath.Join(homeDir, "node.log"))
if err != nil {
t.Fatalf("failed to create log file: %v", err)
}
cmd.Stdout = logFile
cmd.Stderr = logFile

if err := cmd.Start(); err != nil {
t.Fatalf("failed to start background node: %v", err)
}

t.Cleanup(func() {
if err := cmd.Process.Kill(); err != nil {
t.Logf("warning: failed to kill background node: %v", err)
}
if err := logFile.Close(); err != nil {
t.Logf("warning: failed to close log file: %v", err)
}
})

return cmd
}

func waitForMCPAddr(t *testing.T, logPath string) string {
t.Helper()
// Generous under CI load; polling returns as soon as the line appears.
deadline := time.Now().Add(10 * time.Second)
for time.Now().Before(deadline) {
data, _ := os.ReadFile(logPath)
lines := strings.Split(string(data), "\n")
for _, line := range lines {
if strings.Contains(line, "Starting MCP server on TCP address ") {
parts := strings.Split(line, "Starting MCP server on TCP address ")
if len(parts) > 1 {
return strings.TrimSpace(parts[1])
}
}
}
time.Sleep(100 * time.Millisecond)
}
data, _ := os.ReadFile(logPath)
t.Fatalf("timeout waiting for MCP addr in log: %s\n--- log contents ---\n%s", logPath, string(data))
return ""
}

func callMCP(t *testing.T, mcpAddr string, toolName string, params map[string]any) string {
t.Helper()
ctx := context.Background()
Expand Down Expand Up @@ -201,36 +144,6 @@ func (a *authRoundTripper) RoundTrip(req *http.Request) (*http.Response, error)
return a.rt.RoundTrip(clone)
}

func waitForPeerInfoInLog(t *testing.T, logPath string) string {
t.Helper()
deadline := time.Now().Add(5 * time.Second)
for time.Now().Before(deadline) {
data, _ := os.ReadFile(logPath)
lines := strings.Split(string(data), "\n")
var peerID string
var tcpAddr string
for _, line := range lines {
if strings.HasPrefix(line, "PeerID: ") {
peerID = strings.TrimPrefix(line, "PeerID: ")
}
if strings.Contains(line, "Listening on: ") {
parts := strings.Split(line, " ")
for _, p := range parts {
if strings.Contains(p, "/tcp/") {
tcpAddr = strings.Trim(p, "[]")
}
}
}
}
if peerID != "" && tcpAddr != "" {
return tcpAddr + "/p2p/" + peerID
}
time.Sleep(100 * time.Millisecond)
}
t.Fatalf("timeout waiting for peer info in log: %s", logPath)
return ""
}

func TestCatalogRoutingAndFailover(t *testing.T) {
nodeBin := buildBinary(t, "./cmd/sam-node")
_, routerAddr := startMockRouter(t)
Expand All @@ -241,25 +154,25 @@ func TestCatalogRoutingAndFailover(t *testing.T) {

// Start Node A (Client)
t.Log("Starting Node A...")
_ = startBackgroundNode(t, nodeBin, routerAddr, homeA, "--listen", "/ip4/127.0.0.1/udp/0/quic-v1", "--listen", "/ip4/127.0.0.1/tcp/0", "--discovery-interval", "100ms")
nodeA := startBackgroundNode(t, nodeBin, routerAddr, homeA, "--listen", "/ip4/127.0.0.1/udp/0/quic-v1", "--listen", "/ip4/127.0.0.1/tcp/0", "--discovery-interval", "100ms")
t.Log("Node A started.")

// Wait for Node A to start and get its MCP address
mcpAddrA := waitForMCPAddr(t, filepath.Join(homeA, "node.log"))
mcpAddrA := nodeA.waitForAPI(t)

// Start Node B (Provider 1)
t.Log("Starting Node B...")
cmdB := startBackgroundNode(t, nodeBin, routerAddr, homeB, "--listen", "/ip4/127.0.0.1/udp/0/quic-v1", "--listen", "/ip4/127.0.0.1/tcp/0", "--discovery-interval", "100ms")
nodeB := startBackgroundNode(t, nodeBin, routerAddr, homeB, "--listen", "/ip4/127.0.0.1/udp/0/quic-v1", "--listen", "/ip4/127.0.0.1/tcp/0", "--discovery-interval", "100ms")
t.Log("Node B started.")

// Start Node C (Provider 2)
t.Log("Starting Node C...")
_ = startBackgroundNode(t, nodeBin, routerAddr, homeC, "--listen", "/ip4/127.0.0.1/udp/0/quic-v1", "--listen", "/ip4/127.0.0.1/tcp/0", "--discovery-interval", "100ms")
nodeC := startBackgroundNode(t, nodeBin, routerAddr, homeC, "--listen", "/ip4/127.0.0.1/udp/0/quic-v1", "--listen", "/ip4/127.0.0.1/tcp/0", "--discovery-interval", "100ms")
t.Log("Node C started.")

// Wait for Node B and C to start and get their addresses
addrB := waitForPeerInfoInLog(t, filepath.Join(homeB, "node.log"))
addrC := waitForPeerInfoInLog(t, filepath.Join(homeC, "node.log"))
nodeB.waitForAPI(t)
nodeC.waitForAPI(t)
addrB := nodeB.p2pAddr
addrC := nodeC.p2pAddr

// Force Node A to connect to Node B and Node C
connectPeer(t, mcpAddrA, addrB)
Expand Down Expand Up @@ -329,9 +242,7 @@ func TestCatalogRoutingAndFailover(t *testing.T) {
t.Logf("First call response: %s", respData)

// Now kill Node B and assert failover to Node C
if err := cmdB.Process.Kill(); err != nil {
t.Fatalf("failed to kill Node B: %v", err)
}
nodeB.kill()

// Wait a bit for catalog update or failover to happen on next call
time.Sleep(500 * time.Millisecond)
Expand Down
56 changes: 0 additions & 56 deletions tests/integration/cuj_test.go

This file was deleted.

36 changes: 11 additions & 25 deletions tests/integration/datapath_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -48,34 +48,27 @@ func TestIntegrationStdioDatapath(t *testing.T) {

// Start Node A
t.Log("Starting Node A...")
_ = startBackgroundNode(t, nodeBin, routerAddr, homeA,
nodeA := startBackgroundNode(t, nodeBin, routerAddr, homeA,
"--listen", "/ip4/127.0.0.1/udp/0/quic-v1",
"--listen", "/ip4/127.0.0.1/tcp/0",
"--discovery-interval", "100ms",
"--bind-addr", "127.0.0.1:0",
"--api-token-path", tokenPath(t, apiToken),
"--config", cfgA,
)

// Start Node B
t.Log("Starting Node B...")
_ = startBackgroundNode(t, nodeBin, routerAddr, homeB,
nodeB := startBackgroundNode(t, nodeBin, routerAddr, homeB,
"--listen", "/ip4/127.0.0.1/udp/0/quic-v1",
"--listen", "/ip4/127.0.0.1/tcp/0",
"--discovery-interval", "100ms",
"--bind-addr", "127.0.0.1:0",
"--api-token-path", tokenPath(t, apiToken),
)

// Resolve actual addresses from logs
actualApiAddrA := waitForMCPAddr(t, filepath.Join(homeA, "node.log"))
actualApiAddrB := waitForMCPAddr(t, filepath.Join(homeB, "node.log"))
nodeA.waitForAPI(t)
actualApiAddrB := nodeB.waitForAPI(t)

// Wait for nodes to start sidecar API
waitForAPI(t, actualApiAddrA)
waitForAPI(t, actualApiAddrB)

addrA := waitForPeerInfoInLog(t, filepath.Join(homeA, "node.log"))
addrA := nodeA.p2pAddr
peerIDA := getPeerIDFromAddr(addrA)

// Connect Node B to Node A
Expand Down Expand Up @@ -177,36 +170,29 @@ func TestIntegrationHTTPDatapath(t *testing.T) {

// Start Node A
t.Log("Starting Node A...")
_ = startBackgroundNode(t, nodeBin, routerAddr, homeA,
nodeA := startBackgroundNode(t, nodeBin, routerAddr, homeA,
"--listen", "/ip4/127.0.0.1/udp/0/quic-v1",
"--listen", "/ip4/127.0.0.1/tcp/0",
"--discovery-interval", "100ms",
"--bind-addr", "127.0.0.1:0",
"--api-token-path", tokenPath(t, apiToken),
"--config", cfgA,
)

// Start Node B
t.Log("Starting Node B...")
_ = startBackgroundNode(t, nodeBin, routerAddr, homeB,
nodeB := startBackgroundNode(t, nodeBin, routerAddr, homeB,
"--listen", "/ip4/127.0.0.1/udp/0/quic-v1",
"--listen", "/ip4/127.0.0.1/tcp/0",
"--discovery-interval", "100ms",
"--bind-addr", "127.0.0.1:0",
"--api-token-path", tokenPath(t, apiToken),
)

// Resolve actual addresses from logs
actualApiAddrA := waitForMCPAddr(t, filepath.Join(homeA, "node.log"))
actualApiAddrB := waitForMCPAddr(t, filepath.Join(homeB, "node.log"))

// Wait for nodes to start sidecar API
waitForAPI(t, actualApiAddrA)
waitForAPI(t, actualApiAddrB)
nodeA.waitForAPI(t)
actualApiAddrB := nodeB.waitForAPI(t)

addrA := waitForPeerInfoInLog(t, filepath.Join(homeA, "node.log"))
addrA := nodeA.p2pAddr
peerIDA := getPeerIDFromAddr(addrA)
addrB := waitForPeerInfoInLog(t, filepath.Join(homeB, "node.log"))
addrB := nodeB.p2pAddr
peerIDB := getPeerIDFromAddr(addrB)

// Connect Node B to Node A
Expand Down
Loading
Loading