Close transactional deployment failure gaps

3a96e867a7e2 · Christopher David · · parent ef3f17d35f6e

Close transactional deployment failure gaps

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/deployment.ex
  • modified test/openagents/forge/boot_converge_test.exs
  • modified test/openagents/forge/deployment_node_test.exs

Diff

3 files changed, +410 -6

lib/openagents/forge/deployment.ex modified +12 -6

@@ -105,7 +105,8 @@ defmodule OpenAgents.Forge.Deployment do

105 105
               node_results: health_results(results)
106 106
           }}
107 107
        else
108
          {:error, reason} -> {:error, reason, session_with_results(session, results)}
108
          {:error, reason} ->
109
            {:error, reason, %{session | node_results: health_results(results)}}
109 110
        end
110 111
    end
111 112
  end

@@ -232,12 +233,14 @@ defmodule OpenAgents.Forge.Deployment do

232 233
    results = token_fanout(session, :rollback, opts)
233 234
    restored? = all_restored?(results)
234 235
235
    node_results =
236
    rollback_results =
236 237
      Map.new(results, fn
237 238
        {node, {:ok, {:ok, %{"restored" => true}}}} -> {to_string(node), "restored"}
238 239
        {node, result} -> {to_string(node), "rollback_failed:" <> result_code(result)}
239 240
      end)
240 241
242
    node_results = Map.merge(session.node_results, rollback_results)
243
241 244
    {:error, failure_outcome(session, reason, restored?, node_results)}
242 245
  end
243 246

@@ -286,7 +289,13 @@ defmodule OpenAgents.Forge.Deployment do

286 289
  end
287 290
288 291
  defp token_fanout(session, phase, opts) do
289
    nodes = Map.get(session, :internal_nodes, session.expected_nodes)
292
    nodes =
293
      Map.get(
294
        session,
295
        :internal_nodes,
296
        Enum.filter(session.expected_nodes, &Map.has_key?(session.tokens, &1))
297
      )
298
290 299
    fanout_tokens(session, nodes, phase, opts)
291 300
  end
292 301

@@ -392,9 +401,6 @@ defmodule OpenAgents.Forge.Deployment do

392 401
    end)
393 402
  end
394 403
395
  defp session_with_results(session, results),
396
    do: %{session | node_results: merge_results(session, results, "ok")}
397
398 404
  defp membership_delta(expected, current) do
399 405
    %{
400 406
      missing: Enum.map(expected -- current, &to_string/1),
test/openagents/forge/boot_converge_test.exs modified +63

@@ -282,6 +282,69 @@ defmodule OpenAgents.Forge.BootConvergeTest do

282 282
    end)
283 283
  end
284 284
285
  test "the supervised worker refreshes a ready image state" do
286
    previous_enabled = Application.get_env(:openagents, :forge_boot_converge_enabled)
287
    previous_min = Application.get_env(:openagents, :forge_boot_retry_min_ms)
288
    previous_max = Application.get_env(:openagents, :forge_boot_retry_max_ms)
289
290
    Application.put_env(:openagents, :forge_boot_converge_enabled, true)
291
    Application.put_env(:openagents, :forge_boot_retry_min_ms, 10)
292
    Application.put_env(:openagents, :forge_boot_retry_max_ms, 20)
293
294
    name = Module.concat(__MODULE__, "Ready#{System.unique_integer([:positive])}")
295
296
    start_supervised!(
297
      Supervisor.child_spec(
298
        {BootConverge, name: name, repo: @repo},
299
        id: name
300
      )
301
    )
302
303
    assert %{"state" => "image", "ready" => true} = BootConverge.state()
304
    assert :sys.get_state(name).retry_ms == 10
305
306
    send(name, :retry_convergence)
307
    assert :sys.get_state(name).retry_ms == 10
308
309
    send(name, :irrelevant_message)
310
    assert :sys.get_state(name).retry_ms == 10
311
312
    on_exit(fn ->
313
      restore_env(:forge_boot_converge_enabled, previous_enabled)
314
      restore_env(:forge_boot_retry_min_ms, previous_min)
315
      restore_env(:forge_boot_retry_max_ms, previous_max)
316
    end)
317
  end
318
319
  test "image-matching legacy target remains ready without artifact metadata" do
320
    target = insert_target!("live", %{})
321
322
    target
323
    |> Ecto.Changeset.change(%{sha: OpenAgents.BuildInfo.revision()})
324
    |> Repo.update!()
325
326
    assert %{
327
             "state" => "image",
328
             "ready" => true,
329
             "reason" => "image_matches_live",
330
             "sha" => "image"
331
           } = BootConverge.converge(@repo)
332
  end
333
334
  test "an unreadable cache entry degrades with a bounded reason" do
335
    {module, binary} = scratch_beam(OpenAgents.Scratch.BootConvergeUnreadable)
336
    artifact = artifact(module, binary)
337
    artifact_abs = Path.join(Repos.data_dir(), artifact.details["artifact"])
338
    File.mkdir_p!(artifact_abs)
339
    insert_target!("live", artifact.details)
340
341
    assert %{
342
             "state" => "degraded",
343
             "ready" => false,
344
             "reason" => "artifact_cache_read_failed"
345
           } = BootConverge.converge(@repo)
346
  end
347
285 348
  defp restore_env(key, nil), do: Application.delete_env(:openagents, key)
286 349
  defp restore_env(key, value), do: Application.put_env(:openagents, key, value)
287 350
test/openagents/forge/deployment_node_test.exs modified +335

@@ -4,6 +4,7 @@ defmodule OpenAgents.Forge.DeploymentNodeTest do

4 4
  alias OpenAgents.Forge.ArtifactFixtures
5 5
  alias OpenAgents.Forge.BuildArtifact
6 6
  alias OpenAgents.Forge.BuildProtocol
7
  alias OpenAgents.Forge.Deployment
7 8
  alias OpenAgents.Forge.DeploymentNode
8 9
  alias OpenAgents.Forge.Target
9 10

@@ -11,11 +12,13 @@ defmodule OpenAgents.Forge.DeploymentNodeTest do

11 12
    base = Path.join(System.tmp_dir!(), "deployment-node-#{System.unique_integer([:positive])}")
12 13
    previous_data = Application.get_env(:openagents, :forge_data_dir)
13 14
    previous_allowlist = Application.get_env(:openagents, :forge_hot_load_allowlist)
15
    previous_expected = Application.get_env(:openagents, :forge_expected_fleet_size)
14 16
    previous_state = :sys.get_state(DeploymentNode)
15 17
    previous_persisted = :persistent_term.get({DeploymentNode, :state}, :missing)
16 18
17 19
    Application.put_env(:openagents, :forge_data_dir, base)
18 20
    Application.put_env(:openagents, :forge_hot_load_allowlist, ["OpenAgents.Scratch."])
21
    Application.put_env(:openagents, :forge_expected_fleet_size, 1)
19 22
    reset_participant()
20 23
21 24
    on_exit(fn ->

@@ -24,6 +27,7 @@ defmodule OpenAgents.Forge.DeploymentNodeTest do

24 27
25 28
      restore_env(:forge_data_dir, previous_data)
26 29
      restore_env(:forge_hot_load_allowlist, previous_allowlist)
30
      restore_env(:forge_expected_fleet_size, previous_expected)
27 31
      :persistent_term.erase({OpenAgents.Forge.BootConverge, :state})
28 32
      File.rm_rf(base)
29 33
    end)

@@ -82,6 +86,24 @@ defmodule OpenAgents.Forge.DeploymentNodeTest do

82 86
    assert DeploymentNode.health()["ready"]
83 87
  end
84 88
89
  test "rollback removes a candidate whose module was absent before prepare" do
90
    fixture = absent_version("AbsentRollback")
91
    request = request(fixture)
92
93
    refute Code.ensure_loaded?(fixture.module)
94
95
    assert {:ok, %{"token" => token, "prior" => [prior]}} =
96
             DeploymentNode.prepare(request)
97
98
    assert prior == %{"module" => to_string(fixture.module), "state" => "absent"}
99
    assert {:ok, _response} = DeploymentNode.apply_candidate(request.deployment_id, token)
100
    assert fixture.module.revision() == "candidate"
101
102
    assert {:ok, %{"restored" => true}} = DeploymentNode.rollback(request.deployment_id, token)
103
    refute Code.ensure_loaded?(fixture.module)
104
    assert DeploymentNode.health()["ready"]
105
  end
106
85 107
  test "an expired token restores applied code before removing the transaction" do
86 108
    fixture = versions("ExpiredToken")
87 109
    request = request(fixture)

@@ -215,6 +237,264 @@ defmodule OpenAgents.Forge.DeploymentNodeTest do

215 237
    assert DeploymentNode.health()["ready"]
216 238
  end
217 239
240
  test "prepare rejects every malformed request identity before changing state" do
241
    fixture = versions("InvalidRequest")
242
    valid = request(fixture)
243
244
    invalid_requests = [
245
      {:invalid_deployment_request, :not_a_map},
246
      {:unexpected_deployment_fields, Map.delete(valid, :repo)},
247
      {:invalid_deployment_id, %{valid | deployment_id: "invalid"}},
248
      {:invalid_target_id, %{valid | target_id: "invalid"}},
249
      {:invalid_build_id, %{valid | build_id: "invalid"}},
250
      {:invalid_source_sha, %{valid | sha: "invalid"}},
251
      {:invalid_artifact_digest, %{valid | artifact_digest: "invalid"}},
252
      {:invalid_manifest_digest, %{valid | manifest_digest: "invalid"}},
253
      {:invalid_repo, %{valid | repo: nil}},
254
      {:invalid_artifact, %{valid | artifact_bytes: nil}},
255
      {:invalid_expected_nodes, %{valid | expected_nodes: []}}
256
    ]
257
258
    for {reason, invalid} <- invalid_requests do
259
      assert {:error, ^reason} = DeploymentNode.prepare(invalid)
260
    end
261
262
    assert {:error, :manifest_digest_mismatch} =
263
             valid
264
             |> Map.put(:manifest_digest, String.duplicate("0", 64))
265
             |> DeploymentNode.prepare()
266
267
    assert DeploymentNode.health()["ready"]
268
  end
269
270
  test "prepare independently enforces classification, allowlist, and runtime toolchain" do
271
    direct = absent_version("OffAllowlist")
272
    Application.put_env(:openagents, :forge_hot_load_allowlist, [])
273
    assert {:error, :module_not_allowlisted} = DeploymentNode.prepare(request(direct))
274
275
    Application.put_env(:openagents, :forge_hot_load_allowlist, ["OpenAgents.Scratch."])
276
277
    structural =
278
      absent_version("StructuralArtifact", structural_reasons: ["config_changed"])
279
280
    assert {:error, :artifact_not_direct} = DeploymentNode.prepare(request(structural))
281
282
    mismatched_toolchain =
283
      BuildArtifact.current_toolchain()
284
      |> Map.put("otp", "0")
285
286
    wrong_runtime = absent_version("WrongRuntime", toolchain: mismatched_toolchain)
287
288
    assert {:error, :runtime_toolchain_mismatch} =
289
             DeploymentNode.prepare(request(wrong_runtime))
290
  end
291
292
  test "fault injection is bounded and participant phase notifications are content-free" do
293
    fixture = absent_version("FaultBoundary")
294
    request = request(fixture)
295
    test_pid = self()
296
297
    :sys.replace_state(DeploymentNode, fn state ->
298
      %{state | faults: %{prepare: :invalid_fault}}
299
    end)
300
301
    assert {:error, {:invalid_injected_fault, :invalid_fault}} =
302
             DeploymentNode.prepare(request)
303
304
    :sys.replace_state(DeploymentNode, fn state ->
305
      %{state | faults: %{prepare: :timeout}, fault_timeout_ms: 0}
306
    end)
307
308
    assert {:error, :injected_timeout} = DeploymentNode.prepare(request)
309
310
    :sys.replace_state(DeploymentNode, fn state ->
311
      %{state | faults: %{}, notify: test_pid}
312
    end)
313
314
    assert {:ok, %{"token" => token}} = DeploymentNode.prepare(request)
315
    assert_receive {:forge_deployment_node, _node, :prepared}
316
317
    assert {:ok, %{"restored" => true}} = DeploymentNode.rollback(request.deployment_id, token)
318
  end
319
320
  test "phase ordering and deployment identity are fenced by the token" do
321
    fixture = versions("PhaseOrdering")
322
    request = request(fixture)
323
324
    assert {:ok, %{"token" => token}} = DeploymentNode.prepare(request)
325
326
    assert {:error, :deployment_token_mismatch} =
327
             DeploymentNode.apply_candidate(Ecto.UUID.generate(), token)
328
329
    assert {:error, {:invalid_phase, :prepared, :verify}} =
330
             DeploymentNode.verify_candidate(request.deployment_id, token)
331
332
    assert {:error, {:invalid_phase, :prepared, :commit}} =
333
             DeploymentNode.commit(request.deployment_id, token)
334
335
    assert {:error, {:invalid_phase, :prepared, :finalize}} =
336
             DeploymentNode.finalize(request.deployment_id, token)
337
338
    assert {:ok, _response} = DeploymentNode.apply_candidate(request.deployment_id, token)
339
340
    assert {:error, {:invalid_phase, :applied, :apply}} =
341
             DeploymentNode.apply_candidate(request.deployment_id, token)
342
343
    assert {:ok, _response} = DeploymentNode.verify_candidate(request.deployment_id, token)
344
    assert {:ok, _response} = DeploymentNode.verify_candidate(request.deployment_id, token)
345
    assert {:ok, %{"restored" => true}} = DeploymentNode.rollback(request.deployment_id, token)
346
347
    send(DeploymentNode, :irrelevant_message)
348
    _state = :sys.get_state(DeploymentNode)
349
    assert DeploymentNode.health()["ready"]
350
  end
351
352
  test "candidate verification failure restores the exact prior object code" do
353
    fixture = versions("VerificationFailure")
354
    request = request(fixture)
355
356
    assert {:ok, %{"token" => token}} = DeploymentNode.prepare(request)
357
    assert {:ok, _response} = DeploymentNode.apply_candidate(request.deployment_id, token)
358
    unload(fixture.module)
359
360
    assert {:error, {:verification_failed, :candidate_object_code_mismatch}} =
361
             DeploymentNode.verify_candidate(request.deployment_id, token)
362
363
    assert fixture.module.revision() == "prior"
364
    assert DeploymentNode.health()["ready"]
365
  end
366
367
  test "candidate smoke-contract failure removes a previously absent module" do
368
    fixture = invalid_smoke_version("SmokeFailure")
369
    request = request(fixture)
370
371
    assert {:ok, %{"token" => token}} = DeploymentNode.prepare(request)
372
    assert {:ok, _response} = DeploymentNode.apply_candidate(request.deployment_id, token)
373
374
    assert {:error, {:verification_failed, :candidate_smoke_failed}} =
375
             DeploymentNode.verify_candidate(request.deployment_id, token)
376
377
    refute Code.ensure_loaded?(fixture.module)
378
    assert DeploymentNode.health()["ready"]
379
380
    assert {:error, {:verification_failed, :candidate_smoke_failed}} =
381
             DeploymentNode.install_artifact(request)
382
383
    refute Code.ensure_loaded?(fixture.module)
384
    assert DeploymentNode.health()["ready"]
385
  end
386
387
  test "boot installation is exclusive with a transaction and then completes locally" do
388
    first = versions("InstallExclusive")
389
    second = versions("InstallAfterRollback")
390
    first_request = request(first)
391
    second_request = request(second)
392
393
    assert {:ok, %{"token" => token}} = DeploymentNode.prepare(first_request)
394
    assert {:error, :deployment_in_progress} = DeploymentNode.install_artifact(second_request)
395
396
    assert {:ok, %{"restored" => true}} =
397
             DeploymentNode.rollback(first_request.deployment_id, token)
398
399
    assert {:ok, %{"phase" => "live", "revision" => revision}} =
400
             DeploymentNode.install_artifact(second_request)
401
402
    assert revision == second.sha
403
    assert second.module.revision() == "candidate"
404
    assert DeploymentNode.health()["ready"]
405
  end
406
407
  test "participant enforces its bounded transaction capacity" do
408
    fixture = versions("BoundedCapacity")
409
    request = request(fixture)
410
411
    tokens =
412
      for _index <- 1..4 do
413
        assert {:ok, %{"token" => token}} = DeploymentNode.prepare(request)
414
        token
415
      end
416
417
    assert {:error, :deployment_capacity_reached} = DeploymentNode.prepare(request)
418
419
    for token <- tokens do
420
      assert {:ok, %{"restored" => true}} =
421
               DeploymentNode.rollback(request.deployment_id, token)
422
    end
423
424
    assert DeploymentNode.health()["ready"]
425
  end
426
427
  test "single-node coordinator rolls back failures at every transaction boundary" do
428
    fixture = versions("CoordinatorFailures")
429
430
    for {fault, expected_code, expected_result} <- [
431
          {:prepare, "prepare_failed", "failed"},
432
          {:apply, "canary_apply_failed", "failed"},
433
          {:verify, "canary_verify_failed", "reverted"},
434
          {:commit, "fleet_commit_failed", "reverted"}
435
        ] do
436
      set_faults(%{fault => :error})
437
438
      assert {:error, outcome} = run_deployment(fixture)
439
      assert outcome.error_code == expected_code
440
      assert outcome.result == expected_result
441
442
      set_faults(%{})
443
      assert fixture.module.revision() == "prior"
444
      assert DeploymentNode.health()["ready"]
445
    end
446
  end
447
448
  test "coordinator exposes finalize and explicit rollback failures without losing its fence" do
449
    fixture = versions("CoordinatorFinalization")
450
451
    assert {:ok, session} = run_deployment(fixture)
452
    set_faults(%{finalize: :error})
453
454
    assert {:error, {:finalize_failed, node_results}} = Deployment.finalize(session)
455
    assert node_results[to_string(Node.self())] == "injected_failure"
456
    refute DeploymentNode.health()["ready"]
457
458
    set_faults(%{})
459
    assert {:ok, restored} = Deployment.rollback(session)
460
    assert restored[to_string(Node.self())] == "restored"
461
    assert fixture.module.revision() == "prior"
462
463
    assert {:ok, second_session} = run_deployment(fixture)
464
    set_faults(%{rollback: :error})
465
    assert {:error, rollback_results} = Deployment.rollback(second_session)
466
    assert rollback_results[to_string(Node.self())] == "injected_failure"
467
468
    set_faults(%{})
469
    assert {:ok, _restored} = Deployment.rollback(second_session)
470
    assert DeploymentNode.health()["ready"]
471
  end
472
473
  test "coordinator refuses missing, undersized, and unready fleet snapshots" do
474
    fixture = versions("CoordinatorSnapshot")
475
476
    Application.put_env(:openagents, :forge_expected_fleet_size, 0)
477
478
    assert {:error, empty} =
479
             run_deployment(fixture, members: fn -> [] end)
480
481
    assert empty.error_code == "empty_fleet"
482
483
    Application.put_env(:openagents, :forge_expected_fleet_size, 2)
484
    assert {:error, undersized} = run_deployment(fixture)
485
    assert undersized.error_code == "fleet_size_mismatch"
486
487
    Application.put_env(:openagents, :forge_expected_fleet_size, 1)
488
489
    :sys.replace_state(DeploymentNode, fn state ->
490
      %{state | divergence: "test_divergence"}
491
    end)
492
493
    assert {:error, unready} = run_deployment(fixture)
494
    assert unready.error_code == "fleet_not_ready"
495
    assert unready.node_results[to_string(Node.self())] =~ "unhealthy"
496
  end
497
218 498
  defp versions(suffix) do
219 499
    name = "OpenAgents.Scratch.#{suffix}#{System.unique_integer([:positive])}"
220 500
    module = Module.concat([name])

@@ -248,6 +528,36 @@ defmodule OpenAgents.Forge.DeploymentNodeTest do

248 528
    }
249 529
  end
250 530
531
  defp absent_version(suffix, opts \\ []) do
532
    name = "OpenAgents.Scratch.#{suffix}#{System.unique_integer([:positive])}"
533
    module = Module.concat([name])
534
    candidate_binary = compile(name, "candidate")
535
    unload(module)
536
537
    on_exit(fn -> unload(module) end)
538
539
    sha = random_sha()
540
    built = ArtifactFixtures.create!("openagents.com", sha, [{name, candidate_binary}], opts)
541
542
    %{module: module, built: built, sha: sha}
543
  end
544
545
  defp invalid_smoke_version(suffix) do
546
    name = "OpenAgents.Scratch.#{suffix}#{System.unique_integer([:positive])}"
547
    module = Module.concat([name])
548
549
    [{^module, candidate_binary}] =
550
      Code.compile_string("defmodule #{name} do\n  def revision, do: :invalid\nend")
551
552
    unload(module)
553
    on_exit(fn -> unload(module) end)
554
555
    sha = random_sha()
556
    built = ArtifactFixtures.create!("openagents.com", sha, [{name, candidate_binary}])
557
558
    %{module: module, built: built, sha: sha}
559
  end
560
251 561
  defp request(fixture) do
252 562
    manifest_digest =
253 563
      fixture.built.manifest

@@ -267,6 +577,27 @@ defmodule OpenAgents.Forge.DeploymentNodeTest do

267 577
    }
268 578
  end
269 579
580
  defp run_deployment(fixture, opts \\ []) do
581
    {:ok, verified} =
582
      BuildArtifact.verify(fixture.built.bytes,
583
        digest: fixture.built.digest,
584
        repo: "openagents.com",
585
        source_sha: fixture.sha,
586
        build_id: fixture.built.build_id
587
      )
588
589
    build = %{
590
      repo: "openagents.com",
591
      sha: fixture.sha,
592
      target_id: Ecto.UUID.generate(),
593
      build_id: fixture.built.build_id,
594
      modules: verified.modules,
595
      manifest: fixture.built.manifest
596
    }
597
598
    Deployment.run(build, verified, fixture.built.bytes, opts)
599
  end
600
270 601
  defp compile(name, revision) do
271 602
    [{_module, binary}] =
272 603
      Code.compile_string("defmodule #{name} do\n  def revision, do: #{inspect(revision)}\nend")

@@ -300,6 +631,10 @@ defmodule OpenAgents.Forge.DeploymentNodeTest do

300 631
    end)
301 632
  end
302 633
634
  defp set_faults(faults) do
635
    :sys.replace_state(DeploymentNode, fn state -> %{state | faults: faults} end)
636
  end
637
303 638
  defp restore_env(key, nil), do: Application.delete_env(:openagents, key)
304 639
  defp restore_env(key, value), do: Application.put_env(:openagents, key, value)
305 640

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