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
3 changes: 3 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,9 @@ Early scaffolding. In place:
- `Pulso.MCP` — JSON-RPC 2.0 dispatcher (`initialize`, `tools/list`, `tools/call`, `ping`)
- `Pulso.MCP.Tools` — tool registry, currently one read-only tool (`query_logs`)
- `PulsoWeb.MCPController` at `POST /mcp` (handles single and batched JSON-RPC)
- `PulsoWeb.OTLPController` at `POST /v1/logs` — OTLP/HTTP JSON logs ingest
- `PulsoWeb.LokiController` at `POST /loki/api/v1/push` — Loki push JSON ingest (Snappy protobuf not yet supported)
- `PulsoWeb.CompressedBodyReader` — gzip-aware Plug.Parsers body reader, so JSON receivers accept compressed bodies

Not yet built: alerting, ingestion, Mimir/Tempo clients, storage engine, remediation surface, HITL wiring, distribution (Horde/libcluster/ra).

Expand Down
1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@ mix phx.server
```

- OTLP/HTTP JSON logs land at `POST /v1/logs`. Tenant is picked up from `X-Scope-OrgID` (Loki/Cortex convention), defaulting to `default`.
- Loki push JSON lands at `POST /loki/api/v1/push` (same tenant convention). Gzip-encoded bodies are decompressed transparently; Snappy-framed protobuf is not yet supported.
- The MCP endpoint is exposed at `POST /mcp`. It speaks JSON-RPC 2.0 (`initialize`, `tools/list`, `tools/call`, `ping`).

## 🐳 Local S3 (MinIO)
Expand Down
164 changes: 164 additions & 0 deletions lib/pulso/loki/push.ex
Original file line number Diff line number Diff line change
@@ -0,0 +1,164 @@
defmodule Pulso.Loki.Push do
@moduledoc """
Decode a Loki `POST /loki/api/v1/push` JSON body into a flat list of
`Pulso.Record.Log`.

Loki's push wire format is a list of streams. Each stream carries a
map of static labels and a list of `[timestamp_ns, line]` value tuples
(with an optional third element for structured metadata, added in
Loki 3.0):

{
"streams": [
{
"stream": { "service_name": "api", "level": "info" },
"values": [
[ "1700000000000000000", "hello" ],
[ "1700000000000000001", "with meta", {"trace_id": "abc", "user_id": "u1"} ]
]
}
]
}

## Mapping to Pulso.Record.Log

* Stream labels populate `resource` — they identify the emitter and
are constant across the stream, the same role resource attributes
play in OTLP.
* `service_name` (Grafana's OTel-aligned convention) or, as a
fallback, `service`, is lifted to `Log.service`. It stays in
`resource` too so a `resource.service_name` query still sees it.
* `level` (or `detected_level`) is lifted to `Log.severity_text`.
* Structured metadata is per-record — it populates `attributes`.
* `trace_id` / `span_id` from structured metadata are lifted to the
dedicated struct fields and removed from `attributes`, so a
caller querying by trace id sees them in one canonical place
regardless of which ingest path a record came in on.

## Return shape

`{records, rejected}`, mirroring `Pulso.OTLP.Logs.decode/1`. `rejected`
counts value tuples the decoder could not interpret (missing timestamp,
non-string line, malformed shape). A whole-stream drop (e.g. `stream`
is not a map) is counted as `length(values)` rejects for that stream
— the sender should still know that many records did not land.
"""

alias Pulso.Record.Log

@doc """
Decode a parsed JSON payload.

Returns `{records, rejected}`. A payload that is not shaped like a
Loki push request returns `{[], 0}` — there is nothing to count,
because nothing was parseable in the first place.
"""
@spec decode(map()) :: {[Log.t()], non_neg_integer()}
def decode(%{"streams" => streams}) when is_list(streams) do
streams
|> Enum.reduce({[], 0}, fn stream, {records, rejected} ->
{s_records, s_rejected} = decode_stream(stream)
{[s_records | records], rejected + s_rejected}
end)
|> then(fn {records, rejected} ->
{records |> Enum.reverse() |> List.flatten(), rejected}
end)
end

def decode(_), do: {[], 0}

defp decode_stream(%{"stream" => labels, "values" => values}) when is_map(labels) and is_list(values) do
{service, severity_text, resource} = split_labels(labels)

Enum.reduce(values, {[], 0}, fn value, {rs, rj} ->
case decode_value(value, service, severity_text, resource) do
{:ok, record} -> {[record | rs], rj}
:error -> {rs, rj + 1}
end
end)
|> then(fn {records, rejected} -> {Enum.reverse(records), rejected} end)
end

# Missing/malformed `stream` map: every value in this stream is a
# reject, because we can't attribute records without labels.
defp decode_stream(%{"values" => values}) when is_list(values), do: {[], length(values)}
defp decode_stream(_), do: {[], 0}

# Split stream labels into (service, severity_text, resource). The
# resource map keeps *all* labels — including the lifted service and
# level — so structured queries that filter on `resource.level` still
# work, and so a caller looking at the stored record can see the raw
# label set as sent.
defp split_labels(labels) do
service = labels["service_name"] || labels["service"]

severity_text =
case labels["level"] || labels["detected_level"] do
nil -> nil
v when is_binary(v) -> v
v -> to_string(v)
end

{service, severity_text, labels}
end

defp decode_value([ts, line], service, severity_text, resource) when is_binary(line) do
build_record(ts, line, %{}, service, severity_text, resource)
end

defp decode_value([ts, line, meta], service, severity_text, resource) when is_binary(line) and is_map(meta) do
build_record(ts, line, meta, service, severity_text, resource)
end

defp decode_value(_, _, _, _), do: :error

defp build_record(ts, line, meta, service, severity_text, resource) do
with {:ok, ts_ns} <- parse_ts(ts) do
{trace_id, span_id, attributes} = split_metadata(meta)

{:ok,
%Log{
timestamp_ns: ts_ns,
observed_timestamp_ns: nil,
severity_number: nil,
severity_text: severity_text,
service: service,
body: line,
trace_id: trace_id,
span_id: span_id,
attributes: attributes,
resource: resource
}}
end
end

# Loki's spec form is a decimal-string ns integer. Numeric forms show
# up in a handful of client libraries, so accept them too. Anything
# else is a reject; a record with no timestamp has no place in a
# time-partitioned store.
defp parse_ts(ts) when is_integer(ts) and ts >= 0, do: {:ok, ts}

defp parse_ts(ts) when is_binary(ts) do
case Integer.parse(ts) do
{int, ""} when int >= 0 -> {:ok, int}
_ -> :error
end
end

defp parse_ts(_), do: :error

# Structured metadata is a flat string map per the Loki spec. We lift
# `trace_id` and `span_id` into their dedicated struct fields and
# remove them from `attributes` so a caller has exactly one canonical
# place to look. Every other key stays in `attributes`.
defp split_metadata(meta) when is_map(meta) do
trace_id = nil_if_blank(meta["trace_id"])
span_id = nil_if_blank(meta["span_id"])
attributes = meta |> Map.drop(["trace_id", "span_id"])
{trace_id, span_id, attributes}
end

defp nil_if_blank(nil), do: nil
defp nil_if_blank(""), do: nil
defp nil_if_blank(v), do: v
end
48 changes: 48 additions & 0 deletions lib/pulso_web/compressed_body_reader.ex
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
defmodule PulsoWeb.CompressedBodyReader do
@moduledoc """
Body reader for `Plug.Parsers` that transparently decompresses a
`Content-Encoding: gzip` request body before parsing.

Both Loki push and Prometheus `remote_write` agents (and OTLP/HTTP
senders configured to compress) commonly ship gzipped payloads, and
Plug's built-in JSON parser needs plain bytes to hand to the decoder.

Snappy is intentionally not handled here — Loki's snappy-framed
protobuf variant sits on a different content-type and needs its own
decoder anyway.

The reader forwards a `{:more, ...}` return unchanged. If the body
exceeds Plug.Parsers's configured length limit, the parser rejects
it with 413 — a decompressor cannot invent the missing bytes.
"""

@spec read_body(Plug.Conn.t(), keyword()) ::
{:ok, binary(), Plug.Conn.t()}
| {:more, binary(), Plug.Conn.t()}
| {:error, term()}
def read_body(conn, opts) do
case Plug.Conn.read_body(conn, opts) do
{:ok, body, conn} ->
case maybe_decompress(conn, body) do
{:ok, decompressed} -> {:ok, decompressed, conn}
{:error, _} = err -> err
end

other ->
other
end
end

defp maybe_decompress(conn, body) do
case Plug.Conn.get_req_header(conn, "content-encoding") do
["gzip"] -> gunzip(body)
_ -> {:ok, body}
end
end

defp gunzip(body) do
{:ok, :zlib.gunzip(body)}
rescue
_ -> {:error, :invalid_gzip}
end
end
121 changes: 121 additions & 0 deletions lib/pulso_web/controllers/loki_controller.ex
Original file line number Diff line number Diff line change
@@ -0,0 +1,121 @@
defmodule PulsoWeb.LokiController do
use PulsoWeb, :controller

alias Pulso.Auth
alias Pulso.Loki.Push
alias Pulso.Storage

@default_tenant "default"

@doc """
Loki push receiver at `POST /loki/api/v1/push`.

Accepts the JSON push wire format Alloy and Promtail can be configured
to emit. Gzip content encoding is decompressed transparently by
`PulsoWeb.CompressedBodyReader` before this action runs. Snappy-framed
protobuf (Alloy's default push_config) is not yet supported and is
rejected with 415.

Tenant, auth, and idempotency semantics match the OTLP receiver:
tenant from `X-Scope-OrgID` (default `"default"`), verified via
`Pulso.Auth`, and `Idempotency-Key` propagated to the storage backend
so a retry does not duplicate.

On success the response is `204 No Content` per the Loki push
convention — no body, no partial-success envelope. If any records
were rejected during decoding, the count is surfaced in the
`X-Pulso-Rejected-Records` response header so senders that want to
monitor decode drops can, without breaking clients that expect a
bodyless 204.
"""
@tenant_regex ~r/\A[A-Za-z0-9_.\-]{1,128}\z/

def push(conn, params) do
tenant = tenant_from(conn)
opts = append_opts(conn)

with :ok <- validate_content_type(conn),
:ok <- validate_tenant(tenant),
:ok <- Auth.verify(conn, tenant),
{records, rejected} = Push.decode(params),
:ok <- Storage.append(tenant, records, opts) do
conn
|> put_rejected_header(rejected)
|> send_resp(:no_content, "")
else
{:error, :unsupported_content_type} ->
conn
|> put_status(:unsupported_media_type)
|> json(%{error: "unsupported_content_type"})

{:error, {:invalid_tenant, _}} ->
conn
|> put_status(:bad_request)
|> json(%{error: "invalid_tenant"})

{:error, reason} when reason in [:missing_token, :invalid_token, :unknown_tenant] ->
conn
|> put_status(:unauthorized)
|> json(%{error: to_string(reason)})

{:error, {:encode_failed, _}} ->
conn
|> put_status(:bad_request)
|> json(%{error: "encode_failed"})

{:error, {:attribute_key_collision, _}} ->
conn
|> put_status(:bad_request)
|> json(%{error: "attribute_key_collision"})

{:error, reason} ->
conn
|> put_status(:internal_server_error)
|> json(%{error: inspect(reason)})
end
end

# We only speak JSON on this path today. A request with
# `application/x-protobuf` (Alloy/Promtail's default) is rejected
# explicitly rather than silently treated as JSON, which would
# otherwise 400 on the parser — the 415 tells the operator this is a
# missing feature, not a malformed body.
defp validate_content_type(conn) do
case Plug.Conn.get_req_header(conn, "content-type") do
[] -> :ok
[ct | _] -> if json_content_type?(ct), do: :ok, else: {:error, :unsupported_content_type}
end
end

defp json_content_type?(ct) do
ct
|> String.downcase()
|> String.starts_with?("application/json")
end

defp validate_tenant(tenant) do
if Regex.match?(@tenant_regex, tenant),
do: :ok,
else: {:error, {:invalid_tenant, tenant}}
end

defp tenant_from(conn) do
case Plug.Conn.get_req_header(conn, "x-scope-orgid") do
[tenant | _] when is_binary(tenant) and tenant != "" -> tenant
_ -> @default_tenant
end
end

defp append_opts(conn) do
case Plug.Conn.get_req_header(conn, "idempotency-key") do
[key | _] when is_binary(key) and byte_size(key) > 0 -> [idempotency_key: key]
_ -> []
end
end

defp put_rejected_header(conn, 0), do: conn

defp put_rejected_header(conn, rejected) do
Plug.Conn.put_resp_header(conn, "x-pulso-rejected-records", Integer.to_string(rejected))
end
end
3 changes: 2 additions & 1 deletion lib/pulso_web/endpoint.ex
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,8 @@ defmodule PulsoWeb.Endpoint do
plug Plug.Parsers,
parsers: [:urlencoded, :multipart, :json],
pass: ["*/*"],
json_decoder: Phoenix.json_library()
json_decoder: Phoenix.json_library(),
body_reader: {PulsoWeb.CompressedBodyReader, :read_body, []}

plug Plug.MethodOverride
plug Plug.Head
Expand Down
1 change: 1 addition & 0 deletions lib/pulso_web/router.ex
Original file line number Diff line number Diff line change
Expand Up @@ -10,5 +10,6 @@ defmodule PulsoWeb.Router do

post "/mcp", MCPController, :rpc
post "/v1/logs", OTLPController, :logs
post "/loki/api/v1/push", LokiController, :push
end
end
Loading
Loading