Files
self c20a3b4f71 feat: prefix Supertest plugin slug
Rename the Supertest fixture plugin to tribe-one-supertest across manifest, runtime API routes, metrics, e2e fixtures, and tests.
2026-06-17 22:33:39 +02:00

488 lines
14 KiB
Elixir

defmodule TribeOne.TribesPlugin.Supertest.API do
@moduledoc false
@scheduler_control_name "scheduler:control"
@scheduler_singleton_name "scheduler:singleton"
@scheduler_all_nodes_prefix "scheduler:all_nodes:"
def health do
{:ok,
%{
ok: true,
plugin: "tribe-one-supertest",
version: manifest_version()
}}
end
def schema do
{:ok,
%{
ok: true,
schema: %{
version: 1,
capabilities: %{
cluster_pubsub_probe: true,
scheduler_probe: true
},
tables: %{
cases: table_present?("supertest_cases"),
events: table_present?("supertest_events")
}
}
}}
end
def reset(%{"run_id" => run_id}) when is_binary(run_id) and run_id != "" do
with {:ok, events} <-
TribeOne.TribesPlugin.Supertest.list_events_by_run(run_id, authorize?: false),
:ok <-
destroy_all(
events,
&TribeOne.TribesPlugin.Supertest.destroy_event(&1, authorize?: false)
),
{:ok, cases} <-
TribeOne.TribesPlugin.Supertest.list_cases_by_run(run_id, authorize?: false),
:ok <-
destroy_all(
cases,
&TribeOne.TribesPlugin.Supertest.destroy_case(&1, authorize?: false)
) do
{:ok, %{ok: true, deleted: %{cases: length(cases), events: length(events)}}}
end
end
def reset(_params), do: {:error, :invalid_run_id}
def create_case(params) do
attrs = %{
run_id: required_string!(params, "run_id"),
name: required_string!(params, "name"),
counter: integer_param(params, "counter", 0),
status: atom_param(params, "status", :new),
payload: map_param(params, "payload")
}
with {:ok, test_case} <- TribeOne.TribesPlugin.Supertest.create_case(attrs, authorize?: false) do
{:ok, %{ok: true, case: case_json(test_case)}}
end
rescue
error in ArgumentError -> {:error, error.message}
end
def create_event(params) do
attrs = %{
case_id: required_string!(params, "case_id"),
run_id: required_string!(params, "run_id"),
sequence: integer_param(params, "sequence", nil),
body: required_string!(params, "body"),
payload: map_param(params, "payload")
}
with {:ok, event} <- TribeOne.TribesPlugin.Supertest.create_event(attrs, authorize?: false) do
{:ok, %{ok: true, event: event_json(event)}}
end
rescue
error in ArgumentError -> {:error, error.message}
end
def increment_case(%{"id" => id} = params) when is_binary(id) and id != "" do
delta = integer_param(params, "delta", 1)
with {:ok, test_case} <- TribeOne.TribesPlugin.Supertest.get_case(id, authorize?: false),
{:ok, updated} <-
TribeOne.TribesPlugin.Supertest.update_case(
test_case,
%{counter: test_case.counter + delta},
authorize?: false
) do
{:ok, %{ok: true, case: case_json(updated)}}
end
end
def increment_case(_params), do: {:error, :invalid_case_id}
def state(%{"run_id" => run_id}) when is_binary(run_id) and run_id != "" do
with {:ok, cases} <-
TribeOne.TribesPlugin.Supertest.list_cases_by_run(run_id, authorize?: false),
{:ok, events} <-
TribeOne.TribesPlugin.Supertest.list_events_by_run(run_id, authorize?: false) do
{:ok,
%{
ok: true,
run_id: run_id,
cases: Enum.map(cases, &case_json/1),
events: Enum.map(events, &event_json/1)
}}
end
end
def state(_params), do: {:error, :invalid_run_id}
def broadcast_cluster_pubsub(params) do
run_id = required_string!(params, "run_id")
payload = map_param(params, "payload")
with {:ok, message} <-
TribeOne.TribesPlugin.Supertest.ClusterPubSubProbe.broadcast(run_id, payload) do
{:ok, %{ok: true, message: cluster_pubsub_message_json(message)}}
end
rescue
error in ArgumentError -> {:error, error.message}
end
def reset_cluster_pubsub(%{"run_id" => run_id}) when is_binary(run_id) and run_id != "" do
with :ok <- TribeOne.TribesPlugin.Supertest.ClusterPubSubProbe.reset(run_id) do
{:ok, %{ok: true}}
end
end
def reset_cluster_pubsub(_params), do: {:error, :invalid_run_id}
def received_cluster_pubsub(%{"run_id" => run_id}) when is_binary(run_id) and run_id != "" do
with {:ok, messages} <- TribeOne.TribesPlugin.Supertest.ClusterPubSubProbe.received(run_id) do
{:ok,
%{ok: true, run_id: run_id, messages: Enum.map(messages, &cluster_pubsub_message_json/1)}}
end
end
def received_cluster_pubsub(_params), do: {:error, :invalid_run_id}
def emit_metrics_probe(params) do
value = integer_param(params, "value", nil)
:telemetry.execute([:supertest, :metrics, :probe], %{value: value}, %{})
{:ok, %{ok: true, value: value}}
rescue
error in ArgumentError -> {:error, error.message}
end
def start_scheduler_probe(params) do
run_id = required_string!(params, "run_id")
expected_nodes = integer_param(params, "expected_nodes", 1)
with {:ok, control} <-
upsert_scheduler_case(run_id, @scheduler_control_name, %{
counter: expected_nodes,
status: :active,
payload: %{"expected_nodes" => expected_nodes}
}) do
{:ok, %{ok: true, run_id: run_id, control: case_json(control)}}
end
rescue
error in ArgumentError -> {:error, error.message}
end
def scheduler_probe(%{"run_id" => run_id}) when is_binary(run_id) and run_id != "" do
with {:ok, cases} <-
TribeOne.TribesPlugin.Supertest.list_cases_by_run(run_id, authorize?: false) do
{:ok, scheduler_probe_summary(run_id, cases)}
end
end
def scheduler_probe(_params), do: {:error, :invalid_run_id}
def stop_scheduler_probe(%{"run_id" => run_id}) when is_binary(run_id) and run_id != "" do
with {:ok, cases} <-
TribeOne.TribesPlugin.Supertest.list_cases_by_run(run_id, authorize?: false) do
case Enum.find(cases, &(&1.name == @scheduler_control_name)) do
nil ->
{:ok, %{ok: true, run_id: run_id, stopped: false}}
control ->
with {:ok, stopped} <-
TribeOne.TribesPlugin.Supertest.update_case(
control,
%{status: :done},
authorize?: false
) do
{:ok, %{ok: true, run_id: run_id, stopped: true, control: case_json(stopped)}}
end
end
end
end
def stop_scheduler_probe(_params), do: {:error, :invalid_run_id}
def record_scheduler_execution(kind, %Tribes.Plugin.Scheduler.Job{} = job)
when kind in [:singleton, :all_nodes] do
with {:ok, probes} <- active_scheduler_probes() do
Enum.reduce_while(probes, :ok, fn probe, :ok ->
case record_scheduler_execution(probe, kind, job) do
{:ok, _case} -> {:cont, :ok}
{:error, reason} -> {:halt, {:error, reason}}
end
end)
end
end
def management_echo(context, params) do
{:ok,
%{
ok: true,
plugin: context.plugin,
method: context.method,
version: context.version,
admin: context.admin?,
actor_pubkey: context.actor_pubkey,
params: params
}}
end
defp active_scheduler_probes do
with {:ok, cases} <- TribeOne.TribesPlugin.Supertest.list_cases(authorize?: false) do
{:ok,
Enum.filter(cases, fn test_case ->
test_case.name == @scheduler_control_name and test_case.status == :active
end)}
end
end
defp record_scheduler_execution(probe, kind, job) do
run_id = probe.run_id
name = scheduler_result_name(kind, job.node_id)
payload = scheduler_result_payload(kind, job)
with {:ok, cases} <-
TribeOne.TribesPlugin.Supertest.list_cases_by_run(run_id, authorize?: false) do
case Enum.find(cases, &(&1.name == name)) do
nil ->
create_scheduler_result_case(run_id, name, payload)
existing ->
TribeOne.TribesPlugin.Supertest.update_case(
existing,
%{
counter: existing.counter + 1,
status: :active,
payload: payload
},
authorize?: false
)
end
end
end
defp create_scheduler_result_case(run_id, name, payload) do
TribeOne.TribesPlugin.Supertest.create_case(
%{
run_id: run_id,
name: name,
counter: 1,
status: :active,
payload: payload
},
authorize?: false
)
end
defp upsert_scheduler_case(run_id, name, attrs) do
with {:ok, cases} <-
TribeOne.TribesPlugin.Supertest.list_cases_by_run(run_id, authorize?: false) do
case Enum.find(cases, &(&1.name == name)) do
nil ->
TribeOne.TribesPlugin.Supertest.create_case(
Map.merge(%{run_id: run_id, name: name}, attrs),
authorize?: false
)
existing ->
TribeOne.TribesPlugin.Supertest.update_case(existing, attrs, authorize?: false)
end
end
end
defp scheduler_probe_summary(run_id, cases) do
control = Enum.find(cases, &(&1.name == @scheduler_control_name))
singleton = Enum.find(cases, &(&1.name == @scheduler_singleton_name))
all_nodes = Enum.filter(cases, &String.starts_with?(&1.name, @scheduler_all_nodes_prefix))
expected_nodes = expected_scheduler_nodes(control)
all_node_ids = all_nodes |> Enum.map(&scheduler_node_id/1) |> Enum.uniq() |> Enum.sort()
singleton_counter = counter(singleton)
all_nodes_counter =
Enum.reduce(all_nodes, 0, fn test_case, acc -> acc + test_case.counter end)
%{
ok: true,
run_id: run_id,
active: control && control.status == :active,
expected_nodes: expected_nodes,
passed?: singleton_counter >= 1 and length(all_node_ids) >= expected_nodes,
singleton: %{
counter: singleton_counter,
node_ids: scheduler_case_node_ids([singleton]),
case: maybe_case_json(singleton)
},
all_nodes: %{
counter: all_nodes_counter,
node_ids: all_node_ids,
cases: Enum.map(all_nodes, &case_json/1)
}
}
end
defp scheduler_result_name(:singleton, _node_id), do: @scheduler_singleton_name
defp scheduler_result_name(:all_nodes, node_id), do: @scheduler_all_nodes_prefix <> node_id
defp scheduler_result_payload(kind, job) do
%{
"kind" => Atom.to_string(kind),
"job_id" => job.id,
"cron_id" => job.cron_id,
"attempt_id" => job.attempt_id,
"attempt_token" => job.attempt_token,
"node_id" => job.node_id,
"scheduled_at" => datetime_json(job.scheduled_at)
}
end
defp expected_scheduler_nodes(nil), do: 1
defp expected_scheduler_nodes(control) do
case control.payload do
%{"expected_nodes" => expected_nodes}
when is_integer(expected_nodes) and expected_nodes > 0 ->
expected_nodes
_other ->
1
end
end
defp counter(nil), do: 0
defp counter(test_case), do: test_case.counter
defp scheduler_case_node_ids(cases) do
cases
|> Enum.reject(&is_nil/1)
|> Enum.map(&scheduler_node_id/1)
|> Enum.reject(&is_nil/1)
|> Enum.uniq()
|> Enum.sort()
end
defp scheduler_node_id(test_case) do
case test_case.payload do
%{"node_id" => node_id} when is_binary(node_id) and node_id != "" ->
node_id
_other ->
String.trim_leading(test_case.name, @scheduler_all_nodes_prefix)
end
end
defp maybe_case_json(nil), do: nil
defp maybe_case_json(test_case), do: case_json(test_case)
defp destroy_all(entries, fun) do
Enum.reduce_while(entries, :ok, fn entry, :ok ->
case fun.(entry) do
:ok -> {:cont, :ok}
{:ok, _record} -> {:cont, :ok}
{:error, reason} -> {:halt, {:error, reason}}
end
end)
end
defp table_present?(table) do
query = "SELECT to_regclass($1)::text"
case Ecto.Adapters.SQL.query(Tribes.Repo, query, ["public.#{table}"]) do
{:ok, %{rows: [[name]]}} when is_binary(name) -> true
_ -> false
end
end
defp case_json(test_case) do
%{
id: test_case.id,
run_id: test_case.run_id,
name: test_case.name,
counter: test_case.counter,
status: to_string(test_case.status),
payload: test_case.payload,
inserted_at: datetime_json(test_case.inserted_at),
updated_at: datetime_json(test_case.updated_at)
}
end
defp event_json(event) do
%{
id: event.id,
case_id: event.case_id,
run_id: event.run_id,
sequence: event.sequence,
body: event.body,
payload: event.payload,
inserted_at: datetime_json(event.inserted_at),
updated_at: datetime_json(event.updated_at)
}
end
defp cluster_pubsub_message_json(message) do
%{
id: message.id,
run_id: message.run_id,
payload: message.payload,
origin_node: message.origin_node,
sent_at_usec: message.sent_at_usec
}
end
defp datetime_json(nil), do: nil
defp datetime_json(%DateTime{} = value), do: DateTime.to_iso8601(value)
defp required_string!(params, key) do
case Map.get(params, key) do
value when is_binary(value) and value != "" -> value
_other -> raise ArgumentError, "missing or invalid #{key}"
end
end
defp integer_param(params, key, default) do
case Map.get(params, key, default) do
value when is_integer(value) ->
value
value when is_binary(value) ->
case Integer.parse(value) do
{parsed, ""} -> parsed
_other -> raise ArgumentError, "invalid integer #{key}"
end
nil ->
raise ArgumentError, "missing or invalid #{key}"
_other ->
raise ArgumentError, "invalid integer #{key}"
end
end
defp atom_param(params, key, default) do
case Map.get(params, key, default) do
value when value in [:new, :active, :done] -> value
value when value in ["new", "active", "done"] -> String.to_existing_atom(value)
_other -> raise ArgumentError, "invalid status"
end
end
defp map_param(params, key) do
case Map.get(params, key, %{}) do
value when is_map(value) -> value
_other -> raise ArgumentError, "invalid map #{key}"
end
end
defp manifest_version do
[
otp_app: :tribe_one_supertest,
source_manifest_path: Path.expand("../../manifest.json", __DIR__)
]
|> Tribes.Plugin.Base.read_manifest!()
|> Map.get("version")
rescue
_error -> "unknown"
end
end