From 92a3818e4d6a46a1f5304b472bf5d6656ab18c97 Mon Sep 17 00:00:00 2001 From: MikaAK Date: Thu, 23 Apr 2026 12:14:10 -0700 Subject: [PATCH] Use Task.start_link for caller propagation in stream worker The stream worker was spawn_link'd without propagating the caller's process dictionary, so libraries relying on per-PID ownership (e.g. Ecto.Adapters.SQL.Sandbox.allow/3, custom per-test sandboxes keyed by $callers, or process-dictionary-based stubs) could not see the stream consumer as a caller. Switch spawn_link to Task.start_link: Task internally captures and restores $callers and $ancestors, gets OTP-style crash logging for free, and participates correctly in telemetry/tracing trees. No manual Process.put dance and no behaviour change for callers that did not rely on ownership propagation. Includes a regression test asserting the worker sees the consumer PID in both $callers and $ancestors. --- lib/lang_ex/graph/stream.ex | 37 +++++++++++++++++++----------- test/lang_ex/graph/stream_test.exs | 21 +++++++++++++++++ 2 files changed, 44 insertions(+), 14 deletions(-) diff --git a/lib/lang_ex/graph/stream.ex b/lib/lang_ex/graph/stream.ex index fd045f6..3d3ec6a 100644 --- a/lib/lang_ex/graph/stream.ex +++ b/lib/lang_ex/graph/stream.ex @@ -6,6 +6,12 @@ defmodule LangEx.Graph.Stream do executes. Uses a spawned process that runs Pregel and sends events to the consumer via the mailbox. + The streaming worker runs as a `Task`, which propagates `$callers` + and `$ancestors` from the calling process. Libraries that rely on + per-PID ownership (e.g. `Ecto.Adapters.SQL.Sandbox.allow/3`, custom + per-test sandboxes, or process-dictionary-based stubs) see the + consumer in the caller chain. + ## Events - `{:node_start, node_name}` - a node is about to execute @@ -33,21 +39,24 @@ defmodule LangEx.Graph.Stream do defp start_execution(graph, input, opts) do parent = self() - spawn_link(fn -> - state = State.apply_update(graph.initial_state, input, graph.reducers) + {:ok, pid} = + Task.start_link(fn -> + state = State.apply_update(graph.initial_state, input, graph.reducers) + + graph + |> Pregel.run(state, %{ + recursion_limit: Keyword.get(opts, :recursion_limit, 25), + checkpointer: graph.checkpointer, + config: Keyword.get(opts, :config, []), + context: Keyword.get(opts, :context), + resume: nil, + step: 0, + emit_to: parent + }) + |> then(&send(parent, {:lang_ex_stream, {:done, &1}})) + end) - graph - |> Pregel.run(state, %{ - recursion_limit: Keyword.get(opts, :recursion_limit, 25), - checkpointer: graph.checkpointer, - config: Keyword.get(opts, :config, []), - context: Keyword.get(opts, :context), - resume: nil, - step: 0, - emit_to: parent - }) - |> then(&send(parent, {:lang_ex_stream, {:done, &1}})) - end) + pid end defp receive_events(:halted), do: {:halt, :halted} diff --git a/test/lang_ex/graph/stream_test.exs b/test/lang_ex/graph/stream_test.exs index ef22ce7..101aa50 100644 --- a/test/lang_ex/graph/stream_test.exs +++ b/test/lang_ex/graph/stream_test.exs @@ -33,5 +33,26 @@ defmodule LangEx.Graph.StreamTest do assert Enum.any?(events, &match?({:node_start, :upper}, &1)) assert Enum.any?(events, &match?({:node_end, :upper, _}, &1)) end + + test "stream worker inherits $callers and $ancestors from the consumer" do + caller = self() + + Graph.new(value: 0) + |> Graph.add_node(:snapshot, fn state -> + send(caller, {:callers, Process.get(:"$callers")}) + send(caller, {:ancestors, Process.get(:"$ancestors")}) + %{value: state.value + 1} + end) + |> Graph.add_edge(:__start__, :snapshot) + |> Graph.add_edge(:snapshot, :__end__) + |> Graph.compile() + |> LangEx.stream(%{value: 0}) + |> Enum.to_list() + + assert_received {:callers, callers} + assert_received {:ancestors, ancestors} + assert caller in callers + assert caller in ancestors + end end end