Skip to content

refactor(asap-tools): decouple remote_monitor lifecycle and client ownership #696

Description

@milindsrivastava1997

Context

The remote-monitor decoupling design in the .design_docs directory remains relevant and is not implemented on main. The current code still uses execution-mode branches, pgrep-based monitor discovery/lifecycle polling, the PR #522 ingest stop-file plumbing, and the legacy profile_query_engine_pid path.

This issue consolidates the remaining work described in the three design documents attached as comments below. It supersedes the need to track the work only through the older exploratory issues #23 and #53, while directly addressing the still-open concerns in #35 and #91.

Related: #23, #35, #53, #91, and merged PR #522.

Proposed outcome

  • Replace pgrep/stop-file lifecycle handling with a PID-file + command-file protocol and SIGTERM escalation.
  • Decouple remote_monitor.py from QueryClientService.
  • Replace the growing execution_mode interface with a structured monitor request.
  • Consolidate profiler lifecycle code behind a common interface.
  • Migrate the ClickHouse and e2e call sites incrementally, with behavioral verification.

The three source documents are included as comments on this issue.

Activity

  1. milindsrivastava1997 commented on Sep 3, 2026

    @milindsrivastava1997
    ContributorAuthor

    Remote monitor decoupling: staged implementation plan

    Status: Plan drafted, not started
    Design of record: .design_docs/remote-monitor-decoupling-design.md
    PR description (draft): .design_docs/remote-monitor-decoupling-PR-description.md
    Prerequisite: PR #522 merged to main.

    Principles for this refactor

    • Tiny commits, one concern each. Each stage below is a commit (or a small
      handful) that leaves main green and the experiment scripts runnable.
    • Additive-then-subtractive. New mechanisms land alongside the old ones and
      one call site is migrated first; the old path is deleted only in a later stage,
      once nothing calls it. No stage both adds a mechanism and removes its
      predecessor.
    • Manual edits, no scripts. Renames/migrations done by hand, call site by
      call site.
    • Pause after each stage for a manual commit — do not batch.
    • Verification is behavioral, not just typecheck. Where a stage changes
      runtime wiring, the verify step names the actual experiment run to exercise.

    Ordering rationale: earliest stages are pure deletions/internal refactors with
    zero behavior change (lowest risk, build confidence), then additive protocol
    work, then the one genuinely behavioral change (client ownership), then
    subtractive cleanup once the old path is dead.


    Stage 0 — Remove dead profile_query_engine_pid end-to-end

    Why first: completely isolated, provably unreachable, zero behavior change —
    the safest possible starting commit, and it shrinks the surface every later stage
    touches.

    Files:

    • remote_monitor.py: delete the profile_query_engine_pid = None var (lines
      ~312-314) and its pass-through into PrometheusClientService.start (~399).
    • experiment_utils/services/prometheus_client_service.py: drop the
      profile_query_engine_pid param from start / _start_containerized /
      _start_bare_metal and the two if profile_query_engine_pid is not None
      branches.
    • generate_prometheus_client_compose.py: drop the --profile-query-engine-pid
      arg and its template_vars wiring.
    • asap-tools/queriers/prometheus-client/docker-compose.yml.j2: remove the
      --profile_query_engine_pid line.
    • asap-tools/queriers/prometheus-client/main_prometheus_client.py: remove the
      --profile_query_engine_pid argparse arg and the start_query_engine_profiler
      thread block (~688-701); delete start_query_engine_profiler if now unused.

    Verify: grep -rn profile_query_engine_pid returns nothing;
    grep -rn profile_query_engine\b still shows the live --profile_query_engine
    path intact. Bare-metal + containerized prometheus-client still start (dry-run
    the generated compose / command string).

    Closes: issue #23 second bullet.


    Stage 1 — Extract Profiler interface (internal refactor, same behavior)

    Why here: pure reorganization of remote_monitor.py internals; no
    orchestrator or wire-format change; nothing outside remote_monitor.py sees it.

    Files: new classes/profilers.py (or a section of remote_monitor.py):

    class Profiler(ABC):
        def start(self, pids, output_dir) -> Any: ...
        def stop(self, handle, store: bool) -> None: ...
    
    class AsprofProfiler(Profiler)      # from start/stop_profiling_flink_pids
    class FlamegraphProfiler(Profiler)  # from start/stop_profiling_arroyo_pids
    class PerfProfiler(Profiler)        # from start/stop_profiling_query_engine_pids
                                        # + convert_query_engine_perf_data folded into stop()

    remote_monitor.py keeps its current CLI and main() control flow, but the
    if args.profile_flink_pids: ... / arroyo / QE blocks now instantiate and call
    these classes instead of the free functions. Behavior identical.

    Verify: an experiment run with profile_flink/profile_arroyo/
    profile_query_engine enabled produces the same profile artifacts
    (flink_profiles/, arroyo_profiles/, query_engine_profiles/) as before.
    Diff the output tree against a pre-refactor run.

    Leave a check behind: a __main__ self-check in profilers.py asserting
    each profiler builds the same command string it did as a free function (guards
    the copy-paste consolidation).


    Stage 2 — Introduce MonitorRequest + --request arg, additive

    Why here: the wire format is the foundation the protocol stages build on, but
    it can land without touching lifecycle. Old flags stay; --request is an
    alternative entry point migrated one call site at a time.

    Files:

    • New classes/monitor_request.py: the MonitorTarget / ProfilerSpec /
      WaitStrategy / CostExporterSpec / MonitorRequest dataclasses + to_json
      / from_json. (Decision to record: shared module imported by both
      remote_monitor.py and remote_monitor_service.py, vs. duplicated schema —
      see design doc "Not yet decided". Recommend shared module; both run from the
      same experiments/ dir.)
    • remote_monitor.py: add --request arg. When present, build the run from it;
      when absent, fall back to today's flag parsing (a thin adapter that constructs
      a MonitorRequest from the old args, so there is exactly one downstream code
      path). execution_mode maps to WaitStrategy: timed→fixed_seconds,
      interactive→stdin, prometheus_client/ingest→stop_signal.
    • remote_monitor_service.py: no change yet.

    Verify: run the old flag path (unchanged) and a hand-built --request JSON
    for the same scenario; assert identical monitor_output.json. Round-trip test
    MonitorRequest.from_json(r.to_json()) == r.


    Stage 3 — PID file + command file protocol in remote_monitor.py

    Why here: adds the node-side half of the lifecycle contract. Still additive —
    the orchestrator doesn't use it yet, so nothing breaks.

    Files: remote_monitor.py:

    • On start (once output_dir is known): write <output_dir>/remote_monitor.pid
      = os.getpid(), unconditionally (overwrite stale). Register cleanup on clean
      exit.
    • Install a SIGTERM handler that flushes output + removes the PID file.
    • stop_signal wait strategy: each poll interval, check
      <output_dir>/remote_monitor.cmd; on {"action":"stop", "store":...}, break,
      stop profilers with store, write output, remove PID file, exit. (ControlCommand
      dataclass in monitor_request.py or a sibling.)
    • fixed_seconds / stdin strategies also remove the PID file on their normal
      exit (uniform cleanup).

    Verify: launch remote_monitor.py --request with stop_signal directly
    (no orchestrator); confirm .pid appears, kill(pid,0) succeeds; touch/write
    the .cmd file → it exits within one interval and writes output and removes
    .pid. Separately, send it SIGTERM mid-run → it flushes and removes .pid.

    Leave a check behind: small test_control_protocol.py that spawns the
    runner in a subprocess against a temp dir and asserts the .pid/.cmd/output
    lifecycle for both the command-file and SIGTERM paths.


    Stage 4 — MonitorSession orchestrator methods, additive

    Why here: node side (Stage 3) exists, so the orchestrator half can be built
    and unit-tested against it without yet migrating any experiment script.

    Files: remote_monitor_service.py — add to the existing class (rename to
    MonitorSession deferred to Stage 6 to keep this diff additive):

    • start(request: MonitorRequest, manual_mode=False): serialize to JSON, one
      shlex.quote, launch via the existing nohup … & path.
    • wait_for_start(output_dir, timeout=30): poll for .pid existence.
    • wait_for_finish(output_dir, timeout=600): poll for .pid gone.
    • stop(output_dir, store=True, timeout=60): write .cmd; wait_for_finish;
      on timeout read PID from .pid, SIGTERM over SSH, short grace, SIGKILL.

    Old methods (start legacy, wait_for_remote_monitor_to_finish,
    kill_remote_monitor, and PR #522's six) stay untouched and still used by the
    scripts. Both APIs coexist.

    Verify: on a CloudLab node, drive a bare monitor (no client) through
    start/wait_for_start/stop via a throwaway script; confirm .pid lifecycle
    and that stop's SIGTERM escalation fires when the monitor is artificially
    wedged (e.g. kill -STOP the process first).


    Stage 5 — Migrate experiment_run_clickhouse.py to the new session

    Why here: ClickHouse is the smaller / newer call site and already owns its
    client (ClickHouseDataLoaderService), so it's the lower-risk first migration
    and needs no client-ownership change — only swaps the monitor mechanism.

    Files: experiment_run_clickhouse.py:

    • Ingest monitor: replace start_clickhouse_ingest_monitor +
      signal_ingest_monitor_stop + wait_for_remote_monitor_process_exit +
      cleanup_ingest_monitor_stop_file with session.start(MonitorRequest(..., wait=stop_signal)) / wait_for_start / stop.
    • Query workload: same stop_signal session bracket around the existing SQL
      query client launch.

    Verify: full experiment_run_clickhouse run on a node; confirm
    monitor_output_ingest.json and monitor_output.json are both written and
    non-empty, matching a pre-migration baseline run's shape.


    Stage 6 — Migrate experiment_run_e2e.py + client ownership flip

    Why here: the one genuinely behavioral change — experiment_run_e2e.py must
    now launch/await QueryClientService itself instead of remote_monitor.py doing
    it. Isolated to its own stage so it can be reviewed and verified alone.

    Files:

    • experiment_run_e2e.py: around the current remote_monitor_service.start(...)
      • wait_for_remote_monitor_to_finish(...) block: session.start(stop_signal)
        → wait_for_start → launch & await QueryClientService (mirroring the
        ClickHouse data-loader pattern) → session.stop.
    • remote_monitor.py: delete the prometheus_client in-process client launch
      block (the PrometheusClientService import and its start/health-loop/stop) —
      now dead once e2e owns the client.

    Verify: full experiment_run_e2e run (both sketchdb and baseline modes,
    container and bare-metal prometheus-client); confirm query results +
    monitor_output.json match a pre-migration baseline. This is the highest-risk
    stage — run against a known-good prior experiment output and diff.

    Closes: issue #53, issue #23 first bullet.


    Stage 7 — Delete the old path

    Why last: only safe once Stages 5–6 removed every caller.

    Files:

    • remote_monitor.py: remove the legacy flag-parsing adapter and old CLI args;
      --request becomes the only entry point; main() is the straight-line runner.
    • remote_monitor_service.py: delete legacy start, kill_remote_monitor,
      wait_for_remote_monitor_to_finish, and PR feat(asap-tools): separate clickhouse ingest and query cpu/memory monitoring in experiment run clickhouse #522's six methods
      (start_clickhouse_ingest_monitor, is_remote_monitor_running,
      wait_for_remote_monitor_start, wait_for_remote_monitor_process_exit,
      signal_ingest_monitor_stop, cleanup_ingest_monitor_stop_file,
      _remote_monitor_pgrep_pattern); rename class to MonitorSession; remove
      constants.INGEST_MONITOR_STOP_FILE.
    • Remove now-unused keyword-building if/elif tower fed only by the old start.

    Verify: grep confirms no references to the deleted symbols;
    grep -rn pgrep remote_monitor returns nothing; one more full e2e + clickhouse
    run as a final regression check.

    Closes: issue #91 (hardcoded keyword tower gone), issue #35 (no more
    pgrep polling).


    Cross-cutting: process discovery (issue #91)

    Stages 5–6 pass MonitorTarget(pid=...) wherever a service already knows its PID
    (query engine, ClickHouse, Arroyo — several already expose
    get_monitoring_keyword(); extend to hand over the PID). MonitorTarget(keyword=...)

    • node-side resolve_targets() (housing today's get_pids()) remains only as the
      fallback for targets nothing owns. Enumerating exactly which current keywords can
      become PIDs vs. must stay keyword-based is deferred to Stage 5/6 execution — see
      design doc "Not yet decided".

    Suggested verification asset

    Before Stage 5, capture one known-good experiment_run_clickhouse and one
    experiment_run_e2e output tree on main as golden baselines to diff every
    subsequent stage against — the cheapest guard for a refactor this deep into the
    experiment critical path.

  2. milindsrivastava1997 commented on Sep 3, 2026

    @milindsrivastava1997
    ContributorAuthor

    Decoupling remote_monitor.py: design

    Status: Design agreed, staged implementation plan written (.design_docs/remote-monitor-decoupling-implementation-plan.md)
    Touches: asap-tools/experiments/remote_monitor.py, experiment_utils/services/remote_monitor_service.py, classes/process_monitor.py, experiment_run_e2e.py, experiment_run_clickhouse.py, experiment_utils/services/prometheus_client_service.py (QueryClientService)
    Relationship to PR #522: PR #522 (ClickHouse ingest/query monitor split) merges as-is on its own merits. This is a separate follow-up refactor that subsumes and removes the stop-file/pgrep-pattern plumbing PR #522 introduces, replacing it with the general mechanism below.


    Background

    remote_monitor.py runs on a CloudLab node (launched via SSH with nohup … & so it survives the launching SSH exec returning) and currently mixes four concerns: process discovery, CPU/memory sampling, ad hoc profiling (Flink/Arroyo/query-engine), and orchestrating the query client (PrometheusClientService/QueryClientService) itself. experiment_run_e2e.py and experiment_run_clickhouse.py drive it through RemoteMonitorService.

    The load-bearing constraint: why nohup … &

    The reason for nohup … & is to detach the launched process's lifetime from the single SSH exec that started it — the orchestrator's provider.execute_command fires each command as its own SSH exec that returns immediately, so without nohup & the child would get SIGHUP the moment that exec returns. It is not about the orchestrator disconnecting or staying connected. The consequence that shapes this whole design: once detached, the orchestrator holds no process handle and no pipe, so it must learn status and send control out-of-band. This is what rules out live-pipe IPC (e.g. naive multiprocessing across the boundary) and drives the file-based protocol below.

    Where the protocol files physically live, and who touches them across the SSH boundary

    remote_monitor.py runs on the node and uses CloudLabLocalProvider (local exec, no SSH) for its own subprocess needs. The orchestrator (MonitorSession, on the laptop) uses CloudLabProvider (SSH). The .pid and .cmd files live on the node's filesystem under experiment_output_dir:

    So "the orchestrator writes the command file" means an SSH touch/cat > of a node-local path; "the monitor reads it" means a local file read on its poll interval.

    Problems found

    1. Process discovery by string-matching. get_pids() shells out to ps aux | grep -E ... | awk or docker inspect. Every new component adds another keyword threaded through an if/elif tower in remote_monitor_service.py (issue Remove hardcoded keywords from remote_monitor_service #91).
    2. No real lifecycle contract across the SSH boundary. The orchestrator fires remote_monitor.py detached, then learns it's done by polling pgrep -f remote_monitor.py (issue Investigate why remote_monitor takes time to shutdown after experiment is over #35 — the slow-shutdown-detection complaint is this). PR feat(asap-tools): separate clickhouse ingest and query cpu/memory monitoring in experiment run clickhouse #522 needs a way to tell one specific invocation (the ClickHouse ingest monitor) to stop, and since no such primitive exists, it invents one from scratch: a stop-file, a pgrep regex pattern-builder with re.escape/shlex.quote to disambiguate instances, and five new wait/status methods.
    3. execution_mode is a growing, non-orthogonal enum. interactive / timed / prometheus_client, plus PR feat(asap-tools): separate clickhouse ingest and query cpu/memory monitoring in experiment run clickhouse #522's ingest. Each mode has bespoke branches in both main() and remote_monitor_service.start(), even though the real orthogonal axes are: which PIDs, how long to sample, which profilers, what happens when done.
    4. Profiling is three copy-pasted subsystems (start/stop_profiling_{flink,arroyo,query_engine}_pids) — same shape (subprocess per PID, track handle, SIGTERM on stop, optionally convert/store output), different tool per pair.
    5. The wire format is hand-built shell strings. remote_monitor_service.py constructs the entire remote invocation via .format() into a shell command; every new flag touches both the format string and remote_monitor.py's argparse in lockstep.
    6. Components orchestrate each other instead of the master orchestrator controlling everything (issue Re-architect code to use Redis for control messages #23, closed unimplemented): remote_monitor.py imports and directly starts/blocks on QueryClientService just so it knows when to start/stop sampling; main_prometheus_client.py in turn starts/stops query-engine profiling itself via --profile_query_engine_pid — but that value is hardcoded to None at its only live call site (remote_monitor.py:312-314, "unused for Rust QE"), so this second coupling is dead code for the current Rust query engine (which is profiled directly by remote_monitor.py via perf record, through the separate, live --profile_query_engine flag).

    Issue #23 proposed Redis as the control-message transport for "master orchestrator sends control signals to components." We don't need it: the actual requirement is single-node start/stop/status signaling, which a PID file + a small command file achieves with zero new infrastructure.

    Existing precedent for the target shape

    ClickHouseDataLoaderService.start() already has the outer script (experiment_run_clickhouse.py) launch and block on a client directly via provider.execute_command, independent of remote_monitor.py. PR #522's ingest-monitor bracket is this pattern applied to monitoring — it just lacked a clean stop primitive, so it built one (the stop-file) ad hoc. The design below generalizes that pattern and gives it the primitive it was missing, instead of it staying a one-off.


    Goals

    Non-goals

    • Not changing process_monitor.py's MyMonitor/start_monitor/stop_monitor internals — that's an in-process multiprocessing.Pipe boundary between remote_monitor.py and its sampler subprocess, a different and already-clean seam from the SSH-facing protocol below.
    • Not introducing Redis or any other new runtime dependency (considered per issue Re-architect code to use Redis for control messages #23, rejected — see Background).
    • Not changing how the remote process is launched (nohup … & over SSH stays; only how the orchestrator tracks/signals it afterward changes).

    Target design

    Control protocol: PID file + command file + SIGTERM escalation

    Three plain-filesystem pieces, all scoped under the already-unique experiment_output_dir, so no regex disambiguation between instances is ever needed:

    • PID file (<output_dir>/remote_monitor.pid): written unconditionally on start (a stale leftover from a crashed prior run is garbage, not a lock, and gets overwritten). Read by the orchestrator to check liveness (kill(pid, 0)) and as the target for the SIGTERM escalation path.

    • Command file (<output_dir>/remote_monitor.cmd): the general control-signal channel, checked by remote_monitor.py on its existing per-interval poll — no new wake mechanism needed. Generalizes PR feat(asap-tools): separate clickhouse ingest and query cpu/memory monitoring in experiment run clickhouse #522's boolean stop-file into a small structured command:

      @dataclass
      class ControlCommand:
          action: Literal["stop"]   # room to grow: future actions add values here, not new plumbing
          store: bool = True        # consolidates today's per-profiler `store: bool` into one shared field
    • SIGTERM (sent to the PID from the PID file): the escalation/safety-net only, not the primary channel — used when the command file goes unanswered past a timeout, mirroring the existing terminate()/kill() escalation already in process_monitor.stop_monitor.

    Request schema (replaces execution_mode + ~15 CLI flags)

    @dataclass
    class MonitorTarget:
        pid: Optional[int] = None       # preferred: caller's service already knows it
        keyword: Optional[str] = None   # fallback: resolved via ps/docker on the node
        label: str = ""                 # replaces today's "keyword" output tag
    
    @dataclass
    class ProfilerSpec:
        kind: Literal["asprof", "flamegraph", "perf"]
        pids: List[int]
    
    @dataclass
    class WaitStrategy:
        kind: Literal["fixed_seconds", "stdin", "stop_signal"]
        seconds: Optional[int] = None   # only for fixed_seconds
    
    @dataclass
    class CostExporterSpec:
        addr: str
        port: int
        monitors_and_models: dict
    
    @dataclass
    class MonitorRequest:
        targets: List[MonitorTarget]
        monitors: List[str]                      # ["memory_info", "cpu_percent"]
        include_children: bool
        thread_attribution_keyword: Optional[str]
        profilers: List[ProfilerSpec]
        wait: WaitStrategy
        interval_seconds: float
        output_dir: str
        output_file: str
        cost_exporter: Optional[CostExporterSpec]

    Today's four modes become wait values: timed → fixed_seconds, interactive → stdin, both prometheus_client and PR #522's ingest → stop_signal (same shape once the client is owned by the outer script — no reason for them to be different modes).

    Transport: MonitorRequest is serialized as one compact JSON blob passed via a single CLI arg (remote_monitor.py --request '<json>'), replacing today's ~15 hand-formatted flags. No new remote file-push mechanism — strictly less shell-escaping surface than today, not more.

    remote_monitor.py runner shape

    def main(request: MonitorRequest):
        write_pidfile(request.output_dir)             # unconditional overwrite
        install_sigterm_handler(...)                   # escalation path only
    
        pids, labels = resolve_targets(request.targets)   # get_pids() fallback lives here, small
        monitor, ctrl, mpipe = process_monitor.start_monitor(pids, labels, ...)  # unchanged internals
        profilers = [make_profiler(s).start() for s in request.profilers]
    
        wait(request.wait)   # sleep(N) | input() | poll command file each interval until action=="stop"
    
        for p in profilers: p.stop(store=<from command, default True>)
        monitor_info = process_monitor.stop_monitor(monitor, ctrl, mpipe)       # unchanged
        write_output(monitor_info, request.output_dir, request.output_file)
        remove_pidfile(request.output_dir)

    No QueryClientService/PrometheusClientService import. No execution_mode branches in main().

    Profiler interface (replaces 3 copy-pasted function pairs)

    class Profiler(ABC):
        def start(self, pids: List[int], output_dir: str) -> Any: ...
        def stop(self, handle: Any, store: bool) -> None: ...
    
    class AsprofProfiler(Profiler): ...      # today's start/stop_profiling_flink_pids
    class FlamegraphProfiler(Profiler): ...  # today's start/stop_profiling_arroyo_pids
    class PerfProfiler(Profiler): ...        # today's start/stop_profiling_query_engine_pids
                                              # + convert_query_engine_perf_data folded into stop()

    Orchestrator side: MonitorSession (replaces RemoteMonitorService)

    class MonitorSession:
        def __init__(self, provider, node_offset): ...
    
        def start(self, request: MonitorRequest, manual_mode: bool = False) -> None:
            ...  # serialize request, one shell-escape, launch via existing nohup ... &
    
        def wait_for_start(self, output_dir: str, timeout: int = 30) -> None:
            ...  # poll for <output_dir>/remote_monitor.pid existing
    
        def wait_for_finish(self, output_dir: str, timeout: int = 600) -> None:
            ...  # poll for pidfile gone (covers fixed_seconds/stdin self-exit)
    
        def stop(self, output_dir: str, store: bool = True, timeout: int = 60) -> None:
            ...  # write command file {action: stop, store}; wait_for_finish; SIGTERM+SIGKILL escalation on timeout

    This one class replaces RemoteMonitorService plus all six of PR #522's new methods (start_clickhouse_ingest_monitor, is_remote_monitor_running, wait_for_remote_monitor_start, wait_for_remote_monitor_process_exit, signal_ingest_monitor_stop, cleanup_ingest_monitor_stop_file, _remote_monitor_pgrep_pattern) — the PID-file scoping means there's never a need to pattern-match "which remote_monitor.py instance."

    Call sites become uniform

    Every place the outer script owns a client (ClickHouse ingest today, sketchdb query phase after this change) becomes the same four lines:

    session.start(MonitorRequest(..., wait=WaitStrategy("stop_signal")))
    session.wait_for_start(output_dir)
    client_service.start(...)              # data_loader today; query_client_service after this change
    if client_service.use_container:
        while client_service.is_healthy(): time.sleep(5)
        client_service.stop()
    session.stop(output_dir)

    experiment_run_e2e.py changes to launch/await QueryClientService itself here, the way experiment_run_clickhouse.py already does for ClickHouseDataLoaderService — this is the one real behavioral/control-flow change outside remote_monitor.py itself.

    Dead code to remove alongside this

    • profile_query_engine_pid end-to-end: remote_monitor.py, prometheus_client_service.py, generate_prometheus_client_compose.py's flag, docker-compose.yml.j2's template line, main_prometheus_client.py's start_query_engine_profiler thread. Confirmed unreachable — the only call site in this repo hardcodes it to None ("unused for Rust QE").

    Open questions resolved during design

    Question Decision
    Should the outer script always own client launch/await? Yes — matches existing ClickHouseDataLoaderService precedent, resolves issue #53.
    Request transport: one JSON CLI arg vs. pushed file? One JSON CLI arg — no new remote file-push mechanism needed.
    Stale PID files from a crashed prior run? Always overwritten on start; not treated as a lock.
    Redis (issue #23) for control signals? Dropped — PID file + command file + SIGTERM escalation gives the same decoupling with no new dependency, and generalizes to future signal types via ControlCommand.action.

    Not yet decided

    • Staged implementation/commit sequencing (next step).
    • Exact resolve_targets() responsibility split for target types with no owning service object (still falls back to keyword-based ps/docker discovery — scope of that fallback not yet enumerated).
    • Whether MonitorRequest/ControlCommand should be shared Python types imported by both sides, or independently-maintained schemas kept in sync by convention (matters once this is staged into commits).
  3. milindsrivastava1997 commented on Sep 3, 2026

    @milindsrivastava1997
    ContributorAuthor

    refactor(asap-tools): decouple remote_monitor lifecycle, replace pgrep/stop-file with PID+command files

    Draft — no PR open yet. Paste this into the PR when the refactor lands. See
    .design_docs/remote-monitor-decoupling-design.md for the full design.

    Summary

    Reworks how the experiment orchestrator launches, tracks, and stops
    remote_monitor.py on a CloudLab node. remote_monitor.py stops orchestrating
    the query client, execution_mode collapses into a request object, the three
    copy-pasted profilers become one Profiler interface, and pgrep-based status
    polling (plus PR #522's ad hoc stop-file) is replaced by a small, uniform
    control protocol: a PID file, a command file, and SIGTERM
    escalation
    .

    Closes #53, closes #23, closes #91; resolves the shutdown-detection cost in #35.
    Removes the plumbing PR #522 added as a one-off (that PR merged separately on its
    own merits; this subsumes it).

    Motivation

    The orchestrator launches remote_monitor.py detached over SSH (nohup … &, so
    the child survives the launching SSH exec returning). After that it holds no
    process handle and no pipe — the child is a detached process on a remote machine.
    Today it recovers status by pgrep -f remote_monitor.py, which matches leftover
    crashed instances and needs regex disambiguation (exactly what PR #522 had to
    build). And remote_monitor.py itself imports and blocks on the query client
    just to know when to start/stop sampling, so a new backend phase (ClickHouse
    ingest) meant a whole new execution_mode plus a bespoke stop-file.

    The control protocol (the important part)

    Three plain-filesystem pieces give the orchestrator the full lifecycle —
    start / status / graceful stop / forced stop — with no new dependency (no Redis;
    see #23). All scoped under the already-unique experiment_output_dir, so there
    is never any need to pattern-match which monitor instance is meant.

    1. PID file — "is it alive?"

    <output_dir>/remote_monitor.pid, containing one number: the OS PID of the
    remote_monitor.py process.

    • remote_monitor.py writes it (unconditionally — a stale file from a crashed
      prior run is garbage to overwrite, not a lock) as nearly the first thing on
      startup, and deletes it as the last thing on clean exit.
    • The orchestrator reads it to check liveness via kill(pid, 0) (sends no
      signal; just succeeds if the PID is alive). File absent → not started, or
      already finished and cleaned up. File present but kill(pid,0) fails →
      crashed without cleanup.

    This is what lets us delete PR #522's _remote_monitor_pgrep_pattern /
    re.escape / shlex.quote machinery: the path already identifies the instance.

    2. Command file — "stop now" (extensible to other commands)

    <output_dir>/remote_monitor.cmd, written by the orchestrator, read by
    remote_monitor.py. Small structured command, e.g.
    {"action": "stop", "store": true}.

    • The monitor's sampling loop already wakes every interval_seconds to take a
      sample; on each wake (in the stop_signal wait strategy) it also cheaply
      checks for this file. No new thread, no new timer — it piggybacks the existing
      poll cadence.
    • On seeing action: stop, it breaks the loop, stops profilers (honoring
      store), writes the output JSON, removes the PID file, exits.

    Structured action (vs. a bare marker file) is what makes this extensible:
    future commands (pause, rotate, …) are new enum values, not new plumbing.
    store here consolidates the per-profiler store: bool the profiling code
    already threads around ("keep the flamegraph" vs. "discard"). This is the
    generalized form of PR #522's boolean stop-file.

    3. SIGTERM escalation — the safety net

    The command file is cooperative: it only works if the monitor is actually
    running its poll loop. If it's wedged (hung profiler subprocess, stuck syscall),
    it never reads the file. SIGTERM is the fallback, not the normal stop path.
    MonitorSession.stop():

    1. Write {"action":"stop"} to the command file.
    2. Poll the PID file for up to timeout s for the monitor to exit on its own.
    3. Still alive → read PID from the PID file, send SIGTERM over SSH.
      remote_monitor.py installs a SIGTERM handler so even this path flushes
      output.
    4. Still alive after a short grace period → SIGKILL.

    This is the same graceful→forceful ladder already in
    process_monitor.stop_monitor (terminate() → kill()), applied across the
    SSH boundary using the PID from the PID file.

    Need Mechanism Direction
    Has it started / is it alive / done? PID file + kill(pid,0) monitor → orchestrator
    Stop cleanly (extensible) command file (action, store) orchestrator → monitor
    Stop now, unresponsive SIGTERM → SIGKILL to PID orchestrator → monitor (force)

    What changes

    • remote_monitor.py: driven by one --request '<json>' arg
      (MonitorRequest) instead of ~15 flags; main() is a straight line with no
      execution_mode branches; no longer imports QueryClientService. Writes PID
      file, polls command file, escalates on SIGTERM.
    • RemoteMonitorService → MonitorSession: start / wait_for_start /
      wait_for_finish / stop, replacing RemoteMonitorService plus all six
      methods PR feat(asap-tools): separate clickhouse ingest and query cpu/memory monitoring in experiment run clickhouse #522 added.
    • Wait strategy replaces execution_mode: fixed_seconds (was timed),
      stdin (was interactive), stop_signal (was both prometheus_client and
      feat(asap-tools): separate clickhouse ingest and query cpu/memory monitoring in experiment run clickhouse #522's ingest — same shape once the outer script owns the client).
    • Profilers: one Profiler interface with AsprofProfiler /
      FlamegraphProfiler / PerfProfiler implementations, replacing the three
      copy-pasted start/stop_profiling_*_pids pairs.
    • Client ownership: experiment_run_e2e.py launches/awaits
      QueryClientService itself (as experiment_run_clickhouse.py already does for
      ClickHouseDataLoaderService); remote_monitor.py only brackets it.

    Removed (dead code)

    • profile_query_engine_pid end-to-end (remote_monitor.py,
      prometheus_client_service.py, generate_prometheus_client_compose.py,
      docker-compose.yml.j2, main_prometheus_client.py's
      start_query_engine_profiler thread). Its only live call site hardcodes it to
      None ("unused for Rust QE"); the Rust engine is profiled directly by
      remote_monitor.py via perf record through the separate --profile_query_engine
      flag. (issue Re-architect code to use Redis for control messages #23, second bullet.)

    Testing

    • TODO once implemented — end-to-end experiment_run_e2e and
      experiment_run_clickhouse runs on a CloudLab node; verify monitor output JSON
      is written for each wait strategy, and that stop works via both the command
      file (normal) and the SIGTERM path (kill the loop, confirm escalation).
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions