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
5 changes: 2 additions & 3 deletions forge-cli/runtime/runner_jsonrpc_headers_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,6 @@ import (
"net/http"
"strconv"
"testing"
"time"

"github.com/initializ/forge/forge-core/a2a"
"github.com/initializ/forge/forge-core/auth"
Expand Down Expand Up @@ -72,7 +71,7 @@ func TestRunner_JSONRPC_TasksSend_StampsForgeUsageHeaders(t *testing.T) {
go func() { _ = runner.Run(ctx) }()

baseURL := fmt.Sprintf("http://localhost:%d", port)
waitForServer(t, baseURL, 5*time.Second)
baseURL = waitForServer(t, baseURL, serverReadyTimeout)

token, err := auth.LoadToken(dir)
if err != nil {
Expand Down Expand Up @@ -179,7 +178,7 @@ func TestRunner_JSONRPC_WorkflowContextThreadsThroughDispatcher(t *testing.T) {
defer cancel()
go func() { _ = runner.Run(ctx) }()
baseURL := fmt.Sprintf("http://localhost:%d", port)
waitForServer(t, baseURL, 5*time.Second)
baseURL = waitForServer(t, baseURL, serverReadyTimeout)
token, _ := auth.LoadToken(dir)

rpcReq := a2a.JSONRPCRequest{
Expand Down
3 changes: 1 addition & 2 deletions forge-cli/runtime/runner_killswitch_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,6 @@ import (
"fmt"
"net/http"
"testing"
"time"

"github.com/initializ/forge/forge-core/a2a"
"github.com/initializ/forge/forge-core/auth"
Expand Down Expand Up @@ -39,7 +38,7 @@ func TestRunner_KillSwitch_RefusesNewWorkOnEveryIngress(t *testing.T) {
defer cancel()
go func() { _ = runner.Run(ctx) }()
baseURL := fmt.Sprintf("http://localhost:%d", port)
waitForServer(t, baseURL, 5*time.Second)
baseURL = waitForServer(t, baseURL, serverReadyTimeout)
token, _ := auth.LoadToken(dir)

// Trip the kill switch. Idle agent → cancelled=0, but the call must still
Expand Down
7 changes: 3 additions & 4 deletions forge-cli/runtime/runner_message_validation_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,6 @@ import (
"net/http"
"strings"
"testing"
"time"

"github.com/initializ/forge/forge-core/a2a"
"github.com/initializ/forge/forge-core/auth"
Expand Down Expand Up @@ -45,7 +44,7 @@ func TestRunner_JSONRPC_TasksSend_RejectsLegacyTypeDiscriminator(t *testing.T) {
defer cancel()
go func() { _ = runner.Run(ctx) }()
baseURL := fmt.Sprintf("http://localhost:%d", port)
waitForServer(t, baseURL, 5*time.Second)
baseURL = waitForServer(t, baseURL, serverReadyTimeout)
token, _ := auth.LoadToken(dir)

// Exact payload shape from issue #119: parts use `"type"` instead
Expand Down Expand Up @@ -125,7 +124,7 @@ func TestRunner_JSONRPC_TasksSend_SpecCompliantPayloadStillWorks(t *testing.T) {
defer cancel()
go func() { _ = runner.Run(ctx) }()
baseURL := fmt.Sprintf("http://localhost:%d", port)
waitForServer(t, baseURL, 5*time.Second)
baseURL = waitForServer(t, baseURL, serverReadyTimeout)
token, _ := auth.LoadToken(dir)

// Same payload, but with the spec-correct `"kind"` discriminator.
Expand Down Expand Up @@ -192,7 +191,7 @@ func TestRunner_JSONRPC_TasksSend_RejectsEmptyParts(t *testing.T) {
defer cancel()
go func() { _ = runner.Run(ctx) }()
baseURL := fmt.Sprintf("http://localhost:%d", port)
waitForServer(t, baseURL, 5*time.Second)
baseURL = waitForServer(t, baseURL, serverReadyTimeout)
token, _ := auth.LoadToken(dir)

body := []byte(`{
Expand Down
45 changes: 39 additions & 6 deletions forge-cli/runtime/runner_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,8 +6,10 @@ import (
"encoding/json"
"fmt"
"net/http"
"net/url"
"os"
"path/filepath"
"strconv"
"testing"
"time"

Expand Down Expand Up @@ -57,7 +59,7 @@ func TestRunner_MockIntegration(t *testing.T) {

// Wait for server to be ready
baseURL := fmt.Sprintf("http://localhost:%d", port)
waitForServer(t, baseURL, 5*time.Second)
baseURL = waitForServer(t, baseURL, serverReadyTimeout)

// Load the auto-generated auth token.
token, err := auth.LoadToken(dir)
Expand Down Expand Up @@ -385,20 +387,51 @@ func TestExpandEgressDomains(t *testing.T) {
}
}

func waitForServer(t *testing.T, baseURL string, timeout time.Duration) {
// serverReadyTimeout bounds how long a test waits for the runner's HTTP server
// to become reachable. Generous on purpose: it only ever elapses on FAILURE
// (a green run returns as soon as /healthz answers), so a high ceiling costs
// nothing on passing runs while absorbing CI startup jitter that made the
// previous 5s ceiling flaky under load.
const serverReadyTimeout = 20 * time.Second

// waitForServer polls until the runner's HTTP server is reachable and returns
// the base URL it actually came up on.
//
// The server auto-increments its port on conflict (up to 10 attempts, see
// server.Start), and the port findFreePort hands out can be stolen in the gap
// before the runner binds it — so the server may listen on requestedPort+k
// rather than the port the caller put in baseURL. Polling only the requested
// port then times out even though the server is up (the CI flake in
// TestRunner_JSONRPC_WorkflowContextThreadsThroughDispatcher). We therefore
// scan the whole increment window and return the resolved base URL, which the
// caller must use for its subsequent requests: `baseURL = waitForServer(...)`.
func waitForServer(t *testing.T, baseURL string, timeout time.Duration) string {
t.Helper()
u, err := url.Parse(baseURL)
if err != nil {
t.Fatalf("waitForServer: invalid baseURL %q: %v", baseURL, err)
}
basePort, err := strconv.Atoi(u.Port())
if err != nil {
t.Fatalf("waitForServer: no numeric port in baseURL %q: %v", baseURL, err)
}
const portWindow = 10 // matches server.Start's auto-increment attempts
deadline := time.After(timeout)
for {
select {
case <-deadline:
t.Fatalf("server did not start within %v", timeout)
t.Fatalf("server did not start within %v (scanned %s:%d-%d)", timeout, u.Hostname(), basePort, basePort+portWindow-1)
default:
}
resp, err := http.Get(baseURL + "/healthz")
if err == nil {
for off := range portWindow {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Scan is correct for the fix. Note (follow-up #530, non-blocking): probing /healthz across the window has no server-identity check — under heavy parallel overlap it could match a neighbor test's server on a lower port in the window. Rare and strictly better than the pre-fix timeout, but the durable fix is to have the runner expose its resolved port so tests never guess (tracked in #530, alongside surfacing Run's error in the 6 _ = runner.Run(ctx) sites).

cand := fmt.Sprintf("%s://%s:%d", u.Scheme, u.Hostname(), basePort+off)
resp, err := http.Get(cand + "/healthz")
if err != nil {
continue
}
_ = resp.Body.Close()
if resp.StatusCode == http.StatusOK {
return
return cand
}
}
time.Sleep(50 * time.Millisecond)
Expand Down
37 changes: 37 additions & 0 deletions forge-cli/runtime/server_ready_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,37 @@
package runtime

import (
"fmt"
"net"
"net/http"
"testing"
"time"
)

// TestWaitForServer_ResolvesAutoIncrementedPort proves the fix for the CI flake
// in TestRunner_JSONRPC_WorkflowContextThreadsThroughDispatcher: when the port
// findFreePort handed out is stolen before the runner binds, server.Start
// auto-increments and the server comes up on a higher port. waitForServer must
// discover it by scanning the increment window and return the real URL, rather
// than time out polling the requested port.
func TestWaitForServer_ResolvesAutoIncrementedPort(t *testing.T) {
ln, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatal(err)
}
actualPort := ln.Addr().(*net.TCPAddr).Port

mux := http.NewServeMux()
mux.HandleFunc("/healthz", func(w http.ResponseWriter, _ *http.Request) { w.WriteHeader(http.StatusOK) })
srv := &http.Server{Handler: mux, ReadHeaderTimeout: time.Second}
go func() { _ = srv.Serve(ln) }()
t.Cleanup(func() { _ = srv.Close() })

// The caller believes the server is one port lower (as if its requested
// port was taken and Start incremented to actualPort).
requested := fmt.Sprintf("http://127.0.0.1:%d", actualPort-1)
got := waitForServer(t, requested, 3*time.Second)
if want := fmt.Sprintf("http://127.0.0.1:%d", actualPort); got != want {
t.Errorf("waitForServer resolved %q, want %q (should scan up to the incremented port)", got, want)
}
}
3 changes: 2 additions & 1 deletion forge-cli/runtime/tracing_runner_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -87,7 +87,8 @@ func TestRunner_TracingEnabled_InstallsProviderAndShutsDownCleanly(t *testing.T)
go func() { runErrCh <- runner.Run(ctx) }()

baseURL := "http://localhost:" + itoa(port)
waitForServer(t, baseURL, 5*time.Second)
// This test only needs readiness (it cancels next), not the resolved URL.
waitForServer(t, baseURL, serverReadyTimeout)

// At this point Run() has progressed past the tracer install (which
// happens before the executor + HTTP server come up). Cancel and
Expand Down
Loading