Evaluate stack policy against the effective base and emit stack events

efc10d0835e2 · Devin AI · · parent 3323ac3bc986

Evaluate stack policy against the effective base and emit stack events

Every pull request policy evaluation now carries two first-class bases:
the direct base that drives diffs, layer boundaries, and display, and the
effective base — the stack trunk — that drives every policy decision.
OpenAgents.Stacks.Policy resolves configuration such as CODEOWNERS from
the effective base tree, so an unmerged lower layer editing CODEOWNERS
or workflows cannot weaken the rules an upper layer is held to; once the
layer lands on the trunk, later evaluations legitimately see the new
configuration.

Pull request payloads embed stack membership: the stack number, the
layer's position, the active size, the stack health, and the effective
base ref with its live OID. Unstacked pull requests carry an explicit
null.

The event catalog now covers the full lifecycle from docs/stacked-prs.md
section 16: stack created, appended, restructured, dissolved, rebase
started/conflicted/completed, merge started/queued/partially-completed/
completed/failed, and per-pull-request stacked/unstacked/position-changed
events, each carrying the stack number, operation ID where one exists,
version, actor, orderings, and affected head OIDs. Events persist through
the transactional outbox as before; the new
OpenAgents.Stacks.EventDispatcher delivers undelivered rows in insertion
order on a per-repository PubSub topic and marks them delivered, with
at-least-once semantics so consumers deduplicate by event ID.

Also classifies the production deploy runbook's provider resource names
in the Sarah reference allowlist; the runbook landed upstream without an
entry and failed ops/ci/reference-check.sh.

Closes #52

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

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/supervisor.ex
  • modified lib/openagents/stacks.ex
  • added lib/openagents/stacks/event_dispatcher.ex
  • modified lib/openagents/stacks/merge.ex
  • added lib/openagents/stacks/policy.ex
  • modified lib/openagents/stacks/restack.ex
  • modified lib/openagents/stacks/stack_event.ex
  • modified lib/openagents_web/controllers/pull_request_controller.ex
  • modified lib/openagents_web/controllers/pull_request_json.ex
  • modified priv/migration_lineages/prior-2026-08-19.json
  • added priv/repo/migrations/20260823120247_add_stack_event_delivery.exs
  • added test/openagents/stacks/event_dispatcher_test.exs
  • modified test/openagents/stacks/merge_test.exs
  • added test/openagents/stacks/policy_test.exs
  • modified test/openagents/stacks/restack_test.exs
  • modified test/openagents_web/controllers/stack_controller_test.exs

Diff

16 files changed, +889 -22

lib/openagents/forge/supervisor.ex modified +2 -1

@@ -23,7 +23,8 @@ defmodule OpenAgents.Forge.Supervisor do

23 23
      [
24 24
        {OpenAgents.Repositories.Provisioner, []},
25 25
        {OpenAgents.Repositories.ImportWorkspaceJanitor, []},
26
        {OpenAgents.Stacks.OperationWorker, []}
26
        {OpenAgents.Stacks.OperationWorker, []},
27
        {OpenAgents.Stacks.EventDispatcher, []}
27 28
      ]
28 29
    else
29 30
      []
lib/openagents/stacks.ex modified +81 -2

@@ -152,6 +152,63 @@ defmodule OpenAgents.Stacks do

152 152
    end
153 153
  end
154 154
155
  @doc "Loads the stack an active entry belongs to, with active entries."
156
  def get_stack_for_entry!(%StackEntry{stack_id: stack_id}) do
157
    Stack
158
    |> Repo.get!(stack_id)
159
    |> Repo.preload(entries: active_entries_query())
160
  end
161
162
  @doc """
163
  Stack payload contexts for pull requests, keyed by pull request ID.
164
165
  Each context carries the stack number, the entry's position, the active
166
  size, the stack health, and the effective base — the trunk ref with its
167
  live OID — in the shape ordinary pull request payloads embed
168
  (docs/stacked-prs.md section 16). Unstacked pull requests have no key.
169
  """
170
  def payload_contexts(%Repository{} = repository, pull_requests) when is_list(pull_requests) do
171
    ids = Enum.map(pull_requests, & &1.id)
172
173
    entries =
174
      Repo.all(
175
        from entry in StackEntry,
176
          where: entry.pull_request_id in ^ids and is_nil(entry.removed_at),
177
          preload: [:stack]
178
      )
179
180
    stack_ids = entries |> Enum.map(& &1.stack_id) |> Enum.uniq()
181
182
    sizes =
183
      Map.new(
184
        Repo.all(
185
          from entry in StackEntry,
186
            where: entry.stack_id in ^stack_ids and is_nil(entry.removed_at),
187
            group_by: entry.stack_id,
188
            select: {entry.stack_id, count(entry.id)}
189
        )
190
      )
191
192
    trunk_oids =
193
      Map.new(Enum.uniq_by(entries, & &1.stack_id), fn entry ->
194
        case Browse.resolve_commit(repository, entry.stack.trunk_ref) do
195
          {:ok, oid} -> {entry.stack_id, oid}
196
          _other -> {entry.stack_id, nil}
197
        end
198
      end)
199
200
    Map.new(entries, fn entry ->
201
      {entry.pull_request_id,
202
       %{
203
         number: entry.stack.number,
204
         position: entry.position,
205
         size: Map.fetch!(sizes, entry.stack_id),
206
         health: entry.stack.health,
207
         base: %{ref: entry.stack.trunk_ref, sha: Map.fetch!(trunk_oids, entry.stack_id)}
208
       }}
209
    end)
210
  end
211
155 212
  def get_by_number!(%Repository{id: repository_id}, number) when is_integer(number) do
156 213
    Stack
157 214
    |> Repo.get_by!(repository_id: repository_id, number: number)

@@ -421,12 +478,24 @@ defmodule OpenAgents.Stacks do

421 478
      entries = insert_entry_rows!(stack, entry_specs(pull_requests, heads, trunk_oid))
422 479
423 480
      record_event!(stack, "pull_request_stack.created", actor, %{
424
        "number" => stack.number,
481
        "stack_number" => stack.number,
425 482
        "trunk_ref" => stack.trunk_ref,
426 483
        "trunk_oid" => trunk_oid,
484
        "ordering_old" => [],
485
        "ordering_new" => Enum.map(entries, & &1.pull_request.issue.number),
427 486
        "entries" => Enum.map(entries, &event_entry/1)
428 487
      })
429 488
489
      Enum.each(entries, fn entry ->
490
        record_event!(stack, "pull_request.stacked", actor, %{
491
          "stack_number" => stack.number,
492
          "trunk_ref" => stack.trunk_ref,
493
          "pull_request" => entry.pull_request.issue.number,
494
          "position" => entry.position,
495
          "head_oid" => entry.observed_head_oid
496
        })
497
      end)
498
430 499
      record_idempotency!(actor, "stack_create", idempotency_key, request_digest, stack.id)
431 500
      {%{stack | entries: entries}, :created}
432 501
    else

@@ -448,11 +517,21 @@ defmodule OpenAgents.Stacks do

448 517
        insert_entry_row!(stack, pull_request, top.position + 1, top.observed_head_oid, head_oid)
449 518
450 519
      record_event!(stack, "pull_request_stack.appended", actor, %{
451
        "number" => stack.number,
520
        "stack_number" => stack.number,
452 521
        "trunk_ref" => stack.trunk_ref,
522
        "ordering_old" => Enum.map(entries, & &1.pull_request.issue.number),
523
        "ordering_new" => Enum.map(entries ++ [entry], & &1.pull_request.issue.number),
453 524
        "entries" => [event_entry(entry)]
454 525
      })
455 526
527
      record_event!(stack, "pull_request.stacked", actor, %{
528
        "stack_number" => stack.number,
529
        "trunk_ref" => stack.trunk_ref,
530
        "pull_request" => pull_request.issue.number,
531
        "position" => entry.position,
532
        "head_oid" => entry.observed_head_oid
533
      })
534
456 535
      record_idempotency!(actor, "stack_append", idempotency_key, request_digest, stack.id)
457 536
      {%{stack | entries: entries ++ [entry]}, :created}
458 537
    else
lib/openagents/stacks/event_dispatcher.ex added +126

@@ -0,0 +1,126 @@

1
defmodule OpenAgents.Stacks.EventDispatcher do
2
  @moduledoc """
3
  Delivers stack outbox events to PubSub subscribers.
4
5
  Stack mutations write `OpenAgents.Stacks.StackEvent` rows inside the
6
  same metadata transaction (the transactional outbox). This worker polls
7
  for undelivered rows, broadcasts each on the owning repository's stack
8
  event topic in insertion order, and marks it delivered in the same
9
  transaction. Claiming uses `FOR UPDATE SKIP LOCKED`, so multiple nodes
10
  never race on one row — but delivery is still at-least-once (a crash
11
  between broadcast and commit redelivers), so consumers deduplicate by
12
  event ID.
13
  """
14
  use GenServer
15
16
  import Ecto.Query
17
18
  alias OpenAgents.Repo
19
  alias OpenAgents.Stacks.Stack
20
  alias OpenAgents.Stacks.StackEvent
21
22
  @batch_size 100
23
24
  def start_link(options) do
25
    name = Keyword.get(options, :name, __MODULE__)
26
    gen_server_options = if name, do: [name: name], else: []
27
    GenServer.start_link(__MODULE__, options, gen_server_options)
28
  end
29
30
  @doc "Subscribes the caller to a repository's stack events."
31
  def subscribe(repository_id) do
32
    Phoenix.PubSub.subscribe(OpenAgents.PubSub, topic(repository_id))
33
  end
34
35
  @doc "The stack event topic for one repository."
36
  def topic(repository_id), do: "stack_events:#{repository_id}"
37
38
  @doc "Synchronously delivers every pending outbox row."
39
  def drain(server \\ __MODULE__), do: GenServer.call(server, :drain, 30_000)
40
41
  @doc """
42
  Delivers one batch of undelivered events in insertion order.
43
44
  Returns the number of rows delivered. Callers that need everything
45
  flushed call it until it returns `0`.
46
  """
47
  def deliver_pending do
48
    {:ok, count} =
49
      Repo.transaction(fn ->
50
        events =
51
          Repo.all(
52
            from event in StackEvent,
53
              join: stack in Stack,
54
              on: stack.id == event.stack_id,
55
              where: is_nil(event.delivered_at),
56
              order_by: [asc: event.inserted_at, asc: event.id],
57
              limit: @batch_size,
58
              lock: fragment("FOR UPDATE OF ? SKIP LOCKED", event),
59
              select: {event, stack.repository_id}
60
          )
61
62
        now = DateTime.utc_now()
63
64
        Enum.each(events, fn {event, repository_id} ->
65
          broadcast(event, repository_id)
66
67
          event
68
          |> StackEvent.delivered_changeset(now)
69
          |> Repo.update!()
70
        end)
71
72
        length(events)
73
      end)
74
75
    count
76
  end
77
78
  defp broadcast(event, repository_id) do
79
    Phoenix.PubSub.broadcast(
80
      OpenAgents.PubSub,
81
      topic(repository_id),
82
      {:stack_event,
83
       %{
84
         id: event.id,
85
         stack_id: event.stack_id,
86
         event_type: event.event_type,
87
         stack_version: event.stack_version,
88
         actor_user_id: event.actor_user_id,
89
         payload: event.payload,
90
         inserted_at: event.inserted_at
91
       }}
92
    )
93
  end
94
95
  @impl true
96
  def init(options) do
97
    state = %{poll_interval_ms: Keyword.get(options, :poll_interval_ms, poll_interval_ms())}
98
    schedule(state.poll_interval_ms)
99
    {:ok, state}
100
  end
101
102
  @impl true
103
  def handle_call(:drain, _from, state) do
104
    {:reply, {:ok, drain_now(0)}, state}
105
  end
106
107
  @impl true
108
  def handle_info(:poll, state) do
109
    _count = deliver_pending()
110
    schedule(state.poll_interval_ms)
111
    {:noreply, state}
112
  end
113
114
  defp drain_now(total) do
115
    case deliver_pending() do
116
      0 -> total
117
      count -> drain_now(total + count)
118
    end
119
  end
120
121
  defp schedule(interval), do: Process.send_after(self(), :poll, interval)
122
123
  defp poll_interval_ms do
124
    Application.get_env(:openagents, :stack_event_dispatcher_poll_interval_ms, 1_000)
125
  end
126
end
lib/openagents/stacks/merge.ex modified +52 -10

@@ -71,6 +71,16 @@ defmodule OpenAgents.Stacks.Merge do

71 71
            nil ->
72 72
              operation = insert_operation!(stack, actor, key, request, position)
73 73
              set_health!(stack, "operation_in_progress")
74
75
              record_event!(stack, operation, "pull_request_stack.merge_started", %{
76
                "stack_number" => stack.number,
77
                "operation_id" => operation.id,
78
                "trunk_ref" => stack.trunk_ref,
79
                "merge_method" => Map.fetch!(request, "merge_method"),
80
                "pull_request" => Map.fetch!(request, "pull_request_number"),
81
                "target_position" => position
82
              })
83
74 84
              {operation, :created}
75 85
          end
76 86
        else

@@ -685,14 +695,28 @@ defmodule OpenAgents.Stacks.Merge do

685 695
  end
686 696
687 697
  defp record_events!(stack, operation, plan) do
688
    record_event!(stack, operation, "pull_request_stack.merged", %{
698
    record_event!(stack, operation, "pull_request_stack.merge_completed", %{
699
      "stack_number" => stack.number,
689 700
      "operation_id" => operation.id,
701
      "trunk_ref" => stack.trunk_ref,
690 702
      "merge_method" => Map.fetch!(plan, "merge_method"),
691 703
      "trunk_old" => Map.fetch!(plan, "trunk_old"),
692 704
      "trunk_new" => Map.fetch!(plan, "trunk_new"),
705
      "pull_requests" =>
706
        Enum.map(Map.fetch!(plan, "merged"), &Map.fetch!(&1, "pull_request_number")),
693 707
      "merged" => Map.fetch!(plan, "merged")
694 708
    })
695 709
710
    Enum.each(Map.fetch!(plan, "merged"), fn step ->
711
      record_event!(stack, operation, "pull_request.unstacked", %{
712
        "stack_number" => stack.number,
713
        "operation_id" => operation.id,
714
        "pull_request" => Map.fetch!(step, "pull_request_number"),
715
        "reason" => "merged",
716
        "merge_commit_sha" => Map.fetch!(step, "merge_commit_sha")
717
      })
718
    end)
719
696 720
    Enum.each(Map.fetch!(plan, "restacked"), fn step ->
697 721
      record_event!(stack, operation, "pull_request.synchronize", %{
698 722
        "operation_id" => operation.id,

@@ -722,15 +746,26 @@ defmodule OpenAgents.Stacks.Merge do

722 746
  # requests' code is on the trunk, so the operation records exactly what
723 747
  # applied instead of pretending nothing happened.
724 748
  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!()
749
    {:ok, operation} =
750
      Repo.transaction(fn ->
751
        stack = Repo.one!(from s in Stack, where: s.id == ^operation.stack_id, lock: "FOR UPDATE")
752
753
        record_event!(stack, operation, "pull_request_stack.merge_partially_completed", %{
754
          "stack_number" => stack.number,
755
          "operation_id" => operation.id,
756
          "trunk_ref" => stack.trunk_ref,
757
          "error" => error_map(reason)
758
        })
759
760
        operation
761
        |> Operation.transition_changeset(%{
762
          state: "partially_succeeded",
763
          error: error_map(reason),
764
          planned_result: Map.put(plan, "refs_applied", true),
765
          completed_at: DateTime.utc_now()
766
        })
767
        |> Repo.update!()
768
      end)
734 769
735 770
    {:error, operation}
736 771
  end

@@ -741,6 +776,13 @@ defmodule OpenAgents.Stacks.Merge do

741 776
        current = Repo.one!(from s in Stack, where: s.id == ^stack.id, lock: "FOR UPDATE")
742 777
        set_health!(current, failure_health(reason))
743 778
779
        record_event!(current, operation, "pull_request_stack.merge_failed", %{
780
          "stack_number" => current.number,
781
          "operation_id" => operation.id,
782
          "trunk_ref" => current.trunk_ref,
783
          "error" => error_map(reason)
784
        })
785
744 786
        operation
745 787
        |> Operation.transition_changeset(%{
746 788
          state: "failed",
lib/openagents/stacks/policy.ex added +117

@@ -0,0 +1,117 @@

1
defmodule OpenAgents.Stacks.Policy do
2
  @moduledoc """
3
  Policy evaluation bases for pull requests (docs/stacked-prs.md section 11).
4
5
  Every evaluation carries two bases as first-class fields. The direct base
6
  is the pull request's own base branch and drives diff presentation, layer
7
  boundaries, and base display. The effective base is the stack trunk (or
8
  the direct base for an unstacked pull request) and drives every policy
9
  decision: protection rules, required approvals, CODEOWNERS, and merge
10
  configuration.
11
12
  Policy configuration resolves from the effective base tree, never from
13
  the direct parent branch tree, so an unmerged lower layer editing
14
  `CODEOWNERS` or a workflow cannot weaken the rules an upper layer is
15
  held to. Once the lower layer lands on the trunk, the effective base OID
16
  advances and later evaluations legitimately see the new configuration.
17
  """
18
19
  alias OpenAgents.Forge.Browse
20
  alias OpenAgents.PullRequests.PullRequest
21
  alias OpenAgents.Repositories.Repository
22
  alias OpenAgents.Stacks
23
  alias OpenAgents.Stacks.StackEntry
24
25
  @codeowners_paths [".github/CODEOWNERS", "CODEOWNERS", "docs/CODEOWNERS"]
26
27
  @doc """
28
  The policy evaluation context for one pull request.
29
30
  Returns `{:ok, evaluation}` where the evaluation carries the head OID,
31
  the direct base (ref and live OID), the effective base (ref and live
32
  OID), and the stack context (`nil` for an unstacked pull request). Both
33
  base OIDs resolve live, so the evaluation always reflects the current
34
  refs.
35
  """
36
  def evaluation(%Repository{} = repository, %PullRequest{} = pull_request) do
37
    case Stacks.active_entry_for_pull_request(pull_request) do
38
      nil -> unstacked_evaluation(repository, pull_request)
39
      %StackEntry{} = entry -> stacked_evaluation(repository, pull_request, entry)
40
    end
41
  end
42
43
  @doc """
44
  A configuration blob resolved from the evaluation's effective base tree.
45
46
  Reads `path` at the effective base OID — never at the direct parent
47
  branch — and returns `{:ok, %{content, truncated, binary, size}}` or
48
  `{:error, :not_found}`.
49
  """
50
  def configuration_blob(%Repository{} = repository, evaluation, path)
51
      when is_binary(path) do
52
    Browse.blob(repository, evaluation.effective_base.oid, path)
53
  end
54
55
  @doc """
56
  The CODEOWNERS content governing this evaluation.
57
58
  Searches `.github/CODEOWNERS`, `CODEOWNERS`, then `docs/CODEOWNERS` in
59
  the effective base tree and returns
60
  `{:ok, %{path, content, source_oid}}` or `{:error, :not_found}`.
61
  """
62
  def codeowners(%Repository{} = repository, evaluation) do
63
    Enum.find_value(@codeowners_paths, {:error, :not_found}, fn path ->
64
      case configuration_blob(repository, evaluation, path) do
65
        {:ok, blob} ->
66
          {:ok, %{path: path, content: blob.content, source_oid: evaluation.effective_base.oid}}
67
68
        _other ->
69
          nil
70
      end
71
    end)
72
  end
73
74
  defp unstacked_evaluation(repository, pull_request) do
75
    with {:ok, base_oid} <- resolve(repository, pull_request.base_ref) do
76
      base = %{ref: pull_request.base_ref, oid: base_oid}
77
78
      {:ok,
79
       %{
80
         pull_request_id: pull_request.id,
81
         head_oid: pull_request.head_sha,
82
         direct_base: base,
83
         effective_base: base,
84
         stack: nil
85
       }}
86
    end
87
  end
88
89
  defp stacked_evaluation(repository, pull_request, entry) do
90
    stack = Stacks.get_stack_for_entry!(entry)
91
92
    with {:ok, direct_oid} <- resolve(repository, pull_request.base_ref),
93
         {:ok, trunk_oid} <- resolve(repository, stack.trunk_ref) do
94
      {:ok,
95
       %{
96
         pull_request_id: pull_request.id,
97
         head_oid: pull_request.head_sha,
98
         direct_base: %{ref: pull_request.base_ref, oid: direct_oid},
99
         effective_base: %{ref: stack.trunk_ref, oid: trunk_oid},
100
         stack: %{
101
           id: stack.id,
102
           number: stack.number,
103
           position: entry.position,
104
           size: length(stack.entries),
105
           health: stack.health
106
         }
107
       }}
108
    end
109
  end
110
111
  defp resolve(repository, ref) do
112
    case Browse.resolve_commit(repository, ref) do
113
      {:ok, oid} -> {:ok, oid}
114
      _other -> {:error, {:missing_ref, ref}}
115
    end
116
  end
117
end
lib/openagents/stacks/restack.ex modified +20 -1

@@ -59,6 +59,13 @@ defmodule OpenAgents.Stacks.Restack do

59 59
            nil ->
60 60
              operation = insert_operation!(stack, actor, key, request)
61 61
              set_health!(stack, "operation_in_progress")
62
63
              record_event!(stack, operation, "pull_request_stack.rebase_started", %{
64
                "stack_number" => stack.number,
65
                "operation_id" => operation.id,
66
                "trunk_ref" => stack.trunk_ref
67
              })
68
62 69
              {operation, :created}
63 70
          end
64 71
        else

@@ -529,9 +536,12 @@ defmodule OpenAgents.Stacks.Restack do

529 536
  defp record_events!(stack, operation, trunk_tip, steps) do
530 537
    changed = Enum.filter(steps, &(&1.new_head != &1.old_head))
531 538
532
    record_event!(stack, operation, "pull_request_stack.rebased", %{
539
    record_event!(stack, operation, "pull_request_stack.rebase_completed", %{
540
      "stack_number" => stack.number,
533 541
      "operation_id" => operation.id,
542
      "trunk_ref" => stack.trunk_ref,
534 543
      "trunk_oid" => trunk_tip,
544
      "pull_requests" => Enum.map(steps, & &1.pull_request_number),
535 545
      "steps" => Enum.map(steps, &stringify_step/1)
536 546
    })
537 547

@@ -564,6 +574,15 @@ defmodule OpenAgents.Stacks.Restack do

564 574
        stack = Repo.one!(from s in Stack, where: s.id == ^operation.stack_id, lock: "FOR UPDATE")
565 575
        set_health!(stack, "conflicted")
566 576
577
        record_event!(stack, operation, "pull_request_stack.rebase_conflicted", %{
578
          "stack_number" => stack.number,
579
          "operation_id" => operation.id,
580
          "trunk_ref" => stack.trunk_ref,
581
          "pull_request" => Map.fetch!(conflict, "pull_request_number"),
582
          "position" => Map.fetch!(conflict, "position"),
583
          "paths" => Map.fetch!(conflict, "paths")
584
        })
585
567 586
        operation
568 587
        |> Operation.transition_changeset(%{
569 588
          state: "waiting_for_conflict_resolution",
lib/openagents/stacks/stack_event.ex modified +27 -1

@@ -13,7 +13,24 @@ 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_stack.merged pull_request.synchronize)
16
  @event_types ~w(
17
    pull_request_stack.created
18
    pull_request_stack.appended
19
    pull_request_stack.restructured
20
    pull_request_stack.dissolved
21
    pull_request_stack.rebase_started
22
    pull_request_stack.rebase_conflicted
23
    pull_request_stack.rebase_completed
24
    pull_request_stack.merge_started
25
    pull_request_stack.merge_queued
26
    pull_request_stack.merge_partially_completed
27
    pull_request_stack.merge_completed
28
    pull_request_stack.merge_failed
29
    pull_request.stacked
30
    pull_request.unstacked
31
    pull_request.stack_position_changed
32
    pull_request.synchronize
33
  )
17 34
18 35
  schema "pull_request_stack_events" do
19 36
    belongs_to :stack, OpenAgents.Stacks.Stack

@@ -21,9 +38,13 @@ defmodule OpenAgents.Stacks.StackEvent do

21 38
    field :stack_version, :integer
22 39
    belongs_to :actor_user, OpenAgents.Accounts.User
23 40
    field :payload, :map, default: %{}
41
    field :delivered_at, :utc_datetime_usec
24 42
    timestamps()
25 43
  end
26 44
45
  @doc "The full event type catalog (docs/stacked-prs.md section 16)."
46
  def event_types, do: @event_types
47
27 48
  def changeset(event, attrs) do
28 49
    event
29 50
    |> cast(attrs, [:event_type, :stack_version, :payload])

@@ -36,4 +57,9 @@ defmodule OpenAgents.Stacks.StackEvent do

36 57
    |> foreign_key_constraint(:stack_id)
37 58
    |> foreign_key_constraint(:actor_user_id)
38 59
  end
60
61
  @doc "Marks one outbox row delivered."
62
  def delivered_changeset(event, delivered_at) do
63
    change(event, delivered_at: delivered_at)
64
  end
39 65
end
lib/openagents_web/controllers/pull_request_controller.ex modified +15 -2

@@ -3,10 +3,18 @@ defmodule OpenAgentsWeb.PullRequestController do

3 3
4 4
  alias OpenAgents.PullRequests
5 5
  alias OpenAgents.Repositories
6
  alias OpenAgents.Stacks
6 7
7 8
  def index(conn, %{"owner" => owner, "repo" => repo}) do
8 9
    repository = Repositories.get_visible_by_path!(owner, repo, conn.assigns[:current_user])
9
    render(conn, :index, pull_requests: PullRequests.list(repository), owner: owner, repo: repo)
10
    pull_requests = PullRequests.list(repository)
11
12
    render(conn, :index,
13
      pull_requests: pull_requests,
14
      owner: owner,
15
      repo: repo,
16
      stack_contexts: Stacks.payload_contexts(repository, pull_requests)
17
    )
10 18
  rescue
11 19
    Ecto.NoResultsError -> not_found(conn)
12 20
  end

@@ -20,7 +28,12 @@ defmodule OpenAgentsWeb.PullRequestController do

20 28
        OpenAgentsWeb.ControllerHelpers.integer_param!(number)
21 29
      )
22 30
23
    render(conn, :show, pull_request: pull_request, owner: owner, repo: repo)
31
    render(conn, :show,
32
      pull_request: pull_request,
33
      owner: owner,
34
      repo: repo,
35
      stack_contexts: Stacks.payload_contexts(repository, [pull_request])
36
    )
24 37
  rescue
25 38
    Ecto.NoResultsError -> not_found(conn)
26 39
  end
lib/openagents_web/controllers/pull_request_json.ex modified +10

@@ -27,10 +27,20 @@ defmodule OpenAgentsWeb.PullRequestJSON do

27 27
        repo: %{full_name: "#{pr.head_repository.owner}/#{pr.head_repository.name}"}
28 28
      },
29 29
      base: %{ref: pr.base_ref, sha: pr.base_sha, repo: %{full_name: "#{owner}/#{repo}"}},
30
      stack: stack_context(pr, assigns),
30 31
      created_at: pr.inserted_at,
31 32
      updated_at: pr.updated_at,
32 33
      html_url: "#{base_url}/#{owner}/#{repo}/pulls/#{pr.issue.number}",
33 34
      url: "#{base_url}/api/v3/repos/#{owner}/#{repo}/pulls/#{pr.issue.number}"
34 35
    }
35 36
  end
37
38
  # Stack membership: the stack number, this layer's position, the active
39
  # size, the stack health, and the effective base — the stack trunk that
40
  # governs policy — or `nil` for an unstacked pull request.
41
  defp stack_context(pr, assigns) do
42
    assigns
43
    |> Map.get(:stack_contexts, %{})
44
    |> Map.get(pr.id)
45
  end
36 46
end
priv/migration_lineages/prior-2026-08-19.json modified +2 -1

@@ -250,7 +250,8 @@

250 250
    20260823071500,
251 251
    20260823072000,
252 252
    20260823073000,
253
    20260823074000
253
    20260823074000,
254
    20260823120247
254 255
  ],
255 256
  "required_tables": [
256 257
    "users",
priv/repo/migrations/20260823120247_add_stack_event_delivery.exs added +14

@@ -0,0 +1,14 @@

1
defmodule OpenAgents.Repo.Migrations.AddStackEventDelivery do
2
  use Ecto.Migration
3
4
  def change do
5
    alter table(:pull_request_stack_events) do
6
      add :delivered_at, :utc_datetime_usec
7
    end
8
9
    create index(:pull_request_stack_events, [:inserted_at],
10
             where: "delivered_at IS NULL",
11
             name: :pull_request_stack_events_undelivered_index
12
           )
13
  end
14
end
test/openagents/stacks/event_dispatcher_test.exs added +143

@@ -0,0 +1,143 @@

1
defmodule OpenAgents.Stacks.EventDispatcherTest do
2
  @moduledoc """
3
  Outbox delivery (#52): undelivered stack events broadcast in insertion
4
  order on the repository's topic, delivery marks the row so a drained
5
  outbox stays quiet, and redelivery reuses the same event ID so consumers
6
  deduplicate.
7
  """
8
9
  use OpenAgents.DataCase, async: false
10
11
  alias OpenAgents.PullRequests.PullRequest
12
  alias OpenAgents.Repo
13
  alias OpenAgents.Stacks
14
  alias OpenAgents.Stacks.EventDispatcher
15
  alias OpenAgents.Stacks.StackEvent
16
17
  import Ecto.Query
18
  import OpenAgents.AccountsFixtures
19
  import OpenAgents.IssuesFixtures
20
21
  test "delivers pending events in order, marks them, and redelivers with the same ID" do
22
    actor = repository_user_fixture("dispatch-actor")
23
    repository = repository_with_member_fixture(actor)
24
    stack = seed_stack(repository, actor)
25
    record_events(stack, actor)
26
27
    EventDispatcher.subscribe(repository.id)
28
29
    assert EventDispatcher.deliver_pending() == 3
30
31
    assert_receive {:stack_event, %{event_type: "pull_request_stack.created"} = created}
32
    assert_receive {:stack_event, %{event_type: "pull_request.stacked"} = stacked_1}
33
    assert_receive {:stack_event, %{event_type: "pull_request.stacked"} = stacked_2}
34
    refute_receive {:stack_event, _event}
35
36
    for event <- [created, stacked_1, stacked_2] do
37
      assert event.stack_id == stack.id
38
      assert event.stack_version == 1
39
      assert event.actor_user_id == actor.id
40
      assert is_binary(event.id)
41
    end
42
43
    assert created.payload["stack_number"] == stack.number
44
    assert created.payload["ordering_old"] == []
45
    assert created.payload["ordering_new"] == [101, 102]
46
47
    # The outbox is drained: nothing is pending and nothing rebroadcasts.
48
    assert EventDispatcher.deliver_pending() == 0
49
    refute_receive {:stack_event, _event}
50
51
    refute Repo.exists?(from event in StackEvent, where: is_nil(event.delivered_at))
52
53
    # A redelivery (crash between broadcast and commit) reuses the event ID,
54
    # so consumers deduplicate.
55
    created_id = created.id
56
57
    Repo.get!(StackEvent, created_id)
58
    |> StackEvent.delivered_changeset(nil)
59
    |> Repo.update!()
60
61
    assert EventDispatcher.deliver_pending() == 1
62
    assert_receive {:stack_event, %{id: ^created_id}}
63
  end
64
65
  test "the worker drains the outbox on demand" do
66
    actor = repository_user_fixture("dispatch-worker")
67
    repository = repository_with_member_fixture(actor)
68
    stack = seed_stack(repository, actor)
69
    record_events(stack, actor)
70
71
    EventDispatcher.subscribe(repository.id)
72
73
    dispatcher =
74
      start_supervised!({EventDispatcher, name: nil, poll_interval_ms: 3_600_000})
75
76
    assert {:ok, 3} = EventDispatcher.drain(dispatcher)
77
    assert_receive {:stack_event, %{event_type: "pull_request_stack.created"}}
78
    assert {:ok, 0} = EventDispatcher.drain(dispatcher)
79
  end
80
81
  defp seed_stack(repository, actor) do
82
    pr_1 = pull_request(repository, "layer-1", "main")
83
    pr_2 = pull_request(repository, "layer-2", "layer-1")
84
85
    {:ok, stack} = Stacks.create(repository, [pr_1, pr_2], actor)
86
    stack
87
  end
88
89
  # Writes the outbox rows a stack creation records, in insertion order.
90
  defp record_events(stack, actor) do
91
    created = %{
92
      "stack_number" => stack.number,
93
      "trunk_ref" => stack.trunk_ref,
94
      "ordering_old" => [],
95
      "ordering_new" => [101, 102]
96
    }
97
98
    stacked_1 = %{"stack_number" => stack.number, "pull_request" => 101, "position" => 1}
99
    stacked_2 = %{"stack_number" => stack.number, "pull_request" => 102, "position" => 2}
100
101
    for {event_type, payload} <- [
102
          {"pull_request_stack.created", created},
103
          {"pull_request.stacked", stacked_1},
104
          {"pull_request.stacked", stacked_2}
105
        ] do
106
      %StackEvent{}
107
      |> StackEvent.changeset(%{
108
        stack_id: stack.id,
109
        actor_user_id: actor.id,
110
        event_type: event_type,
111
        stack_version: 1,
112
        payload: payload
113
      })
114
      |> Repo.insert!()
115
    end
116
  end
117
118
  defp pull_request(repository, head_ref, base_ref) do
119
    issue = issue_fixture(repository, %{title: "PR #{head_ref}"})
120
121
    {:ok, pull_request} =
122
      %PullRequest{}
123
      |> PullRequest.changeset(%{
124
        repository_id: repository.id,
125
        issue_id: issue.id,
126
        head_repository_id: repository.id,
127
        head_ref: head_ref,
128
        head_sha: sha_for(head_ref),
129
        base_ref: base_ref,
130
        base_sha: sha_for(base_ref),
131
        state: "open"
132
      })
133
      |> Repo.insert()
134
135
    Repo.preload(pull_request, :issue)
136
  end
137
138
  defp sha_for(ref) do
139
    :sha
140
    |> :crypto.hash(ref)
141
    |> Base.encode16(case: :lower)
142
  end
143
end
test/openagents/stacks/merge_test.exs modified +1 -1

@@ -102,7 +102,7 @@ defmodule OpenAgents.Stacks.MergeTest do

102 102
      assert Repo.exists?(
103 103
               from event in StackEvent,
104 104
                 where:
105
                   event.event_type == "pull_request_stack.merged" and
105
                   event.event_type == "pull_request_stack.merge_completed" and
106 106
                     event.stack_id == ^stack.id
107 107
             )
108 108
test/openagents/stacks/policy_test.exs added +232

@@ -0,0 +1,232 @@

1
defmodule OpenAgents.Stacks.PolicyTest do
2
  @moduledoc """
3
  Effective-base policy evaluation (#52): every stacked pull request carries
4
  its direct base and its effective base (the stack trunk), and policy
5
  configuration such as `CODEOWNERS` resolves from the effective base tree —
6
  so an unmerged lower layer cannot weaken the rules an upper layer is held
7
  to.
8
  """
9
10
  use OpenAgents.DataCase, async: false
11
12
  alias OpenAgents.Forge.Repos
13
  alias OpenAgents.PullRequests.PullRequest
14
  alias OpenAgents.Repo
15
  alias OpenAgents.Stacks
16
  alias OpenAgents.Stacks.Policy
17
18
  import OpenAgents.AccountsFixtures
19
  import OpenAgents.IssuesFixtures
20
21
  setup do
22
    base = Path.join(System.tmp_dir!(), "stack-policy-#{System.unique_integer([:positive])}")
23
24
    previous_data = Application.get_env(:openagents, :forge_data_dir)
25
    previous_wal = Application.get_env(:openagents, :forge_wal_dir)
26
    Application.put_env(:openagents, :forge_data_dir, Path.join(base, "data"))
27
    Application.put_env(:openagents, :forge_wal_dir, Path.join(base, "wal"))
28
29
    on_exit(fn ->
30
      restore_env(:forge_data_dir, previous_data)
31
      restore_env(:forge_wal_dir, previous_wal)
32
      File.rm_rf(base)
33
    end)
34
35
    actor = repository_user_fixture("policy-actor")
36
    repository = repository_with_member_fixture(actor)
37
38
    %{actor: actor, repository: repository}
39
  end
40
41
  test "an unstacked pull request evaluates its own base as both bases", context do
42
    %{repository: repository} = context
43
44
    path = Repos.ensure_repo!(repository.storage_key, repository.default_branch)
45
    main = commit(path, nil, "Seed", %{"README.md" => "readme\n"})
46
    {_, 0} = Repos.git(path, ["update-ref", "refs/heads/main", main])
47
    feature = commit(path, main, "Feature", %{"feature.md" => "feature\n"})
48
    {_, 0} = Repos.git(path, ["update-ref", "refs/heads/feature", feature])
49
50
    pull_request = pull_request(repository, "feature", "main", main, feature)
51
52
    assert {:ok, evaluation} = Policy.evaluation(repository, pull_request)
53
    assert evaluation.pull_request_id == pull_request.id
54
    assert evaluation.head_oid == feature
55
    assert evaluation.direct_base == %{ref: "main", oid: main}
56
    assert evaluation.effective_base == %{ref: "main", oid: main}
57
    assert evaluation.stack == nil
58
  end
59
60
  test "a stacked pull request carries the direct base and the trunk as effective base",
61
       context do
62
    %{repository: repository, actor: actor} = context
63
64
    %{oids: oids, stack: stack, pull_requests: [_pr_1, pr_2]} =
65
      seed_stack(repository, actor, %{"CODEOWNERS" => "* @trunk-owners\n"})
66
67
    assert {:ok, evaluation} = Policy.evaluation(repository, pr_2)
68
    assert evaluation.direct_base == %{ref: "layer-1", oid: oids["layer-1"]}
69
    assert evaluation.effective_base == %{ref: "main", oid: oids["main"]}
70
71
    assert evaluation.stack == %{
72
             id: stack.id,
73
             number: stack.number,
74
             position: 2,
75
             size: 2,
76
             health: "healthy"
77
           }
78
  end
79
80
  test "a lower layer editing CODEOWNERS does not weaken an upper layer's policy", context do
81
    %{repository: repository, actor: actor} = context
82
83
    # Layer 1 rewrites CODEOWNERS; layer 2 sits on top of it. Policy for
84
    # layer 2 must keep resolving CODEOWNERS from the trunk, not from the
85
    # unmerged layer-1 tree.
86
    %{path: path, oids: oids, pull_requests: [pr_1, pr_2]} =
87
      seed_stack(repository, actor, %{"CODEOWNERS" => "* @trunk-owners\n"},
88
        layer_1_files: %{"CODEOWNERS" => "* @weakened\n"}
89
      )
90
91
    assert {:ok, evaluation_1} = Policy.evaluation(repository, pr_1)
92
    assert {:ok, evaluation_2} = Policy.evaluation(repository, pr_2)
93
94
    for evaluation <- [evaluation_1, evaluation_2] do
95
      assert {:ok, codeowners} = Policy.codeowners(repository, evaluation)
96
      assert codeowners.path == "CODEOWNERS"
97
      assert codeowners.content == "* @trunk-owners\n"
98
      assert codeowners.source_oid == oids["main"]
99
    end
100
101
    # The direct parent's tree really does carry the weakened file — the
102
    # protection comes from evaluating at the effective base, not from the
103
    # file being absent.
104
    {weakened, 0} = Repos.git(path, ["show", "#{oids["layer-1"]}:CODEOWNERS"])
105
    assert weakened == "* @weakened\n"
106
107
    # Once the lower layer lands on the trunk, the effective base advances
108
    # and later evaluations legitimately see the new configuration.
109
    {_, 0} = Repos.git(path, ["update-ref", "refs/heads/main", oids["layer-1"]])
110
    sync_repository(repository)
111
112
    assert {:ok, landed} = Policy.evaluation(repository, pr_2)
113
    assert landed.effective_base == %{ref: "main", oid: oids["layer-1"]}
114
    assert {:ok, codeowners} = Policy.codeowners(repository, landed)
115
    assert codeowners.content == "* @weakened\n"
116
    assert codeowners.source_oid == oids["layer-1"]
117
  end
118
119
  defp seed_stack(repository, actor, trunk_files, options \\ []) do
120
    path = Repos.ensure_repo!(repository.storage_key, repository.default_branch)
121
122
    main = commit(path, nil, "Seed repository", Map.put(trunk_files, "README.md", "readme\n"))
123
    {_, 0} = Repos.git(path, ["update-ref", "refs/heads/main", main])
124
125
    layer_1_files = Keyword.get(options, :layer_1_files, %{"layer-1.md" => "layer-1\n"})
126
    layer_1 = commit(path, main, "Layer 1", layer_1_files)
127
    {_, 0} = Repos.git(path, ["update-ref", "refs/heads/layer-1", layer_1])
128
129
    layer_2 = commit(path, layer_1, "Layer 2", %{"layer-2.md" => "layer-2\n"})
130
    {_, 0} = Repos.git(path, ["update-ref", "refs/heads/layer-2", layer_2])
131
132
    oids = %{"main" => main, "layer-1" => layer_1, "layer-2" => layer_2}
133
134
    pr_1 = pull_request(repository, "layer-1", "main", main, layer_1)
135
    pr_2 = pull_request(repository, "layer-2", "layer-1", layer_1, layer_2)
136
137
    {:ok, stack} = Stacks.create(repository, [pr_1, pr_2], actor)
138
139
    %{path: path, oids: oids, stack: stack, pull_requests: [pr_1, pr_2]}
140
  end
141
142
  defp pull_request(repository, head_ref, base_ref, base_sha, head_sha) do
143
    issue = issue_fixture(repository, %{title: "PR #{head_ref}"})
144
145
    {:ok, pull_request} =
146
      %PullRequest{}
147
      |> PullRequest.changeset(%{
148
        repository_id: repository.id,
149
        issue_id: issue.id,
150
        head_repository_id: repository.id,
151
        head_ref: head_ref,
152
        head_sha: head_sha,
153
        base_ref: base_ref,
154
        base_sha: base_sha,
155
        state: "open"
156
      })
157
      |> Repo.insert()
158
159
    Repo.preload(pull_request, :issue)
160
  end
161
162
  defp sync_repository(repository) do
163
    OpenAgents.Forge.Sync.ensure_fresh!(repository.storage_key, repository.default_branch)
164
  end
165
166
  # Commits a tree that layers the given files over the parent's tree.
167
  defp commit(path, parent, message, files) do
168
    parent_entries =
169
      if parent do
170
        {listing, 0} = Repos.git(path, ["ls-tree", parent])
171
172
        listing
173
        |> String.split("\n", trim: true)
174
        |> Map.new(fn line ->
175
          [meta, name] = String.split(line, "\t", parts: 2)
176
          {name, meta <> "\t" <> name}
177
        end)
178
      else
179
        %{}
180
      end
181
182
    new_entries =
183
      Map.new(files, fn {name, content} ->
184
        blob = git!(path, ["hash-object", "-w", "--stdin"], content)
185
        {name, "100644 blob #{blob}\t#{name}"}
186
      end)
187
188
    listing =
189
      parent_entries
190
      |> Map.merge(new_entries)
191
      |> Map.values()
192
      |> Enum.map_join("", &(&1 <> "\n"))
193
194
    tree = git!(path, ["mktree"], listing)
195
    parent_args = if parent, do: ["-p", parent], else: []
196
197
    git!(path, ["commit-tree", tree] ++ parent_args ++ ["-m", message], "",
198
      env: [
199
        {"GIT_AUTHOR_NAME", "Test Author"},
200
        {"GIT_AUTHOR_EMAIL", "author@example.test"},
201
        {"GIT_COMMITTER_NAME", "Test Author"},
202
        {"GIT_COMMITTER_EMAIL", "author@example.test"}
203
      ]
204
    )
205
  end
206
207
  defp git!(git_dir, args, input, options \\ []) do
208
    input_path =
209
      Path.join(
210
        System.tmp_dir!(),
211
        "stack-policy-input-#{System.unique_integer([:positive])}"
212
      )
213
214
    File.write!(input_path, input)
215
216
    try do
217
      {output, 0} =
218
        System.cmd(
219
          "sh",
220
          ["-c", ~s(exec git --git-dir "$GIT_DIR" "$@" < "$INPUT"), "sh"] ++ args,
221
          env: [{"GIT_DIR", git_dir}, {"INPUT", input_path}] ++ Keyword.get(options, :env, [])
222
        )
223
224
      String.trim(output)
225
    after
226
      File.rm(input_path)
227
    end
228
  end
229
230
  defp restore_env(key, nil), do: Application.delete_env(:openagents, key)
231
  defp restore_env(key, value), do: Application.put_env(:openagents, key, value)
232
end
test/openagents/stacks/restack_test.exs modified +1 -1

@@ -88,7 +88,7 @@ defmodule OpenAgents.Stacks.RestackTest do

88 88
      assert Repo.exists?(
89 89
               from event in StackEvent,
90 90
                 where:
91
                   event.event_type == "pull_request_stack.rebased" and
91
                   event.event_type == "pull_request_stack.rebase_completed" and
92 92
                     event.stack_id == ^stack.id and event.stack_version == 2
93 93
             )
94 94
test/openagents_web/controllers/stack_controller_test.exs modified +46 -2

@@ -72,8 +72,12 @@ defmodule OpenAgentsWeb.StackControllerTest do

72 72
      assert %{"position" => 3, "observed_head_oid" => head_3} = entry_3
73 73
      assert head_3 == oids["layer-3"]
74 74
75
      assert [%StackEvent{event_type: "pull_request_stack.created", stack_version: 1}] =
76
               Repo.all(StackEvent)
75
      events = Repo.all(from event in StackEvent, order_by: [asc: event.inserted_at])
76
77
      assert [%StackEvent{event_type: "pull_request_stack.created", stack_version: 1} | stacked] =
78
               events
79
80
      assert Enum.map(stacked, & &1.event_type) == List.duplicate("pull_request.stacked", 3)
77 81
78 82
      replay_conn =
79 83
        conn

@@ -526,6 +530,46 @@ defmodule OpenAgentsWeb.StackControllerTest do

526 530
    end
527 531
  end
528 532
533
  describe "stack membership in pull request payloads" do
534
    test "stacked pull requests embed the stack context; unstacked carry nil", %{conn: conn} do
535
      repository = repository_fixture()
536
      oids = seed_chain(repository, ["layer-1", "layer-2"])
537
      [pr_1, pr_2] = pull_request_chain(repository, oids, ["layer-1", "layer-2"])
538
      solo = pull_request(repository, "layer-2", "main", oids["main"], oids["layer-2"])
539
      conn = put_forge_api_token(conn, "stack-payload", repository)
540
541
      create_conn =
542
        conn
543
        |> put_req_header("idempotency-key", "stack-payload-1")
544
        |> post(path(repository), %{trunk_ref: "main", pull_requests: [pr_1, pr_2]})
545
546
      assert %{"number" => stack_number} = json_response(create_conn, 201)
547
548
      pulls = "/api/v3/repos/#{repository.owner}/#{repository.name}/pulls"
549
550
      show_conn = get(conn, "#{pulls}/#{pr_2}")
551
      main_oid = oids["main"]
552
553
      assert %{
554
               "stack" => %{
555
                 "number" => ^stack_number,
556
                 "position" => 2,
557
                 "size" => 2,
558
                 "health" => "healthy",
559
                 "base" => %{"ref" => "main", "sha" => ^main_oid}
560
               }
561
             } = json_response(show_conn, 200)
562
563
      solo_conn = get(conn, "#{pulls}/#{solo}")
564
      assert %{"stack" => nil} = json_response(solo_conn, 200)
565
566
      index_conn = get(conn, pulls)
567
      by_number = Map.new(json_response(index_conn, 200), &{&1["number"], &1})
568
      assert %{"stack" => %{"position" => 1, "size" => 2}} = Map.fetch!(by_number, pr_1)
569
      assert %{"stack" => nil} = Map.fetch!(by_number, solo)
570
    end
571
  end
572
529 573
  defp path(repository), do: "/api/v3/repos/#{repository.owner}/#{repository.name}/stacks"
530 574
531 575
  defp seed_chain(repository, branches) do

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