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