Merge contiguous stack prefixes as one operation

6c086d4a625f · Devin AI · · parent a548f804ef24

Merge contiguous stack prefixes as one operation

Landing a stacked PR merges it and every open layer below it through
one durable stack operation (#51). POST /repos/:owner/:repo/stacks/
:stack_number/merge validates the request, replays idempotency keys,
and returns the operation immediately; the operation worker executes
it by kind.

Preflight validates everything before any ref moves: the contiguous
lowest-open prefix, the ancestry chain against the live trunk, live
branch heads against the recorded entries, caller-supplied expected
heads, the stack version, actor write access, and no conflicting
active operation.

The merge-commit method creates one group merge commit whose tree is
the selected top head's; squash chains one commit per layer on trunk,
preserving each layer's author; rebase replays each layer's unique
commits in order and verifies the final tree equals the top head's.
Layers above the selected prefix restack onto the new trunk result
from their stored boundaries, and only the lowest remaining layer
retargets the trunk. Merged lower branches stay intact.

All generated commits persist under hidden retention refs, then one
batch compare-and-swap moves the trunk and every restacked branch and
drops the retention refs. The planned result persists before the swap,
so a crashed worker reconciles by comparing planned refs with live
refs: applied plans finalize the metadata, untouched plans re-run, and
diverged refs mark the operation partially_succeeded - pull requests
whose refs landed stay merged even when a later step fails.

Co-Authored-By: Christopher David <chris@openagents.com>
Co-Authored-By
Christopher David <chris@openagents.com>

Deploy story

What this commit did to the running system — joined from the forge receipt chain, the part a commit page elsewhere cannot show.

Not deployed through the forge lane

No push, promotion, build, or deploy receipt references this commit (receipts are scanned over a bounded recent window). Changes shipped by full node replacement carry their proof in the release gate receipt instead.

Changed files

  • modified lib/openagents/forge/git_plane.ex
  • added lib/openagents/stacks/merge.ex
  • modified lib/openagents/stacks/operation_worker.ex
  • modified lib/openagents/stacks/stack_event.ex
  • modified lib/openagents_web/api_route_authority.ex
  • modified lib/openagents_web/controllers/stack_controller.ex
  • modified lib/openagents_web/router.ex
  • added test/openagents/stacks/merge_test.exs
  • modified test/openagents_web/controllers/stack_controller_test.exs

Diff

9 files changed, +1944 -9

lib/openagents/forge/git_plane.ex modified +66

@@ -56,6 +56,25 @@ defmodule OpenAgents.Forge.GitPlane do

56 56
    end
57 57
  end
58 58
59
  @doc "The tree OID of a commit."
60
  def tree_of(repo, rev) do
61
    with :ok <- check_rev(rev),
62
         :ok <- Sync.ensure_fresh(repo) do
63
      case git(repo, ["rev-parse", "--verify", "--quiet", "--end-of-options", rev <> "^{tree}"]) do
64
        {output, 0} -> {:ok, String.trim(output)}
65
        _other -> {:error, :not_found}
66
      end
67
    end
68
  end
69
70
  @doc "The author identity of a commit: name, email, and strict ISO date."
71
  def commit_author(repo, rev) do
72
    with :ok <- check_rev(rev),
73
         :ok <- Sync.ensure_fresh(repo) do
74
      author_of(Repos.bare_path(repo), rev)
75
    end
76
  end
77
59 78
  @doc "The parent commit OIDs of a commit, in order (empty for a root commit)."
60 79
  def parents(repo, rev) do
61 80
    with :ok <- check_rev(rev),

@@ -282,6 +301,53 @@ defmodule OpenAgents.Forge.GitPlane do

282 301
    end
283 302
  end
284 303
304
  @doc """
305
  Create one commit object for `tree` with the given `parents` and `message`,
306
  without moving any ref.
307
308
  The new commit is unreachable until a ref moves to it through
309
  `batch_update_refs/3`. Pass `author: %{name: ..., email: ..., date: ...}`
310
  to preserve an original author; the forge service identity is both the
311
  default author and always the committer, and the commit is unsigned
312
  (`docs/stacked-prs.md` section 12.5).
313
  """
314
  def commit_tree(repo, tree, parents, message, opts \\ [])
315
      when is_list(parents) and parents != [] and is_binary(message) do
316
    with :ok <- check_oid(tree),
317
         :ok <- check_parent_oids(parents),
318
         :ok <- Sync.ensure_fresh(repo) do
319
      path = Repos.bare_path(repo)
320
      author = Keyword.get(opts, :author) || forge_author()
321
322
      env = [
323
        {"GIT_AUTHOR_NAME", author.name},
324
        {"GIT_AUTHOR_EMAIL", author.email},
325
        {"GIT_AUTHOR_DATE", author.date},
326
        {"GIT_COMMITTER_NAME", @committer_name},
327
        {"GIT_COMMITTER_EMAIL", @committer_email}
328
      ]
329
330
      parent_args = Enum.flat_map(parents, &["-p", &1])
331
332
      case git_with_stdin(path, ["commit-tree", tree | parent_args], message, env) do
333
        {output, 0} -> {:ok, String.trim(output)}
334
        _other -> {:error, :commit_tree_failed}
335
      end
336
    end
337
  end
338
339
  defp check_parent_oids(parents) do
340
    if Enum.all?(parents, &valid_oid?/1), do: :ok, else: {:error, :not_found}
341
  end
342
343
  defp forge_author do
344
    %{
345
      name: @committer_name,
346
      email: @committer_email,
347
      date: DateTime.utc_now() |> DateTime.to_iso8601()
348
    }
349
  end
350
285 351
  ## Retention refs
286 352
287 353
  @doc """
lib/openagents/stacks/merge.ex added +1002

@@ -0,0 +1,1002 @@

1
defmodule OpenAgents.Stacks.Merge do
2
  @moduledoc """
3
  Contiguous-prefix stack merge (`docs/stacked-prs.md` sections 13 and
4
  5.9-5.13).
5
6
  Selecting a pull request merges it and every open layer below it in one
7
  durable operation. Preflight validates everything — the contiguous
8
  prefix, the ancestry chain, live and expected heads, the stack version —
9
  before any object builds. The merge result and every upper-layer restack
10
  then build as unreachable candidate commits, persist under hidden
11
  retention refs, and land through one batch compare-and-swap that moves
12
  the trunk and every restacked branch together. Merged lower branches
13
  stay untouched, so their history survives until callers delete them.
14
15
  Git refs and SQL metadata are two durable stores, so the operation
16
  records its planned result before the refs move and marks the plan
17
  applied immediately after. A worker that reclaims a crashed operation
18
  reconciles instead of re-merging: when the live refs already equal the
19
  plan it only finalizes the metadata, when nothing moved it re-plans, and
20
  when the refs diverged it marks the operation `partially_succeeded` —
21
  pull requests whose refs landed stay merged even when a later step
22
  fails.
23
24
  An upper-layer restack conflict fails the operation during planning,
25
  before any ref moves, with the conflict persisted on the operation;
26
  resolve it through a stack rebase and retry the merge.
27
  """
28
  import Ecto.Query, warn: false
29
30
  alias OpenAgents.Accounts.User
31
  alias OpenAgents.Forge.GitPlane
32
  alias OpenAgents.Issues.Issue
33
  alias OpenAgents.PullRequests.PullRequest
34
  alias OpenAgents.Repo
35
  alias OpenAgents.Repositories
36
  alias OpenAgents.Repositories.Repository
37
  alias OpenAgents.Stacks.Operation
38
  alias OpenAgents.Stacks.Stack
39
  alias OpenAgents.Stacks.StackEntry
40
  alias OpenAgents.Stacks.StackEvent
41
42
  @merge_methods ~w(merge squash rebase)
43
44
  @doc """
45
  Requests a merge of the contiguous lowest prefix ending at one pull
46
  request.
47
48
  The request names the selected pull request, the merge method, optional
49
  expected heads per layer, and an optional expected stack version. It
50
  inserts one durable `Operation` row in state `pending`; a worker claims
51
  and executes it. A retried idempotency key replays the original
52
  operation.
53
  """
54
  def request_from_api(%Repository{} = repository, number, params, %User{} = actor, key)
55
      when is_integer(number) and is_binary(key) do
56
    with :ok <- authorize(repository, actor),
57
         {:ok, request} <- parse_merge_request(params) do
58
      Repo.transaction(fn ->
59
        lock_repository_stacks(repository.id)
60
61
        with {:ok, stack} <- get_stack_for_update(repository, number),
62
             :ok <- validate_open(stack),
63
             {:ok, replay} <- check_idempotency(stack, key, request),
64
             :ok <- ensure_no_active_operation(stack, replay),
65
             :ok <- validate_expected_version(request["expected_stack_version"], stack),
66
             {:ok, position} <- selected_position(stack, request, replay) do
67
          case replay do
68
            %Operation{} = operation ->
69
              {operation, :replayed}
70
71
            nil ->
72
              operation = insert_operation!(stack, actor, key, request, position)
73
              set_health!(stack, "operation_in_progress")
74
              {operation, :created}
75
          end
76
        else
77
          {:error, reason} -> Repo.rollback(reason)
78
        end
79
      end)
80
    end
81
  end
82
83
  @doc """
84
  Executes one claimed merge operation to a terminal state.
85
86
  The caller (the operation worker) has already marked the row `running`.
87
  A re-execution after a crash reconciles from the persisted plan instead
88
  of repeating work.
89
  """
90
  def execute(%Operation{kind: "merge"} = operation, opts \\ []) do
91
    stack = load_stack(operation.stack_id)
92
    repository = Repo.one!(from r in Repository, where: r.id == ^stack.repository_id)
93
94
    case reconcile_state(repository, operation) do
95
      :fresh ->
96
        run(operation, repository, stack, opts)
97
98
      {:applied, plan} ->
99
        finalize(operation, stack, plan)
100
101
      {:diverged, plan} ->
102
        partially_succeed(operation, plan, :refs_diverged)
103
    end
104
  end
105
106
  defp run(operation, repository, stack, opts) do
107
    with :ok <- validate_executable(operation, stack),
108
         {:ok, trunk_tip} <- resolve_trunk(repository, stack),
109
         {:ok, selected, upper} <- split_entries(stack, operation.target_position),
110
         :ok <- preflight(repository, trunk_tip, selected, upper, operation.request),
111
         operation = record_snapshot!(operation, stack, trunk_tip, selected),
112
         {:ok, plan} <- build_plan(operation, repository, stack, trunk_tip, selected, upper),
113
         operation = record_plan!(operation, plan),
114
         :ok <- apply_refs(operation, repository, plan),
115
         operation = mark_refs_applied!(operation),
116
         :ok <- crash_seam(opts) do
117
      finalize(operation, stack, plan)
118
    else
119
      {:error, reason} -> fail(operation, stack, reason)
120
    end
121
  end
122
123
  # A test-only fault-injection seam: production callers never pass it, so
124
  # the merge proceeds straight to finalization.
125
  defp crash_seam(opts) do
126
    case Keyword.get(opts, :after_refs) do
127
      nil -> :ok
128
      fun when is_function(fun, 0) -> fun.()
129
    end
130
  end
131
132
  ## Reconciliation
133
134
  # The planned result is the recovery record: `refs_applied` flips true
135
  # right after the batch-CAS lands, so a reclaimed operation knows whether
136
  # the git side already moved.
137
  defp reconcile_state(_repository, %Operation{planned_result: nil}), do: :fresh
138
139
  defp reconcile_state(repository, %Operation{planned_result: plan}) do
140
    refs = planned_refs(plan)
141
142
    cond do
143
      Enum.all?(refs, fn ref -> live_oid(repository, ref["ref"]) == ref["new"] end) ->
144
        {:applied, plan}
145
146
      not plan_marked_applied?(plan) and
147
          Enum.all?(refs, fn ref -> live_oid(repository, ref["ref"]) == ref["old"] end) ->
148
        :fresh
149
150
      true ->
151
        {:diverged, plan}
152
    end
153
  end
154
155
  defp plan_marked_applied?(plan), do: Map.get(plan, "refs_applied") == true
156
157
  defp planned_refs(plan) do
158
    trunk = %{
159
      "ref" => Map.fetch!(plan, "trunk_ref"),
160
      "old" => Map.fetch!(plan, "trunk_old"),
161
      "new" => Map.fetch!(plan, "trunk_new")
162
    }
163
164
    restacked =
165
      plan
166
      |> Map.fetch!("restacked")
167
      |> Enum.map(fn step ->
168
        %{
169
          "ref" => Map.fetch!(step, "ref"),
170
          "old" => Map.fetch!(step, "old_head"),
171
          "new" => Map.fetch!(step, "new_head")
172
        }
173
      end)
174
175
    [trunk | restacked]
176
  end
177
178
  defp live_oid(repository, ref) do
179
    case GitPlane.resolve_commit(repository.storage_key, ref) do
180
      {:ok, oid} -> oid
181
      {:error, _reason} -> :absent
182
    end
183
  end
184
185
  ## Preflight
186
187
  defp validate_executable(operation, stack) do
188
    cond do
189
      stack.state != "open" -> {:error, :stack_not_open}
190
      stack.version != operation.expected_stack_version -> {:error, :stale_stack_version}
191
      stack.entries == [] -> {:error, :empty_stack}
192
      true -> :ok
193
    end
194
  end
195
196
  defp resolve_trunk(repository, stack) do
197
    case GitPlane.resolve_commit(repository.storage_key, "refs/heads/" <> stack.trunk_ref) do
198
      {:ok, oid} -> {:ok, oid}
199
      {:error, _reason} -> {:error, {:missing_ref, "refs/heads/" <> stack.trunk_ref}}
200
    end
201
  end
202
203
  defp split_entries(stack, target_position) do
204
    {selected, upper} = Enum.split_with(stack.entries, &(&1.position <= target_position))
205
206
    if selected != [] and Enum.any?(selected, &(&1.position == target_position)) do
207
      {:ok, selected, upper}
208
    else
209
      {:error, :pull_request_not_in_stack}
210
    end
211
  end
212
213
  # Everything validates before any object builds: the selected prefix must
214
  # be open, current on the trunk, an unbroken parent chain, and live on
215
  # the branches the stack recorded — and any caller-supplied expected
216
  # heads must match.
217
  defp preflight(repository, trunk_tip, selected, upper, request) do
218
    entries = selected ++ upper
219
220
    with :ok <- validate_selected_open(selected),
221
         :ok <- validate_chain(trunk_tip, entries),
222
         :ok <- verify_live_heads(repository, entries) do
223
      validate_expected_heads(request["expected_heads"], entries)
224
    end
225
  end
226
227
  defp validate_selected_open(selected) do
228
    Enum.find_value(selected, :ok, fn entry ->
229
      if entry.pull_request.state == "open" do
230
        nil
231
      else
232
        {:error, :pull_request_not_open}
233
      end
234
    end)
235
  end
236
237
  defp validate_chain(trunk_tip, entries) do
238
    entries
239
    |> Enum.reduce_while(trunk_tip, fn entry, expected_boundary ->
240
      if entry.boundary_oid == expected_boundary do
241
        {:cont, entry.observed_head_oid}
242
      else
243
        {:halt, {:error, :needs_rebase}}
244
      end
245
    end)
246
    |> case do
247
      {:error, reason} -> {:error, reason}
248
      _head -> :ok
249
    end
250
  end
251
252
  defp verify_live_heads(repository, entries) do
253
    Enum.find_value(entries, :ok, fn entry ->
254
      ref = "refs/heads/" <> entry.pull_request.head_ref
255
256
      case GitPlane.resolve_commit(repository.storage_key, ref) do
257
        {:ok, oid} when oid == entry.observed_head_oid -> nil
258
        {:ok, actual} -> {:error, {:head_changed, ref, actual}}
259
        {:error, _reason} -> {:error, {:missing_ref, ref}}
260
      end
261
    end)
262
  end
263
264
  defp validate_expected_heads(nil, _entries), do: :ok
265
266
  defp validate_expected_heads(expected, entries) when is_map(expected) do
267
    by_number = Map.new(entries, &{&1.pull_request.issue.number, &1})
268
269
    Enum.find_value(expected, :ok, fn {number_key, oid} ->
270
      with {:ok, number} <- parse_number_key(number_key),
271
           %StackEntry{} = entry <- Map.get(by_number, number, :missing) do
272
        if entry.observed_head_oid == oid, do: nil, else: {:error, :expected_head_mismatch}
273
      else
274
        _invalid -> {:error, :expected_head_mismatch}
275
      end
276
    end)
277
  end
278
279
  defp parse_number_key(key) when is_integer(key), do: {:ok, key}
280
281
  defp parse_number_key(key) when is_binary(key) do
282
    case Integer.parse(key) do
283
      {number, ""} -> {:ok, number}
284
      _other -> :error
285
    end
286
  end
287
288
  ## Snapshot and plan
289
290
  defp record_snapshot!(%Operation{snapshot: nil} = operation, stack, trunk_tip, selected) do
291
    selected_positions = Enum.map(selected, & &1.position)
292
293
    snapshot = %{
294
      "trunk_oid" => trunk_tip,
295
      "stack_version" => stack.version,
296
      "selected_positions" => selected_positions,
297
      "entries" =>
298
        Enum.map(stack.entries, fn entry ->
299
          %{
300
            "position" => entry.position,
301
            "ref" => "refs/heads/" <> entry.pull_request.head_ref,
302
            "boundary_oid" => entry.boundary_oid,
303
            "observed_head_oid" => entry.observed_head_oid
304
          }
305
        end)
306
    }
307
308
    operation
309
    |> Operation.transition_changeset(%{state: operation.state, snapshot: snapshot})
310
    |> Repo.update!()
311
  end
312
313
  defp record_snapshot!(%Operation{} = operation, _stack, _trunk_tip, _selected), do: operation
314
315
  defp build_plan(operation, repository, stack, trunk_tip, selected, upper) do
316
    method = Map.fetch!(operation.request, "merge_method")
317
318
    with {:ok, trunk_new, merged} <-
319
           build_merge_result(repository, stack, trunk_tip, selected, method),
320
         {:ok, restacked} <- plan_upper_restacks(repository, upper, trunk_new) do
321
      {:ok,
322
       %{
323
         "merge_method" => method,
324
         "trunk_ref" => "refs/heads/" <> stack.trunk_ref,
325
         "trunk_old" => trunk_tip,
326
         "trunk_new" => trunk_new,
327
         "merged" => merged,
328
         "restacked" => restacked,
329
         "refs_applied" => false
330
       }}
331
    end
332
  end
333
334
  # One group merge commit: both selected history and the current trunk are
335
  # parents, and the tree is exactly the selected top head's tree.
336
  defp build_merge_result(repository, stack, trunk_tip, selected, "merge") do
337
    top = List.last(selected)
338
    numbers = Enum.map(selected, & &1.pull_request.issue.number)
339
    message = merge_message(stack, numbers)
340
341
    with {:ok, tree} <- tree_of(repository, top.observed_head_oid),
342
         {:ok, merge_commit} <-
343
           commit_tree(repository, tree, [trunk_tip, top.observed_head_oid], message) do
344
      {:ok, merge_commit, Enum.map(selected, &merged_step(&1, merge_commit))}
345
    end
346
  end
347
348
  # One squash commit per pull request, chained on the trunk, so each trunk
349
  # commit diffs to exactly one layer.
350
  defp build_merge_result(repository, _stack, trunk_tip, selected, "squash") do
351
    selected
352
    |> Enum.reduce_while({:ok, trunk_tip, []}, fn entry, {:ok, parent, merged} ->
353
      issue = entry.pull_request.issue
354
      message = "#{issue.title} (##{issue.number})"
355
356
      with {:ok, tree} <- tree_of(repository, entry.observed_head_oid),
357
           {:ok, author} <- commit_author(repository, entry.observed_head_oid),
358
           {:ok, squash} <- commit_tree(repository, tree, [parent], message, author: author) do
359
        {:cont, {:ok, squash, [merged_step(entry, squash) | merged]}}
360
      else
361
        {:error, reason} -> {:halt, {:error, reason}}
362
      end
363
    end)
364
    |> case do
365
      {:ok, trunk_new, merged} -> {:ok, trunk_new, Enum.reverse(merged)}
366
      {:error, reason} -> {:error, reason}
367
    end
368
  end
369
370
  # Replay every selected layer's unique commits in order; the final tree
371
  # must equal the selected top head's tree.
372
  defp build_merge_result(repository, _stack, trunk_tip, selected, "rebase") do
373
    selected
374
    |> Enum.reduce_while({:ok, trunk_tip, []}, fn entry, {:ok, parent, merged} ->
375
      case GitPlane.replay(
376
             repository.storage_key,
377
             entry.boundary_oid,
378
             entry.observed_head_oid,
379
             parent
380
           ) do
381
        {:ok, %{new_head: new_head}} ->
382
          {:cont, {:ok, new_head, [merged_step(entry, new_head) | merged]}}
383
384
        {:conflict, _conflict} ->
385
          {:halt, {:error, {:replay_failed, entry.position, :conflict}}}
386
387
        {:error, reason} ->
388
          {:halt, {:error, {:replay_failed, entry.position, reason}}}
389
      end
390
    end)
391
    |> case do
392
      {:ok, trunk_new, merged} ->
393
        top = List.last(selected)
394
395
        with {:ok, result_tree} <- tree_of(repository, trunk_new),
396
             {:ok, top_tree} <- tree_of(repository, top.observed_head_oid) do
397
          if result_tree == top_tree,
398
            do: {:ok, trunk_new, Enum.reverse(merged)},
399
            else: {:error, :rebase_tree_mismatch}
400
        end
401
402
      {:error, reason} ->
403
        {:error, reason}
404
    end
405
  end
406
407
  defp merged_step(entry, merge_commit_sha) do
408
    %{
409
      "position" => entry.position,
410
      "entry_id" => entry.id,
411
      "pull_request_id" => entry.pull_request_id,
412
      "pull_request_number" => entry.pull_request.issue.number,
413
      "ref" => "refs/heads/" <> entry.pull_request.head_ref,
414
      "old_head" => entry.observed_head_oid,
415
      "merge_commit_sha" => merge_commit_sha
416
    }
417
  end
418
419
  defp merge_message(stack, numbers) do
420
    "Merge pull requests #{Enum.map_join(numbers, ", ", &"##{&1}")} (stack #{stack.number})"
421
  end
422
423
  # Layers above the selected prefix restack onto the merge result using
424
  # their stored boundaries. A conflict here fails the whole operation
425
  # before any ref moves; the persisted error names the conflicting layer.
426
  defp plan_upper_restacks(repository, upper, trunk_new) do
427
    upper
428
    |> Enum.reduce_while({:ok, trunk_new, []}, fn entry, {:ok, parent, steps} ->
429
      case GitPlane.replay(
430
             repository.storage_key,
431
             entry.boundary_oid,
432
             entry.observed_head_oid,
433
             parent
434
           ) do
435
        {:ok, %{new_head: new_head}} ->
436
          step = %{
437
            "position" => entry.position,
438
            "entry_id" => entry.id,
439
            "pull_request_id" => entry.pull_request_id,
440
            "pull_request_number" => entry.pull_request.issue.number,
441
            "ref" => "refs/heads/" <> entry.pull_request.head_ref,
442
            "old_head" => entry.observed_head_oid,
443
            "old_boundary" => entry.boundary_oid,
444
            "new_boundary" => parent,
445
            "new_head" => new_head
446
          }
447
448
          {:cont, {:ok, new_head, [step | steps]}}
449
450
        {:conflict, conflict} ->
451
          {:halt,
452
           {:error,
453
            {:upper_restack_conflict,
454
             %{
455
               "position" => entry.position,
456
               "pull_request_number" => entry.pull_request.issue.number,
457
               "onto" => conflict.onto,
458
               "commit" => conflict.commit,
459
               "paths" => conflict.paths,
460
               "messages" => conflict.messages
461
             }}}}
462
463
        {:error, reason} ->
464
          {:halt, {:error, {:replay_failed, entry.position, reason}}}
465
      end
466
    end)
467
    |> case do
468
      {:ok, _parent, steps} -> {:ok, Enum.reverse(steps)}
469
      {:error, reason} -> {:error, reason}
470
    end
471
  end
472
473
  defp record_plan!(operation, plan) do
474
    operation
475
    |> Operation.transition_changeset(%{state: operation.state, planned_result: plan})
476
    |> Repo.update!()
477
  end
478
479
  defp mark_refs_applied!(operation) do
480
    plan = Map.put(operation.planned_result, "refs_applied", true)
481
482
    operation
483
    |> Operation.transition_changeset(%{state: operation.state, planned_result: plan})
484
    |> Repo.update!()
485
  end
486
487
  ## Ref application
488
489
  defp apply_refs(operation, repository, plan) do
490
    with {:ok, temp_refs} <- retain_new_commits(operation, repository, plan) do
491
      move_public_refs(repository, plan, temp_refs)
492
    end
493
  end
494
495
  defp retain_new_commits(operation, repository, plan) do
496
    targets =
497
      [{"trunk", Map.fetch!(plan, "trunk_new")}] ++
498
        Enum.map(Map.fetch!(plan, "restacked"), fn step ->
499
          {"p#{Map.fetch!(step, "position")}", Map.fetch!(step, "new_head")}
500
        end)
501
502
    temp_refs =
503
      Enum.map(targets, fn {segment, oid} ->
504
        {:ok, ref} =
505
          GitPlane.internal_ref([
506
            "operations",
507
            operation.id,
508
            "a#{operation.attempt_count}",
509
            segment
510
          ])
511
512
        %{ref: ref, expected_old: :absent, new: oid}
513
      end)
514
515
    case GitPlane.batch_update_refs(repository.storage_key, temp_refs, principal(operation)) do
516
      {:ok, _result} -> {:ok, temp_refs}
517
      {:error, reason} -> {:error, {:retention_failed, reason}}
518
    end
519
  end
520
521
  # One atomic batch: the trunk and every restacked branch move together
522
  # and the retention refs delete in the same transaction. Merged lower
523
  # branches do not move and are not deleted.
524
  defp move_public_refs(repository, plan, temp_refs) do
525
    updates =
526
      [
527
        %{
528
          ref: Map.fetch!(plan, "trunk_ref"),
529
          expected_old: Map.fetch!(plan, "trunk_old"),
530
          new: Map.fetch!(plan, "trunk_new")
531
        }
532
      ] ++
533
        Enum.map(Map.fetch!(plan, "restacked"), fn step ->
534
          %{
535
            ref: Map.fetch!(step, "ref"),
536
            expected_old: Map.fetch!(step, "old_head"),
537
            new: Map.fetch!(step, "new_head")
538
          }
539
        end) ++
540
        Enum.map(temp_refs, fn temp ->
541
          %{ref: temp.ref, expected_old: temp.new, new: :delete}
542
        end)
543
544
    case GitPlane.batch_update_refs(repository.storage_key, updates, "stack-merge") do
545
      {:ok, _result} ->
546
        :ok
547
548
      {:error, {:expected_mismatch, ref, actual}} ->
549
        cleanup_temp_refs(repository, temp_refs)
550
        {:error, {:head_changed, ref, actual}}
551
552
      {:error, reason} ->
553
        cleanup_temp_refs(repository, temp_refs)
554
        {:error, {:ref_update_failed, reason}}
555
    end
556
  end
557
558
  defp cleanup_temp_refs(repository, temp_refs) do
559
    deletes = Enum.map(temp_refs, &%{ref: &1.ref, expected_old: &1.new, new: :delete})
560
    _result = GitPlane.batch_update_refs(repository.storage_key, deletes, "stack-merge")
561
    :ok
562
  end
563
564
  defp principal(operation), do: "stack-operation-" <> operation.id
565
566
  ## Finalization
567
568
  defp finalize(operation, stack, plan) do
569
    result =
570
      Repo.transaction(fn ->
571
        lock_repository_stacks(stack.repository_id)
572
        current = Repo.one!(from s in Stack, where: s.id == ^stack.id, lock: "FOR UPDATE")
573
574
        if current.version != operation.expected_stack_version do
575
          Repo.rollback({:version_moved, current.version})
576
        end
577
578
        now = DateTime.utc_now()
579
        Enum.each(Map.fetch!(plan, "merged"), &finalize_merged!(&1, operation, now))
580
        Enum.each(Map.fetch!(plan, "restacked"), &finalize_restacked!(&1, current, plan))
581
582
        stack_after = complete_or_bump!(current, plan)
583
        record_events!(stack_after, operation, plan)
584
585
        operation
586
        |> Operation.transition_changeset(%{
587
          state: "succeeded",
588
          planned_result: Map.put(plan, "refs_applied", true),
589
          completed_at: now
590
        })
591
        |> Repo.update!()
592
      end)
593
594
    case result do
595
      {:ok, operation} -> {:ok, operation}
596
      {:error, reason} -> partially_succeed(operation, plan, reason)
597
    end
598
  end
599
600
  defp finalize_merged!(step, operation, now) do
601
    entry_id = Map.fetch!(step, "entry_id")
602
    pull_request_id = Map.fetch!(step, "pull_request_id")
603
604
    {1, _rows} =
605
      Repo.update_all(
606
        from(entry in StackEntry, where: entry.id == ^entry_id and is_nil(entry.removed_at)),
607
        set: [removed_at: now, updated_at: now]
608
      )
609
610
    pull_request =
611
      Repo.one!(
612
        from pr in PullRequest,
613
          where: pr.id == ^pull_request_id,
614
          lock: "FOR UPDATE"
615
      )
616
617
    {1, _rows} =
618
      Repo.update_all(
619
        from(pr in PullRequest, where: pr.id == ^pull_request_id),
620
        set: [
621
          state: "closed",
622
          merged_at: now,
623
          merged_by_user_id: operation.created_by_user_id,
624
          merge_commit_sha: Map.fetch!(step, "merge_commit_sha"),
625
          updated_at: now
626
        ]
627
      )
628
629
    {1, _rows} =
630
      Repo.update_all(
631
        from(issue in Issue, where: issue.id == ^pull_request.issue_id),
632
        set: [
633
          state: "closed",
634
          state_reason: "completed",
635
          closed_at: DateTime.truncate(now, :second),
636
          updated_at: DateTime.truncate(now, :second)
637
        ]
638
      )
639
640
    :ok
641
  end
642
643
  # The lowest remaining layer retargets the trunk; every higher layer
644
  # keeps its direct base and only its boundary and head advance.
645
  defp finalize_restacked!(step, stack, plan) do
646
    entry_id = Map.fetch!(step, "entry_id")
647
    new_boundary = Map.fetch!(step, "new_boundary")
648
    new_head = Map.fetch!(step, "new_head")
649
    lowest? = Map.fetch!(step, "new_boundary") == Map.fetch!(plan, "trunk_new")
650
651
    Repo.one!(from entry in StackEntry, where: entry.id == ^entry_id, lock: "FOR UPDATE")
652
    |> StackEntry.changeset(%{boundary_oid: new_boundary, observed_head_oid: new_head})
653
    |> Repo.update!()
654
655
    base_ref_set = if lowest?, do: [base_ref: stack.trunk_ref], else: []
656
657
    {1, _rows} =
658
      Repo.update_all(
659
        from(pr in PullRequest, where: pr.id == ^Map.fetch!(step, "pull_request_id")),
660
        set: [head_sha: new_head, base_sha: new_boundary] ++ base_ref_set
661
      )
662
663
    :ok
664
  end
665
666
  defp complete_or_bump!(stack, plan) do
667
    remaining = Map.fetch!(plan, "restacked") != []
668
    state = if remaining, do: "open", else: "completed"
669
670
    {1, [stack_after]} =
671
      Repo.update_all(
672
        from(s in Stack,
673
          where: s.id == ^stack.id and s.version == ^stack.version,
674
          select: s
675
        ),
676
        set: [
677
          version: stack.version + 1,
678
          health: "healthy",
679
          state: state,
680
          updated_at: DateTime.utc_now()
681
        ]
682
      )
683
684
    stack_after
685
  end
686
687
  defp record_events!(stack, operation, plan) do
688
    record_event!(stack, operation, "pull_request_stack.merged", %{
689
      "operation_id" => operation.id,
690
      "merge_method" => Map.fetch!(plan, "merge_method"),
691
      "trunk_old" => Map.fetch!(plan, "trunk_old"),
692
      "trunk_new" => Map.fetch!(plan, "trunk_new"),
693
      "merged" => Map.fetch!(plan, "merged")
694
    })
695
696
    Enum.each(Map.fetch!(plan, "restacked"), fn step ->
697
      record_event!(stack, operation, "pull_request.synchronize", %{
698
        "operation_id" => operation.id,
699
        "pull_request" => Map.fetch!(step, "pull_request_number"),
700
        "ref" => Map.fetch!(step, "ref"),
701
        "before" => Map.fetch!(step, "old_head"),
702
        "after" => Map.fetch!(step, "new_head")
703
      })
704
    end)
705
  end
706
707
  defp record_event!(stack, operation, event_type, payload) do
708
    %StackEvent{}
709
    |> StackEvent.changeset(%{
710
      stack_id: stack.id,
711
      actor_user_id: operation.created_by_user_id,
712
      event_type: event_type,
713
      stack_version: stack.version,
714
      payload: payload
715
    })
716
    |> Repo.insert!()
717
  end
718
719
  ## Terminal transitions
720
721
  # The refs landed but the metadata could not finalize: the merged pull
722
  # requests' code is on the trunk, so the operation records exactly what
723
  # applied instead of pretending nothing happened.
724
  defp partially_succeed(operation, plan, reason) do
725
    operation =
726
      operation
727
      |> Operation.transition_changeset(%{
728
        state: "partially_succeeded",
729
        error: error_map(reason),
730
        planned_result: Map.put(plan, "refs_applied", true),
731
        completed_at: DateTime.utc_now()
732
      })
733
      |> Repo.update!()
734
735
    {:error, operation}
736
  end
737
738
  defp fail(operation, stack, reason) do
739
    {:ok, operation} =
740
      Repo.transaction(fn ->
741
        current = Repo.one!(from s in Stack, where: s.id == ^stack.id, lock: "FOR UPDATE")
742
        set_health!(current, failure_health(reason))
743
744
        operation
745
        |> Operation.transition_changeset(%{
746
          state: "failed",
747
          error: error_map(reason),
748
          completed_at: DateTime.utc_now()
749
        })
750
        |> Repo.update!()
751
      end)
752
753
    {:error, operation}
754
  end
755
756
  defp failure_health({:head_changed, _ref, _actual}), do: "head_changed"
757
  defp failure_health({:missing_ref, _ref}), do: "missing_ref"
758
  defp failure_health({:upper_restack_conflict, _conflict}), do: "conflicted"
759
  defp failure_health(:needs_rebase), do: "needs_rebase"
760
  defp failure_health(_reason), do: "needs_rebase"
761
762
  defp error_map({:head_changed, ref, actual}),
763
    do: %{"code" => "head_changed", "ref" => ref, "actual" => stringify_actual(actual)}
764
765
  defp error_map({:missing_ref, ref}), do: %{"code" => "missing_ref", "ref" => ref}
766
767
  defp error_map({:upper_restack_conflict, conflict}),
768
    do: Map.put(conflict, "code", "upper_restack_conflict")
769
770
  defp error_map({:replay_failed, position, reason}),
771
    do: %{"code" => "replay_failed", "position" => position, "reason" => inspect(reason)}
772
773
  defp error_map({:retention_failed, reason}),
774
    do: %{"code" => "retention_failed", "reason" => inspect(reason)}
775
776
  defp error_map({:ref_update_failed, reason}),
777
    do: %{"code" => "ref_update_failed", "reason" => inspect(reason)}
778
779
  defp error_map({:version_moved, version}),
780
    do: %{"code" => "version_moved", "stack_version" => version}
781
782
  defp error_map(reason) when is_atom(reason), do: %{"code" => Atom.to_string(reason)}
783
  defp error_map(reason), do: %{"code" => "operation_failed", "reason" => inspect(reason)}
784
785
  defp stringify_actual(:absent), do: "absent"
786
  defp stringify_actual(oid), do: oid
787
788
  ## Request validation
789
790
  defp authorize(repository, actor) do
791
    if Repositories.writable?(repository, actor), do: :ok, else: {:error, :forbidden}
792
  end
793
794
  defp parse_merge_request(params) do
795
    with {:ok, number} <- parse_pull_request_number(params),
796
         {:ok, method} <- parse_merge_method(params),
797
         :ok <- parse_merge_action(params),
798
         {:ok, version} <- parse_expected_version(params),
799
         {:ok, expected_heads} <- parse_expected_heads(params) do
800
      {:ok,
801
       %{
802
         "pull_request_number" => number,
803
         "merge_method" => method,
804
         "merge_action" => "direct_merge",
805
         "expected_stack_version" => version,
806
         "expected_heads" => expected_heads
807
       }}
808
    end
809
  end
810
811
  defp parse_pull_request_number(params) do
812
    case Map.get(params, "pull_request_number") do
813
      number when is_integer(number) and number >= 1 -> {:ok, number}
814
      _other -> {:error, :invalid_request}
815
    end
816
  end
817
818
  defp parse_merge_method(params) do
819
    case Map.get(params, "merge_method") do
820
      method when method in @merge_methods -> {:ok, method}
821
      _other -> {:error, :invalid_request}
822
    end
823
  end
824
825
  # The merge queue lands as its own slice; until then only a direct merge
826
  # is available.
827
  defp parse_merge_action(params) do
828
    case Map.get(params, "merge_action", "direct_merge") do
829
      "direct_merge" -> :ok
830
      "queue" -> {:error, :merge_queue_unavailable}
831
      _other -> {:error, :invalid_request}
832
    end
833
  end
834
835
  defp parse_expected_version(params) do
836
    case Map.get(params, "expected_stack_version") do
837
      nil -> {:ok, nil}
838
      version when is_integer(version) and version >= 1 -> {:ok, version}
839
      _other -> {:error, :invalid_request}
840
    end
841
  end
842
843
  defp parse_expected_heads(params) do
844
    case Map.get(params, "expected_heads") do
845
      nil ->
846
        {:ok, nil}
847
848
      heads when is_map(heads) ->
849
        if Enum.all?(heads, fn {_key, oid} -> valid_oid_string?(oid) end),
850
          do: {:ok, heads},
851
          else: {:error, :invalid_request}
852
853
      _other ->
854
        {:error, :invalid_request}
855
    end
856
  end
857
858
  defp valid_oid_string?(oid) when is_binary(oid) and byte_size(oid) in [40, 64] do
859
    match?({:ok, _raw}, Base.decode16(oid, case: :lower))
860
  end
861
862
  defp valid_oid_string?(_oid), do: false
863
864
  defp selected_position(stack, request, replay) do
865
    number =
866
      case replay do
867
        %Operation{request: replayed} -> Map.fetch!(replayed, "pull_request_number")
868
        nil -> Map.fetch!(request, "pull_request_number")
869
      end
870
871
    entries = active_entries(stack)
872
873
    case Enum.find(entries, &(&1.pull_request.issue.number == number)) do
874
      %StackEntry{position: position} -> {:ok, position}
875
      nil -> {:error, :pull_request_not_in_stack}
876
    end
877
  end
878
879
  defp check_idempotency(stack, key, request) do
880
    operation =
881
      Repo.one(
882
        from operation in Operation,
883
          where: operation.stack_id == ^stack.id and operation.idempotency_key == ^key,
884
          lock: "FOR UPDATE"
885
      )
886
887
    case operation do
888
      nil ->
889
        {:ok, nil}
890
891
      %Operation{} = operation ->
892
        if Map.drop(operation.request, ["previous_health"]) == request,
893
          do: {:ok, operation},
894
          else: {:error, :idempotency_conflict}
895
    end
896
  end
897
898
  defp ensure_no_active_operation(_stack, %Operation{}), do: :ok
899
900
  defp ensure_no_active_operation(stack, nil) do
901
    active =
902
      Repo.exists?(
903
        from operation in Operation,
904
          where:
905
            operation.stack_id == ^stack.id and
906
              operation.state in ^Operation.active_states()
907
      )
908
909
    if active, do: {:error, :operation_in_progress}, else: :ok
910
  end
911
912
  defp validate_expected_version(nil, _stack), do: :ok
913
  defp validate_expected_version(version, %Stack{version: version}), do: :ok
914
  defp validate_expected_version(_version, %Stack{}), do: {:error, :stale_stack_version}
915
916
  defp insert_operation!(stack, actor, key, request, position) do
917
    %Operation{}
918
    |> Operation.changeset(%{
919
      stack_id: stack.id,
920
      created_by_user_id: actor.id,
921
      kind: "merge",
922
      state: "pending",
923
      target_position: position,
924
      expected_stack_version: request["expected_stack_version"] || stack.version,
925
      idempotency_key: key,
926
      request: Map.put(request, "previous_health", stack.health),
927
      retry_at: DateTime.utc_now()
928
    })
929
    |> Repo.insert!()
930
  end
931
932
  ## Shared helpers
933
934
  defp tree_of(repository, oid) do
935
    case GitPlane.tree_of(repository.storage_key, oid) do
936
      {:ok, tree} -> {:ok, tree}
937
      {:error, reason} -> {:error, {:tree_read_failed, reason}}
938
    end
939
  end
940
941
  defp commit_author(repository, oid) do
942
    case GitPlane.commit_author(repository.storage_key, oid) do
943
      {:ok, author} -> {:ok, author}
944
      {:error, reason} -> {:error, {:author_read_failed, reason}}
945
    end
946
  end
947
948
  defp commit_tree(repository, tree, parents, message, opts \\ []) do
949
    case GitPlane.commit_tree(repository.storage_key, tree, parents, message, opts) do
950
      {:ok, oid} -> {:ok, oid}
951
      {:error, reason} -> {:error, {:commit_build_failed, reason}}
952
    end
953
  end
954
955
  defp load_stack(stack_id) do
956
    Stack
957
    |> Repo.get!(stack_id)
958
    |> Repo.preload(
959
      entries:
960
        from(entry in StackEntry,
961
          where: is_nil(entry.removed_at),
962
          order_by: [asc: entry.position],
963
          preload: [pull_request: :issue]
964
        )
965
    )
966
  end
967
968
  defp active_entries(%Stack{id: stack_id}) do
969
    Repo.all(
970
      from entry in StackEntry,
971
        where: entry.stack_id == ^stack_id and is_nil(entry.removed_at),
972
        order_by: [asc: entry.position],
973
        preload: [pull_request: :issue]
974
    )
975
  end
976
977
  defp get_stack_for_update(%Repository{id: repository_id}, number) do
978
    case Repo.one(
979
           from stack in Stack,
980
             where: stack.repository_id == ^repository_id and stack.number == ^number,
981
             lock: "FOR UPDATE"
982
         ) do
983
      nil -> {:error, :stack_not_found}
984
      stack -> {:ok, stack}
985
    end
986
  end
987
988
  defp validate_open(%Stack{state: "open"}), do: :ok
989
  defp validate_open(%Stack{}), do: {:error, :stack_not_open}
990
991
  defp lock_repository_stacks(repository_id) do
992
    key = "pull_request_stacks:#{repository_id}"
993
    Repo.query!("SELECT pg_advisory_xact_lock(hashtextextended($1, 0))", [key])
994
    :ok
995
  end
996
997
  defp set_health!(%Stack{} = stack, health) do
998
    stack
999
    |> Stack.changeset(%{health: health})
1000
    |> Repo.update!()
1001
  end
1002
end
lib/openagents/stacks/operation_worker.ex modified +12 -6

@@ -3,10 +3,11 @@ defmodule OpenAgents.Stacks.OperationWorker do

3 3
  Claims and executes durable stack operations.
4 4
5 5
  The worker polls `stack_operations` for pending rows (and running rows
6
  whose lease expired, which recovers a crashed worker) and executes them
7
  through `OpenAgents.Stacks.Restack`. Claiming uses `FOR UPDATE SKIP
8
  LOCKED`, so multiple nodes never execute the same operation twice inside
9
  one lease window.
6
  whose lease expired, which recovers a crashed worker) and dispatches
7
  each row by kind: rebases run through `OpenAgents.Stacks.Restack` and
8
  merges through `OpenAgents.Stacks.Merge`. Claiming uses `FOR UPDATE
9
  SKIP LOCKED`, so multiple nodes never execute the same operation twice
10
  inside one lease window.
10 11
  """
11 12
  use GenServer
12 13

@@ -16,6 +17,7 @@ defmodule OpenAgents.Stacks.OperationWorker do

16 17
17 18
  alias OpenAgents.OperationalLog
18 19
  alias OpenAgents.Repo
20
  alias OpenAgents.Stacks.Merge
19 21
  alias OpenAgents.Stacks.Operation
20 22
  alias OpenAgents.Stacks.Restack
21 23

@@ -30,7 +32,7 @@ defmodule OpenAgents.Stacks.OperationWorker do

30 32
31 33
  def drain(server \\ __MODULE__), do: GenServer.call(server, :drain, 30_000)
32 34
33
  def run_once(executor \\ &Restack.execute/1) when is_function(executor, 1) do
35
  def run_once(executor \\ &execute/1) when is_function(executor, 1) do
34 36
    case claim_next() do
35 37
      nil ->
36 38
        :idle

@@ -41,10 +43,14 @@ defmodule OpenAgents.Stacks.OperationWorker do

41 43
    end
42 44
  end
43 45
46
  @doc "Dispatches one claimed operation to its kind's executor."
47
  def execute(%Operation{kind: "rebase"} = operation), do: Restack.execute(operation)
48
  def execute(%Operation{kind: "merge"} = operation), do: Merge.execute(operation)
49
44 50
  @impl true
45 51
  def init(options) do
46 52
    state = %{
47
      executor: Keyword.get(options, :executor, &Restack.execute/1),
53
      executor: Keyword.get(options, :executor, &execute/1),
48 54
      poll_interval_ms: Keyword.get(options, :poll_interval_ms, poll_interval_ms())
49 55
    }
50 56
lib/openagents/stacks/stack_event.ex modified +1 -1

@@ -13,7 +13,7 @@ defmodule OpenAgents.Stacks.StackEvent do

13 13
  @foreign_key_type :binary_id
14 14
  @timestamps_opts [type: :utc_datetime_usec]
15 15
16
  @event_types ~w(pull_request_stack.created pull_request_stack.appended pull_request_stack.rebased pull_request.synchronize)
16
  @event_types ~w(pull_request_stack.created pull_request_stack.appended pull_request_stack.rebased pull_request_stack.merged pull_request.synchronize)
17 17
18 18
  schema "pull_request_stack_events" do
19 19
    belongs_to :stack, OpenAgents.Stacks.Stack
lib/openagents_web/api_route_authority.ex modified +1

@@ -110,6 +110,7 @@ defmodule OpenAgentsWeb.ApiRouteAuthority do

110 110
      "post /api/v3/repos/:owner/:repo/stacks" => :required_bearer,
111 111
      "post /api/v3/repos/:owner/:repo/stacks/:stack_number/append" => :required_bearer,
112 112
      "post /api/v3/repos/:owner/:repo/stacks/:stack_number/rebase" => :required_bearer,
113
      "post /api/v3/repos/:owner/:repo/stacks/:stack_number/merge" => :required_bearer,
113 114
      "post /api/v3/repos/:owner/:repo/stacks/:stack_number/operations/:operation_id/continue" =>
114 115
        :required_bearer,
115 116
      "post /api/v3/repos/:owner/:repo/stacks/:stack_number/operations/:operation_id/abort" =>
lib/openagents_web/controllers/stack_controller.ex modified +33 -2

@@ -3,6 +3,7 @@ defmodule OpenAgentsWeb.StackController do

3 3
4 4
  alias OpenAgents.Repositories
5 5
  alias OpenAgents.Stacks
6
  alias OpenAgents.Stacks.Merge
6 7
  alias OpenAgents.Stacks.Restack
7 8
  alias OpenAgentsWeb.ControllerHelpers
8 9

@@ -79,6 +80,28 @@ defmodule OpenAgentsWeb.StackController do

79 80
    Ecto.NoResultsError -> not_found(conn)
80 81
  end
81 82
83
  def merge(conn, %{"owner" => owner, "repo" => repo, "stack_number" => number} = params) do
84
    repository = Repositories.get_visible_by_path!(owner, repo, conn.assigns.current_user)
85
86
    with {:ok, idempotency_key} <- idempotency_key(conn),
87
         {:ok, {operation, replay_state}} <-
88
           Merge.request_from_api(
89
             repository,
90
             ControllerHelpers.integer_param!(number),
91
             params,
92
             conn.assigns.current_user,
93
             idempotency_key
94
           ) do
95
      conn
96
      |> put_status(:accepted)
97
      |> render(:operation, operation: operation, replay_state: replay_state)
98
    else
99
      {:error, reason} -> render_error(conn, reason)
100
    end
101
  rescue
102
    Ecto.NoResultsError -> not_found(conn)
103
  end
104
82 105
  def show_operation(conn, %{"owner" => owner, "repo" => repo} = params) do
83 106
    repository = Repositories.get_visible_by_path!(owner, repo, conn.assigns[:current_user])
84 107

@@ -157,7 +180,8 @@ defmodule OpenAgentsWeb.StackController do

157 180
              :stack_not_open,
158 181
              :operation_in_progress,
159 182
              :operation_not_waiting,
160
              :operation_not_abortable
183
              :operation_not_abortable,
184
              :merge_queue_unavailable
161 185
            ],
162 186
       do: conflict(conn, reason)
163 187

@@ -176,7 +200,8 @@ defmodule OpenAgentsWeb.StackController do

176 200
              :already_stacked,
177 201
              :not_stack_top,
178 202
              :resolution_not_found,
179
              :resolution_parent_mismatch
203
              :resolution_parent_mismatch,
204
              :pull_request_not_in_stack
180 205
            ] do
181 206
    conn
182 207
    |> put_status(:unprocessable_entity)

@@ -210,6 +235,12 @@ defmodule OpenAgentsWeb.StackController do

210 235
  defp message(:operation_not_abortable), do: "The operation can no longer be aborted."
211 236
  defp message(:resolution_not_found), do: "The resolution commit does not exist."
212 237
238
  defp message(:merge_queue_unavailable),
239
    do: "The merge queue is not available; use a direct merge."
240
241
  defp message(:pull_request_not_in_stack),
242
    do: "The pull request is not an active entry of this stack."
243
213 244
  defp message(:resolution_parent_mismatch),
214 245
    do: "The resolution commit does not build on the persisted parent."
215 246
lib/openagents_web/router.ex modified +1

@@ -453,6 +453,7 @@ defmodule OpenAgentsWeb.Router do

453 453
    post "/repos/:owner/:repo/stacks", StackController, :create
454 454
    post "/repos/:owner/:repo/stacks/:stack_number/append", StackController, :append
455 455
    post "/repos/:owner/:repo/stacks/:stack_number/rebase", StackController, :rebase
456
    post "/repos/:owner/:repo/stacks/:stack_number/merge", StackController, :merge
456 457
457 458
    post "/repos/:owner/:repo/stacks/:stack_number/operations/:operation_id/continue",
458 459
         StackController,
test/openagents/stacks/merge_test.exs added +725

@@ -0,0 +1,725 @@

1
defmodule OpenAgents.Stacks.MergeTest do
2
  @moduledoc """
3
  Contiguous-prefix stack merge (#51): the merge-commit, squash, and rebase
4
  methods with exact resulting histories, prefix enforcement, upper-layer
5
  restacking after a partial merge, preflight failures, and injected crash
6
  recovery through the persisted plan.
7
  """
8
9
  use OpenAgents.DataCase, async: false
10
11
  alias OpenAgents.Forge.Repos
12
  alias OpenAgents.Issues.Issue
13
  alias OpenAgents.PullRequests.PullRequest
14
  alias OpenAgents.Repo
15
  alias OpenAgents.Stacks
16
  alias OpenAgents.Stacks.Merge
17
  alias OpenAgents.Stacks.Operation
18
  alias OpenAgents.Stacks.OperationWorker
19
  alias OpenAgents.Stacks.Stack
20
  alias OpenAgents.Stacks.StackEntry
21
  alias OpenAgents.Stacks.StackEvent
22
23
  import Ecto.Query
24
  import OpenAgents.AccountsFixtures
25
  import OpenAgents.IssuesFixtures
26
27
  setup do
28
    base = Path.join(System.tmp_dir!(), "stack-merge-#{System.unique_integer([:positive])}")
29
30
    previous_data = Application.get_env(:openagents, :forge_data_dir)
31
    previous_wal = Application.get_env(:openagents, :forge_wal_dir)
32
    Application.put_env(:openagents, :forge_data_dir, Path.join(base, "data"))
33
    Application.put_env(:openagents, :forge_wal_dir, Path.join(base, "wal"))
34
35
    on_exit(fn ->
36
      restore_env(:forge_data_dir, previous_data)
37
      restore_env(:forge_wal_dir, previous_wal)
38
      File.rm_rf(base)
39
    end)
40
41
    actor = repository_user_fixture("merge-actor")
42
    repository = repository_with_member_fixture(actor)
43
44
    %{actor: actor, repository: repository}
45
  end
46
47
  describe "the merge-commit method" do
48
    test "lands one group merge commit whose tree is the selected top head's", context do
49
      %{repository: repository, actor: actor} = context
50
51
      %{path: path, oids: oids, stack: stack, pull_requests: [pr_1, pr_2]} =
52
        seed_stack(repository, actor)
53
54
      {:ok, {operation, :created}} =
55
        Merge.request_from_api(
56
          repository,
57
          stack.number,
58
          %{"pull_request_number" => pr_2.issue.number, "merge_method" => "merge"},
59
          actor,
60
          "merge-1"
61
        )
62
63
      assert operation.state == "pending"
64
      assert operation.target_position == 2
65
      assert reload(Stack, stack.id).health == "operation_in_progress"
66
67
      assert :processed = OperationWorker.run_once()
68
69
      operation = reload(Operation, operation.id)
70
      assert operation.state == "succeeded"
71
72
      merge_commit = show(path, ["rev-parse", "refs/heads/main"])
73
      parents = show(path, ["show", "-s", "--format=%P", merge_commit])
74
      assert parents == "#{oids["main"]} #{oids["layer-2"]}"
75
76
      assert show(path, ["rev-parse", merge_commit <> "^{tree}"]) ==
77
               show(path, ["rev-parse", oids["layer-2"] <> "^{tree}"])
78
79
      # Both selected pull requests merged with the group merge commit.
80
      for pr <- [pr_1, pr_2] do
81
        merged = Repo.one!(from p in PullRequest, where: p.id == ^pr.id)
82
        assert merged.state == "closed"
83
        assert merged.merge_commit_sha == merge_commit
84
        assert merged.merged_by_user_id == actor.id
85
        refute is_nil(merged.merged_at)
86
87
        issue = Repo.one!(from i in Issue, where: i.id == ^pr.issue_id)
88
        assert issue.state == "closed"
89
        assert issue.state_reason == "completed"
90
      end
91
92
      # The merged branches did not move and were not deleted.
93
      assert show(path, ["rev-parse", "refs/heads/layer-1"]) == oids["layer-1"]
94
      assert show(path, ["rev-parse", "refs/heads/layer-2"]) == oids["layer-2"]
95
96
      stack = reload(Stack, stack.id)
97
      assert stack.state == "completed"
98
      assert stack.health == "healthy"
99
      assert stack.version == 2
100
      assert entries(stack.id) == []
101
102
      assert Repo.exists?(
103
               from event in StackEvent,
104
                 where:
105
                   event.event_type == "pull_request_stack.merged" and
106
                     event.stack_id == ^stack.id
107
             )
108
109
      {internal, 0} = Repos.git(path, ["for-each-ref", "refs/internal/"])
110
      assert internal == ""
111
    end
112
  end
113
114
  describe "the squash method" do
115
    test "chains one commit per pull request so each trunk commit is one layer", context do
116
      %{repository: repository, actor: actor} = context
117
118
      %{path: path, oids: oids, stack: stack, pull_requests: [pr_1, pr_2]} =
119
        seed_stack(repository, actor)
120
121
      {:ok, {_operation, :created}} =
122
        Merge.request_from_api(
123
          repository,
124
          stack.number,
125
          %{"pull_request_number" => pr_2.issue.number, "merge_method" => "squash"},
126
          actor,
127
          "squash-1"
128
        )
129
130
      assert :processed = OperationWorker.run_once()
131
132
      squash_2 = show(path, ["rev-parse", "refs/heads/main"])
133
      squash_1 = show(path, ["rev-parse", squash_2 <> "^"])
134
      assert show(path, ["rev-parse", squash_1 <> "^"]) == oids["main"]
135
136
      # Each squash commit's tree is exactly its layer head's tree.
137
      assert show(path, ["rev-parse", squash_1 <> "^{tree}"]) ==
138
               show(path, ["rev-parse", oids["layer-1"] <> "^{tree}"])
139
140
      assert show(path, ["rev-parse", squash_2 <> "^{tree}"]) ==
141
               show(path, ["rev-parse", oids["layer-2"] <> "^{tree}"])
142
143
      # Messages name the pull request and the original author survives.
144
      assert show(path, ["show", "-s", "--format=%s", squash_1]) ==
145
               "PR layer-1 (##{pr_1.issue.number})"
146
147
      assert show(path, ["show", "-s", "--format=%an <%ae>", squash_1]) ==
148
               "Test Author <author@example.test>"
149
150
      merged_1 = Repo.one!(from p in PullRequest, where: p.id == ^pr_1.id)
151
      merged_2 = Repo.one!(from p in PullRequest, where: p.id == ^pr_2.id)
152
      assert merged_1.merge_commit_sha == squash_1
153
      assert merged_2.merge_commit_sha == squash_2
154
    end
155
  end
156
157
  describe "the rebase method" do
158
    test "replays each layer's commits in order onto the trunk", context do
159
      %{repository: repository, actor: actor} = context
160
161
      %{path: path, oids: oids, stack: stack, pull_requests: [_pr_1, pr_2]} =
162
        seed_stack(repository, actor)
163
164
      {:ok, {_operation, :created}} =
165
        Merge.request_from_api(
166
          repository,
167
          stack.number,
168
          %{"pull_request_number" => pr_2.issue.number, "merge_method" => "rebase"},
169
          actor,
170
          "rebase-1"
171
        )
172
173
      assert :processed = OperationWorker.run_once()
174
175
      new_tip = show(path, ["rev-parse", "refs/heads/main"])
176
      new_mid = show(path, ["rev-parse", new_tip <> "^"])
177
      assert show(path, ["rev-parse", new_mid <> "^"]) == oids["main"]
178
179
      # The final tree equals the selected top head's tree, and the replayed
180
      # commits keep their messages.
181
      assert show(path, ["rev-parse", new_tip <> "^{tree}"]) ==
182
               show(path, ["rev-parse", oids["layer-2"] <> "^{tree}"])
183
184
      assert show(path, ["show", "-s", "--format=%s", new_mid]) == "Layer layer-1"
185
      assert show(path, ["show", "-s", "--format=%s", new_tip]) == "Layer layer-2"
186
    end
187
  end
188
189
  describe "prefix enforcement" do
190
    test "selecting the bottom layer merges only it and restacks the layers above", context do
191
      %{repository: repository, actor: actor} = context
192
193
      %{path: path, oids: oids, stack: stack, pull_requests: [pr_1, pr_2]} =
194
        seed_stack(repository, actor)
195
196
      {:ok, {_operation, :created}} =
197
        Merge.request_from_api(
198
          repository,
199
          stack.number,
200
          %{"pull_request_number" => pr_1.issue.number, "merge_method" => "merge"},
201
          actor,
202
          "prefix-1"
203
        )
204
205
      assert :processed = OperationWorker.run_once()
206
207
      trunk_new = show(path, ["rev-parse", "refs/heads/main"])
208
209
      merged_1 = Repo.one!(from p in PullRequest, where: p.id == ^pr_1.id)
210
      assert merged_1.state == "closed"
211
      assert merged_1.merge_commit_sha == trunk_new
212
213
      # The upper layer stays open, restacked onto the merge result, and
214
      # retargets the trunk.
215
      open_2 = Repo.one!(from p in PullRequest, where: p.id == ^pr_2.id)
216
      assert open_2.state == "open"
217
      assert open_2.base_ref == "main"
218
      assert open_2.base_sha == trunk_new
219
220
      new_layer_2 = show(path, ["rev-parse", "refs/heads/layer-2"])
221
      refute new_layer_2 == oids["layer-2"]
222
      assert show(path, ["rev-parse", new_layer_2 <> "^"]) == trunk_new
223
      assert open_2.head_sha == new_layer_2
224
225
      stack = reload(Stack, stack.id)
226
      assert stack.state == "open"
227
      assert stack.health == "healthy"
228
229
      [entry_2] = entries(stack.id)
230
      assert entry_2.boundary_oid == trunk_new
231
      assert entry_2.observed_head_oid == new_layer_2
232
233
      assert Repo.exists?(
234
               from event in StackEvent,
235
                 where:
236
                   event.event_type == "pull_request.synchronize" and
237
                     event.stack_id == ^stack.id
238
             )
239
    end
240
241
    test "a lower-prefix merge restacks every upper layer in order", context do
242
      %{repository: repository, actor: actor} = context
243
244
      %{path: path, oids: oids, stack: stack, pull_requests: [pr_1, pr_2, pr_3, pr_4]} =
245
        seed_stack(repository, actor, ["layer-1", "layer-2", "layer-3", "layer-4"])
246
247
      {:ok, {_operation, :created}} =
248
        Merge.request_from_api(
249
          repository,
250
          stack.number,
251
          %{"pull_request_number" => pr_2.issue.number, "merge_method" => "squash"},
252
          actor,
253
          "prefix-multi-1"
254
        )
255
256
      assert :processed = OperationWorker.run_once()
257
258
      trunk_new = show(path, ["rev-parse", "refs/heads/main"])
259
260
      for pr <- [pr_1, pr_2] do
261
        assert Repo.one!(from p in PullRequest, where: p.id == ^pr.id).state == "closed"
262
      end
263
264
      # Layer 3 restacked onto the squash result, layer 4 onto the new
265
      # layer 3, and only the lowest remaining layer retargets the trunk.
266
      new_layer_3 = show(path, ["rev-parse", "refs/heads/layer-3"])
267
      new_layer_4 = show(path, ["rev-parse", "refs/heads/layer-4"])
268
      refute new_layer_3 == oids["layer-3"]
269
      refute new_layer_4 == oids["layer-4"]
270
      assert show(path, ["rev-parse", new_layer_3 <> "^"]) == trunk_new
271
      assert show(path, ["rev-parse", new_layer_4 <> "^"]) == new_layer_3
272
273
      open_3 = Repo.one!(from p in PullRequest, where: p.id == ^pr_3.id)
274
      assert open_3.state == "open"
275
      assert open_3.base_ref == "main"
276
      assert open_3.base_sha == trunk_new
277
      assert open_3.head_sha == new_layer_3
278
279
      open_4 = Repo.one!(from p in PullRequest, where: p.id == ^pr_4.id)
280
      assert open_4.state == "open"
281
      assert open_4.base_ref == "layer-3"
282
      assert open_4.base_sha == new_layer_3
283
      assert open_4.head_sha == new_layer_4
284
285
      stack = reload(Stack, stack.id)
286
      assert stack.state == "open"
287
      assert stack.health == "healthy"
288
289
      [entry_3, entry_4] = entries(stack.id)
290
      assert entry_3.boundary_oid == trunk_new
291
      assert entry_3.observed_head_oid == new_layer_3
292
      assert entry_4.boundary_oid == new_layer_3
293
      assert entry_4.observed_head_oid == new_layer_4
294
    end
295
296
    test "a pull request outside the stack cannot be selected", context do
297
      %{repository: repository, actor: actor} = context
298
      %{stack: stack} = seed_stack(repository, actor)
299
300
      assert {:error, :pull_request_not_in_stack} =
301
               Merge.request_from_api(
302
                 repository,
303
                 stack.number,
304
                 %{"pull_request_number" => 999_999, "merge_method" => "merge"},
305
                 actor,
306
                 "prefix-missing-1"
307
               )
308
    end
309
  end
310
311
  describe "preflight" do
312
    test "a stack behind its trunk fails with needs_rebase before any ref moves", context do
313
      %{repository: repository, actor: actor} = context
314
315
      %{path: path, oids: oids, stack: stack, pull_requests: [_pr_1, pr_2]} =
316
        seed_stack(repository, actor)
317
318
      trunk_tip = commit(path, oids["main"], "Trunk advance", %{"trunk.md" => "trunk\n"})
319
      {_, 0} = Repos.git(path, ["update-ref", "refs/heads/main", trunk_tip])
320
321
      {:ok, {operation, :created}} =
322
        Merge.request_from_api(
323
          repository,
324
          stack.number,
325
          %{"pull_request_number" => pr_2.issue.number, "merge_method" => "merge"},
326
          actor,
327
          "behind-1"
328
        )
329
330
      assert :processed = OperationWorker.run_once()
331
332
      operation = reload(Operation, operation.id)
333
      assert operation.state == "failed"
334
      assert operation.error["code"] == "needs_rebase"
335
      assert reload(Stack, stack.id).health == "needs_rebase"
336
337
      assert show(path, ["rev-parse", "refs/heads/main"]) == trunk_tip
338
      assert Repo.one!(from p in PullRequest, where: p.id == ^pr_2.id).state == "open"
339
    end
340
341
    test "a concurrent branch push fails the merge and preserves the branch", context do
342
      %{repository: repository, actor: actor} = context
343
344
      %{path: path, oids: oids, stack: stack, pull_requests: [_pr_1, pr_2]} =
345
        seed_stack(repository, actor)
346
347
      {:ok, {operation, :created}} =
348
        Merge.request_from_api(
349
          repository,
350
          stack.number,
351
          %{"pull_request_number" => pr_2.issue.number, "merge_method" => "merge"},
352
          actor,
353
          "race-1"
354
        )
355
356
      pushed = commit(path, oids["layer-1"], "User push", %{"extra.md" => "extra\n"})
357
      {_, 0} = Repos.git(path, ["update-ref", "refs/heads/layer-1", pushed])
358
359
      assert :processed = OperationWorker.run_once()
360
361
      operation = reload(Operation, operation.id)
362
      assert operation.state == "failed"
363
      assert operation.error["code"] == "head_changed"
364
365
      assert show(path, ["rev-parse", "refs/heads/layer-1"]) == pushed
366
      assert show(path, ["rev-parse", "refs/heads/main"]) == oids["main"]
367
    end
368
369
    test "an expected head mismatch fails before any ref moves", context do
370
      %{repository: repository, actor: actor} = context
371
372
      %{path: path, oids: oids, stack: stack, pull_requests: [pr_1, pr_2]} =
373
        seed_stack(repository, actor)
374
375
      wrong = String.duplicate("ab", 20)
376
377
      {:ok, {operation, :created}} =
378
        Merge.request_from_api(
379
          repository,
380
          stack.number,
381
          %{
382
            "pull_request_number" => pr_2.issue.number,
383
            "merge_method" => "merge",
384
            "expected_heads" => %{"#{pr_1.issue.number}" => wrong}
385
          },
386
          actor,
387
          "expected-1"
388
        )
389
390
      assert :processed = OperationWorker.run_once()
391
392
      operation = reload(Operation, operation.id)
393
      assert operation.state == "failed"
394
      assert operation.error["code"] == "expected_head_mismatch"
395
      assert show(path, ["rev-parse", "refs/heads/main"]) == oids["main"]
396
    end
397
398
    test "a stale expected stack version is rejected at request time", context do
399
      %{repository: repository, actor: actor} = context
400
      %{stack: stack, pull_requests: [_pr_1, pr_2]} = seed_stack(repository, actor)
401
402
      assert {:error, :stale_stack_version} =
403
               Merge.request_from_api(
404
                 repository,
405
                 stack.number,
406
                 %{
407
                   "pull_request_number" => pr_2.issue.number,
408
                   "merge_method" => "merge",
409
                   "expected_stack_version" => 42
410
                 },
411
                 actor,
412
                 "stale-1"
413
               )
414
    end
415
416
    test "queueing is not yet available", context do
417
      %{repository: repository, actor: actor} = context
418
      %{stack: stack, pull_requests: [_pr_1, pr_2]} = seed_stack(repository, actor)
419
420
      assert {:error, :merge_queue_unavailable} =
421
               Merge.request_from_api(
422
                 repository,
423
                 stack.number,
424
                 %{
425
                   "pull_request_number" => pr_2.issue.number,
426
                   "merge_method" => "merge",
427
                   "merge_action" => "queue"
428
                 },
429
                 actor,
430
                 "queue-1"
431
               )
432
    end
433
434
    test "a reader cannot request a merge", context do
435
      %{repository: repository, actor: actor} = context
436
      %{stack: stack, pull_requests: [_pr_1, pr_2]} = seed_stack(repository, actor)
437
      reader = repository_user_fixture("merge-reader")
438
439
      assert {:error, :forbidden} =
440
               Merge.request_from_api(
441
                 repository,
442
                 stack.number,
443
                 %{"pull_request_number" => pr_2.issue.number, "merge_method" => "merge"},
444
                 reader,
445
                 "reader-1"
446
               )
447
    end
448
  end
449
450
  describe "idempotency" do
451
    test "a retried key replays the original operation", context do
452
      %{repository: repository, actor: actor} = context
453
      %{stack: stack, pull_requests: [_pr_1, pr_2]} = seed_stack(repository, actor)
454
455
      request = %{"pull_request_number" => pr_2.issue.number, "merge_method" => "merge"}
456
457
      {:ok, {operation, :created}} =
458
        Merge.request_from_api(repository, stack.number, request, actor, "idem-1")
459
460
      assert {:ok, {replayed, :replayed}} =
461
               Merge.request_from_api(repository, stack.number, request, actor, "idem-1")
462
463
      assert replayed.id == operation.id
464
465
      assert {:error, :idempotency_conflict} =
466
               Merge.request_from_api(
467
                 repository,
468
                 stack.number,
469
                 Map.put(request, "merge_method", "squash"),
470
                 actor,
471
                 "idem-1"
472
               )
473
    end
474
  end
475
476
  describe "crash recovery" do
477
    test "a crash between the git refs and the metadata reconciles to success", context do
478
      %{repository: repository, actor: actor} = context
479
      %{path: path, stack: stack, pull_requests: [pr_1, pr_2]} = seed_stack(repository, actor)
480
481
      {:ok, {operation, :created}} =
482
        Merge.request_from_api(
483
          repository,
484
          stack.number,
485
          %{"pull_request_number" => pr_2.issue.number, "merge_method" => "merge"},
486
          actor,
487
          "crash-1"
488
        )
489
490
      # A worker claims the row, moves the refs, and dies before the
491
      # metadata transaction.
492
      {1, _rows} =
493
        Repo.update_all(
494
          from(o in Operation, where: o.id == ^operation.id),
495
          set: [state: "running", claimed_at: DateTime.utc_now(), attempt_count: 1]
496
        )
497
498
      claimed = reload(Operation, operation.id)
499
500
      assert_raise RuntimeError, "injected crash", fn ->
501
        Merge.execute(claimed, after_refs: fn -> raise "injected crash" end)
502
      end
503
504
      # The refs landed but the metadata did not: the pull requests still
505
      # read open and the plan is marked applied.
506
      crashed = reload(Operation, operation.id)
507
      assert crashed.state == "running"
508
      assert crashed.planned_result["refs_applied"] == true
509
      assert show(path, ["rev-parse", "refs/heads/main"]) == crashed.planned_result["trunk_new"]
510
      assert Repo.one!(from p in PullRequest, where: p.id == ^pr_1.id).state == "open"
511
512
      # The lease expires and another worker reconciles instead of
513
      # re-merging.
514
      stale = DateTime.add(DateTime.utc_now(), -600, :second)
515
516
      {1, _rows} =
517
        Repo.update_all(
518
          from(o in Operation, where: o.id == ^operation.id),
519
          set: [claimed_at: stale]
520
        )
521
522
      assert :processed = OperationWorker.run_once()
523
524
      operation = reload(Operation, operation.id)
525
      assert operation.state == "succeeded"
526
      assert operation.attempt_count == 2
527
528
      trunk_new = show(path, ["rev-parse", "refs/heads/main"])
529
      assert trunk_new == operation.planned_result["trunk_new"]
530
531
      for pr <- [pr_1, pr_2] do
532
        merged = Repo.one!(from p in PullRequest, where: p.id == ^pr.id)
533
        assert merged.state == "closed"
534
        assert merged.merge_commit_sha == trunk_new
535
      end
536
537
      assert reload(Stack, stack.id).state == "completed"
538
    end
539
540
    test "a crash after a diverged repository marks the operation partial", context do
541
      %{repository: repository, actor: actor} = context
542
543
      %{path: path, oids: oids, stack: stack, pull_requests: [_pr_1, pr_2]} =
544
        seed_stack(repository, actor)
545
546
      {:ok, {operation, :created}} =
547
        Merge.request_from_api(
548
          repository,
549
          stack.number,
550
          %{"pull_request_number" => pr_2.issue.number, "merge_method" => "merge"},
551
          actor,
552
          "diverge-1"
553
        )
554
555
      {1, _rows} =
556
        Repo.update_all(
557
          from(o in Operation, where: o.id == ^operation.id),
558
          set: [state: "running", claimed_at: DateTime.utc_now(), attempt_count: 1]
559
        )
560
561
      claimed = reload(Operation, operation.id)
562
563
      assert_raise RuntimeError, "injected crash", fn ->
564
        Merge.execute(claimed, after_refs: fn -> raise "injected crash" end)
565
      end
566
567
      # A trunk push lands after the crash, so the live refs no longer match
568
      # the plan.
569
      trunk_new = show(path, ["rev-parse", "refs/heads/main"])
570
      diverged = commit(path, trunk_new, "Post-crash push", %{"post.md" => "post\n"})
571
572
      {:ok, _result} =
573
        OpenAgents.Forge.GitPlane.batch_update_refs(
574
          repository.storage_key,
575
          [%{ref: "refs/heads/main", expected_old: trunk_new, new: diverged}],
576
          "test-push"
577
        )
578
579
      stale = DateTime.add(DateTime.utc_now(), -600, :second)
580
581
      {1, _rows} =
582
        Repo.update_all(
583
          from(o in Operation, where: o.id == ^operation.id),
584
          set: [claimed_at: stale]
585
        )
586
587
      assert :processed = OperationWorker.run_once()
588
589
      operation = reload(Operation, operation.id)
590
      assert operation.state == "partially_succeeded"
591
      assert operation.error["code"] == "refs_diverged"
592
      assert show(path, ["rev-parse", "refs/heads/main"]) == diverged
593
      # The merged prefix stays landed underneath the diverged push.
594
      assert oids["layer-2"] ==
595
               show(path, ["rev-parse", operation.planned_result["trunk_new"] <> "^2"])
596
    end
597
  end
598
599
  ## Fixtures
600
601
  # main ── layer-1 ── layer-2 ── …, each layer adding one file.
602
  defp seed_stack(repository, actor, branches \\ ["layer-1", "layer-2"]) do
603
    path = Repos.ensure_repo!(repository.storage_key, repository.default_branch)
604
605
    main = commit(path, nil, "Seed repository", %{"README.md" => "readme\n"})
606
    {_, 0} = Repos.git(path, ["update-ref", "refs/heads/main", main])
607
608
    {oids, _parent} =
609
      Enum.reduce(branches, {%{"main" => main}, main}, fn branch, {oids, parent} ->
610
        oid = commit(path, parent, "Layer #{branch}", %{"#{branch}.md" => "#{branch}\n"})
611
        {_, 0} = Repos.git(path, ["update-ref", "refs/heads/#{branch}", oid])
612
        {Map.put(oids, branch, oid), oid}
613
      end)
614
615
    pull_requests =
616
      branches
617
      |> Enum.with_index()
618
      |> Enum.map(fn {branch, index} ->
619
        base = Enum.at(["main" | branches], index)
620
        pull_request(repository, branch, base, oids[base], oids[branch])
621
      end)
622
623
    {:ok, stack} = Stacks.create(repository, pull_requests, actor)
624
625
    %{path: path, oids: oids, stack: stack, pull_requests: pull_requests}
626
  end
627
628
  defp pull_request(repository, head_ref, base_ref, base_sha, head_sha) do
629
    issue = issue_fixture(repository, %{title: "PR #{head_ref}"})
630
631
    {:ok, pull_request} =
632
      %PullRequest{}
633
      |> PullRequest.changeset(%{
634
        repository_id: repository.id,
635
        issue_id: issue.id,
636
        head_repository_id: repository.id,
637
        head_ref: head_ref,
638
        head_sha: head_sha,
639
        base_ref: base_ref,
640
        base_sha: base_sha,
641
        state: "open"
642
      })
643
      |> Repo.insert()
644
645
    Repo.preload(pull_request, :issue)
646
  end
647
648
  # Commits a tree that layers the given files over the parent's tree.
649
  defp commit(path, parent, message, files) do
650
    parent_entries =
651
      if parent do
652
        {listing, 0} = Repos.git(path, ["ls-tree", parent])
653
654
        listing
655
        |> String.split("\n", trim: true)
656
        |> Map.new(fn line ->
657
          [meta, name] = String.split(line, "\t", parts: 2)
658
          {name, meta <> "\t" <> name}
659
        end)
660
      else
661
        %{}
662
      end
663
664
    new_entries =
665
      Map.new(files, fn {name, content} ->
666
        blob = git!(path, ["hash-object", "-w", "--stdin"], content)
667
        {name, "100644 blob #{blob}\t#{name}"}
668
      end)
669
670
    listing =
671
      parent_entries
672
      |> Map.merge(new_entries)
673
      |> Map.values()
674
      |> Enum.map_join("", &(&1 <> "\n"))
675
676
    tree = git!(path, ["mktree"], listing)
677
    parent_args = if parent, do: ["-p", parent], else: []
678
679
    git!(path, ["commit-tree", tree] ++ parent_args ++ ["-m", message], "",
680
      env: [
681
        {"GIT_AUTHOR_NAME", "Test Author"},
682
        {"GIT_AUTHOR_EMAIL", "author@example.test"},
683
        {"GIT_COMMITTER_NAME", "Test Author"},
684
        {"GIT_COMMITTER_EMAIL", "author@example.test"}
685
      ]
686
    )
687
  end
688
689
  defp entries(stack_id) do
690
    Repo.all(
691
      from entry in StackEntry,
692
        where: entry.stack_id == ^stack_id and is_nil(entry.removed_at),
693
        order_by: [asc: entry.position]
694
    )
695
  end
696
697
  defp reload(schema, id), do: Repo.get!(schema, id)
698
699
  defp show(path, args) do
700
    {output, 0} = Repos.git(path, args)
701
    String.trim(output)
702
  end
703
704
  defp git!(git_dir, args, input, options \\ []) do
705
    input_path = Path.join(System.tmp_dir!(), "merge-input-#{System.unique_integer([:positive])}")
706
707
    File.write!(input_path, input)
708
709
    try do
710
      {output, 0} =
711
        System.cmd(
712
          "sh",
713
          ["-c", ~s(exec git --git-dir "$GIT_DIR" "$@" < "$INPUT"), "sh"] ++ args,
714
          env: [{"GIT_DIR", git_dir}, {"INPUT", input_path}] ++ Keyword.get(options, :env, [])
715
        )
716
717
      String.trim(output)
718
    after
719
      File.rm(input_path)
720
    end
721
  end
722
723
  defp restore_env(key, nil), do: Application.delete_env(:openagents, key)
724
  defp restore_env(key, value), do: Application.put_env(:openagents, key, value)
725
end
test/openagents_web/controllers/stack_controller_test.exs modified +103

@@ -414,6 +414,109 @@ defmodule OpenAgentsWeb.StackControllerTest do

414 414
    end
415 415
  end
416 416
417
  describe "POST /api/v3/repos/:owner/:repo/stacks/:stack_number/merge" do
418
    test "accepts a merge, exposes the operation, and replays retries", %{conn: conn} do
419
      repository = repository_fixture()
420
      oids = seed_chain(repository, ["layer-1"])
421
      [pr_1] = pull_request_chain(repository, oids, ["layer-1"])
422
      conn = put_forge_api_token(conn, "stack-merge", repository)
423
424
      assert %{"number" => 1} =
425
               conn
426
               |> put_req_header("idempotency-key", "merge-create-1")
427
               |> post(path(repository), %{trunk_ref: "main", pull_requests: [pr_1]})
428
               |> json_response(201)
429
430
      body = %{pull_request_number: pr_1, merge_method: "merge"}
431
432
      merge_conn =
433
        conn
434
        |> put_req_header("idempotency-key", "merge-1")
435
        |> post("#{path(repository)}/1/merge", body)
436
437
      assert %{
438
               "id" => operation_id,
439
               "kind" => "merge",
440
               "state" => "pending",
441
               "replayed" => false
442
             } = json_response(merge_conn, 202)
443
444
      assert %{"health" => "operation_in_progress"} =
445
               json_response(get(conn, "#{path(repository)}/1"), 200)
446
447
      replay_conn =
448
        conn
449
        |> put_req_header("idempotency-key", "merge-1")
450
        |> post("#{path(repository)}/1/merge", body)
451
452
      assert %{"id" => ^operation_id, "replayed" => true} = json_response(replay_conn, 202)
453
454
      show_conn = get(conn, "#{path(repository)}/1/operations/#{operation_id}")
455
      assert %{"id" => ^operation_id, "state" => "pending"} = json_response(show_conn, 200)
456
    end
457
458
    test "rejects a queue action, a stranger PR, and a bad method", %{conn: conn} do
459
      repository = repository_fixture()
460
      oids = seed_chain(repository, ["layer-1"])
461
      [pr_1] = pull_request_chain(repository, oids, ["layer-1"])
462
      conn = put_forge_api_token(conn, "stack-merge-invalid", repository)
463
464
      assert %{"number" => 1} =
465
               conn
466
               |> put_req_header("idempotency-key", "merge-invalid-create")
467
               |> post(path(repository), %{trunk_ref: "main", pull_requests: [pr_1]})
468
               |> json_response(201)
469
470
      queue_conn =
471
        conn
472
        |> put_req_header("idempotency-key", "merge-queue-1")
473
        |> post("#{path(repository)}/1/merge", %{
474
          pull_request_number: pr_1,
475
          merge_method: "merge",
476
          merge_action: "queue"
477
        })
478
479
      assert %{"code" => "merge_queue_unavailable"} = json_response(queue_conn, 409)
480
481
      stranger_conn =
482
        conn
483
        |> put_req_header("idempotency-key", "merge-stranger-1")
484
        |> post("#{path(repository)}/1/merge", %{
485
          pull_request_number: 999_999,
486
          merge_method: "merge"
487
        })
488
489
      assert %{"code" => "pull_request_not_in_stack"} = json_response(stranger_conn, 422)
490
491
      method_conn =
492
        conn
493
        |> put_req_header("idempotency-key", "merge-method-1")
494
        |> post("#{path(repository)}/1/merge", %{
495
          pull_request_number: pr_1,
496
          merge_method: "octopus"
497
        })
498
499
      assert %{"code" => "invalid_request"} = json_response(method_conn, 422)
500
    end
501
502
    test "refuses a caller without write access", %{conn: conn} do
503
      repository = repository_fixture()
504
      oids = seed_chain(repository, ["layer-1"])
505
      [pr_1] = pull_request_chain(repository, oids, ["layer-1"])
506
      conn = put_forge_api_token(conn, "stack-merge-outsider")
507
508
      conn =
509
        conn
510
        |> put_req_header("idempotency-key", "merge-forbidden-1")
511
        |> post("#{path(repository)}/1/merge", %{
512
          pull_request_number: pr_1,
513
          merge_method: "merge"
514
        })
515
516
      assert json_response(conn, 403)
517
    end
518
  end
519
417 520
  describe "GET /api/v3/repos/:owner/:repo/stacks" do
418 521
    test "reads are public for a public repository", %{conn: conn} do
419 522
      repository = repository_fixture()

This page updates live while a promote is in flight · changelog