Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 25 additions & 17 deletions test/datasets_contract_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -206,11 +206,17 @@ defmodule DatasetsContractTest do
File.write!(path, ["q,a\n" | rows] ++ ["bad,row,extra\n"])

caller = spawn(fn -> Datasets.csv(path, [:q]) end)
walk = wait_for_monitored(caller)
walk = wait_for_walk(caller)
ref = Process.monitor(walk)
# Suspended, the walk cannot run to its end and exit on its own, so it
# ends only when something kills it. The suspension is lifted if the test
# process exits.
:erlang.suspend_process(walk)
Process.exit(caller, :kill)

assert_receive {:DOWN, ^ref, :process, ^walk, :killed}, 1_000
receive do
{:DOWN, ^ref, :process, ^walk, reason} -> assert reason == :killed
end
after
cleanup_tmp("killed-caller.csv")
end
Expand Down Expand Up @@ -434,13 +440,15 @@ defmodule DatasetsContractTest do
cleanup_tmp("typed-math.jsonl")
end

# The process `pid` monitors, once it monitors one.
# The caller also monitors other processes for a moment (the file server
# while it reads), so the walk is the monitored process running the CSV code.
defp wait_for_monitored(pid, tries \\ 500) do
walk =
case Process.info(pid, :monitors) do
{:monitors, monitors} ->
# The walk `caller` starts, once it has started it. The caller also monitors
# other processes for a moment (the file server while it reads), so the walk
# is the monitored process running the CSV code. The caller parses the whole
# file before it starts the walk, which takes as long as the machine is slow,
# so this waits for as long as the caller is alive.
defp wait_for_walk(caller) do
case Process.info(caller, :monitors) do
{:monitors, monitors} ->
walk =
Enum.find_value(monitors, fn
{:process, monitored} when is_pid(monitored) ->
if csv_walk?(monitored), do: monitored
Expand All @@ -449,15 +457,15 @@ defmodule DatasetsContractTest do
nil
end)

nil ->
nil
end
if walk do
walk
else
Process.sleep(5)
wait_for_walk(caller)
end

if walk || tries == 0 do
walk
else
Process.sleep(5)
wait_for_monitored(pid, tries - 1)
nil ->
flunk("the CSV caller ended before it started a walk")
end
end

Expand Down
78 changes: 50 additions & 28 deletions test/parity_sidecar_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -76,14 +76,19 @@ defmodule Imp.BenchmarkTruth.ParitySidecarTest do
refute process_group_alive?(group_pid)
end

test "the timeout option ends a command that does not finish" do
assert {:error, :timeout, %Output{}} =
ParitySidecar.run(@python, ["-c", "import time; time.sleep(60)"], timeout: 150)
end

test "timeout returns captured output and terminates the process tree" do
pid_path = tmp_path("timeout")

assert {:error, :timeout, %Output{} = output} =
ParitySidecar.run(@python, process_tree_args(pid_path), timeout: 150)
{result, [parent_pid, child_pid, group_pid]} =
time_out_ready_tree(process_tree_args(pid_path), pid_path)

assert {:error, :timeout, %Output{} = output} = result
assert output.text =~ "tree ready"
[parent_pid, child_pid, group_pid] = await_pids(pid_path)
refute await_alive?(parent_pid)
refute await_alive?(child_pid)
refute process_group_alive?(group_pid)
Expand All @@ -92,14 +97,13 @@ defmodule Imp.BenchmarkTruth.ParitySidecarTest do
test "bounded cleanup kills a fully TERM-resistant process tree" do
pid_path = tmp_path("term-resistant")

assert {:error, :timeout, %Output{}} =
ParitySidecar.run(
@python,
process_tree_args(pid_path, leader_ignore_term: true, child_ignore_term: true),
timeout: 150
)
{result, [parent_pid, child_pid, group_pid]} =
time_out_ready_tree(
process_tree_args(pid_path, leader_ignore_term: true, child_ignore_term: true),
pid_path
)

[parent_pid, child_pid, group_pid] = await_pids(pid_path)
assert {:error, :timeout, %Output{}} = result
refute await_alive?(parent_pid)
refute await_alive?(child_pid)
refute process_group_alive?(group_pid)
Expand All @@ -108,18 +112,14 @@ defmodule Imp.BenchmarkTruth.ParitySidecarTest do
test "Port exit does not hide a TERM-resistant child with redirected stdio" do
pid_path = tmp_path("mixed-tree")

assert {:error, :timeout, %Output{text: text}} =
ParitySidecar.run(
@python,
process_tree_args(pid_path,
child_ignore_term: true,
child_redirect_stdio: true
),
timeout: 150
)
{result, [parent_pid, child_pid, group_pid]} =
time_out_ready_tree(
process_tree_args(pid_path, child_ignore_term: true, child_redirect_stdio: true),
pid_path
)

assert {:error, :timeout, %Output{text: text}} = result
assert text =~ "tree ready"
[parent_pid, child_pid, group_pid] = await_pids(pid_path)
refute await_alive?(parent_pid)
refute await_alive?(child_pid)
refute process_group_alive?(group_pid)
Expand Down Expand Up @@ -203,11 +203,11 @@ defmodule Imp.BenchmarkTruth.ParitySidecarTest do
#{if leader_ignore_term, do: "signal.signal(signal.SIGTERM, signal.SIG_IGN)", else: ""}
child_code = #{inspect(child_code(child_ignore_term))}
child = subprocess.Popen([sys.executable, "-c", child_code]#{child_stdio})
print("tree ready", flush=True)
with open(#{inspect(pid_path)}, "w", encoding="utf-8") as handle:
handle.write(f"{os.getpid()}\\n{child.pid}\\n{os.getpgrp()}\\n")
handle.flush()
os.fsync(handle.fileno())
print("tree ready", flush=True)
time.sleep(60)
"""

Expand All @@ -220,25 +220,47 @@ defmodule Imp.BenchmarkTruth.ParitySidecarTest do

defp child_code(false), do: "import time; time.sleep(60)"

defp await_pids(path, attempts \\ 100)
defp await_pids(_path, 0), do: flunk("sidecar did not publish process ids")
# A timeout is a timer in the process that owns the sidecar's port, counted
# from the start. A timer short enough for a fast test fires, on a loaded
# machine, before Python has started the tree, and a long one makes every
# run wait it out. So these tests run the tree with no timer, wait until it
# is ready, and send its owner the message the timer sends.
defp time_out_ready_tree(args, pid_path) do
task = Task.async(fn -> ParitySidecar.run(@python, args) end)
[parent_pid | _rest] = pids = await_pids(pid_path)
send(port_owner(parent_pid), :command_timeout)
{Task.await(task, :infinity), pids}
end

defp port_owner(os_pid) do
Enum.find_value(Port.list(), fn port ->
with {:os_pid, ^os_pid} <- Port.info(port, :os_pid),
{:connected, owner} <- Port.info(port, :connected) do
owner
else
_other -> nil
end
end) || flunk("no port runs OS process #{os_pid}")
end

defp await_pids(path, attempts) do
# The sidecar writes its process ids once the tree is running, which takes
# as long as the machine is slow.
defp await_pids(path) do
case File.read(path) do
{:ok, contents} ->
case contents |> String.split("\n", trim: true) |> Enum.map(&Integer.parse/1) do
[{parent, ""}, {child, ""}, {group, ""}] -> [parent, child, group]
_other -> retry_pids(path, attempts)
_other -> retry_pids(path)
end

_other ->
retry_pids(path, attempts)
retry_pids(path)
end
end

defp retry_pids(path, attempts) do
defp retry_pids(path) do
Process.sleep(20)
await_pids(path, attempts - 1)
await_pids(path)
end

defp await_alive?(pid, attempts \\ 100)
Expand Down
9 changes: 9 additions & 0 deletions test/test_helper.exs
Original file line number Diff line number Diff line change
Expand Up @@ -78,5 +78,14 @@ external_excludes =

ExUnit.configure(exclude: external_excludes)

# ReqLLM reads its model catalog from LLMDB, which loads it on first use under
# a :global.trans lock. A process that finds the lock taken sleeps a random
# backoff that doubles up to 8 s before it asks again, so when the async tests
# at the start of the suite all make their first model lookup at once, some
# wait several times as long as the load itself (a 1.4 s load, waits of up to
# 4.6 s among 16 first lookups on an idle machine), and on a loaded runner past
# the 5 s of a Task.await. Loading it here, before any test runs, leaves nothing to contend for.
{:ok, _catalog} = LLMDB.load()

Imp.Test.OwnLog.install()
ExUnit.start()
26 changes: 16 additions & 10 deletions test/timed_out_connection_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -23,21 +23,24 @@ defmodule Imp.TimedOutConnectionTest do
}

test "the call after a timed-out call is answered on a new connection" do
test_pid = self()
{:ok, seen} = Agent.start_link(fn -> [] end)

# Bandit serves each HTTP/1 connection from one process, so the handler's
# pid names the connection a request came in on. The first request is held
# past the client's timeout; every later one is answered at once.
# pid names the connection a request came in on. The first request is not
# answered while the test runs, so the first call times out however slow
# the machine is, and its connection stays busy: a request written behind
# it would not be answered either. Every later request is answered at once.
base_url =
Imp.Test.LocalHTTP.start(fn _request ->
connection = self()
count = Agent.get_and_update(seen, &{length(&1) + 1, &1 ++ [connection]})

if count == 1 do
test_ref = Process.monitor(test_pid)

receive do
:release -> :ok
after
2_000 -> :ok
{:DOWN, ^test_ref, :process, _pid, _reason} -> :ok
end
end

Expand All @@ -61,13 +64,16 @@ defmodule Imp.TimedOutConnectionTest do
# connection only when it draws the same shard. Choosing the first shard
# for every request makes that meeting certain.
finch = [name: ReqLLM.Application.finch_name(), pool_strategy: &hd/1]
opts = [timeout: 200, max_retries: 0, req_http_options: [finch: finch]]
opts = [max_retries: 0, req_http_options: [finch: finch]]
messages = [%{role: :user, content: "ping"}]

assert {:error, %Imp.LMError{retryable: true}} = Imp.LM.generate(lm, messages, opts)
# Only the first call has to time out. The second gets room to be answered
# on a loaded machine; landing on the stuck connection still fails below.
assert {:ok, _reply} = Imp.LM.generate(lm, messages, Keyword.put(opts, :timeout, 5_000))
assert {:error, %Imp.LMError{retryable: true}} =
Imp.LM.generate(lm, messages, Keyword.put(opts, :timeout, 200))

# The second call keeps ReqLLM's own receive timeout. On a new connection
# it is answered at once; on the stuck one it is not answered, and fails
# when that timeout runs out.
assert {:ok, _reply} = Imp.LM.generate(lm, messages, opts)

assert [first, second] = Agent.get(seen, & &1)
refute first == second
Expand Down
Loading