| 1 | package workspacestate |
| 2 | |
| 3 | import ( |
| 4 | "bufio" |
| 5 | "context" |
| 6 | "errors" |
| 7 | "fmt" |
| 8 | "io" |
| 9 | "os" |
| 10 | "os/exec" |
| 11 | "path/filepath" |
| 12 | "strconv" |
| 13 | "strings" |
| 14 | "testing" |
| 15 | "time" |
| 16 | ) |
| 17 | |
| 18 | type lifecycleProcess struct { |
| 19 | cmd *exec.Cmd |
| 20 | stdin io.WriteCloser |
| 21 | scan *bufio.Scanner |
| 22 | stderr *strings.Builder |
| 23 | } |
| 24 | |
| 25 | func startLifecycleProcess(t *testing.T, path, action string, expected uint64) *lifecycleProcess { |
| 26 | t.Helper() |
| 27 | ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) |
| 28 | t.Cleanup(cancel) |
| 29 | cmd := exec.CommandContext(ctx, os.Args[0], "-test.run=^TestWorkspaceLifecycleProcessHelper$") |
| 30 | cmd.Env = append(os.Environ(), |
| 31 | "REASONIX_WORKSPACE_PROCESS_HELPER=1", |
| 32 | "REASONIX_WORKSPACE_PROCESS_PATH="+path, |
| 33 | "REASONIX_WORKSPACE_PROCESS_ACTION="+action, |
| 34 | "REASONIX_WORKSPACE_PROCESS_EXPECTED="+strconv.FormatUint(expected, 10), |
| 35 | ) |
| 36 | stdin, err := cmd.StdinPipe() |
| 37 | if err != nil { |
| 38 | t.Fatal(err) |
| 39 | } |
| 40 | stdout, err := cmd.StdoutPipe() |
| 41 | if err != nil { |
| 42 | t.Fatal(err) |
| 43 | } |
| 44 | var stderr strings.Builder |
| 45 | cmd.Stderr = &stderr |
| 46 | if err := cmd.Start(); err != nil { |
| 47 | t.Fatal(err) |
| 48 | } |
| 49 | child := &lifecycleProcess{cmd: cmd, stdin: stdin, scan: bufio.NewScanner(stdout), stderr: &stderr} |
| 50 | if !child.scan.Scan() || child.scan.Text() != "ready" { |
| 51 | _ = cmd.Process.Kill() |
| 52 | t.Fatalf("%s child did not become ready: line=%q err=%v stderr=%s", action, child.scan.Text(), child.scan.Err(), stderr.String()) |
| 53 | } |
| 54 | return child |
| 55 | } |
| 56 | |
| 57 | func (p *lifecycleProcess) run(t *testing.T) string { |
| 58 | t.Helper() |
| 59 | if _, err := io.WriteString(p.stdin, "go\n"); err != nil { |
| 60 | t.Fatal(err) |
| 61 | } |
| 62 | if err := p.stdin.Close(); err != nil { |
| 63 | t.Fatal(err) |
| 64 | } |
| 65 | if !p.scan.Scan() { |
| 66 | _ = p.cmd.Process.Kill() |
| 67 | t.Fatalf("child result missing: err=%v stderr=%s", p.scan.Err(), p.stderr.String()) |
| 68 | } |
| 69 | result := p.scan.Text() |
| 70 | if err := p.cmd.Wait(); err != nil { |
| 71 | t.Fatalf("child failed: %v stderr=%s", err, p.stderr.String()) |
| 72 | } |
| 73 | return result |
| 74 | } |
| 75 | |
| 76 | func seedArchivedProcessState(t *testing.T) (*Store, uint64) { |
| 77 | t.Helper() |
| 78 | store := NewStore(filepath.Join(t.TempDir(), "state.json")) |
| 79 | ctx := t.Context() |
| 80 | if err := store.EnsureWorkspace(ctx, Workspace{ID: GlobalWorkspaceID, Visible: true}); err != nil { |
| 81 | t.Fatal(err) |
| 82 | } |
| 83 | if err := store.AttachSession(ctx, "", GlobalWorkspaceID, "victim", ""); err != nil { |
| 84 | t.Fatal(err) |
| 85 | } |
| 86 | if err := store.ArchiveSession(ctx, "victim"); err != nil { |
| 87 | t.Fatal(err) |
| 88 | } |
| 89 | state, err := store.Load(ctx) |
| 90 | if err != nil { |
| 91 | t.Fatal(err) |
| 92 | } |
| 93 | return store, state.Generation |
| 94 | } |
| 95 | |
| 96 | func TestWorkspaceLifecycleCommitOrderAcrossProcesses(t *testing.T) { |
| 97 | t.Run("restore wins", func(t *testing.T) { |
| 98 | store, expected := seedArchivedProcessState(t) |
| 99 | restore := startLifecycleProcess(t, store.Path(), "restore", expected) |
| 100 | purge := startLifecycleProcess(t, store.Path(), "purge", expected) |
| 101 | if result := restore.run(t); result != "ok" { |
| 102 | t.Fatalf("restore result = %q", result) |
| 103 | } |
| 104 | if result := purge.run(t); result != "conflict" { |
| 105 | t.Fatalf("purge result = %q", result) |
| 106 | } |
| 107 | state, err := store.Load(t.Context()) |
| 108 | if err != nil { |
| 109 | t.Fatal(err) |
| 110 | } |
| 111 | if state.SessionStates["victim"].Lifecycle != Active || ClassifyPurge(state, "victim") != PurgeAbsent { |
| 112 | t.Fatalf("restore-first state = lifecycle=%+v purge=%v", state.SessionStates["victim"], ClassifyPurge(state, "victim")) |
| 113 | } |
| 114 | }) |
| 115 | |
| 116 | t.Run("purge wins", func(t *testing.T) { |
| 117 | store, expected := seedArchivedProcessState(t) |
| 118 | restore := startLifecycleProcess(t, store.Path(), "restore", expected) |
| 119 | purge := startLifecycleProcess(t, store.Path(), "purge", expected) |
| 120 | if result := purge.run(t); result != "ok" { |
| 121 | t.Fatalf("purge result = %q", result) |
| 122 | } |
| 123 | if result := restore.run(t); result != "conflict" { |
| 124 | t.Fatalf("restore result = %q", result) |
| 125 | } |
| 126 | state, err := store.Load(t.Context()) |
| 127 | if err != nil { |
| 128 | t.Fatal(err) |
| 129 | } |
| 130 | if state.SessionStates["victim"].Lifecycle != Deleted || ClassifyPurge(state, "victim") != PurgeTombstoned { |
| 131 | t.Fatalf("purge-first state = lifecycle=%+v purge=%v", state.SessionStates["victim"], ClassifyPurge(state, "victim")) |
| 132 | } |
| 133 | }) |
| 134 | } |
| 135 | |
| 136 | func TestWorkspaceLifecycleProcessHelper(t *testing.T) { |
| 137 | if os.Getenv("REASONIX_WORKSPACE_PROCESS_HELPER") != "1" { |
| 138 | return |
| 139 | } |
| 140 | path := os.Getenv("REASONIX_WORKSPACE_PROCESS_PATH") |
| 141 | action := os.Getenv("REASONIX_WORKSPACE_PROCESS_ACTION") |
| 142 | expected, err := strconv.ParseUint(os.Getenv("REASONIX_WORKSPACE_PROCESS_EXPECTED"), 10, 64) |
| 143 | if err != nil { |
| 144 | t.Fatal(err) |
| 145 | } |
| 146 | store := NewStore(path) |
| 147 | var observed Operation |
| 148 | if action == "resume" { |
| 149 | state, loadErr := store.Load(t.Context()) |
| 150 | if loadErr != nil { |
| 151 | t.Fatal(loadErr) |
| 152 | } |
| 153 | observed = state.PendingOperations["purge-victim"] |
| 154 | } |
| 155 | fmt.Println("ready") |
| 156 | if _, err := bufio.NewReader(os.Stdin).ReadString('\n'); err != nil { |
| 157 | t.Fatal(err) |
| 158 | } |
| 159 | switch action { |
| 160 | case "restore": |
| 161 | err = store.RestoreSession(t.Context(), "victim") |
| 162 | case "purge": |
| 163 | err = store.BeginPurge(t.Context(), "victim", expected) |
| 164 | case "resume": |
| 165 | err = store.ResumePurge(t.Context(), "victim", observed) |
| 166 | default: |
| 167 | err = fmt.Errorf("unknown action %q", action) |
| 168 | } |
| 169 | switch { |
| 170 | case err == nil: |
| 171 | fmt.Println("ok") |
| 172 | case errors.Is(err, ErrMutationConflict): |
| 173 | fmt.Println("conflict") |
| 174 | default: |
| 175 | fmt.Printf("error:%v\n", err) |
| 176 | } |
| 177 | } |
| 178 | |
| 179 | func TestWorkspaceOldPurgeProcessCannotReplaceNewOperation(t *testing.T) { |
| 180 | store, expected := seedArchivedProcessState(t) |
| 181 | if err := store.mutate(t.Context(), func(s *State) error { |
| 182 | s.PendingOperations["purge-victim"] = Operation{ID: "purge-victim", Kind: "purge", Phase: "prepared", Lifecycle: Deleted, SessionIDs: []string{"victim"}, ExpectedGeneration: expected} |
| 183 | return nil |
| 184 | }); err != nil { |
| 185 | t.Fatal(err) |
| 186 | } |
| 187 | old := startLifecycleProcess(t, store.Path(), "resume", expected) |
| 188 | if err := store.RestoreSession(t.Context(), "victim"); err != nil { |
| 189 | t.Fatal(err) |
| 190 | } |
| 191 | if err := store.ArchiveSession(t.Context(), "victim"); err != nil { |
| 192 | t.Fatal(err) |
| 193 | } |
| 194 | state, err := store.Load(t.Context()) |
| 195 | if err != nil { |
| 196 | t.Fatal(err) |
| 197 | } |
| 198 | if err := store.BeginPurge(t.Context(), "victim", state.Generation); err != nil { |
| 199 | t.Fatal(err) |
| 200 | } |
| 201 | if got := old.run(t); got != "conflict" { |
| 202 | t.Fatalf("old replay=%s", got) |
| 203 | } |
| 204 | state, err = store.Load(t.Context()) |
| 205 | if err != nil { |
| 206 | t.Fatal(err) |
| 207 | } |
| 208 | if ClassifyPurge(state, "victim") != PurgeTombstoned { |
| 209 | t.Fatal("old process changed new deletion") |
| 210 | } |
| 211 | } |
| 212 |