test: prove native reconnect recovery

This commit is contained in:
2026-09-11 08:08:32 +00:00
parent e79f882993
commit f86abcecb5
9 changed files with 158 additions and 17 deletions
+11 -4
View File
@@ -35,6 +35,10 @@ type RunnerOptions struct {
OnSignal func(context.Context, Session, *rvboxv1.SignalCommand) error
OnScriptReady func(context.Context, Session, domain.UUID) error
OnTerminate func(context.Context, domain.UUID) error
// OnSessionError observes one failed dial, handshake, protocol, or active
// transport session before normal reconnect backoff. It must not block; the
// durable spool and retry policy remain owned by Run.
OnSessionError func(error)
// EventReady wakes the active session after a supervisor worker appends a
// durable event. The network loop remains the sole writer; a reconnect can
// safely ignore a stale notification because replay reads the spool again.
@@ -95,10 +99,13 @@ func Run(ctx context.Context, options RunnerOptions) error {
return nil
}
sessionStarted := options.Now()
_ = runOnce(ctx, options)
sessionErr := runOnce(ctx, options)
if ctx.Err() != nil {
return nil
}
if sessionErr != nil && options.OnSessionError != nil {
options.OnSessionError(sessionErr)
}
if options.Now().Sub(sessionStarted) >= options.Backoff.StableReset {
failures = 0
}
@@ -344,17 +351,17 @@ func serveActive(ctx context.Context, transport Transport, options RunnerOptions
switch {
case envelope.GetCommandDispatch() != nil:
if err := handleDispatch(ctx, transport, options, session, envelope.GetCommandDispatch(), limits); err != nil {
return err
return fmt.Errorf("handle command dispatch %s: %w", envelope.GetCommandDispatch().GetIssueUuid(), err)
}
issue, parseErr := domain.ParseUUIDv7(envelope.GetCommandDispatch().GetIssueUuid())
if parseErr == nil {
if err := flushEvents(ctx, transport, options.Store, session, issue, limits, sent); err != nil {
return err
return fmt.Errorf("flush dispatched command %s events: %w", issue, err)
}
}
case envelope.GetEventAck() != nil:
if err := ApplyEventAck(ctx, options.Store, envelope.GetEventAck()); err != nil {
return err
return fmt.Errorf("apply event acknowledgement for %s: %w", envelope.GetEventAck().GetIssueUuid(), err)
}
case envelope.GetScriptChunk() != nil:
if err := handleScriptChunk(ctx, transport, options.Store, session, envelope.GetScriptChunk(), limits, options.Now, sent); err != nil {
+30
View File
@@ -72,6 +72,36 @@ func TestRunnerOptionsRejectMissingJitter_HP_RUNTIME_02(t *testing.T) {
}
}
func TestRunReportsRetryableSessionError_HP_RUNTIME_04(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
store, err := spool.Open(ctx, spool.Options{DataDir: filepath.Join(t.TempDir(), "spool"), BusyTimeout: time.Second})
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = store.Close() })
want := errors.New("dial refused")
reported := make(chan error, 1)
hello := &rvboxv1.ClientHello{ClientId: "runner-client", SupportedProtocol: &rvboxv1.ProtocolRange{Major: 1, MinMinor: 0, MaxMinor: 0}, DaemonVersion: "test", Platform: rvboxv1.Platform_PLATFORM_WINDOWS, Architecture: "amd64", DaemonCwd: `C:\\`, SupportedShells: []rvboxv1.ShellType{rvboxv1.ShellType_SHELL_POWERSHELL}, ClientInstanceId: store.ClientInstanceID().String(), MaxRunningCommands: 1, MaxQueuedCommands: 1, SentAt: timestamppb.Now()}
err = Run(ctx, RunnerOptions{
Store: store, Dial: func(context.Context) (Transport, error) { return nil, want }, Hello: hello,
Limits: agentproto.DefaultLimits(), Backoff: BackoffOptions{Initial: time.Millisecond, Maximum: time.Millisecond, StableReset: time.Second},
Jitter: func(time.Duration) time.Duration { return 0 }, Now: func() time.Time { return time.Now().UTC() },
OnSessionError: func(got error) { reported <- got; cancel() },
})
if err != nil {
t.Fatalf("Run error = %v, want graceful cancellation", err)
}
select {
case got := <-reported:
if !errors.Is(got, want) {
t.Fatalf("reported error = %v, want %v", got, want)
}
default:
t.Fatal("retryable session error was not reported")
}
}
func TestFlushEventsAssignsAndSendsOnlyUnacknowledgedRows_HP_RUNTIME_03(t *testing.T) {
ctx := context.Background()
store, err := spool.Open(ctx, spool.Options{DataDir: filepath.Join(t.TempDir(), "spool"), BusyTimeout: time.Second})