diff --git a/.github/workflows/benchmark-tmpl.yml b/.github/workflows/benchmark-tmpl.yml index 6920b9ea51..0ef227fe1c 100644 --- a/.github/workflows/benchmark-tmpl.yml +++ b/.github/workflows/benchmark-tmpl.yml @@ -2,6 +2,10 @@ name: Template - Benchmark on: workflow_call: secrets: + SRT_STATUS_ENDPOINT: + required: false + SRTCTL_STATUS_TOKEN: + required: false INFERENCEX_OFFICIAL_RO_HF_TOKEN: required: true MODAL_TOKEN_ID: @@ -278,6 +282,13 @@ jobs: submodules: true persist-credentials: false + - name: Download the patched Tachometer build + if: hashFiles('inferencex-e2e/runners/srt-slurm/patches/539.patch') != '' + uses: actions/download-artifact@3e5f45b2cfb9172054b4087a40e8e0b5a5461e7c # v8.0.1 + with: + name: srt-pr539-tachometer + path: srt-streaming-build + - name: Checkout result tooling uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 with: @@ -299,6 +310,9 @@ jobs: - name: Launch job script env: + SRT_STATUS_ENDPOINT: ${{ secrets.SRT_STATUS_ENDPOINT }} # zizmor: ignore[secrets-outside-env] + SRTCTL_STATUS_TOKEN: ${{ secrets.SRTCTL_STATUS_TOKEN }} # zizmor: ignore[secrets-outside-env] + SRT_TACHOMETER_BUILD_DIR: ${{ github.workspace }}/srt-streaming-build RUNNER_NAME: ${{ runner.name }} RUNNER_TYPE: ${{ inputs.runner }} RESULT_FILENAME_BASE: ${{ env.EXP_NAME }}_${{ env.PRECISION }}_${{ env.FRAMEWORK }}_tp${{ env.TP }}-pp${{ env.PP_SIZE }}-dcp${{ env.DCP_SIZE }}-pcp${{ env.PCP_SIZE }}-ep${{ env.EP_SIZE }}-dpa${{ env.DP_ATTENTION }}_disagg-${{ env.DISAGG }}_spec-${{ env.SPEC_DECODING }}_conc${{ env.CONC }}_${{ runner.name }} diff --git a/.github/workflows/run-sweep.yml b/.github/workflows/run-sweep.yml index c3d5c0b3e0..ae58220616 100644 --- a/.github/workflows/run-sweep.yml +++ b/.github/workflows/run-sweep.yml @@ -404,6 +404,35 @@ jobs: ['Reused benchmark source', process.env.REUSE_SHA], ]).write(); + build-srt-tachometer: + name: Build patched Tachometer + needs: setup + if: ${{ !cancelled() && needs.setup.result == 'success' && needs.setup.outputs.reuse-enabled != 'true' }} + runs-on: ubuntu-24.04 + timeout-minutes: 30 + steps: + - uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 + with: + ref: ${{ github.event.pull_request.head.sha || github.sha }} + submodules: true + persist-credentials: false + - name: Apply PR 539 and build its atomic writer + working-directory: inferencex-e2e/utils/srt-slurm + run: | + git apply ../../runners/srt-slurm/patches/539.patch + cargo build --release --locked --bin tachometer-scraper + mkdir -p "$GITHUB_WORKSPACE/srt-streaming-build" + cp target/release/tachometer-scraper "$GITHUB_WORKSPACE/srt-streaming-build/" + git rev-parse HEAD > "$GITHUB_WORKSPACE/srt-streaming-build/base-commit.txt" + cd "$GITHUB_WORKSPACE/srt-streaming-build" + sha256sum tachometer-scraper > tachometer.sha256 + - uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a # v7.0.1 + with: + name: srt-pr539-tachometer + path: srt-streaming-build + if-no-files-found: error + retention-days: 7 + canary-select: name: canary-select needs: setup @@ -462,7 +491,7 @@ jobs: } >> "$GITHUB_OUTPUT" canary-sweep: - needs: canary-select + needs: [canary-select, build-srt-tachometer] if: ${{ needs.canary-select.outputs.canary-config != '' && needs.canary-select.outputs.canary-config != '[]' }} uses: $/.github/workflows/benchmark-tmpl.yml name: canary / @@ -471,6 +500,8 @@ jobs: matrix: config: ${{ fromJson(needs.canary-select.outputs.canary-config) }} secrets: + SRT_STATUS_ENDPOINT: ${{ secrets.SRT_STATUS_ENDPOINT }} + SRTCTL_STATUS_TOKEN: ${{ secrets.SRTCTL_STATUS_TOKEN }} INFERENCEX_OFFICIAL_RO_HF_TOKEN: ${{ secrets.INFERENCEX_OFFICIAL_RO_HF_TOKEN }} MODAL_TOKEN_ID: ${{ secrets.MODAL_TOKEN_ID }} MODAL_TOKEN_SECRET: ${{ secrets.MODAL_TOKEN_SECRET }} @@ -575,11 +606,12 @@ jobs: with: *multi-node-inputs sweep-single-node-1k1k: - needs: [setup, canary-select, canary-sweep, canary-multi-node-sweep] + needs: [setup, canary-select, canary-sweep, canary-multi-node-sweep, build-srt-tachometer] if: >- ${{ !cancelled() && needs.setup.result == 'success' && + needs.build-srt-tachometer.result == 'success' && needs.setup.outputs.reuse-enabled != 'true' && (needs.canary-sweep.result == 'success' || needs.canary-sweep.result == 'skipped') && (needs.canary-multi-node-sweep.result == 'success' || needs.canary-multi-node-sweep.result == 'skipped') && @@ -593,6 +625,8 @@ jobs: matrix: config: ${{ fromJson((needs.canary-sweep.result == 'success' && needs.canary-select.outputs.remaining-search-space-config) || needs.setup.outputs.search-space-config).single_node['1k1k'] }} secrets: + SRT_STATUS_ENDPOINT: ${{ secrets.SRT_STATUS_ENDPOINT }} + SRTCTL_STATUS_TOKEN: ${{ secrets.SRTCTL_STATUS_TOKEN }} INFERENCEX_OFFICIAL_RO_HF_TOKEN: ${{ secrets.INFERENCEX_OFFICIAL_RO_HF_TOKEN }} MODAL_TOKEN_ID: ${{ secrets.MODAL_TOKEN_ID }} MODAL_TOKEN_SECRET: ${{ secrets.MODAL_TOKEN_SECRET }} @@ -608,11 +642,12 @@ jobs: run-eval: ${{ matrix.config.run-eval }} sweep-single-node-8k1k: - needs: [setup, canary-select, canary-sweep, canary-multi-node-sweep] + needs: [setup, canary-select, canary-sweep, canary-multi-node-sweep, build-srt-tachometer] if: >- ${{ !cancelled() && needs.setup.result == 'success' && + needs.build-srt-tachometer.result == 'success' && needs.setup.outputs.reuse-enabled != 'true' && (needs.canary-sweep.result == 'success' || needs.canary-sweep.result == 'skipped') && (needs.canary-multi-node-sweep.result == 'success' || needs.canary-multi-node-sweep.result == 'skipped') && @@ -626,17 +661,20 @@ jobs: matrix: config: ${{ fromJson((needs.canary-sweep.result == 'success' && needs.canary-select.outputs.remaining-search-space-config) || needs.setup.outputs.search-space-config).single_node['8k1k'] }} secrets: + SRT_STATUS_ENDPOINT: ${{ secrets.SRT_STATUS_ENDPOINT }} + SRTCTL_STATUS_TOKEN: ${{ secrets.SRTCTL_STATUS_TOKEN }} INFERENCEX_OFFICIAL_RO_HF_TOKEN: ${{ secrets.INFERENCEX_OFFICIAL_RO_HF_TOKEN }} MODAL_TOKEN_ID: ${{ secrets.MODAL_TOKEN_ID }} MODAL_TOKEN_SECRET: ${{ secrets.MODAL_TOKEN_SECRET }} with: *single-node-inputs sweep-agentic: - needs: [setup, canary-select, canary-sweep, canary-multi-node-sweep] + needs: [setup, canary-select, canary-sweep, canary-multi-node-sweep, build-srt-tachometer] if: >- ${{ !cancelled() && needs.setup.result == 'success' && + needs.build-srt-tachometer.result == 'success' && needs.setup.outputs.reuse-enabled != 'true' && (needs.canary-sweep.result == 'success' || needs.canary-sweep.result == 'skipped') && (needs.canary-multi-node-sweep.result == 'success' || needs.canary-multi-node-sweep.result == 'skipped') && @@ -650,6 +688,8 @@ jobs: matrix: config: ${{ fromJson((needs.canary-sweep.result == 'success' && needs.canary-select.outputs.remaining-search-space-config) || needs.setup.outputs.search-space-config).single_node['agentic'] }} secrets: + SRT_STATUS_ENDPOINT: ${{ secrets.SRT_STATUS_ENDPOINT }} + SRTCTL_STATUS_TOKEN: ${{ secrets.SRTCTL_STATUS_TOKEN }} INFERENCEX_OFFICIAL_RO_HF_TOKEN: ${{ secrets.INFERENCEX_OFFICIAL_RO_HF_TOKEN }} MODAL_TOKEN_ID: ${{ secrets.MODAL_TOKEN_ID }} MODAL_TOKEN_SECRET: ${{ secrets.MODAL_TOKEN_SECRET }} @@ -705,11 +745,12 @@ jobs: agentx-fast: ${{ contains(github.event.pull_request.labels.*.name, 'agentx-fast') }} sweep-evals: - needs: [setup, canary-select, canary-sweep, canary-multi-node-sweep] + needs: [setup, canary-select, canary-sweep, canary-multi-node-sweep, build-srt-tachometer] if: >- ${{ !cancelled() && needs.setup.result == 'success' && + needs.build-srt-tachometer.result == 'success' && needs.setup.outputs.reuse-enabled != 'true' && (needs.canary-sweep.result == 'success' || needs.canary-sweep.result == 'skipped') && (needs.canary-multi-node-sweep.result == 'success' || needs.canary-multi-node-sweep.result == 'skipped') && @@ -723,6 +764,8 @@ jobs: matrix: config: ${{ fromJson(needs.setup.outputs.search-space-config).evals }} secrets: + SRT_STATUS_ENDPOINT: ${{ secrets.SRT_STATUS_ENDPOINT }} + SRTCTL_STATUS_TOKEN: ${{ secrets.SRTCTL_STATUS_TOKEN }} INFERENCEX_OFFICIAL_RO_HF_TOKEN: ${{ secrets.INFERENCEX_OFFICIAL_RO_HF_TOKEN }} MODAL_TOKEN_ID: ${{ secrets.MODAL_TOKEN_ID }} MODAL_TOKEN_SECRET: ${{ secrets.MODAL_TOKEN_SECRET }} @@ -743,11 +786,12 @@ jobs: # dispatched with sweep-agentic's inputs rather than sweep-evals' fixed-seq-len # inputs (isl/osl/max-model-len, which agentic rows don't have). sweep-agentic-evals: - needs: [setup, canary-select, canary-sweep, canary-multi-node-sweep] + needs: [setup, canary-select, canary-sweep, canary-multi-node-sweep, build-srt-tachometer] if: >- ${{ !cancelled() && needs.setup.result == 'success' && + needs.build-srt-tachometer.result == 'success' && needs.setup.outputs.reuse-enabled != 'true' && (needs.canary-sweep.result == 'success' || needs.canary-sweep.result == 'skipped') && (needs.canary-multi-node-sweep.result == 'success' || needs.canary-multi-node-sweep.result == 'skipped') && @@ -761,6 +805,8 @@ jobs: matrix: config: ${{ fromJson(needs.setup.outputs.search-space-config).agentic_evals }} secrets: + SRT_STATUS_ENDPOINT: ${{ secrets.SRT_STATUS_ENDPOINT }} + SRTCTL_STATUS_TOKEN: ${{ secrets.SRTCTL_STATUS_TOKEN }} INFERENCEX_OFFICIAL_RO_HF_TOKEN: ${{ secrets.INFERENCEX_OFFICIAL_RO_HF_TOKEN }} MODAL_TOKEN_ID: ${{ secrets.MODAL_TOKEN_ID }} MODAL_TOKEN_SECRET: ${{ secrets.MODAL_TOKEN_SECRET }} diff --git a/inferencex-e2e/benchmarks/single_node/srt-slurm-recipes/glm5.2/sglang/b300-fp8-mtp/agentic.yaml b/inferencex-e2e/benchmarks/single_node/srt-slurm-recipes/glm5.2/sglang/b300-fp8-mtp/agentic.yaml index cfcbf7d3f1..a94003c38c 100644 --- a/inferencex-e2e/benchmarks/single_node/srt-slurm-recipes/glm5.2/sglang/b300-fp8-mtp/agentic.yaml +++ b/inferencex-e2e/benchmarks/single_node/srt-slurm-recipes/glm5.2/sglang/b300-fp8-mtp/agentic.yaml @@ -18,7 +18,8 @@ base: observability: enabled: false tachometer: - enabled: false + enabled: true + collect_interval_ms: 1000 engine: sglang roles: agg: diff --git a/inferencex-e2e/docs/configuration-procedures.md b/inferencex-e2e/docs/configuration-procedures.md index 701b54cef1..71ea4857c6 100644 --- a/inferencex-e2e/docs/configuration-procedures.md +++ b/inferencex-e2e/docs/configuration-procedures.md @@ -691,3 +691,16 @@ runtime directories stay out of `/workspace`. The MI300X launcher also raises it allocation from 180 to 480 minutes for this checkpoint: the HF cache there is node-local, so the first arm on each node downloads 511 GB before serving. GPU sweep and eval evidence is required before calling either arm validated. + +### Testing SRT raw streaming on B300 + +PR #3591 applies NVIDIA/srt-slurm#539 to the pinned submodule in each job's +checkout. The normal sweep builds the patched Tachometer on a GitHub-hosted runner +once, verifies its checksum and base commit, and installs it before `make setup`; +the released binary does not include the required atomic Arrow writer. + +The existing `glm5.2-fp8-b300-sglang-agentic-mtp` recipe collects Tachometer +metrics every second and sends raw logs/captures every five seconds. The workflow +forwards `SRT_STATUS_ENDPOINT` and `SRTCTL_STATUS_TOKEN` repository secrets into +the B300 profile. Start its normal sweep with `non-canary-full-sweep-enabled` +after the Dash collector is deployed. diff --git a/inferencex-e2e/docs/configuration-procedures_zh.md b/inferencex-e2e/docs/configuration-procedures_zh.md index bc29939872..e2e5597e55 100644 --- a/inferencex-e2e/docs/configuration-procedures_zh.md +++ b/inferencex-e2e/docs/configuration-procedures_zh.md @@ -592,3 +592,15 @@ python -m pytest infx/tests/matrix/ -v 检查点将仓库挂载到 `/ix` 并重写 `RESULT_DIR`,使 AgentX 运行目录不落在 `/workspace` 下。MI300X launcher 还为该检查点将 Slurm 分配时长从 180 分钟提高到 480 分钟:那里的 HF 缓存为节点本地, 每个节点上的首次运行需先下载 511 GB。在获得 GPU sweep 与 eval 证据之前,不得将任一配方视为已验证。 + +### 在 B300 上测试 SRT 原始数据流 + +PR #3591 在每个作业的独立检出中,将 NVIDIA/srt-slurm#539 补丁应用到固定的 +子模块提交。常规 sweep 在 GitHub 托管 runner 上构建一次补丁版 Tachometer, +验证校验和及基础提交,并在 `make setup` 前安装;已发布的二进制不包含所需的 +原子 Arrow 写入实现。 + +现有 `glm5.2-fp8-b300-sglang-agentic-mtp` 配方每秒采集 Tachometer 指标, +每五秒上传原始日志及采集文件。工作流将仓库 Secret `SRT_STATUS_ENDPOINT` 和 +`SRTCTL_STATUS_TOKEN` 传给 B300 配置。部署 Dash 收集 API 后,添加 +`non-canary-full-sweep-enabled` 标签即可启动常规 sweep。 diff --git a/inferencex-e2e/perf-changelog.yaml b/inferencex-e2e/perf-changelog.yaml index cdd9f07f34..a6c4fa949a 100644 --- a/inferencex-e2e/perf-changelog.yaml +++ b/inferencex-e2e/perf-changelog.yaml @@ -9104,3 +9104,11 @@ - "Restore the DCP8 LMCache bands on the native srt-slurm recipe: concurrency 14 and 16 (DSpark 3, ReplaySSM) and 48, 56 and 72 (no draft) run ATOM's in-process lmcache_offload connector through roles.agg.args.extra-kv-connectors (srt-slurm patch 507), with 128 GB/rank up to 48 and 192 GB/rank at 56 and 72. Concurrency 1 and 4 stay GPU-resident." - "No change to the Inferact/Kimi-K3-DSpark draft's precision: online_quant_config still excludes every draft linear (layers.*, context_proj), so its weights and activations stay BF16, and it keeps the target's FP8 KV cache (kv_cache_dtype fp8). FlyDSL FP8 prefill attention applies only to the target, since the draft runs its block pass as decode attention." pr-link: https://github.com/SemiAnalysisAI/InferenceX/pull/3407 + +- config-keys: + - glm5.2-fp8-b300-sglang-agentic-mtp + scenario-type: + - agentic-coding + description: + - "Test SRT Slurm PR 539 raw log streaming and Tachometer capture during the existing B300 GLM-5.2 FP8 SGLang AgentX sweep. Enable 1-second metric collection and 5-second HTTP uploads; serving settings and concurrency points are unchanged." + pr-link: https://github.com/SemiAnalysisAI/InferenceX/pull/3591 diff --git a/inferencex-e2e/runners/slurm_utils.sh b/inferencex-e2e/runners/slurm_utils.sh index 50370edfeb..f5e9956787 100644 --- a/inferencex-e2e/runners/slurm_utils.sh +++ b/inferencex-e2e/runners/slurm_utils.sh @@ -28,6 +28,7 @@ write_srt_cluster_config() { --var SLURM_ACCOUNT "$SLURM_ACCOUNT" --var SLURM_PARTITION "$SLURM_PARTITION" \ --var SRTCTL_ROOT "$SRTCTL_ROOT" --var SQUASH_FILE "$SQUASH_FILE" \ --var NGINX_SQUASH_FILE "$NGINX_SQUASH_FILE" --var IMAGE "$IMAGE" \ + --var SRT_STATUS_ENDPOINT "${SRT_STATUS_ENDPOINT:-}" \ "$@" "${power_args[@]}" } @@ -87,6 +88,15 @@ PYENV if [[ "$uses_power" == "1" ]]; then cp "$GITHUB_WORKSPACE/srt-slurm-sha.txt" "$GITHUB_WORKSPACE/power-producer-sha.txt" || return 1 fi + if [[ -n "${SRT_TACHOMETER_BUILD_DIR:-}" ]]; then + [[ "$(cat "$SRT_TACHOMETER_BUILD_DIR/base-commit.txt")" == "$SRT_SLURM_COMMIT" ]] || return 1 + (cd "$SRT_TACHOMETER_BUILD_DIR" && sha256sum --check tachometer.sha256) || return 1 + mkdir -p bin || return 1 + install -m755 "$SRT_TACHOMETER_BUILD_DIR/tachometer-scraper" bin/tachometer-scraper || return 1 + fi + if [[ -n "${SRT_STATUS_ENDPOINT:-}" ]]; then + check_env_vars SRTCTL_STATUS_TOKEN SRT_TACHOMETER_BUILD_DIR + fi mkdir -p recipes benchmarks/multi_node || return 1 cp -R "$GITHUB_WORKSPACE/benchmarks/multi_node/srt-slurm-recipes/." recipes/ || return 1 # Both CONFIG_FILE spellings currently occur in master configs. @@ -191,6 +201,7 @@ launch_srt_single_node() { --var SRTCTL_ROOT "$SRTCTL_ROOT" --var SQUASH_FILE "$SRT_CONTAINER" \ --var IMAGE "$IMAGE" --var NGINX_SQUASH_FILE nginx:1.27.4 \ --var SRT_DEFAULT_TIME_LIMIT "$SALLOC_TIME_LIMIT" \ + --var SRT_STATUS_ENDPOINT "${SRT_STATUS_ENDPOINT:-}" \ --model "hf:$MODEL" "$SRT_MODEL_PATH" --container "$IMAGE" "$SRT_CONTAINER" \ --mount "$HF_HUB_CACHE_MOUNT" "$HF_HUB_CACHE" --exclusive "$@" run_srt_setup "ARCH=${SRT_SETUP_ARCH:-x86_64}" diff --git a/inferencex-e2e/runners/srt-slurm/b300-dsxe.yaml b/inferencex-e2e/runners/srt-slurm/b300-dsxe.yaml index 06e6434e4a..58f7436c19 100644 --- a/inferencex-e2e/runners/srt-slurm/b300-dsxe.yaml +++ b/inferencex-e2e/runners/srt-slurm/b300-dsxe.yaml @@ -36,3 +36,9 @@ default_sbatch_directives: cpus-per-task: '192' # gpu-16 retains a foreign 1.63 TB tmpfs allocation; rechecked 2026-09-21. exclude: dsxe-sa-b300-prd0-gpu-16 + +reporting: + status: + endpoint: ${SRT_STATUS_ENDPOINT} + token_env: SRTCTL_STATUS_TOKEN + logging-stream-interval: 5 diff --git a/inferencex-e2e/runners/srt-slurm/patches/539.patch b/inferencex-e2e/runners/srt-slurm/patches/539.patch new file mode 100644 index 0000000000..07b2380cec --- /dev/null +++ b/inferencex-e2e/runners/srt-slurm/patches/539.patch @@ -0,0 +1,1736 @@ +diff --git a/docs/schema-reference.md b/docs/schema-reference.md +index 9ee6090d..ada9ca1c 100644 +--- a/docs/schema-reference.md ++++ b/docs/schema-reference.md +@@ -516,6 +516,7 @@ Status reporting configuration. + | `endpoint` | str \| None | `None` | | + | `endpoints` | list[str] \| None | `None` | | + | `token_env` | str \| None | `None` | Name of the environment variable holding the bearer token the reporter sends as ``Authorization: Bearer`` on every request (default SRTCTL_STATUS_TOKEN). Only the variable name belongs in a recipe: the resolved config is written to the lockfile and the log directory, so a literal token there would leak. | ++| `logging-stream-interval` | float \| None | `None` | Seconds between uploads of raw logs and Tachometer captures to every endpoint. Unset disables streaming; lifecycle events are unaffected. | + + ### AIAnalysisConfig + +diff --git a/docs/status-api-spec.md b/docs/status-api-spec.md +index c8ba9a2c..593e6729 100644 +--- a/docs/status-api-spec.md ++++ b/docs/status-api-spec.md +@@ -16,6 +16,8 @@ reporting: + - "https://status.example.com" + # Optional: which environment variable holds the bearer token (default SRTCTL_STATUS_TOKEN) + token_env: SRTCTL_STATUS_TOKEN ++ # Optional: push new log and metric output every N seconds (off when unset) ++ logging-stream-interval: 10 + ``` + + If not configured, status reporting is disabled and jobs run normally. +@@ -253,9 +255,71 @@ Incremental event feed for one job. Events carry a monotonically increasing `id` + + Same as above across every job, with an optional `job_id` filter. This is the feed for dashboards and agents that want to react to job transitions without polling each job. + ++### POST /api/jobs/{job_id}/logs ++ ++Streaming is opt-in with `reporting.status.logging-stream-interval`. The sweep ++uploads new bytes from `.out`, `.err`, `.log`, `.csv` and `.jsonl` files, using ++persistent HTTP connections and chunks of at most 1 MiB. Shutdown gives the ++worker up to five seconds for a final best-effort flush. ++ ++```http ++POST /api/jobs/12345/logs?file=worker.out&offset=4096&cluster=b200 ++Content-Type: application/octet-stream ++Authorization: Bearer ++ ++ ++``` ++ ++`file` is a relative path and `offset` is the source byte position. The API ++stores the bytes and decodes UTF-8 when read, including characters split across ++chunks. `final=1` marks EOF at shutdown; an empty body can finalize bytes already ++sent. Response: `{"job_id": "12345", "stored": 1}` (`stored: 0` for a resend). ++Exact retries are idempotent; conflicting overlaps return 409. The legacy JSON ++`{"chunks": [{"file", "offset", "size", "data"}], "metadata": {"cluster": "b200"}}` ++format is also accepted. ++ ++Files must be append-only. Uploader offsets are in memory; delivery across ++uploader restarts is not guaranteed. Logs remain on disk if uploads fail. ++ ++### POST /api/jobs/{job_id}/captures ++ ++Tachometer captures are sent unchanged from `tachometer/local`: ++ ++```http ++POST /api/jobs/12345/captures?file=tachometer/local/current.arrow&generation=&offset=0&total=8192&cluster=b200 ++Content-Type: application/octet-stream ++Authorization: Bearer ++ ++ ++``` ++ ++`generation` identifies one file version; `total` is its byte length. Chunks are ++sequential and exact retries are idempotent. The response is ++`{"job_id": "12345", "next_offset": 8192, "complete": true}` once the complete ++capture has been processed. Incomplete generations never produce metric rows. ++ ++The collector decodes captures, removes observations repeated by snapshots or ++compaction, and exposes the result as `tachometer_rows.jsonl` through the log ++read API. The cluster does no decoding, row hashing, JSON conversion, or local ++indexing for streaming. Tachometer publishes Arrow snapshots atomically so an ++open upload can finish even when the next snapshot replaces the file. ++Unchanged captures are skipped; changed Arrow snapshots are uploaded in full, ++so polling still uses disk bandwidth and network bandwidth. ++ ++When configured, `cluster` accompanies every raw upload. Shared collectors must ++use `(cluster, job_id)` for identity; the built-in collector retains its existing ++job-ID scope. Custom collectors must implement both binary upload routes before ++enabling streaming. These routes use the existing write bearer token. ++ ++### GET /api/jobs/{job_id}/logs ++ ++Without `file`: `{"job_id": "12345", "files": [{"file": "...", "size": 4608, "updated_at": "..."}]}`, `size` being the bytes received so far. ++ ++With `file` (and optional `offset`, default 0): the contiguous content from `offset`, up to about 1 MiB, as `{"job_id", "file", "offset", "next_offset", "data"}`. Tail a file by polling with `offset = next_offset`. ++ + ### DELETE /api/jobs/{job_id} + +-Remove a job and its events. `200 {"deleted": true, "job_id": ...}` or `404`. ++Remove a job, its events and its streamed logs. `200 {"deleted": true, "job_id": ...}` or `404`. + + ### GET /api/health + +diff --git a/src/srtctl/cli/do_sweep.py b/src/srtctl/cli/do_sweep.py +index ceba6ad6..d7538497 100644 +--- a/src/srtctl/cli/do_sweep.py ++++ b/src/srtctl/cli/do_sweep.py +@@ -45,7 +45,7 @@ from srtctl.core.resource_snapshot import record_resource_snapshot + from srtctl.core.runtime import RuntimeContext + from srtctl.core.schema import SrtConfig + from srtctl.core.slurm import get_slurm_job_id, start_srun_process +-from srtctl.core.status import JobStage, JobStatus, StatusReporter ++from srtctl.core.status import JobStage, JobStatus, LogStreamer, StatusReporter + from srtctl.core.topology import Endpoint, NodePortAllocator, Process, allocate_endpoints_het + from srtctl.logging_utils import setup_logging + from srtctl.ports import ( +@@ -646,6 +646,17 @@ class SweepOrchestrator( + + exit_code = 1 + ++ # Live log/metric streaming to the status API (reporting.status.logging-stream-interval) ++ observability = self.config.observability ++ tachometer_dir = ( ++ self.runtime.log_dir / observability.tachometer.storage_subdir / "local" ++ if observability.tachometer_enabled ++ else None ++ ) ++ log_streamer = LogStreamer.from_config(self.config.reporting, reporter, self.runtime.log_dir, tachometer_dir) ++ if log_streamer is not None: ++ log_streamer.start() ++ + try: + # Stage 0: Bare-host node setup (GPU clocks, kernel modules). Runs + # before anything containerized so workers see the prepared node. +@@ -778,6 +789,8 @@ class SweepOrchestrator( + # push logs_url to the status API. Runs before report_completed so + # the final PUT can reassert the artifact pointer. + self.run_postprocess(exit_code, reporter=reporter) ++ if log_streamer is not None: ++ log_streamer.stop() + reporter.report_completed( + exit_code, + logs_url=getattr(self, "_last_logs_url", None), +diff --git a/src/srtctl/contract/__init__.py b/src/srtctl/contract/__init__.py +index 929f15da..d56185f0 100644 +--- a/src/srtctl/contract/__init__.py ++++ b/src/srtctl/contract/__init__.py +@@ -17,15 +17,18 @@ Usage (server, e.g. srtctl.status_server): + """ + + from srtctl.contract.enums import JobStage, JobStatus +-from srtctl.contract.requests import JobCreatePayload, JobUpdatePayload ++from srtctl.contract.requests import JobCreatePayload, JobUpdatePayload, LogAppendPayload, LogChunk + from srtctl.contract.responses import ( + EventFeedResponse, + JobDetail, + JobEventListResponse, + JobEventRecord, + JobListResponse, ++ JobLogFilesResponse, ++ JobLogResponse, + JobResponse, + JobSummary, ++ LogFileSummary, + ) + + __all__ = [ +@@ -35,9 +38,14 @@ __all__ = [ + "JobEventListResponse", + "JobEventRecord", + "JobListResponse", ++ "JobLogFilesResponse", ++ "JobLogResponse", + "JobResponse", + "JobStage", + "JobStatus", + "JobSummary", + "JobUpdatePayload", ++ "LogAppendPayload", ++ "LogChunk", ++ "LogFileSummary", + ] +diff --git a/src/srtctl/contract/requests.py b/src/srtctl/contract/requests.py +index eaa5b58e..ba8d40b5 100644 +--- a/src/srtctl/contract/requests.py ++++ b/src/srtctl/contract/requests.py +@@ -3,7 +3,7 @@ + + """Request payload models for the Status API contract.""" + +-from pydantic import BaseModel, Field ++from pydantic import BaseModel, Field, field_validator, model_validator + + + class JobCreatePayload(BaseModel): +@@ -31,3 +31,32 @@ class JobUpdatePayload(BaseModel): + benchmark_results: dict | None = Field(None, description="Parsed benchmark results") + artifacts: dict | None = Field(None, description="Collector-side artifact pointers to merge") + metadata: dict | None = Field(None, description="Additional metadata to merge") ++ ++ ++class LogChunk(BaseModel): ++ """New bytes of one file under the run's log directory.""" ++ ++ file: str = Field(..., min_length=1, max_length=4096, description="Path relative to the run's log directory") ++ offset: int = Field(..., ge=0, le=(1 << 63) - 1, description="Byte offset of the chunk in the file") ++ size: int = Field(..., gt=0, le=1 << 20, description="Byte length of the chunk in the file") ++ data: str = Field(..., min_length=1, description="The chunk decoded as UTF-8 (invalid bytes replaced)") ++ ++ @field_validator("file") ++ @classmethod ++ def relative_file(cls, value: str) -> str: ++ if "\x00" in value or any(part in ("", ".", "..") for part in value.split("/")): ++ raise ValueError("file must be a relative path without empty, '.' or '..' components") ++ return value ++ ++ @model_validator(mode="after") ++ def bounded_range(self) -> "LogChunk": ++ if self.offset + self.size > (1 << 63) - 1: ++ raise ValueError("chunk end exceeds the maximum byte offset") ++ return self ++ ++ ++class LogAppendPayload(BaseModel): ++ """Payload for POST /api/jobs/{job_id}/logs.""" ++ ++ chunks: list[LogChunk] = Field(..., min_length=1, description="Chunks to store; an exact resend is ignored") ++ metadata: dict | None = Field(None, description="Additional metadata, including the source cluster") +diff --git a/src/srtctl/contract/responses.py b/src/srtctl/contract/responses.py +index 22c239f1..c618cec0 100644 +--- a/src/srtctl/contract/responses.py ++++ b/src/srtctl/contract/responses.py +@@ -83,3 +83,31 @@ class EventFeedResponse(BaseModel): + + events: list[JobEventRecord] + next_cursor: int | None = None ++ ++ ++class LogFileSummary(BaseModel): ++ """One streamed file: bytes received so far and when the last chunk arrived.""" ++ ++ file: str ++ size: int ++ updated_at: str ++ ++ ++class JobLogFilesResponse(BaseModel): ++ """GET /api/jobs/{job_id}/logs: every streamed file of a job.""" ++ ++ job_id: str ++ files: list[LogFileSummary] ++ ++ ++class JobLogResponse(BaseModel): ++ """GET /api/jobs/{job_id}/logs?file=...: contiguous content from ``offset``. ++ ++ ``next_offset`` is the ``offset`` to pass on the next poll. ++ """ ++ ++ job_id: str ++ file: str ++ offset: int ++ next_offset: int ++ data: str +diff --git a/src/srtctl/core/schema.py b/src/srtctl/core/schema.py +index aba6de72..651efbfb 100755 +--- a/src/srtctl/core/schema.py ++++ b/src/srtctl/core/schema.py +@@ -101,9 +101,23 @@ class ReportingStatusConfig: + # variable name belongs in a recipe: the resolved config is written to the lockfile + # and the log directory, so a literal token there would leak. + token_env: str | None = None ++ # Seconds between uploads of raw logs and Tachometer captures to every endpoint. ++ # Unset disables streaming; lifecycle events are unaffected. ++ logging_stream_interval: float | None = field( ++ default=None, ++ metadata={ ++ "marshmallow_field": fields.Float(data_key="logging-stream-interval", load_default=None, allow_none=True) ++ }, ++ ) + + Schema: ClassVar[type[Schema]] = Schema + ++ def __post_init__(self) -> None: ++ if self.logging_stream_interval is not None and not _is_finite_positive(self.logging_stream_interval): ++ raise ValidationError( ++ f"reporting.status.logging-stream-interval must be positive, got {self.logging_stream_interval!r}" ++ ) ++ + + @dataclass(frozen=True) + class ReportingConfig: +diff --git a/src/srtctl/core/status.py b/src/srtctl/core/status.py +index f43c504e..44bfdb36 100644 +--- a/src/srtctl/core/status.py ++++ b/src/srtctl/core/status.py +@@ -27,13 +27,23 @@ Configuration (in srtslurm.yaml or recipe YAML): + endpoint: "https://status.example.com" + endpoints: + - "https://status2.example.com" ++ ++ # Also push new log and metric output every 10 seconds ++ reporting: ++ status: ++ endpoint: "https://status.example.com" ++ logging-stream-interval: 10 + """ + + import logging + import os ++import threading ++import time ++import uuid + from dataclasses import dataclass + from datetime import datetime, timezone +-from typing import TYPE_CHECKING ++from pathlib import Path ++from typing import TYPE_CHECKING, BinaryIO + + import requests + +@@ -427,3 +437,227 @@ def create_job_record( + break + + return any_success ++ ++ ++# Append-only logs and structured output; contents are interpreted by the API. ++STREAM_SUFFIXES = (".out", ".err", ".log", ".csv", ".jsonl") ++TACHOMETER_STREAM_FILE = "tachometer_rows.jsonl" # generated by the collector ++STREAM_REQUEST_BYTES = 1 << 20 ++STREAM_SHUTDOWN_SECONDS = 5.0 ++ ++ ++def _file_version(stat: os.stat_result) -> tuple[int, int, int]: ++ return stat.st_ino, stat.st_size, stat.st_mtime_ns ++ ++ ++@dataclass ++class _CaptureUpload: ++ source: BinaryIO ++ version: tuple[int, int, int] ++ generation: str ++ offset: int = 0 ++ pending: bytes | None = None ++ ++ ++class LogStreamer: ++ """Upload raw log deltas and binary captures; the collector does all decoding. ++ ++ Memory is bounded to one chunk per pending file. Capture file descriptors ++ survive compaction unlinking them. Rewritten snapshots get a new generation ++ and are never combined with bytes from an earlier snapshot. ++ """ ++ ++ def __init__(self, reporter: StatusReporter, log_dir: Path, interval: float, tachometer_dir: Path | None = None): ++ self.reporter = reporter ++ self.log_dir = log_dir ++ self.interval = interval ++ self.tachometer_dir = tachometer_dir ++ self._offsets: dict[tuple[str, str], int] = {} ++ self._pending: dict[tuple[str, str], bytes] = {} ++ self._finalized: set[tuple[str, str]] = set() ++ self._captures: dict[tuple[str, str], _CaptureUpload] = {} ++ self._capture_done: dict[tuple[str, str], tuple[int, int, int]] = {} ++ self._cluster = _cluster_setting() ++ self._session = requests.Session() ++ self._shutdown_deadline: float | None = None ++ self._stop = threading.Event() ++ self._thread = threading.Thread(target=self._run, name="status-log-stream", daemon=True) ++ ++ @classmethod ++ def from_config( ++ cls, ++ reporting: "ReportingConfig | None", ++ reporter: StatusReporter, ++ log_dir: Path, ++ tachometer_dir: Path | None = None, ++ ) -> "LogStreamer | None": ++ interval = reporting.status.logging_stream_interval if reporting and reporting.status else None ++ if not reporter.enabled or interval is None: ++ return None ++ return cls(reporter, log_dir, interval, tachometer_dir) ++ ++ def start(self) -> None: ++ self._thread.start() ++ ++ def stop(self) -> None: ++ """Give the worker a bounded final flush; telemetry cannot delay completion indefinitely.""" ++ self._shutdown_deadline = time.monotonic() + STREAM_SHUTDOWN_SECONDS ++ self._stop.set() ++ if self._thread.ident is not None: ++ self._thread.join(timeout=STREAM_SHUTDOWN_SECONDS) ++ if self._thread.is_alive(): ++ logger.warning( ++ "Log streaming did not finish within %.1fs; local logs are retained", STREAM_SHUTDOWN_SECONDS ++ ) ++ ++ def _run(self) -> None: ++ try: ++ while not self._stop.wait(self.interval): ++ self.flush() ++ self.flush() ++ finally: ++ for upload in self._captures.values(): ++ upload.source.close() ++ self._captures.clear() ++ self._session.close() ++ ++ def _expired(self) -> bool: ++ return self._shutdown_deadline is not None and time.monotonic() >= self._shutdown_deadline ++ ++ def flush(self) -> None: ++ """Upload finite file snapshots, retaining failed chunks verbatim for retry.""" ++ try: ++ if self._expired(): ++ return ++ files = sorted(p for p in self.log_dir.rglob("*") if p.suffix in STREAM_SUFFIXES and p.is_file()) ++ for endpoint in self.reporter.api_endpoints: ++ if self._flush_logs(endpoint, files): ++ self._flush_captures(endpoint) ++ except Exception: ++ logger.warning("Log streaming flush failed; will retry", exc_info=True) ++ ++ def _flush_logs(self, endpoint: str, files: list[Path]) -> bool: ++ for path in files: ++ if self._expired(): ++ return False ++ rel = path.relative_to(self.log_dir).as_posix() ++ key = (endpoint, rel) ++ try: ++ with path.open("rb") as source: ++ end = os.fstat(source.fileno()).st_size ++ while not self._expired(): ++ offset = self._offsets.get(key, 0) ++ data = self._pending.get(key) ++ if data is None: ++ if offset >= end: ++ if self._stop.is_set() and key not in self._finalized: ++ ack = self._post( ++ endpoint, "logs", {"file": rel, "offset": str(offset), "final": "1"}, b"" ++ ) ++ if ack is None or ack.get("stored") not in (0, 1): ++ return False ++ self._finalized.add(key) ++ break ++ source.seek(offset) ++ data = source.read(min(STREAM_REQUEST_BYTES, end - offset)) ++ if not data: ++ break ++ self._pending[key] = data ++ params = {"file": rel, "offset": str(offset)} ++ if self._stop.is_set() and offset + len(data) == end: ++ params["final"] = "1" ++ ack = self._post(endpoint, "logs", params, data) ++ if ack is None or ack.get("stored") not in (0, 1): ++ return False ++ self._offsets[key] = offset + len(data) ++ del self._pending[key] ++ if params.get("final") == "1": ++ self._finalized.add(key) ++ except OSError as exc: ++ logger.debug("Log stream skipped %s: %s", rel, exc) ++ return True ++ ++ def _flush_captures(self, endpoint: str) -> None: ++ if self.tachometer_dir is None: ++ return ++ # File discovery does not import Arrow or inspect any metric rows. ++ files = sorted(p for p in self.tachometer_dir.glob("*") if p.suffix in (".arrow", ".parquet")) ++ paths = {p.relative_to(self.log_dir).as_posix(): p for p in files} ++ for ep, rel in self._captures: ++ if ep == endpoint: ++ paths.setdefault(rel, self.log_dir / rel) ++ for rel, path in paths.items(): ++ if self._expired(): ++ return ++ key = (endpoint, rel) ++ upload = self._captures.get(key) ++ try: ++ if upload is None: ++ version = _file_version(path.stat()) ++ if version == self._capture_done.get(key) or version[1] == 0: ++ continue ++ source = path.open("rb") ++ upload = _CaptureUpload(source, _file_version(os.fstat(source.fileno())), uuid.uuid4().hex) ++ self._captures[key] = upload ++ while upload.offset < upload.version[1] and not self._expired(): ++ if upload.pending is None: ++ upload.pending = upload.source.read( ++ min(STREAM_REQUEST_BYTES, upload.version[1] - upload.offset) ++ ) ++ # current.arrow is truncated in place. Reject a changed source ++ # before sending the last chunk that would commit a mixed file. ++ if _file_version(os.fstat(upload.source.fileno())) != upload.version or not upload.pending: ++ upload.source.close() ++ del self._captures[key] ++ break ++ params = { ++ "file": rel, ++ "generation": upload.generation, ++ "offset": str(upload.offset), ++ "total": str(upload.version[1]), ++ } ++ ack = self._post(endpoint, "captures", params, upload.pending) ++ next_offset = upload.offset + len(upload.pending) ++ if ( ++ ack is None ++ or ack.get("next_offset") != next_offset ++ or ack.get("complete") is not (next_offset == upload.version[1]) ++ ): ++ return ++ upload.offset = next_offset ++ upload.pending = None ++ if upload.offset == upload.version[1]: ++ self._capture_done[key] = upload.version ++ upload.source.close() ++ del self._captures[key] ++ except OSError as exc: ++ logger.debug("Capture upload skipped %s: %s", rel, exc) ++ if key in self._captures: ++ self._captures.pop(key).source.close() ++ ++ def _post(self, endpoint: str, route: str, params: dict[str, str], data: bytes) -> dict | None: ++ if self._cluster: ++ params["cluster"] = self._cluster ++ timeout = self.reporter.timeout ++ if self._shutdown_deadline is not None: ++ timeout = min(timeout, self._shutdown_deadline - time.monotonic()) ++ if timeout <= 0: ++ return None ++ try: ++ response = self._session.post( ++ f"{endpoint}/api/jobs/{self.reporter.job_id}/{route}", ++ params=params, ++ data=data, ++ headers={**_auth_headers(self.reporter.token_env), "Content-Type": "application/octet-stream"}, ++ timeout=timeout, ++ allow_redirects=False, ++ ) ++ if response.status_code != 200: ++ _log_rejection("Log stream", endpoint, response.status_code, self.reporter.token_env) ++ return None ++ ack = response.json() ++ if isinstance(ack, dict) and ack.get("job_id") == self.reporter.job_id: ++ return ack ++ except (requests.exceptions.RequestException, ValueError) as exc: ++ logger.debug("Log stream to %s failed: %s", endpoint, exc) ++ return None +diff --git a/src/srtctl/status_server/server.py b/src/srtctl/status_server/server.py +index 6f795727..f96e787d 100644 +--- a/src/srtctl/status_server/server.py ++++ b/src/srtctl/status_server/server.py +@@ -56,13 +56,16 @@ from srtctl.contract import ( + JobDetail, + JobEventListResponse, + JobListResponse, ++ JobLogFilesResponse, ++ JobLogResponse, + JobResponse, + JobStage, + JobStatus, + JobSummary, + JobUpdatePayload, ++ LogAppendPayload, + ) +-from srtctl.status_server.store import StatusStore ++from srtctl.status_server.store import LogChunkConflict, StatusStore + + logger = logging.getLogger(__name__) + +@@ -71,9 +74,7 @@ DEFAULT_PORT = 8080 + DEFAULT_DB_PATH = Path("~/.local/state/srtctl/status.db") + DEFAULT_TOKEN_ENV = "SRTCTL_STATUS_TOKEN" + DEFAULT_READ_TOKEN_ENV = "SRTCTL_STATUS_READ_TOKEN" +-# The largest legitimate body is the started-metadata PUT, a few KB. Anything +-# beyond this is rejected before it is read so a public endpoint cannot be +-# used to fill the disk. ++# Bound each raw streaming chunk and JSON request before reading its body. + MAX_BODY_BYTES = 1 << 20 + HEALTH_PATH = "/api/health" + # The single-page UI. It is static and reveals nothing, so it is served without +@@ -83,6 +84,8 @@ UI_PATHS = frozenset({"/", "/index.html"}) + + _JOB_ROUTE = re.compile(r"^/api/jobs/(?P[^/]+)$") + _JOB_EVENTS_ROUTE = re.compile(r"^/api/jobs/(?P[^/]+)/events$") ++_JOB_CAPTURES_ROUTE = re.compile(r"^/api/jobs/(?P[^/]+)/captures$") ++_JOB_LOGS_ROUTE = re.compile(r"^/api/jobs/(?P[^/]+)/logs$") + + Response = tuple[HTTPStatus, dict[str, Any]] + +@@ -251,6 +254,11 @@ def route(store: StatusStore, method: str, raw_path: str, body: dict[str, Any] | + return _event_feed(store, query) + if (match := _JOB_EVENTS_ROUTE.match(path)) and method == "GET": + return _job_events(store, match["job_id"], query) ++ if match := _JOB_LOGS_ROUTE.match(path): ++ if method == "POST": ++ return _append_logs(store, match["job_id"], body) ++ if method == "GET": ++ return _job_logs(store, match["job_id"], query) + if match := _JOB_ROUTE.match(path): + if method == "GET": + return _get_job(store, match["job_id"]) +@@ -325,6 +333,63 @@ def _job_events(store: StatusStore, job_id: str, query: dict[str, str]) -> Respo + return HTTPStatus.OK, response.model_dump() + + ++def _append_logs(store: StatusStore, job_id: str, body: dict[str, Any] | None) -> Response: ++ payload = LogAppendPayload.model_validate(body or {}) ++ try: ++ stored = store.append_logs(job_id, [chunk.model_dump() for chunk in payload.chunks]) ++ except LogChunkConflict as exc: ++ raise ApiError(HTTPStatus.CONFLICT, str(exc)) from exc ++ return HTTPStatus.OK, {"job_id": job_id, "stored": stored} ++ ++ ++def _raw_upload(store: StatusStore, path: str, raw: bytes) -> Response: ++ url = urlparse(path) ++ query = {key: values[-1] for key, values in parse_qs(url.query).items()} ++ file = query.get("file", "") ++ if not file or len(file) > 4096 or "\x00" in file or any(part in ("", ".", "..") for part in file.split("/")): ++ raise ApiError(HTTPStatus.UNPROCESSABLE_ENTITY, "file must be a relative path without traversal") ++ if len(query.get("cluster", "")) > 256: ++ raise ApiError(HTTPStatus.UNPROCESSABLE_ENTITY, "cluster is too long") ++ offset = _int_param(query, "offset", 0, minimum=0, maximum=(1 << 63) - 1) ++ if offset + len(raw) > (1 << 63) - 1: ++ raise ApiError(HTTPStatus.UNPROCESSABLE_ENTITY, "chunk exceeds maximum offset") ++ normalized = url.path.rstrip("/") ++ try: ++ if match := _JOB_LOGS_ROUTE.fullmatch(normalized): ++ final = _int_param(query, "final", 0, minimum=0, maximum=1) ++ if not raw and not final: ++ raise ApiError(HTTPStatus.UNPROCESSABLE_ENTITY, "Empty log chunk requires final=1") ++ stored = store.append_raw_log(match["job_id"], file, offset, raw, final=bool(final)) ++ return HTTPStatus.OK, {"job_id": match["job_id"], "stored": stored} ++ if match := _JOB_CAPTURES_ROUTE.fullmatch(normalized): ++ generation = query.get("generation", "") ++ if not re.fullmatch(r"[0-9a-f]{32}", generation): ++ raise ApiError(HTTPStatus.UNPROCESSABLE_ENTITY, "generation must be 32 lowercase hex characters") ++ if Path(file).suffix not in (".arrow", ".parquet"): ++ raise ApiError(HTTPStatus.UNPROCESSABLE_ENTITY, "capture must be an Arrow or Parquet file") ++ total = _int_param(query, "total", 0, minimum=1, maximum=(1 << 63) - 1) ++ if not raw or offset + len(raw) > total: ++ raise ApiError(HTTPStatus.UNPROCESSABLE_ENTITY, "capture chunk must be nonempty and within total") ++ return HTTPStatus.OK, store.append_capture(match["job_id"], file, generation, offset, total, raw) ++ except LogChunkConflict as exc: ++ raise ApiError(HTTPStatus.CONFLICT, str(exc)) from exc ++ raise ApiError(HTTPStatus.NOT_FOUND, "No raw upload route") ++ ++ ++def _job_logs(store: StatusStore, job_id: str, query: dict[str, str]) -> Response: ++ """File list without ``file``; with it, content from ``offset`` (poll with ``offset = next_offset``).""" ++ file = query.get("file") ++ if file is None: ++ files = store.list_log_files(job_id) ++ if not files and store.get_job(job_id) is None: ++ raise ApiError(HTTPStatus.NOT_FOUND, "Job not found") ++ return HTTPStatus.OK, JobLogFilesResponse(job_id=job_id, files=files).model_dump() ++ offset = _int_param(query, "offset", 0, minimum=0) ++ data, next_offset = store.read_log(job_id, file, offset=offset) ++ response = JobLogResponse(job_id=job_id, file=file, offset=offset, next_offset=next_offset, data=data) ++ return HTTPStatus.OK, response.model_dump() ++ ++ + def _event_feed(store: StatusStore, query: dict[str, str]) -> Response: + after = _int_param(query, "after", 0, minimum=0) + limit = _int_param(query, "limit", 100, minimum=1, maximum=1000) +@@ -434,7 +499,10 @@ def _handler_class(store: StatusStore, auth: AuthPolicy, cors: CorsPolicy) -> ty + self._send(HTTPStatus.OK, page, "text/html; charset=utf-8", cors_headers, head_only) + return + auth.check(effective, path, self.headers.get("Authorization")) +- status, body = route(store, effective, self.path, _parse_json(raw)) ++ if effective == "POST" and self.headers.get_content_type() == "application/octet-stream": ++ status, body = _raw_upload(store, self.path, raw or b"") ++ else: ++ status, body = route(store, effective, self.path, _parse_json(raw)) + except ApiError as exc: + if exc.status in (HTTPStatus.UNAUTHORIZED, HTTPStatus.FORBIDDEN): + logger.info("%s %s %s from %s", exc.status.value, method, self.path, self.address_string()) +@@ -455,13 +523,20 @@ def _handler_class(store: StatusStore, auth: AuthPolicy, cors: CorsPolicy) -> ty + except ValueError: + self.close_connection = True + raise ApiError(HTTPStatus.BAD_REQUEST, "Content-Length must be an integer") from None ++ if length < 0: ++ self.close_connection = True ++ raise ApiError(HTTPStatus.BAD_REQUEST, "Content-Length must be nonnegative") + if length > MAX_BODY_BYTES: + # Not read, so the connection cannot be reused for a keep-alive request. + self.close_connection = True + raise ApiError(HTTPStatus.REQUEST_ENTITY_TOO_LARGE, f"Body larger than {MAX_BODY_BYTES} bytes") + if length == 0: + return None +- return self.rfile.read(length) ++ data = self.rfile.read(length) ++ if len(data) != length: ++ self.close_connection = True ++ raise ApiError(HTTPStatus.BAD_REQUEST, "Incomplete request body") ++ return data + + def _send( + self, status: HTTPStatus, data: bytes, content_type: str, headers: dict[str, str], head_only: bool +diff --git a/src/srtctl/status_server/store.py b/src/srtctl/status_server/store.py +index 3bd59d55..838a0f95 100644 +--- a/src/srtctl/status_server/store.py ++++ b/src/srtctl/status_server/store.py +@@ -5,6 +5,8 @@ + + from __future__ import annotations + ++import codecs ++import hashlib + import json + import sqlite3 + from collections.abc import Iterator +@@ -12,6 +14,7 @@ from contextlib import contextmanager + from dataclasses import dataclass + from datetime import datetime, timezone + from pathlib import Path ++from tempfile import NamedTemporaryFile + from typing import Any + + from srtctl.contract import JobStatus +@@ -45,6 +48,39 @@ CREATE TABLE IF NOT EXISTS job_events ( + created_at TEXT NOT NULL + ); + ++-- Streamed log and metric output. A chunk is keyed by where it sits in its file, ++-- so a resend after a lost response is a no-op. ++CREATE TABLE IF NOT EXISTS job_logs ( ++ job_id TEXT NOT NULL, ++ file TEXT NOT NULL, ++ offset INTEGER NOT NULL, ++ size INTEGER NOT NULL, ++ data TEXT NOT NULL, ++ created_at TEXT NOT NULL, ++ PRIMARY KEY (job_id, file, offset) ++); ++ ++-- Raw uploads remain staged until every byte of one immutable generation arrives. ++CREATE TABLE IF NOT EXISTS job_captures ( ++ job_id TEXT NOT NULL, file TEXT NOT NULL, generation TEXT NOT NULL, ++ total INTEGER NOT NULL, next_offset INTEGER NOT NULL DEFAULT 0, ++ processed INTEGER NOT NULL DEFAULT 0, ++ PRIMARY KEY (job_id, file, generation) ++); ++CREATE TABLE IF NOT EXISTS capture_chunks ( ++ job_id TEXT NOT NULL, file TEXT NOT NULL, generation TEXT NOT NULL, ++ offset INTEGER NOT NULL, size INTEGER NOT NULL, digest BLOB NOT NULL, data BLOB, ++ PRIMARY KEY (job_id, file, generation, offset) ++); ++CREATE TABLE IF NOT EXISTS metric_observations ( ++ job_id TEXT NOT NULL, identity BLOB NOT NULL, ++ PRIMARY KEY (job_id, identity) ++) WITHOUT ROWID; ++CREATE TABLE IF NOT EXISTS job_log_ends ( ++ job_id TEXT NOT NULL, file TEXT NOT NULL, end INTEGER NOT NULL, ++ PRIMARY KEY (job_id, file) ++); ++ + CREATE INDEX IF NOT EXISTS idx_jobs_status ON jobs(status); + CREATE INDEX IF NOT EXISTS idx_jobs_cluster ON jobs(cluster); + CREATE INDEX IF NOT EXISTS idx_jobs_submitted_at ON jobs(submitted_at DESC); +@@ -69,6 +105,10 @@ def placeholder_name(job_id: str) -> str: + return f"job-{job_id}" + + ++class LogChunkConflict(ValueError): ++ """An incoming chunk disagrees with bytes already stored for the same file.""" ++ ++ + def _decode_job(row: sqlite3.Row) -> dict[str, Any]: + job = dict(row) + for key in _JSON_COLUMNS: +@@ -271,8 +311,154 @@ class StatusStore: + with self._transaction() as conn: + deleted = conn.execute("DELETE FROM jobs WHERE job_id = ?", (job_id,)).rowcount + conn.execute("DELETE FROM job_events WHERE job_id = ?", (job_id,)) ++ conn.execute("DELETE FROM job_logs WHERE job_id = ?", (job_id,)) ++ for table in ("job_captures", "capture_chunks", "metric_observations", "job_log_ends"): ++ conn.execute(f"DELETE FROM {table} WHERE job_id = ?", (job_id,)) + return deleted > 0 + ++ def append_logs(self, job_id: str, chunks: list[dict[str, Any]]) -> int: ++ """Store chunks atomically, accepting exact retries and rejecting conflicting ranges.""" ++ now = now_iso() ++ with self._transaction() as conn: ++ stored = 0 ++ for chunk in chunks: ++ overlapping = conn.execute( ++ """SELECT offset, size, data FROM job_logs ++ WHERE job_id = ? AND file = ? AND offset < ? ORDER BY offset DESC LIMIT 1""", ++ (job_id, chunk["file"], chunk["offset"] + chunk["size"]), ++ ).fetchone() ++ if overlapping is not None and overlapping["offset"] + overlapping["size"] > chunk["offset"]: ++ if all(overlapping[key] == chunk[key] for key in ("offset", "size", "data")): ++ continue ++ raise LogChunkConflict(f"Conflicting log chunk for {chunk['file']} at offset {chunk['offset']}") ++ conn.execute( ++ "INSERT INTO job_logs (job_id, file, offset, size, data, created_at) VALUES (?, ?, ?, ?, ?, ?)", ++ (job_id, chunk["file"], chunk["offset"], chunk["size"], chunk["data"], now), ++ ) ++ stored += 1 ++ return stored ++ ++ def append_raw_log(self, job_id: str, file: str, offset: int, data: bytes, *, final: bool = False) -> int: ++ """Keep raw bytes intact; UTF-8 decoding happens when the API is read.""" ++ stored = ( ++ self.append_logs(job_id, [{"file": file, "offset": offset, "size": len(data), "data": data}]) if data else 0 ++ ) ++ if final: ++ with self._transaction() as conn: ++ conn.execute( ++ "INSERT INTO job_log_ends VALUES (?, ?, ?) " ++ "ON CONFLICT(job_id, file) DO UPDATE SET end = MAX(end, excluded.end)", ++ (job_id, file, offset + len(data)), ++ ) ++ return stored ++ ++ def append_capture( ++ self, job_id: str, file: str, generation: str, offset: int, total: int, data: bytes ++ ) -> dict[str, Any]: ++ """Stage raw chunks durably, then decode complete captures on the collector. ++ ++ Digests remain after processing so a lost response can be retried exactly. ++ A new generation removes abandoned partial uploads for that same file. ++ """ ++ key = (job_id, file, generation) ++ digest = hashlib.sha256(data).digest() ++ with self._transaction() as conn: ++ manifest = conn.execute( ++ "SELECT * FROM job_captures WHERE job_id = ? AND file = ? AND generation = ?", key ++ ).fetchone() ++ if manifest is None: ++ if offset != 0: ++ raise LogChunkConflict("Capture generation must start at offset 0") ++ abandoned = conn.execute( ++ "SELECT generation FROM job_captures WHERE job_id = ? AND file = ? AND next_offset < total", ++ (job_id, file), ++ ).fetchall() ++ for old in abandoned: ++ for table in ("job_captures", "capture_chunks"): ++ conn.execute( ++ f"DELETE FROM {table} WHERE job_id = ? AND file = ? AND generation = ?", ++ (job_id, file, old["generation"]), ++ ) ++ conn.execute("INSERT INTO job_captures VALUES (?, ?, ?, ?, 0, 0)", (*key, total)) ++ next_offset, processed = 0, False ++ else: ++ if manifest["total"] != total: ++ raise LogChunkConflict("Capture generation total changed") ++ next_offset, processed = manifest["next_offset"], bool(manifest["processed"]) ++ existing = conn.execute( ++ "SELECT size, digest FROM capture_chunks WHERE job_id = ? AND file = ? AND generation = ? AND offset = ?", ++ (*key, offset), ++ ).fetchone() ++ if existing is not None: ++ if existing["size"] != len(data) or existing["digest"] != digest: ++ raise LogChunkConflict("Capture retry differs from stored bytes") ++ else: ++ if offset != next_offset or processed: ++ raise LogChunkConflict("Capture chunks must be contiguous") ++ conn.execute( ++ "INSERT INTO capture_chunks VALUES (?, ?, ?, ?, ?, ?, ?)", (*key, offset, len(data), digest, data) ++ ) ++ next_offset += len(data) ++ conn.execute( ++ "UPDATE job_captures SET next_offset = ? WHERE job_id = ? AND file = ? AND generation = ?", ++ (next_offset, *key), ++ ) ++ if next_offset == total and not processed: ++ self._process_capture(key) ++ return {"job_id": job_id, "next_offset": offset + len(data), "complete": offset + len(data) == total} ++ ++ def _process_capture(self, key: tuple[str, str, str]) -> None: ++ """Decode with bounded Python batches and atomically deduplicate each output batch. ++ ++ Retrying after a decoder or database failure cannot duplicate committed ++ observations. Raw bytes are deleted only after the entire decode succeeds. ++ """ ++ from srtctl.dsight.metrics import batches # lazy: pyarrow is only needed by the collector ++ ++ with NamedTemporaryFile(suffix=Path(key[1]).suffix) as capture: ++ with self._connect() as conn: ++ manifest = conn.execute( ++ "SELECT processed FROM job_captures WHERE job_id = ? AND file = ? AND generation = ?", key ++ ).fetchone() ++ if manifest is None or manifest["processed"]: ++ return ++ for row in conn.execute( ++ "SELECT data FROM capture_chunks WHERE job_id = ? AND file = ? AND generation = ? ORDER BY offset", ++ key, ++ ): ++ capture.write(row["data"]) ++ capture.flush() ++ for batch in batches(Path(capture.name)): ++ for start in range(0, batch.num_rows, 4096): ++ with self._transaction() as conn: ++ output = [] ++ for row in batch.slice(start, 4096).to_pylist(): ++ row.pop("metric_name_clean", None) ++ line = json.dumps(row, sort_keys=True, separators=(",", ":")) + "\n" ++ inserted = conn.execute( ++ "INSERT OR IGNORE INTO metric_observations VALUES (?, ?)", ++ (key[0], hashlib.sha256(line.encode()).digest()), ++ ).rowcount ++ if inserted: ++ output.append(line) ++ if output: ++ text = "".join(output) ++ offset = conn.execute( ++ "SELECT COALESCE(MAX(offset + size), 0) FROM job_logs WHERE job_id = ? AND file = ?", ++ (key[0], "tachometer_rows.jsonl"), ++ ).fetchone()[0] ++ conn.execute( ++ "INSERT INTO job_logs VALUES (?, ?, ?, ?, ?, ?)", ++ (key[0], "tachometer_rows.jsonl", offset, len(text.encode()), text, now_iso()), ++ ) ++ with self._transaction() as conn: ++ conn.execute( ++ "UPDATE job_captures SET processed = 1 WHERE job_id = ? AND file = ? AND generation = ?", key ++ ) ++ conn.execute( ++ "UPDATE capture_chunks SET data = NULL WHERE job_id = ? AND file = ? AND generation = ?", key ++ ) ++ + # ------------------------------------------------------------------- reads + + def get_job(self, job_id: str) -> dict[str, Any] | None: +@@ -330,3 +516,52 @@ class StatusStore: + with self._connect() as conn: + rows = conn.execute(query, params).fetchall() + return [dict(row) for row in rows] ++ ++ def list_log_files(self, job_id: str) -> list[dict[str, Any]]: ++ """Every streamed file of a job with the bytes received so far.""" ++ with self._connect() as conn: ++ rows = conn.execute( ++ """ ++ SELECT file, MAX(offset + size) AS size, MAX(created_at) AS updated_at ++ FROM job_logs WHERE job_id = ? GROUP BY file ORDER BY file ++ """, ++ (job_id,), ++ ).fetchall() ++ return [dict(row) for row in rows] ++ ++ def read_log(self, job_id: str, file: str, *, offset: int = 0, max_bytes: int = 1 << 20) -> tuple[str, int]: ++ """Contiguous content of ``file`` from ``offset``, stopping at a gap or after ``max_bytes``. ++ ++ Returns ``(data, next_offset)``. ``offset`` is a chunk boundary, i.e. 0 or a ++ ``next_offset`` from an earlier read. ++ """ ++ parts: list[str] = [] ++ cursor = offset ++ decoder = codecs.getincrementaldecoder("utf-8")(errors="replace") ++ with self._connect() as conn: ++ ending = conn.execute( ++ "SELECT end FROM job_log_ends WHERE job_id = ? AND file = ?", (job_id, file) ++ ).fetchone() ++ rows = conn.execute( ++ "SELECT offset, size, data FROM job_logs WHERE job_id = ? AND file = ? AND offset + size > ? ORDER BY offset", ++ (job_id, file, offset), ++ ) ++ for row in rows: ++ if row["offset"] > cursor or cursor - offset >= max_bytes: ++ break ++ if isinstance(row["data"], bytes): ++ raw = row["data"][cursor - row["offset"] : cursor - row["offset"] + max_bytes - (cursor - offset)] ++ parts.append(decoder.decode(raw)) ++ cursor += len(raw) ++ else: ++ # Legacy JSON chunks use source byte sizes, which can differ ++ # from the UTF-8 size of already-decoded replacement characters. ++ if row["offset"] != cursor: ++ break ++ parts.append(row["data"]) ++ cursor += row["size"] ++ if ending is not None and cursor == ending["end"]: ++ parts.append(decoder.decode(b"", final=True)) ++ else: ++ cursor -= len(decoder.getstate()[0]) ++ return "".join(parts), cursor +diff --git a/src/tachometer/tachometer-writer/src/writer.rs b/src/tachometer/tachometer-writer/src/writer.rs +index 60c1a769..8b97c528 100644 +--- a/src/tachometer/tachometer-writer/src/writer.rs ++++ b/src/tachometer/tachometer-writer/src/writer.rs +@@ -5,7 +5,7 @@ use arrow::ipc::writer::StreamWriter; + use arrow::record_batch::RecordBatch; + use log::{error, info}; + use std::fs::File; +-use std::io::BufWriter; ++use std::io::{BufWriter, Write}; + use std::path::{Path, PathBuf}; + use std::sync::Arc; + use tokio::sync::Mutex; +@@ -338,26 +338,21 @@ async fn periodic_save_task( + + /// Write an Arrow IPC stream file to local disk. + fn write_arrow_file_local(local_dir: &Path, filename: &str, batch: &RecordBatch) -> Result<()> { +- let path = local_dir.join(filename); +- let file = File::create(&path).map_err(|e| { +- crate::NoMoreError::Io(std::io::Error::other(format!( +- "Failed to create file {}: {}", +- path.display(), +- e +- ))) +- })?; +- let mut writer = BufWriter::new(file); +- +- let mut stream_writer = StreamWriter::try_new(&mut writer, batch.schema().as_ref()) +- .map_err(|e| crate::NoMoreError::Arrow(e.to_string()))?; +- stream_writer +- .write(batch) +- .map_err(|e| crate::NoMoreError::Arrow(e.to_string()))?; +- stream_writer +- .finish() +- .map_err(|e| crate::NoMoreError::Arrow(e.to_string()))?; +- +- Ok(()) ++ publish_capture_file(local_dir, filename, |file| { ++ let mut writer = BufWriter::new(file); ++ { ++ let mut stream_writer = StreamWriter::try_new(&mut writer, batch.schema().as_ref()) ++ .map_err(|e| crate::NoMoreError::Arrow(e.to_string()))?; ++ stream_writer ++ .write(batch) ++ .map_err(|e| crate::NoMoreError::Arrow(e.to_string()))?; ++ stream_writer ++ .finish() ++ .map_err(|e| crate::NoMoreError::Arrow(e.to_string()))?; ++ } ++ writer.flush()?; ++ Ok(()) ++ }) + } + + /// Write a Parquet file to local disk. +@@ -371,7 +366,7 @@ fn write_parquet_file_local(local_dir: &Path, filename: &str, batch: &RecordBatc + .set_write_batch_size(100_000) + .build(); + +- publish_parquet_file(local_dir, filename, |file| { ++ publish_capture_file(local_dir, filename, |file| { + let mut writer = ArrowWriter::try_new(file, batch.schema(), Some(props)) + .map_err(|e| crate::NoMoreError::Parquet(e.to_string()))?; + writer +@@ -384,9 +379,9 @@ fn write_parquet_file_local(local_dir: &Path, filename: &str, batch: &RecordBatc + }) + } + +-/// Keep in-progress files outside the compactor's out-*.parquet scan. The +-/// callback must close its Parquet writer successfully before publication. +-fn publish_parquet_file( ++/// Publish complete captures atomically. Open readers keep the previous snapshot ++/// during replacement, and the compactor cannot observe unfinished output. ++fn publish_capture_file( + local_dir: &Path, + filename: &str, + write: impl FnOnce(&mut File) -> Result<()>, +@@ -434,6 +429,27 @@ mod tests { + buffer.to_record_batch().unwrap() + } + ++ #[test] ++ fn arrow_snapshot_readers_survive_replacement() { ++ use std::io::Read; ++ ++ let local = tempfile::tempdir().unwrap(); ++ let path = local.path().join("current.arrow"); ++ write_arrow_file_local(local.path(), "current.arrow", &batch(0)).unwrap(); ++ let expected = std::fs::read(&path).unwrap(); ++ let mut open_snapshot = File::open(&path).unwrap(); ++ ++ write_arrow_file_local(local.path(), "current.arrow", &batch(3)).unwrap(); ++ let mut previous = Vec::new(); ++ open_snapshot.read_to_end(&mut previous).unwrap(); ++ assert_eq!(previous, expected); ++ assert_ne!(std::fs::read(&path).unwrap(), previous); ++ let mut reader = ++ arrow::ipc::reader::StreamReader::try_new(File::open(&path).unwrap(), None).unwrap(); ++ assert_eq!(reader.next().unwrap().unwrap(), batch(3)); ++ assert!(reader.next().is_none()); ++ } ++ + #[test] + fn compaction_cannot_observe_an_unfinished_parquet_file() { + let local = tempfile::tempdir().unwrap(); +@@ -446,7 +462,7 @@ mod tests { + .unwrap(); + write_parquet_file_local(local.path(), "out-1.parquet", &batch(0)).unwrap(); + +- publish_parquet_file(local.path(), "out-2.parquet", |file| { ++ publish_capture_file(local.path(), "out-2.parquet", |file| { + let next = batch(3); + let mut writer = ArrowWriter::try_new(file, next.schema(), None).unwrap(); + writer.write(&next).unwrap(); +@@ -526,7 +542,7 @@ mod tests { + #[test] + fn failed_parquet_write_does_not_publish_or_leave_a_partial_file() { + let local = tempfile::tempdir().unwrap(); +- let result = publish_parquet_file(local.path(), "out-1.parquet", |file| { ++ let result = publish_capture_file(local.path(), "out-1.parquet", |file| { + file.write_all(b"PAR1unfinished parquet")?; + Err(crate::NoMoreError::Io(std::io::Error::other( + "injected write failure", +diff --git a/tests/test_status.py b/tests/test_status.py +index b75a6b1b..2c7d2ee9 100644 +--- a/tests/test_status.py ++++ b/tests/test_status.py +@@ -598,3 +598,236 @@ class TestJobStageEnum: + assert JobStage.FRONTEND.value == "frontend" + assert JobStage.BENCHMARK.value == "benchmark" + assert JobStage.CLEANUP.value == "cleanup" ++ ++ ++class TestTachometerStreaming: ++ """The collector owns capture decoding, deduplication and text decoding.""" ++ ++ @staticmethod ++ def _store(tmp_path): ++ from srtctl.status_server.store import StatusStore ++ ++ store = StatusStore(tmp_path / "collector.sqlite3") ++ store.init() ++ return store ++ ++ @staticmethod ++ def _capture(tmp_path, rows, suffix="arrow"): ++ import pyarrow as pa ++ import pyarrow.parquet as pq ++ from pyarrow import ipc ++ ++ path = tmp_path / f"capture.{suffix}" ++ table = pa.Table.from_pylist(rows) ++ if suffix == "parquet": ++ pq.write_table(table, path) ++ else: ++ with ipc.new_stream(path, table.schema) as writer: ++ writer.write_table(table) ++ return path.read_bytes() ++ ++ @staticmethod ++ def _upload(store, data, file="tachometer/local/current.arrow", generation=None): ++ from uuid import uuid4 ++ ++ generation = generation or uuid4().hex ++ result = None ++ for offset in range(0, len(data), 127): ++ chunk = data[offset : offset + 127] ++ result = store.append_capture("8", file, generation, offset, len(data), chunk) ++ assert result["next_offset"] == offset + len(chunk) ++ assert result["complete"] ++ return generation ++ ++ @staticmethod ++ def _rows(store): ++ import json ++ ++ text, _ = store.read_log("8", "tachometer_rows.jsonl") ++ return [json.loads(line) for line in text.splitlines()] ++ ++ def test_snapshots_rotation_and_compaction_deduplicate_on_collector(self, tmp_path): ++ store = self._store(tmp_path) ++ rows = [ ++ {"timestamp_ns": 20, "metric_name": "z", "metric_value": 1.0}, ++ {"timestamp_ns": 20, "metric_name": "a", "metric_value": 2.0}, ++ {"timestamp_ns": 10, "metric_name": "a", "metric_value": 3.0}, ++ {"timestamp_ns": 30, "metric_name": "a", "metric_value": 4.0}, ++ ] ++ for end in (1, 2): ++ self._upload(store, self._capture(tmp_path, rows[:end])) ++ self._upload(store, self._capture(tmp_path, rows[:3], "parquet"), "tachometer/local/out-1.parquet") ++ self._upload(store, self._capture(tmp_path, rows[3:])) ++ compacted = [{**row, "metric_name_clean": row["metric_name"]} for row in reversed(rows)] ++ self._upload(store, self._capture(tmp_path, compacted, "parquet"), "tachometer/local/final.parquet") ++ assert self._rows(store) == rows ++ with store._connect() as conn: ++ assert conn.execute("SELECT COUNT(*) FROM capture_chunks WHERE data IS NOT NULL").fetchone()[0] == 0 ++ ++ def test_exact_retry_after_lost_final_response_never_reprocesses(self, tmp_path): ++ store = self._store(tmp_path) ++ data = self._capture(tmp_path, [{"metric_name": "a", "metric_value": 1}]) ++ generation = "a" * 32 ++ result = store.append_capture("8", "current.arrow", generation, 0, len(data), data) ++ with patch.object(type(store), "_process_capture", side_effect=AssertionError("Already processed")): ++ assert store.append_capture("8", "current.arrow", generation, 0, len(data), data) == result ++ assert len(self._rows(store)) == 1 ++ ++ def test_incomplete_generation_is_never_decoded_and_is_pruned(self, tmp_path): ++ store = self._store(tmp_path) ++ data = self._capture(tmp_path, [{"metric_name": "a", "metric_value": 1}]) ++ first = store.append_capture("8", "current.arrow", "a" * 32, 0, len(data), data[:100]) ++ assert first == {"job_id": "8", "next_offset": 100, "complete": False} ++ assert self._rows(store) == [] ++ store.append_capture("8", "current.arrow", "b" * 32, 0, len(data), data[:100]) ++ with store._connect() as conn: ++ assert conn.execute("SELECT COUNT(*) FROM job_captures").fetchone()[0] == 1 ++ assert conn.execute("SELECT COUNT(*) FROM capture_chunks").fetchone()[0] == 1 ++ store.append_capture("8", "current.arrow", "b" * 32, 100, len(data), data[100:]) ++ assert len(self._rows(store)) == 1 ++ ++ def test_capture_conflicts_and_out_of_order_chunks(self, tmp_path): ++ import pytest ++ ++ from srtctl.status_server.store import LogChunkConflict ++ ++ store = self._store(tmp_path) ++ args = ("8", "current.arrow", "a" * 32) ++ store.append_capture(*args, 0, 1000, b"first") ++ assert store.append_capture(*args, 0, 1000, b"first")["next_offset"] == 5 ++ for offset, total, data in ((0, 1000, b"other"), (0, 1001, b"first"), (9, 1000, b"gap"), (2, 1000, b"overlap")): ++ with pytest.raises(LogChunkConflict): ++ store.append_capture(*args, offset, total, data) ++ assert self._rows(store) == [] ++ ++ def test_decode_failure_retries_without_duplicate_output(self, tmp_path): ++ import pyarrow as pa ++ import pytest ++ ++ store = self._store(tmp_path) ++ rows = [{"metric_name": name, "metric_value": 1} for name in ("a", "b")] ++ data = self._capture(tmp_path, rows) ++ args = ("8", "current.arrow", "a" * 32, 0, len(data), data) ++ ++ def failing_batches(path): ++ yield pa.RecordBatch.from_pylist(rows[:1]) ++ raise OSError("temporary decoder failure") ++ ++ with patch("srtctl.dsight.metrics.batches", failing_batches), pytest.raises(OSError): ++ store.append_capture(*args) ++ assert self._rows(store) == rows[:1] ++ store.append_capture(*args) ++ assert self._rows(store) == rows ++ ++ def test_raw_log_decodes_split_utf8_and_final_invalid_tail(self, tmp_path): ++ store = self._store(tmp_path) ++ raw = "hello €!".encode() ++ assert store.append_raw_log("8", "worker.log", 0, raw[:7]) == 1 ++ assert store.read_log("8", "worker.log") == ("hello ", 6) ++ assert store.append_raw_log("8", "worker.log", 0, raw[:7]) == 0 ++ store.append_raw_log("8", "worker.log", 7, raw[7:] + b"\xe2") ++ assert store.read_log("8", "worker.log", offset=6) == ("€!", len(raw)) ++ store.append_raw_log("8", "worker.log", len(raw) + 1, b"", final=True) ++ assert store.read_log("8", "worker.log", offset=len(raw)) == ("�", len(raw) + 1) ++ ++ def test_raw_log_partial_reads_resume_within_chunk(self, tmp_path): ++ store = self._store(tmp_path) ++ raw = "a€b".encode() ++ store.append_raw_log("8", "worker.log", 0, raw) ++ assert store.read_log("8", "worker.log", max_bytes=3) == ("a", 1) ++ assert store.read_log("8", "worker.log", offset=1, max_bytes=4) == ("€b", 5) ++ ++ def test_delete_removes_capture_state_and_observation_index(self, tmp_path): ++ store = self._store(tmp_path) ++ store.create_job("8", "test") ++ self._upload(store, self._capture(tmp_path, [{"metric_name": "a", "metric_value": 1}])) ++ store.append_raw_log("8", "worker.log", 0, b"done", final=True) ++ assert store.delete_job("8") ++ with store._connect() as conn: ++ for table in ("job_captures", "capture_chunks", "metric_observations", "job_logs", "job_log_ends"): ++ assert conn.execute(f"SELECT COUNT(*) FROM {table}").fetchone()[0] == 0 ++ ++ def test_raw_routes_validate_bounds_and_paths(self, tmp_path): ++ from urllib.parse import urlencode ++ ++ import pytest ++ ++ from srtctl.status_server.server import ApiError, _raw_upload ++ ++ store = self._store(tmp_path) ++ params = {"file": "current.arrow", "generation": "a" * 32, "offset": 0, "total": 10} ++ for invalid in ( ++ {"file": "../current.arrow"}, ++ {"file": "/current.arrow"}, ++ {"file": "worker.log"}, ++ {"generation": "bad"}, ++ {"offset": -1}, ++ {"offset": 10}, ++ {"total": 0}, ++ {"total": 1 << 63}, ++ ): ++ with pytest.raises(ApiError) as error: ++ _raw_upload(store, "/api/jobs/8/captures?" + urlencode({**params, **invalid}), b"data") ++ assert error.value.status == 422 ++ ++ def test_binary_http_upload_requires_write_token_and_serves_rows(self, tmp_path): ++ import socket ++ import threading ++ from urllib.parse import urlencode ++ ++ import requests ++ ++ from srtctl.status_server.server import AuthPolicy, make_server ++ ++ store = self._store(tmp_path) ++ server = make_server(store, port=0, auth=AuthPolicy(write_token="secret", read_token="viewer")) ++ thread = threading.Thread(target=server.serve_forever, daemon=True) ++ thread.start() ++ try: ++ base = f"http://127.0.0.1:{server.server_port}/api/jobs/8" ++ data = self._capture(tmp_path, [{"metric_name": "a", "metric_value": 1}]) ++ path = ( ++ base ++ + "/captures?" ++ + urlencode( ++ { ++ "file": "tachometer/local/current.arrow", ++ "generation": "a" * 32, ++ "offset": 0, ++ "total": len(data), ++ "cluster": "test-cluster", ++ } ++ ) ++ ) ++ headers = {"Content-Type": "application/octet-stream"} ++ assert requests.post(path, data=data, headers=headers, timeout=5).status_code == 401 ++ assert ( ++ requests.post( ++ path, data=data, headers={**headers, "Authorization": "Bearer viewer"}, timeout=5 ++ ).status_code ++ == 403 ++ ) ++ # A disconnected sender must not leave a shorter acknowledged chunk. ++ with socket.create_connection(("127.0.0.1", server.server_port), timeout=5) as connection: ++ connection.sendall( ++ b"POST /api/jobs/8/logs?file=truncated.log&offset=0 HTTP/1.1\r\n" ++ b"Host: localhost\r\nAuthorization: Bearer secret\r\n" ++ b"Content-Type: application/octet-stream\r\nContent-Length: 10\r\n\r\nabc" ++ ) ++ connection.shutdown(socket.SHUT_WR) ++ assert b"400" in connection.recv(4096).split(b"\r\n")[0] ++ assert store.read_log("8", "truncated.log") == ("", 0) ++ headers["Authorization"] = "Bearer secret" ++ result = requests.post(path, data=data, headers=headers, timeout=5) ++ assert result.status_code == 200 ++ assert result.json() == {"job_id": "8", "next_offset": len(data), "complete": True} ++ result = requests.get(base + "/logs?file=tachometer_rows.jsonl", headers=headers, timeout=5) ++ assert result.status_code == 200 ++ assert '"metric_name":"a"' in result.json()["data"] ++ log = base + "/logs?file=worker.log&offset=0&final=1" ++ assert requests.post(log, data=b"raw log", headers=headers, timeout=5).json()["stored"] == 1 ++ assert requests.get(base + "/logs?file=worker.log", headers=headers, timeout=5).json()["data"] == "raw log" ++ finally: ++ server.shutdown() ++ server.server_close() ++ thread.join(timeout=5) +diff --git a/tests/test_status_server.py b/tests/test_status_server.py +index 45b563a6..1580cb36 100644 +--- a/tests/test_status_server.py ++++ b/tests/test_status_server.py +@@ -12,6 +12,7 @@ agree with each other with nothing patched in between. + from __future__ import annotations + + import http.client ++import json + import logging + import sys + import threading +@@ -868,3 +869,383 @@ class TestCreateJobRecordRetry: + assert [r.getMessage() for r in caplog.records if r.levelno == logging.WARNING] == [ + "Status report to https://collector.example lost after 2 attempts: timed out" + ] ++ ++ ++class TestLogCollector: ++ def test_exact_retry_is_ignored_and_conflicting_batch_is_rolled_back(self, base_url, store): ++ url = f"{base_url}/api/jobs/7/logs" ++ original = {"file": "worker.log", "offset": 0, "size": 3, "data": "abc"} ++ payload = {"chunks": [original], "metadata": {"cluster": "test-cluster"}} ++ assert requests.post(url, json=payload, timeout=5).json()["stored"] == 1 ++ assert requests.post(url, json=payload, timeout=5).json()["stored"] == 0 ++ ++ response = requests.post( ++ url, ++ json={ ++ "chunks": [ ++ {"file": "other.log", "offset": 0, "size": 3, "data": "new"}, ++ {**original, "size": 6, "data": "abcdef"}, ++ ] ++ }, ++ timeout=5, ++ ) ++ assert response.status_code == 409 ++ assert store.read_log("7", "worker.log") == ("abc", 3) ++ assert store.list_log_files("7")[0]["file"] == "worker.log" ++ assert len(store.list_log_files("7")) == 1 ++ ++ @pytest.mark.parametrize( ++ "conflict", ++ [ ++ {"offset": 3, "size": 3, "data": "xyz"}, ++ {"offset": 2, "size": 2, "data": "cd"}, ++ {"offset": 5, "size": 2, "data": "fg"}, ++ {"offset": 0, "size": 9, "data": "abcdefghi"}, ++ ], ++ ) ++ def test_overlapping_ranges_are_rejected(self, base_url, store, conflict): ++ url = f"{base_url}/api/jobs/7/logs" ++ original = {"file": "worker.log", "offset": 3, "size": 3, "data": "def"} ++ assert requests.post(url, json={"chunks": [original]}, timeout=5).status_code == 200 ++ response = requests.post(url, json={"chunks": [{"file": "worker.log", **conflict}]}, timeout=5) ++ assert response.status_code == 409 ++ # Out-of-order, non-overlapping chunks remain supported. ++ prefix = {"file": "worker.log", "offset": 0, "size": 3, "data": "abc"} ++ assert requests.post(url, json={"chunks": [prefix]}, timeout=5).status_code == 200 ++ assert store.read_log("7", "worker.log") == ("abcdef", 6) ++ ++ @pytest.mark.parametrize( ++ "invalid", ++ [ ++ {"size": 0}, ++ {"size": (1 << 20) + 1}, ++ {"offset": -1}, ++ {"offset": 1 << 63}, ++ {"offset": (1 << 63) - 1}, ++ {"file": "../outside.log"}, ++ {"file": "/absolute.log"}, ++ {"file": "worker/../outside.log"}, ++ {"data": ""}, ++ ], ++ ) ++ def test_invalid_chunks_are_rejected_before_storage(self, base_url, store, invalid): ++ chunk = {"file": "worker.log", "offset": 0, "size": 3, "data": "abc", **invalid} ++ response = requests.post(f"{base_url}/api/jobs/7/logs", json={"chunks": [chunk]}, timeout=5) ++ assert response.status_code == 422 ++ assert store.list_log_files("7") == [] ++ ++ def test_log_routes_require_the_existing_tokens(self, auth_url): ++ url = f"{auth_url}/api/jobs/7/logs" ++ payload = {"chunks": [{"file": "worker.log", "offset": 0, "size": 3, "data": "abc"}]} ++ assert requests.post(url, json=payload, timeout=5).status_code == 401 ++ assert requests.post(url, json=payload, headers=_bearer(READ), timeout=5).status_code == 403 ++ assert requests.post(url, json=payload, headers=_bearer(WRITE), timeout=5).status_code == 200 ++ assert requests.get(url, timeout=5).status_code == 401 ++ assert requests.get(url, headers=_bearer(READ), timeout=5).status_code == 200 ++ ++ ++class TestLogStreaming: ++ def test_streamed_files_are_stored_as_deltas_and_read_back(self, base_url, tmp_path): ++ from srtctl.core.status import LogStreamer ++ ++ log_dir = tmp_path / "logs" ++ (log_dir / "power").mkdir(parents=True) ++ worker = log_dir / "node-01_decode_w0.out" ++ worker.write_text("loading\n") ++ (log_dir / "power" / "samples.csv").write_text("t,w\n0,700\n") ++ (log_dir / "benchmark-rollup.json").write_text("{}") # not an append-only suffix ++ ++ reporter = StatusReporter(job_id="7", api_endpoints=(base_url,)) ++ streamer = LogStreamer(reporter, log_dir, interval=60) ++ streamer.flush() ++ with worker.open("a") as f: ++ f.write("ready\n") ++ streamer.flush() ++ streamer.flush() # nothing new ++ ++ files = _get(base_url, "/api/jobs/7/logs").json()["files"] ++ assert [(f["file"], f["size"]) for f in files] == [("node-01_decode_w0.out", 14), ("power/samples.csv", 10)] ++ ++ log = _get(base_url, "/api/jobs/7/logs?file=node-01_decode_w0.out").json() ++ assert (log["data"], log["next_offset"]) == ("loading\nready\n", 14) ++ tail = _get(base_url, "/api/jobs/7/logs?file=node-01_decode_w0.out&offset=8").json() ++ assert tail["data"] == "ready\n" ++ ++ # A resend of a stored chunk is a no-op. ++ resent = requests.post( ++ f"{base_url}/api/jobs/7/logs", ++ params={"file": "node-01_decode_w0.out", "offset": 0}, ++ data=b"loading\n", ++ headers={"Content-Type": "application/octet-stream"}, ++ timeout=5, ++ ) ++ assert resent.json()["stored"] == 0 ++ ++ def test_unknown_job_logs_404(self, base_url): ++ assert _get(base_url, "/api/jobs/nope/logs").status_code == 404 ++ ++ def test_tachometer_rows_stream_once_as_json_lines(self, base_url, tmp_path): ++ import pyarrow as pa ++ from pyarrow import ipc ++ ++ from srtctl.core.status import TACHOMETER_STREAM_FILE, LogStreamer ++ ++ log_dir, capture = tmp_path / "logs", tmp_path / "logs" / "tachometer" / "local" ++ capture.mkdir(parents=True) ++ ++ def snapshot(values: list[float]) -> None: # tachometer rewrites current.arrow with its whole buffer ++ table = pa.table( ++ { ++ "metric_name": ["tok_s"] * len(values), ++ "metric_value": values, ++ "timestamp_ns": list(range(1, len(values) + 1)), ++ } ++ ) ++ with ipc.new_file(capture / "current.arrow", table.schema) as writer: ++ writer.write_table(table) ++ ++ streamer = LogStreamer(StatusReporter(job_id="8", api_endpoints=(base_url,)), log_dir, 60, capture) ++ snapshot([1.0]) ++ streamer.flush() ++ snapshot([1.0, 2.0]) ++ streamer.flush() ++ ++ log = _get(base_url, f"/api/jobs/8/logs?file={TACHOMETER_STREAM_FILE}").json() ++ assert [json.loads(line)["metric_value"] for line in log["data"].splitlines()] == [1.0, 2.0] ++ assert not (log_dir / TACHOMETER_STREAM_FILE).exists() ++ assert not (log_dir / ".tachometer-stream.sqlite3").exists() ++ ++ def test_lost_response_retries_identical_bytes_after_file_grows(self, base_url, tmp_path, monkeypatch): ++ from srtctl.core.status import LogStreamer ++ ++ path = tmp_path / "worker.log" ++ path.write_text("abc") ++ streamer = LogStreamer(StatusReporter(job_id="retry", api_endpoints=(base_url,)), tmp_path, 60) ++ real_post = requests.Session.post ++ sent = [] ++ ++ def lose_first_response(session, url, **kwargs): ++ sent.append({"offset": int(kwargs["params"]["offset"]), "data": kwargs["data"].decode()}) ++ response = real_post(session, url, **kwargs) ++ if len(sent) == 1: ++ raise requests.Timeout("response lost after commit") ++ return response ++ ++ monkeypatch.setattr(requests.Session, "post", lose_first_response) ++ streamer.flush() ++ with path.open("a") as out: ++ out.write("def") ++ streamer.flush() ++ with path.open("a") as out: ++ out.write("ghi") ++ streamer.flush() ++ log = _get(base_url, "/api/jobs/retry/logs?file=worker.log").json() ++ assert (log["data"], log["next_offset"]) == ("abcdefghi", 9) ++ assert [(chunk["offset"], chunk["data"]) for chunk in sent] == [ ++ (0, "abc"), ++ (0, "abc"), ++ (3, "def"), ++ (6, "ghi"), ++ ] ++ ++ def test_utf8_survives_chunk_and_append_boundaries(self, base_url, tmp_path): ++ from srtctl.core.status import STREAM_REQUEST_BYTES, LogStreamer ++ ++ path = tmp_path / "worker.log" ++ prefix = "a" * (STREAM_REQUEST_BYTES - 1) ++ path.write_bytes(prefix.encode() + "€".encode()[:2]) ++ streamer = LogStreamer(StatusReporter(job_id="utf8", api_endpoints=(base_url,)), tmp_path, 60) ++ streamer.flush() ++ with path.open("ab") as out: ++ out.write("€".encode()[2:] + "文\n".encode()) ++ streamer.flush() ++ log = _get(base_url, "/api/jobs/utf8/logs?file=worker.log").json() ++ tail = _get(base_url, f"/api/jobs/utf8/logs?file=worker.log&offset={log['next_offset']}").json() ++ assert log["data"] + tail["data"] == prefix + "€文\n" ++ assert tail["next_offset"] == len(path.read_bytes()) ++ ++ def test_stream_authenticates_and_includes_cluster_on_each_chunk(self, auth_url, tmp_path, monkeypatch): ++ from srtctl.core.status import LogStreamer ++ ++ monkeypatch.setenv("SRTCTL_STATUS_TOKEN", WRITE) ++ monkeypatch.setattr("srtctl.core.status._cluster_setting", lambda: "mi300x-amd") ++ path = tmp_path / "worker.log" ++ path.write_text("first\n") ++ real_post = requests.Session.post ++ payloads = [] ++ ++ def record(session, url, **kwargs): ++ payloads.append(dict(kwargs["params"])) ++ return real_post(session, url, **kwargs) ++ ++ monkeypatch.setattr(requests.Session, "post", record) ++ streamer = LogStreamer(StatusReporter(job_id="cluster", api_endpoints=(auth_url,)), tmp_path, 60) ++ streamer.flush() ++ with path.open("a") as out: ++ out.write("second\n") ++ streamer.flush() ++ assert [payload["cluster"] for payload in payloads] == ["mi300x-amd", "mi300x-amd"] ++ response = requests.get( ++ f"{auth_url}/api/jobs/cluster/logs?file=worker.log", ++ headers=_bearer(READ), ++ timeout=5, ++ ) ++ assert response.json()["data"] == "first\nsecond\n" ++ ++ def test_stop_does_not_wait_indefinitely_for_a_collector(self, tmp_path, monkeypatch, caplog): ++ from srtctl.core.status import LogStreamer ++ ++ entered, release = threading.Event(), threading.Event() ++ (tmp_path / "worker.log").write_text("data") ++ ++ def stalled_post(*args, **kwargs): ++ entered.set() ++ release.wait(timeout=5) ++ raise requests.Timeout("stalled collector") ++ ++ monkeypatch.setattr(requests.Session, "post", stalled_post) ++ monkeypatch.setattr("srtctl.core.status.STREAM_SHUTDOWN_SECONDS", 0.02) ++ streamer = LogStreamer(StatusReporter(job_id="stalled", api_endpoints=("https://collector",)), tmp_path, 0.001) ++ streamer.start() ++ try: ++ assert entered.wait(timeout=2) ++ streamer.stop() ++ assert not release.is_set() ++ assert "local logs are retained" in caplog.text ++ finally: ++ release.set() ++ streamer._thread.join(timeout=2) ++ assert not streamer._thread.is_alive() ++ ++ def test_stop_flushes_final_bytes_before_first_interval(self, base_url, tmp_path): ++ from srtctl.core.status import LogStreamer ++ ++ (tmp_path / "worker.log").write_bytes(b"last line\n\xe2") ++ streamer = LogStreamer(StatusReporter(job_id="final", api_endpoints=(base_url,)), tmp_path, 60) ++ streamer.start() ++ streamer.stop() ++ log = _get(base_url, "/api/jobs/final/logs?file=worker.log").json() ++ assert log["data"] == "last line\n\ufffd" ++ assert log["next_offset"] == 11 ++ ++ def test_shutdown_finalizes_an_already_uploaded_partial_character(self, base_url, tmp_path): ++ from srtctl.core.status import LogStreamer ++ ++ (tmp_path / "worker.log").write_bytes(b"done\xe2") ++ streamer = LogStreamer(StatusReporter(job_id="eof", api_endpoints=(base_url,)), tmp_path, 60) ++ streamer.flush() ++ first = _get(base_url, "/api/jobs/eof/logs?file=worker.log").json() ++ assert (first["data"], first["next_offset"]) == ("done", 4) ++ streamer.start() ++ streamer.stop() ++ final = _get(base_url, "/api/jobs/eof/logs?file=worker.log&offset=4").json() ++ assert (final["data"], final["next_offset"]) == ("\ufffd", 5) ++ ++ def test_capture_retry_keeps_raw_bytes_after_compaction_unlinks_source(self, base_url, tmp_path, monkeypatch): ++ import pyarrow as pa ++ import pyarrow.parquet as pq ++ ++ from srtctl.core.status import LogStreamer ++ ++ capture = tmp_path / "tachometer" / "local" ++ capture.mkdir(parents=True) ++ path = capture / "out-1.parquet" ++ pq.write_table(pa.table({"metric_name": ["x"], "metric_value": [3.0]}), path) ++ original = path.read_bytes() ++ sent = [] ++ real_post = requests.Session.post ++ ++ def lose_response(session, url, **kwargs): ++ if url.endswith("/captures"): ++ sent.append((dict(kwargs["params"]), kwargs["data"])) ++ response = real_post(session, url, **kwargs) ++ if len(sent) == 1: ++ path.unlink() ++ raise requests.Timeout("lost capture acknowledgment") ++ return response ++ return real_post(session, url, **kwargs) ++ ++ monkeypatch.setattr(requests.Session, "post", lose_response) ++ streamer = LogStreamer(StatusReporter(job_id="capture-retry", api_endpoints=(base_url,)), tmp_path, 60, capture) ++ streamer.flush() ++ streamer.flush() ++ assert len(sent) == 2 ++ assert sent[0] == sent[1] ++ assert sent[0][1] == original ++ log = _get(base_url, "/api/jobs/capture-retry/logs?file=tachometer_rows.jsonl").json() ++ assert len(log["data"].splitlines()) == 1 ++ assert not streamer._captures ++ ++ def test_capture_replacement_finishes_old_generation_without_copying(self, base_url, tmp_path, monkeypatch): ++ import pyarrow as pa ++ from pyarrow import ipc ++ ++ from srtctl.core.status import LogStreamer ++ ++ capture = tmp_path / "tachometer" / "local" ++ capture.mkdir(parents=True) ++ path = capture / "current.arrow" ++ ++ def write_snapshot(destination, value): ++ table = pa.table({"metric_name": ["x"], "metric_value": [value]}) ++ with ipc.new_stream(destination, table.schema) as writer: ++ writer.write_table(table) ++ ++ write_snapshot(path, 1.0) ++ original = path.read_bytes() ++ replacement = capture / "next.tmp" ++ write_snapshot(replacement, 2.0) ++ real_post = requests.Session.post ++ raw_chunks = [] ++ ++ def replace_during_upload(session, url, **kwargs): ++ if url.endswith("/captures"): ++ raw_chunks.append(kwargs["data"]) ++ if len(raw_chunks) == 1: ++ replacement.replace(path) ++ return real_post(session, url, **kwargs) ++ ++ monkeypatch.setattr("srtctl.core.status.STREAM_REQUEST_BYTES", 128) ++ monkeypatch.setattr(requests.Session, "post", replace_during_upload) ++ streamer = LogStreamer(StatusReporter(job_id="replace", api_endpoints=(base_url,)), tmp_path, 60, capture) ++ streamer.flush() ++ assert b"".join(raw_chunks) == original ++ streamer.flush() ++ log = _get(base_url, "/api/jobs/replace/logs?file=tachometer_rows.jsonl").json() ++ assert [json.loads(row)["metric_value"] for row in log["data"].splitlines()] == [1.0, 2.0] ++ ++ def test_rewritten_capture_never_commits_mixed_generations(self, base_url, tmp_path, monkeypatch): ++ import pyarrow as pa ++ from pyarrow import ipc ++ ++ from srtctl.core.status import LogStreamer ++ ++ capture = tmp_path / "tachometer" / "local" ++ capture.mkdir(parents=True) ++ path = capture / "current.arrow" ++ ++ def write_snapshot(value): ++ table = pa.table({"metric_name": ["x"], "metric_value": [value]}) ++ with ipc.new_stream(path, table.schema) as writer: ++ writer.write_table(table) ++ ++ write_snapshot(1.0) ++ real_post = requests.Session.post ++ sent = [] ++ ++ def rewrite_during_upload(session, url, **kwargs): ++ if url.endswith("/captures"): ++ sent.append(dict(kwargs["params"])) ++ if len(sent) == 1: ++ write_snapshot(2.0) ++ return real_post(session, url, **kwargs) ++ ++ monkeypatch.setattr("srtctl.core.status.STREAM_REQUEST_BYTES", 128) ++ monkeypatch.setattr(requests.Session, "post", rewrite_during_upload) ++ streamer = LogStreamer(StatusReporter(job_id="rewrite", api_endpoints=(base_url,)), tmp_path, 60, capture) ++ streamer.flush() ++ assert _get(base_url, "/api/jobs/rewrite/logs?file=tachometer_rows.jsonl").json()["data"] == "" ++ streamer.flush() ++ assert sent[0]["generation"] != sent[1]["generation"] ++ log = _get(base_url, "/api/jobs/rewrite/logs?file=tachometer_rows.jsonl").json() ++ assert [json.loads(row)["metric_value"] for row in log["data"].splitlines()] == [2.0] diff --git a/inferencex-e2e/runners/srt-slurm/patches/README.md b/inferencex-e2e/runners/srt-slurm/patches/README.md index feb1dab1c5..121ee38728 100644 --- a/inferencex-e2e/runners/srt-slurm/patches/README.md +++ b/inferencex-e2e/runners/srt-slurm/patches/README.md @@ -8,3 +8,4 @@ Each patch is a temporary fix for an open upstream PR. When the PR merges and th | Patch | Upstream PR | Fix | |-------|-------------|-----| +| `539.patch` | [NVIDIA/srt-slurm#539](https://github.com/NVIDIA/srt-slurm/pull/539) | Stream raw logs and Tachometer captures; use the source-built atomic writer from the sweep build job. |