From a4836eaaf6ec3b27549fe74ea4e1360b1382ccff Mon Sep 17 00:00:00 2001 From: Joe Freeman Date: Sat, 7 Mar 2026 20:26:29 +0000 Subject: [PATCH] Make argument waiting recursive --- server/lib/coflux/orchestration/server.ex | 49 ++++++++++++++++------- 1 file changed, 35 insertions(+), 14 deletions(-) diff --git a/server/lib/coflux/orchestration/server.ex b/server/lib/coflux/orchestration/server.ex index 748484a5..5a0f2697 100644 --- a/server/lib/coflux/orchestration/server.ex +++ b/server/lib/coflux/orchestration/server.ex @@ -4339,25 +4339,46 @@ defmodule Coflux.Orchestration.Server do nil -> [] end - Enum.all?(references, fn - {:execution, run_ext, step_num, attempt} -> - case Runs.get_execution_id(db, run_ext, step_num, attempt) do - {:ok, {execution_id}} when not is_nil(execution_id) -> + all_references_ready?(db, references, MapSet.new()) + end) + end + + defp all_references_ready?(db, references, seen) do + Enum.all?(references, fn + {:execution, run_ext, step_num, attempt} -> + case Runs.get_execution_id(db, run_ext, step_num, attempt) do + {:ok, {execution_id}} when not is_nil(execution_id) -> + if MapSet.member?(seen, execution_id) do + true + else case resolve_result(db, execution_id) do - {:ok, _} -> true - {:pending, _} -> false + {:ok, {:value, value}} -> + inner_refs = + case value do + {:raw, _, refs} -> refs + {:blob, _, _, refs} -> refs + _ -> [] + end + + all_references_ready?(db, inner_refs, MapSet.put(seen, execution_id)) + + {:ok, _} -> + true + + {:pending, _} -> + false end + end - _ -> - false - end + _ -> + false + end - {:fragment, _format, _blob_key, _size, _metadata} -> - true + {:fragment, _format, _blob_key, _size, _metadata} -> + true - {:asset, _external_id} -> - true - end) + {:asset, _external_id} -> + true end) end