You've already forked tribes-plugin-supertest
forked from tribes/tribes-plugin-template
c20a3b4f71
Rename the Supertest fixture plugin to tribe-one-supertest across manifest, runtime API routes, metrics, e2e fixtures, and tests.
488 lines
14 KiB
Elixir
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
|