| 1 | package serve |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "encoding/json" |
| 6 | "io" |
| 7 | "net/http" |
| 8 | "net/http/httptest" |
| 9 | "path/filepath" |
| 10 | "sync/atomic" |
| 11 | "testing" |
| 12 | "time" |
| 13 | |
| 14 | "reasonix/internal/agent" |
| 15 | "reasonix/internal/boot" |
| 16 | "reasonix/internal/config" |
| 17 | "reasonix/internal/control" |
| 18 | "reasonix/internal/jobs" |
| 19 | ) |
| 20 | |
| 21 | func modelApplicationServer(t *testing.T) (*Server, *jobs.SessionBackgroundScope, *atomic.Int32) { |
| 22 | t.Helper() |
| 23 | bc := NewBroadcaster() |
| 24 | scope := jobs.NewSessionBackgroundScope(jobs.NewManager(bc), nil) |
| 25 | options := control.Options{Sink: bc, BackgroundSink: bc, Jobs: scope.Manager, BackgroundScope: scope, SessionDir: t.TempDir(), ModelRef: "p/m", ModelSettingsRevision: "old", ModelSettingsCurrent: func() (string, error) { return "new", nil }, ModelSettingsContinuation: func() error { return nil }} |
| 26 | options.SessionPath = filepath.Join(options.SessionDir, "conversation.jsonl") |
| 27 | old := control.New(options) |
| 28 | old.PublishBackgroundScope() |
| 29 | s := New(old, bc, config.ServeConfig{AuthMode: "none"}) |
| 30 | t.Cleanup(s.Close) |
| 31 | builds := &atomic.Int32{} |
| 32 | s.buildControllerWithOptions = func(_ context.Context, _ string, opts boot.Options) (*control.Controller, error) { |
| 33 | builds.Add(1) |
| 34 | if err := opts.BackgroundScope.Acquire(); err != nil { |
| 35 | return nil, err |
| 36 | } |
| 37 | next := options |
| 38 | next.Sink, next.BackgroundSink = opts.Sink, opts.Sink |
| 39 | next.BackgroundScope = opts.BackgroundScope |
| 40 | next.SessionTemp, next.PersistentShell = opts.SessionTemp, opts.PersistentShell |
| 41 | next.ModelSettingsRevision = "new" |
| 42 | return control.New(next), nil |
| 43 | } |
| 44 | return s, scope, builds |
| 45 | } |
| 46 | |
| 47 | func TestModelApplicationServePreservesProcess(t *testing.T) { |
| 48 | s, scope, builds := modelApplicationServer(t) |
| 49 | entered := make(chan struct{}) |
| 50 | job := scope.Manager.StartSessionProcess(agent.BranchID(s.ctl().SessionPath()), "bash", "gateway", func(ctx context.Context, out io.Writer) (string, error) { |
| 51 | close(entered) |
| 52 | <-ctx.Done() |
| 53 | return "", ctx.Err() |
| 54 | }) |
| 55 | <-entered |
| 56 | old := s.ctl() |
| 57 | s.bindMu.Lock() |
| 58 | err := s.refreshRunModelSettingsLocked(t.Context()) |
| 59 | s.bindMu.Unlock() |
| 60 | if err != nil { |
| 61 | t.Fatal(err) |
| 62 | } |
| 63 | if old == s.ctl() || builds.Load() != 1 || len(scope.Manager.Running()) != 1 { |
| 64 | t.Fatal("model replacement did not preserve gateway") |
| 65 | } |
| 66 | if !s.ctl().(*control.Controller).CancelJob(job.ID) { |
| 67 | t.Fatal("replacement cannot cancel gateway") |
| 68 | } |
| 69 | scope.Manager.Wait(t.Context(), []string{job.ID}, 2) |
| 70 | if len(scope.Manager.Running()) != 0 { |
| 71 | t.Fatal("cancel did not drain") |
| 72 | } |
| 73 | } |
| 74 | |
| 75 | func TestModelApplicationServeBlockedChoiceAndEventDrivenApply(t *testing.T) { |
| 76 | s, scope, builds := modelApplicationServer(t) |
| 77 | built := make(chan struct{}) |
| 78 | factory := s.buildControllerWithOptions |
| 79 | s.buildControllerWithOptions = func(ctx context.Context, ref string, opts boot.Options) (*control.Controller, error) { |
| 80 | defer close(built) |
| 81 | return factory(ctx, ref, opts) |
| 82 | } |
| 83 | entered, unwind := make(chan struct{}), make(chan struct{}) |
| 84 | job := scope.Manager.StartForSession(agent.BranchID(s.ctl().SessionPath()), "task", "dependent", func(ctx context.Context, _ io.Writer) (string, error) { |
| 85 | close(entered) |
| 86 | <-ctx.Done() |
| 87 | <-unwind |
| 88 | return "", ctx.Err() |
| 89 | }) |
| 90 | <-entered |
| 91 | s.bindMu.Lock() |
| 92 | w := httptest.NewRecorder() |
| 93 | if s.admitModelSettingsRunLocked(w, httptest.NewRequest(http.MethodPost, "/submit", nil)) { |
| 94 | t.Error("dependent task admitted rebuild") |
| 95 | } |
| 96 | var response struct { |
| 97 | Data struct { |
| 98 | Outcome string `json:"submissionOutcome"` |
| 99 | Details control.ModelApplicationDetails `json:"modelApplication"` |
| 100 | } `json:"data"` |
| 101 | } |
| 102 | if err := json.Unmarshal(w.Body.Bytes(), &response); err != nil { |
| 103 | t.Fatal(err) |
| 104 | } |
| 105 | d := response.Data.Details |
| 106 | choice := control.ModelApplicationChoice{Mode: "applied_once", ExpectedRuntimeIdentity: d.RuntimeIdentity, ExpectedAppliedRevision: d.AppliedRevision, ExpectedDesiredRevision: d.DesiredRevision} |
| 107 | if response.Data.Outcome != "not_accepted" || len(d.BlockingJobs) != 1 || s.validateAppliedModelChoiceLocked(t.Context(), choice) != nil { |
| 108 | t.Error("missing recovery contract") |
| 109 | } |
| 110 | s.bindMu.Unlock() |
| 111 | scope.Manager.Kill(job.ID) |
| 112 | if !control.ModelReplacementBlocked(s.ctl()) { |
| 113 | t.Error("cancelling task no longer blocks") |
| 114 | } |
| 115 | close(unwind) |
| 116 | scope.Manager.Wait(t.Context(), []string{job.ID}, 2) |
| 117 | // A lifecycle event wakes the existing coalesced owner; no polling loop. |
| 118 | s.kickModelApplication() |
| 119 | select { |
| 120 | case <-built: |
| 121 | case <-time.After(3 * time.Second): |
| 122 | t.Fatal("deferred apply did not start") |
| 123 | } |
| 124 | s.bindMu.Lock() |
| 125 | defer s.bindMu.Unlock() |
| 126 | if builds.Load() != 1 { |
| 127 | t.Fatalf("builds=%d", builds.Load()) |
| 128 | } |
| 129 | if s.validateAppliedModelChoiceLocked(t.Context(), choice) == nil { |
| 130 | t.Fatal("old confirmation survived replacement") |
| 131 | } |
| 132 | } |
| 133 |