From 2d83b393f7eada6da6317638891b999502271550 Mon Sep 17 00:00:00 2001 From: SimpleTest Date: Mon, 20 Jul 2026 07:12:34 +0300 Subject: [PATCH] fix: right-size workers and isolate cluster ingress --- .env.e2e.example | 2 + .env.example | 2 + .env.load.example | 2 + compose.yaml | 2 + config/config.exs | 2 +- config/runtime.exs | 12 + .../who-need-help/templates/deployments.yaml | 4 + .../who-need-help/templates/migrate-job.yaml | 4 + .../templates/network-policy.yaml | 25 ++ deploy/helm/who-need-help/values.yaml | 5 + docs/architecture.md | 19 +- docs/operations.md | 32 ++ docs/performance.md | 37 ++ docs/verification.md | 27 ++ .../workers/local_burst_probe.ex | 27 ++ scripts/clean-deploy-verify.sh | 2 + scripts/db-scale-benchmark.sh | 7 +- scripts/ensure-local-e2e-env.sh | 2 + scripts/ensure-local-load-env.sh | 26 +- scripts/oban-burst-run.sh | 320 ++++++++++++++++++ scripts/quality.sh | 5 + .../workers/local_burst_probe_test.exs | 23 ++ 22 files changed, 579 insertions(+), 8 deletions(-) create mode 100644 deploy/helm/who-need-help/templates/network-policy.yaml create mode 100644 lib/who_need_help/workers/local_burst_probe.ex create mode 100755 scripts/oban-burst-run.sh create mode 100644 test/who_need_help/workers/local_burst_probe_test.exs diff --git a/.env.e2e.example b/.env.e2e.example index e42a29e..aef42c5 100644 --- a/.env.e2e.example +++ b/.env.e2e.example @@ -23,6 +23,8 @@ DATABASE_URL=ecto://postgres:GENERATE_URL_SAFE_PASSWORD@db/who_need_help_e2e WEB_POOL_SIZE=4 WORKER_POOL_SIZE=2 MIGRATE_POOL_SIZE=2 +OBAN_MAINTENANCE_CONCURRENCY=2 +OBAN_PUSH_CONCURRENCY=1 WEB_REPLICAS=2 WORKER_REPLICAS=2 ERLANG_PORT_LIMIT=65536 diff --git a/.env.example b/.env.example index ee2811b..50818b6 100644 --- a/.env.example +++ b/.env.example @@ -67,6 +67,8 @@ DATABASE_URL=ecto://postgres:replace-with-url-encoded-password@db/who_need_help WEB_POOL_SIZE=4 WORKER_POOL_SIZE=2 MIGRATE_POOL_SIZE=2 +OBAN_MAINTENANCE_CONCURRENCY=2 +OBAN_PUSH_CONCURRENCY=1 WEB_REPLICAS=2 WORKER_REPLICAS=2 # Maximum simultaneously existing Erlang ports (files, sockets and drivers). diff --git a/.env.load.example b/.env.load.example index d402f29..6c23372 100644 --- a/.env.load.example +++ b/.env.load.example @@ -24,6 +24,8 @@ DATABASE_URL=GENERATE_DATABASE_URL WEB_POOL_SIZE=4 WORKER_POOL_SIZE=2 MIGRATE_POOL_SIZE=2 +OBAN_MAINTENANCE_CONCURRENCY=2 +OBAN_PUSH_CONCURRENCY=1 WEB_REPLICAS=2 WORKER_REPLICAS=2 ERLANG_PORT_LIMIT=65536 diff --git a/compose.yaml b/compose.yaml index a3765fa..97854bf 100644 --- a/compose.yaml +++ b/compose.yaml @@ -43,6 +43,8 @@ x-app-environment: &app-environment PUSH_HTTP_RECEIVE_TIMEOUT_MS: ${PUSH_HTTP_RECEIVE_TIMEOUT_MS:-} PUSH_HTTP_CONNECT_TIMEOUT_MS: ${PUSH_HTTP_CONNECT_TIMEOUT_MS:-} PUSH_HTTP_RETRY_DELAY_MS: ${PUSH_HTTP_RETRY_DELAY_MS:-} + OBAN_MAINTENANCE_CONCURRENCY: ${OBAN_MAINTENANCE_CONCURRENCY:-2} + OBAN_PUSH_CONCURRENCY: ${OBAN_PUSH_CONCURRENCY:-1} services: docker-api-proxy: diff --git a/config/config.exs b/config/config.exs index b4118f9..ffffd2a 100644 --- a/config/config.exs +++ b/config/config.exs @@ -39,7 +39,7 @@ config :geo_postgis, json_library: Jason config :who_need_help, Oban, repo: WhoNeedHelp.Repo, - queues: [default: 10, maintenance: 2, push: 1], + queues: [maintenance: 2, push: 1], plugins: [ {Oban.Plugins.Pruner, max_age: 86_400}, {Oban.Plugins.Cron, crontab: [{"* * * * *", WhoNeedHelp.Workers.ExpireRequests}]} diff --git a/config/runtime.exs b/config/runtime.exs index b58f412..7cfc673 100644 --- a/config/runtime.exs +++ b/config/runtime.exs @@ -46,6 +46,18 @@ required_positive_integer = fn name -> end end +positive_integer_with_default = fn name, default -> + optional_positive_integer.(name) || default +end + +if config_env() == :prod do + config :who_need_help, Oban, + queues: [ + maintenance: positive_integer_with_default.("OBAN_MAINTENANCE_CONCURRENCY", 2), + push: positive_integer_with_default.("OBAN_PUSH_CONCURRENCY", 1) + ] +end + required_non_negative_integer = fn name -> case System.get_env(name) do value when value in [nil, ""] -> diff --git a/deploy/helm/who-need-help/templates/deployments.yaml b/deploy/helm/who-need-help/templates/deployments.yaml index 51bb809..e2550b4 100644 --- a/deploy/helm/who-need-help/templates/deployments.yaml +++ b/deploy/helm/who-need-help/templates/deployments.yaml @@ -85,6 +85,10 @@ spec: value: {{ $root.Values.app.port | quote }} - name: POOL_SIZE value: {{ $settings.poolSize | quote }} + - name: OBAN_MAINTENANCE_CONCURRENCY + value: {{ $root.Values.worker.maintenanceConcurrency | quote }} + - name: OBAN_PUSH_CONCURRENCY + value: {{ $root.Values.worker.pushConcurrency | quote }} - name: SMTP_RELAY value: {{ $root.Values.app.smtpRelay | quote }} - name: SMTP_PORT diff --git a/deploy/helm/who-need-help/templates/migrate-job.yaml b/deploy/helm/who-need-help/templates/migrate-job.yaml index 0f6752b..b8a8a72 100644 --- a/deploy/helm/who-need-help/templates/migrate-job.yaml +++ b/deploy/helm/who-need-help/templates/migrate-job.yaml @@ -39,6 +39,10 @@ spec: value: {{ .Values.app.clusterInterface | quote }} - name: POOL_SIZE value: {{ .Values.app.migratePoolSize | quote }} + - name: OBAN_MAINTENANCE_CONCURRENCY + value: {{ .Values.worker.maintenanceConcurrency | quote }} + - name: OBAN_PUSH_CONCURRENCY + value: {{ .Values.worker.pushConcurrency | quote }} securityContext: allowPrivilegeEscalation: false readOnlyRootFilesystem: true diff --git a/deploy/helm/who-need-help/templates/network-policy.yaml b/deploy/helm/who-need-help/templates/network-policy.yaml new file mode 100644 index 0000000..6ea0539 --- /dev/null +++ b/deploy/helm/who-need-help/templates/network-policy.yaml @@ -0,0 +1,25 @@ +{{- if .Values.networkPolicy.enabled }} +apiVersion: networking.k8s.io/v1 +kind: NetworkPolicy +metadata: + name: {{ include "who-need-help.fullname" . }} + labels: + {{- include "who-need-help.labels" . | nindent 4 }} +spec: + podSelector: + matchLabels: + {{- include "who-need-help.selectorLabels" . | nindent 6 }} + policyTypes: + - Ingress + ingress: + # Web and worker nodes need unrestricted pod-to-pod Erlang distribution. + - from: + - podSelector: + matchLabels: + {{- include "who-need-help.selectorLabels" . | nindent 14 }} + # Public web traffic, kubelet probes, and authenticated metrics use the + # same fixed listener. All other inbound ports stay isolated. + - ports: + - protocol: TCP + port: {{ .Values.app.port }} +{{- end }} diff --git a/deploy/helm/who-need-help/values.yaml b/deploy/helm/who-need-help/values.yaml index a465a4f..97e89ad 100644 --- a/deploy/helm/who-need-help/values.yaml +++ b/deploy/helm/who-need-help/values.yaml @@ -10,6 +10,8 @@ web: worker: replicas: 2 poolSize: "2" + maintenanceConcurrency: "2" + pushConcurrency: "1" service: type: ClusterIP @@ -54,6 +56,9 @@ ingress: pathType: Prefix tls: [] +networkPolicy: + enabled: true + podAnnotations: {} podLabels: {} nodeSelector: {} diff --git a/docs/architecture.md b/docs/architecture.md index 558b6e1..e6de257 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -77,7 +77,10 @@ Phoenix PubSub and Presence handle transient fan-out. Phoenix's generated `dns_cluster` dependency discovers nodes using DNS polling: - Compose nodes discover the `web` service on the shared internal network; - workers join through those nodes and distributed Erlang forms the full mesh. + the internal-only `cluster-web` alias and `CLUSTER_INTERFACE=eth0` keep node + names on that network even when the container also has ingress or egress + interfaces. Workers join through those nodes and distributed Erlang forms + the full mesh. - Kubernetes uses a headless service. No sticky session is required for authenticated requests. Session cookies are @@ -88,7 +91,10 @@ Web and migrate application processes start a producerless Oban client with no queues, plugins, or peer leadership so transactions can insert unique jobs. Only the worker role starts queue consumers and scheduled-job plugins. PostgreSQL coordinates queues and leadership, so no Redis dependency is -introduced. +introduced. The worker runs only the queues used by product code: +`maintenance` for expiry/probes and `push` for provider-neutral delivery. +Their per-worker concurrency is configured independently; no unused default +queue is started. ## Geospatial data @@ -137,7 +143,14 @@ not create database high availability. The Helm chart contains separate web and worker Deployments, Services, a headless cluster-discovery Service, a migration Job, Secret interfaces, probes, -and disruption-aware rolling updates. It does not invent +disruption-aware rolling updates, and an ingress NetworkPolicy. The policy +allows unrestricted Erlang distribution only between pods belonging to the +same chart instance, exposes the configured HTTP listener, and isolates other +inbound pod ports. Kubernetes enforces this only when the cluster's network +plugin implements NetworkPolicy. It does not restrict egress because the +production database, SMTP, OAuth, map, and push-provider destinations are not +known yet; those rules must be added from verified deployment-specific +addresses rather than invented in the chart. The chart does not invent CPU/memory limits or an HPA threshold before measurements exist. The bundled kind path is for reproducible local verification. A production diff --git a/docs/operations.md b/docs/operations.md index 21a6fc5..17fe637 100644 --- a/docs/operations.md +++ b/docs/operations.md @@ -284,6 +284,38 @@ The rollout timeout, probe interval/timeout/retry count, and cluster-join timeout are experiment inputs. They are not production SLOs or resource requirements. +## Isolated Oban burst measurement + +The worker role consumes only `maintenance` and `push`. Configure their +per-worker limits with `OBAN_MAINTENANCE_CONCURRENCY` and +`OBAN_PUSH_CONCURRENCY`; multiplying either value by the number of worker +replicas gives the configured cluster-wide concurrency for that queue. + +After starting the isolated load project, run an explicitly sized experiment: + +```bash +./scripts/load-stack-up.sh +./scripts/oban-burst-run.sh local-oban-burst 1000 120 +``` + +The three required arguments are an evidence label, job count, and recorded +experiment timeout in seconds. They are not capacity thresholds. The script +refuses the ordinary `who_need_help` project, verifies the expected worker +replicas, inserts only confirmed `LocalBurstProbe` jobs with a unique run id, +records each worker's effective queue configuration and container samples, +requires every job to complete, compares domain-table row counts, and deletes +exactly its own jobs. The probe worker has no product side effects and no +product flow enqueues it. + +## Kubernetes ingress isolation + +`networkPolicy.enabled=true` renders the chart's ingress NetworkPolicy. +Enforcement is a property of the cluster CNI, not of the YAML object alone. +Before relying on it, verify that the target cluster uses a NetworkPolicy- +capable plugin and run positive HTTP/cluster checks plus negative blocked-port +checks there. Egress remains intentionally unrestricted until the actual +database and external-service destinations are known. + ## BEAM runtime memory guard Compose sets `ERL_ZFLAGS="+Q ${ERLANG_PORT_LIMIT}"`; Helm renders the same flag diff --git a/docs/performance.md b/docs/performance.md index 2e9c2fa..c543bf5 100644 --- a/docs/performance.md +++ b/docs/performance.md @@ -544,3 +544,40 @@ application commit `b96d443`. ultimately returned HTTP 200; two required one retry during local NodePort endpoint replacement. Evidence is `output/resilience/final-kind-rollout-fixed-20260720`. + +## Runtime hardening and queue replay on 2026-07-20 + +The ordinary and isolated load releases were rebuilt after separating +internal cluster, ingress, and outbound networks. All web and worker node names +used internal `172.16.52.x` or `172.16.55.x` addresses respectively, every +node observed the expected peers, and cross-node PubSub passed. The +worker runtime no longer starts the unused `default` queue. Both measured load +workers reported exactly `maintenance: 2` and `push: 1`; their BEAM process +counts were 574 after the change, compared with the earlier observation of +577/578. This three-to-four-process difference is an observation, not a memory +or capacity guarantee. + +The isolated Oban burst experiment inserted 1,000 side-effect-free jobs in one +batch. With two workers and configured maintenance concurrency two per worker, +all 1,000 reached `completed` in an observed 4,516 ms. Three resource samples +per measured container recorded maxima of 80.80 MiB for PostgreSQL and +218.80/186.90 MiB for the two worker cgroups. No error matched the run-scoped +logs; the database reported zero deadlocks, temporary files, or conflicts. +The script removed exactly 1,000 probe rows, verified zero remaining probe +jobs, and found byte-identical before/after row-count documents for all domain +tables. These short local samples do not establish a production throughput, +steady-state memory plateau, or minimum server size. Evidence is +`output/performance/oban-burst-1000-final-20260720`. + +The database-scale harness was repaired to build and pass its own hardened +PostGIS image rather than relying on an image variable owned by another +script. Its fresh PostgreSQL 18.4 replay again used 50,000 configured rows per +large table and identical before/after table counts. Every asserted list plan +selected its intended cursor index. Observed examples were 8.806 to 0.123 ms +for urgent-help discovery, 8.348 to 0.061 ms for Activity discovery, and 9.084 +to 0.048 ms for visible reviews. The aggregate leaderboard and single hot +helper reputation plans still scanned their 50,000-row completed-assignment +working sets and took 53.949 and 41.247 ms in this deliberately concentrated +fixture. That is a measured future optimization target, not evidence of a +current failure or a portable latency. Evidence is +`output/db-scale/20260720035817-2416501`. diff --git a/docs/verification.md b/docs/verification.md index 5801576..94f6983 100644 --- a/docs/verification.md +++ b/docs/verification.md @@ -810,6 +810,33 @@ state through the idempotent settings page instead of replaying the token. observability and container-image security gates. Evidence is `output/regression/final-single-use-quality-20260720.log`. +## Queue and ingress-policy hardening replay + +The 2026-07-20 hardening pass removed the unused Oban `default` consumer and +made the two real queue limits explicit in Compose, generated local +environments, and Helm. After rebuilding the isolated load release, both +worker nodes reported exactly `maintenance: [limit: 2]` and +`push: [limit: 1]`; the five-node cluster and cross-node PubSub probe passed. +The explicit 1,000-job burst completed 1,000/1,000 jobs in the observed 4,516 +ms, produced no matched application/database error, removed exactly its 1,000 +run-scoped rows, left zero probe jobs, and retained identical domain-table +counts. Evidence is +`output/performance/oban-burst-1000-final-20260720`. + +The Helm chart now renders one ingress NetworkPolicy for the chart instance. +Helm lint and the rendered-manifest Trivy scan pass. The manifest permits +chart-instance pod-to-pod Erlang distribution and the configured HTTP +listener while isolating other inbound pod ports. Actual packet enforcement +has not been claimed because it depends on the target cluster's CNI; that must +be verified on the eventual deployment environment. + +The same source state passed `scripts/quality.sh`: 186/186 ExUnit tests, +format and warnings-as-errors compilation, xref, Credo, Sobelow, Dialyzer, +Hex/npm audits, ShellCheck, Hadolint, actionlint, every Compose render, Helm +lint, observability configuration, rendered-manifest scanning, and all +configured runtime image scans. Every reported HIGH/CRITICAL vulnerability +count was zero. + ## Known work before a public production launch - Replace the temporary staging origin with the production-owned domain and diff --git a/lib/who_need_help/workers/local_burst_probe.ex b/lib/who_need_help/workers/local_burst_probe.ex new file mode 100644 index 0000000..2847a56 --- /dev/null +++ b/lib/who_need_help/workers/local_burst_probe.ex @@ -0,0 +1,27 @@ +defmodule WhoNeedHelp.Workers.LocalBurstProbe do + @moduledoc """ + A side-effect-free worker for the isolated local Oban burst measurement. + + No application flow enqueues this worker. The measurement script requires an + explicit confirmation value, tags every job with a unique run id, and removes + only those exact rows after collecting evidence. + """ + + use Oban.Worker, queue: :maintenance, max_attempts: 1, tags: ["local-burst-probe"] + + @confirmation "isolated-local-oban-burst-probe" + + @impl Oban.Worker + def perform(%Oban.Job{ + args: %{ + "confirmation" => @confirmation, + "run_id" => run_id, + "sequence" => sequence + } + }) + when is_binary(run_id) and run_id != "" and is_integer(sequence) and sequence > 0 do + :ok + end + + def perform(_job), do: {:cancel, :invalid_local_burst_probe} +end diff --git a/scripts/clean-deploy-verify.sh b/scripts/clean-deploy-verify.sh index a3d66d3..bb50865 100755 --- a/scripts/clean-deploy-verify.sh +++ b/scripts/clean-deploy-verify.sh @@ -163,6 +163,8 @@ DATABASE_URL=ecto://$postgres_user:$postgres_password@db/$postgres_db WEB_POOL_SIZE=4 WORKER_POOL_SIZE=2 MIGRATE_POOL_SIZE=2 +OBAN_MAINTENANCE_CONCURRENCY=2 +OBAN_PUSH_CONCURRENCY=1 WEB_REPLICAS=2 WORKER_REPLICAS=2 ERLANG_PORT_LIMIT=65536 diff --git a/scripts/db-scale-benchmark.sh b/scripts/db-scale-benchmark.sh index c72f158..af5b1bc 100755 --- a/scripts/db-scale-benchmark.sh +++ b/scripts/db-scale-benchmark.sh @@ -21,18 +21,20 @@ fi run_id="$(date -u +%Y%m%d%H%M%S)-$$" project="wnh_db_scale_$(printf '%s' "$run_id" | tr -d '-')" image="who-need-help:db-scale-$run_id" +postgis_image="who-need-help:db-scale-postgis-$run_id" output="$ROOT/output/db-scale/$run_id" umask 077 QUALITY_POSTGRES_USER="wnh_scale_$(openssl rand -hex 6)" QUALITY_POSTGRES_PASSWORD=$(openssl rand -base64 48 | tr -d '\n') -export QUALITY_POSTGRES_USER QUALITY_POSTGRES_PASSWORD +QUALITY_POSTGIS_IMAGE=$postgis_image +export QUALITY_POSTGRES_USER QUALITY_POSTGRES_PASSWORD QUALITY_POSTGIS_IMAGE compose="docker compose -p $project -f $ROOT/compose.quality.yaml" cleanup() { $compose down --volumes --remove-orphans >/dev/null 2>&1 || true - docker image rm "$image" >/dev/null 2>&1 || true + docker image rm "$image" "$postgis_image" >/dev/null 2>&1 || true } trap cleanup EXIT HUP INT TERM @@ -40,6 +42,7 @@ mkdir -p "$output" echo "Building the isolated database-scale image" docker build --target test --tag "$image" . +docker build --tag "$postgis_image" --file Dockerfile.postgis . echo "Starting an isolated PostgreSQL/PostGIS 18 database" $compose up --detach --wait db diff --git a/scripts/ensure-local-e2e-env.sh b/scripts/ensure-local-e2e-env.sh index 384cb9c..7d2c991 100755 --- a/scripts/ensure-local-e2e-env.sh +++ b/scripts/ensure-local-e2e-env.sh @@ -63,6 +63,8 @@ DATABASE_URL=ecto://postgres:$postgres_password@db/who_need_help_e2e WEB_POOL_SIZE=4 WORKER_POOL_SIZE=2 MIGRATE_POOL_SIZE=2 +OBAN_MAINTENANCE_CONCURRENCY=2 +OBAN_PUSH_CONCURRENCY=1 WEB_REPLICAS=2 WORKER_REPLICAS=2 ERLANG_PORT_LIMIT=65536 diff --git a/scripts/ensure-local-load-env.sh b/scripts/ensure-local-load-env.sh index 20339ce..6c5fb1b 100755 --- a/scripts/ensure-local-load-env.sh +++ b/scripts/ensure-local-load-env.sh @@ -16,6 +16,8 @@ if [ -f "$ENV_FILE" ]; then chmod 600 "$ENV_FILE" needs_fixture_password=true + needs_oban_maintenance_concurrency=true + needs_oban_push_concurrency=true needs_resilience_timeout=true needs_resilience_interval=true needs_resilience_request_timeout=true @@ -40,6 +42,10 @@ if [ -f "$ENV_FILE" ]; then needs_backup_interruption_interval=true grep -q '^LOAD_FIXTURE_PASSWORD=' "$ENV_FILE" && needs_fixture_password=false + grep -q '^OBAN_MAINTENANCE_CONCURRENCY=' "$ENV_FILE" && + needs_oban_maintenance_concurrency=false + grep -q '^OBAN_PUSH_CONCURRENCY=' "$ENV_FILE" && + needs_oban_push_concurrency=false grep -q '^LOAD_RESILIENCE_RECOVERY_TIMEOUT_SECONDS=' "$ENV_FILE" && needs_resilience_timeout=false grep -q '^LOAD_RESILIENCE_PROBE_INTERVAL_SECONDS=' "$ENV_FILE" && @@ -85,6 +91,8 @@ if [ -f "$ENV_FILE" ]; then needs_backup_interruption_interval=false if [ "$needs_fixture_password" = false ] && + [ "$needs_oban_maintenance_concurrency" = false ] && + [ "$needs_oban_push_concurrency" = false ] && [ "$needs_resilience_timeout" = false ] && [ "$needs_resilience_interval" = false ] && [ "$needs_resilience_request_timeout" = false ] && @@ -147,6 +155,19 @@ if [ -f "$ENV_FILE" ]; then printf 'LOAD_FIXTURE_PASSWORD=%s\n' "$load_fixture_password" fi + if [ "$needs_oban_maintenance_concurrency" = true ] || + [ "$needs_oban_push_concurrency" = true ]; then + printf '\n# Added by the worker-queue configuration upgrade.\n' + fi + + if [ "$needs_oban_maintenance_concurrency" = true ]; then + printf 'OBAN_MAINTENANCE_CONCURRENCY=2\n' + fi + + if [ "$needs_oban_push_concurrency" = true ]; then + printf 'OBAN_PUSH_CONCURRENCY=1\n' + fi + if [ "$needs_resilience_timeout" = true ] || [ "$needs_resilience_interval" = true ]; then printf '\n# Added by the local resilience-profile upgrade.\n' @@ -268,7 +289,8 @@ if [ -f "$ENV_FILE" ]; then chmod 600 "$ENV_FILE" unset load_fixture_password observability_grafana_admin_password \ backup_minio_root_user backup_minio_root_password backup_restic_password \ - needs_fixture_password needs_resilience_timeout \ + needs_fixture_password needs_oban_maintenance_concurrency \ + needs_oban_push_concurrency needs_resilience_timeout \ needs_resilience_interval needs_resilience_request_timeout \ needs_traefik_retry_attempts needs_observability_prometheus_port \ needs_observability_alertmanager_port needs_observability_grafana_port \ @@ -280,7 +302,7 @@ if [ -f "$ENV_FILE" ]; then needs_backup_bucket_prefix needs_backup_timeout \ needs_backup_interruption_chunks needs_backup_interruption_chunk_bytes \ needs_backup_interruption_interval - echo "Added missing load/resilience/observability/backup inputs to ignored .env.load." + echo "Added missing queue/load/resilience/observability/backup inputs to ignored .env.load." exit 0 fi diff --git a/scripts/oban-burst-run.sh b/scripts/oban-burst-run.sh new file mode 100755 index 0000000..72c6d8e --- /dev/null +++ b/scripts/oban-burst-run.sh @@ -0,0 +1,320 @@ +#!/usr/bin/env bash +set -euo pipefail + +ROOT=$(CDPATH='' cd -- "$(dirname -- "$0")/.." && pwd) +ENV_FILE="$ROOT/.env.load" + +if [[ $# -ne 3 ]]; then + echo "Usage: $0 LABEL JOB_COUNT TIMEOUT_SECONDS" >&2 + echo "All three values are recorded experiment inputs, not capacity thresholds." >&2 + exit 2 +fi + +label=$1 +job_count=$2 +timeout_seconds=$3 + +if [[ ! "$label" =~ ^[A-Za-z0-9._-]+$ ]]; then + echo "LABEL may contain only letters, numbers, dot, underscore, and dash." >&2 + exit 2 +fi + +for value_name in job_count timeout_seconds; do + value=${!value_name} + + if [[ ! "$value" =~ ^[1-9][0-9]*$ ]]; then + echo "${value_name^^} must be a positive integer." >&2 + exit 2 + fi +done + +if [[ ! -f "$ENV_FILE" ]]; then + echo "Missing $ENV_FILE. Run scripts/ensure-local-load-env.sh first." >&2 + exit 1 +fi + +set -a +# shellcheck source=/dev/null +. "$ENV_FILE" +set +a + +for name in LOAD_PROJECT LOAD_WORKER_REPLICAS POSTGRES_DB; do + if [[ -z "${!name:-}" ]]; then + echo "$name is missing from .env.load." >&2 + exit 1 + fi +done + +if [[ "$LOAD_PROJECT" == "who_need_help" ]]; then + echo "The Oban burst measurement must not use the ordinary Compose project." >&2 + exit 1 +fi + +if [[ ! "$LOAD_WORKER_REPLICAS" =~ ^[1-9][0-9]*$ ]]; then + echo "LOAD_WORKER_REPLICAS must be a positive integer." >&2 + exit 1 +fi + +compose=( + docker compose + --env-file "$ENV_FILE" + -p "$LOAD_PROJECT" + -f compose.yaml + -f compose.load.yaml +) + +mapfile -t worker_containers < <("${compose[@]}" ps -q worker) + +if [[ "${#worker_containers[@]}" -ne "$LOAD_WORKER_REPLICAS" ]]; then + echo "Expected $LOAD_WORKER_REPLICAS workers; observed ${#worker_containers[@]}." >&2 + exit 1 +fi + +db_container=$("${compose[@]}" ps -q db) + +if [[ -z "$db_container" ]]; then + echo "The isolated load database is not running." >&2 + exit 1 +fi + +output_dir="$ROOT/output/performance/$label" + +if [[ -e "$output_dir" ]]; then + echo "Output already exists: $output_dir" >&2 + exit 1 +fi + +mkdir -p "$output_dir" +chmod 700 "$ROOT/output" "$ROOT/output/performance" "$output_dir" + +run_id="${label}-$(date -u +%Y%m%dT%H%M%SZ)-$$" +worker_name="WhoNeedHelp.Workers.LocalBurstProbe" +sampler_pid= +sampling_marker="$output_dir/.sampling" +cleanup_required=false + +database_psql() { + # The variables are intentionally expanded inside the database container. + # shellcheck disable=SC2016 + "${compose[@]}" exec -T db sh -c \ + 'psql --no-psqlrc --quiet --tuples-only --no-align --set ON_ERROR_STOP=1 \ + --username "$POSTGRES_USER" --dbname "$POSTGRES_DB"' +} + +exact_job_count() { + database_psql <>'run_id' = '$run_id'; +SQL +} + +delete_exact_jobs() { + database_psql <>'run_id' = '$run_id' + RETURNING id +) +SELECT count(*) FROM deleted; +SQL +} + +snapshot_domain_counts() { + database_psql <<'SQL' +SELECT jsonb_pretty( + jsonb_object_agg(table_name, row_count ORDER BY table_name) +) +FROM ( + SELECT 'abuse_signals' AS table_name, count(*) AS row_count FROM abuse_signals + UNION ALL SELECT 'activities', count(*) FROM activities + UNION ALL SELECT 'activity_messages', count(*) FROM activity_messages + UNION ALL SELECT 'activity_participants', count(*) FROM activity_participants + UNION ALL SELECT 'audit_events', count(*) FROM audit_events + UNION ALL SELECT 'blocks', count(*) FROM blocks + UNION ALL SELECT 'categories', count(*) FROM categories + UNION ALL SELECT 'category_proposals', count(*) FROM category_proposals + UNION ALL SELECT 'category_votes', count(*) FROM category_votes + UNION ALL SELECT 'help_assignments', count(*) FROM help_assignments + UNION ALL SELECT 'help_requests', count(*) FROM help_requests + UNION ALL SELECT 'messages', count(*) FROM messages + UNION ALL SELECT 'reports', count(*) FROM reports + UNION ALL SELECT 'reviews', count(*) FROM reviews + UNION ALL SELECT 'social_identities', count(*) FROM social_identities + UNION ALL SELECT 'tracking_positions', count(*) FROM tracking_positions + UNION ALL SELECT 'tracking_sessions', count(*) FROM tracking_sessions + UNION ALL SELECT 'users', count(*) FROM users + UNION ALL SELECT 'users_tokens', count(*) FROM users_tokens +) AS counts; +SQL +} + +cleanup() { + trap - EXIT HUP INT TERM + + rm -f "$sampling_marker" + + if [[ -n "$sampler_pid" ]]; then + kill "$sampler_pid" >/dev/null 2>&1 || true + wait "$sampler_pid" >/dev/null 2>&1 || true + fi + + if [[ "$cleanup_required" == true ]]; then + delete_exact_jobs >"$output_dir/job-cleanup-on-exit.txt" 2>&1 || true + cleanup_required=false + fi +} +trap cleanup EXIT HUP INT TERM + +if [[ "$(exact_job_count)" != "0" ]]; then + echo "The exact burst run id already exists in oban_jobs." >&2 + exit 1 +fi + +snapshot_domain_counts >"$output_dir/domain-counts-before.json" + +{ + printf 'label=%s\n' "$label" + printf 'run_id=%s\n' "$run_id" + printf 'job_count=%s\n' "$job_count" + printf 'timeout_seconds=%s\n' "$timeout_seconds" + printf 'started_at=%s\n' "$(date -u +%Y-%m-%dT%H:%M:%SZ)" + printf 'load_project=%s\n' "$LOAD_PROJECT" + printf 'worker_replicas=%s\n' "$LOAD_WORKER_REPLICAS" + printf 'maintenance_concurrency_per_worker=%s\n' \ + "${OBAN_MAINTENANCE_CONCURRENCY:-2}" + printf 'push_concurrency_per_worker=%s\n' "${OBAN_PUSH_CONCURRENCY:-1}" +} >"$output_dir/inputs.txt" + +for container in "${worker_containers[@]}"; do + docker exec "$container" /app/bin/who_need_help rpc \ + 'IO.inspect(%{node: node(), queues: Oban.config().queues})' \ + >>"$output_dir/worker-queue-config.txt" +done + +measured_containers=("$db_container" "${worker_containers[@]}") +touch "$sampling_marker" +( + while [[ -e "$sampling_marker" ]]; do + captured_at=$(date -u +%Y-%m-%dT%H:%M:%S.%3NZ) + + docker stats --no-stream --format json "${measured_containers[@]}" | + jq --compact-output --arg captured_at "$captured_at" \ + '. + {captured_at: $captured_at}' + done +) >"$output_dir/docker-stats.jsonl" & +sampler_pid=$! + +insert_started_ns=$(date +%s%N) +docker exec \ + "${worker_containers[0]}" \ + /app/bin/who_need_help rpc " + run_id = \"$run_id\" + count = $job_count + + jobs = + 1..count + |> Enum.map(fn sequence -> + WhoNeedHelp.Workers.LocalBurstProbe.new(%{ + \"confirmation\" => \"isolated-local-oban-burst-probe\", + \"run_id\" => run_id, + \"sequence\" => sequence + }) + end) + |> Oban.insert_all() + + IO.puts(length(jobs)) + " >"$output_dir/job-insert.txt" +cleanup_required=true + +if [[ "$(tail -n 1 "$output_dir/job-insert.txt")" != "$job_count" ]]; then + echo "Oban.insert_all did not return the requested number of jobs." >&2 + exit 1 +fi + +deadline=$((SECONDS + timeout_seconds)) + +while :; do + database_psql >"$output_dir/job-state-latest.json" <>'run_id' = '$run_id' + GROUP BY state +) AS grouped; +SQL + + total=$(jq -r '.total // 0' "$output_dir/job-state-latest.json") + completed=$(jq -r '.states.completed // 0' "$output_dir/job-state-latest.json") + + if [[ "$total" == "$job_count" && "$completed" == "$job_count" ]]; then + break + fi + + if ((SECONDS >= deadline)); then + cp "$output_dir/job-state-latest.json" "$output_dir/job-state-timeout.json" + echo "The recorded experiment timeout elapsed before all jobs completed." >&2 + exit 1 + fi + + sleep 0.1 +done + +completed_ns=$(date +%s%N) +cp "$output_dir/job-state-latest.json" "$output_dir/job-state-final.json" + +rm -f "$sampling_marker" +wait "$sampler_pid" >/dev/null 2>&1 || true +sampler_pid= + +elapsed_ms=$(((completed_ns - insert_started_ns) / 1000000)) +printf '%s\n' "$elapsed_ms" >"$output_dir/elapsed-ms.txt" + +snapshot_domain_counts >"$output_dir/domain-counts-after.json" + +if ! cmp -s \ + "$output_dir/domain-counts-before.json" \ + "$output_dir/domain-counts-after.json"; then + echo "Domain row counts changed during the isolated burst measurement." >&2 + exit 1 +fi + +deleted=$(delete_exact_jobs) +printf '%s\n' "$deleted" >"$output_dir/job-cleanup.txt" + +if [[ "$deleted" != "$job_count" || "$(exact_job_count)" != "0" ]]; then + echo "The script did not remove exactly its own burst jobs." >&2 + exit 1 +fi + +cleanup_required=false + +jq -s ' + group_by(.Name) + | map({ + container: .[0].Name, + samples: length, + max_memory_bytes: + (map(.MemUsage | split(" / ")[0]) | map( + capture("(?[0-9.]+)(?[KMG]iB)") + | (.value | tonumber) * + (if .unit == "KiB" then 1024 + elif .unit == "MiB" then 1048576 + else 1073741824 end) + ) | max), + max_cpu_percent: + (map(.CPUPerc | rtrimstr("%") | tonumber) | max) + }) +' "$output_dir/docker-stats.jsonl" >"$output_dir/resource-summary.json" + +printf 'Completed %s isolated Oban jobs in %s ms. Evidence: %s\n' \ + "$job_count" "$elapsed_ms" "$output_dir" diff --git a/scripts/quality.sh b/scripts/quality.sh index b392605..0cd216c 100755 --- a/scripts/quality.sh +++ b/scripts/quality.sh @@ -91,6 +91,8 @@ docker compose --env-file .env.example -f compose.yaml config --format json | and $root.services.web.environment.POOL_SIZE == "4" and $root.services.worker.environment.POOL_SIZE == "2" and $root.services.migrate.environment.POOL_SIZE == "2" + and $root.services.worker.environment.OBAN_MAINTENANCE_CONCURRENCY == "2" + and $root.services.worker.environment.OBAN_PUSH_CONCURRENCY == "1" and $root.services.web.deploy.replicas == 2 and $root.services.worker.deploy.replicas == 2 and ($root.services.proxy.networks | keys | sort) == ["docker-api", "edge", "ingress"] @@ -224,6 +226,9 @@ test "$( test "$( grep -c 'value: "+Q 65536"' "$scan_dir/rendered-helm.yaml" )" -eq 5 +test "$( + grep -c '^kind: NetworkPolicy$' "$scan_dir/rendered-helm.yaml" +)" -eq 1 mkdir -p "$ROOT/.tools/trivy-cache" docker run --rm \ --volume "$scan_dir:/scan:ro" \ diff --git a/test/who_need_help/workers/local_burst_probe_test.exs b/test/who_need_help/workers/local_burst_probe_test.exs new file mode 100644 index 0000000..ae535db --- /dev/null +++ b/test/who_need_help/workers/local_burst_probe_test.exs @@ -0,0 +1,23 @@ +defmodule WhoNeedHelp.Workers.LocalBurstProbeTest do + use ExUnit.Case, async: true + use Oban.Testing, repo: WhoNeedHelp.Repo + + alias WhoNeedHelp.Workers.LocalBurstProbe + + test "accepts an explicitly confirmed isolated burst job" do + assert :ok = + perform_job(LocalBurstProbe, %{ + "confirmation" => "isolated-local-oban-burst-probe", + "run_id" => "unit-burst", + "sequence" => 1 + }) + end + + test "cancels jobs without the exact confirmation" do + assert {:cancel, :invalid_local_burst_probe} = + perform_job(LocalBurstProbe, %{ + "run_id" => "unit-burst", + "sequence" => 1 + }) + end +end