Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
36 changes: 36 additions & 0 deletions packages/cli/src/telemetry/events.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,9 +2,20 @@ import { describe, expect, it, vi, beforeEach } from "vitest";

const trackEvent = vi.fn();
const flush = vi.fn(() => Promise.resolve());
const shouldTrack = vi.fn(() => true);
vi.mock("./client.js", () => ({
trackEvent: (...args: unknown[]) => trackEvent(...args),
flush: () => flush(),
shouldTrack: () => shouldTrack(),
}));

// Power state shells out to `pmset`; spy so tests can assert it is NOT
// sampled for opted-out installs (the fields are built at the call site,
// before trackEvent's own shouldTrack guard).
const getPowerState = vi.fn(() => ({ on_battery: true, low_power_mode: false }));
vi.mock("./system.js", async () => ({
...(await vi.importActual<typeof import("./system.js")>("./system.js")),
getPowerState: () => getPowerState(),
}));

// identifyUser reads the install anonymousId; pin it so the $identify alias is
Expand Down Expand Up @@ -706,3 +717,28 @@ describe("auth login telemetry events", () => {
expect(trackEvent).not.toHaveBeenCalled();
});
});

describe("power-state sampling respects the telemetry opt-out", () => {
beforeEach(() => {
getPowerState.mockClear();
shouldTrack.mockReturnValue(true);
});

it("samples power state for a tracked render", () => {
trackRenderComplete({ durationMs: 1, fps: 30, quality: "high", docker: false, gpu: false });
expect(getPowerState).toHaveBeenCalled();
const props = trackEvent.mock.calls.at(-1)?.[1] as Record<string, unknown>;
expect(props.on_battery).toBe(true);
expect(props.low_power_mode).toBe(false);
});

it("does NOT spawn pmset when telemetry is disabled", () => {
// Regression: powerStateFields() is spread into the properties object at
// the call site, so it runs BEFORE trackEvent's `if (!shouldTrack())`
// guard — an opted-out install would otherwise pay two blocking
// subprocess spawns per render for an event that is then discarded.
shouldTrack.mockReturnValue(false);
trackRenderComplete({ durationMs: 1, fps: 30, quality: "high", docker: false, gpu: false });
expect(getPowerState).not.toHaveBeenCalled();
});
});
26 changes: 25 additions & 1 deletion packages/cli/src/telemetry/events.ts
Original file line number Diff line number Diff line change
@@ -1,8 +1,30 @@
import { redactTelemetryString, type OutputResolutionIssueKind } from "@hyperframes/core";
import type { SubTimelineWaitOutcome } from "@hyperframes/engine";
import { FEEDBACK_RATING_SCALE } from "../utils/feedbackRating.js";
import { flush, trackEvent } from "./client.js";
import { flush, shouldTrack, trackEvent } from "./client.js";
import { readConfig } from "./config.js";
import { getPowerState } from "./system.js";

// Power state is volatile (a laptop docks/undocks mid-session), so it is
// sampled per render event rather than cached with SystemMeta. Attached to
// render_complete AND render_error: the DE fleet is macOS laptops whose
// power management shifts render perf ~1.8x with no other telemetry signal,
// and perf/soak analysis needs to segment by it (see getPowerState).
//
// shouldTrack() is checked HERE, not just inside trackEvent: this helper is
// spread into the properties object at the CALL SITE, so it runs before
// trackEvent's own `if (!shouldTrack()) return` guard. Without this an
// opted-out install would still pay two blocking `pmset` subprocess spawns
// per render for an event that is then discarded (review finding).
// shouldTrack() memoizes, so this costs nothing on the tracked path.
function powerStateFields(): { on_battery?: boolean; low_power_mode?: boolean } {
if (!shouldTrack()) return {};
const power = getPowerState();
return {
on_battery: power.on_battery ?? undefined,
low_power_mode: power.low_power_mode ?? undefined,
};
}

// run_id is attached only when the orchestrator set HYPERFRAMES_RUN_ID — an
// absent property, never null/"" (PostHog treats those as real values).
Expand Down Expand Up @@ -300,6 +322,7 @@ export function trackRenderComplete(
de_blank_recaptures: props.deBlankRecaptures,
de_boundary_frames: props.deBoundaryFrames,
de_ncpr_fallbacks: props.deNcprFallbacks,
...powerStateFields(),
source: props.source ?? "cli",
composition_duration_ms: props.compositionDurationMs,
composition_width: props.compositionWidth,
Expand Down Expand Up @@ -379,6 +402,7 @@ export function trackRenderError(
elapsed_ms: props.elapsedMs,
peak_memory_mb: props.peakMemoryMb,
memory_free_mb: props.memoryFreeMb,
...powerStateFields(),
// gpu_renderer arrives via renderObservabilityEventProperties below:
// on the failure path perfSummary is never built, so live capture
// observability is the only source. Backend attribution matters MOST
Expand Down
51 changes: 51 additions & 0 deletions packages/cli/src/telemetry/system.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -110,3 +110,54 @@ describe("getAvailableMemoryMb", () => {
expect(getAvailableMemoryMb()).toBe(6144);
});
});

describe("power state (laptop fleet segmentation)", () => {
it("parsePmsetPowerSource reads battery vs AC", async () => {
const { parsePmsetPowerSource } = await import("./system.js");
expect(parsePmsetPowerSource("Now drawing from 'Battery Power'\n -InternalBattery-0")).toBe(
true,
);
expect(parsePmsetPowerSource("Now drawing from 'AC Power'\n -InternalBattery-0")).toBe(false);
expect(parsePmsetPowerSource("garbage output")).toBe(null);
});

it("getPowerState samples pmset on darwin", async () => {
vi.doMock("node:os", async () => ({
...(await vi.importActual<typeof import("node:os")>("node:os")),
platform: () => "darwin",
}));
vi.doMock("node:child_process", async () => ({
...(await vi.importActual<typeof import("node:child_process")>("node:child_process")),
execSync: (cmd: string) =>
cmd === "pmset -g batt"
? "Now drawing from 'Battery Power'\n"
: "SleepDisabled 0\n lowpowermode 1\n",
}));
const { getPowerState } = await import("./system.js");
expect(getPowerState()).toEqual({ on_battery: true, low_power_mode: true });
});

it("getPowerState returns nulls off-darwin (no guessing)", async () => {
vi.doMock("node:os", async () => ({
...(await vi.importActual<typeof import("node:os")>("node:os")),
platform: () => "linux",
}));
const { getPowerState } = await import("./system.js");
expect(getPowerState()).toEqual({ on_battery: null, low_power_mode: null });
});

it("getPowerState survives pmset failure with nulls", async () => {
vi.doMock("node:os", async () => ({
...(await vi.importActual<typeof import("node:os")>("node:os")),
platform: () => "darwin",
}));
vi.doMock("node:child_process", async () => ({
...(await vi.importActual<typeof import("node:child_process")>("node:child_process")),
execSync: () => {
throw new Error("pmset: command failed");
},
}));
const { getPowerState } = await import("./system.js");
expect(getPowerState()).toEqual({ on_battery: null, low_power_mode: null });
});
});
54 changes: 54 additions & 0 deletions packages/cli/src/telemetry/system.ts
Original file line number Diff line number Diff line change
Expand Up @@ -184,6 +184,60 @@ export function getFreeDiskMb(path: string = "."): number | null {
}
}

export interface PowerState {
/** true = running on battery, false = external power, null = undetectable. */
on_battery: boolean | null;
/** macOS Low Power Mode; null off-darwin or undetectable. */
low_power_mode: boolean | null;
}

/**
* Parse `pmset -g batt` output for the power source line. Exported for tests.
* Example first line: `Now drawing from 'Battery Power'`.
*/
export function parsePmsetPowerSource(raw: string): boolean | null {
const m = raw.match(/Now drawing from '([^']+)'/);
if (!m) return null;
return m[1] === "Battery Power";
}

/**
* Sample the machine's power state. Volatile — sample at event time, never
* cache alongside SystemMeta.
*
* Why this exists: the DE fast path only engages on macOS + hardware GPU, so
* the render fleet is overwhelmingly laptops, and laptop perf is
* power-managed — bench sweeps on an M4 Pro showed the SAME render flipping
* between ~9.6 and ~17.2 ms/frame regimes with no load/thermal signal to
* explain it. Without a power-state dimension on render telemetry those
* regimes are indistinguishable noise; with it, perf distributions (and the
* DE parallel-router soak) can be segmented by the machine state real users
* actually render in.
*/
export function getPowerState(): PowerState {
if (platform() !== "darwin") {
// Linux laptops exist but the DE fleet is darwin; don't guess elsewhere.
return { on_battery: null, low_power_mode: null };
}
let on_battery: boolean | null = null;
let low_power_mode: boolean | null = null;
try {
on_battery = parsePmsetPowerSource(
execSync("pmset -g batt", { encoding: "utf-8", timeout: 2000 }),
);
} catch {
// pmset missing/slow — leave null rather than fail telemetry.
}
try {
const raw = execSync("pmset -g", { encoding: "utf-8", timeout: 2000 });
const m = raw.match(/lowpowermode\s+(\d)/);
if (m) low_power_mode = m[1] === "1";
} catch {
// Same: absence of the reading is itself acceptable.
}
return { on_battery, low_power_mode };
}

/**
* Get available memory in MB, accounting for OS-level page caching.
*
Expand Down
11 changes: 11 additions & 0 deletions packages/producer/src/services/renderOrchestrator.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1843,6 +1843,7 @@ describe("shouldPreferParallelDrawElement (DE parallel router)", () => {
probeDeGated: false,
experimentalParallelDeOptIn: false,
routerEnabled: true,
parallelStreamingAvailable: true,
totalMemoryMb: 32768,
minMemoryMb: 24576,
};
Expand All @@ -1851,6 +1852,16 @@ describe("shouldPreferParallelDrawElement (DE parallel router)", () => {
expect(shouldPreferParallelDrawElement(eligible)).toBe(true);
});

it("withholds the bet when parallel streaming can't run (e.g. over the duration cap)", () => {
// The router pins workerCount to 3 and skips calibration to serve the
// verified parallel DE STREAMING path. If streaming is off for this
// render — the >240s duration cap is the common case — firing would pay
// the whole cost of the pin for none of the benefit.
expect(
shouldPreferParallelDrawElement({ ...eligible, parallelStreamingAvailable: false }),
).toBe(false);
});

it("withholds the parallel bet below the RAM floor (16 GB black-slab report)", () => {
expect(shouldPreferParallelDrawElement({ ...eligible, totalMemoryMb: 16384 })).toBe(false);
});
Expand Down
50 changes: 42 additions & 8 deletions packages/producer/src/services/renderOrchestrator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1340,9 +1340,14 @@ export function resolveInversionRetryPlan(args: {
* but the decision itself stays gated behind its own flag pending the
* telemetry soak (revert rate, de_verify_min_db distribution) on real wild
* traffic — there is currently none, since nothing routes here by default.
* Takes priority over the single-worker inversion when both would fire (a
* higher minFrames than HF_DE_SINGLE_MIN_FRAMES is the intended shape: this
* only picks up the long tail the inversion's own benchmark didn't cover).
* Takes priority over the single-worker inversion when both would fire.
* Re-calibrated 2026-07-27: a controlled crossover sweep (three content
* profiles including a genuinely init-expensive 24-sub-composition comp;
* worker counts and capture modes verified per run) found par3 > single at
* every size from 350f up in every profile — workers init concurrently, so
* per-worker init duplication costs CPU, not wall-clock. minFrames therefore
* dropped below the inversion's threshold (700 vs 900): where both fire,
* parallel wins over the inversion's single-worker pick (+17–21% at 700f).
*/
export function shouldPreferParallelDrawElement(args: {
workerCount: number;
Expand All @@ -1361,13 +1366,24 @@ export function shouldPreferParallelDrawElement(args: {
experimentalParallelDeOptIn: boolean;
/** HF_DE_PARALLEL_ROUTER === "true" — the router's own kill switch, default off. */
routerEnabled: boolean;
/**
* Whether verified parallel DE STREAMING can actually run for this render
* (`shouldUseStreamingEncode` at the router's worker count with
* forceParallelStream). The router's entire value is that path; without it
* firing would pin workerCount to 3 and skip calibration while delivering
* none of the benefit — e.g. a composition longer than
* `streamingEncodeMaxDurationSeconds` (240 s default), where the duration
* cap disables streaming before the router's force flag is consulted.
*/
parallelStreamingAvailable: boolean;
/** Machine RAM (os.totalmem, MB). */
totalMemoryMb: number;
/** RAM floor for routing; <=0 disables the guard. */
minMemoryMb: number;
}): boolean {
return (
args.routerEnabled &&
args.parallelStreamingAvailable &&
args.workerCount > 1 &&
typeof args.requestedWorkers !== "number" &&
args.useDrawElement &&
Expand Down Expand Up @@ -2414,18 +2430,27 @@ async function executeRenderPipeline(input: {
process.env.HF_DE_PARALLEL_STREAM === "true",
});
// DE parallel-router eligibility — see shouldPreferParallelDrawElement.
// Default-off (HF_DE_PARALLEL_ROUTER); HF_DE_PARALLEL_MIN_FRAMES defaults
// higher than the single-worker inversion's threshold since it targets
// the long tail the inversion's own benchmark didn't cover.
// Default-off (HF_DE_PARALLEL_ROUTER); HF_DE_PARALLEL_MIN_FRAMES default
// 700, re-calibrated 2026-07-27 from the original safe-high 2000. A
// controlled frame-count sweep (fixed content-per-frame, three synthetic
// profiles × {350..3000f} × {single,par2,par3} × 3 reps, worker counts +
// capture modes verified per run) showed par3 beating single at EVERY
// size in every profile — +17–21% at 700f rising to +28–34% at 3000f —
// including a 24-sub-composition profile built to reproduce the
// "workers re-pay init" failure (92k tweens, ~2.5s
// pollSubCompositionTimelines per worker): workers init CONCURRENTLY, so
// duplicated init costs CPU, not wall-clock. Below ~700f the win thins
// toward ~+10% while still paying 3 hardware-GPU browsers, so the floor
// stays. Harness: plans/drawelement-fast-capture/de-crossover-bench.sh.
const deParallelRouterEnabled = process.env.HF_DE_PARALLEL_ROUTER === "true";
const deParallelMinFramesRaw = process.env.HF_DE_PARALLEL_MIN_FRAMES;
const deParallelMinFramesNum =
deParallelMinFramesRaw === undefined || deParallelMinFramesRaw.trim() === ""
? 2000
? 700
: Number(deParallelMinFramesRaw);
const deParallelMinFrames = Number.isFinite(deParallelMinFramesNum)
? deParallelMinFramesNum
: 2000;
: 700;
// RAM floor default 24 GB: the wild black-slab report was a 16 GB
// machine; every clean routed cohort in telemetry so far is >=24 GB.
// HF_DE_PARALLEL_MIN_MEM_MB overrides (0 disables the guard).
Expand Down Expand Up @@ -2456,6 +2481,15 @@ async function executeRenderPipeline(input: {
process.env.PRODUCER_EXPERIMENTAL_FAST_CAPTURE === "true" ||
process.env.HF_DE_PARALLEL_STREAM === "true",
routerEnabled: deParallelRouterEnabled,
// Router pins 3 workers for the streaming path; don't pin when the
// duration cap (or any other streaming gate) would turn that path off.
parallelStreamingAvailable: shouldUseStreamingEncode(
cfg,
outputFormat,
3,
job.duration,
true,
),
totalMemoryMb: Math.round(totalmem() / (1024 * 1024)),
minMemoryMb: deParallelMinMemoryMb,
});
Expand Down
Loading