Retry work recovery during cluster startup

58046b6c39d3 · AtlantisPleb · · parent c82c4a18e6cb

Retry work recovery during cluster startup

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/work.ex
  • modified test/openagents/work_job_test.exs

Diff

2 files changed, +116 -7

lib/openagents/work.ex modified +81 -7

@@ -12,6 +12,7 @@ defmodule OpenAgents.Work do

12 12
  """
13 13
14 14
  import Ecto.Query
15
  require Logger
15 16
16 17
  alias Ecto.Multi
17 18
  alias OpenAgents.Conversations

@@ -537,28 +538,101 @@ defmodule OpenAgents.Work do

537 538
  the live ACP session by its durably checkpointed id. On a fleet node this is
538 539
  also safe for jobs still running elsewhere — the Horde singleton answers
539 540
  `already_started` and nothing is disturbed. A restart is still recorded as a
540
  degraded incident (queryable evidence of "what happened here"), and only a
541
  job whose worker cannot even start is finalized `interrupted` (the old
542
  behavior, now the last resort).
541
  degraded incident (queryable evidence of "what happened here"). A transient
542
  worker startup failure retries in the supervised background while Horde and
543
  Raft converge. Only a job that exhausts the bounded recovery window becomes
544
  `interrupted`.
543 545
  """
544
  def recover_interrupted_jobs do
546
  def recover_interrupted_jobs(options \\ []) do
547
    start_worker = Keyword.get(options, :start_worker, &start_worker/2)
548
    schedule_retry = Keyword.get(options, :schedule_retry, &schedule_recovery_retry/3)
545 549
    jobs = Repo.all(from(job in Job, where: job.status in ^@active_statuses))
546 550
547 551
    Enum.each(jobs, fn job ->
548 552
      _incident = record_interruption_incident(job)
549 553
550
      case start_worker(worker_module(job.kind), job.id) do
554
      case start_worker.(worker_module(job.kind), job.id) do
551 555
        {:ok, _pid_or_ignore} ->
552 556
          :ok
553 557
554
        {:error, _reason} ->
555
          {:ok, _job} = finish_job(job.id, "interrupted", error_code: "runtime_restarted")
558
        {:error, reason} ->
559
          schedule_retry.(job, reason, start_worker)
556 560
      end
557 561
    end)
558 562
559 563
    :ok
560 564
  end
561 565
566
  @recovery_retry_delays_ms [250, 500, 1_000, 2_000, 4_000, 8_000, 15_000]
567
568
  defp schedule_recovery_retry(job, reason, start_worker) do
569
    Logger.warning(
570
      "work_recovery_retry_scheduled job_id=#{job.id} kind=#{job.kind} " <>
571
        "error_code=#{recovery_error_code(reason)}"
572
    )
573
574
    case Task.Supervisor.start_child(OpenAgents.ProviderTaskSupervisor, fn ->
575
           retry_recovery(job.id, job.kind, start_worker, @recovery_retry_delays_ms)
576
         end) do
577
      {:ok, _pid} ->
578
        :ok
579
580
      {:error, task_reason} ->
581
        Logger.error(
582
          "work_recovery_retry_start_failed job_id=#{job.id} kind=#{job.kind} " <>
583
            "error_code=#{recovery_error_code(task_reason)}"
584
        )
585
586
        {:ok, _job} = finish_job(job.id, "interrupted", error_code: "runtime_restarted")
587
        :ok
588
    end
589
  end
590
591
  defp retry_recovery(job_id, kind, start_worker, [delay_ms | remaining_delays]) do
592
    receive do
593
    after
594
      delay_ms -> :ok
595
    end
596
597
    case get_job(job_id) do
598
      %Job{status: status} when status in @active_statuses ->
599
        case start_worker.(worker_module(kind), job_id) do
600
          {:ok, _pid_or_ignore} ->
601
            Logger.info("work_recovery_resumed job_id=#{job_id} kind=#{kind}")
602
            :ok
603
604
          {:error, reason} when remaining_delays != [] ->
605
            Logger.warning(
606
              "work_recovery_retry job_id=#{job_id} kind=#{kind} " <>
607
                "error_code=#{recovery_error_code(reason)}"
608
            )
609
610
            retry_recovery(job_id, kind, start_worker, remaining_delays)
611
612
          {:error, reason} ->
613
            Logger.error(
614
              "work_recovery_exhausted job_id=#{job_id} kind=#{kind} " <>
615
                "error_code=#{recovery_error_code(reason)}"
616
            )
617
618
            {:ok, _job} =
619
              finish_job(job_id, "interrupted", error_code: "runtime_restarted")
620
621
            :ok
622
        end
623
624
      _terminal_or_missing ->
625
        :ok
626
    end
627
  end
628
629
  defp recovery_error_code(reason) when is_atom(reason), do: Atom.to_string(reason)
630
631
  defp recovery_error_code({reason, _detail}) when is_atom(reason),
632
    do: Atom.to_string(reason)
633
634
  defp recovery_error_code(_reason), do: "worker_start_failed"
635
562 636
  defp worker_module("delegation"), do: OpenAgents.Work.DelegationServer
563 637
  defp worker_module(_kind), do: OpenAgents.Work.JobServer
564 638
test/openagents/work_job_test.exs modified +35

@@ -153,6 +153,41 @@ defmodule OpenAgents.WorkJobTest do

153 153
    refute final.status == "interrupted"
154 154
  end
155 155
156
  test "startup recovery schedules transient cluster startup failures without burying the job" do
157
    {_conversation, job} = create_job("work-recovery-retry", "recover after cluster startup")
158
    job_id = job.id
159
    test_process = self()
160
161
    start_worker = fn _worker_module, job_id ->
162
      send(test_process, {:recovery_start_failed, job_id})
163
      {:error, :cluster_starting}
164
    end
165
166
    schedule_retry = fn scheduled_job, reason, scheduled_start_worker ->
167
      send(
168
        test_process,
169
        {:recovery_retry_scheduled, scheduled_job.id, reason, scheduled_start_worker}
170
      )
171
172
      :ok
173
    end
174
175
    assert :ok =
176
             Work.recover_interrupted_jobs(
177
               start_worker: start_worker,
178
               schedule_retry: schedule_retry
179
             )
180
181
    assert_receive {:recovery_start_failed, ^job_id}
182
183
    assert_receive {:recovery_retry_scheduled, ^job_id, :cluster_starting, ^start_worker}
184
185
    recovered = Work.get_job!(job.id)
186
    assert recovered.status == "queued"
187
    assert recovered.error_code == nil
188
    assert recovered.report == nil
189
  end
190
156 191
  defp wait_for_terminal(job_id, timeout_ms \\ 10_000) do
157 192
    job = Work.get_job!(job_id)
158 193

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