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
53 changes: 47 additions & 6 deletions internal/harness/universe/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -357,14 +357,55 @@ func check(cond bool, format string, args ...any) {
// should not depend on a live model deciding to send another notification.
func a2aReachable(ctx context.Context, base, agent string) error {
probe := "A2A reachability probe only. Reply with the words concierge reachable. Do not call tools or send notifications."
reply, err := a2a.NewClient(base+"/agents/"+agent).Send(ctx, probe)
if err != nil {
return err
deadline, ok := ctx.Deadline()
if !ok {
deadline = time.Now().Add(10 * time.Second)
}
if strings.TrimSpace(reply) == "" {
return fmt.Errorf("empty A2A reply")

var lastErr error
for attempt := 1; ; attempt++ {
if err := ctx.Err(); err != nil {
if lastErr != nil {
return fmt.Errorf("A2A reachability probe failed after %d attempt(s): %w", attempt-1, lastErr)
}
return err
}

remaining := time.Until(deadline)
if remaining <= 0 {
if lastErr != nil {
return fmt.Errorf("A2A reachability probe failed after %d attempt(s): %w", attempt-1, lastErr)
}
return context.DeadlineExceeded
}

attemptTimeout := 4 * time.Second
if remaining < attemptTimeout {
attemptTimeout = remaining
}
attemptCtx, cancel := context.WithTimeout(ctx, attemptTimeout)
reply, err := a2a.NewClient(base+"/agents/"+agent).Send(attemptCtx, probe)
cancel()
if err == nil && strings.TrimSpace(reply) != "" {
return nil
}
if err == nil {
err = fmt.Errorf("empty A2A reply")
}
lastErr = err

if time.Until(deadline) <= 0 {
return fmt.Errorf("A2A reachability probe failed after %d attempt(s): %w", attempt, lastErr)
}
time.Sleep(minDuration(200*time.Millisecond*time.Duration(attempt), time.Until(deadline)))
}
return nil
}

func minDuration(a, b time.Duration) time.Duration {
if a < b {
return a
}
return b
}

func providerKey(provider string) string {
Expand Down
28 changes: 28 additions & 0 deletions internal/harness/universe/main_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,9 @@ package main
import (
"context"
"errors"
"fmt"
"net/http"
"net/http/httptest"
"sync/atomic"
"testing"
"time"
Expand All @@ -25,6 +28,31 @@ func TestUniverseHarnessContract(t *testing.T) {
}
}

func TestA2AReachableRetriesTransientTimeout(t *testing.T) {
var calls int64
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if got, want := r.URL.Path, "/agents/concierge"; got != want {
t.Fatalf("path = %q, want %q", got, want)
}
Comment on lines +34 to +36
if atomic.AddInt64(&calls, 1) == 1 {
time.Sleep(5 * time.Second)
return
}
Comment on lines +37 to +40
w.Header().Set("Content-Type", "application/json")
fmt.Fprint(w, `{"jsonrpc":"2.0","id":1,"result":{"kind":"task","id":"task-1","contextId":"ctx-1","status":{"state":"completed"},"artifacts":[{"artifactId":"artifact-1","parts":[{"kind":"text","text":"concierge reachable"}]}]}}`)
}))
defer srv.Close()

ctx, cancel := context.WithTimeout(context.Background(), 6*time.Second)
defer cancel()
if err := a2aReachable(ctx, srv.URL, "concierge"); err != nil {
t.Fatalf("a2aReachable returned error: %v", err)
}
if got := atomic.LoadInt64(&calls); got < 2 {
t.Fatalf("A2A calls = %d, want retry after transient timeout", got)
}
}

func TestNotifyStepCompletesAfterObservedSideEffectTimeout(t *testing.T) {
ntf := new(Notify)
before := atomic.LoadInt64(&ntf.sent)
Expand Down
Loading