diff --git a/CHANGELOG.md b/CHANGELOG.md index e9aca020..761933de 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Fixed +- **WDA keeps enough idle connections for its own parallel reads.** The driver reads an element's name, rect, text and displayed at once, and a tap looks an element up four ways at once, but Go's default transport keeps two idle connections per host, so every burst closed two connections and opened two new ones. Through a forwarded port to a physical iPhone, new connections opened together fail with EOF and are sent again, which costs time on every step. The WDA client now keeps up to eight. +- **`retry` counts retries, not attempts, as Maestro does.** `maxRetries: 1` now runs the commands twice (once, then one retry), an unset `maxRetries` means one retry, and the value is capped at 3. The runner ran exactly `maxRetries` attempts, three when unset and with no cap, so `maxRetries: 1` never retried. A value that is not an integer is logged and read as 1 instead of failing the step. + ## [1.1.28] - 2026-09-30 This release is about **running a Maestro suite and getting Maestro's answer**. Most of it came from running real Maestro suites (DuckDuckGo Android and iOS, React Native's RNTester, React Navigation) side by side with Maestro and fixing each place the two disagreed. The largest of those is selector matching: text and id selectors now match the whole value, as Maestro does, which is a behaviour change — see below. `--driver devicelab` on iOS now runs a new agent, on simulators and on real iPhones. DeviceLab becomes the default Android driver. Every flow gets `DEVICE_UDID`, `-e` stops splitting values at commas, and the Appium driver closes a long list of gaps with the native drivers. diff --git a/pkg/driver/wda/client.go b/pkg/driver/wda/client.go index 56a0d35b..5078bc4e 100644 --- a/pkg/driver/wda/client.go +++ b/pkg/driver/wda/client.go @@ -34,11 +34,32 @@ func NewClient(port uint16) *Client { return &Client{ baseURL: fmt.Sprintf("http://127.0.0.1:%d", port), httpClient: &http.Client{ - Timeout: 60 * time.Second, + Timeout: 60 * time.Second, + Transport: newWDATransport(), }, } } +// wdaIdleConnsPerHost is how many idle connections the client keeps open to WebDriverAgent. +// The driver sends up to four requests at once (an element's name, rect, text and displayed; +// a tap's four lookups), and Go keeps two idle connections per host by default, so every such +// burst closed two connections and opened two new ones. On a real iPhone reached through a +// forward (SSH, then iproxy), a new connection took about 300 ms, and the two opened together +// failed at once with EOF in 297 of 996 bursts of one 44-flow run, always both of them and +// never a kept connection. Kept for every request of a burst, no new ones are needed. +const wdaIdleConnsPerHost = 8 + +// newWDATransport is Go's default transport with room to keep a whole burst's connections. +func newWDATransport() *http.Transport { + base, ok := http.DefaultTransport.(*http.Transport) + if !ok { + base = &http.Transport{Proxy: http.ProxyFromEnvironment} + } + t := base.Clone() + t.MaxIdleConnsPerHost = wdaIdleConnsPerHost + return t +} + // Session management // CreateSession creates a new WDA session. diff --git a/pkg/driver/wda/kept_connections_test.go b/pkg/driver/wda/kept_connections_test.go new file mode 100644 index 00000000..ea0b8652 --- /dev/null +++ b/pkg/driver/wda/kept_connections_test.go @@ -0,0 +1,115 @@ +package wda + +import ( + "fmt" + "net" + "net/http" + "net/http/httptest" + "strconv" + "strings" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/devicelab-dev/maestro-runner/pkg/core" +) + +// keptConnServer answers an element's four property reads, each after a short delay so the +// four of a burst are in flight together, as they are on a phone. It counts the connections +// it accepted, and with dropNewAfter > 0 closes, without an answer, every new connection past +// that many: a forward whose new connections fail, while kept ones work. +type keptConnServer struct { + server *httptest.Server + accepted int32 + dropped int32 + dropNewAfter int32 +} + +func newKeptConnServer(t *testing.T, dropNewAfter int32) *keptConnServer { + t.Helper() + es := &keptConnServer{dropNewAfter: dropNewAfter} + handler := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + time.Sleep(20 * time.Millisecond) + switch { + case strings.HasSuffix(r.URL.Path, "/rect"): + jsonResponse(w, map[string]interface{}{"value": map[string]interface{}{"x": 10, "y": 20, "width": 100, "height": 40}}) + case strings.HasSuffix(r.URL.Path, "/displayed"): + jsonResponse(w, map[string]interface{}{"value": true}) + default: + jsonResponse(w, map[string]interface{}{"value": "Continue"}) + } + }) + es.server = httptest.NewUnstartedServer(handler) + es.server.Config.ConnState = func(c net.Conn, state http.ConnState) { + if state != http.StateNew { + return + } + n := atomic.AddInt32(&es.accepted, 1) + if es.dropNewAfter > 0 && n > es.dropNewAfter { + atomic.AddInt32(&es.dropped, 1) + _ = c.Close() + } + } + es.server.Start() + t.Cleanup(es.server.Close) + return es +} + +// driverFor is a driver whose client is built the way the runner builds it. +func (es *keptConnServer) driverFor(t *testing.T) *Driver { + t.Helper() + _, portText, err := net.SplitHostPort(strings.TrimPrefix(es.server.URL, "http://")) + if err != nil { + t.Fatalf("server address %q: %v", es.server.URL, err) + } + port, _ := strconv.Atoi(portText) + client := NewClient(uint16(port)) + client.baseURL = es.server.URL // the listener is on 127.0.0.1, not localhost's first address + client.sessionID = "s1" + return &Driver{client: client, info: &core.PlatformInfo{Platform: "ios", ScreenWidth: 390, ScreenHeight: 844}} +} + +// Fifty elements' reads, four at a time, open four connections and then keep them, instead of +// closing two and opening two for every element. +func TestElementReadsKeepTheirConnections(t *testing.T) { + es := newKeptConnServer(t, 0) + d := es.driverFor(t) + for i := 0; i < 50; i++ { + if _, err := d.getElementInfo(fmt.Sprintf("E%d", i)); err != nil { + t.Fatalf("element %d: %v", i, err) + } + } + if n := atomic.LoadInt32(&es.accepted); n > 4 { + t.Errorf("50 bursts of four reads opened %d connections, want at most 4", n) + } +} + +// Where a new connection fails (the phone's forward), only the first burst opens any, so the +// reads after it never meet a failing one. +func TestElementReadsSurviveAForwardWhoseNewConnectionsFail(t *testing.T) { + es := newKeptConnServer(t, 4) + d := es.driverFor(t) + var mu sync.Mutex + failed := 0 + for i := 0; i < 50; i++ { + if _, err := d.getElementInfo(fmt.Sprintf("E%d", i)); err != nil { + mu.Lock() + failed++ + mu.Unlock() + } + } + if dropped := atomic.LoadInt32(&es.dropped); dropped != 0 || failed != 0 { + t.Errorf("after the first burst, %d new connections were needed and dropped, and %d reads failed; want none", dropped, failed) + } +} + +func TestTheWDATransportKeepsABurstsConnections(t *testing.T) { + tr := newWDATransport() + if tr.MaxIdleConnsPerHost < 4 { + t.Errorf("MaxIdleConnsPerHost = %d, fewer than the driver's four requests at once", tr.MaxIdleConnsPerHost) + } + if tr.Proxy == nil || tr.DialContext == nil { + t.Error("the WDA transport lost the default transport's proxy or dialer") + } +} diff --git a/pkg/executor/flow_runner.go b/pkg/executor/flow_runner.go index 77a86bf2..a9cb8f3a 100644 --- a/pkg/executor/flow_runner.go +++ b/pkg/executor/flow_runner.go @@ -5,6 +5,7 @@ import ( "fmt" "os" "path/filepath" + "strconv" "strings" "time" @@ -1091,16 +1092,53 @@ func (fr *FlowRunner) executeRepeat(step *flow.RepeatStep) *core.CommandResult { } } -// executeRetry handles retry step execution. -func (fr *FlowRunner) executeRetry(step *flow.RetryStep) *core.CommandResult { - maxRetries, err := fr.script.ParseIntStrict(step.MaxRetries, 3) - if err != nil { +// maxRetriesAllowed caps retry's maxRetries, as Maestro's MAX_RETRIES_ALLOWED +// does (Orchestra.kt:1844). +const maxRetriesAllowed = 3 + +// retryAttempts is how many times a retry runs its commands: the first run +// plus one per retry. maxRetries is read as Maestro reads it (Orchestra.kt:934 +// and 939-955): unset or not an integer is 1, more than 3 is 3, and a negative +// value runs nothing, because Maestro's `while (attempt <= maxRetries)` never +// starts. +func (fr *FlowRunner) retryAttempts(raw string) int { + maxRetries := 1 + if expanded := fr.script.ExpandVariables(raw); expanded != "" { + if n, err := strconv.Atoi(expanded); err == nil { + maxRetries = n + } else { + logger.Warn("retry: maxRetries %q is not an integer, so it is 1, as in Maestro", expanded) + } + } + if maxRetries > maxRetriesAllowed { + maxRetries = maxRetriesAllowed + } + if maxRetries < 0 { + return 0 + } + return maxRetries + 1 +} + +// retryExhausted is the result of a retry that ran out of attempts. With no +// attempt at all (a negative maxRetries) Maestro's command completes without +// running anything, so that one passes. +func retryExhausted(lastErr error, attempts int) *core.CommandResult { + if attempts == 0 { return &core.CommandResult{ - Success: false, - Error: err, - Message: fmt.Sprintf("retry: invalid 'maxRetries' value: %v", err), + Success: true, + Message: "Retry ran no attempts (maxRetries is negative)", } } + return &core.CommandResult{ + Success: false, + Error: lastErr, + Message: fmt.Sprintf("Retry failed after %d attempts", attempts), + } +} + +// executeRetry handles retry step execution. +func (fr *FlowRunner) executeRetry(step *flow.RetryStep) *core.CommandResult { + attempts := fr.retryAttempts(step.MaxRetries) // Apply env variables with restore defer fr.script.withEnvVars(step.Env)() @@ -1116,12 +1154,12 @@ func (fr *FlowRunner) executeRetry(step *flow.RetryStep) *core.CommandResult { Message: fmt.Sprintf("Failed to parse flow file: %s", filePath), } } - return fr.executeSubFlowWithRetry(*subFlow, maxRetries) + return fr.executeSubFlowWithRetry(*subFlow, attempts) } // Execute inline steps with retry var lastErr error - for attempt := 1; attempt <= maxRetries; attempt++ { + for attempt := 1; attempt <= attempts; attempt++ { if fr.ctx.Err() != nil { return &core.CommandResult{ Success: false, @@ -1148,11 +1186,7 @@ func (fr *FlowRunner) executeRetry(step *flow.RetryStep) *core.CommandResult { } } - return &core.CommandResult{ - Success: false, - Error: lastErr, - Message: fmt.Sprintf("Retry failed after %d attempts", maxRetries), - } + return retryExhausted(lastErr, attempts) } // executeRunFlow handles runFlow step execution. @@ -1721,11 +1755,11 @@ func (fr *FlowRunner) executeSubFlow(subFlow flow.Flow) *core.CommandResult { } } -// executeSubFlowWithRetry executes a sub-flow with retry logic. -func (fr *FlowRunner) executeSubFlowWithRetry(subFlow flow.Flow, maxRetries int) *core.CommandResult { +// executeSubFlowWithRetry runs a sub-flow up to attempts times, until it passes. +func (fr *FlowRunner) executeSubFlowWithRetry(subFlow flow.Flow, attempts int) *core.CommandResult { var lastErr error - for attempt := 1; attempt <= maxRetries; attempt++ { + for attempt := 1; attempt <= attempts; attempt++ { if fr.ctx.Err() != nil { return &core.CommandResult{ Success: false, @@ -1744,11 +1778,7 @@ func (fr *FlowRunner) executeSubFlowWithRetry(subFlow flow.Flow, maxRetries int) lastErr = result.Error } - return &core.CommandResult{ - Success: false, - Error: lastErr, - Message: fmt.Sprintf("Retry failed after %d attempts", maxRetries), - } + return retryExhausted(lastErr, attempts) } // captureArtifacts captures the step screenshot and, when captureHierarchy is diff --git a/pkg/executor/retry_attempts_test.go b/pkg/executor/retry_attempts_test.go new file mode 100644 index 00000000..6fbc9495 --- /dev/null +++ b/pkg/executor/retry_attempts_test.go @@ -0,0 +1,133 @@ +package executor + +import ( + "os" + "path/filepath" + "testing" + + "github.com/devicelab-dev/maestro-runner/pkg/core" + "github.com/devicelab-dev/maestro-runner/pkg/flow" + "github.com/devicelab-dev/maestro-runner/pkg/report" +) + +// Maestro reads maxRetries as retries after the first run, not as attempts: +// `(maxRetries ?: 1).coerceAtMost(3)`, then `while (attempt <= maxRetries)` +// from 0 (Orchestra.kt:934-955). +func TestRetryAttempts_CountsRetriesAsMaestroDoes(t *testing.T) { + se := NewScriptEngine() + defer se.Close() + se.SetVariable("TWO", "2") + fr := &FlowRunner{script: se} + + for _, tc := range []struct { + maxRetries string + want int + }{ + {"", 2}, // unset: one retry + {"0", 1}, // no retry, one run + {"1", 2}, // one retry + {"3", 4}, // the most there can be + {"10", 4}, // capped at three retries + {"-1", 0}, // the loop never starts + {"${TWO}", 3}, // an expression is evaluated first + {"${1 + 1}", 3}, // and so is arithmetic + {"${UNSET_VAR}", 2}, // evaluates to nothing: one retry + {"five", 2}, // not an integer: one retry, not an error + {"1.5", 2}, // toIntOrNull rejects a decimal + } { + if got := fr.retryAttempts(tc.maxRetries); got != tc.want { + t.Errorf("maxRetries %q: %d attempts, want %d", tc.maxRetries, got, tc.want) + } + } +} + +// failingTaps is a driver whose taps fail the first failures times and then +// pass; every other step passes. It counts the taps. +func failingTaps(failures int, taps *int) *mockDriver { + return &mockDriver{executeFunc: func(step flow.Step) *core.CommandResult { + if _, ok := step.(*flow.TapOnStep); ok { + *taps++ + if *taps <= failures { + return &core.CommandResult{Success: false, Error: &testError{msg: "not there yet"}} + } + } + return &core.CommandResult{Success: true} + }} +} + +func retryAroundTap(maxRetries string) *flow.RetryStep { + return &flow.RetryStep{ + BaseStep: flow.BaseStep{StepType: flow.StepRetry}, + MaxRetries: maxRetries, + Steps: []flow.Step{&flow.TapOnStep{BaseStep: flow.BaseStep{StepType: flow.StepTapOn}}}, + } +} + +// `maxRetries: 1` ran the commands once and never retried. +func TestRetry_MaxRetriesOneRetriesOnce(t *testing.T) { + taps := 0 + result := runOneFlow(t, failingTaps(1, &taps), flow.Flow{ + SourcePath: "test.yaml", + Config: flow.Config{Name: "retry once"}, + Steps: []flow.Step{retryAroundTap("1")}, + }) + if result.Status != report.StatusPassed { + t.Errorf("status = %v, want passed: the retry should run the tap a second time", result.Status) + } + if taps != 2 { + t.Errorf("tap ran %d times, want 2", taps) + } +} + +func TestRetry_AttemptsWhenEveryRunFails(t *testing.T) { + for _, tc := range []struct { + maxRetries string + taps int + status report.Status + }{ + {"", 2, report.StatusFailed}, + {"0", 1, report.StatusFailed}, + {"1", 2, report.StatusFailed}, + {"10", 4, report.StatusFailed}, + // A negative maxRetries runs nothing, and Maestro's command then + // completes without an error. + {"-1", 0, report.StatusPassed}, + } { + taps := 0 + result := runOneFlow(t, failingTaps(1000, &taps), flow.Flow{ + SourcePath: "test.yaml", + Config: flow.Config{Name: "retry " + tc.maxRetries}, + Steps: []flow.Step{retryAroundTap(tc.maxRetries)}, + }) + if taps != tc.taps { + t.Errorf("maxRetries %q: tap ran %d times, want %d", tc.maxRetries, taps, tc.taps) + } + if result.Status != tc.status { + t.Errorf("maxRetries %q: status = %v, want %v", tc.maxRetries, result.Status, tc.status) + } + } +} + +// The file form counts the same way. +func TestRetry_FileFormRetriesOnce(t *testing.T) { + dir := t.TempDir() + if err := os.WriteFile(filepath.Join(dir, "tap.yaml"), []byte("appId: com.example.app\n---\n- tapOn: OK\n"), 0o644); err != nil { + t.Fatal(err) + } + taps := 0 + result := runOneFlow(t, failingTaps(1000, &taps), flow.Flow{ + SourcePath: filepath.Join(dir, "main.yaml"), + Config: flow.Config{Name: "retry file"}, + Steps: []flow.Step{&flow.RetryStep{ + BaseStep: flow.BaseStep{StepType: flow.StepRetry}, + MaxRetries: "1", + File: "tap.yaml", + }}, + }) + if taps != 2 { + t.Errorf("tap ran %d times, want 2 (one run and one retry)", taps) + } + if result.Status != report.StatusFailed { + t.Errorf("status = %v, want failed", result.Status) + } +} diff --git a/pkg/executor/scripting_test.go b/pkg/executor/scripting_test.go index 4fc9211e..5e4eaf71 100644 --- a/pkg/executor/scripting_test.go +++ b/pkg/executor/scripting_test.go @@ -660,20 +660,6 @@ func TestExecuteRepeat_InvalidTimes(t *testing.T) { } } -func TestExecuteRetry_InvalidMaxRetries(t *testing.T) { - se := NewScriptEngine() - defer se.Close() - fr := &FlowRunner{ctx: context.Background(), driver: &mockDriver{}, script: se} - - result := fr.executeRetry(&flow.RetryStep{MaxRetries: "five"}) - if result.Success { - t.Error("expected failure for non-numeric maxRetries, got success") - } - if !strings.Contains(result.Message, "invalid 'maxRetries'") { - t.Errorf("message = %q, want it to mention invalid 'maxRetries'", result.Message) - } -} - func TestScriptEngine_withEnvVars(t *testing.T) { se := NewScriptEngine() defer se.Close()