| 1 | package draftstate |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "errors" |
| 6 | "os" |
| 7 | "os/exec" |
| 8 | "testing" |
| 9 | ) |
| 10 | |
| 11 | func TestWorkerLeaseIndependentProcess(t *testing.T) { |
| 12 | if path := os.Getenv("REASONIX_DRAFT_LOCK_FIXTURE"); path != "" { |
| 13 | s := New(path) |
| 14 | defer s.Close() |
| 15 | if release, err := s.WorkerLease("same-session"); err == nil { |
| 16 | release() |
| 17 | t.Fatal("child acquired parent's worker lease") |
| 18 | } |
| 19 | return |
| 20 | } |
| 21 | s := testStore(t) |
| 22 | release, err := s.WorkerLease("same-session") |
| 23 | if err != nil { |
| 24 | t.Fatal(err) |
| 25 | } |
| 26 | defer release() |
| 27 | child := exec.Command(os.Args[0], "-test.run=^TestWorkerLeaseIndependentProcess$") |
| 28 | child.Env = append(os.Environ(), "REASONIX_DRAFT_LOCK_FIXTURE="+s.path) |
| 29 | if output, err := child.CombinedOutput(); err != nil { |
| 30 | t.Fatalf("child: %v %s", err, output) |
| 31 | } |
| 32 | } |
| 33 | |
| 34 | func TestRequestIdentitySurvivesConversionAndRejectsChangedSnapshot(t *testing.T) { |
| 35 | s := testStore(t) |
| 36 | ctx := context.Background() |
| 37 | d, _, err := s.Open(ctx, "workspace", "global", "", "draft", `{}`) |
| 38 | if err != nil { |
| 39 | t.Fatal(err) |
| 40 | } |
| 41 | digest, err := SnapshotDigest(d.ContentJSON, d.SettingsJSON) |
| 42 | if err != nil { |
| 43 | t.Fatal(err) |
| 44 | } |
| 45 | request := Operation{ID: "operation", RequestID: "request", SourceDigest: digest, DraftID: d.ID, DraftRevision: d.Revision, WorkspaceID: d.WorkspaceID, SessionID: "session", SubmissionID: "submission", Fingerprint: "original", RequestJSON: `{}`} |
| 46 | op, created, err := s.BeginOperation(ctx, request) |
| 47 | if err != nil || !created { |
| 48 | t.Fatalf("begin: %+v %v %v", op, created, err) |
| 49 | } |
| 50 | if _, err = s.SetOperationPhase(ctx, op.ID, "dispatching", ""); err != nil { |
| 51 | t.Fatal(err) |
| 52 | } |
| 53 | if _, err = s.AcceptAndConvert(ctx, d.ID, op.ID); err != nil { |
| 54 | t.Fatal(err) |
| 55 | } |
| 56 | request.ID = "must-not-create" |
| 57 | again, created, err := s.BeginOperation(ctx, request) |
| 58 | if err != nil || created || again.ID != op.ID { |
| 59 | t.Fatalf("retry: %+v %v %v", again, created, err) |
| 60 | } |
| 61 | request.Fingerprint = "changed" |
| 62 | if _, _, err = s.BeginOperation(ctx, request); !errors.Is(err, ErrOperationConflict) { |
| 63 | t.Fatalf("changed request: %v", err) |
| 64 | } |
| 65 | } |
| 66 | |
| 67 | func TestSnapshotDigestTreatsInheritedModelAsLiveCompatibilityMirror(t *testing.T) { |
| 68 | inheritedA, err := SnapshotDigest(`{}`, `{"model":"fixture/a","modelSource":"default"}`) |
| 69 | if err != nil { |
| 70 | t.Fatal(err) |
| 71 | } |
| 72 | inheritedB, err := SnapshotDigest(`{}`, `{"model":"fixture/b","modelSource":"default"}`) |
| 73 | if err != nil { |
| 74 | t.Fatal(err) |
| 75 | } |
| 76 | if inheritedA != inheritedB { |
| 77 | t.Fatalf("inherited model mirror changed digest: %s != %s", inheritedA, inheritedB) |
| 78 | } |
| 79 | explicitA, err := SnapshotDigest(`{}`, `{"model":"fixture/a","modelSource":"explicit"}`) |
| 80 | if err != nil { |
| 81 | t.Fatal(err) |
| 82 | } |
| 83 | explicitB, err := SnapshotDigest(`{}`, `{"model":"fixture/b","modelSource":"explicit"}`) |
| 84 | if err != nil { |
| 85 | t.Fatal(err) |
| 86 | } |
| 87 | if explicitA == explicitB { |
| 88 | t.Fatal("explicit model was omitted from the draft digest") |
| 89 | } |
| 90 | } |
| 91 | |
| 92 | func TestOperationResumeCASAndWorkerExclusion(t *testing.T) { |
| 93 | s := testStore(t) |
| 94 | ctx := context.Background() |
| 95 | d, _, err := s.Open(ctx, "workspace", "global", "", "draft", `{}`) |
| 96 | if err != nil { |
| 97 | t.Fatal(err) |
| 98 | } |
| 99 | op, _, err := s.BeginOperation(ctx, Operation{ID: "op", DraftID: d.ID, DraftRevision: d.Revision, WorkspaceID: d.WorkspaceID, SessionID: "session", SubmissionID: "submission", Fingerprint: "f", RequestJSON: `{}`}) |
| 100 | if err != nil { |
| 101 | t.Fatal(err) |
| 102 | } |
| 103 | op, err = s.SetOperationPhase(ctx, op.ID, "resume_required", "") |
| 104 | if err != nil { |
| 105 | t.Fatal(err) |
| 106 | } |
| 107 | resumed, err := s.ResumeOperation(ctx, op.ID, op.Revision) |
| 108 | if err != nil || resumed.Revision <= op.Revision { |
| 109 | t.Fatalf("resume: %+v %v", resumed, err) |
| 110 | } |
| 111 | if _, err = s.ResumeOperation(ctx, op.ID, op.Revision); !errors.Is(err, ErrOperationConflict) { |
| 112 | t.Fatalf("stale resume: %v", err) |
| 113 | } |
| 114 | release, err := s.WorkerLease(op.SessionID) |
| 115 | if err != nil { |
| 116 | t.Fatal(err) |
| 117 | } |
| 118 | defer release() |
| 119 | other := New(s.path) |
| 120 | defer other.Close() |
| 121 | if unlock, err := other.WorkerLease(op.SessionID); err == nil { |
| 122 | unlock() |
| 123 | t.Fatal("second worker acquired live execution ownership") |
| 124 | } |
| 125 | } |
| 126 |