diff --git a/.credo.exs b/.credo.exs index 3d043d9..678b812 100644 --- a/.credo.exs +++ b/.credo.exs @@ -142,6 +142,16 @@ {Credo.Check.Warning.BoolOperationOnSameValues, []}, {Credo.Check.Warning.Dbg, []}, {Credo.Check.Warning.ExpensiveEmptyEnumCheck, []}, + # Pulso uses Elixir's built-in `JSON` module (Elixir 1.18+); the + # `Jason` dependency is only present transitively for optional deps + # of other libraries. Any direct `Jason.*` reference in our code + # should fail CI. See AGENTS.md > Conventions. + {Credo.Check.Warning.ForbiddenModule, + [ + modules: [ + {Jason, "Use Elixir's built-in `JSON` module instead of Jason (see AGENTS.md)."} + ] + ]}, {Credo.Check.Warning.IExPry, []}, {Credo.Check.Warning.IoInspect, []}, {Credo.Check.Warning.MissedMetadataKeyInLoggerConfig, []}, diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 02c7b86..92d7dbf 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -74,11 +74,15 @@ jobs: env: MIX_ENV: test PULSO_INTEGRATION: "1" - PULSO_MINIO_ENDPOINT: "http://localhost:9000" - PULSO_MINIO_BUCKET: "pulso" - PULSO_MINIO_REGION: "us-east-1" - PULSO_MINIO_ACCESS_KEY_ID: "minioadmin" - PULSO_MINIO_SECRET_ACCESS_KEY: "minioadmin" + # RustFS in docker-compose binds to these host ports by default (see + # mise/utilities/dev_instance_env.sh for the local per-worktree scheme). + PULSO_RUSTFS_API_PORT: "11100" + PULSO_RUSTFS_CONSOLE_PORT: "12100" + PULSO_S3_ENDPOINT: "http://localhost:11100" + PULSO_S3_BUCKET: "pulso" + PULSO_S3_REGION: "us-east-1" + PULSO_S3_ACCESS_KEY_ID: "rustfsadmin" + PULSO_S3_SECRET_ACCESS_KEY: "rustfsadmin" steps: - uses: actions/checkout@v4 @@ -105,29 +109,27 @@ jobs: restore-keys: | ${{ runner.os }}-mix-deps- - - name: Start MinIO + - name: Start RustFS via docker compose run: | - docker run -d --name minio -p 9000:9000 \ - -e MINIO_ROOT_USER=minioadmin \ - -e MINIO_ROOT_PASSWORD=minioadmin \ - quay.io/minio/minio:latest server /data - for _ in $(seq 1 30); do - if curl -sf http://localhost:9000/minio/health/live > /dev/null; then - echo "MinIO is up" + docker compose up -d rustfs + for _ in $(seq 1 60); do + if curl -sf "http://localhost:${PULSO_RUSTFS_API_PORT}/" > /dev/null 2>&1 \ + || curl -sf -o /dev/null -w '%{http_code}' "http://localhost:${PULSO_RUSTFS_API_PORT}/" | grep -qE '^(200|403|404)$'; then + echo "RustFS is up" exit 0 fi sleep 1 done - echo "MinIO did not start in time" - docker logs minio + echo "RustFS did not start in time" + docker compose logs rustfs exit 1 - name: Create bucket env: - AWS_ACCESS_KEY_ID: minioadmin - AWS_SECRET_ACCESS_KEY: minioadmin - AWS_DEFAULT_REGION: us-east-1 - run: aws --endpoint-url http://localhost:9000 s3 mb s3://pulso + AWS_ACCESS_KEY_ID: ${{ env.PULSO_S3_ACCESS_KEY_ID }} + AWS_SECRET_ACCESS_KEY: ${{ env.PULSO_S3_SECRET_ACCESS_KEY }} + AWS_DEFAULT_REGION: ${{ env.PULSO_S3_REGION }} + run: aws --endpoint-url "http://localhost:${PULSO_RUSTFS_API_PORT}" s3 mb "s3://${PULSO_S3_BUCKET}" - name: Install dependencies run: mix deps.get diff --git a/.gitignore b/.gitignore index 86caa1f..66795bf 100644 --- a/.gitignore +++ b/.gitignore @@ -29,3 +29,9 @@ pulso-*.tar # shared object into priv/native/. Both are build artifacts. native/*/target/ /priv/native/ + +# Per-worktree dev-instance suffix written by mise/utilities/dev_instance_env.sh. +# Persisted inside .git/worktrees/*/pulso-dev-instance by design; the +# top-level fallback is tracked here as a safety net for setups where the +# git-scoped path is unwritable. +/.pulso-dev-instance diff --git a/AGENTS.md b/AGENTS.md index 9c0d9c3..bdf0031 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -80,7 +80,7 @@ Storage backend URLs are read from `config :pulso, Pulso.Loki, base_url: ...` an ## Conventions - **HTTP client**: use `Req`. Never `HTTPoison`, `Tesla`, `:httpc`, or `Finch` directly. -- **JSON**: `Jason` (Phoenix's configured library). +- **JSON**: use Elixir's built-in `JSON` module (Elixir 1.18+), never `Jason`. Phoenix's `:json_library` is set to `JSON` in `config/config.exs`, and a Credo rule (`Credo.Check.Warning.ForbiddenModule`) fails CI on any direct `Jason.*` reference. Jason may still appear as a transitive dep of `phoenix` or a dev dep, but no code in `lib/`, `config/`, or `test/` may call it. - **New backends** go under `Pulso.` (e.g. `Pulso.Mimir`, `Pulso.Tempo`), with the same read-only-first shape as `Pulso.Loki`. Every read function must accept a `:base_url` override in opts. - **MCP tools** live in `Pulso.MCP.Tools`. Each tool has an `inputSchema`, and its `call/2` clause returns `{:ok, [content_block]}` or `{:error, reason}`. Content blocks follow the MCP shape: `%{"type" => "text", "text" => "..."}`. - **Alerting** (when added): each rule is its own supervised process, cluster-wide singleton via Horde. Rules that require exactly-once firing route through `ra`. diff --git a/config/config.exs b/config/config.exs index 29688eb..99b6a8d 100644 --- a/config/config.exs +++ b/config/config.exs @@ -7,13 +7,21 @@ # General application configuration import Config +alias Pulso.Auth.Open + # Configure Elixir's Logger config :logger, :default_formatter, format: "$time $metadata[$level] $message\n", metadata: [:request_id] -# Use Jason for JSON parsing in Phoenix -config :phoenix, :json_library, Jason +# Use Elixir's built-in JSON module for Phoenix (avoids the Jason dependency; +# see AGENTS.md conventions). +config :phoenix, :json_library, JSON + +# Explicit auth default. Environment-specific configs override; prod requires +# a runtime override to `Pulso.Auth.SharedSecret` via runtime.exs — an unset +# release still raises rather than falling back to open access. +config :pulso, Pulso.Auth, module: Open # Configure the endpoint config :pulso, PulsoWeb.Endpoint, diff --git a/config/prod.exs b/config/prod.exs index 30393dd..385c668 100644 --- a/config/prod.exs +++ b/config/prod.exs @@ -3,6 +3,12 @@ import Config # Do not print debug messages in production config :logger, level: :info +# A sentinel that runtime.exs is expected to overwrite. If a release starts +# without runtime.exs having populated the real auth module, `Pulso.Auth` +# raises rather than serving requests with the Open (accept-everything) +# fallback that ships in config.exs. +config :pulso, Pulso.Auth, module: :must_configure_at_runtime + config :pulso, PulsoWeb.Endpoint, force_ssl: [ rewrite_on: [:x_forwarded_proto], diff --git a/config/runtime.exs b/config/runtime.exs index 5da5ce0..b6589fb 100644 --- a/config/runtime.exs +++ b/config/runtime.exs @@ -16,12 +16,90 @@ import Config # # Alternatively, you can use `mix phx.gen.release` to generate a `bin/server` # script that automatically sets the env var above. +alias Pulso.Auth.SharedSecret +alias Pulso.Storage.S3 + if System.get_env("PHX_SERVER") do config :pulso, PulsoWeb.Endpoint, server: true end config :pulso, PulsoWeb.Endpoint, http: [port: String.to_integer(System.get_env("PORT", "4000"))] +# Log storage adapter. Tests keep the in-memory adapter (see config/test.exs); +# dev and prod use the S3 adapter against any S3-compatible endpoint. RustFS +# runs locally via docker-compose.yml — the dev defaults below match its +# out-of-the-box credentials. Prod requires the env vars to be set explicitly. +case config_env() do + :dev -> + config :pulso, Pulso.Storage, adapter: S3 + + config :pulso, S3, + bucket: System.get_env("PULSO_S3_BUCKET", "pulso"), + # mise/utilities/dev_instance_env.sh sets PULSO_S3_ENDPOINT per worktree. + # The fallback matches the docker-compose default host port when mise + # is not in the loop. + endpoint: System.get_env("PULSO_S3_ENDPOINT", "http://localhost:11100"), + region: System.get_env("PULSO_S3_REGION", "us-east-1"), + access_key_id: System.get_env("PULSO_S3_ACCESS_KEY_ID", "rustfsadmin"), + secret_access_key: System.get_env("PULSO_S3_SECRET_ACCESS_KEY", "rustfsadmin"), + allow_http: System.get_env("PULSO_S3_ALLOW_HTTP", "true") in ["1", "true", "yes"] + + :prod -> + require_env = fn name -> + case System.get_env(name) do + value when is_binary(value) and value != "" -> + value + + _ -> + raise """ + environment variable #{name} is missing or empty. + Pulso.Storage.S3 requires bucket/region/credentials in prod. + """ + end + end + + # Tenant tokens. Expected shape: a JSON object mapping tenant name to + # "sha256$". Deployments compute the hash offline + # and store only the digest in env, never the plaintext token. + tokens = + case System.get_env("PULSO_TENANT_TOKENS") do + blob when is_binary(blob) and blob != "" -> + case JSON.decode(blob) do + {:ok, map} when is_map(map) -> + map + + {:ok, _} -> + raise "PULSO_TENANT_TOKENS must decode to a JSON object" + + {:error, reason} -> + raise "PULSO_TENANT_TOKENS is not valid JSON: #{inspect(reason)}" + end + + _ -> + raise """ + environment variable PULSO_TENANT_TOKENS is missing or empty. + Pulso.Auth.SharedSecret requires at least one tenant token in prod. + """ + end + + config :pulso, Pulso.Auth, + module: SharedSecret, + tokens: tokens + + config :pulso, Pulso.Storage, adapter: S3 + + config :pulso, S3, + bucket: require_env.("PULSO_S3_BUCKET"), + endpoint: System.get_env("PULSO_S3_ENDPOINT"), + region: require_env.("PULSO_S3_REGION"), + access_key_id: require_env.("PULSO_S3_ACCESS_KEY_ID"), + secret_access_key: require_env.("PULSO_S3_SECRET_ACCESS_KEY"), + allow_http: System.get_env("PULSO_S3_ALLOW_HTTP", "false") in ["1", "true", "yes"] + + :test -> + :noop +end + if config_env() == :prod do # The secret key base is used to sign/encrypt cookies and other secrets. # A default value is used in config/dev.exs and config/test.exs but you diff --git a/docker-compose.yml b/docker-compose.yml index 650dd88..608fc0f 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,39 +1,66 @@ # Local development stack for Pulso. # -# `docker compose up -d` gives you MinIO on http://localhost:9000 with a -# preseeded bucket named `pulso`. The console is on http://localhost:9001 -# (user/password: minioadmin/minioadmin). +# `docker compose up -d` gives you a RustFS (S3-compatible) server with a +# preseeded bucket named `pulso`. The published ports are derived from +# PULSO_RUSTFS_API_PORT / PULSO_RUSTFS_CONSOLE_PORT, which +# mise/utilities/dev_instance_env.sh scopes per git worktree so parallel +# worktrees do not collide. Defaults (9195 / 9198) apply when the vars are +# not set. Credentials for the console: rustfsadmin / rustfsadmin. +# +# RustFS is used instead of MinIO because MinIO's community edition is no +# longer actively maintained. Pulso only talks to the S3 API, so any +# S3-compatible backend works; RustFS is a drop-in that stays maintained. services: - minio: - image: quay.io/minio/minio:latest - command: server /data --console-address ":9001" + rustfs: + image: rustfs/rustfs:latest environment: - MINIO_ROOT_USER: minioadmin - MINIO_ROOT_PASSWORD: minioadmin + # RUSTFS_VOLUMES is required — without it RustFS exits at startup. We + # run single-volume in dev and CI; the 4-volume erasure layout that + # RustFS ships in its own compose expects four distinct physical disks + # and refuses to start when they resolve to the same st_dev (any + # laptop or CI runner). + RUSTFS_VOLUMES: /data/rustfs0 + RUSTFS_ADDRESS: 0.0.0.0:9000 + RUSTFS_CONSOLE_ADDRESS: 0.0.0.0:9001 + RUSTFS_CONSOLE_ENABLE: "true" + RUSTFS_ACCESS_KEY: rustfsadmin + RUSTFS_SECRET_KEY: rustfsadmin + # Bind to loopback only so the well-known dev credentials cannot be reached + # from another host on the LAN. Host ports come from the per-worktree env + # to keep multiple checkouts from fighting over 9000/9001. ports: - - "9000:9000" - - "9001:9001" + - "127.0.0.1:${PULSO_RUSTFS_API_PORT:-11100}:9000" + - "127.0.0.1:${PULSO_RUSTFS_CONSOLE_PORT:-12100}:9001" volumes: - - minio-data:/data - healthcheck: - test: ["CMD", "curl", "-f", "http://localhost:9000/minio/health/live"] - interval: 5s - timeout: 3s - retries: 10 + - rustfs-data:/data - minio-init: - image: quay.io/minio/mc:latest + rustfs-init: + image: amazon/aws-cli:2 depends_on: - minio: - condition: service_healthy - entrypoint: > - /bin/sh -c " - mc alias set local http://minio:9000 minioadmin minioadmin; - mc mb --ignore-existing local/pulso; - mc anonymous set none local/pulso; - echo 'pulso bucket ready'; - " + - rustfs + environment: + AWS_ACCESS_KEY_ID: rustfsadmin + AWS_SECRET_ACCESS_KEY: rustfsadmin + AWS_DEFAULT_REGION: us-east-1 + entrypoint: + - /bin/sh + - -c + - | + for i in $$(seq 1 60); do + if aws --endpoint-url http://rustfs:9000 s3 mb s3://pulso 2>/dev/null; then + echo "pulso bucket created" + exit 0 + fi + if aws --endpoint-url http://rustfs:9000 s3api head-bucket --bucket pulso 2>/dev/null; then + echo "pulso bucket already exists" + exit 0 + fi + echo "waiting for rustfs..." + sleep 2 + done + echo "gave up waiting for rustfs" >&2 + exit 1 volumes: - minio-data: + rustfs-data: diff --git a/lib/pulso/application.ex b/lib/pulso/application.ex index 3ac731a..a671c7e 100644 --- a/lib/pulso/application.ex +++ b/lib/pulso/application.ex @@ -9,13 +9,12 @@ defmodule Pulso.Application do @impl true def start(_type, _args) do - children = [ - PulsoWeb.Telemetry, - {DNSCluster, query: Application.get_env(:pulso, :dns_cluster_query) || :ignore}, - {Phoenix.PubSub, name: Pulso.PubSub}, - Memory, - PulsoWeb.Endpoint - ] + children = + [ + PulsoWeb.Telemetry, + {DNSCluster, query: Application.get_env(:pulso, :dns_cluster_query) || :ignore}, + {Phoenix.PubSub, name: Pulso.PubSub} + ] ++ storage_children() ++ [PulsoWeb.Endpoint] # See https://elixir.hexdocs.pm/Supervisor.html # for other strategies and supported options @@ -30,4 +29,21 @@ defmodule Pulso.Application do PulsoWeb.Endpoint.config_change(changed, removed) :ok end + + # `Memory` only runs when it is the configured adapter — in `mix test` no + # adapter is set, so `Pulso.Storage.adapter/0` falls back to it. In dev and + # prod the S3 adapter is configured and Memory would just be dead weight. + defp storage_children do + adapter = + case Application.get_env(:pulso, Pulso.Storage) do + nil -> nil + env -> Keyword.get(env, :adapter) + end + + case adapter do + nil -> [Memory] + Memory -> [Memory] + _ -> [] + end + end end diff --git a/lib/pulso/auth.ex b/lib/pulso/auth.ex new file mode 100644 index 0000000..c568de4 --- /dev/null +++ b/lib/pulso/auth.ex @@ -0,0 +1,60 @@ +defmodule Pulso.Auth do + @moduledoc """ + Boundary for tenant authorization checks. + + `verify/2` is called from ingest and read paths before Pulso attributes a + request to a tenant. Implementations decide whether the caller is allowed + to speak for the tenant they named; they do **not** decide which tenant a + request belongs to (that is a routing concern, e.g. `X-Scope-OrgID`). + + Two implementations ship in-tree: + + * `Pulso.Auth.Open` — accepts everything. Default for dev and test. + * `Pulso.Auth.SharedSecret` — requires an `Authorization: Bearer ` + header that matches a per-tenant token from application env. Default + for prod. + + The active implementation is read from `Application.get_env(:pulso, + Pulso.Auth)[:module]` at call time so tests can swap it without + recompiling. + """ + + alias Plug.Conn + + @type reason :: :missing_token | :invalid_token | :unknown_tenant | term() + + @callback verify(Conn.t(), tenant :: String.t()) :: :ok | {:error, reason()} + + @spec verify(Conn.t(), String.t()) :: :ok | {:error, reason()} + def verify(conn, tenant) when is_binary(tenant) do + module().verify(conn, tenant) + end + + @spec module() :: module() + def module do + case Application.get_env(:pulso, __MODULE__) do + nil -> + raise """ + Pulso.Auth is not configured. This should never happen — config/config.exs + sets `Pulso.Auth.Open` as the compile-time default. Refusing to accept + traffic rather than fail open. + """ + + env -> + case Keyword.get(env, :module) do + nil -> + raise "Pulso.Auth :module key is not set. See config/config.exs." + + :must_configure_at_runtime -> + raise """ + Pulso.Auth is still set to the prod sentinel `:must_configure_at_runtime`. + runtime.exs must set `config :pulso, Pulso.Auth, module: Pulso.Auth.SharedSecret, tokens: %{...}` + before the app accepts traffic. + """ + + mod when is_atom(mod) -> + mod + end + end + end +end diff --git a/lib/pulso/auth/open.ex b/lib/pulso/auth/open.ex new file mode 100644 index 0000000..d19c02e --- /dev/null +++ b/lib/pulso/auth/open.ex @@ -0,0 +1,14 @@ +defmodule Pulso.Auth.Open do + @moduledoc """ + No-auth implementation of `Pulso.Auth`. Accepts every request. + + Default for dev and test where the loopback-bound ingest port makes a + real credential check net negative. Never suitable for a prod deployment + reachable from the network. + """ + + @behaviour Pulso.Auth + + @impl Pulso.Auth + def verify(_conn, _tenant), do: :ok +end diff --git a/lib/pulso/auth/shared_secret.ex b/lib/pulso/auth/shared_secret.ex new file mode 100644 index 0000000..fdf0a3a --- /dev/null +++ b/lib/pulso/auth/shared_secret.ex @@ -0,0 +1,67 @@ +defmodule Pulso.Auth.SharedSecret do + @moduledoc """ + Per-tenant Bearer-token authentication. A minimal but real auth surface + suitable for prod until Pulso grows a real accounts service. + + Tokens are configured as: + + config :pulso, Pulso.Auth, + module: Pulso.Auth.SharedSecret, + tokens: %{"acme" => "sha256$", "beta" => "sha256$"} + + Values are of the form `"$"`. Only `sha256` is accepted today. + Comparison is constant-time (`Plug.Crypto.secure_compare/2`) so token + presence cannot be inferred from response timing. + + A tenant with no configured token is rejected as `:unknown_tenant`. Empty + tokens are rejected as `:invalid_token` so an accidentally-blank env var + cannot turn into "auth off for this tenant". + """ + + @behaviour Pulso.Auth + + alias Plug.Conn + + @impl Pulso.Auth + def verify(conn, tenant) when is_binary(tenant) do + with {:ok, presented} <- extract_token(conn), + {:ok, stored_hash} <- fetch_stored_hash(tenant) do + compare_hash(presented, stored_hash) + end + end + + defp extract_token(conn) do + case Conn.get_req_header(conn, "authorization") do + ["Bearer " <> token] when byte_size(token) > 0 -> {:ok, token} + ["bearer " <> token] when byte_size(token) > 0 -> {:ok, token} + _ -> {:error, :missing_token} + end + end + + defp fetch_stored_hash(tenant) do + tokens = + case Application.get_env(:pulso, Pulso.Auth) do + nil -> %{} + env -> Keyword.get(env, :tokens, %{}) + end + + case Map.fetch(tokens, tenant) do + {:ok, "sha256$" <> hex} when byte_size(hex) == 64 -> {:ok, hex} + {:ok, _malformed} -> {:error, :invalid_token} + :error -> {:error, :unknown_tenant} + end + end + + defp compare_hash(presented, stored_hex) do + computed_hex = + :sha256 + |> :crypto.hash(presented) + |> Base.encode16(case: :lower) + + if Plug.Crypto.secure_compare(computed_hex, stored_hex) do + :ok + else + {:error, :invalid_token} + end + end +end diff --git a/lib/pulso/mcp.ex b/lib/pulso/mcp.ex index e24890f..83b97e5 100644 --- a/lib/pulso/mcp.ex +++ b/lib/pulso/mcp.ex @@ -13,27 +13,36 @@ defmodule Pulso.MCP do @protocol_version "2025-06-18" @server_info %{"name" => "pulso", "version" => "0.1.0"} + @type context :: %{optional(:conn) => Plug.Conn.t()} + @doc """ Dispatch a single JSON-RPC message. + `context` carries per-request state that individual tools need to run + authorization or other checks. Today only `:conn` is populated (from + `PulsoWeb.MCPController`); the shape is intentionally open for future + fields (tenant hint, feature flags). + Returns `{:reply, response}` for requests and `:noreply` for notifications (messages without an `id`). """ - @spec dispatch(map()) :: {:reply, map()} | :noreply - def dispatch(%{"method" => method} = msg) do + @spec dispatch(map(), context()) :: {:reply, map()} | :noreply + def dispatch(msg, context \\ %{}) + + def dispatch(%{"method" => method} = msg, context) when is_map(context) do id = Map.get(msg, "id") params = Map.get(msg, "params", %{}) - case {id, handle(method, params)} do + case {id, handle(method, params, context)} do {nil, _} -> :noreply {id, {:ok, result}} -> {:reply, ok(id, result)} {id, {:error, code, message}} -> {:reply, error(id, code, message)} end end - def dispatch(_), do: {:reply, error(nil, -32_600, "Invalid Request")} + def dispatch(_, _), do: {:reply, error(nil, -32_600, "Invalid Request")} - defp handle("initialize", _params) do + defp handle("initialize", _params, _context) do {:ok, %{ "protocolVersion" => @protocol_version, @@ -42,14 +51,14 @@ defmodule Pulso.MCP do }} end - defp handle("tools/list", _params) do + defp handle("tools/list", _params, _context) do {:ok, %{"tools" => Tools.list()}} end - defp handle("tools/call", %{"name" => name} = params) do + defp handle("tools/call", %{"name" => name} = params, context) do arguments = Map.get(params, "arguments", %{}) - case Tools.call(name, arguments) do + case Tools.call(name, arguments, context) do {:ok, content} -> {:ok, %{"content" => content, "isError" => false}} @@ -62,8 +71,8 @@ defmodule Pulso.MCP do end end - defp handle("ping", _params), do: {:ok, %{}} - defp handle(_unknown, _params), do: {:error, -32_601, "Method not found"} + defp handle("ping", _params, _context), do: {:ok, %{}} + defp handle(_unknown, _params, _context), do: {:error, -32_601, "Method not found"} defp ok(id, result), do: %{"jsonrpc" => "2.0", "id" => id, "result" => result} diff --git a/lib/pulso/mcp/tools.ex b/lib/pulso/mcp/tools.ex index ef0b659..5c7c620 100644 --- a/lib/pulso/mcp/tools.ex +++ b/lib/pulso/mcp/tools.ex @@ -7,6 +7,7 @@ defmodule Pulso.MCP.Tools do `docs/architecture.md`. """ + alias Pulso.Auth alias Pulso.Record.Log alias Pulso.Storage @@ -43,8 +44,10 @@ defmodule Pulso.MCP.Tools do @spec list() :: [map()] def list, do: @tools - @spec call(String.t(), map()) :: {:ok, [map()]} | {:error, term()} - def call("query_logs", %{"tenant" => tenant} = args) when is_binary(tenant) do + @spec call(String.t(), map(), Pulso.MCP.context()) :: {:ok, [map()]} | {:error, term()} + def call(name, args, context \\ %{}) + + def call("query_logs", %{"tenant" => tenant} = args, context) when is_binary(tenant) do opts = [] |> put_opt(:start_ts, args["start_ts_ns"]) @@ -52,13 +55,32 @@ defmodule Pulso.MCP.Tools do |> put_opt(:limit, args["limit"]) |> put_opt(:service, args["service"]) - with {:ok, records} <- Storage.query(tenant, opts) do - {:ok, [%{"type" => "text", "text" => Jason.encode!(Enum.map(records, &encode_record/1))}]} + with :ok <- verify(context, tenant), + {:ok, records} <- Storage.query(tenant, opts) do + {:ok, [%{"type" => "text", "text" => JSON.encode!(Enum.map(records, &encode_record/1))}]} end end - def call("query_logs", _args), do: {:error, {:invalid_arguments, "tenant is required"}} - def call(name, _args), do: {:error, {:unknown_tool, name}} + def call("query_logs", _args, _context), do: {:error, {:invalid_arguments, "tenant is required"}} + def call(name, _args, _context), do: {:error, {:unknown_tool, name}} + + # Every tool that names a tenant runs it through `Pulso.Auth.verify/2`. + # MCP is the same JSON-RPC transport for read and write; without this hop + # the ingest boundary's auth check would be bypassable via the read path. + # + # A missing conn falls through to a fresh `%Plug.Conn{}`. When the active + # auth module is `Pulso.Auth.Open` (dev/test default) that still returns + # :ok. When it is `Pulso.Auth.SharedSecret` (prod) it fails, closed — + # there is no in-process caller that legitimately reaches this path + # without a conn under real auth. + defp verify(context, tenant) do + conn = Map.get(context, :conn) || %Plug.Conn{} + + case Auth.verify(conn, tenant) do + :ok -> :ok + {:error, reason} -> {:error, {:unauthorized, reason}} + end + end defp put_opt(opts, _key, nil), do: opts defp put_opt(opts, key, value), do: Keyword.put(opts, key, value) diff --git a/lib/pulso/otlp/logs.ex b/lib/pulso/otlp/logs.ex index b324c27..b5e2b90 100644 --- a/lib/pulso/otlp/logs.ex +++ b/lib/pulso/otlp/logs.ex @@ -13,73 +13,102 @@ defmodule Pulso.OTLP.Logs do alias Pulso.Record.Log @doc """ - Decode a parsed JSON payload. Returns the flat list of log records; malformed - entries are skipped rather than aborting the whole batch. + Decode a parsed JSON payload. + + Returns `{records, rejected}`: + + * `records` — the flat list of `Pulso.Record.Log` that survived decoding. + * `rejected` — the count of `LogRecord` entries that could not be + decoded (missing or malformed `timeUnixNano`, wrong shape, etc.). + The OTLP receiver surfaces this to the sender via + `ExportLogsPartialSuccess.rejected_log_records` per the OTLP spec so + the sender knows some records did not make it into storage. """ - @spec decode(map()) :: [Log.t()] + @spec decode(map()) :: {[Log.t()], non_neg_integer()} def decode(%{"resourceLogs" => resource_logs}) when is_list(resource_logs) do - Enum.flat_map(resource_logs, &decode_resource_logs/1) + resource_logs + |> Enum.reduce({[], 0}, fn rl, {records, rejected} -> + {rl_records, rl_rejected} = decode_resource_logs(rl) + {[rl_records | records], rejected + rl_rejected} + end) + |> then(fn {records, rejected} -> {records |> Enum.reverse() |> List.flatten(), rejected} end) end - def decode(_), do: [] + # A top-level shape that is not an ExportLogsServiceRequest is not a valid + # OTLP body at all — treat every record as rejected (well, zero counted, + # since we can't count what we couldn't parse) and return empty. + def decode(_), do: {[], 0} defp decode_resource_logs(%{"scopeLogs" => scope_logs} = resource_logs) when is_list(scope_logs) do resource_attrs = attributes(resource_logs["resource"]) service = resource_attrs["service.name"] - Enum.flat_map(scope_logs, fn scope_logs -> - records = scope_logs["logRecords"] || [] - Enum.flat_map(records, &decode_log_record(&1, resource_attrs, service)) - end) + scope_logs + |> Enum.reduce({[], 0}, fn sl, acc -> decode_scope_logs(sl, resource_attrs, service, acc) end) + |> then(fn {records, rejected} -> {records |> Enum.reverse() |> List.flatten(), rejected} end) end - defp decode_resource_logs(_), do: [] + defp decode_resource_logs(_), do: {[], 0} - defp decode_log_record(%{} = record, resource_attrs, service) do - case timestamp(record["timeUnixNano"]) do - {:ok, ts} -> - [ - %Log{ - timestamp_ns: ts, - observed_timestamp_ns: nano(record["observedTimeUnixNano"]), - severity_number: record["severityNumber"], - severity_text: record["severityText"], - service: service, - body: any_value(record["body"]), - trace_id: nil_if_empty(record["traceId"]), - span_id: nil_if_empty(record["spanId"]), - attributes: attributes(record), - resource: resource_attrs - } - ] - - _ -> - [] - end - end + defp decode_scope_logs(scope_logs, resource_attrs, service, {records, rejected}) do + raw = scope_logs["logRecords"] || [] - defp decode_log_record(_, _, _), do: [] + {sl_records, sl_rejected} = + Enum.reduce(raw, {[], 0}, fn record, acc -> + decode_and_collect(record, resource_attrs, service, acc) + end) - defp timestamp(nil), do: :error + {[Enum.reverse(sl_records) | records], rejected + sl_rejected} + end - defp timestamp(value) do - case nano(value) do - nil -> :error - ns -> {:ok, ns} + defp decode_and_collect(record, resource_attrs, service, {rs, rj}) do + case decode_log_record(record, resource_attrs, service) do + {:ok, r} -> {[r | rs], rj} + :error -> {rs, rj + 1} end end - defp nano(nil), do: nil - defp nano(value) when is_integer(value), do: value + # Per the OTLP logs data model, both `time_unix_nano` and + # `observed_time_unix_nano` MAY be absent, and a value of 0 explicitly + # means "unknown". The decoder preserves the caller's intent verbatim + # (nil for absent or 0). The storage adapter fills a stored timestamp + # later so that (a) the caller_content_hash is stable across retries + # even when both timestamps are absent, and (b) records with only an + # observed time are not lost or misplaced at the Unix epoch. + defp decode_log_record(%{} = record, resource_attrs, service) do + {:ok, + %Log{ + timestamp_ns: nano_or_nil(record["timeUnixNano"]), + observed_timestamp_ns: nano_or_nil(record["observedTimeUnixNano"]), + severity_number: record["severityNumber"], + severity_text: record["severityText"], + service: service, + body: any_value(record["body"]), + trace_id: nil_if_empty(record["traceId"]), + span_id: nil_if_empty(record["spanId"]), + attributes: attributes(record), + resource: resource_attrs + }} + end + + defp decode_log_record(_, _, _), do: :error + + # OTLP: a `*_unix_nano` value of 0 signals "unknown", identical in + # meaning to the field being absent. Fold both into nil so downstream + # code has one shape to reason about. + defp nano_or_nil(nil), do: nil + defp nano_or_nil(0), do: nil + defp nano_or_nil(value) when is_integer(value), do: value - defp nano(value) when is_binary(value) do + defp nano_or_nil(value) when is_binary(value) do case Integer.parse(value) do + {0, ""} -> nil {int, ""} -> int _ -> nil end end - defp nano(_), do: nil + defp nano_or_nil(_), do: nil defp attributes(%{"attributes" => kvs}) when is_list(kvs) do Map.new(kvs, fn diff --git a/lib/pulso/record/log.ex b/lib/pulso/record/log.ex index b764343..a24aefa 100644 --- a/lib/pulso/record/log.ex +++ b/lib/pulso/record/log.ex @@ -7,7 +7,12 @@ defmodule Pulso.Record.Log do everything else lives in `attributes`. """ - @enforce_keys [:timestamp_ns] + # `timestamp_ns` may be nil at decode time — the OTLP spec allows both + # `time_unix_nano` and `observed_time_unix_nano` to be absent, and it + # explicitly treats a value of 0 as "unknown". The storage adapter fills + # in a stored timestamp from the observed value or the wall clock; the + # nil is what lets a retry keep an idempotent fingerprint (the caller + # sent the same batch with no timestamps, so both retries hash the same). defstruct [ :timestamp_ns, :observed_timestamp_ns, @@ -22,7 +27,7 @@ defmodule Pulso.Record.Log do ] @type t :: %__MODULE__{ - timestamp_ns: non_neg_integer(), + timestamp_ns: non_neg_integer() | nil, observed_timestamp_ns: non_neg_integer() | nil, severity_number: non_neg_integer() | nil, severity_text: String.t() | nil, diff --git a/lib/pulso/storage.ex b/lib/pulso/storage.ex index 765a420..8f8a5f8 100644 --- a/lib/pulso/storage.ex +++ b/lib/pulso/storage.ex @@ -3,15 +3,26 @@ defmodule Pulso.Storage do Behaviour for log storage backends and the runtime dispatcher. The active adapter is read from application configuration at call time so it - can be swapped in tests without recompiling. In step 1 the default adapter is - `Pulso.Storage.Memory`; step 2 will introduce an S3-backed adapter and demote - Memory to a test-only backend. + can be swapped in tests without recompiling. In dev and prod the default + adapter is `Pulso.Storage.S3` (step 2). Tests fall back to + `Pulso.Storage.Memory` because no adapter is configured. Step 3 will replace + the flat NDJSON layout of the S3 adapter with columnar segments and a + manifest that supports conditional writes. """ alias Pulso.Record.Log alias Pulso.Storage.Memory @type tenant :: String.t() + @type append_opts :: [ + # Opt-in idempotency: two `append` calls with the same tenant and + # the same idempotency_key resolve to the same underlying object + # so a retry does not duplicate. Callers who want distinct writes + # for identical payloads (e.g. two producers with genuinely + # different events that happen to serialize the same) simply + # omit the key. + {:idempotency_key, String.t()} + ] @type query_opts :: [ {:start_ts, non_neg_integer()} | {:end_ts, non_neg_integer()} @@ -19,11 +30,11 @@ defmodule Pulso.Storage do | {:service, String.t()} ] - @callback append(tenant, [Log.t()]) :: :ok | {:error, term()} + @callback append(tenant, [Log.t()], append_opts) :: :ok | {:error, term()} @callback query(tenant, query_opts) :: {:ok, [Log.t()]} | {:error, term()} - @spec append(tenant, [Log.t()]) :: :ok | {:error, term()} - def append(tenant, records), do: adapter().append(tenant, records) + @spec append(tenant, [Log.t()], append_opts) :: :ok | {:error, term()} + def append(tenant, records, opts \\ []), do: adapter().append(tenant, records, opts) @spec query(tenant, query_opts) :: {:ok, [Log.t()]} | {:error, term()} def query(tenant, opts \\ []), do: adapter().query(tenant, opts) diff --git a/lib/pulso/storage/memory.ex b/lib/pulso/storage/memory.ex index 927b3d8..8837cbe 100644 --- a/lib/pulso/storage/memory.ex +++ b/lib/pulso/storage/memory.ex @@ -2,10 +2,10 @@ defmodule Pulso.Storage.Memory do @moduledoc """ In-memory log storage backed by a public ETS table. - Step 1 default adapter. Meant to prove the ingest → storage → query spine - end to end without dragging in Rust, Parquet, or S3. Step 2 will replace it - with a real S3-backed adapter; this module will move to `test/support/` and - keep serving the test suite. + Test-only adapter as of step 2. Meant to prove the ingest → storage → query + spine end to end without dragging in Rust, Parquet, or S3. Dev and prod use + `Pulso.Storage.S3`; this module stays wired as the default in `mix test` + because `config/test.exs` sets no adapter override. Records for each tenant are kept in a private list-per-tenant, appended to as batches arrive and scanned linearly on query. That is deliberately naive: @@ -18,6 +18,7 @@ defmodule Pulso.Storage.Memory do use GenServer alias Pulso.Record.Log + alias Pulso.Storage.SortOrder @table __MODULE__ @@ -27,13 +28,13 @@ defmodule Pulso.Storage.Memory do end @impl Pulso.Storage - def append(tenant, records) when is_binary(tenant) and is_list(records) do - now = System.system_time(:nanosecond) - - normalized = - for %Log{} = record <- records do - %{record | observed_timestamp_ns: record.observed_timestamp_ns || now} - end + def append(tenant, records, _opts \\ []) when is_binary(tenant) and is_list(records) do + # Memory ignores :idempotency_key — it's a test adapter. Storage does + # NOT backfill timestamps: injecting `now` would defeat the retry + # story that the S3 adapter relies on for idempotency (see + # `Pulso.Storage.S3` docstring). Callers that need a wall-clock + # timestamp set it themselves at ingest. + normalized = records existing = case :ets.lookup(@table, tenant) do @@ -57,7 +58,7 @@ defmodule Pulso.Storage.Memory do records |> filter_by_time(Keyword.get(opts, :start_ts), Keyword.get(opts, :end_ts)) |> filter_by_service(Keyword.get(opts, :service)) - |> Enum.sort_by(& &1.timestamp_ns, :desc) + |> SortOrder.sort() |> take_limit(Keyword.get(opts, :limit)) {:ok, filtered} @@ -83,7 +84,13 @@ defmodule Pulso.Storage.Memory do defp filter_by_time(records, start_ts, end_ts) do Enum.filter(records, fn %Log{timestamp_ns: ts} -> - (start_ts == nil or ts >= start_ts) and (end_ts == nil or ts <= end_ts) + # A nil timestamp does not fit inside a time-bounded range. Elixir's + # term ordering puts atoms greater than numbers, so `nil >= 5` is + # true without an explicit guard — leaving nil-ts records leaking + # through every time filter. + is_integer(ts) and + (start_ts == nil or ts >= start_ts) and + (end_ts == nil or ts <= end_ts) end) end diff --git a/lib/pulso/storage/s3.ex b/lib/pulso/storage/s3.ex new file mode 100644 index 0000000..5a3a0da --- /dev/null +++ b/lib/pulso/storage/s3.ex @@ -0,0 +1,403 @@ +defmodule Pulso.Storage.S3 do + @moduledoc """ + S3-backed log storage. Step 2 adapter. + + Each `append/3` writes one NDJSON object under a tenant-scoped prefix + (`tenants//logs/-.ndjson`). Tenant isolation is + enforced by key construction: tenant names are validated against a + conservative charset so a batch cannot land outside its own prefix. + + The key suffix depends on whether the caller supplied an + `idempotency_key`: + + * **With `idempotency_key`** — the suffix is a deterministic hash of + `tenant || idempotency_key`. Two `append` calls with the same key + resolve to the same object, so a lost-response retry does not + duplicate. Two producers that pick the same idempotency key are + explicitly claiming "these are the same write" — semantics that + match Stripe's `Idempotency-Key` and RFC 9457. + * **Without `idempotency_key`** — the suffix mixes a truncated content + hash with random bytes. Distinct calls always produce distinct + objects (no accidental collapse of two identical-content batches). + Retries in this mode duplicate — callers who need dedup must opt in. + + `query/2` lists the tenant prefix, downloads every object, decodes NDJSON, + filters, and sorts. Order is shared with `Pulso.Storage.Memory` via + `Pulso.Storage.SortOrder`. A `NotFound` for a key that was listed but + disappeared before the fetch (concurrent retention, compaction, another + process deleting) is skipped rather than aborting the query. + + Known limits, deferred to step 3 (segments + manifest): + + * **Unbounded query work when `limit` is set.** Without per-batch time + metadata this adapter cannot safely skip objects: a batch written + recently may contain an old-timestamp record, so scanning every + object is required to preserve the sort semantics. Step 3's segment + manifest will carry min/max `timestamp_ns` per segment and let the + query short-circuit. + * **No columnar layout.** Records go on the wire as NDJSON, not Parquet. + + ## Key format stability + + Every object lives under a versioned prefix (`tenants//v1/…`). + Within one schema version, the exact suffix format is deliberately not + a public API. It is derived from `:erlang.term_to_binary(_, + [:deterministic])`, which is stable within an OTP release but is not + guaranteed to survive a major OTP upgrade (per the erts release notes, + the algorithm can change intentionally). Any change to the fingerprint, + the delimiter, or the sort-key width bumps the schema version — new + writes go to `v2/`, old objects stay at `v1/`, and a compaction job + migrates at its own pace. The reader can be taught to look at both + during the migration window. + """ + + @behaviour Pulso.Storage + + alias Pulso.ObjectStore + alias Pulso.Record.Log + alias Pulso.Storage.SortOrder + + @tenant_regex ~r/\A[A-Za-z0-9_.\-]{1,128}\z/ + # 20 decimal digits fits a u64 nanosecond timestamp (max ~1.84e19). Zero-padding + # keeps S3's UTF-8 list order chronological by write time (`sort_ns`). + @sort_key_width 20 + # 16 hex chars = 64 bits from SHA-256. Collision probability is negligible + # for the volumes any single tenant will produce in a step-2 adapter. + @content_hash_width 16 + + @impl Pulso.Storage + def append(tenant, records, opts \\ []) + + def append(tenant, [], _opts) when is_binary(tenant) do + # Validate even on empty so an adversarial tenant name is rejected on the + # first attempt, not only once a real record survives OTLP decoding. + validate_tenant(tenant) + end + + def append(tenant, records, opts) when is_binary(tenant) and is_list(records) do + idempotency_key = Keyword.get(opts, :idempotency_key) + + with :ok <- validate_tenant(tenant), + # Both the fingerprint AND the sort-key prefix are derived from the + # caller-provided records BEFORE `normalize/1` fills any wall-clock + # timestamps. A legitimate retry then produces the same object key + # in full — the sort_ns prefix and the idempotency suffix are both + # stable, so the second PUT overwrites the first as intended. + {:ok, caller_hash} = caller_content_hash(records), + sort_ns = caller_sort_ns(records), + {:ok, normalized} <- normalize(records), + {:ok, payload} <- encode(normalized) do + key = object_key(tenant, sort_ns, caller_hash, idempotency_key) + ObjectStore.put(config!(), key, payload) + end + end + + @impl Pulso.Storage + def query(tenant, opts) when is_binary(tenant) and is_list(opts) do + with :ok <- validate_tenant(tenant), + config = config!(), + {:ok, keys} <- ObjectStore.list(config, prefix(tenant)), + {:ok, records} <- fetch_records(config, keys) do + filtered = + records + |> filter_by_time(Keyword.get(opts, :start_ts), Keyword.get(opts, :end_ts)) + |> filter_by_service(Keyword.get(opts, :service)) + |> SortOrder.sort() + |> take_limit(Keyword.get(opts, :limit)) + + {:ok, filtered} + end + end + + # -- helpers ----------------------------------------------------------------- + + defp validate_tenant(tenant) do + if Regex.match?(@tenant_regex, tenant) do + :ok + else + {:error, {:invalid_tenant, tenant}} + end + end + + defp normalize(records) do + Enum.reduce_while(records, {:ok, []}, fn + %Log{} = record, {:ok, acc} -> + with {:ok, attrs} <- sanitize_map(record.attributes || %{}), + {:ok, resource} <- sanitize_map(record.resource || %{}) do + # Deliberately no wall-clock backfill here. Injecting `now` for a + # nil timestamp would make retries under the same idempotency key + # overwrite the first stored record with a later timestamp, so an + # already-acknowledged log would disappear from its original time + # range and re-emerge in a later one. Step 3's conditional PUT + # (write-only-if-absent) will allow first-write-wins backfill. + normalized = %{record | attributes: attrs, resource: resource} + {:cont, {:ok, [normalized | acc]}} + else + err -> {:halt, err} + end + end) + |> case do + {:ok, records} -> {:ok, Enum.reverse(records)} + err -> err + end + end + + # Coerce every attribute/resource map key to a string, recursively. OTLP + # decoding already produces string keys, but a caller building `%Log{}` + # directly (or a future backend surface) could hand us atoms or integers. + # If two logical keys coerce to the same string (`%{1 => a, "1" => b}`) we + # refuse — silently dropping either value would surprise a reader looking + # at either the original struct or the JSON-encoded record. + @doc false + @spec sanitize_map(map()) :: {:ok, map()} | {:error, {:attribute_key_collision, [String.t()]}} + def sanitize_map(map) when is_map(map) do + Enum.reduce_while(map, {:ok, %{}}, &insert_sanitized/2) + end + + defp insert_sanitized({k, v}, {:ok, acc}) do + string_key = stringify_key(k) + + if Map.has_key?(acc, string_key) do + {:halt, {:error, {:attribute_key_collision, [string_key]}}} + else + put_sanitized(acc, string_key, v) + end + end + + defp put_sanitized(acc, key, value) do + case sanitize_value(value) do + {:ok, sanitized} -> {:cont, {:ok, Map.put(acc, key, sanitized)}} + err -> {:halt, err} + end + end + + defp sanitize_value(v) when is_map(v), do: sanitize_map(v) + defp sanitize_value(v) when is_list(v), do: sanitize_list(v) + defp sanitize_value(v), do: {:ok, v} + + defp sanitize_list(list) do + list + |> Enum.reduce_while({:ok, []}, &prepend_sanitized/2) + |> case do + {:ok, sanitized} -> {:ok, Enum.reverse(sanitized)} + err -> err + end + end + + defp prepend_sanitized(value, {:ok, acc}) do + case sanitize_value(value) do + {:ok, sanitized} -> {:cont, {:ok, [sanitized | acc]}} + err -> {:halt, err} + end + end + + defp stringify_key(k) when is_binary(k), do: k + defp stringify_key(k) when is_atom(k), do: Atom.to_string(k) + defp stringify_key(k) when is_integer(k), do: Integer.to_string(k) + defp stringify_key(k), do: inspect(k) + + # Pick the smallest caller-supplied `timestamp_ns` for the object key + # prefix. Runs on records BEFORE normalization so a retry with identical + # caller input produces the same sort_ns. A record with no timestamp + # contributes 0, which parks the object at the head of the tenant + # listing — good enough for the fallback case and, importantly, + # deterministic across retries. + defp caller_sort_ns(records) do + records + |> Enum.map(fn %Log{timestamp_ns: ts} -> ts || 0 end) + |> Enum.min() + end + + defp encode(records) do + encoded = + Enum.reduce_while(records, {:ok, []}, fn record, {:ok, acc} -> + try do + line = JSON.encode!(Map.from_struct(record)) + {:cont, {:ok, [[line, "\n"] | acc]}} + rescue + e -> {:halt, {:error, {:encode_failed, e}}} + end + end) + + with {:ok, lines} <- encoded do + {:ok, lines |> Enum.reverse() |> IO.iodata_to_binary()} + end + end + + defp fetch_records(config, keys) do + # Accumulate batches as a list of lists then flatten once, so a large + # tenant does not pay O(n^2) list concatenation. A `:not_found` for a + # key that vanished after `list` is treated as "raced with a delete" and + # skipped; anything else halts the query so an outage is not hidden. + Enum.reduce_while(keys, {:ok, []}, fn key, {:ok, batches} -> + case ObjectStore.get(config, key) do + {:ok, blob} -> {:cont, {:ok, [decode(blob) | batches]}} + {:error, :not_found} -> {:cont, {:ok, batches}} + {:error, _} = err -> {:halt, err} + end + end) + |> case do + {:ok, batches} -> {:ok, batches |> Enum.reverse() |> List.flatten()} + {:error, _} = err -> err + end + end + + defp decode(blob) do + blob + |> String.split("\n", trim: true) + |> Enum.map(&decode_line/1) + end + + defp decode_line(line) do + map = JSON.decode!(line) + + %Log{ + timestamp_ns: Map.fetch!(map, "timestamp_ns"), + observed_timestamp_ns: map["observed_timestamp_ns"], + severity_number: map["severity_number"], + severity_text: map["severity_text"], + service: map["service"], + body: map["body"], + trace_id: map["trace_id"], + span_id: map["span_id"], + attributes: map["attributes"] || %{}, + resource: map["resource"] || %{} + } + end + + # Schema version segment. Baked into every object key so a future change + # to the key format (a new fingerprint algorithm, a different sort key + # width) can coexist with v1 objects rather than orphan them. Bump the + # version, keep readers that recognize both, and let a compaction job + # migrate the old prefix at leisure. + @schema_version "v1" + + defp prefix(tenant), do: "tenants/#{tenant}/#{@schema_version}/logs/" + + # `caller_hash` is a 16-hex fingerprint of the pre-normalization records + # from `caller_content_hash/1`. The pre-normalization form matters: + # `normalize/1` fills in a fresh wall-clock `observed_timestamp_ns` on + # every call, so hashing after normalization would make identical retries + # produce different keys even under an idempotency key. + @doc false + @spec object_key(String.t(), non_neg_integer(), String.t(), String.t() | nil) :: String.t() + def object_key(tenant, sort_ns, caller_hash, idempotency_key) when is_binary(tenant) and is_binary(caller_hash) do + suffix = + case idempotency_key do + <> when byte_size(key) > 0 -> + # Mixes tenant, key, and caller_hash. A client that accidentally + # reuses an idempotency key with different content produces a + # different object (no silent overwrite). Same content + same key + # collapses onto one object, which is the point of idempotency. + "idem-" <> stable_hash(tenant <> "\0" <> key <> "\0" <> caller_hash) + + _ -> + # No idempotency key: the caller accepts duplicates on retry. A + # random suffix guarantees distinct writes even when payload and + # timestamp collide. + "rand-" <> caller_hash <> "-" <> rand_hex() + end + + "#{prefix(tenant)}#{zero_pad(sort_ns)}-#{suffix}.ndjson" + end + + # Fingerprint of the raw caller records. Two properties hold: + # + # 1. Every field the caller controls (including `observed_timestamp_ns` + # when they set it) is included — a caller who legitimately changes + # that field on a "retry" is signalling a distinct write, and gets a + # distinct object. + # 2. The encoding is canonical across runtime, GC, and Jason versions. + # Two identical records always hash to the same byte sequence, in + # this process and in any future release. That is what makes cross- + # version idempotency safe. + # + # `:erlang.term_to_binary/2` with `:deterministic` gives us the canonical + # form for free within an OTP release: map keys are sorted, atoms and + # integers are encoded canonically, and the same term always produces + # the same bytes. That is stronger than a JSON encoder (map keys emit + # in `Map.to_list/1` order, which is not canonical and can shift when a + # small map promotes to a hash map). It is NOT guaranteed across major + # OTP upgrades — see the module docstring's "Key format stability" + # section for how a version bump migrates old objects when that happens. + @doc false + @spec caller_content_hash([Log.t()]) :: {:ok, String.t()} + def caller_content_hash(records) when is_list(records) do + canonical = + Enum.map(records, fn %Log{} = r -> + %{ + timestamp_ns: r.timestamp_ns, + observed_timestamp_ns: r.observed_timestamp_ns, + severity_number: r.severity_number, + severity_text: r.severity_text, + service: r.service, + body: normalize_hash_value(r.body), + trace_id: r.trace_id, + span_id: r.span_id, + attributes: r.attributes, + resource: r.resource + } + end) + + digest = + :sha256 + |> :crypto.hash(:erlang.term_to_binary(canonical, [:deterministic])) + |> Base.encode16(case: :lower) + |> binary_part(0, @content_hash_width) + + {:ok, digest} + end + + # `body` can arrive as a non-string term via OTLP AnyValue (int, bool, + # list). `:erlang.term_to_binary` handles all of these fine; nothing to + # normalize. This hook exists so a future callsite that wants a stable + # representation can add one without changing the fingerprint contract. + defp normalize_hash_value(v), do: v + + defp zero_pad(ns) when is_integer(ns) and ns >= 0 do + ns + |> Integer.to_string() + |> String.pad_leading(@sort_key_width, "0") + end + + defp stable_hash(bin) do + :sha256 + |> :crypto.hash(bin) + |> Base.encode16(case: :lower) + |> binary_part(0, @content_hash_width) + end + + defp rand_hex do + :crypto.strong_rand_bytes(8) |> Base.encode16(case: :lower) + end + + defp filter_by_time(records, nil, nil), do: records + + defp filter_by_time(records, start_ts, end_ts) do + Enum.filter(records, fn %Log{timestamp_ns: ts} -> + # A record with a nil timestamp has no place inside a time-bounded + # range. Elixir's term ordering puts atoms greater than numbers, so + # `nil >= 5` is true — without the `is_integer` guard, nil-ts + # records would leak through every time filter. + is_integer(ts) and + (start_ts == nil or ts >= start_ts) and + (end_ts == nil or ts <= end_ts) + end) + end + + defp filter_by_service(records, nil), do: records + defp filter_by_service(records, service), do: Enum.filter(records, &(&1.service == service)) + + defp take_limit(records, nil), do: records + defp take_limit(records, limit) when is_integer(limit) and limit > 0, do: Enum.take(records, limit) + + defp config! do + case Application.get_env(:pulso, __MODULE__) do + nil -> + raise "Pulso.Storage.S3 is not configured. Set `config :pulso, Pulso.Storage.S3, bucket: ..., endpoint: ..., region: ..., access_key_id: ..., secret_access_key: ..., allow_http: ...`" + + config -> + Map.new(config) + end + end +end diff --git a/lib/pulso/storage/sort_order.ex b/lib/pulso/storage/sort_order.ex new file mode 100644 index 0000000..9c130bc --- /dev/null +++ b/lib/pulso/storage/sort_order.ex @@ -0,0 +1,46 @@ +defmodule Pulso.Storage.SortOrder do + @moduledoc """ + Canonical sort order for a query result. Every storage adapter must apply + this so a client that switches adapters — or fans a query out to more than + one — cannot observe the tiebreaker changing. + + Primary key: `timestamp_ns` descending (newest first). + Tiebreakers, in order: `observed_timestamp_ns` desc, `trace_id`, `span_id`, + `body`. `nil` sorts last within its position. + """ + + alias Pulso.Record.Log + + @spec sort([Log.t()]) :: [Log.t()] + def sort(records) do + Enum.sort_by(records, &sort_key/1, &compare_desc/2) + end + + # Each string-shaped tiebreaker becomes `{presence_flag, value}` so a nil + # sorts *after* every real string in descending order, and — crucially — + # never collapses with an empty string. `1` for present, `0` for nil: + # descending order places `1 > 0` first, so real strings win the tie and + # a nil-vs-"" comparison sees a real difference in the first element of + # the pair. + defp sort_key(%Log{} = r) do + { + r.timestamp_ns || 0, + r.observed_timestamp_ns || 0, + presence_pair(r.trace_id), + presence_pair(r.span_id), + presence_pair(r.body) + } + end + + # `body` is `String.t() | nil` by the internal Log spec, but the OTLP + # decoder can produce integers, booleans, or lists when an incoming + # record uses those OTLP `AnyValue` shapes. Rather than crash the sort, + # coerce any non-string, non-nil term to a string so it still tiebreaks + # deterministically. `inspect/1` gives a stable, bounded representation + # for every Elixir term. + defp presence_pair(nil), do: {0, ""} + defp presence_pair(value) when is_binary(value), do: {1, value} + defp presence_pair(value), do: {1, inspect(value)} + + defp compare_desc(a, b), do: a >= b +end diff --git a/lib/pulso_web/controllers/mcp_controller.ex b/lib/pulso_web/controllers/mcp_controller.ex index ce475fb..63434fc 100644 --- a/lib/pulso_web/controllers/mcp_controller.ex +++ b/lib/pulso_web/controllers/mcp_controller.ex @@ -2,13 +2,15 @@ defmodule PulsoWeb.MCPController do use PulsoWeb, :controller def rpc(conn, params) when is_map(params) do - respond(conn, Pulso.MCP.dispatch(params)) + respond(conn, Pulso.MCP.dispatch(params, %{conn: conn})) end def rpc(conn, params) when is_list(params) do + context = %{conn: conn} + responses = params - |> Enum.map(&Pulso.MCP.dispatch/1) + |> Enum.map(&Pulso.MCP.dispatch(&1, context)) |> Enum.flat_map(fn {:reply, response} -> [response] :noreply -> [] diff --git a/lib/pulso_web/controllers/otlp_controller.ex b/lib/pulso_web/controllers/otlp_controller.ex index 23144b9..51f1720 100644 --- a/lib/pulso_web/controllers/otlp_controller.ex +++ b/lib/pulso_web/controllers/otlp_controller.ex @@ -1,6 +1,7 @@ defmodule PulsoWeb.OTLPController do use PulsoWeb, :controller + alias Pulso.Auth alias Pulso.OTLP.Logs alias Pulso.Storage @@ -10,23 +11,89 @@ defmodule PulsoWeb.OTLPController do OTLP/HTTP JSON logs receiver at `POST /v1/logs`. Tenant is taken from the `X-Scope-OrgID` header (Loki/Cortex convention) and - defaults to `"default"` when absent. The success response is the empty - `ExportLogsServiceResponse` object per the OTLP spec. + defaults to `"default"` when absent. The caller is then verified against + the configured `Pulso.Auth` module: `Pulso.Auth.Open` in dev/test accepts + everything; `Pulso.Auth.SharedSecret` in prod requires a bearer token. + + On success, the response body is the empty `ExportLogsServiceResponse` + object per the OTLP spec. """ + # A conservative tenant charset; mirrors `Pulso.Storage.S3`. Validating + # here (before auth) is what makes a `bad/name` tenant come back as 400 + # rather than 401 or 500 — the storage layer would still reject it, but + # by then the request has already spent an auth check on a value we know + # is invalid. + @tenant_regex ~r/\A[A-Za-z0-9_.\-]{1,128}\z/ + def logs(conn, params) do tenant = tenant_from(conn) - records = Logs.decode(params) + opts = append_opts(conn) + + with :ok <- validate_tenant(tenant), + :ok <- Auth.verify(conn, tenant), + {records, rejected} = Logs.decode(params), + :ok <- Storage.append(tenant, records, opts) do + # OTLP requires the receiver to signal partial success via a + # top-level `partialSuccess` block instead of a plain success. That + # lets the sender know some records did not make it into storage + # without turning the whole batch into a retry. + json(conn, partial_success_body(rejected, length(records))) + else + {: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)}) - case Storage.append(tenant, records) do - :ok -> json(conn, %{}) - {:error, reason} -> conn |> put_status(:internal_server_error) |> json(%{error: inspect(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 + 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 partial_success_body(0, _accepted), do: %{} + + defp partial_success_body(rejected, accepted) do + %{ + "partialSuccess" => %{ + "rejectedLogRecords" => rejected, + "errorMessage" => "#{rejected} log record(s) rejected; #{accepted} accepted. Cause: malformed logRecord entry." + } + } + end end diff --git a/mise.toml b/mise.toml index 4fe72bc..664d28c 100644 --- a/mise.toml +++ b/mise.toml @@ -1,3 +1,10 @@ [tools] erlang = "29.1" elixir = "1.20.4-otp-29" + +[env] +# Sourced by mise on every invocation inside the project (and each worktree). +# Exports PULSO_DEV_INSTANCE, PORT, PULSO_RUSTFS_API_PORT, +# PULSO_RUSTFS_CONSOLE_PORT, and PULSO_S3_ENDPOINT so multiple worktrees do +# not fight over the same TCP ports. +_.source = "{{config_root}}/mise/utilities/dev_instance_env.sh" diff --git a/mise/utilities/dev_instance_env.sh b/mise/utilities/dev_instance_env.sh new file mode 100644 index 0000000..6b96355 --- /dev/null +++ b/mise/utilities/dev_instance_env.sh @@ -0,0 +1,151 @@ +# Per-worktree port and instance suffix scoping. Sourced by mise so every +# `mise` invocation inside this project (or a linked worktree) exports a +# stable instance suffix and a set of ports derived from it. +# +# Adapted from tuist/tuist. The persisted suffix lives inside git's per- +# worktree state directory, so two worktrees pointing at the same repo get +# distinct suffixes and never collide on Phoenix, RustFS, or console ports. + +if [[ -n "${BASH_SOURCE[0]:-}" ]]; then + SCRIPT_PATH="${BASH_SOURCE[0]}" +elif [[ -n "${ZSH_VERSION:-}" ]]; then + SCRIPT_PATH="${(%):-%x}" +else + SCRIPT_PATH="${0}" +fi + +SCRIPT_DIR="$(cd "$(dirname "${SCRIPT_PATH}")" && pwd)" +PROJECT_ROOT="$(cd "${SCRIPT_DIR}/../.." && pwd)" +ROOT_INSTANCE_FILE="${PROJECT_ROOT}/.pulso-dev-instance" + +resolve_git_path() { + local target_name="$1" + local fallback_path="$2" + local git_path="" + + if command -v git >/dev/null 2>&1 && git -C "${PROJECT_ROOT}" rev-parse --is-inside-work-tree >/dev/null 2>&1; then + git_path="$( + git -C "${PROJECT_ROOT}" rev-parse --path-format=absolute --git-path "${target_name}" 2>/dev/null || + git -C "${PROJECT_ROOT}" rev-parse --git-path "${target_name}" 2>/dev/null || + true + )" + + if [[ -n "${git_path}" && "${git_path}" != /* ]]; then + git_path="${PROJECT_ROOT}/${git_path#./}" + fi + fi + + if [[ -n "${git_path}" ]]; then + printf '%s' "${git_path}" + else + printf '%s' "${fallback_path}" + fi +} + +INSTANCE_FILE="$(resolve_git_path "pulso-dev-instance" "${ROOT_INSTANCE_FILE}")" + +validate_suffix() { + local suffix="$1" + [[ "$suffix" =~ ^[0-9]+$ ]] || return 1 + (( suffix >= 1 && suffix <= 999 )) +} + +persist_suffix() { + local suffix="$1" + local target="$2" + + mkdir -p "$(dirname "${target}")" 2>/dev/null || return 1 + printf '%s' "${suffix}" | tee "${target}" >/dev/null 2>&1 +} + +collect_used_suffixes() { + # Suffixes already claimed by the main checkout and every linked worktree, + # so a freshly generated one can dodge collisions. + local common_dir="" f + if command -v git >/dev/null 2>&1 && git -C "${PROJECT_ROOT}" rev-parse --is-inside-work-tree >/dev/null 2>&1; then + common_dir="$(git -C "${PROJECT_ROOT}" rev-parse --path-format=absolute --git-common-dir 2>/dev/null || true)" + fi + [[ -n "${common_dir}" && -d "${common_dir}" ]] || return 0 + + for f in "${common_dir}/pulso-dev-instance" "${common_dir}"/worktrees/*/pulso-dev-instance; do + [[ -s "${f}" ]] || continue + [[ "${f}" -ef "${INSTANCE_FILE}" ]] 2>/dev/null && continue + tr -d '[:space:]' < "${f}" + printf '\n' + done +} + +generate_suffix() { + # Pick a suffix in [100, 999] not used by any other instance. Seed awk's RNG + # with the PID so worktrees bootstrapped within the same second diverge + # instead of sharing awk's default time(0) seed. + local used + used="$(collect_used_suffixes | tr '\n' ' ')" + awk -v used="${used}" -v seed="$$" ' + BEGIN { + srand(seed) + n = split(used, list, " ") + for (i = 1; i <= n; i++) taken[list[i]] = 1 + for (attempt = 0; attempt < 100000; attempt++) { + candidate = int(100 + rand() * 900) + if (!(candidate in taken)) { print candidate; exit 0 } + } + exit 1 + } + ' +} + +ensure_suffix() { + local suffix="" + + # This instance's own persisted suffix wins over everything else. Nested + # worktrees would otherwise inherit the parent's PULSO_DEV_INSTANCE. + if [[ -s "${INSTANCE_FILE}" ]]; then + suffix="$(tr -d '[:space:]' < "${INSTANCE_FILE}")" + elif [[ -n "${PULSO_DEV_INSTANCE:-}" ]] && + { [[ "${PULSO_DEV_INSTANCE_ROOT:-}" == "${PROJECT_ROOT}" ]] || [[ -z "${PULSO_DEV_INSTANCE_ROOT:-}" ]]; }; then + suffix="${PULSO_DEV_INSTANCE}" + elif [[ -s "${ROOT_INSTANCE_FILE}" ]]; then + suffix="$(tr -d '[:space:]' < "${ROOT_INSTANCE_FILE}")" + else + suffix="$(generate_suffix)" + fi + + validate_suffix "${suffix}" || { + echo "Invalid dev instance suffix '${suffix}'. Expected an integer between 1 and 999." >&2 + return 1 + } + + if ! persist_suffix "${suffix}" "${INSTANCE_FILE}"; then + if [[ "${INSTANCE_FILE}" != "${ROOT_INSTANCE_FILE}" ]] && + persist_suffix "${suffix}" "${ROOT_INSTANCE_FILE}"; then + INSTANCE_FILE="${ROOT_INSTANCE_FILE}" + else + echo "Failed to persist dev instance suffix '${suffix}'." >&2 + return 1 + fi + fi + + printf '%s' "${suffix}" +} + +suffix="$(ensure_suffix)" + +export PULSO_DEV_INSTANCE="${suffix}" +export PULSO_DEV_INSTANCE_ROOT="${PROJECT_ROOT}" + +# Phoenix endpoint. 4100..4999 — clear of the default Phoenix dev port (4000) +# and the ExUnit port (4002). +export PORT="$((4000 + suffix))" + +# RustFS via docker-compose. Ranges are 1000 apart so two suffixes N and +# N+3 cannot accidentally claim the same TCP port. Bases picked to avoid +# common developer defaults on macOS: ClickHouse (9000), Node.js debugger +# (9229), Prometheus (9090), Grafana (3000), gRPC dev (50051). +# API: 11100..11999 +# Console: 12100..12999 +export PULSO_RUSTFS_API_PORT="$((11000 + suffix))" +export PULSO_RUSTFS_CONSOLE_PORT="$((12000 + suffix))" + +# What the Elixir app reads. runtime.exs reads PULSO_S3_ENDPOINT verbatim. +export PULSO_S3_ENDPOINT="http://localhost:${PULSO_RUSTFS_API_PORT}" diff --git a/mix.exs b/mix.exs index f540b9d..3230c4e 100644 --- a/mix.exs +++ b/mix.exs @@ -42,7 +42,6 @@ defmodule Pulso.MixProject do {:phoenix, "~> 1.8.14"}, {:telemetry_metrics, "~> 1.0"}, {:telemetry_poller, "~> 1.0"}, - {:jason, "~> 1.2"}, {:dns_cluster, "~> 0.2.0"}, {:bandit, "~> 1.5"}, {:req, "~> 0.5"}, diff --git a/native/pulso_object_store/src/lib.rs b/native/pulso_object_store/src/lib.rs index 70d74af..a95cf87 100644 --- a/native/pulso_object_store/src/lib.rs +++ b/native/pulso_object_store/src/lib.rs @@ -11,7 +11,7 @@ use bytes::Bytes; use futures::TryStreamExt; use object_store::aws::AmazonS3Builder; use object_store::path::Path; -use object_store::{ObjectStore, PutPayload}; +use object_store::{Error as ObjectStoreError, ObjectStore, PutPayload}; use once_cell::sync::Lazy; use rustler::{Atom, Binary, Env, Error, NifResult, OwnedBinary}; use std::sync::Arc; @@ -46,6 +46,17 @@ fn nif_error(err: E) -> Error { Error::Term(Box::new(err.to_string())) } +// A NotFound response from the object store surfaces as the atom +// `:not_found` on the Elixir side. Everything else stays as a string +// message so the caller keeps the underlying context (permission denied, +// throttling, timeout, etc.). +fn map_object_store_error(err: ObjectStoreError) -> Error { + match err { + ObjectStoreError::NotFound { .. } => Error::Term(Box::new(atoms::not_found())), + other => Error::Term(Box::new(other.to_string())), + } +} + fn build_store(config: &StoreConfig) -> Result, Error> { let mut builder = AmazonS3Builder::new() .with_bucket_name(&config.bucket) @@ -87,7 +98,7 @@ fn get<'a>(env: Env<'a>, config: StoreConfig, key: String) -> NifResult<(Atom, B let obj = store.get(&path).await?; obj.bytes().await }) - .map_err(nif_error)?; + .map_err(map_object_store_error)?; let mut owned = OwnedBinary::new(bytes.len()) .ok_or_else(|| Error::Term(Box::new("failed to allocate binary")))?; @@ -103,7 +114,7 @@ fn delete(config: StoreConfig, key: String) -> NifResult { RUNTIME .block_on(async { store.delete(&path).await }) - .map_err(nif_error)?; + .map_err(map_object_store_error)?; Ok(atoms::ok()) } diff --git a/test/pulso/auth_test.exs b/test/pulso/auth_test.exs new file mode 100644 index 0000000..33aa470 --- /dev/null +++ b/test/pulso/auth_test.exs @@ -0,0 +1,93 @@ +defmodule Pulso.AuthTest do + use ExUnit.Case, async: false + + alias Pulso.Auth + alias Pulso.Auth.Open + alias Pulso.Auth.SharedSecret + + setup do + # Restore the compile-time default (module: Pulso.Auth.Open, set in + # config/config.exs) after each test so we do not poison subsequent + # tests that rely on the default. + on_exit(fn -> Application.put_env(:pulso, Auth, module: Open) end) + :ok + end + + defp conn(headers \\ []) do + Enum.reduce(headers, %Plug.Conn{}, fn {k, v}, c -> + Plug.Conn.put_req_header(c, k, v) + end) + end + + describe "Pulso.Auth.module/0" do + test "raises when nothing is configured — refuses to silently fail open" do + Application.delete_env(:pulso, Auth) + # config.exs sets an explicit default. Reaching this branch means + # someone deleted it; the correct response is a loud crash, not a + # quiet accept-everything. + assert_raise RuntimeError, ~r/Pulso.Auth is not configured/, fn -> + Auth.module() + end + end + + test "raises when the prod sentinel is still in place" do + Application.put_env(:pulso, Auth, module: :must_configure_at_runtime) + assert_raise RuntimeError, ~r/must_configure_at_runtime/, fn -> Auth.module() end + end + + test "returns the configured module" do + Application.put_env(:pulso, Auth, module: SharedSecret) + assert Auth.module() == SharedSecret + end + end + + describe "Pulso.Auth.Open" do + test "accepts any tenant on any conn" do + assert Open.verify(conn(), "acme") == :ok + assert Open.verify(conn(), "") == :ok + end + end + + describe "Pulso.Auth.SharedSecret" do + setup do + token = "the-secret" + hex = Base.encode16(:crypto.hash(:sha256, token), case: :lower) + Application.put_env(:pulso, Auth, tokens: %{"acme" => "sha256$#{hex}"}) + {:ok, token: token} + end + + test "accepts the correct bearer token", %{token: token} do + assert SharedSecret.verify(conn([{"authorization", "Bearer #{token}"}]), "acme") == :ok + end + + test "accepts lowercase 'bearer' too", %{token: token} do + assert SharedSecret.verify(conn([{"authorization", "bearer #{token}"}]), "acme") == :ok + end + + test "rejects a wrong token with :invalid_token" do + assert SharedSecret.verify(conn([{"authorization", "Bearer wrong"}]), "acme") == + {:error, :invalid_token} + end + + test "rejects a missing header with :missing_token" do + assert SharedSecret.verify(conn(), "acme") == {:error, :missing_token} + end + + test "rejects an empty bearer with :missing_token" do + assert SharedSecret.verify(conn([{"authorization", "Bearer "}]), "acme") == + {:error, :missing_token} + end + + test "rejects a tenant with no configured token as :unknown_tenant", %{token: token} do + assert SharedSecret.verify(conn([{"authorization", "Bearer #{token}"}]), "other") == + {:error, :unknown_tenant} + end + + test "rejects a malformed stored value as :invalid_token" do + Application.put_env(:pulso, Auth, tokens: %{"acme" => "plaintext-not-hashed"}) + + assert SharedSecret.verify(conn([{"authorization", "Bearer whatever"}]), "acme") == + {:error, :invalid_token} + end + end +end diff --git a/test/pulso/mcp/tools_test.exs b/test/pulso/mcp/tools_test.exs index 930debc..948ba00 100644 --- a/test/pulso/mcp/tools_test.exs +++ b/test/pulso/mcp/tools_test.exs @@ -1,6 +1,8 @@ defmodule Pulso.MCP.ToolsTest do use ExUnit.Case, async: false + alias Pulso.Auth.Open + alias Pulso.Auth.SharedSecret alias Pulso.MCP.Tools alias Pulso.Record.Log alias Pulso.Storage @@ -25,7 +27,7 @@ defmodule Pulso.MCP.ToolsTest do ]) assert {:ok, [%{"type" => "text", "text" => text}]} = Tools.call("query_logs", %{"tenant" => "acme"}) - assert [%{"body" => "two"}, %{"body" => "one"}] = Jason.decode!(text) + assert [%{"body" => "two"}, %{"body" => "one"}] = JSON.decode!(text) end test "query_logs applies service and limit filters" do @@ -39,10 +41,66 @@ defmodule Pulso.MCP.ToolsTest do assert {:ok, [%{"text" => text}]} = Tools.call("query_logs", %{"tenant" => "acme", "service" => "api", "limit" => 1}) - assert [%{"body" => "c", "service" => "api"}] = Jason.decode!(text) + assert [%{"body" => "c", "service" => "api"}] = JSON.decode!(text) end test "query_logs errors when tenant is missing" do assert {:error, {:invalid_arguments, _}} = Tools.call("query_logs", %{}) end + + describe "auth on the read path" do + setup do + hex = Base.encode16(:crypto.hash(:sha256, "the-token"), case: :lower) + + Application.put_env(:pulso, Pulso.Auth, + module: SharedSecret, + tokens: %{"acme" => "sha256$#{hex}"} + ) + + on_exit(fn -> + Application.put_env(:pulso, Pulso.Auth, module: Open) + end) + + :ok + end + + test "rejects a query_logs call with no bearer token" do + Storage.append("acme", [%Log{timestamp_ns: 1}]) + + assert {:error, {:unauthorized, :missing_token}} = + Tools.call("query_logs", %{"tenant" => "acme"}, %{conn: %Plug.Conn{}}) + end + + test "fails closed when no conn is passed at all" do + # Previous version had a "no conn = allow" fallback for in-process + # callers. That was a bypass: any code path that forgot the context + # would silently read another tenant's data under shared-secret + # auth. Now the fallback constructs an empty %Plug.Conn{}, which + # SharedSecret.verify sees as :missing_token. + Storage.append("acme", [%Log{timestamp_ns: 1}]) + + assert {:error, {:unauthorized, :missing_token}} = + Tools.call("query_logs", %{"tenant" => "acme"}, %{}) + end + + test "rejects a query_logs call with a bad token" do + Storage.append("acme", [%Log{timestamp_ns: 1}]) + + conn = %Plug.Conn{} |> Plug.Conn.put_req_header("authorization", "Bearer wrong") + + assert {:error, {:unauthorized, :invalid_token}} = + Tools.call("query_logs", %{"tenant" => "acme"}, %{conn: conn}) + end + + test "accepts a query_logs call with the correct token" do + Storage.append("acme", [%Log{timestamp_ns: 1, body: "ok"}]) + + conn = %Plug.Conn{} |> Plug.Conn.put_req_header("authorization", "Bearer the-token") + + assert {:ok, [%{"text" => text}]} = + Tools.call("query_logs", %{"tenant" => "acme"}, %{conn: conn}) + + assert [%{"body" => "ok"}] = JSON.decode!(text) + end + end end diff --git a/test/pulso/object_store_test.exs b/test/pulso/object_store_test.exs index 77076a6..d0fa72a 100644 --- a/test/pulso/object_store_test.exs +++ b/test/pulso/object_store_test.exs @@ -1,8 +1,8 @@ defmodule Pulso.ObjectStoreTest do - # This module talks to a live MinIO endpoint via the Rust NIF. - # It only runs when the caller opts in with PULSO_INTEGRATION=1 (see - # test/test_helper.exs), so plain `mix test` on a machine without - # docker-compose still passes. + # This module talks to a live S3-compatible endpoint (RustFS by default via + # docker-compose.yml) through the Rust NIF. It only runs when the caller + # opts in with PULSO_INTEGRATION=1 (see test/test_helper.exs), so plain + # `mix test` on a machine without docker-compose still passes. use ExUnit.Case, async: false @@ -12,11 +12,11 @@ defmodule Pulso.ObjectStoreTest do setup do config = %{ - bucket: System.get_env("PULSO_MINIO_BUCKET", "pulso"), - endpoint: System.get_env("PULSO_MINIO_ENDPOINT", "http://localhost:9000"), - region: System.get_env("PULSO_MINIO_REGION", "us-east-1"), - access_key_id: System.get_env("PULSO_MINIO_ACCESS_KEY_ID", "minioadmin"), - secret_access_key: System.get_env("PULSO_MINIO_SECRET_ACCESS_KEY", "minioadmin"), + bucket: System.get_env("PULSO_S3_BUCKET", "pulso"), + endpoint: System.get_env("PULSO_S3_ENDPOINT", "http://localhost:11100"), + region: System.get_env("PULSO_S3_REGION", "us-east-1"), + access_key_id: System.get_env("PULSO_S3_ACCESS_KEY_ID", "rustfsadmin"), + secret_access_key: System.get_env("PULSO_S3_SECRET_ACCESS_KEY", "rustfsadmin"), allow_http: true } diff --git a/test/pulso/otlp/logs_test.exs b/test/pulso/otlp/logs_test.exs index b17ce46..fe7f626 100644 --- a/test/pulso/otlp/logs_test.exs +++ b/test/pulso/otlp/logs_test.exs @@ -4,9 +4,9 @@ defmodule Pulso.OTLP.LogsTest do alias Pulso.OTLP.Logs alias Pulso.Record.Log - test "returns [] for a payload without resourceLogs" do - assert Logs.decode(%{}) == [] - assert Logs.decode(%{"resourceLogs" => "not-a-list"}) == [] + test "returns {[], 0} for a payload without resourceLogs" do + assert Logs.decode(%{}) == {[], 0} + assert Logs.decode(%{"resourceLogs" => "not-a-list"}) == {[], 0} end test "decodes a full record with resource, attributes, and body" do @@ -42,23 +42,27 @@ defmodule Pulso.OTLP.LogsTest do ] } - assert [ - %Log{ - timestamp_ns: 1_700_000_000_000_000_000, - observed_timestamp_ns: 1_700_000_000_000_000_001, - severity_number: 9, - severity_text: "INFO", - service: "api", - body: "hello", - trace_id: "abc", - span_id: "def", - attributes: %{"user.id" => "u1"}, - resource: %{"service.name" => "api", "deploy.env" => "prod"} - } - ] = Logs.decode(payload) + assert {[ + %Log{ + timestamp_ns: 1_700_000_000_000_000_000, + observed_timestamp_ns: 1_700_000_000_000_000_001, + severity_number: 9, + severity_text: "INFO", + service: "api", + body: "hello", + trace_id: "abc", + span_id: "def", + attributes: %{"user.id" => "u1"}, + resource: %{"service.name" => "api", "deploy.env" => "prod"} + } + ], 0} = Logs.decode(payload) end - test "skips records without a timestamp" do + test "preserves absent timestamps as nil so storage can backfill deterministically" do + # OTLP allows both timestamps to be absent, and a zero value means + # "unknown". The decoder does NOT invent a wall-clock value here — + # that would defeat the caller_content_hash for a retry. The storage + # adapter fills in a stored timestamp later. payload = %{ "resourceLogs" => [ %{ @@ -69,7 +73,77 @@ defmodule Pulso.OTLP.LogsTest do ] } - assert Logs.decode(payload) == [] + assert {[%Log{timestamp_ns: nil, observed_timestamp_ns: nil, body: "no ts"}], 0} = + Logs.decode(payload) + end + + test "preserves an observed-only timestamp verbatim" do + payload = %{ + "resourceLogs" => [ + %{ + "scopeLogs" => [ + %{ + "logRecords" => [ + %{ + "observedTimeUnixNano" => "42", + "body" => %{"stringValue" => "only observed"} + } + ] + } + ] + } + ] + } + + assert {[%Log{timestamp_ns: nil, observed_timestamp_ns: 42, body: "only observed"}], 0} = + Logs.decode(payload) + end + + test "folds a zero timestamp to nil per the OTLP spec" do + # `time_unix_nano: 0` means "unknown", identical in meaning to the + # field being absent. If the decoder left it as 0, the record would + # be stored at the Unix epoch and vanish from time-bounded queries. + payload = %{ + "resourceLogs" => [ + %{ + "scopeLogs" => [ + %{ + "logRecords" => [ + %{ + "timeUnixNano" => "0", + "observedTimeUnixNano" => "1700000000000000000", + "body" => %{"stringValue" => "zero ts"} + } + ] + } + ] + } + ] + } + + assert {[%Log{timestamp_ns: nil, observed_timestamp_ns: 1_700_000_000_000_000_000}], 0} = + Logs.decode(payload) + end + + test "counts a non-map logRecords entry as rejected" do + # A logRecord entry that is not a map is genuinely malformed — the + # decoder cannot invent fields, so it counts toward rejected. + payload = %{ + "resourceLogs" => [ + %{ + "scopeLogs" => [ + %{ + "logRecords" => [ + "not-a-map", + %{"timeUnixNano" => "1", "body" => %{"stringValue" => "ok"}} + ] + } + ] + } + ] + } + + assert {[%Log{body: "ok"}], 1} = Logs.decode(payload) end test "decodes AnyValue variants in attributes" do @@ -111,7 +185,7 @@ defmodule Pulso.OTLP.LogsTest do ] } - assert [%Log{attributes: attrs}] = Logs.decode(payload) + assert {[%Log{attributes: attrs}], 0} = Logs.decode(payload) assert attrs["s"] == "x" assert attrs["i"] == 42 assert attrs["b"] == true @@ -141,7 +215,7 @@ defmodule Pulso.OTLP.LogsTest do ] } - records = Logs.decode(payload) + assert {records, 0} = Logs.decode(payload) assert length(records) == 4 assert Enum.map(records, & &1.service) == ["a", "a", "a", "b"] assert Enum.map(records, & &1.timestamp_ns) == [1, 2, 3, 4] diff --git a/test/pulso/storage/memory_test.exs b/test/pulso/storage/memory_test.exs index b756c15..6c12aec 100644 --- a/test/pulso/storage/memory_test.exs +++ b/test/pulso/storage/memory_test.exs @@ -53,12 +53,42 @@ defmodule Pulso.Storage.MemoryTest do assert Enum.map(records, & &1.timestamp_ns) == [4, 3] end - test "populates observed_timestamp_ns when the record does not carry one" do - before_append = System.system_time(:nanosecond) - :ok = Storage.append("t", [record(1)]) - after_append = System.system_time(:nanosecond) + test "nil-timestamp records do not leak through time-bounded queries" do + # Guards a subtle Elixir term-ordering trap: nil >= 5 returns true + # because atoms sort above integers. Without an explicit is_integer + # guard in filter_by_time, nil-ts records would slip past every + # start_ts filter. + :ok = Storage.append("t", [%Log{timestamp_ns: nil}, record(100)]) - assert {:ok, [%Log{observed_timestamp_ns: observed}]} = Storage.query("t") - assert observed >= before_append and observed <= after_append + assert {:ok, [%Log{timestamp_ns: 100}]} = Storage.query("t", start_ts: 50) + assert {:ok, [%Log{timestamp_ns: 100}]} = Storage.query("t", end_ts: 200) + + # But an unbounded query still returns them. + assert {:ok, records} = Storage.query("t") + assert Enum.any?(records, &is_nil(&1.timestamp_ns)) + end + + test "stores records verbatim without wall-clock backfill" do + # Prior versions injected `now` when a timestamp was nil. That was + # dropped because injecting a fresh timestamp per call breaks the + # retry-idempotency story in the S3 adapter — the second PUT under + # the same idempotency key would overwrite the first with a later + # timestamp. Memory is the test-only adapter and mirrors the same + # contract for consistency. + :ok = Storage.append("t", [%Log{timestamp_ns: nil, observed_timestamp_ns: nil}]) + assert {:ok, [%Log{timestamp_ns: nil, observed_timestamp_ns: nil}]} = Storage.query("t") + end + + test "equal timestamps are broken by observed_timestamp_ns then trace_id" do + # Guard against a limit response depending on adapter-internal insertion + # order. Two adapters must sort ties the same way; both delegate to + # Pulso.Storage.SortOrder. + a = %Log{timestamp_ns: 10, observed_timestamp_ns: 100, trace_id: "aaa"} + b = %Log{timestamp_ns: 10, observed_timestamp_ns: 300, trace_id: "aaa"} + c = %Log{timestamp_ns: 10, observed_timestamp_ns: 200, trace_id: "bbb"} + + :ok = Storage.append("t", [a, b, c]) + assert {:ok, sorted} = Storage.query("t") + assert Enum.map(sorted, & &1.observed_timestamp_ns) == [300, 200, 100] end end diff --git a/test/pulso/storage/s3_test.exs b/test/pulso/storage/s3_test.exs new file mode 100644 index 0000000..2d8f414 --- /dev/null +++ b/test/pulso/storage/s3_test.exs @@ -0,0 +1,196 @@ +defmodule Pulso.Storage.S3Test do + # Round-trips logs through the real S3-compatible endpoint (RustFS via + # docker-compose). Only runs with PULSO_INTEGRATION=1; plain `mix test` + # skips it. See test/test_helper.exs. + + use ExUnit.Case, async: false + + alias Pulso.ObjectStore + alias Pulso.Record.Log + alias Pulso.Storage.S3 + + @moduletag :integration + + setup do + config = %{ + bucket: System.get_env("PULSO_S3_BUCKET", "pulso"), + endpoint: System.get_env("PULSO_S3_ENDPOINT", "http://localhost:11100"), + region: System.get_env("PULSO_S3_REGION", "us-east-1"), + access_key_id: System.get_env("PULSO_S3_ACCESS_KEY_ID", "rustfsadmin"), + secret_access_key: System.get_env("PULSO_S3_SECRET_ACCESS_KEY", "rustfsadmin"), + allow_http: true + } + + Application.put_env(:pulso, S3, config) + + tenant = "test-#{System.unique_integer([:positive])}" + + on_exit(fn -> + # The adapter writes objects under `tenants//v1/logs/`; clean up so + # a re-run starts empty. + case ObjectStore.list(config, "tenants/#{tenant}/v1/logs/") do + {:ok, keys} -> Enum.each(keys, &ObjectStore.delete(config, &1)) + _ -> :ok + end + end) + + {:ok, config: config, tenant: tenant} + end + + defp record(ts, opts \\ []) do + %Log{ + timestamp_ns: ts, + severity_text: Keyword.get(opts, :severity_text), + service: Keyword.get(opts, :service), + body: Keyword.get(opts, :body), + attributes: Keyword.get(opts, :attributes, %{}), + resource: Keyword.get(opts, :resource, %{}) + } + end + + test "append then query round-trips log records", %{tenant: tenant} do + assert :ok = + S3.append(tenant, [ + record(10, service: "api", body: "hello"), + record(20, service: "web", body: "world") + ]) + + assert {:ok, records} = S3.query(tenant, []) + assert Enum.map(records, & &1.timestamp_ns) == [20, 10] + assert Enum.map(records, & &1.service) == ["web", "api"] + assert Enum.map(records, & &1.body) == ["world", "hello"] + end + + test "records for one tenant are invisible to another", %{tenant: tenant, config: config} do + other = "test-other-#{System.unique_integer([:positive])}" + + on_exit(fn -> + case ObjectStore.list(config, "tenants/#{other}/v1/logs/") do + {:ok, keys} -> Enum.each(keys, &ObjectStore.delete(config, &1)) + _ -> :ok + end + end) + + assert :ok = S3.append(tenant, [record(1)]) + assert :ok = S3.append(other, [record(2)]) + + assert {:ok, [%Log{timestamp_ns: 1}]} = S3.query(tenant, []) + assert {:ok, [%Log{timestamp_ns: 2}]} = S3.query(other, []) + end + + test "filters by time range and service", %{tenant: tenant} do + assert :ok = + S3.append(tenant, [ + record(10, service: "api"), + record(20, service: "web"), + record(30, service: "api"), + record(40, service: "api") + ]) + + assert {:ok, records} = S3.query(tenant, start_ts: 15, end_ts: 35, service: "api") + assert Enum.map(records, & &1.timestamp_ns) == [30] + end + + test "applies limit", %{tenant: tenant} do + assert :ok = S3.append(tenant, [record(1), record(2), record(3), record(4)]) + assert {:ok, records} = S3.query(tenant, limit: 2) + assert length(records) == 2 + assert Enum.map(records, & &1.timestamp_ns) == [4, 3] + end + + test "stores records verbatim without injecting a wall-clock timestamp", %{tenant: tenant} do + # Wall-clock backfill would defeat retry idempotency (the second call + # under the same idempotency key would overwrite the first with a + # later observed_ts). Records that arrive without timestamps are + # stored as they came in; queries with time bounds naturally skip + # them, unbounded queries return them. + assert :ok = + S3.append(tenant, [ + %Log{timestamp_ns: nil, observed_timestamp_ns: nil, body: "no ts"} + ]) + + assert {:ok, [%Log{timestamp_ns: nil, observed_timestamp_ns: nil, body: "no ts"}]} = + S3.query(tenant, []) + end + + test "preserves NDJSON-hostile bodies through the round trip", %{tenant: tenant} do + tricky = "line1\nline2\t\"quoted\"\r\nline3" + assert :ok = S3.append(tenant, [record(1, body: tricky)]) + assert {:ok, [%Log{body: ^tricky}]} = S3.query(tenant, []) + end + + test "append with an empty batch is a no-op", %{tenant: tenant, config: config} do + assert :ok = S3.append(tenant, []) + assert {:ok, keys} = ObjectStore.list(config, "tenants/#{tenant}/v1/logs/") + assert keys == [] + end + + test "rejects tenant names that could escape the prefix" do + for bad <- ["../evil", "foo/bar", "foo bar", "", String.duplicate("a", 200)] do + assert {:error, {:invalid_tenant, ^bad}} = S3.append(bad, [record(1)]) + assert {:error, {:invalid_tenant, ^bad}} = S3.query(bad, []) + end + end + + test "same idempotency_key deduplicates identical retries", %{ + tenant: tenant, + config: config + } do + # Opt-in idempotency: a caller that wants retry safety passes the same + # idempotency_key on the retry. The object key becomes deterministic, + # so the second PUT overwrites the first with identical content and no + # duplicate appears at query time. + batch = [record(1, body: "same", service: "api")] + + assert :ok = S3.append(tenant, batch, idempotency_key: "req-1") + assert :ok = S3.append(tenant, batch, idempotency_key: "req-1") + + assert {:ok, keys} = ObjectStore.list(config, "tenants/#{tenant}/v1/logs/") + assert length(keys) == 1 + + assert {:ok, records} = S3.query(tenant, []) + assert length(records) == 1 + end + + test "without an idempotency_key identical batches remain distinct writes", %{ + tenant: tenant, + config: config + } do + # If a caller doesn't opt in, two producers with byte-identical payloads + # must not silently collapse — that would drop data. + batch = [record(1, body: "same", service: "api")] + + assert :ok = S3.append(tenant, batch) + assert :ok = S3.append(tenant, batch) + + assert {:ok, keys} = ObjectStore.list(config, "tenants/#{tenant}/v1/logs/") + assert length(keys) == 2 + end + + test "a key deleted after listing does not fail the query", %{ + tenant: tenant, + config: config + } do + assert :ok = S3.append(tenant, [record(1), record(2)]) + assert :ok = S3.append(tenant, [record(3)]) + + # Delete one of the objects between our own list and get, mimicking a + # compaction / retention job racing with a query. + assert {:ok, [first | _]} = ObjectStore.list(config, "tenants/#{tenant}/v1/logs/") + assert :ok = ObjectStore.delete(config, first) + + # Query should still return the surviving records, not error. + assert {:ok, remaining} = S3.query(tenant, []) + assert remaining != [] + end + + test "equal timestamps sort deterministically across adapters", %{tenant: tenant} do + a = %Log{timestamp_ns: 10, observed_timestamp_ns: 100, trace_id: "aaa"} + b = %Log{timestamp_ns: 10, observed_timestamp_ns: 300, trace_id: "aaa"} + c = %Log{timestamp_ns: 10, observed_timestamp_ns: 200, trace_id: "bbb"} + + assert :ok = S3.append(tenant, [a, b, c]) + assert {:ok, sorted} = S3.query(tenant, []) + assert Enum.map(sorted, & &1.observed_timestamp_ns) == [300, 200, 100] + end +end diff --git a/test/pulso/storage/s3_unit_test.exs b/test/pulso/storage/s3_unit_test.exs new file mode 100644 index 0000000..80f8cfe --- /dev/null +++ b/test/pulso/storage/s3_unit_test.exs @@ -0,0 +1,217 @@ +defmodule Pulso.Storage.S3UnitTest do + # Pure-Elixir tests of Pulso.Storage.S3's pre-network guards and key + # construction. Anything that actually talks to an S3 endpoint lives in + # s3_test.exs behind the `:integration` tag. + + use ExUnit.Case, async: true + + alias Pulso.Record.Log + alias Pulso.Storage.S3 + + describe "tenant validation" do + test "rejects adversarial tenant names before touching the object store" do + for bad <- ["", "../evil", "foo/bar", "foo bar", String.duplicate("a", 200)] do + assert {:error, {:invalid_tenant, ^bad}} = S3.append(bad, [%Log{timestamp_ns: 1}]) + assert {:error, {:invalid_tenant, ^bad}} = S3.query(bad, []) + end + end + + test "rejects adversarial tenant names even for empty batches" do + # Prior version bailed early on `[]` and returned :ok, letting an + # attacker probe the auth surface with a no-op payload. + assert {:error, {:invalid_tenant, "../evil"}} = S3.append("../evil", []) + assert {:error, {:invalid_tenant, ""}} = S3.append("", []) + end + + test "accepts common tenant name shapes" do + for good <- ["default", "customer-42", "team.alpha", "svc_web"] do + assert :ok = S3.append(good, []) + end + end + end + + describe "encode failures" do + test "returns an error tuple instead of raising on non-UTF-8 body bytes" do + record = %Log{timestamp_ns: 1, body: <<0xFF, 0xFE>>} + + assert {:error, {:encode_failed, _}} = S3.append("default", [record]) + end + end + + describe "object key construction" do + test "with an idempotency key, the same call maps to the same key" do + k1 = S3.object_key("acme", 100, "payload", "req-1") + k2 = S3.object_key("acme", 100, "payload", "req-1") + assert k1 == k2 + end + + test "with an idempotency key, different keys map to different objects" do + k1 = S3.object_key("acme", 100, "payload", "req-1") + k2 = S3.object_key("acme", 100, "payload", "req-2") + refute k1 == k2 + end + + test "without an idempotency key, identical payloads still map to distinct objects" do + # This is the "two legitimate producers happen to have identical bytes" + # case Codex flagged in review 2: content-addressing alone would + # silently drop one. With no idempotency key we always take a random + # suffix so distinct calls always land on distinct objects. + k1 = S3.object_key("acme", 100, "payload", nil) + k2 = S3.object_key("acme", 100, "payload", nil) + refute k1 == k2 + end + + test "an empty-string idempotency key is treated as absent" do + # Prevents a client that sends `Idempotency-Key: ` from accidentally + # collapsing every batch onto one key. + k1 = S3.object_key("acme", 100, "payload", "") + k2 = S3.object_key("acme", 100, "payload", "") + refute k1 == k2 + end + + test "the same idempotency key with different content produces different objects" do + # A caller who reuses an idempotency key with a different payload is + # almost certainly buggy. We must not silently overwrite the earlier + # write with the newer one; distinct content should land on distinct + # keys. + k1 = S3.object_key("acme", 100, "payload-a", "req-1") + k2 = S3.object_key("acme", 100, "payload-b", "req-1") + refute k1 == k2 + end + + test "the idempotency key is scoped by tenant" do + # A shared idempotency key across tenants must not collide. + k1 = S3.object_key("alpha", 100, "p", "req-1") + k2 = S3.object_key("beta", 100, "p", "req-1") + refute String.replace(k1, "alpha", "beta") == k2 + end + + test "different tenants never share a prefix" do + k1 = S3.object_key("alpha", 100, "p", "req-1") + k2 = S3.object_key("beta", 100, "p", "req-1") + refute String.starts_with?(k1, "tenants/beta/") + refute String.starts_with?(k2, "tenants/alpha/") + end + + test "the key path carries a schema version segment" do + # v1 is the current write format. A future format change bumps to + # v2/ so old objects can be migrated at their own pace rather than + # orphaned. Guarding this at the key level catches an accidental + # removal of the versioning. + assert String.starts_with?(S3.object_key("acme", 100, "p", "req-1"), "tenants/acme/v1/logs/") + end + + test "keys sort chronologically by sort_ns within a tenant" do + k_early = S3.object_key("acme", 100, "p", "req-1") + k_late = S3.object_key("acme", 200, "p", "req-1") + assert k_early < k_late + end + end + + describe "caller_content_hash" do + test "identical caller-provided records produce the same fingerprint" do + records = [ + %Log{timestamp_ns: 1, service: "api", body: "hello"}, + %Log{timestamp_ns: 2, service: "web", body: "world"} + ] + + assert {:ok, h1} = S3.caller_content_hash(records) + assert {:ok, h2} = S3.caller_content_hash(records) + assert h1 == h2 + end + + test "two calls with observed_timestamp_ns left nil produce the same fingerprint" do + # This is the retry-safety case: OTLP callers commonly omit + # `observedTimeUnixNano`. Two retries both arrive with nil, both hash + # to the same value, and the idempotency-key path deduplicates. The + # normalizer fills in a wall clock later, but that happens AFTER the + # fingerprint is computed. + r = [%Log{timestamp_ns: 1, body: "same"}] + + assert {:ok, h1} = S3.caller_content_hash(r) + assert {:ok, h2} = S3.caller_content_hash(r) + assert h1 == h2 + end + + test "a caller-set observed_timestamp_ns IS part of the fingerprint" do + # If the caller explicitly declares an observed timestamp, that is + # part of the record they authored. A "retry" that changes it is a + # distinct write; silently overwriting the earlier value would lose + # data. + no_obs = [%Log{timestamp_ns: 1, body: "same"}] + with_obs_a = [%Log{timestamp_ns: 1, observed_timestamp_ns: 100, body: "same"}] + with_obs_b = [%Log{timestamp_ns: 1, observed_timestamp_ns: 200, body: "same"}] + + {:ok, h_none} = S3.caller_content_hash(no_obs) + {:ok, h_a} = S3.caller_content_hash(with_obs_a) + {:ok, h_b} = S3.caller_content_hash(with_obs_b) + + refute h_none == h_a + refute h_a == h_b + end + + test "every caller-controlled field influences the fingerprint" do + base = [%Log{timestamp_ns: 1, service: "api", body: "same"}] + diff_body = [%Log{timestamp_ns: 1, service: "api", body: "different"}] + diff_service = [%Log{timestamp_ns: 1, service: "web", body: "same"}] + diff_ts = [%Log{timestamp_ns: 2, service: "api", body: "same"}] + + {:ok, h_base} = S3.caller_content_hash(base) + {:ok, h_body} = S3.caller_content_hash(diff_body) + {:ok, h_service} = S3.caller_content_hash(diff_service) + {:ok, h_ts} = S3.caller_content_hash(diff_ts) + + refute h_base == h_body + refute h_base == h_service + refute h_base == h_ts + end + + test "a map with keys inserted in different orders produces the same fingerprint" do + # `:erlang.term_to_binary(_, [:deterministic])` sorts map keys before + # encoding. Without that, two logically-equal records could hash + # differently just because the caller inserted attributes in a + # different order — which would defeat idempotent retries across + # heterogeneous producers. + order_a = %{"a" => 1, "b" => 2, "c" => 3} + order_b = order_a |> Map.delete("a") |> Map.put("a", 1) + + r_a = [%Log{timestamp_ns: 1, attributes: order_a}] + r_b = [%Log{timestamp_ns: 1, attributes: order_b}] + + {:ok, h_a} = S3.caller_content_hash(r_a) + {:ok, h_b} = S3.caller_content_hash(r_b) + assert h_a == h_b + end + + test "a non-UTF-8 body still fingerprints without crashing" do + # `:erlang.term_to_binary` handles any Elixir term, including binaries + # that are not valid UTF-8. The stored payload's `encode/1` still + # surfaces `:encode_failed` at write time — we just don't need to + # trip that path here. + assert {:ok, _} = S3.caller_content_hash([%Log{timestamp_ns: 1, body: <<255>>}]) + end + end + + describe "attribute sanitization" do + test "coerces non-string map keys to strings recursively" do + assert {:ok, sanitized} = + S3.sanitize_map(%{ + :status => 200, + "nested" => %{404 => "missing", :ref => "abc"}, + "list" => [%{true => 1}] + }) + + assert Map.has_key?(sanitized, "status") + assert sanitized["nested"] == %{"404" => "missing", "ref" => "abc"} + assert sanitized["list"] == [%{"true" => 1}] + end + + test "rejects rather than collapses keys that stringify to the same value" do + # A previous version silently dropped one value on collision. Codex + # flagged the risk (map size decreases without a diagnostic). Now we + # return an explicit error so the caller can surface it. + assert {:error, {:attribute_key_collision, _}} = + S3.sanitize_map(%{1 => "int", "1" => "string"}) + end + end +end diff --git a/test/pulso/storage/sort_order_test.exs b/test/pulso/storage/sort_order_test.exs new file mode 100644 index 0000000..de23623 --- /dev/null +++ b/test/pulso/storage/sort_order_test.exs @@ -0,0 +1,79 @@ +defmodule Pulso.Storage.SortOrderTest do + # SortOrder is the shared tie-breaker for every storage adapter. A client + # that swaps adapters must never see the order change when timestamps are + # equal, so this coverage runs against the pure sort — no adapter needed. + + use ExUnit.Case, async: true + + alias Pulso.Record.Log + alias Pulso.Storage.SortOrder + + defp log(ts, opts \\ []) do + %Log{ + timestamp_ns: ts, + observed_timestamp_ns: Keyword.get(opts, :observed), + trace_id: Keyword.get(opts, :trace_id), + span_id: Keyword.get(opts, :span_id), + body: Keyword.get(opts, :body) + } + end + + test "primary sort is timestamp_ns descending" do + result = SortOrder.sort([log(10), log(30), log(20)]) + assert Enum.map(result, & &1.timestamp_ns) == [30, 20, 10] + end + + test "equal timestamp_ns falls through to observed_timestamp_ns descending" do + result = + SortOrder.sort([ + log(10, observed: 100), + log(10, observed: 300), + log(10, observed: 200) + ]) + + assert Enum.map(result, & &1.observed_timestamp_ns) == [300, 200, 100] + end + + test "trace_id and body break further ties deterministically" do + a = log(10, observed: 5, trace_id: "aaa", body: "one") + b = log(10, observed: 5, trace_id: "aaa", body: "two") + c = log(10, observed: 5, trace_id: "bbb", body: "one") + + # The precise order matters less than the fact that repeated sorts of + # the same input produce the same output. + assert SortOrder.sort([a, b, c]) == SortOrder.sort([c, b, a]) + assert SortOrder.sort([a, b, c]) == SortOrder.sort([b, a, c]) + end + + test "nil fields sort last within their position" do + with_trace = log(10, observed: 5, trace_id: "aaa") + without_trace = log(10, observed: 5, trace_id: nil) + + result = SortOrder.sort([without_trace, with_trace]) + assert Enum.map(result, & &1.trace_id) == ["aaa", nil] + end + + test "a non-string body sorts without crashing" do + # OTLP `AnyValue` can produce a body that is not a String — an + # integer, a boolean, a list. The sort layer must not crash on those; + # sorting is for stability, not for user-facing order. + int_body = %Log{timestamp_ns: 10, observed_timestamp_ns: 5, body: 42} + list_body = %Log{timestamp_ns: 10, observed_timestamp_ns: 5, body: [1, 2, 3]} + bool_body = %Log{timestamp_ns: 10, observed_timestamp_ns: 5, body: true} + string_body = %Log{timestamp_ns: 10, observed_timestamp_ns: 5, body: "z"} + nil_body = %Log{timestamp_ns: 10, observed_timestamp_ns: 5, body: nil} + + assert [_, _, _, _, _] = SortOrder.sort([int_body, list_body, bool_body, string_body, nil_body]) + end + + test "nil and empty string in the same position are distinguishable" do + # Prior version collapsed `nil` and `""` into the same sort key, so two + # otherwise-identical records could reorder unpredictably. Now `nil` + # sorts after every real string, including `""`. + with_empty = log(10, observed: 5, trace_id: "", body: "z") + without = log(10, observed: 5, trace_id: nil, body: "z") + + result = SortOrder.sort([without, with_empty]) + assert Enum.map(result, & &1.trace_id) == ["", nil] + end +end diff --git a/test/pulso_web/controllers/otlp_controller_test.exs b/test/pulso_web/controllers/otlp_controller_test.exs index da6d7c3..6012cbc 100644 --- a/test/pulso_web/controllers/otlp_controller_test.exs +++ b/test/pulso_web/controllers/otlp_controller_test.exs @@ -1,6 +1,8 @@ defmodule PulsoWeb.OTLPControllerTest do use PulsoWeb.ConnCase, async: false + alias Pulso.Auth.Open + alias Pulso.Auth.SharedSecret alias Pulso.Record.Log alias Pulso.Storage alias Pulso.Storage.Memory @@ -66,4 +68,122 @@ defmodule PulsoWeb.OTLPControllerTest do assert json_response(conn, 200) == %{} assert {:ok, []} = Storage.query("default") end + + test "POST /v1/logs returns 400 for a tenant name that would escape the prefix", + %{conn: conn} do + conn = + conn + |> put_req_header("content-type", "application/json") + |> put_req_header("x-scope-orgid", "bad/name") + |> post(~p"/v1/logs", payload()) + + assert json_response(conn, 400) == %{"error" => "invalid_tenant"} + end + + test "POST /v1/logs surfaces rejected records via partialSuccess per the OTLP spec", + %{conn: conn} do + # One valid record + one malformed (non-map) entry. The receiver must + # signal the drop rather than acknowledging silently — otherwise the + # sender never retries the record that was never stored. + payload = %{ + "resourceLogs" => [ + %{ + "scopeLogs" => [ + %{ + "logRecords" => [ + %{"timeUnixNano" => "1700000000000000000", "body" => %{"stringValue" => "ok"}}, + "not-a-map" + ] + } + ] + } + ] + } + + conn = + conn + |> put_req_header("content-type", "application/json") + |> post(~p"/v1/logs", payload) + + body = json_response(conn, 200) + assert %{"partialSuccess" => %{"rejectedLogRecords" => 1}} = body + end + + test "POST /v1/logs passes the Idempotency-Key header through", %{conn: conn} do + # The Idempotency-Key header should end up in Storage.append opts. Two + # POSTs with the same key + same payload are the same write to the store, + # so downstream queries see one record even if two arrived over the wire. + # (The Memory adapter ignores :idempotency_key, so we only assert the + # controller accepts the header; S3 adapter coverage of the actual + # dedup lives in s3_test.exs.) + conn = + conn + |> put_req_header("content-type", "application/json") + |> put_req_header("idempotency-key", "req-1") + |> post(~p"/v1/logs", payload()) + + assert json_response(conn, 200) == %{} + end + + describe "with Pulso.Auth.SharedSecret enabled" do + setup do + hex = Base.encode16(:crypto.hash(:sha256, "the-token"), case: :lower) + + Application.put_env(:pulso, Pulso.Auth, + module: SharedSecret, + tokens: %{"acme" => "sha256$#{hex}"} + ) + + on_exit(fn -> + Application.put_env(:pulso, Pulso.Auth, module: Open) + end) + + :ok + end + + test "accepts the correct token", %{conn: conn} do + conn = + conn + |> put_req_header("content-type", "application/json") + |> put_req_header("x-scope-orgid", "acme") + |> put_req_header("authorization", "Bearer the-token") + |> post(~p"/v1/logs", payload()) + + assert json_response(conn, 200) == %{} + assert {:ok, [%Log{service: "api", body: "hello"}]} = Storage.query("acme") + end + + test "rejects a request with a bad token", %{conn: conn} do + conn = + conn + |> put_req_header("content-type", "application/json") + |> put_req_header("x-scope-orgid", "acme") + |> put_req_header("authorization", "Bearer wrong") + |> post(~p"/v1/logs", payload()) + + assert json_response(conn, 401) == %{"error" => "invalid_token"} + assert {:ok, []} = Storage.query("acme") + end + + test "rejects a tenant with no configured token", %{conn: conn} do + conn = + conn + |> put_req_header("content-type", "application/json") + |> put_req_header("x-scope-orgid", "unknown") + |> put_req_header("authorization", "Bearer the-token") + |> post(~p"/v1/logs", payload()) + + assert json_response(conn, 401) == %{"error" => "unknown_tenant"} + end + + test "rejects a request with no authorization header", %{conn: conn} do + conn = + conn + |> put_req_header("content-type", "application/json") + |> put_req_header("x-scope-orgid", "acme") + |> post(~p"/v1/logs", payload()) + + assert json_response(conn, 401) == %{"error" => "missing_token"} + end + end end