diff --git a/mkdocs/blog/posts/agentic-orchestration.md b/mkdocs/blog/posts/agentic-orchestration.md
index 15e2c1b8e..42406c77d 100644
--- a/mkdocs/blog/posts/agentic-orchestration.md
+++ b/mkdocs/blog/posts/agentic-orchestration.md
@@ -250,12 +250,11 @@ $ dstack event --within-run train-qwen
```shell
$ dstack metrics train-qwen
- UTILIZATION MEMORY
- cpu ▅▄▄▆▆▆▆▆▆▆▆▆▆▆▆▆▆▅▆▆▆▆▆▆▆▆▆ 91% of 32 ▃▃▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄ 116GB/200GB
+ UTILIZATION MEMORY
+ job=0 cpu ▂▂▁▁▁▁▁▁▁▁▁▁▂▂▂▂▂▂▂▂▂▂▂▂▁▁ 15% ▃▃▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄ 116GB/200GB
+ gpu=0 ▃▃▃▃▃▃▃▃▃▃▃▃▃▄▄▄▄▅▅▄▄▄▃▃▃▃ 43% ▄▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 70GB/80GB
- gpu=0 ▁▂▃▆▆▆▆▆▆▆▆▆▆▆▆▆▆▁▆▆▆▆▆▆▆▆▆ 92% ▄▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 70GB/80GB
-
- 4 Aug 14:10 ┄┄┄┄┄┄┄┄┄┄┄ now 4 Aug 14:10 ┄┄┄┄┄┄┄┄┄┄┄ now
+ 6 Aug 13:14 ┄┄┄┄┄┄┄┄┄┄ now 6 Aug 13:14 ┄┄┄┄┄┄┄┄┄┄ now
```
diff --git a/mkdocs/blog/posts/dstack-metrics.md b/mkdocs/blog/posts/dstack-metrics.md
index eeccdc45b..026a0cb7c 100644
--- a/mkdocs/blog/posts/dstack-metrics.md
+++ b/mkdocs/blog/posts/dstack-metrics.md
@@ -21,15 +21,14 @@ for monitoring container metrics, including GPU usage for `NVIDIA`, `AMD`, and o
```shell
$ dstack metrics llama-70b-sft
- UTILIZATION MEMORY
- cpu ▅▄▄▄▃▃▃▃▃▃▃▃▃▃▃▃▅▅▄▂▃▃▃▃▃▃▃ 39% of 64 ▃▃▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄ 297GB/480GB
-
- gpu=0 ▁▂▃▆▆▆▆▆▆▆▆▆▆▆▆▆▆▁▆▆▆▆▆▆▆▆▆ 89% ▄▅▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 67GB/80GB
- gpu=1 ▁▂▆▆▆▅▆▆▆▆▆▆▆▆▆▆▁▁▅▆▆▅▆▆▆▆▆ 84% ▄▅▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 67GB/80GB
- gpu=2 ▁▂▆▆▆▆▆▆▆▆▆▆▆▆▆▆▁▆▆▆▆▆▆▆▆▆▆ 87% ▄▅▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 67GB/80GB
- gpu=3 ▂▃▆▅▅▅▅▅▅▆▅▅▆▆▆▆▁▅▅▅▅▅▆▅▅▅▅ 82% ▄▅▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 67GB/80GB
-
- 4 Aug 14:10 ┄┄┄┄┄┄┄┄┄┄┄ now 4 Aug 14:10 ┄┄┄┄┄┄┄┄┄┄┄ now
+ UTILIZATION MEMORY
+ job=0 cpu ▁▂▂▂▂▂▂▂▂▂▂▂▂▂▂▂▁▂▂▂▂▂▂▂▂▂ 33% ▃▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄ 296GB/480GB
+ gpu=0 ▁▄▆▆▆▆▆▆▆▆▆▆▆▆▆▆▁▆▆▆▆▆▆▆▆▆ 93% ▃▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 67GB/80GB
+ gpu=1 ▂▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▁▆▆▆▆▆▆▆▆▆ 93% ▃▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 67GB/80GB
+ gpu=2 ▃▆▆▆▆▆▆▆▆▆▆▆▆▆▆▁▆▆▆▆▆▆▆▆▆▆ 93% ▃▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 67GB/80GB
+ gpu=3 ▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▁▆▆▆▆▆▆▆▆▆▆ 93% ▃▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 67GB/80GB
+
+ 6 Aug 13:14 ┄┄┄┄┄┄┄┄┄┄ now 6 Aug 13:14 ┄┄┄┄┄┄┄┄┄┄ now
```
@@ -47,9 +46,6 @@ difference is that `dstack stats` includes GPU VRAM usage and GPU utilization pe
Similar to `kubectl top`, if a run consists of multiple jobs (such as distributed training or an auto-scalable service),
`dstack stats` will display metrics per job.
-> Note, `dstack metrics` now shows one job at a time, like `dstack logs`. Use `--replica` and `--job` to
-> choose it; both default to `0`.
-
!!! info "HTTP API"
In addition to the `dstack stats` CLI commands, metrics can be obtained via the
[`/api/project/{project_name}/metrics/job/{run_name}`](../../docs/reference/http/metrics.md) HTTP endpoint.
diff --git a/mkdocs/docs/concepts/metrics.md b/mkdocs/docs/concepts/metrics.md
index 8e916bd98..bf91901db 100644
--- a/mkdocs/docs/concepts/metrics.md
+++ b/mkdocs/docs/concepts/metrics.md
@@ -19,34 +19,31 @@ This tab displays key CPU, memory, and GPU metrics collected during the last hou
## CLI
As an alternative to the UI, you can track essential metrics via the CLI.
-The `dstack metrics` command charts CPU, memory, and GPU utilization over the last hour of the
-job, with the latest value beside each chart.
+The `dstack metrics` command charts CPU, memory, and GPU utilization over the last hour.
```shell
dstack metrics gentle-mayfly-1
- UTILIZATION MEMORY
- cpu ▅▄▄▄▃▃▃▃▃▃▃▃▃▃▃▃▅▅▄▃▃▃▃▃▃▃▃ 41% of 128 ▃▃▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄ 581GB/960GB
-
- gpu=0 ▁▂▃▆▆▆▆▆▆▆▆▆▆▆▆▆▆▁▆▆▆▆▆▆▆▆▆ 89% ▄▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 71GB/80GB
- gpu=1 ▁▂▆▆▆▅▆▆▆▆▆▆▆▆▆▆▁▁▅▆▆▅▆▆▆▆▆ 84% ▄▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 71GB/80GB
- gpu=2 ▁▂▆▆▆▆▆▆▆▆▆▆▆▆▆▆▁▆▆▆▆▆▆▆▆▆▆ 87% ▄▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 71GB/80GB
- gpu=3 ▂▃▆▅▅▅▅▅▅▆▅▅▆▆▆▆▁▅▅▅▅▅▆▅▅▅▅ 82% ▄▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 71GB/80GB
- gpu=4 ▂▆▆▆▆▆▆▆▆▆▆▆▆▆▆▁▁▆▆▆▆▆▆▆▆▆▆ 90% ▄▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 71GB/80GB
- gpu=5 ▂▆▆▆▆▆▆▆▆▆▆▆▆▆▆▁▆▆▆▆▆▆▆▆▆▆▆ 85% ▄▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 71GB/80GB
- gpu=6 ▃▆▆▅▅▆▆▆▆▆▆▆▆▆▆▁▅▅▅▅▆▆▆▅▅▅▅ 83% ▄▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 71GB/80GB
- gpu=7 ▃▆▆▆▆▆▆▆▆▆▆▆▆▆▁▁▆▆▆▆▆▆▆▆▆▆▆ 88% ▄▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 71GB/80GB
-
- 4 Aug 14:10 ┄┄┄┄┄┄┄┄┄┄┄ now 4 Aug 14:10 ┄┄┄┄┄┄┄┄┄┄┄ now
+ UTILIZATION MEMORY
+ job=0 cpu ▁▂▂▂▂▂▂▂▂▂▂▂▂▂▂▂▁▂▂▂▂▂▂▂▂▂ 33% ▃▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄▄ 579GB/960GB
+ gpu=0 ▁▄▆▆▆▆▆▆▆▆▆▆▆▆▆▆▁▆▆▆▆▆▆▆▆▆ 93% ▄▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 71GB/80GB
+ gpu=1 ▂▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▁▆▆▆▆▆▆▆▆▆ 93% ▄▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 71GB/80GB
+ gpu=2 ▃▆▆▆▆▆▆▆▆▆▆▆▆▆▆▁▆▆▆▆▆▆▆▆▆▆ 93% ▄▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 71GB/80GB
+ gpu=3 ▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▁▆▆▆▆▆▆▆▆▆▆ 93% ▄▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 71GB/80GB
+ gpu=4 ▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 93% ▄▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 71GB/80GB
+ gpu=5 ▆▆▆▆▆▆▆▆▆▆▆▆▆▆▁▆▆▆▆▆▆▆▆▆▆▆ 93% ▄▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 71GB/80GB
+ gpu=6 ▆▆▆▆▆▆▆▆▆▆▆▆▆▆▁▆▆▆▆▆▆▆▆▆▆▆ 93% ▄▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 71GB/80GB
+ gpu=7 ▆▆▆▆▆▆▆▆▆▆▆▆▆▁▆▆▆▆▆▆▆▆▆▆▆▆ 93% ▄▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 71GB/80GB
+
+ 6 Aug 13:14 ┄┄┄┄┄┄┄┄┄┄ now 6 Aug 13:14 ┄┄┄┄┄┄┄┄┄┄ now
```
-Like [`dstack logs`](../reference/cli/dstack/logs.md), the command shows a single job.
-Use `--replica` and `--job` to select one; both default to `0`.
-Pass `-w` to keep the charts updating.
+By default, metrics are shown for all jobs and replicas. Use `--replica` or `--job` to
+show a single one, and `-w` to keep the charts updating.
## Prometheus
diff --git a/mkdocs/docs/guides/migration/slurm.md b/mkdocs/docs/guides/migration/slurm.md
index fbe2dc3c3..9edf863e9 100644
--- a/mkdocs/docs/guides/migration/slurm.md
+++ b/mkdocs/docs/guides/migration/slurm.md
@@ -1476,12 +1476,11 @@ Check real-time metrics:
```shell
$ dstack metrics training-job
- UTILIZATION MEMORY
- cpu ▅▄▄▄▃▃▃▃▃▃▃▃▃▃▃▃▅▅▄▃▃▃▃▃▃▃▃ 45% of 32 ▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁ 16GB/200GB
+ UTILIZATION MEMORY
+ job=0 cpu ▁▂▂▂▂▂▂▂▂▂▂▂▂▂▂▂▁▂▂▂▂▂▂▂▂▂ 33% ▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁▁ 16GB/200GB
+ gpu=0 ▁▄▆▆▆▆▆▆▆▆▆▆▆▆▆▆▁▆▆▆▆▆▆▆▆▆ 93% ▄▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 71GB/80GB
- gpu=0 ▁▂▃▆▆▆▆▆▆▆▆▆▆▆▆▆▆▁▆▆▆▆▆▆▆▆▆ 90% ▄▅▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆▆ 71GB/80GB
-
- 4 Aug 14:10 ┄┄┄┄┄┄┄┄┄┄┄ now 4 Aug 14:10 ┄┄┄┄┄┄┄┄┄┄┄ now
+ 6 Aug 13:14 ┄┄┄┄┄┄┄┄┄┄ now 6 Aug 13:14 ┄┄┄┄┄┄┄┄┄┄ now
```
diff --git a/src/dstack/_internal/cli/commands/metrics.py b/src/dstack/_internal/cli/commands/metrics.py
index 3166640df..d2f3285f0 100644
--- a/src/dstack/_internal/cli/commands/metrics.py
+++ b/src/dstack/_internal/cli/commands/metrics.py
@@ -1,5 +1,6 @@
import argparse
import time
+from typing import Optional
from rich.live import Live
@@ -36,50 +37,45 @@ def _register(self):
)
self._parser.add_argument(
"--replica",
- help="The replica number. Defaults to 0.",
+ help="Show only this replica. By default, all jobs are shown.",
type=int,
- default=0,
)
self._parser.add_argument(
"--job",
- help="The job number inside the replica. Defaults to 0.",
+ help="Show only this job number. By default, all jobs are shown.",
type=int,
- default=0,
)
def _command(self, args: argparse.Namespace):
super()._command(args)
- job, metrics = self._fetch(args)
+ jobs, metrics = self._fetch(args)
if not args.watch:
- console.print(get_metrics_table(job, metrics))
+ console.print(get_metrics_table(jobs, metrics))
return
try:
with Live(console=console, refresh_per_second=LIVE_TABLE_REFRESH_RATE_PER_SEC) as live:
while True:
- live.update(get_metrics_table(job, metrics))
+ live.update(get_metrics_table(jobs, metrics))
time.sleep(WATCH_INTERVAL_SECONDS)
- job, metrics = self._fetch(args)
+ jobs, metrics = self._fetch(args)
except KeyboardInterrupt:
pass
- def _fetch(self, args: argparse.Namespace) -> tuple[Job, JobMetrics]:
+ def _fetch(self, args: argparse.Namespace) -> tuple[list[Job], list[JobMetrics]]:
run = self.api.runs.get(run_name=args.run_name)
if run is None:
raise CLIError(f"Run {args.run_name} not found")
- job = _get_job(run, args.replica, args.job)
- return job, _get_job_metrics(self.api, run, job)
-
-
-def _get_job(run: Run, replica_num: int, job_num: int) -> Job:
- for job in run._run.jobs:
- if job.job_spec.replica_num == replica_num and job.job_spec.job_num == job_num:
- return job
- raise CLIError(
- f"Run {run.name} has no replica={replica_num} job={job_num}."
- " Use --replica and --job to select one."
- )
+ jobs = select_jobs(run._run.jobs, args.replica, args.job)
+ if not jobs:
+ wanted = " ".join(
+ f"{name}={value}"
+ for name, value in (("replica", args.replica), ("job", args.job))
+ if value is not None
+ )
+ raise CLIError(f"Run {args.run_name} has no job matching {wanted}")
+ return jobs, [_get_job_metrics(self.api, run, job) for job in jobs]
def _get_job_metrics(api: Client, run: Run, job: Job) -> JobMetrics:
@@ -92,3 +88,12 @@ def _get_job_metrics(api: Client, run: Run, job: Job) -> JobMetrics:
job_num=job.job_spec.job_num,
limit=MAX_SAMPLES,
)
+
+
+def select_jobs(jobs: list[Job], replica: Optional[int], job_num: Optional[int]) -> list[Job]:
+ return [
+ job
+ for job in jobs
+ if (replica is None or job.job_spec.replica_num == replica)
+ and (job_num is None or job.job_spec.job_num == job_num)
+ ]
diff --git a/src/dstack/_internal/cli/utils/metrics.py b/src/dstack/_internal/cli/utils/metrics.py
index ccf840a99..22ce839a0 100644
--- a/src/dstack/_internal/cli/utils/metrics.py
+++ b/src/dstack/_internal/cli/utils/metrics.py
@@ -1,5 +1,5 @@
-from datetime import datetime
-from typing import Any, List, Optional
+from datetime import datetime, timedelta
+from typing import Any, List, Optional, Sequence
from rich.console import RenderableType
from rich.table import Table
@@ -22,50 +22,139 @@
MIN_SPARK_WIDTH = 10
MAX_SPARK_WIDTH = 80
+RETENTION = timedelta(hours=1)
+"""What the server keeps for a running job, and so the widest window there can be."""
+
AXIS_RULE = "┄"
-_FIXED_COLUMNS = 34
-"""Labels, numbers and padding. Hand-measured against a `589GB/1480GB`-sized label; a
-wider one (a 2000GB host prints `1218GB/2000GB`) overflows and Rich ellipsizes the row
-labels rather than shrinking the sparklines."""
+_FIXED_COLUMNS = 30
+"""Everything but the sparklines and the job label: the `gpu=N` column, both numbers, and
+the table's padding. Hand-measured against a `589GB/1480GB`-sized number; a wider one
+overflows and Rich ellipsizes the labels rather than shrinking the sparklines."""
_SPARKLINE_COLUMNS = 2
-def _spark_width(console_width: int) -> int:
- budget = console_width - _FIXED_COLUMNS
+def _spark_width(console_width: int, label_width: int = 0) -> int:
+ budget = console_width - _FIXED_COLUMNS - label_width
return max(MIN_SPARK_WIDTH, min(MAX_SPARK_WIDTH, budget // _SPARKLINE_COLUMNS))
def get_metrics_table(
- job: Job, metrics: JobMetrics, console_width: Optional[int] = None
+ jobs: Sequence[Job],
+ metrics: Sequence[JobMetrics],
+ console_width: Optional[int] = None,
) -> RenderableType:
- resources = _get_resources(job)
- width = _spark_width(console_width or console.width)
+ labels = job_labels(jobs)
+ label_width = max((len(label) for label in labels), default=0)
+ width = _spark_width(console_width or console.width, label_width)
+ span = _span(metrics)
table = Table(box=None)
- # no header: every cell in this column already reads `cpu` or `gpu=N`
+ # no headers: the cells read `replica=0` and `gpu=1`, which need no naming
+ table.add_column("", no_wrap=True)
table.add_column("", style="secondary", no_wrap=True)
table.add_column("UTILIZATION", no_wrap=True)
table.add_column("MEMORY", no_wrap=True)
+ for index, (job, job_metrics) in enumerate(zip(jobs, metrics)):
+ if index:
+ table.add_row("", "", "", "")
+ _add_job(table, job, job_metrics, width, labels[index], span)
+
+ if span is not None:
+ table.add_row("", "", "", "")
+ # the axis spans the widest chart drawn: a job with fewer samples than cells draws
+ # one cell per sample and cannot fill its share
+ axis = _axis(max(_drawn(m, span, width) for m in metrics), *span)
+ table.add_row("", "", axis, axis)
+ return table
+
+
+def job_labels(jobs: Sequence[Job]) -> List[str]:
+ """`replica=`/`group=` only where they distinguish something, as `dstack ps` does --
+ one replica across four nodes is `job=0..3`, not `replica=0 job=0..3`.
+
+ Unlike `ps`, `job=` is always printed. This table is keyed by job, so every row names
+ one; `replica=` joins it only where there is more than one replica to tell apart.
+ """
+ groups = {job.job_spec.replica_group for job in jobs}
+ show_group = len(groups) > 1
+ show_replica = len({job.job_spec.replica_num for job in jobs}) > 1
+
+ labels, last_group = [], None
+ for job in jobs:
+ parts = []
+ if show_group:
+ # as `ps`: name the group where it changes, and indent the replicas under it
+ group = job.job_spec.replica_group
+ parts.append(f"group={group}" if group != last_group else " " * len(f"group={group}"))
+ last_group = group
+ if show_replica:
+ parts.append(f"replica={job.job_spec.replica_num}")
+ parts.append(f"job={job.job_spec.job_num}")
+ labels.append(" ".join(parts))
+ return labels
+
+
+def _add_job(
+ table: Table,
+ job: Job,
+ metrics: JobMetrics,
+ width: int,
+ label: str,
+ span: Optional[tuple[datetime, datetime]],
+) -> None:
+ resources = _get_resources(job)
+ lead = _lead(metrics, span, width)
+ cells = width - lead
table.add_row(
+ label,
"cpu",
- _cpu_cell(metrics, resources, width),
- _memory_cell(metrics, resources, width),
+ _pad(_cpu_cell(metrics, resources, cells), lead),
+ _pad(_memory_cell(metrics, resources, cells), lead),
)
- table.add_row("", "", "") # host and devices are different things; separate them
for index in range(_gpus_num(metrics, resources)):
table.add_row(
+ "",
f"gpu={index}",
- _gpu_util_cell(metrics, index, width),
- _gpu_memory_cell(metrics, resources, index, width),
+ _pad(_gpu_util_cell(metrics, index, cells), lead),
+ _pad(_gpu_memory_cell(metrics, resources, index, cells), lead),
)
+
+
+def _span(metrics: Sequence[JobMetrics]) -> Optional[tuple[datetime, datetime]]:
+ """The window every job is drawn against: always the full retention hour.
+
+ Fixed rather than fitted to the data, so a row means the same thing in every
+ invocation and across every job. A job younger than the hour fills only its share of
+ the row and the rest is blank -- which is the fact worth seeing about a replica that
+ started two minutes ago.
+ """
+ windows = [w for w in (_window(m) for m in metrics) if w is not None]
+ if not windows:
+ return None
+ latest, earliest = max(w[1] for w in windows), min(w[0] for w in windows)
+ return min(earliest, latest - RETENTION), latest
+
+
+def _lead(metrics: JobMetrics, span: Optional[tuple[datetime, datetime]], width: int) -> int:
+ """Cells before this job's first sample -- time it was not running for."""
window = _window(metrics)
- if window is not None:
- axis = _axis(min(width, _samples_num(metrics)), *window)
- table.add_row("", "", "")
- table.add_row("", axis, axis)
- return table
+ if window is None or span is None:
+ return 0
+ total = (span[1] - span[0]).total_seconds()
+ if total <= 0:
+ return 0
+ return min(width - 1, max(0, round((window[0] - span[0]).total_seconds() / total * width)))
+
+
+def _drawn(metrics: JobMetrics, span: Optional[tuple[datetime, datetime]], width: int) -> int:
+ lead = _lead(metrics, span, width)
+ return lead + min(width - lead, _samples_num(metrics))
+
+
+def _pad(cell: Text, lead: int) -> Text:
+ return cell if lead <= 0 else Text.assemble(Text(" " * lead), cell)
def _cpu_cell(job_metrics: JobMetrics, resources: Optional[Resources], width: int) -> Text:
@@ -75,10 +164,9 @@ def _cpu_cell(job_metrics: JobMetrics, resources: Optional[Resources], width: in
cpus = resources.cpus if resources else None
if cpus:
values = [v / cpus for v in values]
- label = f"{values[-1]:.0f}%"
- if cpus:
- label += f" of {cpus}"
- return _cell(sparkline(values, width, HOST_RAMP), label)
+ # no core count: the value is already normalised to it, and unlike memory there is no
+ # total to give the number meaning
+ return _cell(sparkline(values, width, HOST_RAMP), f"{values[-1]:.0f}%")
def _memory_cell(job_metrics: JobMetrics, resources: Optional[Resources], width: int) -> Text:
diff --git a/src/tests/_internal/cli/commands/test_metrics.py b/src/tests/_internal/cli/commands/test_metrics.py
index 2309a480f..9016edc26 100644
--- a/src/tests/_internal/cli/commands/test_metrics.py
+++ b/src/tests/_internal/cli/commands/test_metrics.py
@@ -2,9 +2,8 @@
import pytest
-from dstack._internal.cli.commands.metrics import _get_job, _get_job_metrics
+from dstack._internal.cli.commands.metrics import _get_job_metrics, select_jobs
from dstack._internal.cli.utils.metrics import MAX_SAMPLES
-from dstack._internal.core.errors import CLIError
from dstack._internal.core.models.metrics import JobMetrics
@@ -22,17 +21,14 @@ def _run(replicas: int = 1, jobs_per_replica: int = 1):
class TestJobSelection:
- def test_defaults_to_the_first_job_of_the_first_replica(self):
- job = _get_job(_run(replicas=3), replica_num=0, job_num=0)
- assert (job.job_spec.replica_num, job.job_spec.job_num) == (0, 0)
-
- def test_selects_by_replica_and_job(self):
- job = _get_job(_run(replicas=3, jobs_per_replica=2), replica_num=2, job_num=1)
- assert (job.job_spec.replica_num, job.job_spec.job_num) == (2, 1)
-
- def test_unknown_job_is_an_error(self):
- with pytest.raises(CLIError, match="replica=7"):
- _get_job(_run(replicas=3), replica_num=7, job_num=0)
+ @pytest.mark.parametrize(
+ "replica,job_num,expected",
+ [(None, None, 4), (0, None, 2), (None, 1, 2), (0, 1, 1), (9, None, 0)],
+ ids=["all", "one-replica", "one-node", "both", "no-match"],
+ )
+ def test_filters(self, replica, job_num, expected):
+ jobs = _run(replicas=2, jobs_per_replica=2)._run.jobs
+ assert len(select_jobs(jobs, replica, job_num)) == expected
class TestMetricsRequest:
diff --git a/src/tests/_internal/cli/utils/test_metrics.py b/src/tests/_internal/cli/utils/test_metrics.py
index 2f9d1e6b4..17a81244b 100644
--- a/src/tests/_internal/cli/utils/test_metrics.py
+++ b/src/tests/_internal/cli/utils/test_metrics.py
@@ -7,7 +7,11 @@
from rich.console import Console
from rich.theme import Theme
-from dstack._internal.cli.utils.metrics import format_memory, get_metrics_table
+from dstack._internal.cli.utils.metrics import (
+ format_memory,
+ get_metrics_table,
+ job_labels,
+)
from dstack._internal.cli.utils.sparkline import SPARKS
from dstack._internal.core.models.metrics import JobMetrics, Metric
@@ -30,6 +34,9 @@ def make_run(
state: str = "running",
gpus: int = 1,
cpus: int = 8,
+ replica: int = 0,
+ job_num: int = 0,
+ group: str = "default",
) -> Tuple[MagicMock, JobMetrics]:
"""A job and its metrics. `state` decides whether the newest sample reads as `now`."""
newest = datetime.now(timezone.utc)
@@ -67,6 +74,8 @@ def series(fn) -> List[float]:
)
job = MagicMock()
+ job.job_spec.replica_num, job.job_spec.job_num = replica, job_num
+ job.job_spec.replica_group = group
submission = MagicMock()
resources = MagicMock()
resources.cpus, resources.memory_mib = cpus, 32 * 1024
@@ -76,7 +85,10 @@ def series(fn) -> List[float]:
return job, JobMetrics(metrics=metrics)
-def render(job, metrics: JobMetrics, width: int = 200, color: bool = False) -> str:
+def render(jobs, metrics, width: int = 200, color: bool = False) -> str:
+ """`jobs`/`metrics` may be a single pair, as most tests use, or whole lists."""
+ if not isinstance(jobs, list):
+ jobs, metrics = [jobs], [metrics]
console = Console(
width=width,
theme=Theme({"secondary": "grey58"}),
@@ -85,7 +97,7 @@ def render(job, metrics: JobMetrics, width: int = 200, color: bool = False) -> s
color_system="truecolor" if color else None,
)
with console.capture() as capture:
- console.print(get_metrics_table(job, metrics, console_width=width))
+ console.print(get_metrics_table(jobs, metrics, console_width=width))
return capture.get()
@@ -154,18 +166,17 @@ def test_no_data_is_not_zero(self):
class TestWindow:
- @pytest.mark.parametrize("samples", [12, 360], ids=["two-minutes", "an-hour"])
@pytest.mark.parametrize("state", ["running", "terminated"])
- def test_draws_only_what_was_measured(self, samples: int, state: str):
- """A young run fills part of the row and the timeline stops with it. Drawn to the
- full width it would claim a span nothing was measured over, and Rich would widen
- the column to fit, pulling MEMORY out of line."""
- job, metrics = make_run("ramp", samples=samples, state=state)
- output = render(job, metrics, width=200)
- drawn = len(bars(row(output, "cpu")))
- assert drawn == min(samples, 80) # 80 is MAX_SPARK_WIDTH
+ def test_an_hour_old_run_fills_the_row(self, state: str):
+ job, metrics = make_run("ramp", samples=360, state=state)
+ assert len(bars(row(render(job, metrics), "job=0"))) == 80 # MAX_SPARK_WIDTH
+
+ def test_a_young_run_fills_only_its_share(self):
+ job, metrics = make_run("ramp", samples=12)
+ output = render(job, metrics, width=120)
+ assert len(bars(row(output, "job=0"))) < 5
axis = lines(output)[-1]
- assert [len(segment) for segment in re.split(r"\s{3,}", axis.strip())] == [drawn, drawn]
+ assert axis.endswith("now")
@pytest.mark.parametrize("state,live", [("running", True), ("terminated", False)])
def test_a_finished_run_cannot_look_live(self, state: str, live: bool):
@@ -176,13 +187,46 @@ def test_a_finished_run_cannot_look_live(self, state: str, live: bool):
assert ":" in axis # a real clock time, not an age
-@pytest.mark.parametrize("width", [80, 100, 140, 190, 240])
-def test_fits_every_terminal_width(width: int):
- """Nothing wraps or gets truncated, on the widest realistic row: eight GPUs."""
- job, metrics = make_run("saturated", gpus=8)
- output = render(job, metrics, width=width)
- assert max(len(line.rstrip()) for line in output.splitlines()) <= width
- assert "…" not in output
+class TestJobs:
+ def test_every_job_is_shown_and_keyed(self):
+ run = [make_run(replica=r, gpus=1) for r in range(3)]
+ output = render([j for j, _ in run], [m for _, m in run])
+ assert [ln.split()[0] for ln in lines(output) if ln.startswith(" replica")] == [
+ "replica=0",
+ "replica=1",
+ "replica=2",
+ ]
+
+ @pytest.mark.parametrize(
+ "jobs,expected",
+ [
+ ([(0, 0, "default")], ["job=0"]),
+ ([(r, 0, "default") for r in range(2)], ["replica=0 job=0", "replica=1 job=0"]),
+ ([(0, n, "default") for n in range(2)], ["job=0", "job=1"]),
+ (
+ [(0, 0, "spot"), (1, 0, "spot"), (2, 0, "on-demand")],
+ [
+ "group=spot replica=0 job=0",
+ "replica=1 job=0", # same group: named once, replicas indented under it
+ "group=on-demand replica=2 job=0",
+ ],
+ ),
+ ],
+ ids=["one-job", "replicas", "nodes", "groups"],
+ )
+ def test_labels_name_only_what_distinguishes(self, jobs, expected):
+ built = [make_run(replica=r, job_num=j, group=g)[0] for r, j, g in jobs]
+ assert [" ".join(label.split()) for label in job_labels(built)] == expected
+
+ def test_a_late_job_starts_where_it_started(self):
+ old, young = make_run(samples=360, replica=0), make_run(samples=12, replica=1)
+ output = render([old[0], young[0]], [old[1], young[1]])
+ first = {
+ label: min(row(output, label).index(g) for g in SPARKS if g in row(output, label))
+ for label in ("replica=0", "replica=1")
+ }
+ assert len(bars(row(output, "replica=0"))) > len(bars(row(output, "replica=1")))
+ assert first["replica=1"] > first["replica=0"] # pushed right by the blank
@pytest.mark.parametrize(