Queue stacked merges as ordered logical items

4da7eff07b63 · Devin AI · · parent c17b6e8b6141

Queue stacked merges as ordered logical items

No merge queue exists on the forge yet, so this fixes the stack
contract any queue implementation must satisfy
(docs/stacked-prs.md section 14). A selected contiguous stack prefix
enters the queue as one logical item with ordered members, expected
head OIDs, the queue base OID, the stack version, and the policy
version. Speculation applies the current queue base, then every member
bottom to top. The queue may reorder and interleave independent items
but never the members of one stack, and ejecting a member ejects every
selected member above it: their prerequisite is no longer present. A
speculative result invalidates when the queue base SHA, the stack
version, any selected head SHA, or the policy version changes, and
every trigger reports.

Closes #54

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

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/stacks/merge.ex
  • added lib/openagents/stacks/merge_queue.ex
  • added lib/openagents/stacks/queue_item.ex
  • added test/openagents/stacks/merge_queue_test.exs

Diff

4 files changed, +428 -2

lib/openagents/stacks/merge.ex modified +3 -2

@@ -864,8 +864,9 @@ defmodule OpenAgents.Stacks.Merge do

864 864
    end
865 865
  end
866 866
867
  # The merge queue lands as its own slice; until then only a direct merge
868
  # is available.
867
  # No merge queue exists on the forge yet. OpenAgents.Stacks.MergeQueue
868
  # fixes the stack contract any queue implementation must satisfy; until
869
  # one lands, only a direct merge is available.
869 870
  defp parse_merge_action(params) do
870 871
    case Map.get(params, "merge_action", "direct_merge") do
871 872
      "direct_merge" -> :ok
lib/openagents/stacks/merge_queue.ex added +161

@@ -0,0 +1,161 @@

1
defmodule OpenAgents.Stacks.MergeQueue do
2
  @moduledoc """
3
  The stack contract for merge queues (docs/stacked-prs.md section 14).
4
5
  No merge queue exists on the forge yet; this module defines the contract
6
  any queue implementation must satisfy for stacked pull requests, so the
7
  behavior stays fixed before an implementation lands:
8
9
  - A selected stack prefix enters the queue as one logical item
10
    (`OpenAgents.Stacks.QueueItem`) with ordered members, expected heads,
11
    the queue base OID, the stack version, and the policy version.
12
  - Speculation applies the current queue base, then each member bottom
13
    to top — never another order.
14
  - The queue may reorder independent items freely, but it must never
15
    reorder members within one item.
16
  - Ejecting a member ejects every selected member above it in the same
17
    speculative group: their prerequisite is no longer present.
18
  - A speculative result invalidates when the queue base SHA, the stack
19
    version, any selected head SHA, or the policy version changes.
20
  """
21
22
  alias OpenAgents.Stacks.QueueItem
23
  alias OpenAgents.Stacks.Stack
24
  alias OpenAgents.Stacks.StackEntry
25
26
  @doc """
27
  Build one logical queue item from a selected contiguous prefix of a
28
  stack's active entries.
29
30
  `entries` must be the active entries selected for the merge, bottom
31
  first; the selection must be a contiguous prefix starting at position 1,
32
  matching the merge contract. `queue_base_oid` is the base the queue will
33
  speculate onto, and `policy_version` is an opaque identity of the policy
34
  configuration in force (for example the workflow definition blob OID).
35
  """
36
  def item(%Stack{} = stack, entries, queue_base_oid, policy_version)
37
      when is_list(entries) and entries != [] and is_binary(queue_base_oid) do
38
    with :ok <- check_prefix(entries) do
39
      {:ok,
40
       %QueueItem{
41
         stack_id: stack.id,
42
         stack_number: stack.number,
43
         stack_version: stack.version,
44
         queue_base_oid: queue_base_oid,
45
         policy_version: policy_version,
46
         members:
47
           Enum.map(entries, fn %StackEntry{} = entry ->
48
             %{
49
               position: entry.position,
50
               pull_request_id: entry.pull_request_id,
51
               expected_head_oid: entry.observed_head_oid
52
             }
53
           end)
54
       }}
55
    end
56
  end
57
58
  @doc """
59
  The ordered application plan for one item's speculative merge group:
60
  the current queue base, then every member head bottom to top.
61
  """
62
  def application_plan(%QueueItem{} = item) do
63
    [item.queue_base_oid | Enum.map(item.members, & &1.expected_head_oid)]
64
  end
65
66
  @doc """
67
  Whether a proposed global member order is a legal reordering of the
68
  given items.
69
70
  The queue may interleave and reorder independent items however it wants,
71
  but the members of one item must keep their relative order. `proposed`
72
  is the full flattened order as `{stack_id, position}` pairs; it must
73
  contain exactly the members of `items`.
74
  """
75
  def valid_order?(items, proposed) when is_list(items) and is_list(proposed) do
76
    expected =
77
      items
78
      |> Enum.flat_map(fn item ->
79
        Enum.map(item.members, &{item.stack_id, &1.position})
80
      end)
81
      |> Enum.sort()
82
83
    same_members = Enum.sort(proposed) == expected
84
85
    ordered_within_each_item =
86
      proposed
87
      |> Enum.group_by(fn {stack_id, _position} -> stack_id end)
88
      |> Enum.all?(fn {_stack_id, members} ->
89
        positions = Enum.map(members, fn {_stack_id, position} -> position end)
90
        positions == Enum.sort(positions)
91
      end)
92
93
    same_members and ordered_within_each_item
94
  end
95
96
  @doc """
97
  Eject the member at `position` from the item's speculative group.
98
99
  The ejected member's prerequisite chain breaks for everything selected
100
  above it, so every higher member leaves the group with it. Returns
101
  `{:ok, %{ejected: members, remaining: members}}` with both halves in
102
  order, or `{:error, :not_a_member}` when the position is not selected.
103
  """
104
  def eject(%QueueItem{} = item, position) when is_integer(position) do
105
    if Enum.any?(item.members, &(&1.position == position)) do
106
      {remaining, ejected} = Enum.split_while(item.members, &(&1.position < position))
107
      {:ok, %{ejected: ejected, remaining: remaining}}
108
    else
109
      {:error, :not_a_member}
110
    end
111
  end
112
113
  @doc """
114
  Every reason the item's speculative result is no longer valid against
115
  the observed current state.
116
117
  `observed` carries `queue_base_oid`, `stack_version`, `policy_version`,
118
  and `head_oids` (a map of position to current head OID). An empty list
119
  means the speculation is still keyed to reality; any entry —
120
  `:queue_base_changed`, `:stack_version_changed`, `:policy_changed`, or
121
  `{:head_changed, position}` — invalidates it.
122
  """
123
  def invalidations(%QueueItem{} = item, observed) when is_map(observed) do
124
    base =
125
      if Map.fetch!(observed, :queue_base_oid) == item.queue_base_oid,
126
        do: [],
127
        else: [:queue_base_changed]
128
129
    version =
130
      if Map.fetch!(observed, :stack_version) == item.stack_version,
131
        do: [],
132
        else: [:stack_version_changed]
133
134
    policy =
135
      if Map.fetch!(observed, :policy_version) == item.policy_version,
136
        do: [],
137
        else: [:policy_changed]
138
139
    head_oids = Map.fetch!(observed, :head_oids)
140
141
    heads =
142
      item.members
143
      |> Enum.filter(&(Map.get(head_oids, &1.position) != &1.expected_head_oid))
144
      |> Enum.map(&{:head_changed, &1.position})
145
146
    base ++ version ++ policy ++ heads
147
  end
148
149
  @doc "Whether the item's speculative result is still valid."
150
  def valid?(%QueueItem{} = item, observed), do: invalidations(item, observed) == []
151
152
  defp check_prefix(entries) do
153
    positions = Enum.map(entries, & &1.position)
154
155
    if positions == Enum.to_list(1..length(entries)) do
156
      :ok
157
    else
158
      {:error, :not_a_prefix}
159
    end
160
  end
161
end
lib/openagents/stacks/queue_item.ex added +22

@@ -0,0 +1,22 @@

1
defmodule OpenAgents.Stacks.QueueItem do
2
  @moduledoc """
3
  One logical merge-queue item for a selected stack prefix
4
  (docs/stacked-prs.md section 14).
5
6
  The queue treats the whole prefix as one item with ordered internal
7
  members. Each member records the position, pull request, and expected
8
  head OID it was selected with; the item records the queue base OID the
9
  speculation applies onto, the stack version, and the policy version, so
10
  a speculative result stays keyed to everything that produced it.
11
  """
12
13
  @enforce_keys [:stack_id, :stack_number, :stack_version, :queue_base_oid, :members]
14
  defstruct [
15
    :stack_id,
16
    :stack_number,
17
    :stack_version,
18
    :queue_base_oid,
19
    :policy_version,
20
    :members
21
  ]
22
end
test/openagents/stacks/merge_queue_test.exs added +242

@@ -0,0 +1,242 @@

1
defmodule OpenAgents.Stacks.MergeQueueTest do
2
  @moduledoc """
3
  The stack merge-queue contract (#54): logical item grouping with expected
4
  heads and the queue base, bottom-to-top speculation order, reordering
5
  bounds, cascade ejection, and every invalidation trigger — against a fake
6
  queue base, since no queue implementation exists on the forge yet.
7
  """
8
9
  use ExUnit.Case, async: true
10
11
  alias OpenAgents.Stacks.MergeQueue
12
  alias OpenAgents.Stacks.QueueItem
13
  alias OpenAgents.Stacks.Stack
14
  alias OpenAgents.Stacks.StackEntry
15
16
  @queue_base String.duplicate("a", 40)
17
  @policy String.duplicate("f", 40)
18
19
  describe "item/4" do
20
    test "groups a selected prefix with expected heads, queue base, and versions" do
21
      stack = stack(version: 4)
22
      entries = entries(["1111", "2222", "3333"])
23
24
      {:ok, item} = MergeQueue.item(stack, entries, @queue_base, @policy)
25
26
      assert item.stack_id == stack.id
27
      assert item.stack_number == stack.number
28
      assert item.stack_version == 4
29
      assert item.queue_base_oid == @queue_base
30
      assert item.policy_version == @policy
31
32
      assert Enum.map(item.members, & &1.position) == [1, 2, 3]
33
34
      assert Enum.map(item.members, & &1.expected_head_oid) ==
35
               Enum.map(entries, & &1.observed_head_oid)
36
37
      assert Enum.map(item.members, & &1.pull_request_id) ==
38
               Enum.map(entries, & &1.pull_request_id)
39
    end
40
41
    test "rejects a selection that is not a contiguous prefix" do
42
      stack = stack(version: 1)
43
      [entry_1, _entry_2, entry_3] = entries(["1111", "2222", "3333"])
44
45
      assert {:error, :not_a_prefix} =
46
               MergeQueue.item(stack, [entry_1, entry_3], @queue_base, @policy)
47
48
      assert {:error, :not_a_prefix} =
49
               MergeQueue.item(stack, [entry_3], @queue_base, @policy)
50
    end
51
  end
52
53
  describe "application_plan/1" do
54
    test "applies the current queue base, then every member bottom to top" do
55
      {:ok, item} =
56
        MergeQueue.item(stack(version: 1), entries(["1111", "2222"]), @queue_base, @policy)
57
58
      assert MergeQueue.application_plan(item) == [
59
               @queue_base,
60
               oid("1111"),
61
               oid("2222")
62
             ]
63
    end
64
  end
65
66
  describe "valid_order?/2" do
67
    test "independent items may reorder and interleave" do
68
      {:ok, item_a} =
69
        MergeQueue.item(stack(id: "a"), entries(["1111", "2222"]), @queue_base, @policy)
70
71
      {:ok, item_b} =
72
        MergeQueue.item(stack(id: "b"), entries(["3333", "4444"]), @queue_base, @policy)
73
74
      assert MergeQueue.valid_order?([item_a, item_b], [
75
               {"b", 1},
76
               {"a", 1},
77
               {"b", 2},
78
               {"a", 2}
79
             ])
80
    end
81
82
    test "members of one stack never reorder" do
83
      {:ok, item_a} =
84
        MergeQueue.item(stack(id: "a"), entries(["1111", "2222"]), @queue_base, @policy)
85
86
      {:ok, item_b} = MergeQueue.item(stack(id: "b"), entries(["3333"]), @queue_base, @policy)
87
88
      refute MergeQueue.valid_order?([item_a, item_b], [
89
               {"a", 2},
90
               {"b", 1},
91
               {"a", 1}
92
             ])
93
    end
94
95
    test "the proposed order must contain exactly the selected members" do
96
      {:ok, item} =
97
        MergeQueue.item(stack(id: "a"), entries(["1111", "2222"]), @queue_base, @policy)
98
99
      refute MergeQueue.valid_order?([item], [{"a", 1}])
100
      refute MergeQueue.valid_order?([item], [{"a", 1}, {"a", 2}, {"a", 3}])
101
    end
102
  end
103
104
  describe "eject/2" do
105
    test "ejecting a lower member ejects every selected member above it" do
106
      {:ok, item} =
107
        MergeQueue.item(
108
          stack(version: 1),
109
          entries(["1111", "2222", "3333"]),
110
          @queue_base,
111
          @policy
112
        )
113
114
      {:ok, %{ejected: ejected, remaining: remaining}} = MergeQueue.eject(item, 2)
115
116
      assert Enum.map(ejected, & &1.position) == [2, 3]
117
      assert Enum.map(remaining, & &1.position) == [1]
118
    end
119
120
    test "ejecting the bottom member empties the group" do
121
      {:ok, item} =
122
        MergeQueue.item(stack(version: 1), entries(["1111", "2222"]), @queue_base, @policy)
123
124
      {:ok, %{ejected: ejected, remaining: []}} = MergeQueue.eject(item, 1)
125
      assert Enum.map(ejected, & &1.position) == [1, 2]
126
    end
127
128
    test "an unselected position is not a member" do
129
      {:ok, item} = MergeQueue.item(stack(version: 1), entries(["1111"]), @queue_base, @policy)
130
131
      assert {:error, :not_a_member} = MergeQueue.eject(item, 2)
132
    end
133
  end
134
135
  describe "invalidations/2" do
136
    test "a matching observation keeps the speculative result valid" do
137
      {:ok, item} =
138
        MergeQueue.item(stack(version: 3), entries(["1111", "2222"]), @queue_base, @policy)
139
140
      observed = observed(item)
141
142
      assert MergeQueue.invalidations(item, observed) == []
143
      assert MergeQueue.valid?(item, observed)
144
    end
145
146
    test "a queue base advance invalidates" do
147
      {:ok, item} = MergeQueue.item(stack(version: 1), entries(["1111"]), @queue_base, @policy)
148
149
      observed = %{observed(item) | queue_base_oid: oid("9999")}
150
151
      assert MergeQueue.invalidations(item, observed) == [:queue_base_changed]
152
      refute MergeQueue.valid?(item, observed)
153
    end
154
155
    test "a stack version move invalidates" do
156
      {:ok, item} = MergeQueue.item(stack(version: 1), entries(["1111"]), @queue_base, @policy)
157
158
      observed = %{observed(item) | stack_version: 2}
159
160
      assert MergeQueue.invalidations(item, observed) == [:stack_version_changed]
161
    end
162
163
    test "a policy version change invalidates" do
164
      {:ok, item} = MergeQueue.item(stack(version: 1), entries(["1111"]), @queue_base, @policy)
165
166
      observed = %{observed(item) | policy_version: oid("8888")}
167
168
      assert MergeQueue.invalidations(item, observed) == [:policy_changed]
169
    end
170
171
    test "any selected head move invalidates that member" do
172
      {:ok, item} =
173
        MergeQueue.item(stack(version: 1), entries(["1111", "2222"]), @queue_base, @policy)
174
175
      observed =
176
        Map.update!(observed(item), :head_oids, &Map.put(&1, 2, oid("7777")))
177
178
      assert MergeQueue.invalidations(item, observed) == [{:head_changed, 2}]
179
    end
180
181
    test "every trigger reports together" do
182
      {:ok, item} = MergeQueue.item(stack(version: 1), entries(["1111"]), @queue_base, @policy)
183
184
      observed = %{
185
        queue_base_oid: oid("9999"),
186
        stack_version: 5,
187
        policy_version: oid("8888"),
188
        head_oids: %{1 => oid("7777")}
189
      }
190
191
      assert MergeQueue.invalidations(item, observed) == [
192
               :queue_base_changed,
193
               :stack_version_changed,
194
               :policy_changed,
195
               {:head_changed, 1}
196
             ]
197
    end
198
  end
199
200
  describe "the queue item shape" do
201
    test "is the documented logical item" do
202
      {:ok, item} = MergeQueue.item(stack(version: 1), entries(["1111"]), @queue_base, @policy)
203
204
      assert %QueueItem{} = item
205
    end
206
  end
207
208
  ## Fakes
209
210
  defp stack(attrs) do
211
    %Stack{
212
      id: Keyword.get(attrs, :id, "stack-id"),
213
      number: 7,
214
      trunk_ref: "main",
215
      version: Keyword.get(attrs, :version, 1)
216
    }
217
  end
218
219
  defp entries(seeds) do
220
    seeds
221
    |> Enum.with_index(1)
222
    |> Enum.map(fn {seed, position} ->
223
      %StackEntry{
224
        position: position,
225
        pull_request_id: "pr-#{position}",
226
        boundary_oid: oid("0000"),
227
        observed_head_oid: oid(seed)
228
      }
229
    end)
230
  end
231
232
  defp observed(item) do
233
    %{
234
      queue_base_oid: item.queue_base_oid,
235
      stack_version: item.stack_version,
236
      policy_version: item.policy_version,
237
      head_oids: Map.new(item.members, &{&1.position, &1.expected_head_oid})
238
    }
239
  end
240
241
  defp oid(seed), do: String.duplicate(seed, div(40, byte_size(seed)))
242
end

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