towerops/lib/towerops_web/controllers/api/agent_controller.ex

210 lines
5.5 KiB
Elixir

defmodule ToweropsWeb.Api.AgentController do
@moduledoc """
API controller for remote agent communication.
Provides endpoints for agents to:
- Fetch polling configuration
- Submit metrics
- Send heartbeats
"""
use ToweropsWeb, :controller
alias Towerops.Agent.MetricBatch
alias Towerops.Agents
alias Towerops.Snmp
@doc """
GET /api/v1/agent/config
Returns polling configuration for all equipment assigned to the authenticated agent.
"""
def get_config(conn, _params) do
agent_token = conn.assigns.current_agent_token
equipment_list = Agents.list_agent_equipment(agent_token.id)
config = %{
version: "1.0",
poll_interval_seconds: 60,
equipment: Enum.map(equipment_list, &build_equipment_config/1)
}
json(conn, config)
end
@doc """
POST /api/v1/agent/metrics
Accepts a batch of metrics from the agent and processes them asynchronously.
Supports both JSON (legacy) and Protocol Buffers (efficient).
"""
def submit_metrics(conn, params) do
agent_token = conn.assigns.current_agent_token
metrics =
case get_req_header(conn, "content-type") do
["application/x-protobuf" | _] ->
# Decode protobuf
{:ok, body, _conn} = Plug.Conn.read_body(conn)
batch = MetricBatch.decode(body)
convert_protobuf_metrics(batch.metrics)
_ ->
# JSON format (fallback)
Map.get(params, "metrics", [])
end
_ =
Task.start(fn ->
process_metrics(agent_token, metrics)
end)
json(conn, %{status: "accepted", received: length(metrics)})
end
@doc """
POST /api/v1/agent/heartbeat
Updates the agent's last_seen_at timestamp and metadata.
"""
def heartbeat(conn, params) do
agent_token = conn.assigns.current_agent_token
ip = to_string(:inet_parse.ntoa(conn.remote_ip))
metadata = %{
"version" => Map.get(params, "version"),
"hostname" => Map.get(params, "hostname"),
"uptime_seconds" => Map.get(params, "uptime_seconds")
}
# Update synchronously in the heartbeat endpoint
Agents.update_agent_token_heartbeat(agent_token.id, ip, metadata)
json(conn, %{status: "ok"})
end
# Private helpers
defp build_equipment_config(equipment) do
device = equipment.snmp_device
%{
id: equipment.id,
name: equipment.name,
ip_address: equipment.ip_address,
snmp: %{
enabled: equipment.snmp_enabled,
version: equipment.snmp_version,
community: equipment.snmp_community,
port: equipment.snmp_port || 161
},
poll_interval_seconds: equipment.check_interval_seconds || 60,
sensors: build_sensor_config(device),
interfaces: build_interface_config(device)
}
end
defp build_sensor_config(nil), do: []
defp build_sensor_config(device) do
Enum.map(device.sensors, fn sensor ->
%{
id: sensor.id,
type: sensor.sensor_type,
oid: sensor.sensor_oid,
divisor: sensor.sensor_divisor,
unit: sensor.sensor_unit,
metadata: sensor.metadata
}
end)
end
defp build_interface_config(nil), do: []
defp build_interface_config(device) do
Enum.map(device.interfaces, fn interface ->
%{
id: interface.id,
if_index: interface.if_index,
if_name: interface.if_name
}
end)
end
defp process_metrics(_agent_token, metrics) do
Enum.each(metrics, fn metric ->
case metric do
%{"type" => "sensor_reading"} = m ->
Snmp.create_sensor_reading(%{
sensor_id: m["sensor_id"],
value: m["value"],
status: m["status"] || "ok",
checked_at: parse_timestamp(m["timestamp"])
})
%{"type" => "interface_stat"} = m ->
Snmp.create_interface_stat(%{
interface_id: m["interface_id"],
if_in_octets: m["if_in_octets"],
if_out_octets: m["if_out_octets"],
if_in_errors: m["if_in_errors"],
if_out_errors: m["if_out_errors"],
if_in_discards: m["if_in_discards"],
if_out_discards: m["if_out_discards"],
checked_at: parse_timestamp(m["timestamp"])
})
_ ->
:ok
end
end)
end
defp convert_protobuf_metrics(metrics) do
metrics
|> Enum.map(fn metric ->
case metric.metric_type do
{:sensor_reading, sr} ->
%{
"type" => "sensor_reading",
"sensor_id" => sr.sensor_id,
"value" => sr.value,
"status" => sr.status,
"timestamp" => sr.timestamp
}
{:interface_stat, is} ->
%{
"type" => "interface_stat",
"interface_id" => is.interface_id,
"if_in_octets" => is.if_in_octets,
"if_out_octets" => is.if_out_octets,
"if_in_errors" => is.if_in_errors,
"if_out_errors" => is.if_out_errors,
"if_in_discards" => is.if_in_discards,
"if_out_discards" => is.if_out_discards,
"timestamp" => is.timestamp
}
_ ->
nil
end
end)
|> Enum.reject(&is_nil/1)
end
defp parse_timestamp(nil), do: DateTime.truncate(DateTime.utc_now(), :second)
defp parse_timestamp(ts) when is_integer(ts) do
ts
|> DateTime.from_unix!(:second)
|> DateTime.truncate(:second)
end
defp parse_timestamp(ts) when is_binary(ts) do
case DateTime.from_iso8601(ts) do
{:ok, dt, _} -> DateTime.truncate(dt, :second)
_ -> DateTime.truncate(DateTime.utc_now(), :second)
end
end
end