Support candidate readiness during rolling updates

06952374819b · AtlantisPleb · · parent ac40f633af5c

Support candidate readiness during rolling updates

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/rolling_node_probe.ex
  • modified lib/openagents/forge/rolling_provider/gcp.ex
  • added test/openagents/forge/rolling_node_probe_test.exs
  • modified test/openagents/forge/rolling_provider/gcp_test.exs

Diff

4 files changed, +124 -11

lib/openagents/forge/rolling_node_probe.ex modified +21 -3

@@ -10,14 +10,32 @@ defmodule OpenAgents.Forge.RollingNodeProbe do

10 10
  @doc "Return node health and Ra quorum for an expected fleet size."
11 11
  def status(expected_fleet_size)
12 12
      when is_integer(expected_fleet_size) and expected_fleet_size > 0 do
13
    status(expected_fleet_size, nil, nil)
14
  end
15
16
  @doc "Return node health against an operator-approved rolling image identity."
17
  def status(expected_fleet_size, expected_sha, expected_image_digest)
18
      when is_integer(expected_fleet_size) and expected_fleet_size > 0 do
13 19
    report = OpenAgents.Cluster.local_report()
14 20
    ra_members = OpenAgents.Cluster.Ra.members()
15 21
22
    candidate_image? =
23
      report["revision"] == expected_sha and report["image_digest"] == expected_image_digest
24
25
    boot_converged? = report["boot_converged"] == true or candidate_image?
26
27
    candidate_ready? =
28
      candidate_image? and
29
        report["deployment_ready"] == true and
30
        report["admission_ready"] == true
31
32
    database_ready? = database_ready?()
33
16 34
    %{
17 35
      member: true,
18
      ready: report["ready"] == true,
19
      boot_converged: report["boot_converged"] == true,
20
      database_ready: database_ready?(),
36
      ready: (report["ready"] == true or candidate_ready?) and database_ready?,
37
      boot_converged: boot_converged?,
38
      database_ready: database_ready?,
21 39
      sha: report["revision"],
22 40
      image_digest: report["image_digest"],
23 41
      ra_quorum: length(ra_members) * 2 > expected_fleet_size
lib/openagents/forge/rolling_provider/gcp.ex modified +15 -5

@@ -57,7 +57,7 @@ defmodule OpenAgents.Forge.RollingProvider.Gcp do

57 57
  @impl true
58 58
  def capacity(nodes, context) do
59 59
    with {:ok, config} <- config(),
60
         {:ok, probes} <- probes(config, nodes, length(context.expected_nodes)) do
60
         {:ok, probes} <- probes(config, nodes, context) do
61 61
      ready = Enum.count(probes, & &1.ready)
62 62
      majority = div(length(context.expected_nodes), 2) + 1
63 63
      ra_quorum? = Enum.any?(probes, & &1.ra_quorum)

@@ -76,7 +76,7 @@ defmodule OpenAgents.Forge.RollingProvider.Gcp do

76 76
  @impl true
77 77
  def status(node, context) do
78 78
    with {:ok, config} <- config(),
79
         {:ok, probe} <- probe(config, node, length(context.expected_nodes)) do
79
         {:ok, probe} <- probe(config, node, context) do
80 80
      {:ok,
81 81
       Map.take(probe, [:member, :ready, :boot_converged, :database_ready, :sha, :image_digest])}
82 82
    end

@@ -113,11 +113,11 @@ defmodule OpenAgents.Forge.RollingProvider.Gcp do

113 113
114 114
  def validate_config(_config), do: {:error, :invalid_provider_config}
115 115
116
  defp probes(config, nodes, expected_fleet_size) do
116
  defp probes(config, nodes, context) do
117 117
    results =
118 118
      Task.async_stream(
119 119
        nodes,
120
        &probe(config, &1, expected_fleet_size),
120
        &probe(config, &1, context),
121 121
        ordered: false,
122 122
        timeout: timeout(config)
123 123
      )

@@ -133,7 +133,17 @@ defmodule OpenAgents.Forge.RollingProvider.Gcp do

133 133
    end
134 134
  end
135 135
136
  defp probe(config, node, expected_fleet_size) do
136
  defp probe(config, node, context) do
137
    arguments = [length(context.expected_nodes), context.sha, context.image_digest]
138
139
    case rpc(config, node, RollingNodeProbe, :status, arguments) do
140
      %{member: true} = result -> {:ok, result}
141
      {:error, _reason} -> legacy_probe(config, node, length(context.expected_nodes))
142
      other -> {:error, {:invalid_node_probe, other}}
143
    end
144
  end
145
146
  defp legacy_probe(config, node, expected_fleet_size) do
137 147
    case rpc(config, node, RollingNodeProbe, :status, [expected_fleet_size]) do
138 148
      %{member: true} = result -> {:ok, result}
139 149
      {:error, _transport_reason} -> {:ok, unavailable_probe()}
test/openagents/forge/rolling_node_probe_test.exs added +65

@@ -0,0 +1,65 @@

1
defmodule OpenAgents.Forge.RollingNodeProbeTest do
2
  use OpenAgents.DataCase, async: false
3
4
  alias OpenAgents.Forge.BootConverge
5
  alias OpenAgents.Forge.RollingNodeProbe
6
7
  setup do
8
    previous_boot = :persistent_term.get({BootConverge, :state}, :missing)
9
    previous_digest = Application.get_env(:openagents, :image_digest)
10
    digest = "sha256:" <> String.duplicate("f", 64)
11
12
    Application.put_env(:openagents, :image_digest, digest)
13
14
    :persistent_term.put(
15
      {BootConverge, :state},
16
      %{
17
        "schema" => "openagents.forge.boot-convergence.v2",
18
        "state" => "degraded",
19
        "ready" => false,
20
        "reason" => "previous_live_target",
21
        "sha" => OpenAgents.BuildInfo.revision(),
22
        "artifact_digest" => nil,
23
        "manifest_digest" => nil,
24
        "modules" => 0,
25
        "attempts" => 1,
26
        "retry_in_ms" => nil
27
      }
28
    )
29
30
    on_exit(fn ->
31
      restore_boot(previous_boot)
32
      restore_env(:image_digest, previous_digest)
33
    end)
34
35
    %{digest: digest}
36
  end
37
38
  test "accepts the exact rolling image while boot convergence references the previous target", %{
39
    digest: digest
40
  } do
41
    assert %{
42
             ready: true,
43
             boot_converged: true,
44
             database_ready: true,
45
             sha: revision,
46
             image_digest: ^digest
47
           } = RollingNodeProbe.status(1, OpenAgents.BuildInfo.revision(), digest)
48
49
    assert revision == OpenAgents.BuildInfo.revision()
50
  end
51
52
  test "rejects a rolling image digest mismatch", %{digest: digest} do
53
    refute RollingNodeProbe.status(
54
             1,
55
             OpenAgents.BuildInfo.revision(),
56
             String.replace_suffix(digest, "f", "e")
57
           ).ready
58
  end
59
60
  defp restore_boot(:missing), do: :persistent_term.erase({BootConverge, :state})
61
  defp restore_boot(state), do: :persistent_term.put({BootConverge, :state}, state)
62
63
  defp restore_env(key, nil), do: Application.delete_env(:openagents, key)
64
  defp restore_env(key, value), do: Application.put_env(:openagents, key, value)
65
end
test/openagents/forge/rolling_provider/gcp_test.exs modified +23 -3

@@ -51,6 +51,7 @@ defmodule OpenAgents.Forge.RollingProvider.GcpTest do

51 51
    assert :ok = Gcp.remove_readiness(hd(@nodes), context)
52 52
    assert {:ok, 0} = Gcp.drain(hd(@nodes), context)
53 53
    assert {:ok, %{ready: 2, quorum: true}} = Gcp.capacity(tl(@nodes), context)
54
    assert_receive {:rpc, _, RollingNodeProbe, :status, [3, @sha, @digest]}
54 55
55 56
    assert {:ok,
56 57
            %{

@@ -97,7 +98,7 @@ defmodule OpenAgents.Forge.RollingProvider.GcpTest do

97 98
  end
98 99
99 100
  test "fails capacity closed when Ra quorum is absent" do
100
    rpc = fn _node, RollingNodeProbe, :status, [_expected], _timeout ->
101
    rpc = fn _node, RollingNodeProbe, :status, [_expected, @sha, @digest], _timeout ->
101 102
      Map.put(probe(hd(@nodes)), :ra_quorum, false)
102 103
    end
103 104

@@ -106,7 +107,7 @@ defmodule OpenAgents.Forge.RollingProvider.GcpTest do

106 107
  end
107 108
108 109
  test "reports a rebooting node as unavailable when Erlang distribution disconnects" do
109
    rpc = fn _node, RollingNodeProbe, :status, [_expected], _timeout ->
110
    rpc = fn _node, RollingNodeProbe, :status, [_expected, @sha, @digest], _timeout ->
110 111
      :erlang.error({:erpc, :noconnection})
111 112
    end
112 113

@@ -124,7 +125,7 @@ defmodule OpenAgents.Forge.RollingProvider.GcpTest do

124 125
  end
125 126
126 127
  test "reports other reboot transport exits as unavailable" do
127
    rpc = fn _node, RollingNodeProbe, :status, [_expected], _timeout ->
128
    rpc = fn _node, RollingNodeProbe, :status, [_expected, @sha, @digest], _timeout ->
128 129
      exit(:nodedown)
129 130
    end
130 131

@@ -133,6 +134,25 @@ defmodule OpenAgents.Forge.RollingProvider.GcpTest do

133 134
    assert {:ok, %{ready: 0, quorum: false}} = Gcp.capacity([hd(@nodes)], context())
134 135
  end
135 136
137
  test "uses the bounded legacy probe while replacing an older fleet image" do
138
    owner = self()
139
140
    rpc = fn node, RollingNodeProbe, :status, arguments, _timeout ->
141
      send(owner, {:probe_arguments, arguments})
142
143
      case arguments do
144
        [3, @sha, @digest] -> {:error, :undef}
145
        [3] -> probe(node)
146
      end
147
    end
148
149
    put_config(rpc)
150
151
    assert {:ok, %{ready: 2, quorum: true}} = Gcp.capacity(tl(@nodes), context())
152
    assert_receive {:probe_arguments, [3, @sha, @digest]}
153
    assert_receive {:probe_arguments, [3]}
154
  end
155
136 156
  test "reports only connected nodes in the configured fleet inventory" do
137 157
    rpc = fn _node, _module, _function, _arguments, _timeout -> :ok end
138 158

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