diff --git a/.github/workflows/checks.yml b/.github/workflows/checks.yml index 93961db..96fe970 100644 --- a/.github/workflows/checks.yml +++ b/.github/workflows/checks.yml @@ -52,19 +52,33 @@ jobs: run: | python3 -m pip install ruff==0.16.5 ruff check --target-version py312 \ + install/lib/prometheus_agent.py \ tools/seismic-node.py \ tools/seismic_node \ tools/checkpoint-start/summit-checkpoint-runner.py \ tests ruff format --check --target-version py312 \ + install/lib/prometheus_agent.py \ tools/seismic-node.py \ tools/seismic_node \ tools/checkpoint-start/summit-checkpoint-runner.py \ tests + - name: Fetch checksum-pinned Prometheus for isolated agent tests + run: | + read -r version checksum < <(python3 -c 'import runpy; m = runpy.run_path("install/lib/prometheus_agent.py"); print(m["VERSION"], m["CHECKSUMS"]["amd64"])') + archive="$RUNNER_TEMP/prometheus.tar.gz" + curl -fsSL --max-time 120 "https://github.com/prometheus/prometheus/releases/download/v$version/prometheus-$version.linux-amd64.tar.gz" -o "$archive" + printf '%s %s\n' "$checksum" "$archive" | sha256sum -c - + tar -xzf "$archive" -C "$RUNNER_TEMP" "prometheus-$version.linux-amd64/prometheus" "prometheus-$version.linux-amd64/promtool" + echo "PROMETHEUS_AGENT_BIN=$RUNNER_TEMP/prometheus-$version.linux-amd64/prometheus" >> "$GITHUB_ENV" + echo "PROMETHEUS_AGENT_PROMTOOL=$RUNNER_TEMP/prometheus-$version.linux-amd64/promtool" >> "$GITHUB_ENV" + python3 -m pip install supervisor==4.2.5 setuptools==80.9.0 + - name: Check Python compilation and tests run: | python3 -m py_compile \ + install/lib/prometheus_agent.py \ tools/seismic-node.py \ tools/seismic_node/*.py \ tools/checkpoint-start/summit-checkpoint-runner.py diff --git a/install/OBSERVER.md b/install/OBSERVER.md index 73ec936..68e9786 100644 --- a/install/OBSERVER.md +++ b/install/OBSERVER.md @@ -14,6 +14,11 @@ after installation. For a validator node, use `install/install-validator.sh` and follow the separate **[Validator Installer and First-Start Guide](README.md)**. +Optional authenticated push monitoring is shared with the validator installer: +see [Prometheus Agent](PROMETHEUS_AGENT.md). Register observer push nodes with +role `observer` on the monitoring server so they do not affect validator quorum +alerts. + ## Safety model The installer is designed to avoid replacing persistent observer state: @@ -28,7 +33,8 @@ The installer is designed to avoid replacing persistent observer state: - Existing observer Custodian root keys are preserved for verification against the parent Custodian. - Service binaries are installed as root-owned, non-writable executables. -- Supervisor programs use `autostart=false` and `autorestart=false`. +- Node Supervisor programs use `autostart=false` and `autorestart=false`. The + optional independent Prometheus Agent uses `autorestart=true` once started. - The installer does not start, enable, reread, update, reload, or restart observer services. @@ -154,10 +160,10 @@ Summit, seismic-reth, summit-checkpointer, and Centralized Custodian support: The current source-build defaults are: ```text -Summit: main -seismic-reth: feat/purpose-key-rotation-reth -Checkpointer: main -Custodian: d/centralized-custodian +Summit: internal-testnet-v0 (tag) +seismic-reth: internal-testnet-v0 (tag) +Checkpointer: main (branch) +Custodian: internal-testnet-v0 (tag, enclave repository) ``` A prebuilt or already-present deferred summit-checkpointer must support @@ -748,11 +754,27 @@ On a normal rerun, the installer: Changing a persistent path does not migrate existing data. It creates or uses a separate store. -For source installations, new checkouts fetch all remote branches. On a rerun, -the installer validates the origin and clean working tree, configures `origin` -to fetch all branches, fetches and prunes remote references, checks out or -creates the configured local branch, and merges `origin/` with -`--ff-only` before rebuilding. Dirty, diverged, or force-pushed checkouts are +For source installations, refs are configured in `install/lib/configuration.sh`. +Use `refs/tags/` for a release tag; unqualified names (or +`refs/heads/`) select branches. Summit, seismic-reth, and Custodian +default to the explicit tag `refs/tags/internal-testnet-v0`; these tags must be +published in their respective repositories before installation. Checkpointer +remains on the `main` branch. + +New checkouts fetch all remote branches. Every tag installation fetches the +exact tag without forcing or pruning tags, resolves it to a commit, and checks +it out with detached HEAD. Both lightweight and annotated tags are supported. +Reruns stay on that release even when branches advance. A missing/deleted tag or +a remote tag that differs from the existing local tag is rejected, rather than +silently switching releases. Publish a new tag for a new release; do not move an +existing release tag. This is not signature verification or an independent +commit-SHA pin. + +On a rerun, the installer validates the origin and clean working tree before +switching refs. Existing branch checkouts can migrate to a tag without deleting +their local branches. Branch-mode updates expand legacy single-branch fetch +configuration, fetch/prune remote branches, and merge `origin/` with +`--ff-only`. Dirty working trees and non-fast-forward branch updates are rejected rather than reset. The installer does not update its own `seismic-node-ops` checkout. diff --git a/install/PROMETHEUS_AGENT.md b/install/PROMETHEUS_AGENT.md new file mode 100644 index 0000000..e8a34eb --- /dev/null +++ b/install/PROMETHEUS_AGENT.md @@ -0,0 +1,176 @@ +# Optional Prometheus Agent + +Both node installers can install an optional, independent Prometheus Agent. +Fresh installations default to disabled. Installing the component starts no +services. Silencing and maintenance exclusions are not part of this feature. + +## Prerequisites + +- Enable authenticated remote-write ingress on the monitoring VM and provision + its DNS, TLS certificate, and TCP 443 access. +- Issue a token there with + `sudo monitoring-token add --output /root/.token`. The helper + lives in the separate deploy repository under + `monitoring/monitoring-token.py`; the node installer does not issue tokens. +- Transfer the handoff securely to the node, owned by root with mode `0600`. + Never pass the token itself on the command line. The installer asks for the + **file path**, not its contents. +- Choose a stable lowercase hostname-style node label. Use the same name in the + monitoring server's expected push inventory. + +## Installer configuration + +Run `install/install-validator.sh` or `install/install-observer.sh` normally and +opt in at the Prometheus Agent prompt. The configuration review includes an +agent edit action. + +Inputs: + +- Remote-write URL: `https://metrics.example.com/api/v1/write`. Port 443, valid + certificate verification, no URL credentials, query strings, or fragments. +- Stable node name: for example `internal-0.seismictest.net`. +- Root-only token handoff file. +- Dedicated WAL directory, default `/var/lib/seismic-prometheus-agent`. It must + not overlap node data, keys, checkpoint directories, or managed configuration. + +The installer downloads official Prometheus **3.5.0** for Linux amd64 or arm64, +checks a pinned SHA-256 checksum, and installs Prometheus and promtool under: + +```text +/usr/local/lib/seismic/prometheus-agent/3.5.0/ +``` + +Other paths: + +```text +/etc/seismic/prometheus-agent/prometheus.yml +/etc/seismic/prometheus-agent/token +/etc/seismic/prometheus-agent/installation.json +/etc/supervisor/conf.d/prometheus-agent.conf +/var/log/seismic-prometheus-agent/ +``` + +The dedicated `seismic-prometheus` account runs the agent. Configuration and +credentials are root-owned and group-readable only by that account. The settings +JSON is root-only and contains no credential. Non-secret settings are also +recorded in an optional `[monitoring]` table in the node installation inventory. +Older inventories without this table remain supported. + +## Collection and labels + +The agent scrapes every 15 seconds: + +| Target | Local address | Job | +| ------ | ---------------- | ------------------------- | +| Summit | `127.0.0.1:9090` | `summit/` | +| Reth | `127.0.0.1:9001` | `reth/` | +| Agent | `127.0.0.1:9091` | `prometheus-agent/` | + +Summit and Reth use path `/`, matching the existing OpenResty proxy upstreams; +the agent's own metrics use `/metrics`. + +Each job has `node=`, `instance=`, and `role=validator` or +`role=observer`. No local OpenResty JWT is needed. The agent initiates outbound +HTTPS to the remote-write URL; its port 9091 is a **local** status/metrics +listener, not the destination. No new inbound firewall rule is required. + +The bearer token authorizes ingestion; it does not enforce ownership of metric +labels. These credentials are for trusted node operators, not tenant isolation. + +## Lifecycle + +After successful validator or observer startup, `seismic-node` ensures a +configured agent is running. This covers normal and checkpoint start commands, +and the actual startup phase of validator onboarding. Waiting, refused +onboarding, identity failures, and checkpoint installation alone do not start +it. An already-running agent is left alone. A startup/verification failure emits +a warning without rolling back node startup or requiring remote connectivity. + +The agent uses `autostart=false` and `autorestart=true`. It does not start +merely because the installer loads configuration. Once started, Supervisor +restarts it on unexpected exits, subject to Supervisor's startup retry limits. + +**Node stop and startup rollback do not stop the agent.** It can continue +sending failed scrape results and draining buffered samples. Explicit commands: + +```bash +sudo ./tools/seismic-node.py monitoring status --role validator +sudo ./tools/seismic-node.py monitoring start --role validator +sudo ./tools/seismic-node.py monitoring stop --role validator +``` + +Use `--role observer` for observers and `--inventory /absolute/path.toml` for a +custom inventory. Monitoring start updates only the agent's Supervisor group; +status and stop do not reload configuration or change node programs. A later +node start will ensure the agent runs again if it remains enabled in inventory. + +## Reinstallation, rotation, and disabling + +An existing installation offers **keep**, **update**, or **disable**. Keep is +the default and leaves the binary, configuration, token, WAL, and running +process untouched. The dedicated agent settings are reused, not the old node +inventory. + +Before updating or rotating the token, explicitly stop the agent. Update refuses +running or indeterminate Supervisor states. If Supervisor is unavailable, +inspect and restore its control interface first; the installer will not assume +an agent is stopped. Choose update, provide the new handoff path (or retain the +installed token), and start the agent afterwards. A nonempty unrelated WAL +directory is never adopted. Generated configuration is checked with +`promtool check config --agent` before replacing credentials/configuration. + +Disable explicitly stops the agent and removes its managed Supervisor file, +preserving the token and WAL. Successful installation records monitoring as +disabled. Disabling locally or deleting a handoff does **not** revoke the token; +revoke it separately on the monitoring VM when appropriate. + +Do not run node installers, node starts, or agent updates concurrently. + +## Buffering and monitoring-side rollout + +The WAL is compressed. Minimum retention is 5 minutes; samples older than 6 +hours may be forcibly deleted during truncation, even if they were never sent. +This is **not a disk-size limit or a six-hour delivery guarantee**. Monitor free +space, especially during outages. Remote-write concurrency is limited to two +shards and HTTP 429 responses are retried with backoff. + +On the monitoring VM, replace that node's pull `NODE` entry with: + +```text +PUSH_NODE internal-0.seismictest.net validator +PUSH_NODE observer-0.seismictest.net observer +``` + +The role defaults to validator if omitted. Other validators can remain on pull. +Do not have both collectors send the same node's Summit/Reth series. Arrange the +cutover to stop central pulling before starting the agent; temporary missing +telemetry alerts are possible during the transition. + +The monitoring deployment generates expected agent, Summit, and Reth streams. +Missing data is eligible after a 2-minute freshness window and alerts after an +additional 2 minutes continuously missing (plus evaluation/notification delay). +A never-connected expected node alerts after the 2-minute pending period. Fresh +`up=0` uses the existing scrape-failure alerts instead. Observer failures notify +Slack but do not contribute to validator fleet counts. + +## Verification + +From the repository root: + +```bash +python3 -m unittest discover -s tests -p 'test_prometheus_agent*.py' -v +``` + +To additionally test the actual pinned binary, provide locally downloaded, +checksum-verified executables: + +```bash +PROMETHEUS_AGENT_BIN=/path/to/prometheus-3.5.0.linux-amd64/prometheus \ +PROMETHEUS_AGENT_PROMTOOL=/path/to/prometheus-3.5.0.linux-amd64/promtool \ +python3 -m unittest discover -s tests -p 'test_prometheus_agent*.py' -v +``` + +The opt-in integration test uses ephemeral loopback ports, a temporary trusted +TLS certificate, fake exporters, and a temporary Prometheus receiver. It does +not install software, invoke Supervisor, or contact live nodes. Monitoring-side +push inventory and alert tests live in the deploy repository. diff --git a/install/README.md b/install/README.md index e42159e..c9fde42 100644 --- a/install/README.md +++ b/install/README.md @@ -14,6 +14,11 @@ after installation. For an observer node, use `install/install-observer.sh` and follow the separate **[Observer Installer and First-Start Guide](OBSERVER.md)**. +Optional authenticated push monitoring is shared with the observer installer: +see [Prometheus Agent](PROMETHEUS_AGENT.md). When configured, `seismic-node` +starts the agent after successful node startup and leaves it running on node +stop. + ## Safety model The installer is designed to avoid replacing persistent validator state: @@ -24,7 +29,8 @@ The installer is designed to avoid replacing persistent validator state: - An incomplete Summit key pair causes the installer to stop rather than regenerate either key. - Service binaries are installed as root-owned, non-writable executables. -- Supervisor programs use `autostart=false` and `autorestart=false`. +- Node Supervisor programs use `autostart=false` and `autorestart=false`. The + optional independent Prometheus Agent uses `autorestart=true` once started. - The installer does not start, enable, reread, update, reload, or restart validator services. @@ -132,10 +138,10 @@ supported installation modes are: The current source-build defaults are: ```text -Summit: main -seismic-reth: feat/purpose-key-rotation-reth -Checkpointer: main -Custodian: d/centralized-custodian +Summit: internal-testnet-v0 (tag) +seismic-reth: internal-testnet-v0 (tag) +Checkpointer: main (branch) +Custodian: internal-testnet-v0 (tag, enclave repository) ``` Deferred binaries must be installed at the configured target paths before the @@ -363,13 +369,56 @@ coordinate startup: ```bash sudo ./tools/seismic-node.py validator onboard \ - --deposit-signature /root/deposit-signature.json \ --summit-rpc-url https://trusted-validator.example/summit \ --snapshot-api-url https://snapshot.example/checkpointer \ --snapshot-bearer-token-file /root/snapshot-token \ --weak-subjectivity-rpc-url https://independent-validator.example/summit ``` +Onboarding reads `summit_keys_dir` from the installation inventory and uses +`/usr/local/bin/summit keys show` to derive the installed node public key +locally. The default inventory is `/etc/seismic/validator-installation.toml`; +select a custom file with `--inventory /absolute/path/to/installation.toml`. +Only the public key is sent to the trusted Summit RPC. No signing endpoint is +started. + +The deposit-signature file is not required for onboarding. Existing commands may +still pass `--deposit-signature /root/deposit-signature.json` as an optional +identity cross-check; a mismatch with the installed node key is rejected before +polling or checkpoint installation. The installed public key is checked again +before startup so a key change during the wait cannot start a different +identity. Generating and submitting the deposit remains a separate prerequisite. + +#### Onboard without a checkpoint + +To wait for `Joining` and then start from local state without installing or +requiring a checkpoint: + +```bash +sudo ./tools/seismic-node.py validator onboard \ + --mode normal \ + --summit-rpc-url https://trusted-validator.example/summit \ + --pre-joining-policy wait +``` + +This can be run after submitting the deposit transaction. It waits while the +account is `NotFound` or `Inactive`, starts when it reaches `Joining`, and also +allows `Active` with a warning. By default it waits indefinitely; use +`--validator-wait-timeout SECONDS` to bound the wait. Status is checked again +immediately before startup, using the same wait deadline. + +Normal mode starts `summit` instead of `summit-checkpoint` and does not read +checkpoint-start configuration or replace local state. It rejects checkpoint +source and installation options, including `--yes`. It does not guarantee that a +freshly reset node can synchronize with the running network from genesis. Unlike +`validator start --mode normal`, it applies the onboarding lifecycle checks. + +Checkpoint mode remains the default (`--mode checkpoint`). In that mode, +omitting source options uses an already-installed checkpoint; it does not switch +to normal startup. + +#### Startup behavior + When startup is authorized, `validator onboard` automatically runs: ```bash @@ -378,7 +427,8 @@ sudo supervisorctl reread sudo supervisorctl update ``` -It then starts Custodian when configured, Reth, `summit-checkpoint`, and +It then starts Custodian when configured, Reth, the selected Summit program +(`summit` in normal mode or `summit-checkpoint` in checkpoint mode), and summit-checkpointer when configured. It does not start or reload OpenResty. For unattended installation, pass `--yes` to skip the interactive confirmation, @@ -660,11 +710,27 @@ On a normal rerun, the installer preserves existing validator keys and state, then replaces its generated OpenResty and Supervisor configuration. It does not start or reload those services. -For source installations, new checkouts fetch all remote branches. On a rerun, -the installer validates the origin and clean working tree, configures `origin` -to fetch all branches, fetches and prunes remote references, checks out or -creates the configured local branch, and merges `origin/` with -`--ff-only` before rebuilding. Dirty, diverged, or force-pushed checkouts are +For source installations, refs are configured in `install/lib/configuration.sh`. +Use `refs/tags/` for a release tag; unqualified names (or +`refs/heads/`) select branches. Summit, seismic-reth, and Custodian +default to the explicit tag `refs/tags/internal-testnet-v0`; these tags must be +published in their respective repositories before installation. Checkpointer +remains on the `main` branch. + +New checkouts fetch all remote branches. Every tag installation fetches the +exact tag without forcing or pruning tags, resolves it to a commit, and checks +it out with detached HEAD. Both lightweight and annotated tags are supported. +Reruns stay on that release even when branches advance. A missing/deleted tag or +a remote tag that differs from the existing local tag is rejected, rather than +silently switching releases. Publish a new tag for a new release; do not move an +existing release tag. This is not signature verification or an independent +commit-SHA pin. + +On a rerun, the installer validates the origin and clean working tree before +switching refs. Existing branch checkouts can migrate to a tag without deleting +their local branches. Branch-mode updates expand legacy single-branch fetch +configuration, fetch/prune remote branches, and merge `origin/` with +`--ff-only`. Dirty working trees and non-fast-forward branch updates are rejected rather than reset. The installer does not update its own `seismic-node-ops` checkout. diff --git a/install/install-observer.sh b/install/install-observer.sh index cb0833e..2c55242 100755 --- a/install/install-observer.sh +++ b/install/install-observer.sh @@ -42,12 +42,16 @@ source "$SCRIPT_DIR/lib/supervisor.sh" source "$SCRIPT_DIR/lib/observer-supervisor.sh" # shellcheck source=lib/configuration.sh source "$SCRIPT_DIR/lib/configuration.sh" +# shellcheck source=lib/prometheus-agent.sh +source "$SCRIPT_DIR/lib/prometheus-agent.sh" # shellcheck source=lib/observer-configuration.sh source "$SCRIPT_DIR/lib/observer-configuration.sh" # shellcheck source=lib/observer-instructions.sh source "$SCRIPT_DIR/lib/observer-instructions.sh" write_observer_installation_inventory() { + local monitoring + monitoring=$(write_prometheus_agent_inventory) || die "Could not record agent configuration." install_installation_inventory "$INSTALLATION_INVENTORY_PATH" <." ;; + *) + ref_kind=branch + ref_name=$source_ref + full_ref="refs/heads/$ref_name" + clone_options+=(--branch "$ref_name") + ;; + esac + run_as_service_user git check-ref-format "$full_ref" \ + >>"$LOG_FILE" 2>&1 \ + || die "$description source ref is invalid: $source_ref" + [[ "$ref_name" != -* && "$ref_name" != HEAD ]] \ + || die "$description source ref is invalid: $source_ref" + [[ ! -L "$source_dir" ]] \ + || die "$description source directory must not be a symbolic link: $source_dir" prepare_source_root if [[ ! -e "$source_dir" ]]; then - info "Cloning $description branch $branch into $source_dir..." + info "Cloning $description $ref_kind $ref_name into $source_dir..." if ! run_as_service_user git clone \ - --branch "$branch" \ + "${clone_options[@]}" \ "$repository" \ "$source_dir" >>"$LOG_FILE" 2>&1; then die "Could not clone $description; see $LOG_FILE" fi + # git clone --branch also accepts tag names; branch mode must not. + if [[ "$ref_kind" == branch ]]; then + run_as_service_user git -C "$source_dir" show-ref \ + --verify --quiet "$full_ref" \ + || die "$description branch $ref_name was not found; use refs/tags/ for a tag." + fi else - [[ ! -L "$source_dir" ]] \ - || die "$description source directory must not be a symbolic link: $source_dir" [[ -d "$source_dir/.git" ]] \ || die "$description source path exists but is not a Git checkout: $source_dir" @@ -196,35 +232,51 @@ prepare_source_checkout() { die "$description checkout has local changes; refusing to overwrite them: $source_dir" fi - info "Updating $description branch $branch with fast-forward only..." - run_as_service_user git -C "$source_dir" remote \ - set-branches origin '*' \ - >>"$LOG_FILE" 2>&1 \ - || die "Could not configure $description to fetch all origin branches." - run_as_service_user git -C "$source_dir" fetch --prune origin \ - >>"$LOG_FILE" 2>&1 \ - || die "Could not fetch $description branches; see $LOG_FILE" - - if run_as_service_user git -C "$source_dir" show-ref \ - --verify --quiet "refs/heads/$branch"; then - run_as_service_user git -C "$source_dir" checkout "$branch" \ + if [[ "$ref_kind" == branch ]]; then + info "Updating $description branch $ref_name with fast-forward only..." + run_as_service_user git -C "$source_dir" remote \ + set-branches origin '*' \ >>"$LOG_FILE" 2>&1 \ - || die "Could not check out $description branch $branch." - else - run_as_service_user git -C "$source_dir" checkout \ - -b "$branch" --track "origin/$branch" \ + || die "Could not configure $description to fetch all origin branches." + run_as_service_user git -C "$source_dir" fetch --no-tags --prune origin \ >>"$LOG_FILE" 2>&1 \ - || die "Could not create local $description branch $branch." + || die "Could not fetch $description branches; see $LOG_FILE" + + if run_as_service_user git -C "$source_dir" show-ref \ + --verify --quiet "$full_ref"; then + run_as_service_user git -C "$source_dir" checkout "$ref_name" -- \ + >>"$LOG_FILE" 2>&1 \ + || die "Could not check out $description branch $ref_name." + else + run_as_service_user git -C "$source_dir" checkout \ + -b "$ref_name" --track "origin/$ref_name" \ + >>"$LOG_FILE" 2>&1 \ + || die "Could not create local $description branch $ref_name." + fi + + run_as_service_user git -C "$source_dir" merge \ + --ff-only "refs/remotes/origin/$ref_name" >>"$LOG_FILE" 2>&1 \ + || die "$description branch cannot be updated with fast-forward only." fi + fi - run_as_service_user git -C "$source_dir" merge \ - --ff-only "origin/$branch" >>"$LOG_FILE" 2>&1 \ - || die "$description branch cannot be updated with fast-forward only." + if [[ "$ref_kind" == tag ]]; then + info "Checking $description release tag $ref_name..." + # Fetch the exact tag, even in a legacy --single-branch/--no-tags clone. + # Never force or prune tags: a moved/deleted release must fail closed. + run_as_service_user git -C "$source_dir" fetch --no-tags --no-prune origin \ + "$full_ref:$full_ref" >>"$LOG_FILE" 2>&1 \ + || die "Could not fetch $description tag $ref_name unchanged; it may be missing or moved. See $LOG_FILE" + commit=$(run_as_service_user git -C "$source_dir" rev-parse --verify "$full_ref^{commit}") \ + || die "$description tag $ref_name does not resolve to a commit." + run_as_service_user git -C "$source_dir" checkout --detach "$commit" -- \ + >>"$LOG_FILE" 2>&1 \ + || die "Could not check out $description tag $ref_name." fi commit=$(run_as_service_user git -C "$source_dir" rev-parse HEAD) \ || die "Could not determine the installed $description source revision." - success "$description source ready: $branch at $commit" + success "$description source ready: $source_ref at $commit" } install_binary_target() { diff --git a/install/lib/configuration.sh b/install/lib/configuration.sh index fd28ce0..f1751b8 100644 --- a/install/lib/configuration.sh +++ b/install/lib/configuration.sh @@ -454,8 +454,9 @@ configure_node_software() { SUMMIT_TARGET_BIN="/usr/local/bin/summit" RETH_TARGET_BIN="/usr/local/bin/seismic-reth" - SUMMIT_SOURCE_REF="main" - RETH_SOURCE_REF="feat/purpose-key-rotation-reth" + # Explicit tag refs cannot be mistaken for same-named branches. + SUMMIT_SOURCE_REF="refs/tags/internal-testnet-v0" + RETH_SOURCE_REF="refs/tags/internal-testnet-v0" SUMMIT_INSTALL_METHOD="" RETH_INSTALL_METHOD="" SUMMIT_BINARY="" @@ -791,9 +792,9 @@ configure_custodian() { COUNCIL_ADDRESS="" CUSTODIAN_CHAIN_ID="" PARENT_CUSTODIAN="" - CUSTODIAN_REQUIRED_SUMMIT_REF="main" - CUSTODIAN_REQUIRED_RETH_REF="feat/purpose-key-rotation-reth" - CUSTODIAN_SOURCE_REF="d/centralized-custodian" + CUSTODIAN_REQUIRED_SUMMIT_REF="internal-testnet-v0" + CUSTODIAN_REQUIRED_RETH_REF="internal-testnet-v0" + CUSTODIAN_SOURCE_REF="refs/tags/internal-testnet-v0" if ! confirm "Enable Centralized Custodian?"; then _out "Centralized Custodian: $INSTALL_CUSTODIAN" @@ -1035,6 +1036,7 @@ print_configuration_summary() { _out " Enabled: false" fi + print_prometheus_agent_plan print_system_package_plan } @@ -1054,6 +1056,7 @@ review_configuration() { _out " 7) Edit Centralized Custodian" _out " 8) Accept configuration" _out " 9) Cancel" + _out " 10) Edit Prometheus Agent" prompt selection "Select an action" "8" case "$selection" in @@ -1079,7 +1082,8 @@ review_configuration() { info "Configuration cancelled; no installation changes were made." exit 0 ;; - *) error "Select a number from 1 through 9." ;; + 10) configure_prometheus_agent ;; + *) error "Select a number from 1 through 10." ;; esac done } @@ -1094,5 +1098,6 @@ configure() { configure_validator_software configure_checkpointer configure_custodian + configure_prometheus_agent review_configuration } diff --git a/install/lib/observer-configuration.sh b/install/lib/observer-configuration.sh index bfabbc3..0014a86 100644 --- a/install/lib/observer-configuration.sh +++ b/install/lib/observer-configuration.sh @@ -192,6 +192,7 @@ print_observer_configuration_summary() { fi fi + print_prometheus_agent_plan print_system_package_plan } @@ -212,6 +213,7 @@ review_observer_configuration() { _out " 8) Edit Centralized Custodian" _out " 9) Accept configuration" _out " 10) Cancel" + _out " 11) Edit Prometheus Agent" prompt selection "Select an action" "9" case "$selection" in @@ -238,7 +240,8 @@ review_observer_configuration() { info "Configuration cancelled; no installation changes were made." exit 0 ;; - *) error "Select a number from 1 through 10." ;; + 11) configure_prometheus_agent ;; + *) error "Select a number from 1 through 11." ;; esac done } @@ -254,5 +257,6 @@ configure_observer() { configure_node_software configure_checkpointer configure_custodian + configure_prometheus_agent review_observer_configuration } diff --git a/install/lib/prometheus-agent.sh b/install/lib/prometheus-agent.sh new file mode 100644 index 0000000..6d14aec --- /dev/null +++ b/install/lib/prometheus-agent.sh @@ -0,0 +1,117 @@ +#!/usr/bin/env bash + +# Shared optional telemetry component. Never source credentials or saved settings. +PROMETHEUS_AGENT_HELPER="$SCRIPT_DIR/lib/prometheus_agent.py" +PROMETHEUS_AGENT_ACTION=keep +PROMETHEUS_AGENT_ENABLED=false + +agent_saved_field() { + python3 "$PROMETHEUS_AGENT_HELPER" show --field "$1" +} + +agent_arguments() { + AGENT_ARGUMENTS=( + --node "$PROMETHEUS_AGENT_NODE" + --role "$NODE_ROLE" + --url "$PROMETHEUS_AGENT_URL" + --data-dir "$PROMETHEUS_AGENT_DATA_DIR" + --token-source "$PROMETHEUS_AGENT_TOKEN_SOURCE" + ) + local name + for name in RETH_DATA_DIR RETH_P2P_KEY_PATH SUMMIT_DATA_DIR SUMMIT_KEYS_DIR VALIDATOR_KEYS_DIR CHECKPOINTS_DIR CUSTODIAN_DATA_DIR; do + [[ -z "${!name:-}" ]] || AGENT_ARGUMENTS+=(--exclude-path "${!name}") + done +} + +configure_prometheus_agent() { + local existing selection saved_role + section "Optional Prometheus Agent" + existing=$(agent_saved_field enabled) || die "Could not read existing agent settings." + if [[ -n "$existing" ]]; then + saved_role=$(agent_saved_field role) + [[ "$saved_role" == "$NODE_ROLE" ]] || die "Existing agent belongs to a different node role." + PROMETHEUS_AGENT_ENABLED=$existing + _out "Existing agent enabled: $existing" + prompt selection "Agent action: keep, update, disable" "keep" + case "$selection" in + keep) + PROMETHEUS_AGENT_ACTION=keep + return + ;; + disable) + confirm "Stop and disable this agent? Credentials and WAL will be preserved." || return + PROMETHEUS_AGENT_ACTION=disable + PROMETHEUS_AGENT_ENABLED=false + return + ;; + update) ;; + *) die "Select keep, update, or disable." ;; + esac + elif ! confirm "Install Prometheus Agent for authenticated remote write?"; then + PROMETHEUS_AGENT_ENABLED=false + PROMETHEUS_AGENT_ACTION=keep + return + fi + PROMETHEUS_AGENT_NODE=$(agent_saved_field node) + PROMETHEUS_AGENT_URL=$(agent_saved_field url) + PROMETHEUS_AGENT_DATA_DIR=$(agent_saved_field data_dir) + prompt PROMETHEUS_AGENT_NODE "Stable monitoring node name (lowercase hostname)" "${PROMETHEUS_AGENT_NODE:-$(hostname -f)}" + prompt PROMETHEUS_AGENT_URL "Remote-write HTTPS URL ending in /api/v1/write" "$PROMETHEUS_AGENT_URL" + prompt PROMETHEUS_AGENT_DATA_DIR "Agent WAL directory" "${PROMETHEUS_AGENT_DATA_DIR:-/var/lib/seismic-prometheus-agent}" + local default_token="" + [[ -z "$existing" ]] || default_token=/etc/seismic/prometheus-agent/token + prompt PROMETHEUS_AGENT_TOKEN_SOURCE "Root-only token handoff file (contents will not be displayed)" "$default_token" + agent_arguments + python3 "$PROMETHEUS_AGENT_HELPER" check "${AGENT_ARGUMENTS[@]}" || die "Invalid agent configuration." + PROMETHEUS_AGENT_ENABLED=true + PROMETHEUS_AGENT_ACTION=install +} + +print_prometheus_agent_plan() { + _out "Prometheus Agent: $PROMETHEUS_AGENT_ENABLED (action: $PROMETHEUS_AGENT_ACTION)" + if [[ "$PROMETHEUS_AGENT_ACTION" == install ]]; then + _out " Version: 3.5.0; node: $PROMETHEUS_AGENT_NODE; role: $NODE_ROLE" + _out " Remote write: $PROMETHEUS_AGENT_URL" + _out " WAL: $PROMETHEUS_AGENT_DATA_DIR (maximum sample age 6h; not a disk-size limit)" + _out " No services start during installation. Token contents are never printed." + fi +} + +validate_prometheus_agent_plan() { + local name + if [[ "$PROMETHEUS_AGENT_ACTION" == install ]]; then + agent_arguments + python3 "$PROMETHEUS_AGENT_HELPER" check "${AGENT_ARGUMENTS[@]}" + elif [[ "$PROMETHEUS_AGENT_ACTION" == keep ]]; then + local excluded=() + for name in RETH_DATA_DIR RETH_P2P_KEY_PATH SUMMIT_DATA_DIR SUMMIT_KEYS_DIR VALIDATOR_KEYS_DIR CHECKPOINTS_DIR CUSTODIAN_DATA_DIR; do + [[ -z "${!name:-}" ]] || excluded+=(--exclude-path "${!name}") + done + python3 "$PROMETHEUS_AGENT_HELPER" check-existing "${excluded[@]}" + fi +} + +install_prometheus_agent() { + case "$PROMETHEUS_AGENT_ACTION" in + keep) return ;; + disable) python3 "$PROMETHEUS_AGENT_HELPER" disable ;; + install) + agent_arguments + python3 "$PROMETHEUS_AGENT_HELPER" install "${AGENT_ARGUMENTS[@]}" + ;; + esac +} + +write_prometheus_agent_inventory() { + python3 "$PROMETHEUS_AGENT_HELPER" inventory +} + +print_prometheus_agent_instructions() { + [[ "$PROMETHEUS_AGENT_ENABLED" == true ]] || return 0 + _out "Prometheus Agent is configured. Successful seismic-node $NODE_ROLE startup ensures it is running." + _out "Node stop leaves monitoring running. Explicit management:" + _out " sudo ./tools/seismic-node.py monitoring status --role $NODE_ROLE" + _out " sudo ./tools/seismic-node.py monitoring start --role $NODE_ROLE" + _out " sudo ./tools/seismic-node.py monitoring stop --role $NODE_ROLE" + _out "Register this node as PUSH_NODE on the monitoring VM; do not also pull-scrape it." +} diff --git a/install/lib/prometheus_agent.py b/install/lib/prometheus_agent.py new file mode 100755 index 0000000..e9956b5 --- /dev/null +++ b/install/lib/prometheus_agent.py @@ -0,0 +1,473 @@ +#!/usr/bin/env python3 +"""Root-side optional agent installation. No implicit service activation. + +JSON settings contain no credentials and are also the installer's rerun state. +The node inventory contains the same non-secret settings. Keep, update, and +explicit disable are distinct operations; only disable may stop a service. +""" + +from __future__ import annotations + +import argparse +import grp +import hashlib +import json +import os +import platform +import pwd +import re +import shutil +import stat +import subprocess +import tarfile +import tempfile +import time +import urllib.parse +import urllib.request +from pathlib import Path + +VERSION = "3.5.0" +CHECKSUMS = { + "amd64": "e811827af26d822afb09a4f28314f61b618b12cff5369835a67f674d8b46f39a", + "arm64": "173389cc42bf09c4e6e54cb53fa07a5a835d7c261e14775d2183181d6e385d1c", +} +HOME = Path("/etc/seismic/prometheus-agent") +SETTINGS = HOME / "installation.json" +TOKEN = HOME / "token" +CONFIG = HOME / "prometheus.yml" +BIN_DIR = Path(f"/usr/local/lib/seismic/prometheus-agent/{VERSION}") +SUPERVISOR_CONFIG = Path("/etc/supervisor/conf.d/prometheus-agent.conf") +LOG_DIR = Path("/var/log/seismic-prometheus-agent") +USER = "seismic-prometheus" +TEMPLATES = Path(__file__).resolve().parents[1] / "templates" +KEYS = {"enabled", "node", "role", "url", "data_dir", "version"} +NODE = re.compile(r"[a-z0-9](?:[a-z0-9.-]{0,251}[a-z0-9])?\Z") + + +class AgentError(Exception): + pass + + +def safe_path(path: Path) -> None: + """Reject symlinks and writable/non-root parents, including missing paths.""" + if not path.is_absolute() or str(path) != os.path.normpath(path): + raise AgentError("Agent paths must be absolute and normalized") + for part in [*reversed(path.parents), path]: + try: + metadata = part.lstat() + except FileNotFoundError: + continue + if stat.S_ISLNK(metadata.st_mode): + raise AgentError("Agent paths must not contain symbolic links") + if part != path and ( + not stat.S_ISDIR(metadata.st_mode) + or metadata.st_uid != 0 + or metadata.st_mode & 0o022 + ): + raise AgentError( + "Agent path parents must be root-owned and not writable by others" + ) + + +def root_file(path: Path) -> os.stat_result: + safe_path(path) + meta = path.lstat() + if not stat.S_ISREG(meta.st_mode) or meta.st_uid != 0 or meta.st_mode & 0o022: + raise AgentError("Expected a root-managed regular file") + return meta + + +def validate(settings: dict, excluded: list[Path] = ()) -> dict: + if set(settings) != KEYS or type(settings["enabled"]) is not bool: + raise AgentError("Invalid agent settings schema") + if settings["version"] != VERSION or settings["role"] not in ( + "validator", + "observer", + ): + raise AgentError("Unsupported agent version or node role") + node = settings["node"] + if not isinstance(node, str) or not NODE.fullmatch(node) or ".." in node: + raise AgentError("Agent node must be a lowercase hostname-style label") + url = settings["url"] + try: + parsed = urllib.parse.urlsplit(url) + valid_url = ( + isinstance(url, str) + and url.isascii() + and not re.search(r"[\s\\\"'<>]", url) + and parsed.scheme == "https" + and parsed.hostname + and parsed.port in (None, 443) + and parsed.username is None + and parsed.password is None + and parsed.path == "/api/v1/write" + and not parsed.query + and not parsed.fragment + and "?" not in url + and "#" not in url + ) + except (ValueError, TypeError): + valid_url = False + if not valid_url: + raise AgentError( + "Remote write requires HTTPS port 443 and /api/v1/write, without credentials, query or fragment" + ) + data = settings["data_dir"] + if not isinstance(data, str) or not re.fullmatch(r"/[A-Za-z0-9_./-]+", data): + raise AgentError("Unsafe agent WAL path") + path = Path(data) + if str(path) != data or os.path.normpath(data) != data: + raise AgentError("Agent WAL path must be normalized") + # WAL must not overlap node state/keys or any managed agent config/binaries. + for other in [HOME, BIN_DIR.parent, LOG_DIR, Path("/etc"), Path("/usr"), *excluded]: + if path == other or path in other.parents or other in path.parents: + raise AgentError( + "Agent WAL must be separate from node state, keys and configuration" + ) + return settings + + +def load_settings() -> dict | None: + if not SETTINGS.exists() and not SETTINGS.is_symlink(): + return None + root_file(SETTINGS) + try: + return validate(json.loads(SETTINGS.read_text())) + except (ValueError, TypeError, KeyError): + raise AgentError("Invalid existing agent settings") from None + + +def read_token(path: Path) -> bytes: + meta = root_file(path) + if meta.st_size > 128 or meta.st_mode & 0o007: + raise AgentError("Token file must be private and contain one opaque token") + if path != TOKEN and meta.st_mode & 0o077: + raise AgentError("Token handoff file must be root-only (0600)") + token = path.read_bytes().strip() + if not re.fullmatch(rb"[0-9a-f]{64}", token): + raise AgentError("Expected a 64-character opaque monitoring token") + return token + b"\n" + + +def atomic_write(path: Path, data: bytes, mode: int, gid: int = 0) -> None: + safe_path(path) + if path.exists(): + root_file(path) + fd, name = tempfile.mkstemp(prefix=f".{path.name}.", dir=path.parent) + try: + with os.fdopen(fd, "wb") as stream: + os.fchown(stream.fileno(), 0, gid) + os.fchmod(stream.fileno(), mode) + stream.write(data) + stream.flush() + os.fsync(stream.fileno()) + os.replace(name, path) + finally: + Path(name).unlink(missing_ok=True) + + +def directory(path: Path, uid: int = 0, gid: int = 0, mode: int = 0o755) -> None: + safe_path(path) + if path.exists(): + meta = path.lstat() + if not stat.S_ISDIR(meta.st_mode) or meta.st_uid not in (0, uid): + raise AgentError("Refusing to take over an existing agent directory") + path.mkdir(parents=True, exist_ok=True) + os.chown(path, uid, gid) + os.chmod(path, mode) + + +def architecture(machine: str) -> str: + try: + return {"x86_64": "amd64", "aarch64": "arm64"}[machine] + except KeyError: + raise AgentError( + "Prometheus Agent supports only Linux amd64 and arm64" + ) from None + + +def install_release(tmp: Path) -> None: + arch = architecture(platform.machine()) + filename = f"prometheus-{VERSION}.linux-{arch}.tar.gz" + url = f"https://github.com/prometheus/prometheus/releases/download/v{VERSION}/{filename}" + archive = tmp / filename + digest = hashlib.sha256() + size = 0 + with ( + urllib.request.urlopen(url, timeout=60) as response, + archive.open("wb") as output, + ): + while chunk := response.read(1024 * 1024): + size += len(chunk) + if size > 400 * 1024 * 1024: + raise AgentError("Prometheus release exceeds download limit") + digest.update(chunk) + output.write(chunk) + if digest.hexdigest() != CHECKSUMS[arch]: + raise AgentError("Prometheus release SHA-256 mismatch") + directory(BIN_DIR) + with tarfile.open(archive, "r:gz") as tar: + for name in ("prometheus", "promtool"): + member_name = f"prometheus-{VERSION}.linux-{arch}/{name}" + matches = [m for m in tar.getmembers() if m.name == member_name] + if ( + len(matches) != 1 + or not matches[0].isfile() + or matches[0].size > 512 * 1024 * 1024 + ): + raise AgentError("Invalid Prometheus release executable") + # Never extract archive paths, links or permissions onto the host. + with tar.extractfile(matches[0]) as stream: + atomic_write(BIN_DIR / name, stream.read(), 0o755) + + +def supervisor_state() -> str: + result = subprocess.run( + ["/usr/bin/supervisorctl", "status", "prometheus-agent"], + capture_output=True, + text=True, + check=False, + timeout=15, + ) + output = result.stdout.strip() + if result.returncode in (0, 3): + parts = output.split() + if len(parts) >= 2 and parts[0] == "prometheus-agent": + return parts[1] + if "no such process" in output.lower(): + return "MISSING" + # Never print subprocess output, which may contain operator configuration. + raise AgentError( + "Cannot determine agent state; inspect Supervisor before updating/disabling" + ) + + +def require_stopped() -> None: + if (SUPERVISOR_CONFIG.exists() or SETTINGS.exists()) and supervisor_state() not in ( + "MISSING", + "STOPPED", + "FATAL", + ): + raise AgentError("Stop prometheus-agent before changing its configuration") + + +def render(settings: dict, token_path: Path = TOKEN) -> tuple[str, str]: + validate(settings) + config = (TEMPLATES / "prometheus-agent.yml").read_text() + for key, value in { + "NODE_PLACEHOLDER": settings["node"], + "ROLE_PLACEHOLDER": settings["role"], + "URL_PLACEHOLDER": json.dumps(settings["url"]), + "TOKEN_PLACEHOLDER": json.dumps(str(token_path)), + }.items(): + config = config.replace(key, value) + service = (TEMPLATES / "supervisor/prometheus-agent.conf").read_text() + service = service.replace("BINARY_PLACEHOLDER", str(BIN_DIR / "prometheus")) + service = service.replace("DATA_PLACEHOLDER", settings["data_dir"]) + return config, service + + +def service_account(): + """Never grant token access to an existing shared primary group.""" + try: + account = pwd.getpwnam(USER) + except KeyError: + subprocess.run( + [ + "useradd", + "--system", + "--user-group", + "--no-create-home", + "--shell", + "/usr/sbin/nologin", + USER, + ], + check=True, + ) + account = pwd.getpwnam(USER) + if account.pw_uid == 0 or account.pw_gid == 0: + raise AgentError("Agent account must be unprivileged") + group = grp.getgrgid(account.pw_gid) + if ( + group.gr_name != USER + or set(group.gr_mem) - {USER, "root"} + or any(p.pw_gid == account.pw_gid and p.pw_name != USER for p in pwd.getpwall()) + ): + raise AgentError("Agent credentials require a dedicated private service group") + return account + + +def require_managed_or_new() -> None: + if load_settings() is None and ( + SUPERVISOR_CONFIG.exists() or SUPERVISOR_CONFIG.is_symlink() + ): + raise AgentError("Refusing to overwrite an unmanaged prometheus-agent program") + + +def install(settings: dict, token_source: Path) -> None: + require_managed_or_new() + token = read_token(token_source) + require_stopped() + account = service_account() + directory(HOME, gid=account.pw_gid, mode=0o750) + data = Path(settings["data_dir"]) + safe_path(data) + previous = load_settings() + if ( + data.exists() + and any(data.iterdir()) + and (previous is None or previous["data_dir"] != str(data)) + ): + raise AgentError("Refusing to adopt an unrelated nonempty WAL directory") + directory(data, account.pw_uid, account.pw_gid, 0o750) + directory(LOG_DIR) + # Private staging protects the handoff even during validation failures. + with tempfile.TemporaryDirectory(prefix=".install-", dir=HOME) as tmp_name: + tmp = Path(tmp_name) + install_release(tmp) + staged_token = tmp / "token" + staged_token.write_bytes(token) + staged_token.chmod(0o600) + config, service = render(settings, staged_token) + staged_config = tmp / "prometheus.yml" + staged_config.write_text(config) + result = subprocess.run( + [ + str(BIN_DIR / "promtool"), + "check", + "config", + "--agent", + str(staged_config), + ], + capture_output=True, + check=False, + timeout=30, + ) + if result.returncode: + raise AgentError( + "Generated Prometheus agent configuration did not validate" + ) + final_config, _ = render(settings) + atomic_write(TOKEN, token, 0o640, account.pw_gid) + atomic_write(CONFIG, final_config.encode(), 0o640, account.pw_gid) + atomic_write(SUPERVISOR_CONFIG, service.encode(), 0o644) + atomic_write(SETTINGS, (json.dumps(settings, indent=2) + "\n").encode(), 0o600) + print("Prometheus Agent installed; no services were started.") + + +def disable() -> None: + settings = load_settings() + if settings is None: + if SUPERVISOR_CONFIG.exists(): + raise AgentError( + "Agent config exists without managed settings; refusing to remove it" + ) + return + deadline = time.monotonic() + 60 + while supervisor_state() not in ("MISSING", "STOPPED", "FATAL"): + # EXITED may restart automatically; a NOT_RUNNING response can race + # that transition. Only STOPPED/FATAL/MISSING are stable terminal states. + result = subprocess.run( + ["/usr/bin/supervisorctl", "stop", "prometheus-agent"], + capture_output=True, + check=False, + timeout=60, + ) + if result.returncode not in (0, 7) or time.monotonic() >= deadline: + raise AgentError( + "Could not confirm agent stop; disabling was not completed" + ) + time.sleep(0.1) + require_stopped() + if SUPERVISOR_CONFIG.exists(): + root_file(SUPERVISOR_CONFIG) + SUPERVISOR_CONFIG.unlink() + settings["enabled"] = False + atomic_write(SETTINGS, (json.dumps(settings, indent=2) + "\n").encode(), 0o600) + print("Agent disabled; credentials and WAL preserved.") + + +def main() -> None: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument( + "action", + choices=("show", "check", "check-existing", "install", "disable", "inventory"), + ) + parser.add_argument("--field", choices=sorted(KEYS)) + parser.add_argument("--node") + parser.add_argument("--role", choices=("validator", "observer")) + parser.add_argument("--url") + parser.add_argument("--data-dir") + parser.add_argument("--token-source", type=Path) + parser.add_argument("--exclude-path", action="append", type=Path, default=[]) + args = parser.parse_args() + if os.geteuid() != 0: + raise AgentError("Run the installer as root") + if args.action in ("show", "inventory"): + settings = load_settings() + if args.action == "show": + if settings is not None: + print( + str(settings[args.field]).lower() + if isinstance(settings[args.field], bool) + else settings[args.field] + ) + else: + print("\n[monitoring]") + for key, value in (settings or {"enabled": False}).items(): + print(f"{key} = {json.dumps(value)}") + return + if args.action == "disable": + disable() + return + if args.action == "check-existing": + settings = load_settings() + if settings and settings["enabled"]: + validate(settings, args.exclude_path) + safe_path(Path(settings["data_dir"])) + read_token(TOKEN) + root_file(CONFIG) + root_file(SUPERVISOR_CONFIG) + return + require_managed_or_new() + settings = validate( + { + "enabled": True, + "node": args.node, + "role": args.role, + "url": args.url, + "data_dir": args.data_dir, + "version": VERSION, + }, + args.exclude_path, + ) + safe_path(Path(settings["data_dir"])) + if args.token_source is None: + raise AgentError("A root-only token source file is required") + if any( + args.token_source == path or path in args.token_source.parents + for path in args.exclude_path + ): + raise AgentError("Token handoff must be outside node data and key directories") + read_token(args.token_source) + if args.action == "install": + install(settings, args.token_source) + + +if __name__ == "__main__": + try: + main() + except AgentError as error: + import sys + + print(f"Agent: {error}", file=sys.stderr) + raise SystemExit(1) from None + except (OSError, ValueError, subprocess.SubprocessError, shutil.Error): + # Avoid echoing token contents, URLs or command output from exceptions. + import sys + + print( + "Agent operation refused or failed; check paths, permissions, settings and Supervisor state. No services were started.", + file=sys.stderr, + ) + raise SystemExit(1) from None diff --git a/install/templates/prometheus-agent.yml b/install/templates/prometheus-agent.yml new file mode 100644 index 0000000..e6986d4 --- /dev/null +++ b/install/templates/prometheus-agent.yml @@ -0,0 +1,47 @@ +global: + scrape_interval: 15s + scrape_timeout: 10s + +scrape_configs: + - job_name: "summit/NODE_PLACEHOLDER" + # Match the existing /prom-summit proxy's upstream path. + metrics_path: / + static_configs: + - targets: ["127.0.0.1:9090"] + labels: + node: "NODE_PLACEHOLDER" + instance: "NODE_PLACEHOLDER" + role: "ROLE_PLACEHOLDER" + - job_name: "reth/NODE_PLACEHOLDER" + # Match the existing /prom-reth proxy's upstream path. + metrics_path: / + static_configs: + - targets: ["127.0.0.1:9001"] + labels: + node: "NODE_PLACEHOLDER" + instance: "NODE_PLACEHOLDER" + role: "ROLE_PLACEHOLDER" + - job_name: "prometheus-agent/NODE_PLACEHOLDER" + static_configs: + - targets: ["127.0.0.1:9091"] + labels: + node: "NODE_PLACEHOLDER" + instance: "NODE_PLACEHOLDER" + role: "ROLE_PLACEHOLDER" + +remote_write: + - name: seismic-monitoring + url: URL_PLACEHOLDER + authorization: + type: Bearer + credentials_file: TOKEN_PLACEHOLDER + remote_timeout: 30s + queue_config: + min_shards: 1 + max_shards: 2 + capacity: 10000 + max_samples_per_send: 2000 + batch_send_deadline: 5s + min_backoff: 1s + max_backoff: 30s + retry_on_http_429: true diff --git a/install/templates/supervisor/prometheus-agent.conf b/install/templates/supervisor/prometheus-agent.conf new file mode 100644 index 0000000..00574c3 --- /dev/null +++ b/install/templates/supervisor/prometheus-agent.conf @@ -0,0 +1,26 @@ +[program:prometheus-agent] +command=BINARY_PLACEHOLDER + --agent + --config.file=/etc/seismic/prometheus-agent/prometheus.yml + --storage.agent.path=DATA_PLACEHOLDER + --storage.agent.retention.min-time=5m + --storage.agent.retention.max-time=6h + --storage.agent.wal-compression + --storage.remote.flush-deadline=30s + --web.listen-address=127.0.0.1:9091 +user=seismic-prometheus +umask=0027 +autostart=false +autorestart=true +startsecs=3 +startretries=5 +priority=400 +stopwaitsecs=40 +stopasgroup=true +killasgroup=true +stdout_logfile=/var/log/seismic-prometheus-agent/agent.log +stderr_logfile=/var/log/seismic-prometheus-agent/agent.err +stdout_logfile_maxbytes=20MB +stderr_logfile_maxbytes=20MB +stdout_logfile_backups=3 +stderr_logfile_backups=3 diff --git a/internal_testnet_genesis.toml b/internal_testnet_genesis.toml index f388a9e..5ec7895 100644 --- a/internal_testnet_genesis.toml +++ b/internal_testnet_genesis.toml @@ -9,10 +9,13 @@ namespace = "_SUMMIT" validator_minimum_stake = 32000000000 blocks_per_epoch = 5000 allowed_timestamp_future_ms = 10000 +treasury_address = "0x0000000000000000000000000000000000000000" max_deposits_per_epoch = 3 max_withdrawals_per_epoch = 16 observers_per_validator = 5 +max_validator_count = 128 minimum_validator_count = 3 +invalid_deposit_tax = 0 max_pending_withdrawals_per_validator = 5 [[validators]] diff --git a/tests/test_node_onboarding.py b/tests/test_node_onboarding.py index 6ac1773..b817545 100644 --- a/tests/test_node_onboarding.py +++ b/tests/test_node_onboarding.py @@ -2,6 +2,7 @@ from __future__ import annotations +import argparse import contextlib import hashlib import importlib.util @@ -614,7 +615,7 @@ def test_pre_joining_leave_stopped_does_not_authorize_start(self) -> None: ) def test_checkpoint_start_prepares_supervisor_before_programs(self) -> None: - args = SimpleNamespace(startup_timeout=30.0) + args = SimpleNamespace(startup_timeout=30.0, inventory=None, mode="checkpoint") events: list[str] = [] with ( mock.patch.object( @@ -623,6 +624,10 @@ def test_checkpoint_start_prepares_supervisor_before_programs(self) -> None: return_value=validator.StartDecision(start=True), ), mock.patch.object(checkpoint, "validate_checkpoint_start_configuration"), + mock.patch.object(checkpoint, "load_inventory"), + mock.patch.object( + validator, "installed_node_public_key", return_value="11" * 32 + ), mock.patch.object( supervisor, "prepare_supervisor", @@ -634,7 +639,7 @@ def test_checkpoint_start_prepares_supervisor_before_programs(self) -> None: side_effect=lambda *args, **kwargs: events.append("start"), ) as start_node, ): - validator.start_checkpoint_validator( + validator.start_onboarded_validator( args, "11" * 32, allow_pre_joining_start=False, @@ -669,6 +674,442 @@ def test_cli_accepts_validator_stop_without_onboarding_options(self) -> None: self.assertIsNone(args.inventory) +class InstalledValidatorIdentityTests(unittest.TestCase): + def setUp(self) -> None: + temporary = tempfile.TemporaryDirectory() + self.addCleanup(temporary.cleanup) + self.keys_dir = Path(temporary.name) + for name in ("node_key.pem", "consensus_key.pem"): + (self.keys_dir / name).write_text("fixture-only-not-a-real-key") + self.inventory = {"summit_keys_dir": self.keys_dir} + + def test_read_only_summit_command_returns_normalized_public_key(self) -> None: + for prefix in ("", "0x"): + with self.subTest(prefix=prefix): + output = ( + f"Node Public Key (ed25519): {prefix}{'AB' * 32}\n" + f"Consensus Public Key (BLS): {'CD' * 48}\n" + ) + with mock.patch.object( + validator.subprocess, + "run", + return_value=subprocess.CompletedProcess([], 0, stdout=output), + ) as run: + public_key = validator.installed_node_public_key(self.inventory) + self.assertEqual(public_key, "ab" * 32) + run.assert_called_once_with( + [ + str(validator.SUMMIT), + "keys", + "show", + "--key-store-path", + str(self.keys_dir), + ], + stdin=subprocess.DEVNULL, + stdout=subprocess.PIPE, + stderr=subprocess.DEVNULL, + text=True, + check=True, + timeout=30.0, + ) + + def test_invalid_or_ambiguous_public_output_is_rejected(self) -> None: + valid = f"Node Public Key (ed25519): {'ab' * 32}\n" + for output in ( + "", + "Consensus Public Key (BLS): " + "ab" * 48, + "Node Public Key (ed25519): " + "ab" * 31, + "Node Public Key (ed25519): " + "zz" * 32, + valid + valid, + valid.rstrip() + " extra", + "secret-fixture-not-a-public-key", + ): + with ( + self.subTest(output=output), + mock.patch.object( + validator.subprocess, + "run", + return_value=subprocess.CompletedProcess([], 0, stdout=output), + ), + self.assertRaisesRegex( + checkpoint.CheckpointError, "exactly one valid" + ) as raised, + ): + validator.installed_node_public_key(self.inventory) + self.assertNotIn("secret-fixture", str(raised.exception)) + + def test_command_failures_do_not_expose_output(self) -> None: + secret = "secret-fixture-never-print" + for error in ( + FileNotFoundError(secret), + subprocess.CalledProcessError(1, ["summit"], output=secret, stderr=secret), + subprocess.TimeoutExpired(["summit"], 30, output=secret, stderr=secret), + UnicodeError(secret), + ): + with ( + self.subTest(error=type(error).__name__), + mock.patch.object(validator.subprocess, "run", side_effect=error), + self.assertRaisesRegex( + checkpoint.CheckpointError, "Could not read" + ) as raised, + ): + validator.installed_node_public_key(self.inventory) + self.assertNotIn(secret, str(raised.exception)) + + def test_symlinked_key_directory_is_rejected(self) -> None: + link = self.keys_dir / "symlink" + link.symlink_to(self.keys_dir, target_is_directory=True) + with ( + mock.patch.object(validator.subprocess, "run") as run, + self.assertRaises(checkpoint.CheckpointError), + ): + validator.installed_node_public_key({"summit_keys_dir": link}) + run.assert_not_called() + + def test_missing_empty_and_symlinked_key_files_are_rejected(self) -> None: + for name in ("node_key.pem", "consensus_key.pem"): + path = self.keys_dir / name + for kind in ("missing", "empty", "symlink"): + with self.subTest(name=name, kind=kind): + path.unlink() + if kind == "empty": + path.touch() + elif kind == "symlink": + path.symlink_to(self.keys_dir / "missing") + with ( + mock.patch.object(validator.subprocess, "run") as run, + self.assertRaises(checkpoint.CheckpointError), + ): + validator.installed_node_public_key(self.inventory) + run.assert_not_called() + if path.is_symlink() or path.exists(): + path.unlink() + path.write_text("fixture-only-not-a-real-key") + + +class ValidatorOnboardTests(unittest.TestCase): + def setUp(self) -> None: + self.stack = contextlib.ExitStack() + self.addCleanup(self.stack.close) + self.output = io.StringIO() + self.stack.enter_context(contextlib.redirect_stdout(self.output)) + self.inventory = self.patch(checkpoint, "load_inventory") + self.identity = self.patch( + validator, "installed_node_public_key", return_value="11" * 32 + ) + self.deposit = self.patch( + validator, "load_deposit_response", return_value=({}, "11" * 32) + ) + self.validate_checkpoint = self.patch( + checkpoint, "validate_checkpoint_start_configuration" + ) + self.install = self.patch(node_cli, "install_from_resolved_inputs") + self.prepare = self.patch(supervisor, "prepare_supervisor") + self.start = self.patch(supervisor, "start_node") + self.account = self.patch( + validator, "validator_account", return_value=self.status("Joining") + ) + self.sleep = self.patch(validator.time, "sleep") + self.patch(rpc, "read_bearer_token", return_value=None) + + def patch(self, target: object, name: str, **kwargs: object) -> mock.Mock: + return self.stack.enter_context(mock.patch.object(target, name, **kwargs)) + + @staticmethod + def status(name: str) -> dict[str, object]: + return {"status": name, "balance": 32_000_000_000, "joining_epoch": 14} + + def args(self, *options: str) -> argparse.Namespace: + with mock.patch.object( + sys, + "argv", + [ + "seismic-node.py", + "validator", + "onboard", + "--summit-rpc-url", + "https://network.example/summit", + "--pre-joining-policy", + "wait", + *options, + ], + ): + return node_cli.parse_args() + + def assert_no_start(self) -> None: + self.prepare.assert_not_called() + self.start.assert_not_called() + + def test_normal_waits_for_joining_then_rechecks_before_start(self) -> None: + self.account.side_effect = [ + None, + self.status("Inactive"), + self.status("Joining"), + self.status("Joining"), + ] + events: list[str] = [] + + def check_waiting(_: float) -> None: + self.assert_no_start() + self.install.assert_not_called() + events.append("wait") + + self.sleep.side_effect = check_waiting + self.prepare.side_effect = lambda: events.append("prepare") + self.start.side_effect = lambda *a, **kw: events.append("start") + node_cli.handle_validator( + self.args("--mode", "normal", "--inventory", "/etc/seismic/custom.toml") + ) + self.assertEqual(events, ["wait", "wait", "prepare", "start"]) + self.assertEqual(self.account.call_count, 4) + self.inventory.assert_has_calls( + [mock.call("validator", Path("/etc/seismic/custom.toml"))] * 2 + ) + self.start.assert_called_once_with( + "summit", "summit-checkpoint", startup_timeout=30.0 + ) + self.validate_checkpoint.assert_not_called() + self.install.assert_not_called() + + def test_both_modes_use_installed_identity_without_deposit_file(self) -> None: + for mode in ("normal", "checkpoint"): + with self.subTest(mode=mode): + args = self.args("--mode", mode) + self.assertIsNone(args.deposit_signature) + self.account.reset_mock() + self.identity.reset_mock() + node_cli.handle_validator(args) + self.assertEqual(self.account.call_count, 2) + for call in self.account.call_args_list: + self.assertEqual(call.args[1], "11" * 32) + self.assertEqual(self.identity.call_count, 2) + self.identity.assert_called_with(self.inventory.return_value) + self.deposit.assert_not_called() + + def test_optional_deposit_response_is_checked_against_installed_identity( + self, + ) -> None: + path = Path("/root/deposit-signature.json") + node_cli.handle_validator( + self.args("--mode", "normal", "--deposit-signature", str(path)) + ) + self.deposit.assert_called_once_with(path) + self.start.assert_called_once() + + def test_invalid_optional_deposit_response_is_not_ignored(self) -> None: + self.deposit.side_effect = checkpoint.CheckpointError( + "Invalid deposit response" + ) + with self.assertRaisesRegex( + checkpoint.CheckpointError, "Invalid deposit response" + ): + node_cli.handle_validator( + self.args("--deposit-signature", "/root/deposit-signature.json") + ) + self.account.assert_not_called() + self.install.assert_not_called() + self.assert_no_start() + + def test_mismatched_deposit_identity_fails_before_polling_or_install(self) -> None: + self.deposit.return_value = ({}, "22" * 32) + with self.assertRaisesRegex(checkpoint.CheckpointError, "does not match"): + node_cli.handle_validator( + self.args("--deposit-signature", "/root/deposit-signature.json") + ) + self.account.assert_not_called() + self.install.assert_not_called() + self.assert_no_start() + + def test_identity_read_failure_fails_before_polling_or_install(self) -> None: + self.identity.side_effect = checkpoint.CheckpointError("Cannot read identity") + with self.assertRaisesRegex(checkpoint.CheckpointError, "Cannot read identity"): + node_cli.handle_validator(self.args("--mode", "normal")) + self.account.assert_not_called() + self.install.assert_not_called() + self.assert_no_start() + + def test_key_change_while_waiting_refuses_start_in_both_modes(self) -> None: + for mode in ("normal", "checkpoint"): + with self.subTest(mode=mode): + self.identity.side_effect = ["11" * 32, "22" * 32] + with self.assertRaisesRegex( + checkpoint.CheckpointError, "changed during" + ): + node_cli.handle_validator(self.args("--mode", mode)) + self.assert_no_start() + + def test_checkpoint_remains_default_and_uses_existing_inputs(self) -> None: + args = self.args() + self.assertEqual(args.mode, "checkpoint") + node_cli.handle_validator(args) + self.validate_checkpoint.assert_has_calls([mock.call("validator")] * 2) + self.install.assert_not_called() + self.start.assert_called_once_with( + "summit-checkpoint", "summit", startup_timeout=30.0 + ) + + def test_checkpoint_download_follows_authorization(self) -> None: + self.account.side_effect = [ + None, + self.status("Joining"), + self.status("Joining"), + ] + self.sleep.side_effect = lambda _: self.install.assert_not_called() + self.install.side_effect = lambda *a, **kw: self.assert_no_start() + node_cli.handle_validator( + self.args("--snapshot-api-url", "https://snapshot.example") + ) + self.install.assert_called_once() + self.start.assert_called_once_with( + "summit-checkpoint", "summit", startup_timeout=30.0 + ) + + def test_normal_rejects_checkpoint_sources_and_install_modifiers(self) -> None: + for options in ( + ["--archive", "/tmp/archive"], + ["--manifest", "/tmp/manifest"], + ["--snapshot-api-url", "https://snapshot.example"], + ["--snapshot-bearer-token-file", "/root/token"], + ["--checkpoint-epoch", "14"], + ["--checkpoint-policy", "exact"], + ["--weak-subjectivity-path", "/tmp/anchor"], + ["--weak-subjectivity-rpc-url", "https://anchor.example"], + ["--checkpoint-path", "/persistence/checkpoint"], + ["--backup-root", "/persistence/backup"], + ["--installed-weak-subjectivity-path", "/etc/seismic/custom-anchor"], + ["--allow-same-origin-weak-subjectivity"], + ["--yes"], + ): + with ( + self.subTest(options=options), + self.assertRaisesRegex(checkpoint.CheckpointError, "--mode normal"), + ): + node_cli.handle_validator(self.args("--mode", "normal", *options)) + self.account.assert_not_called() + self.inventory.assert_not_called() + self.install.assert_not_called() + self.assert_no_start() + + def test_normal_active_starts_with_warning(self) -> None: + self.account.return_value = self.status("Active") + stderr = io.StringIO() + with redirect_stderr(stderr): + node_cli.handle_validator(self.args("--mode", "normal")) + self.assertIn("already Active", stderr.getvalue()) + self.start.assert_called_once() + + def test_normal_refuses_unsafe_status_on_either_check(self) -> None: + for status in ("SubmittedExitRequest", "FullPayoutPending", "Unknown"): + for initial_joining in (False, True): + with self.subTest(status=status, initial_joining=initial_joining): + self.account.side_effect = ( + [self.status("Joining"), self.status(status)] + if initial_joining + else [self.status(status)] + ) + with self.assertRaisesRegex(checkpoint.CheckpointError, "Refusing"): + node_cli.handle_validator(self.args("--mode", "normal")) + self.assert_no_start() + self.install.assert_not_called() + + def test_normal_leave_stopped_on_either_check(self) -> None: + for initial_joining in (False, True): + with self.subTest(initial_joining=initial_joining): + self.account.side_effect = ( + [self.status("Joining"), None] if initial_joining else [None] + ) + node_cli.handle_validator( + self.args( + "--mode", "normal", "--pre-joining-policy", "leave-stopped" + ) + ) + self.assert_no_start() + self.assertIn("No checkpoint was installed", self.output.getvalue()) + self.assertNotIn("checkpoint remains installed", self.output.getvalue().lower()) + self.install.assert_not_called() + + def test_normal_preserves_interactive_early_start_authorization(self) -> None: + args = self.args("--mode", "normal") + args.pre_joining_policy = None + self.account.return_value = None + with mock.patch.object( + validator, "choose_pre_joining_action", return_value="start" + ) as choose: + node_cli.handle_validator(args) + choose.assert_called_once_with("NotFound", None) + self.assertEqual(self.account.call_count, 2) + self.start.assert_called_once() + + def test_normal_shared_wait_deadline_applies_to_second_check(self) -> None: + self.account.side_effect = [self.status("Joining"), None] + with ( + mock.patch.object(validator.time, "monotonic", side_effect=[10.0, 12.0]), + self.assertRaisesRegex(checkpoint.CheckpointError, "Timed out"), + ): + node_cli.handle_validator( + self.args("--mode", "normal", "--validator-wait-timeout", "1") + ) + self.assert_no_start() + self.install.assert_not_called() + + def test_normal_waits_again_if_status_regresses_before_startup(self) -> None: + self.account.side_effect = [ + self.status("Joining"), + self.status("Inactive"), + self.status("Joining"), + ] + self.sleep.side_effect = lambda _: self.assert_no_start() + node_cli.handle_validator(self.args("--mode", "normal")) + self.sleep.assert_called_once() + self.assertEqual(self.account.call_count, 3) + self.start.assert_called_once() + + def test_normal_rpc_error_never_starts_services(self) -> None: + self.account.side_effect = checkpoint.CheckpointError("Invalid RPC response") + with self.assertRaisesRegex(checkpoint.CheckpointError, "Invalid RPC response"): + node_cli.handle_validator(self.args("--mode", "normal")) + self.assert_no_start() + self.install.assert_not_called() + + def test_normal_invalid_inventory_fails_before_polling(self) -> None: + self.inventory.side_effect = checkpoint.CheckpointError("Invalid inventory") + with self.assertRaisesRegex(checkpoint.CheckpointError, "Invalid inventory"): + node_cli.handle_validator(self.args("--mode", "normal")) + self.account.assert_not_called() + self.assert_no_start() + + def test_normal_requires_trusted_rpc(self) -> None: + args = self.args("--mode", "normal") + args.summit_rpc_url = None + with self.assertRaisesRegex(checkpoint.CheckpointError, "--summit-rpc-url"): + node_cli.handle_validator(args) + self.account.assert_not_called() + self.assert_no_start() + + def test_checkpoint_startup_failure_still_reports_rollback(self) -> None: + backup = Path("/persistence/rollback/fixture") + self.install.return_value = backup + self.start.side_effect = supervisor.SupervisorError("Startup failed") + with ( + mock.patch.object(node_cli, "print_startup_rollback") as rollback, + self.assertRaisesRegex(supervisor.SupervisorError, "Startup failed"), + ): + node_cli.handle_validator( + self.args("--snapshot-api-url", "https://snapshot.example") + ) + rollback.assert_called_once_with(backup) + + def test_normal_startup_failure_has_no_checkpoint_rollback(self) -> None: + self.start.side_effect = supervisor.SupervisorError("Startup failed") + with ( + mock.patch.object(node_cli, "print_startup_rollback") as rollback, + self.assertRaisesRegex(supervisor.SupervisorError, "Startup failed"), + ): + node_cli.handle_validator(self.args("--mode", "normal")) + rollback.assert_not_called() + self.install.assert_not_called() + + class ValidatorStartTests(unittest.TestCase): def start_args(self, mode: str) -> SimpleNamespace: return SimpleNamespace(inventory=None, mode=mode, startup_timeout=30.0) diff --git a/tests/test_prometheus_agent.py b/tests/test_prometheus_agent.py new file mode 100644 index 0000000..8ad3517 --- /dev/null +++ b/tests/test_prometheus_agent.py @@ -0,0 +1,707 @@ +"""Optional agent tests: local fixtures only, never run the root installer.""" + +from __future__ import annotations + +import argparse +import configparser +import contextlib +import hashlib +import importlib.util +import io +import json +import os +import shutil +import subprocess +import sys +import tarfile +import tempfile +import unittest +from pathlib import Path +from types import SimpleNamespace +from unittest import mock + +ROOT = Path(__file__).resolve().parents[1] +sys.path.insert(0, str(ROOT / "tools")) +from seismic_node import checkpoint, monitoring, observer, supervisor, validator + + +def load(name, path): + spec = importlib.util.spec_from_file_location(name, path) + module = importlib.util.module_from_spec(spec) + spec.loader.exec_module(module) + return module + + +agent = load("agent_installer", ROOT / "install/lib/prometheus_agent.py") +cli = load("agent_node_cli", ROOT / "tools/seismic-node.py") +PROMTOOL = os.environ.get("PROMETHEUS_AGENT_PROMTOOL") or shutil.which("promtool") + + +def settings(**changes): + return { + "enabled": True, + "node": "validator-0.example.com", + "role": "validator", + "url": "https://metrics.example.com/api/v1/write", + "data_dir": "/var/lib/seismic-prometheus-agent", + "version": "3.5.0", + **changes, + } + + +class SettingsTests(unittest.TestCase): + def test_architectures_and_checksums(self): + self.assertEqual(agent.architecture("x86_64"), "amd64") + self.assertEqual(agent.architecture("aarch64"), "arm64") + with self.assertRaises(agent.AgentError): + agent.architecture("i686") + self.assertEqual(len(agent.CHECKSUMS), 2) + for checksum in agent.CHECKSUMS.values(): + self.assertRegex(checksum, r"^[a-f0-9]{64}$") + + def test_reject_unsafe_endpoint_and_label(self): + for url in ( + "http://metrics.example.com/api/v1/write", + "https://user:secret@metrics.example.com/api/v1/write", + "https://metrics.example.com/api/v1/write?x=1", + "https://metrics.example.com/api/v1/write?", + "https://metrics.example.com/api/v1/write#", + "https://metrics.example.com/", + "https://metrics.example.com:9090/api/v1/write", + "https://metrics.example.com/\napi/v1/write", + ): + with self.subTest(url=url), self.assertRaises(agent.AgentError): + agent.validate(settings(url=url)) + for node in ("../key", "bad\nlabel", 'bad"', "UPPER", "a..b", ""): + with self.subTest(node=node), self.assertRaises(agent.AgentError): + agent.validate(settings(node=node)) + + def test_wal_must_be_separate_and_normalized(self): + for path in ( + "/", + "/etc/agent", + "/usr/data", + "/var/lib/../data", + "/data/", + "/data/a b", + "/data/node", + "/data/node/agent", + ): + with self.subTest(path=path), self.assertRaises(agent.AgentError): + agent.validate(settings(data_dir=path), [Path("/data/node")]) + agent.validate(settings(data_dir="/data/agent"), [Path("/data/node")]) + + def test_token_source_in_node_keys_is_rejected_without_reading_it(self): + arguments = [ + "agent", + "check", + "--node", + "v0", + "--role", + "validator", + "--url", + "https://metrics.example.com/api/v1/write", + "--data-dir", + "/var/lib/seismic-prometheus-agent", + "--token-source", + "/keys/summit/node_key.pem", + "--exclude-path", + "/keys/summit", + ] + with ( + mock.patch.object(sys, "argv", arguments), + mock.patch.object(agent.os, "geteuid", return_value=0), + mock.patch.object(agent, "require_managed_or_new"), + mock.patch.object(agent, "safe_path"), + mock.patch.object(agent, "read_token") as read, + ): + with self.assertRaisesRegex(agent.AgentError, "outside node data"): + agent.main() + read.assert_not_called() + + def test_root_file_rejects_untrusted_parents(self): + with tempfile.TemporaryDirectory() as tmp: + path = Path(tmp) / "token" + path.write_text("a" * 64) + with self.assertRaises(agent.AgentError): + agent.root_file(path) + + def test_symlink_and_writable_parent_are_refused(self): + with tempfile.TemporaryDirectory() as tmp: + directory = Path(tmp) + (directory / "real").mkdir() + (directory / "link").symlink_to( + directory / "real", target_is_directory=True + ) + with self.assertRaises(agent.AgentError): + agent.safe_path(directory / "link" / "token") + + def test_existing_account_must_have_a_private_group(self): + account = SimpleNamespace(pw_uid=998, pw_gid=998, pw_name=agent.USER) + with ( + mock.patch.object(agent.pwd, "getpwnam", return_value=account), + mock.patch.object(agent.pwd, "getpwall", return_value=[account]), + mock.patch.object( + agent.grp, + "getgrgid", + return_value=SimpleNamespace(gr_name=agent.USER, gr_mem=[]), + ) as group, + ): + self.assertEqual(agent.service_account(), account) + group.return_value.gr_mem = ["unrelated-user"] + with self.assertRaises(agent.AgentError): + agent.service_account() + group.return_value.gr_mem = [] + group.return_value.gr_name = "shared-users" + with self.assertRaises(agent.AgentError): + agent.service_account() + + def test_private_token_and_bad_token_redaction(self): + with tempfile.TemporaryDirectory() as tmp: + path = Path(tmp) / "token" + path.write_text("a" * 64 + "\n") + path.chmod(0o600) + # Only bypass host-root ownership, retaining real mode and bytes. + with mock.patch.object(agent, "root_file", side_effect=lambda p: p.stat()): + self.assertEqual(agent.read_token(path), b"a" * 64 + b"\n") + path.chmod(0o644) + with self.assertRaises(agent.AgentError): + agent.read_token(path) + path.chmod(0o600) + path.write_text("PRIVATE-CREDENTIAL") + with self.assertRaises(agent.AgentError) as caught: + agent.read_token(path) + self.assertNotIn("PRIVATE-CREDENTIAL", str(caught.exception)) + + def test_render_both_roles(self): + for role in ("validator", "observer"): + config, service = agent.render(settings(role=role)) + self.assertNotIn("PLACEHOLDER", config + service) + self.assertIn(f'role: "{role}"', config) + for stream in ("summit", "reth", "prometheus-agent"): + self.assertIn(f'job_name: "{stream}/validator-0.example.com"', config) + for port in (9090, 9001, 9091): + self.assertIn(f"127.0.0.1:{port}", config) + self.assertIn("credentials_file:", config) + self.assertIn("retry_on_http_429: true", config) + self.assertNotIn("insecure_skip_verify", config) + parser = configparser.ConfigParser(interpolation=None) + parser.read_string(service) + program = parser["program:prometheus-agent"] + self.assertFalse(program.getboolean("autostart")) + self.assertTrue(program.getboolean("autorestart")) + self.assertEqual(program["user"], "seismic-prometheus") + self.assertIn("--web.listen-address=127.0.0.1:9091", program["command"]) + self.assertIn("--storage.agent.retention.max-time=6h", program["command"]) + + def test_supervisor_parses_generated_service(self): + try: + from supervisor.options import ServerOptions + except ModuleNotFoundError as error: + if error.name != "supervisor": + raise + self.skipTest("supervisor package unavailable") + with tempfile.TemporaryDirectory() as tmp: + root = Path(tmp) + service = agent.render(settings())[1] + service = service.replace(str(agent.LOG_DIR), str(root)) + service = service.replace("user=seismic-prometheus", f"user={os.getuid()}") + (root / "agent.conf").write_text(service) + wrapper = root / "supervisord.conf" + wrapper.write_text( + f"[supervisord]\nlogfile={root}/supervisor.log\npidfile={root}/pid\n" + f"[include]\nfiles={root}/agent.conf\n" + ) + options = ServerOptions() + options.realize(["-c", str(wrapper)]) + self.assertEqual(len(options.process_group_configs), 1) + group = options.process_group_configs[0] + self.assertEqual(group.name, "prometheus-agent") + self.assertFalse(group.process_configs[0].autostart) + self.assertIn("--agent", group.process_configs[0].command) + + @unittest.skipUnless(PROMTOOL, "promtool unavailable") + def test_config_validated_in_agent_mode(self): + with tempfile.TemporaryDirectory() as tmp: + token = Path(tmp) / "token" + token.write_text("a" * 64) + path = Path(tmp) / "prometheus.yml" + for role in ("validator", "observer"): + path.write_text(agent.render(settings(role=role), token)[0]) + result = subprocess.run( + [PROMTOOL, "check", "config", "--agent", str(path)], + check=False, + capture_output=True, + text=True, + ) + self.assertEqual(result.returncode, 0, result.stdout + result.stderr) + + def test_checksum_failure_writes_no_executables(self): + with ( + tempfile.TemporaryDirectory() as tmp, + mock.patch.object(agent.platform, "machine", return_value="x86_64"), + mock.patch.object( + agent.urllib.request, + "urlopen", + return_value=io.BytesIO(b"wrong-release"), + ), + mock.patch.object(agent, "directory") as mkdir, + mock.patch.object(agent, "atomic_write") as write, + ): + with self.assertRaisesRegex(agent.AgentError, "SHA-256"): + agent.install_release(Path(tmp)) + mkdir.assert_not_called() + write.assert_not_called() + + def test_only_expected_regular_executables_are_installed(self): + data = io.BytesIO() + with tarfile.open(fileobj=data, mode="w:gz") as tar: + for name in ("prometheus", "promtool", "../../escape"): + member = tarfile.TarInfo(f"prometheus-3.5.0.linux-amd64/{name}") + member.size = 3 + tar.addfile(member, io.BytesIO(b"bin")) + archive = data.getvalue() + with ( + tempfile.TemporaryDirectory() as tmp, + mock.patch.object(agent.platform, "machine", return_value="x86_64"), + mock.patch.dict( + agent.CHECKSUMS, {"amd64": hashlib.sha256(archive).hexdigest()} + ), + mock.patch.object( + agent.urllib.request, "urlopen", return_value=io.BytesIO(archive) + ), + mock.patch.object(agent, "directory"), + mock.patch.object(agent, "atomic_write") as write, + ): + agent.install_release(Path(tmp)) + self.assertEqual( + [c.args[0].name for c in write.call_args_list], + ["prometheus", "promtool"], + ) + self.assertFalse((Path(tmp) / "escape").exists()) + + +class InstallationTests(unittest.TestCase): + def test_first_disabled_install_never_controls_services(self): + with ( + mock.patch.object(agent, "load_settings", return_value=None), + mock.patch.object(agent, "SUPERVISOR_CONFIG") as conf, + mock.patch.object(agent.subprocess, "run") as run, + ): + conf.exists.return_value = False + agent.disable() + run.assert_not_called() + + def test_disabled_mode_ignores_but_never_adopts_unmanaged_agent(self): + with tempfile.TemporaryDirectory() as tmp: + root = Path(tmp) + conf = root / "agent.conf" + conf.write_text("unrelated program") + with ( + mock.patch.object(agent, "SETTINGS", root / "absent.json"), + mock.patch.object(agent, "SUPERVISOR_CONFIG", conf), + mock.patch.object(agent.subprocess, "run") as run, + ): + self.assertIsNone(agent.load_settings()) + with self.assertRaisesRegex(agent.AgentError, "unmanaged"): + agent.install(settings(), root / "token") + run.assert_not_called() + self.assertEqual(conf.read_text(), "unrelated program") + + def test_update_refuses_running_agent(self): + with ( + mock.patch.object(agent, "read_token", return_value=b"a" * 64), + mock.patch.object(agent, "load_settings", return_value=settings()), + mock.patch.object(agent, "SUPERVISOR_CONFIG") as conf, + mock.patch.object(agent, "supervisor_state", return_value="RUNNING"), + mock.patch.object(agent, "install_release") as release, + ): + conf.exists.return_value = True + with self.assertRaisesRegex(agent.AgentError, "Stop prometheus-agent"): + agent.install(settings(), Path("/root/token")) + release.assert_not_called() + + def test_disable_stops_only_agent_and_preserves_files(self): + with tempfile.TemporaryDirectory() as tmp: + root = Path(tmp) + conf = root / "supervisor.conf" + conf.write_text("managed") + token = root / "token" + token.write_text("secret") + wal = root / "wal" + wal.write_text("buffered") + with ( + mock.patch.object(agent, "load_settings", return_value=settings()), + mock.patch.object( + agent, + "supervisor_state", + side_effect=["EXITED", "RUNNING", "STOPPED", "STOPPED"], + ), + mock.patch.object(agent, "SUPERVISOR_CONFIG", conf), + mock.patch.object(agent, "root_file"), + mock.patch.object(agent, "atomic_write") as write, + mock.patch.object( + agent.subprocess, "run", return_value=SimpleNamespace(returncode=0) + ) as run, + ): + agent.disable() + self.assertEqual( + run.call_args.args[0], + ["/usr/bin/supervisorctl", "stop", "prometheus-agent"], + ) + self.assertFalse(json.loads(write.call_args.args[1])["enabled"]) + self.assertFalse(conf.exists()) + self.assertEqual(token.read_text(), "secret") + self.assertEqual(wal.read_text(), "buffered") + + def test_failed_disable_preserves_enabled_config(self): + with ( + mock.patch.object(agent, "load_settings", return_value=settings()), + mock.patch.object(agent, "supervisor_state", return_value="RUNNING"), + mock.patch.object( + agent.subprocess, "run", return_value=SimpleNamespace(returncode=1) + ), + mock.patch.object(agent, "atomic_write") as write, + ): + with self.assertRaises(agent.AgentError): + agent.disable() + write.assert_not_called() + + @contextlib.contextmanager + def fixture(self): + with tempfile.TemporaryDirectory() as tmp, contextlib.ExitStack() as stack: + root = Path(tmp) + paths = { + "HOME": root / "config", + "SETTINGS": root / "config/installation.json", + "TOKEN": root / "config/token", + "CONFIG": root / "config/prometheus.yml", + "SUPERVISOR_CONFIG": root / "supervisor.conf", + "LOG_DIR": root / "logs", + "BIN_DIR": root / "binaries/3.5.0", + } + for key, value in paths.items(): + stack.enter_context(mock.patch.object(agent, key, value)) + # Fixture-only bypass of host-root ownership; all writes remain in tmp. + stack.enter_context(mock.patch.object(agent, "safe_path")) + stack.enter_context( + mock.patch.object(agent, "root_file", side_effect=lambda p: p.stat()) + ) + stack.enter_context(mock.patch.object(agent.os, "chown")) + stack.enter_context(mock.patch.object(agent.os, "fchown")) + stack.enter_context( + mock.patch.object( + agent, + "service_account", + return_value=SimpleNamespace( + pw_uid=os.getuid() or 1000, pw_gid=os.getgid() or 1000 + ), + ) + ) + stack.enter_context(mock.patch.object(agent, "install_release")) + stack.enter_context( + mock.patch.object(agent, "supervisor_state", return_value="STOPPED") + ) + run = stack.enter_context( + mock.patch.object( + agent.subprocess, "run", return_value=SimpleNamespace(returncode=0) + ) + ) + token = root / "handoff" + token.write_text("a" * 64) + token.chmod(0o600) + yield root, token, run + + def test_install_and_reinstall_preserve_token_wal_and_do_not_activate(self): + with self.fixture() as (root, token, run): + config = settings(data_dir=str(root / "wal")) + agent.install(config, token) + (root / "wal/buffered").write_text("samples") + self.assertEqual(agent.TOKEN.stat().st_mode & 0o777, 0o640) + self.assertEqual(agent.SETTINGS.stat().st_mode & 0o777, 0o600) + self.assertEqual(agent.CONFIG.stat().st_mode & 0o777, 0o640) + self.assertNotIn("a" * 64, agent.CONFIG.read_text()) + self.assertNotIn("a" * 64, agent.SETTINGS.read_text()) + # Existing fixture directories are owned by the unprivileged tester; + # preserve the real directory contents while bypassing chown policy. + with mock.patch.object( + agent, + "directory", + side_effect=lambda p, *a, **kw: p.mkdir(parents=True, exist_ok=True), + ): + agent.install(config, agent.TOKEN) + self.assertEqual(agent.TOKEN.read_text(), "a" * 64 + "\n") + self.assertEqual((root / "wal/buffered").read_text(), "samples") + for call in run.call_args_list: + self.assertEqual(call.args[0][1:4], ["check", "config", "--agent"]) + + def test_validation_failure_does_not_replace_existing_configuration(self): + with self.fixture() as (root, token, run): + config = settings(data_dir=str(root / "wal")) + agent.install(config, token) + before = { + p: p.read_bytes() + for p in ( + agent.CONFIG, + agent.TOKEN, + agent.SETTINGS, + agent.SUPERVISOR_CONFIG, + ) + } + token.write_text("b" * 64) + run.return_value.returncode = 1 + with ( + mock.patch.object( + agent, + "directory", + side_effect=lambda p, *a, **kw: p.mkdir( + parents=True, exist_ok=True + ), + ), + self.assertRaises(agent.AgentError), + ): + agent.install(config, token) + self.assertEqual({p: p.read_bytes() for p in before}, before) + + def test_installer_wiring_and_keep_is_noop(self): + for role in ("validator", "observer"): + source = (ROOT / f"install/install-{role}.sh").read_text() + self.assertIn('source "$SCRIPT_DIR/lib/prometheus-agent.sh"', source) + self.assertIn("monitoring=$(write_prometheus_agent_inventory)", source) + self.assertLess( + source.index(" validate_prometheus_agent_plan"), + source.index(" install_system_packages"), + ) + result = subprocess.run( + [ + "bash", + "-c", + f'set -eu; SCRIPT_DIR="{ROOT}/install"; source "$SCRIPT_DIR/lib/prometheus-agent.sh"; python3() {{ exit 99; }}; install_prometheus_agent', + ], + check=False, + capture_output=True, + ) + self.assertEqual(result.returncode, 0) + + +class LifecycleTests(unittest.TestCase): + def test_all_node_start_modes_ensure_agent_after_success(self): + for role, module, function in ( + ("validator", validator, validator.start_validator), + ("observer", observer, observer.start_observer), + ): + for mode in ("normal", "checkpoint"): + inventory = {"monitoring": settings(role=role)} + args = argparse.Namespace( + inventory=Path("/custom/inventory.toml"), + mode=mode, + startup_timeout=3, + ) + events = [] + with ( + self.subTest(role=role, mode=mode), + mock.patch.object( + checkpoint, "load_inventory", return_value=inventory + ), + mock.patch.object( + checkpoint, "validate_checkpoint_start_configuration" + ), + mock.patch.object(supervisor, "prepare_supervisor"), + mock.patch.object( + supervisor, + "start_node", + side_effect=lambda *a, events=events, **kw: events.append( + "node" + ), + ), + mock.patch.object( + supervisor, + "start_program", + side_effect=lambda *a, events=events, **kw: events.append(a[0]), + ), + ): + function(args) + self.assertEqual(events, ["node", "prometheus-agent"]) + + def test_disabled_and_old_inventories_are_noops(self): + for inventory in ({}, {"monitoring": {"enabled": False}}): + with mock.patch.object(supervisor, "start_program") as start: + monitoring.ensure_running(inventory, "validator", 3) + start.assert_not_called() + + def test_agent_failure_warns_without_node_rollback(self): + inventory = {"monitoring": settings()} + for error in ( + supervisor.SupervisorError("private output"), + OSError("private output"), + ): + with ( + mock.patch.object(supervisor, "start_program", side_effect=error), + mock.patch.object(supervisor, "stop_program") as stop, + contextlib.redirect_stderr(io.StringIO()) as err, + ): + monitoring.ensure_running(inventory, "validator", 3) + self.assertIn("WARNING", err.getvalue()) + self.assertNotIn("private output", err.getvalue()) + stop.assert_not_called() + with contextlib.redirect_stderr(io.StringIO()) as err: + monitoring.ensure_running({"monitoring": "broken"}, "validator", 3) + self.assertIn("WARNING", err.getvalue()) + + def test_node_start_failure_does_not_start_or_stop_agent(self): + for role, start in ( + ("validator", validator.start_validator), + ("observer", observer.start_observer), + ): + args = argparse.Namespace(inventory=None, mode="normal", startup_timeout=3) + with ( + mock.patch.object( + checkpoint, + "load_inventory", + return_value={"monitoring": settings(role=role)}, + ), + mock.patch.object(supervisor, "prepare_supervisor"), + mock.patch.object( + supervisor, + "start_node", + side_effect=supervisor.SupervisorError("node failed"), + ), + mock.patch.object(monitoring, "ensure_running") as ensure, + ): + with self.assertRaises(supervisor.SupervisorError): + start(args) + ensure.assert_not_called() + + def test_refused_onboarding_does_not_start_agent(self): + args = argparse.Namespace(mode="normal") + with ( + mock.patch.object( + validator, + "wait_for_start_authorization", + return_value=validator.StartDecision(False), + ), + mock.patch.object(validator, "start_validator") as start, + mock.patch.object(monitoring, "ensure_running") as ensure, + ): + validator.start_onboarded_validator( + args, "a" * 64, allow_pre_joining_start=False + ) + start.assert_not_called() + ensure.assert_not_called() + + def test_agent_stop_does_not_treat_exited_as_quiescent(self): + states = [ + supervisor.ProgramStatus("prometheus-agent", state, "") + for state in ("EXITED", "RUNNING", "STOPPED") + ] + with ( + mock.patch.object(supervisor, "status", side_effect=states), + mock.patch.object( + supervisor, + "run_supervisorctl", + side_effect=[ + SimpleNamespace(returncode=7), + SimpleNamespace(returncode=0), + ], + ) as control, + mock.patch.object(supervisor.time, "sleep"), + ): + supervisor.stop_autorestarting_program("prometheus-agent") + self.assertEqual( + control.call_args_list, [mock.call("stop", "prometheus-agent")] * 2 + ) + + def test_node_stop_leaves_agent_alone(self): + with mock.patch.object(supervisor, "stop_program", return_value=True) as stop: + supervisor.stop_node(("summit", "summit-checkpoint")) + self.assertNotIn(mock.call("prometheus-agent"), stop.call_args_list) + + def test_explicit_commands_and_parser(self): + for action in ("start", "status", "stop"): + with mock.patch.object( + sys, + "argv", + [ + "seismic-node.py", + "monitoring", + action, + "--role", + "observer", + "--inventory", + "/custom.toml", + ], + ): + args = cli.parse_args() + with ( + mock.patch.object( + checkpoint, + "read_toml", + return_value={"monitoring": settings(role="observer")}, + ) as read, + mock.patch.object(monitoring, "prepare_agent") as prepare, + mock.patch.object(supervisor, "start_program") as start, + mock.patch.object(supervisor, "stop_autorestarting_program") as stop, + mock.patch.object( + supervisor, "status", return_value=SimpleNamespace(state="RUNNING") + ), + ): + monitoring.handle(args) + self.assertEqual(read.call_args.args[0], Path("/custom.toml")) + self.assertEqual(start.called, action == "start") + self.assertEqual(stop.called, action == "stop") + self.assertEqual(prepare.called, action == "start") + + def test_explicit_start_updates_only_agent_group(self): + result = SimpleNamespace(returncode=0, stdout="", stderr="") + with ( + mock.patch.object(supervisor, "run_systemctl", return_value=result), + mock.patch.object( + supervisor, "run_supervisorctl", return_value=result + ) as control, + ): + monitoring.prepare_agent() + self.assertEqual( + control.call_args_list, + [mock.call("reread"), mock.call("update", "prometheus-agent")], + ) + + def test_inventory_accepts_optional_monitoring_without_weakening_base(self): + data = { + "schema_version": 1, + "reth_data_dir": "/data/reth", + "reth_p2p_key_path": "/keys/reth-key", + "summit_data_dir": "/data/summit", + "summit_keys_dir": "/keys/summit", + "monitoring": settings(), + } + with ( + mock.patch.object(checkpoint, "read_toml", return_value=data), + mock.patch.object(checkpoint, "validate_installed_identity") as identity, + ): + loaded = checkpoint.load_inventory("validator", Path("/inventory.toml")) + self.assertTrue(monitoring.configured(loaded, "validator")) + identity.assert_called_once() + data["unexpected"] = True + with ( + mock.patch.object(checkpoint, "read_toml", return_value=data), + self.assertRaises(checkpoint.CheckpointError), + ): + checkpoint.load_inventory("validator", Path("/inventory.toml")) + + def test_already_running_is_not_restarted(self): + with ( + mock.patch.object( + supervisor, + "status", + return_value=supervisor.ProgramStatus( + "prometheus-agent", "RUNNING", "" + ), + ), + mock.patch.object(supervisor, "run_supervisorctl") as control, + ): + monitoring.ensure_running({"monitoring": settings()}, "validator", 3) + control.assert_not_called() + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_prometheus_agent_integration.py b/tests/test_prometheus_agent_integration.py new file mode 100644 index 0000000..a68147a --- /dev/null +++ b/tests/test_prometheus_agent_integration.py @@ -0,0 +1,257 @@ +"""Opt-in real Prometheus 3.5 agent -> HTTPS fixture -> receiver integration. + +All listeners use ephemeral loopback ports, all data is temporary, and only the +binary explicitly supplied in PROMETHEUS_AGENT_BIN is executed. No installer, +Supervisor or live validator is contacted. +""" + +from __future__ import annotations + +import configparser +import json +import os +import shlex +import shutil +import socket +import ssl +import subprocess +import tempfile +import threading +import time +import unittest +import urllib.error +import urllib.parse +import urllib.request +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from pathlib import Path + +from test_prometheus_agent import agent, settings + +BINARY = os.environ.get("PROMETHEUS_AGENT_BIN") + + +def free_port(): + with socket.socket() as sock: + sock.bind(("127.0.0.1", 0)) + return sock.getsockname()[1] + + +@unittest.skipUnless( + BINARY and shutil.which("openssl"), + "set PROMETHEUS_AGENT_BIN for isolated integration", +) +class AgentIntegrationTests(unittest.TestCase): + def test_forward_labels_authentication_and_failed_scrapes(self): + received_auth = [] + exporter_ok = threading.Event() + exporter_ok.set() + receiver_port, agent_port = free_port(), free_port() + + class Exporter(BaseHTTPRequestHandler): + def do_GET(self): + self.send_response( + 200 if exporter_ok.is_set() and self.path == "/" else 503 + ) + self.send_header("Content-Type", "text/plain; version=0.0.4") + self.end_headers() + self.wfile.write( + b"# TYPE fixture_counter_total counter\nfixture_counter_total 1\n" + ) + + def log_message(self, *args): + pass + + class Ingress(BaseHTTPRequestHandler): + def do_POST(self): + received_auth.append(self.headers.get("Authorization")) + if ( + self.path != "/api/v1/write" + or self.headers.get("Authorization") != "Bearer " + "a" * 64 + ): + self.send_error(401) + return + payload = self.rfile.read(int(self.headers["Content-Length"])) + headers = { + key: self.headers[key] + for key in ( + "Content-Type", + "Content-Encoding", + "X-Prometheus-Remote-Write-Version", + ) + if key in self.headers + } + request = urllib.request.Request( + f"http://127.0.0.1:{receiver_port}/api/v1/write", + data=payload, + headers=headers, + method="POST", + ) + try: + with urllib.request.urlopen(request, timeout=10) as response: + code = response.status + except urllib.error.HTTPError as error: + code = error.code + self.send_response(code) + self.end_headers() + + def log_message(self, *args): + pass + + def query(expression): + query_string = urllib.parse.urlencode({"query": expression}) + with urllib.request.urlopen( + f"http://127.0.0.1:{receiver_port}/api/v1/query?{query_string}", + timeout=2, + ) as response: + return json.load(response)["data"]["result"] + + def wait_for(predicate): + deadline = time.monotonic() + 45 + while time.monotonic() < deadline: + try: + if predicate(): + return + except (OSError, KeyError, urllib.error.URLError): + pass + time.sleep(0.25) + self.fail("Isolated agent/receiver did not reach the expected state") + + with tempfile.TemporaryDirectory(prefix="seismic-agent-integration-") as tmp: + root = Path(tmp) + subprocess.run( + [ + "openssl", + "req", + "-x509", + "-newkey", + "rsa:2048", + "-nodes", + "-keyout", + str(root / "key.pem"), + "-out", + str(root / "cert.pem"), + "-days", + "1", + "-subj", + "/CN=localhost", + "-addext", + "subjectAltName=IP:127.0.0.1", + ], + check=True, + capture_output=True, + ) + exporter = ThreadingHTTPServer(("127.0.0.1", 0), Exporter) + ingress = ThreadingHTTPServer(("127.0.0.1", 0), Ingress) + tls = ssl.SSLContext(ssl.PROTOCOL_TLS_SERVER) + tls.load_cert_chain(root / "cert.pem", root / "key.pem") + ingress.socket = tls.wrap_socket(ingress.socket, server_side=True) + threads = [ + threading.Thread(target=server.serve_forever, daemon=True) + for server in (exporter, ingress) + ] + for thread in threads: + thread.start() + processes = [] + handles = [] + try: + receiver_config = root / "receiver.yml" + receiver_config.write_text("scrape_configs: []\n") + commands = [ + [ + BINARY, + f"--config.file={receiver_config}", + f"--storage.tsdb.path={root / 'receiver-data'}", + f"--web.listen-address=127.0.0.1:{receiver_port}", + "--web.enable-remote-write-receiver", + ] + ] + token = root / "token" + token.write_text("a" * 64) + config, service = agent.render(settings(), token) + # Production requires HTTPS:443; fixture transport is HTTPS on a + # high loopback port with an explicitly trusted temporary CA. + config = config.replace( + "https://metrics.example.com/api/v1/write", + f"https://127.0.0.1:{ingress.server_port}/api/v1/write", + ) + for port in (9090, 9001): + config = config.replace( + f"127.0.0.1:{port}", f"127.0.0.1:{exporter.server_port}" + ) + config = config.replace("127.0.0.1:9091", f"127.0.0.1:{agent_port}") + config = config.replace( + "scrape_interval: 15s", "scrape_interval: 1s" + ).replace("scrape_timeout: 10s", "scrape_timeout: 1s") + config = config.replace( + " authorization:", + f" tls_config:\n ca_file: {root / 'cert.pem'}\n authorization:", + ) + config_path = root / "agent.yml" + config_path.write_text(config) + parser = configparser.ConfigParser(interpolation=None) + parser.read_string(service) + command = shlex.split(parser["program:prometheus-agent"]["command"]) + command[0] = BINARY + command = [ + part.replace(str(agent.CONFIG), str(config_path)) + .replace(settings()["data_dir"], str(root / "agent-data")) + .replace("127.0.0.1:9091", f"127.0.0.1:{agent_port}") + for part in command + ] + commands.append(command) + for number, command in enumerate(commands): + output = (root / f"process-{number}.log").open("wb") + handles.append(output) + processes.append( + subprocess.Popen( + command, stdout=output, stderr=subprocess.STDOUT + ) + ) + expression = 'up{node="validator-0.example.com"}' + wait_for(lambda: len(query(expression)) == 3) + result = query(expression) + self.assertEqual( + {sample["metric"]["job"] for sample in result}, + { + f"{component}/validator-0.example.com" + for component in ("summit", "reth", "prometheus-agent") + }, + ) + self.assertTrue( + all( + sample["metric"]["instance"] == "validator-0.example.com" + for sample in result + ) + ) + self.assertTrue( + all(sample["metric"]["role"] == "validator" for sample in result) + ) + self.assertTrue(received_auth) + self.assertEqual(set(received_auth), {"Bearer " + "a" * 64}) + exporter_ok.clear() + wait_for(lambda: len(query(expression + " == 0")) == 2) + self.assertEqual( + len( + query('up{job="prometheus-agent/validator-0.example.com"} == 1') + ), + 1, + ) + finally: + for process in reversed(processes): + process.terminate() + try: + process.wait(timeout=40) + except subprocess.TimeoutExpired: + process.kill() + process.wait(timeout=5) + for output in handles: + output.close() + for server in (exporter, ingress): + server.shutdown() + server.server_close() + for thread in threads: + thread.join(timeout=5) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_source_checkout.py b/tests/test_source_checkout.py new file mode 100644 index 0000000..b52153c --- /dev/null +++ b/tests/test_source_checkout.py @@ -0,0 +1,383 @@ +"""Exercise installer Git refs using local repositories, never live hosts.""" + +import os +import subprocess +import tempfile +import unittest +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[1] +TAG = "internal-testnet-v0" + + +class SourceCheckoutTests(unittest.TestCase): + def setUp(self): + temporary = tempfile.TemporaryDirectory() + self.addCleanup(temporary.cleanup) + self.root = Path(temporary.name) + self.env = {k: v for k, v in os.environ.items() if not k.startswith("GIT_")} + self.env.update( + GIT_CONFIG_NOSYSTEM="1", + GIT_CONFIG_GLOBAL=os.devnull, + GIT_TERMINAL_PROMPT="0", + GIT_AUTHOR_NAME="Fixture", + GIT_AUTHOR_EMAIL="fixture@example.invalid", + GIT_COMMITTER_NAME="Fixture", + GIT_COMMITTER_EMAIL="fixture@example.invalid", + ) + self.remote = self.root / "remote.git" + self.seed = self.root / "seed" + self.checkout = self.root / "service/src/component" + self.log = self.root / "installer.log" + self.git("init", "--bare", "--initial-branch=main", str(self.remote)) + self.git("init", "--initial-branch=main", str(self.seed)) + self.first = self.commit("initial") + self.git("branch", "other", cwd=self.seed) + self.git("remote", "add", "origin", self.remote.as_uri(), cwd=self.seed) + self.git("push", "origin", "main", "other", cwd=self.seed) + self.publish_tag(TAG) + + def git(self, *args, cwd=None, check=True): + return subprocess.run( + ["git", *args], + cwd=cwd or self.root, + env=self.env, + capture_output=True, + text=True, + check=check, + timeout=20, + ) + + def commit(self, text): + (self.seed / "payload").write_text(text + "\n") + self.git("add", "payload", cwd=self.seed) + self.git("commit", "-m", text, cwd=self.seed) + return self.git("rev-parse", "HEAD", cwd=self.seed).stdout.strip() + + def publish_tag(self, name, annotated=False): + args = ["tag"] + if annotated: + args.extend(["-a", "-m", "fixture release"]) + self.git(*args, name, cwd=self.seed) + self.git("push", "origin", f"refs/tags/{name}", cwd=self.seed) + + def prepare(self, ref=f"refs/tags/{TAG}", expected=True): + result = subprocess.run( + [ + "bash", + "-c", + """set -euo pipefail +source "$1/install/lib/binaries.sh" +SERVICE_HOME="$2/service" +LOG_FILE="$2/installer.log" +run_as_service_user() { "$@"; } +prepare_source_root() { mkdir -p "$SERVICE_HOME/src"; } +info() { printf '%s\\n' "$*"; } +success() { printf '%s\\n' "$*"; } +die() { printf '%s\\n' "$*" >&2; exit 1; } +prepare_source_checkout "Fixture" "$3" "$4" "$5" +""", + "fixture", + str(ROOT), + str(self.root), + self.remote.as_uri(), + str(self.checkout), + ref, + ], + env=self.env, + capture_output=True, + text=True, + check=False, + timeout=30, + ) + diagnostic = result.stdout + result.stderr + if self.log.exists(): + diagnostic += self.log.read_text() + if expected: + self.assertEqual(result.returncode, 0, diagnostic) + else: + self.assertNotEqual(result.returncode, 0, diagnostic) + return result + + def head(self): + return self.git("rev-parse", "HEAD", cwd=self.checkout).stdout.strip() + + def assert_detached(self, commit): + self.assertEqual(self.head(), commit) + self.assertEqual( + self.git( + "symbolic-ref", "-q", "HEAD", cwd=self.checkout, check=False + ).returncode, + 1, + ) + self.assertEqual( + self.git("status", "--porcelain", cwd=self.checkout).stdout, "" + ) + + def test_fresh_lightweight_tag_and_rerun_stay_pinned(self): + self.prepare() + self.assert_detached(self.first) + newer = self.commit("branch advanced") + self.git("push", "origin", "main", cwd=self.seed) + self.prepare() + self.assert_detached(self.first) + self.assertNotEqual(self.head(), newer) + self.git("show-ref", "--verify", "refs/remotes/origin/other", cwd=self.checkout) + + def test_annotated_tag_resolves_commit_and_survives_rerun(self): + self.publish_tag("annotated-release", annotated=True) + self.prepare("refs/tags/annotated-release") + self.assert_detached(self.first) + self.assertNotEqual( + self.git( + "rev-parse", "refs/tags/annotated-release", cwd=self.checkout + ).stdout.strip(), + self.first, + ) + self.prepare("refs/tags/annotated-release") + self.assert_detached(self.first) + + def test_legacy_single_branch_clone_can_migrate_to_tag(self): + self.checkout.parent.mkdir(parents=True) + self.git( + "clone", + "--single-branch", + "--no-tags", + "--branch", + "other", + self.remote.as_uri(), + str(self.checkout), + ) + newer = self.commit("new release") + self.publish_tag("next-release", annotated=True) + self.prepare("refs/tags/next-release") + self.assert_detached(newer) + self.assertEqual( + self.git("rev-parse", "refs/heads/other", cwd=self.checkout).stdout.strip(), + self.first, + ) + + def test_same_named_branch_does_not_override_explicit_tag(self): + newer = self.commit("branch sharing the tag name") + self.git("branch", TAG, newer, cwd=self.seed) + self.git("push", "origin", f"refs/heads/{TAG}", cwd=self.seed) + self.prepare() + self.assert_detached(self.first) + self.prepare() + self.assert_detached(self.first) + + def test_moved_tag_refused_without_changing_head_or_local_tag(self): + self.prepare() + old_tag = self.git("rev-parse", f"refs/tags/{TAG}", cwd=self.checkout).stdout + self.commit("retagged release") + self.git("tag", "-f", TAG, cwd=self.seed) + self.git("push", "--force", "origin", f"refs/tags/{TAG}", cwd=self.seed) + result = self.prepare(expected=False) + self.assertIn("missing or moved", result.stderr) + self.assert_detached(self.first) + self.assertEqual( + self.git("rev-parse", f"refs/tags/{TAG}", cwd=self.checkout).stdout, old_tag + ) + + def test_replaced_annotated_tag_is_refused_even_at_same_commit(self): + self.publish_tag("annotated-release", annotated=True) + self.prepare("refs/tags/annotated-release") + self.git( + "tag", + "-f", + "-a", + "-m", + "changed annotation", + "annotated-release", + cwd=self.seed, + ) + self.git( + "push", "--force", "origin", "refs/tags/annotated-release", cwd=self.seed + ) + self.prepare("refs/tags/annotated-release", expected=False) + self.assert_detached(self.first) + + def test_deleted_tag_is_not_reused_or_replaced_by_same_named_branch(self): + self.prepare() + self.git("branch", TAG, cwd=self.seed) + self.git( + "push", "origin", f"refs/heads/{TAG}", f":refs/tags/{TAG}", cwd=self.seed + ) + # Even user-configured pruning must not remove the local release pin. + self.git("config", "fetch.prune", "true", cwd=self.checkout) + self.git("config", "fetch.pruneTags", "true", cwd=self.checkout) + self.prepare(expected=False) + self.assert_detached(self.first) + self.git("show-ref", "--verify", f"refs/tags/{TAG}", cwd=self.checkout) + + def test_missing_tag_does_not_fall_back_to_branch(self): + self.git("branch", "branch-only", cwd=self.seed) + self.git("push", "origin", "branch-only", cwd=self.seed) + result = self.prepare("refs/tags/branch-only", expected=False) + self.assertIn("missing or moved", result.stderr) + + def test_non_commit_tag_is_rejected(self): + blob = self.git("rev-parse", "HEAD:payload", cwd=self.seed).stdout.strip() + self.git("tag", "blob-tag", blob, cwd=self.seed) + self.git("push", "origin", "refs/tags/blob-tag", cwd=self.seed) + result = self.prepare("refs/tags/blob-tag", expected=False) + self.assertIn("does not resolve to a commit", result.stderr) + + def test_dirty_tracked_and_untracked_files_are_preserved(self): + self.prepare() + for name in ("payload", "untracked"): + with self.subTest(name=name): + path = self.checkout / name + original = path.read_bytes() if path.exists() else None + path.write_text("operator changes") + result = self.prepare(expected=False) + self.assertIn("local changes", result.stderr) + self.assertEqual(self.head(), self.first) + self.assertEqual(path.read_text(), "operator changes") + if original is None: + path.unlink() + else: + path.write_bytes(original) + + def test_new_tag_can_be_selected_explicitly(self): + self.prepare() + newer = self.commit("next release") + self.publish_tag("internal-testnet-v1") + self.prepare("refs/tags/internal-testnet-v1") + self.assert_detached(newer) + self.assertEqual( + self.git("rev-parse", f"refs/tags/{TAG}", cwd=self.checkout).stdout.strip(), + self.first, + ) + + def test_origin_mismatch_refused(self): + self.prepare() + self.git( + "remote", + "set-url", + "origin", + (self.root / "wrong.git").as_uri(), + cwd=self.checkout, + ) + result = self.prepare(expected=False) + self.assertIn("unexpected origin", result.stderr) + self.assert_detached(self.first) + + def test_dangling_source_symlink_refused(self): + self.checkout.parent.mkdir(parents=True) + target = self.root / "must-not-create" + self.checkout.symlink_to(target) + result = self.prepare(expected=False) + self.assertIn("symbolic link", result.stderr) + self.assertFalse(target.exists()) + + def test_branch_selector_does_not_accept_a_tag(self): + self.publish_tag("tag-only") + result = self.prepare("tag-only", expected=False) + self.assertIn("use refs/tags/", result.stderr) + + def test_branch_clone_and_fast_forward_still_work(self): + self.prepare("main") + self.assertEqual( + self.git("branch", "--show-current", cwd=self.checkout).stdout.strip(), + "main", + ) + self.git("show-ref", "--verify", "refs/remotes/origin/other", cwd=self.checkout) + newer = self.commit("fast forward") + self.git("push", "origin", "main", cwd=self.seed) + self.prepare("refs/heads/main") + self.assertEqual(self.head(), newer) + + def test_legacy_single_branch_clone_can_change_branch_and_return_from_tag(self): + self.checkout.parent.mkdir(parents=True) + self.git( + "clone", + "--single-branch", + "--branch", + "other", + self.remote.as_uri(), + str(self.checkout), + ) + self.prepare("main") + self.assertEqual( + self.git( + "config", "--get", "remote.origin.fetch", cwd=self.checkout + ).stdout.strip(), + "+refs/heads/*:refs/remotes/origin/*", + ) + self.prepare() + self.assert_detached(self.first) + self.prepare("main") + self.assertEqual( + self.git("branch", "--show-current", cwd=self.checkout).stdout.strip(), + "main", + ) + + def test_non_fast_forward_branch_update_refused(self): + self.prepare("main") + (self.checkout / "local").write_text("local commit") + self.git("add", "local", cwd=self.checkout) + self.git("commit", "-m", "local commit", cwd=self.checkout) + local = self.head() + self.commit("remote commit") + self.git("push", "origin", "main", cwd=self.seed) + result = self.prepare("main", expected=False) + self.assertIn("fast-forward only", result.stderr) + self.assertEqual(self.head(), local) + + def test_force_pushed_branch_history_refused(self): + self.prepare("main") + self.git("checkout", "--orphan", "replacement", cwd=self.seed) + self.commit("replacement history") + self.git("push", "--force", "origin", "HEAD:refs/heads/main", cwd=self.seed) + self.prepare("main", expected=False) + self.assertEqual(self.head(), self.first) + + def test_invalid_refs_refused_before_creating_checkout(self): + for ref in ( + "", + "--option", + "HEAD", + "refs/tags/", + "refs/tags/v0^{commit}", + "refs/remotes/origin/main", + "refs/tags/a:b", + ): + with self.subTest(ref=ref): + self.prepare(ref, expected=False) + self.assertFalse(self.checkout.exists()) + + def test_configured_release_refs_are_explicit_tags(self): + result = subprocess.run( + [ + "bash", + "-c", + """set -euo pipefail +source "$1/install/lib/configuration.sh" +section() { :; } +_out() { :; } +configure_component_installation() { :; } +print_component_installation() { :; } +confirm() { return 1; } +configure_node_software +configure_checkpointer +configure_custodian +printf '%s\\n' "$SUMMIT_SOURCE_REF" "$RETH_SOURCE_REF" "$CUSTODIAN_SOURCE_REF" "$CHECKPOINTER_SOURCE_REF" "$CUSTODIAN_REQUIRED_SUMMIT_REF" "$CUSTODIAN_REQUIRED_RETH_REF" +""", + "fixture", + str(ROOT), + ], + env=self.env, + capture_output=True, + text=True, + check=False, + timeout=10, + ) + self.assertEqual(result.returncode, 0, result.stderr) + self.assertEqual( + result.stdout.splitlines(), [f"refs/tags/{TAG}"] * 3 + ["main", TAG, TAG] + ) + + +if __name__ == "__main__": + unittest.main() diff --git a/tools/seismic-node.py b/tools/seismic-node.py index dfb6777..05e8626 100755 --- a/tools/seismic-node.py +++ b/tools/seismic-node.py @@ -17,7 +17,15 @@ from pathlib import Path from typing import NoReturn -from seismic_node import checkpoint, download, observer, rpc, supervisor, validator +from seismic_node import ( + checkpoint, + download, + monitoring, + observer, + rpc, + supervisor, + validator, +) def add_inventory_argument(parser: argparse.ArgumentParser) -> None: @@ -150,7 +158,7 @@ def parse_args() -> argparse.Namespace: validator_onboard = validator_commands.add_parser( "onboard", allow_abbrev=False, - help="install or validate a checkpoint and start validator services", + help="wait for validator lifecycle authorization and start services", formatter_class=argparse.RawDescriptionHelpFormatter, epilog=( "Manual checkpoint-start equivalent after lifecycle authorization:\n" @@ -160,11 +168,24 @@ def parse_args() -> argparse.Namespace: " sudo supervisorctl start custodian # when configured\n" " sudo supervisorctl start reth\n" " sudo supervisorctl start summit-checkpoint\n" - " sudo supervisorctl start checkpointer # when configured" + " sudo supervisorctl start checkpointer # when configured\n" + "For --mode normal, start summit instead of summit-checkpoint.\n" + "Normal mode waits for lifecycle authorization without installing " + "or requiring a checkpoint." ), ) + validator_onboard.add_argument( + "--mode", + choices=("normal", "checkpoint"), + default="checkpoint", + help="startup mode (default: checkpoint); normal skips checkpoint installation", + ) add_inventory_argument(validator_onboard) - validator_onboard.add_argument("--deposit-signature", type=Path, required=True) + validator_onboard.add_argument( + "--deposit-signature", + type=Path, + help="optional deposit response to cross-check against the installed node identity", + ) validator_onboard.add_argument( "--pre-joining-policy", choices=("wait", "start", "leave-stopped"), @@ -268,6 +289,17 @@ def parse_args() -> argparse.Namespace: ) add_inventory_argument(observer_stop) + monitoring_parser = commands.add_parser("monitoring", allow_abbrev=False) + monitoring_commands = monitoring_parser.add_subparsers( + dest="monitoring_command", required=True + ) + for action in ("start", "stop", "status"): + command = monitoring_commands.add_parser(action, allow_abbrev=False) + command.add_argument("--role", choices=("validator", "observer"), required=True) + add_inventory_argument(command) + if action == "start": + add_startup_argument(command) + return parser.parse_args() @@ -407,10 +439,25 @@ def handle_validator(args: argparse.Namespace) -> None: return if args.summit_rpc_url is None: raise checkpoint.CheckpointError("validator onboard requires --summit-rpc-url") - _, node_public_key = validator.load_deposit_response(args.deposit_signature) source_requested = checkpoint_source_requested(args) + if args.mode == "normal" and ( + source_requested or checkpoint_install_options_requested(args) + ): + raise checkpoint.CheckpointError( + "Checkpoint source and installation options cannot be used with " + "validator onboard --mode normal" + ) require_checkpoint_source_for_install_options(args, source_requested) - if not source_requested: + inventory_path = args.inventory or checkpoint.DEFAULT_INVENTORY_PATHS["validator"] + inventory = checkpoint.load_inventory("validator", inventory_path) + node_public_key = validator.installed_node_public_key(inventory) + if args.deposit_signature is not None: + _, deposit_public_key = validator.load_deposit_response(args.deposit_signature) + if deposit_public_key != node_public_key: + raise checkpoint.CheckpointError( + "Deposit-signature node public key does not match the installed validator" + ) + if args.mode == "checkpoint" and not source_requested: checkpoint.validate_checkpoint_start_configuration("validator") # Use one deadline across both lifecycle checks so installation cannot reset # the operator's total wait budget. @@ -438,16 +485,18 @@ def handle_validator(args: argparse.Namespace) -> None: if not start_decision.start: if source_requested: print("Checkpoint was installed. All validator services remain stopped.") - else: + elif args.mode == "checkpoint": print( "Existing checkpoint remains installed. " "All validator services remain stopped." ) + else: + print("All validator services remain stopped. No checkpoint was installed.") return try: - # Recheck status after installation because lifecycle state may have - # changed while a large checkpoint was downloaded and installed. - validator.start_checkpoint_validator( + # Recheck immediately before startup, including after any checkpoint + # download and installation that may have taken a long time. + validator.start_onboarded_validator( args, node_public_key, allow_pre_joining_start=start_decision.pre_joining, @@ -496,6 +545,8 @@ def main() -> NoReturn: handle_checkpoint(args) elif args.command == "validator": handle_validator(args) + elif args.command == "monitoring": + monitoring.handle(args) else: handle_observer(args) raise SystemExit(0) diff --git a/tools/seismic_node/README.md b/tools/seismic_node/README.md index f8a018a..d280737 100644 --- a/tools/seismic_node/README.md +++ b/tools/seismic_node/README.md @@ -19,12 +19,34 @@ boundary can be reviewed and tested independently. | `download.py` | Selects a completed checkpoint epoch, polls remote manifests, downloads verified archives, and resolves local, URL, or Summit-RPC weak-subjectivity anchors. | | `rpc.py` | Provides strict standard-library HTTP and JSON-RPC helpers with URL, redirect, response-size, token-file, archive-size, and SHA-256 checks. | | `supervisor.py` | Starts selected Supervisor programs in dependency order and reverses only the start requests made by the current command after a partial failure. | -| `validator.py` | Generates the exact deposit-signature response, applies validator lifecycle policy before checkpoint startup, and coordinates lifecycle-free validator restarts and shutdown. | +| `validator.py` | Generates the exact deposit-signature response, applies validator lifecycle policy before normal or checkpoint startup, and coordinates lifecycle-free validator restarts and shutdown. | | `observer.py` | Selects normal or checkpoint observer startup and coordinates observer shutdown without applying validator lifecycle rules. | `__init__.py` only identifies the internal package. It intentionally performs no startup work or global configuration. +## Optional monitoring + +`monitoring.py` manages the independent Prometheus Agent configured by the +[installers](../../install/PROMETHEUS_AGENT.md). A shared best-effort helper +runs only after successful validator/observer startup; onboarding authorization, +identity checks, checkpoint validation, and node-start rollback are unchanged. +Old inventories without `[monitoring]` remain supported. Invalid agent metadata +or agent startup failure warns without rolling back successful node startup. + +Node stop and rollback never stop the agent. Explicit commands: + +```bash +sudo ./tools/seismic-node.py monitoring start --role validator +sudo ./tools/seismic-node.py monitoring status --role validator +sudo ./tools/seismic-node.py monitoring stop --role validator +``` + +Use `--role observer` or `--inventory /absolute/path.toml` as appropriate. +Explicit monitoring start updates only the agent Supervisor group, not node +groups. Status and stop do not reload Supervisor. Monitoring commands read the +protected inventory without requiring usable node databases or keys. + ## Workflow boundaries ### Checkpoint acquisition @@ -87,8 +109,20 @@ automatically, and non-interactive use requires `--backup`. ### Validator onboarding -The validator flow uses the node public key from a strictly validated -`deposit-signature.json` response to call `getValidatorAccount`. +The validator flow reads `summit_keys_dir` from the validated installation +inventory, then invokes `/usr/local/bin/summit keys show --key-store-path DIR` +to derive the installed node public key. Only that public key is sent to +`getValidatorAccount`; private keys remain local and no deposit RPC is started. +The command is read-only, has a bounded runtime, and its raw output is never +included in errors. The parsed identity must be exactly one 32-byte ed25519 +public key. + +`--deposit-signature` is optional. When supplied, the response is strictly +validated and its node public key must match the installed identity before any +lifecycle polling or checkpoint installation. The installed identity is derived +again before startup and must still match the public key whose lifecycle was +checked. `--inventory` overrides the default +`/etc/seismic/validator-installation.toml` without requiring a signature file. - `Joining` starts normally. - `Active` starts with a late-onboarding warning. @@ -97,9 +131,28 @@ The validator flow uses the node public key from a strictly validated - `SubmittedExitRequest`, `FullPayoutPending`, malformed responses, and unknown statuses are refused. -The lifecycle state is checked before preparation and again immediately before -startup. A confirmed early-start decision is carried across the second check so -an interactive operator is not prompted twice. +`validator onboard --mode checkpoint` is the default: it installs a selected +checkpoint or validates existing checkpoint-start inputs before startup. +`validator onboard --mode normal` instead starts from local state without +installing or requiring a checkpoint. Normal mode rejects checkpoint sources and +installation modifiers, including `--yes`. + +For lifecycle-gated normal startup: + +```bash +sudo ./tools/seismic-node.py validator onboard \ + --mode normal \ + --summit-rpc-url https://trusted-validator.example/summit \ + --pre-joining-policy wait +``` + +Both modes validate the installation inventory and apply the same lifecycle +policy. The lifecycle state is checked before preparation and again immediately +before startup. A confirmed early-start decision is carried across the second +check so an interactive operator is not prompted twice. One wait deadline spans +both checks; `--validator-wait-timeout 0` (the default) waits indefinitely. +Normal startup does not guarantee synchronization from empty state. +`validator start --mode normal` remains the lifecycle-free startup command. ### Supervisor startup diff --git a/tools/seismic_node/checkpoint.py b/tools/seismic_node/checkpoint.py index 48d953e..2b471cc 100644 --- a/tools/seismic_node/checkpoint.py +++ b/tools/seismic_node/checkpoint.py @@ -305,6 +305,11 @@ def load_inventory(role: str, path: Path) -> dict[str, Any]: expected = COMMON_INVENTORY_KEYS | ( OBSERVER_INVENTORY_KEYS if role == "observer" else set() ) + # Optional telemetry metadata must not make old inventories incompatible. + # Its validation is deferred to best-effort monitoring startup so a broken + # agent configuration cannot prevent otherwise valid node startup. + if "monitoring" in inventory: + expected = expected | {"monitoring"} require_exact_keys(inventory, expected, "Installation inventory") if inventory["schema_version"] != INVENTORY_VERSION: raise CheckpointError( diff --git a/tools/seismic_node/monitoring.py b/tools/seismic_node/monitoring.py new file mode 100644 index 0000000..2397cb3 --- /dev/null +++ b/tools/seismic_node/monitoring.py @@ -0,0 +1,84 @@ +"""Independent, optional Prometheus Agent lifecycle. + +Automatic startup is best effort and happens only after node startup succeeds. +Monitoring is deliberately absent from node shutdown/rollback sequences. +""" + +from __future__ import annotations + +import sys +from typing import Any + +from . import checkpoint, supervisor + +PROGRAM = "prometheus-agent" + + +def configured(inventory: dict[str, Any], role: str) -> bool: + settings = inventory.get("monitoring", {"enabled": False}) + if not isinstance(settings, dict) or type(settings.get("enabled")) is not bool: + raise checkpoint.CheckpointError("Invalid monitoring inventory settings") + if not settings["enabled"]: + return False + required = {"enabled", "node", "role", "url", "data_dir", "version"} + if set(settings) != required or settings["role"] != role: + raise checkpoint.CheckpointError( + "Monitoring inventory does not match node role" + ) + if any( + not isinstance(settings[key], str) or not settings[key] + for key in required - {"enabled"} + ): + raise checkpoint.CheckpointError("Invalid monitoring inventory values") + return True + + +def ensure_running(inventory: dict[str, Any], role: str, timeout: float) -> None: + """Do not let telemetry errors undo authorized node startup.""" + try: + if configured(inventory, role): + supervisor.start_program(PROGRAM, timeout) + print("Prometheus Agent is running.") + except (checkpoint.CheckpointError, supervisor.SupervisorError, OSError): + print( + "WARNING: Prometheus Agent could not be started or verified. " + "Node startup was not rolled back. Inspect monitoring status and agent logs.", + file=sys.stderr, + ) + + +def prepare_agent() -> None: + """Load only this program; a monitoring command must not update node groups.""" + supervisor.require_command_success( + supervisor.run_systemctl("enable", "--now", "supervisor"), + "Enabling and starting Supervisor", + ) + supervisor.require_command_success( + supervisor.run_supervisorctl("reread"), "Rereading Supervisor configuration" + ) + supervisor.require_command_success( + supervisor.run_supervisorctl("update", PROGRAM), "Updating agent configuration" + ) + + +def handle(args: Any) -> None: + inventory_path = args.inventory or checkpoint.DEFAULT_INVENTORY_PATHS[args.role] + # Monitoring can be inspected even when node keys/state are unavailable. + inventory = checkpoint.read_toml( + inventory_path, "Installation inventory", root_managed=True + ) + if not configured(inventory, args.role): + if args.monitoring_command == "status": + print("Prometheus Agent is not configured/enabled in this inventory.") + return + raise checkpoint.CheckpointError("Prometheus Agent is not configured/enabled") + if args.monitoring_command == "start": + prepare_agent() + supervisor.start_program(PROGRAM, args.startup_timeout) + print("Prometheus Agent is running.") + elif args.monitoring_command == "stop": + supervisor.stop_autorestarting_program(PROGRAM) + print("Prometheus Agent stopped. Node services were not changed.") + else: + value = supervisor.status(PROGRAM) + print(f"{PROGRAM}: {value.state}") diff --git a/tools/seismic_node/observer.py b/tools/seismic_node/observer.py index 99f4a2b..064dfb1 100644 --- a/tools/seismic_node/observer.py +++ b/tools/seismic_node/observer.py @@ -8,13 +8,13 @@ from typing import Any -from . import checkpoint, supervisor +from . import checkpoint, monitoring, supervisor def start_observer(args: Any) -> None: """Validate the selected mode and delegate ordered startup to Supervisor.""" inventory_path = args.inventory or checkpoint.DEFAULT_INVENTORY_PATHS["observer"] - checkpoint.load_inventory("observer", inventory_path) + inventory = checkpoint.load_inventory("observer", inventory_path) if args.mode == "checkpoint": checkpoint.validate_checkpoint_start_configuration("observer") summit_program = "summit-observer-checkpoint" @@ -30,6 +30,7 @@ def start_observer(args: Any) -> None: conflicting_program, startup_timeout=args.startup_timeout, ) + monitoring.ensure_running(inventory, "observer", args.startup_timeout) print(f"Observer {args.mode} startup requested successfully.") @@ -39,4 +40,4 @@ def stop_observer(args: Any) -> None: checkpoint.load_inventory("observer", inventory_path) supervisor.stop_node(("summit-observer", "summit-observer-checkpoint")) print("Observer services stopped successfully.") - print("Supervisor and OpenResty remain running.") + print("Supervisor, OpenResty and any running Prometheus Agent remain running.") diff --git a/tools/seismic_node/supervisor.py b/tools/seismic_node/supervisor.py index 1fa664c..48d7f6d 100644 --- a/tools/seismic_node/supervisor.py +++ b/tools/seismic_node/supervisor.py @@ -202,6 +202,25 @@ def stop_program(name: str) -> bool: return True +def stop_autorestarting_program(name: str, timeout: float = 60.0) -> None: + """EXITED is not quiescent for autorestart=true programs. + + Supervisor may report NOT_RUNNING while between EXITED and a restart. + Retry until an administrative stop is confirmed; do not change node rules. + """ + deadline = time.monotonic() + timeout + while True: + current = status(name) + if not current.exists or current.state in {"STOPPED", "FATAL"}: + return + result = run_supervisorctl("stop", name) + if result.returncode not in (0, 7): + raise SupervisorError(f"Could not stop {name}; inspect Supervisor") + if time.monotonic() >= deadline: + raise SupervisorError(f"Timed out confirming {name} is stopped") + time.sleep(0.1) + + def optional_programs() -> tuple[bool, bool]: """Detect optional Custodian and checkpointer Supervisor programs.""" return status("custodian").exists, status("checkpointer").exists diff --git a/tools/seismic_node/validator.py b/tools/seismic_node/validator.py index 3527c77..37ab11d 100644 --- a/tools/seismic_node/validator.py +++ b/tools/seismic_node/validator.py @@ -1,24 +1,27 @@ """Validator deposit-signature generation, lifecycle-gated startup, and shutdown. Deposit signing is intentionally isolated behind the temporary loopback-only -Summit deposit RPC. Onboarding then derives the validator identity from that -exact response and asks a trusted Summit RPC whether startup is safe. +Summit deposit RPC. Onboarding derives identity from the installed keys via +Summit's read-only keys command and asks a trusted RPC whether startup is safe. """ from __future__ import annotations import json +import re +import subprocess import sys import time from dataclasses import dataclass from pathlib import Path from typing import Any -from . import checkpoint, rpc, supervisor +from . import checkpoint, monitoring, rpc, supervisor DEPOSIT_AMOUNT_GWEI = 32_000_000_000 WITHDRAWAL_ADDRESS = "0xd412c5ecd343e264381ff15afc0ad78a67b79f35" DEPOSIT_RPC_URL = "http://127.0.0.1:3031" +SUMMIT = Path("/usr/local/bin/summit") PRE_JOINING_STATUSES = {"NotFound", "Inactive"} REFUSED_STATUSES = {"SubmittedExitRequest", "FullPayoutPending"} @@ -109,6 +112,41 @@ def load_deposit_response(path: Path) -> tuple[dict[str, Any], str]: return response, validate_deposit_response(response) +def installed_node_public_key(inventory: dict[str, Any]) -> str: + """Read public identity using Summit's read-only command, never print secrets.""" + keys_dir = inventory["summit_keys_dir"] + checkpoint.require_directory(keys_dir, "Summit keys directory") + for name in ("node_key.pem", "consensus_key.pem"): + checkpoint.require_regular_file(keys_dir / name, f"Summit {name}") + try: + result = subprocess.run( + [str(SUMMIT), "keys", "show", "--key-store-path", str(keys_dir)], + stdin=subprocess.DEVNULL, + stdout=subprocess.PIPE, + stderr=subprocess.DEVNULL, + text=True, + check=True, + timeout=30.0, + ) + except (OSError, subprocess.SubprocessError, UnicodeError): + # Do not surface subprocess output: key-decoding failures could include + # sensitive input. Only the validated public key may leave this helper. + raise checkpoint.CheckpointError( + f"Could not read installed validator identity with {SUMMIT} keys show; " + "check the Summit executable and installed key files" + ) from None + public_keys = re.findall( + r"^[ \t]*Node Public Key \(ed25519\):[ \t]*(?:0[xX])?([0-9a-fA-F]{64})[ \t]*$", + result.stdout, + re.MULTILINE, + ) + if len(public_keys) != 1: + raise checkpoint.CheckpointError( + "Summit keys show did not return exactly one valid ed25519 node public key" + ) + return public_keys[0].lower() + + def prompt_operator(message: str) -> str: """Read an interactive answer or explain which explicit option is needed.""" try: @@ -377,8 +415,8 @@ def wait_for_start_authorization( f"Refusing unknown validator lifecycle state {current_status!r}" ) - # A second status check occurs after checkpoint installation. Preserve a - # previously confirmed early-start choice instead of prompting twice. + # A second status check occurs before startup in either mode. Preserve + # a previously confirmed early-start choice instead of prompting twice. if allow_pre_joining_start: action = "start" else: @@ -399,7 +437,7 @@ def wait_for_start_authorization( time.sleep(args.validator_poll_interval) -def start_checkpoint_validator( +def start_onboarded_validator( args: Any, node_public_key: str, *, @@ -414,19 +452,17 @@ def start_checkpoint_validator( wait_deadline=wait_deadline, ) if not decision.start: - print("Checkpoint remains installed. All validator services remain stopped.") + if args.mode == "checkpoint": + print( + "Checkpoint remains installed. All validator services remain stopped." + ) + else: + print("All validator services remain stopped. No checkpoint was installed.") return - checkpoint.validate_checkpoint_start_configuration("validator") - supervisor.prepare_supervisor() - supervisor.start_node( - "summit-checkpoint", - "summit", - startup_timeout=args.startup_timeout, - ) - print("Validator checkpoint startup requested successfully.") + start_validator(args, expected_node_public_key=node_public_key) -def start_validator(args: Any) -> None: +def start_validator(args: Any, *, expected_node_public_key: str | None = None) -> None: """Start validator services without lifecycle checks or checkpoint downloads. ``validator onboard`` remains the lifecycle-gated entry point; this command @@ -434,7 +470,15 @@ def start_validator(args: Any) -> None: from the installed checkpoint inputs. """ inventory_path = args.inventory or checkpoint.DEFAULT_INVENTORY_PATHS["validator"] - checkpoint.load_inventory("validator", inventory_path) + inventory = checkpoint.load_inventory("validator", inventory_path) + if ( + expected_node_public_key is not None + and installed_node_public_key(inventory) != expected_node_public_key + ): + raise checkpoint.CheckpointError( + "Installed validator node public key changed during onboarding; " + "refusing to start" + ) if args.mode == "checkpoint": checkpoint.validate_checkpoint_start_configuration("validator") summit_program = "summit-checkpoint" @@ -450,6 +494,7 @@ def start_validator(args: Any) -> None: conflicting_program, startup_timeout=args.startup_timeout, ) + monitoring.ensure_running(inventory, "validator", args.startup_timeout) print(f"Validator {args.mode} startup requested successfully.") @@ -459,4 +504,4 @@ def stop_validator(args: Any) -> None: checkpoint.load_inventory("validator", inventory_path) supervisor.stop_node(("summit-deposit-rpc", "summit", "summit-checkpoint")) print("Validator services stopped successfully.") - print("Supervisor and OpenResty remain running.") + print("Supervisor, OpenResty and any running Prometheus Agent remain running.")