Skip to content
Merged
Show file tree
Hide file tree
Changes from 5 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
2 changes: 2 additions & 0 deletions apps/desktop/src/settings/DesktopClientSettings.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,8 @@ const clientSettings: ClientSettings = {
sidebarThreadSortOrder: "created_at",
sidebarThreadPreviewCount: 6,
legacySidebarEnabled: false,
loadBalancingEnabled: false,
loadBalancingWeights: { "environment-1": 75, "environment-2": 0 },
timestampFormat: "24-hour",
wordWrap: true,
};
Expand Down
1 change: 1 addition & 0 deletions apps/server/src/auth/RpcAuthorization.ts
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@ export const RPC_REQUIRED_SCOPES = {
[WS_METHODS.serverDiscoverSourceControl]: AuthOrchestrationReadScope,
[WS_METHODS.serverGetTraceDiagnostics]: AuthOrchestrationReadScope,
[WS_METHODS.serverGetProcessDiagnostics]: AuthOrchestrationReadScope,
[WS_METHODS.serverGetHostResources]: AuthOrchestrationReadScope,
[WS_METHODS.serverGetProcessResourceHistory]: AuthOrchestrationReadScope,
[WS_METHODS.serverGetResourceTelemetryHistory]: AuthOrchestrationReadScope,
[WS_METHODS.serverRetryResourceTelemetry]: AuthOrchestrationOperateScope,
Expand Down
88 changes: 88 additions & 0 deletions apps/server/src/resourceTelemetry/HostResources.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,88 @@
import * as NodeOS from "node:os";
import type { HostResourcesSnapshot } from "@t3tools/contracts";
import { HostProcessPlatform } from "@t3tools/shared/hostProcess";
import * as Context from "effect/Context";
import * as DateTime from "effect/DateTime";
import * as Effect from "effect/Effect";
import * as FileSystem from "effect/FileSystem";
import * as Layer from "effect/Layer";
import { ChildProcess, ChildProcessSpawner } from "effect/unstable/process";

export class HostResources extends Context.Service<
HostResources,
{ readonly read: Effect.Effect<HostResourcesSnapshot> }
>()("t3/resourceTelemetry/HostResources") {}

function readCpu() {
const cpus = NodeOS.cpus();
const cpu = cpus.reduce(
(sum, { times }) => ({
idle: sum.idle + times.idle,
total: sum.total + times.user + times.nice + times.sys + times.idle + times.irq,
}),
{ idle: 0, total: 0 },
);
return { ...cpu, count: cpus.length };
}

function darwinAvailableMemory(output: string): number | null {
const pageSize = /page size of (\d+) bytes/.exec(output)?.[1];
const free = /^Pages free:\s+(\d+)\./m.exec(output)?.[1];
const inactive = /^Pages inactive:\s+(\d+)\./m.exec(output)?.[1];
const speculative = /^Pages speculative:\s+(\d+)\./m.exec(output)?.[1];
if (!pageSize || !free || !inactive || !speculative) return null;
// vm_stat subtracts speculative pages from its printed "Pages free" count.
// Adding them here counts each reclaimable page once; purgeable pages overlap.
const available = (Number(free) + Number(inactive) + Number(speculative)) * Number(pageSize);
return Number.isSafeInteger(available) && Number(pageSize) > 0 ? available : null;
}

export const make = Effect.fn("makeHostResources")(function* () {
const fs = yield* FileSystem.FileSystem;
const platform = yield* HostProcessPlatform;
const spawner = yield* ChildProcessSpawner.ChildProcessSpawner;

const sample = Effect.fn("HostResources.sample")(function* () {
const previousCpu = readCpu();
// CPU counters need two readings; idle servers do no polling or process scans.
yield* Effect.sleep("200 millis");
const cpu = readCpu();
const totalDelta = cpu.total - previousCpu.total;
const idleDelta = cpu.idle - previousCpu.idle;
const cpuUtilization =
previousCpu.count === cpu.count && totalDelta > 0 && idleDelta >= 0
? Math.min(1, Math.max(0, 1 - idleDelta / totalDelta))
: null;
const totalMemoryBytes = NodeOS.totalmem();
// On Windows libuv returns GlobalMemoryStatusEx.ullAvailPhys, including standby memory.
let availableMemoryBytes = NodeOS.freemem();
if (platform === "linux") {
const meminfo = yield* fs
.readFileString("/proc/meminfo")
.pipe(Effect.catch(() => Effect.succeed("")));
const available = /^MemAvailable:\s+(\d+)\s+kB$/m.exec(meminfo)?.[1];
if (available) availableMemoryBytes = Number(available) * 1024;
} else if (platform === "darwin") {
const output = yield* spawner
.string(ChildProcess.make("/usr/bin/vm_stat", [], { stdin: "ignore", stderr: "ignore" }))
.pipe(
Effect.timeout("1 second"),
Effect.catch(() => Effect.succeed("")),
);
availableMemoryBytes = darwinAvailableMemory(output) ?? availableMemoryBytes;
}
return {
sampledAt: DateTime.toEpochMillis(yield* DateTime.now),
cpuUtilization,
cpuCount: cpu.count,
availableMemoryBytes: Math.min(totalMemoryBytes, Math.max(0, availableMemoryBytes)),
totalMemoryBytes,
};
});

// One server-lifetime cache deduplicates simultaneous requests from all sockets.
const read = yield* Effect.cachedWithTTL(sample(), "5 seconds");
return HostResources.of({ read });
Comment thread
cursor[bot] marked this conversation as resolved.
Outdated
});

export const layer = Layer.effect(HostResources, make());
62 changes: 60 additions & 2 deletions apps/server/src/server.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -161,6 +161,7 @@ import * as PairingGrantStore from "./auth/PairingGrantStore.ts";
import * as CloudManagedEndpointRuntime from "./cloud/ManagedEndpointRuntime.ts";
import * as CloudCliTokenManager from "./cloud/CliTokenManager.ts";
import * as ProcessDiagnostics from "./diagnostics/ProcessDiagnostics.ts";
import * as HostResources from "./resourceTelemetry/HostResources.ts";
import * as ProcessResourceMonitor from "./diagnostics/ProcessResourceMonitor.ts";
import * as TraceDiagnostics from "./diagnostics/TraceDiagnostics.ts";
import * as DesktopTelemetryReceiver from "./resourceTelemetry/DesktopTelemetryReceiver.ts";
Expand Down Expand Up @@ -829,7 +830,8 @@ const buildAppUnderTest = (options?: {
}),
}),
),
Layer.provide(
Layer.provide([
HostResources.layer,
Layer.mock(ProcessResourceMonitor.ProcessResourceMonitor)({
readHistory: (input) =>
Effect.succeed({
Expand All @@ -844,7 +846,7 @@ const buildAppUnderTest = (options?: {
error: Option.none(),
}),
}),
),
]),
Layer.provide(
Layer.mock(TraceDiagnostics.TraceDiagnostics)({
read: () =>
Expand Down Expand Up @@ -6040,6 +6042,62 @@ it.layer(NodeServices.layer)("server router seam", (it) => {
}).pipe(Effect.provide(NodeHttpServer.layerTest)),
);

it.effect("returns cached whole-host resources over websocket", () =>
Effect.gen(function* () {
yield* buildAppUnderTest();
const wsUrl = yield* getWsServerUrl("/ws");
const [first, second] = yield* Effect.scoped(
withWsRpcClient(wsUrl, (client) =>
Effect.all(
[
client[WS_METHODS.serverGetHostResources]({}),
client[WS_METHODS.serverGetHostResources]({}),
],
{ concurrency: "unbounded" },
),
),
);
assert.deepEqual(first, second);
assert.isAtLeast(first.sampledAt, 0);
assert.isAbove(first.cpuCount, 0);
assert.isAbove(first.totalMemoryBytes, 0);
assert.isAtLeast(first.availableMemoryBytes, 0);
assert.isAtMost(first.availableMemoryBytes, first.totalMemoryBytes);
if (first.cpuUtilization !== null) {
assert.isAtLeast(first.cpuUtilization, 0);
assert.isAtMost(first.cpuUtilization, 1);
}
}).pipe(Effect.provide(NodeHttpServer.layerTest), TestClock.withLive),
);

it.effect("counts macOS reclaimable memory once and shares concurrent samples", () =>
Effect.gen(function* () {
const commandCalls = yield* Ref.make(0);
const hostResources = yield* HostResources.make().pipe(
Effect.provideService(HostProcessPlatform, "darwin"),
Effect.provide(
Layer.mock(ChildProcessSpawner.ChildProcessSpawner)({
string: () =>
Ref.update(commandCalls, (count) => count + 1).pipe(
Effect.as(
"Mach Virtual Memory Statistics: (page size of 16384 bytes)\n" +
"Pages free: 10.\nPages inactive: 20.\nPages speculative: 5.\n" +
"Pages purgeable: 999.\n",
),
),
}),
),
);
const [first, second] = yield* Effect.all([hostResources.read, hostResources.read], {
concurrency: "unbounded",
});
assert.equal(first.availableMemoryBytes, 35 * 16384);
assert.deepEqual(first, second);
assert.deepEqual(yield* hostResources.read, first);
assert.equal(yield* Ref.get(commandCalls), 1);
}).pipe(TestClock.withLive),
);

it.effect("routes websocket resource telemetry through the subscription", () =>
Effect.gen(function* () {
yield* buildAppUnderTest();
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -115,6 +115,7 @@ import * as ServerSelfUpdate from "./cloud/selfUpdate.ts";
import * as DesktopAppUpdate from "./desktopUpdate/DesktopAppUpdate.ts";
import * as ServiceLauncherClient from "./cloud/serviceLauncherClient.ts";
import * as ProcessDiagnostics from "./diagnostics/ProcessDiagnostics.ts";
import * as HostResources from "./resourceTelemetry/HostResources.ts";
import * as ProcessResourceMonitor from "./diagnostics/ProcessResourceMonitor.ts";
import * as TraceDiagnostics from "./diagnostics/TraceDiagnostics.ts";
import * as DesktopTelemetryReceiver from "./resourceTelemetry/DesktopTelemetryReceiver.ts";
Expand Down Expand Up @@ -199,6 +200,7 @@ const BackgroundLayerLive = BackgroundPolicy.layer.pipe(
const UsageLayerLive = UsageService.layer.pipe(Layer.provide(ServerSettingsLayerLive));

const ResourceDiagnosticsLayerLive = Layer.mergeAll(
HostResources.layer,
ResourceTelemetryLayerLive,
ProcessDiagnostics.layer.pipe(Layer.provide(ResourceTelemetryLayerLive)),
ProcessResourceMonitor.layer.pipe(Layer.provide(ResourceTelemetryLayerLive)),
Expand Down
6 changes: 6 additions & 0 deletions apps/server/src/ws.ts
Original file line number Diff line number Diff line change
Expand Up @@ -127,6 +127,7 @@ import { requiredScopeForRpcMethod } from "./auth/RpcAuthorization.ts";
import * as ProcessDiagnostics from "./diagnostics/ProcessDiagnostics.ts";
import * as ProcessResourceMonitor from "./diagnostics/ProcessResourceMonitor.ts";
import * as ResourceTelemetry from "./resourceTelemetry/ResourceTelemetry.ts";
import * as HostResources from "./resourceTelemetry/HostResources.ts";
import * as AnalyticsService from "./telemetry/AnalyticsService.ts";
import * as UsageLimitSources from "./usage/UsageLimitSources.ts";
import * as UsageService from "./usage/UsageService.ts";
Expand Down Expand Up @@ -595,6 +596,7 @@ const makeWsRpcLayer = (
const bootstrapCredentials = yield* PairingGrantStore.PairingGrantStore;
const sessions = yield* SessionStore.SessionStore;
const processDiagnostics = yield* ProcessDiagnostics.ProcessDiagnostics;
const hostResources = yield* HostResources.HostResources;
const processResourceMonitor = yield* ProcessResourceMonitor.ProcessResourceMonitor;
const resourceTelemetry = yield* ResourceTelemetry.ResourceTelemetry;
const usage = yield* UsageService.UsageService;
Expand Down Expand Up @@ -1998,6 +2000,10 @@ const makeWsRpcLayer = (
observeRpcEffect(WS_METHODS.serverGetProcessDiagnostics, processDiagnostics.read, {
"rpc.aggregate": "server",
}),
[WS_METHODS.serverGetHostResources]: (_input) =>
observeRpcEffect(WS_METHODS.serverGetHostResources, hostResources.read, {
"rpc.aggregate": "server",
}),
[WS_METHODS.serverGetProcessResourceHistory]: (input) =>
observeRpcEffect(
WS_METHODS.serverGetProcessResourceHistory,
Expand Down
34 changes: 31 additions & 3 deletions apps/web/src/components/BranchToolbar.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,8 @@ interface BranchToolbarProps {
onActiveThreadBranchOverrideChange?: (branch: string | null) => void;
startFromOrigin: boolean;
onStartFromOriginChange: (startFromOrigin: boolean) => void;
autoEnvironmentLabel?: string | undefined;
onAutoEnvironment?: (() => void) | undefined;
envLocked: boolean;
onCheckoutPullRequestRequest?: (reference: string) => void;
onComposerFocusRequest?: () => void;
Expand All @@ -66,6 +68,8 @@ interface BranchToolbarProps {
}

interface MobileRunContextSelectorProps {
autoEnvironmentLabel?: string | undefined;
onAutoEnvironment?: (() => void) | undefined;
envLocked: boolean;
envModeLocked: boolean;
environmentId: EnvironmentId;
Expand All @@ -81,6 +85,8 @@ interface MobileRunContextSelectorProps {
}

const MobileRunContextSelector = memo(function MobileRunContextSelector({
autoEnvironmentLabel,
onAutoEnvironment,
envLocked,
envModeLocked,
environmentId,
Expand Down Expand Up @@ -134,7 +140,8 @@ const MobileRunContextSelector = memo(function MobileRunContextSelector({
data-composer-label-motion
className="block w-full min-w-0 max-w-[240px] origin-left truncate transition-[opacity,transform] duration-180 ease-[cubic-bezier(0.32,0.72,0,1)] group-data-[compact]/composer-context:[transform:translateX(-0.25rem)_scaleX(0.95)] group-data-[compact]/composer-context:opacity-0 motion-reduce:transform-none motion-reduce:transition-opacity"
>
{showEnvironmentIndicator ? (activeEnvironment?.label ?? "Run on") : workspaceLabel}
{autoEnvironmentLabel ??
(showEnvironmentIndicator ? (activeEnvironment?.label ?? "Run on") : workspaceLabel)}
</span>
</span>
</>
Expand Down Expand Up @@ -167,9 +174,24 @@ const MobileRunContextSelector = memo(function MobileRunContextSelector({
<MenuGroup>
<MenuGroupLabel>Run on</MenuGroupLabel>
<MenuRadioGroup
value={environmentId}
onValueChange={(value) => onEnvironmentChange(value as EnvironmentId)}
value={autoEnvironmentLabel ? "auto" : environmentId}
onValueChange={(value) =>
value === "auto"
? onAutoEnvironment?.()
: onEnvironmentChange(value as EnvironmentId)
}
>
{onAutoEnvironment && (
<MenuRadioItem
value="auto"
disabled={envLocked}
onClick={() => {
if (autoEnvironmentLabel) onAutoEnvironment?.();
}}
>
{autoEnvironmentLabel ?? "Balance load"}
</MenuRadioItem>
)}
{availableEnvironments.map((env) => (
<MenuRadioItem
key={env.environmentId}
Expand Down Expand Up @@ -421,6 +443,8 @@ export const BranchToolbar = memo(function BranchToolbar({
onActiveThreadBranchOverrideChange,
startFromOrigin,
onStartFromOriginChange,
autoEnvironmentLabel,
onAutoEnvironment,
envLocked,
onCheckoutPullRequestRequest,
onComposerFocusRequest,
Expand Down Expand Up @@ -518,6 +542,8 @@ export const BranchToolbar = memo(function BranchToolbar({
{showGitControls ? (
<div className="contents @3xl/composer-surface:hidden">
<MobileRunContextSelector
autoEnvironmentLabel={autoEnvironmentLabel}
onAutoEnvironment={onAutoEnvironment}
envLocked={envLocked}
envModeLocked={envModeLocked}
environmentId={environmentId}
Expand All @@ -544,6 +570,8 @@ export const BranchToolbar = memo(function BranchToolbar({
{showEnvironmentIndicator && availableEnvironments && (
<>
<BranchToolbarEnvironmentSelector
autoEnvironmentLabel={autoEnvironmentLabel}
onAutoEnvironment={onAutoEnvironment}
envLocked={envLocked}
environmentId={environmentId}
availableEnvironments={availableEnvironments}
Expand Down
6 changes: 4 additions & 2 deletions apps/web/src/components/BranchToolbarBranchSelector.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -154,7 +154,7 @@ export function BranchToolbarBranchSelector({
// Thread branch mutation (colocated — only this component calls it)
// ---------------------------------------------------------------------------
const setThreadBranch = useCallback(
(branch: string | null, worktreePath: string | null) => {
(branch: string | null, worktreePath: string | null, automatic = false) => {
if (!activeThreadId || !activeProject) return;
if (serverSession && worktreePath !== activeWorktreePath) {
void stopThreadSession({
Expand Down Expand Up @@ -185,6 +185,7 @@ export function BranchToolbarBranchSelector({
branch,
worktreePath,
envMode: nextDraftEnvMode,
environmentSelection: automatic ? (draftThread?.environmentSelection ?? "auto") : "manual",
Comment thread
cursor[bot] marked this conversation as resolved.
projectRef: scopeProjectRef(environmentId, activeProject.id),
});
},
Expand All @@ -200,6 +201,7 @@ export function BranchToolbarBranchSelector({
threadRef,
environmentId,
effectiveEnvMode,
draftThread?.environmentSelection,
stopThreadSession,
updateThreadMetadata,
],
Expand Down Expand Up @@ -506,7 +508,7 @@ export function BranchToolbarBranchSelector({
) {
return;
}
setThreadBranch(worktreeBaseBranchCandidate, null);
setThreadBranch(worktreeBaseBranchCandidate, null, true);
}, [
activeThreadBranch,
activeWorktreePath,
Expand Down
Loading
Loading