diff --git a/apps/server/src/diagnostics/ProcessDiagnostics.test.ts b/apps/server/src/diagnostics/ProcessDiagnostics.test.ts index 2efa3375d275..4e9f676379f5 100644 --- a/apps/server/src/diagnostics/ProcessDiagnostics.test.ts +++ b/apps/server/src/diagnostics/ProcessDiagnostics.test.ts @@ -19,7 +19,7 @@ function makeNativeSnapshot( processes: ResourceMonitorSnapshotEvent["processes"], ): ResourceMonitorSnapshotEvent { return { - version: 2, + version: 4, type: "snapshot", sequence: 1, sampledAtUnixMs: DateTime.toEpochMillis(DateTime.makeUnsafe("2026-05-05T10:00:00.000Z")), diff --git a/apps/server/src/preview/PortScanner.test.ts b/apps/server/src/preview/PortScanner.test.ts index 7fa15defeca9..d88399c8252e 100644 --- a/apps/server/src/preview/PortScanner.test.ts +++ b/apps/server/src/preview/PortScanner.test.ts @@ -22,6 +22,7 @@ import { expect } from "vite-plus/test"; import { FetchHttpClient } from "effect/unstable/http"; import * as ProcessRunner from "../processRunner.ts"; +import * as NativeTelemetryClient from "../resourceTelemetry/NativeTelemetryClient.ts"; import * as PortScanner from "./PortScanner.ts"; const processProbeFailure: ProcessRunner.ProcessRunner["Service"]["run"] = (input) => Effect.fail( @@ -41,6 +42,7 @@ const processProbeFailure: ProcessRunner.ProcessRunner["Service"]["run"] = (inpu const TestProcessRunner = Layer.succeed(ProcessRunner.ProcessRunner, { run: processProbeFailure, }); +const TestNativeTelemetry = NativeTelemetryClient.layerTest(); let integrationListeningPort: number | null = null; @@ -68,6 +70,7 @@ const makeProbeFailureLayer = ( findAvailablePort: (preferred) => Effect.succeed(preferred), }), Layer.succeed(HostProcessPlatform, "linux"), + TestNativeTelemetry, FetchHttpClient.layer.pipe(Layer.provide(Layer.succeed(FetchHttpClient.Fetch, fetch))), ), ), @@ -79,6 +82,7 @@ const TestPortDiscoveryLive = PortScanner.layer.pipe( TestProcessRunner, TestIntegrationNet, Layer.succeed(HostProcessPlatform, "win32"), + TestNativeTelemetry, FetchHttpClient.layer, ), ), @@ -114,6 +118,7 @@ const makeLsofScannerLayer = (input: { findAvailablePort: (preferred) => Effect.succeed(preferred), }), Layer.succeed(HostProcessPlatform, "linux"), + TestNativeTelemetry, FetchHttpClient.layer.pipe( Layer.provide(Layer.succeed(FetchHttpClient.Fetch, input.fetch)), ), @@ -121,6 +126,31 @@ const makeLsofScannerLayer = (input: { ), ); +const makeWindowsScannerLayer = (input: { + readonly windowsListeners: NativeTelemetryClient.NativeTelemetryClient["Service"]["windowsListeners"]; + readonly run: ProcessRunner.ProcessRunner["Service"]["run"]; + readonly fetch?: typeof globalThis.fetch; +}) => + PortScanner.layer.pipe( + Layer.provide( + Layer.mergeAll( + Layer.succeed(ProcessRunner.ProcessRunner, { run: input.run }), + Layer.succeed(Net.NetService, { + canListenOnHost: () => Effect.succeed(true), + isPortAvailableOnLoopback: () => Effect.succeed(true), + hasListenerOnHost: () => Effect.succeed(false), + reserveLoopbackPort: () => Effect.succeed(40_000), + findAvailablePort: (preferred) => Effect.succeed(preferred), + }), + Layer.succeed(HostProcessPlatform, "win32"), + NativeTelemetryClient.layerTest({ windowsListeners: input.windowsListeners }), + FetchHttpClient.layer.pipe( + Layer.provide(Layer.succeed(FetchHttpClient.Fetch, input.fetch ?? globalThis.fetch)), + ), + ), + ), + ); + const openServer = ( port: number, onConnection: (socket: NodeNet.Socket) => void, @@ -244,6 +274,104 @@ effectIt.layer(TestPortDiscoveryLive)("PortDiscovery integration (TCP probe fall ); }); +effectIt.effect("uses native Windows listeners without spawning PowerShell", () => { + let fallbackRuns = 0; + const layer = makeWindowsScannerLayer({ + windowsListeners: Effect.succeed([{ port: LSOF_TEST_PORT, pid: 4_242, processName: "node" }]), + run: (input) => { + fallbackRuns += 1; + return processProbeFailure(input); + }, + fetch: ((_input: Parameters[0]) => + Promise.resolve( + new Response("app", { headers: { "content-type": "text/html" } }), + )) as typeof globalThis.fetch, + }); + + return Effect.gen(function* () { + const scanner = yield* PortScanner.PortDiscovery; + yield* scanner.registerTerminalProcesses({ + threadId: "thread-1", + terminalId: "default", + processIds: [4_242], + }); + const servers = yield* scanner.scan(); + + expect(fallbackRuns).toBe(0); + expect(servers).toEqual([ + { + host: "localhost", + port: LSOF_TEST_PORT, + url: `http://localhost:${LSOF_TEST_PORT}`, + processName: "node", + pid: 4_242, + terminal: { threadId: "thread-1", terminalId: "default" }, + }, + ]); + }).pipe(Effect.provide(layer)); +}); + +effectIt.effect("backs off every failed Windows PowerShell fallback", () => { + let fallbackRuns = 0; + const layer = makeWindowsScannerLayer({ + windowsListeners: Effect.fail( + new NativeTelemetryClient.NativeTelemetryUnavailable({ reason: "test" }), + ), + run: (input) => { + fallbackRuns += 1; + return processProbeFailure(input); + }, + }); + + return Effect.gen(function* () { + const scanner = yield* PortScanner.PortDiscovery; + yield* scanner.scan(); + yield* scanner.scan(); + expect(fallbackRuns).toBe(1); + + yield* TestClock.adjust(Duration.seconds(3)); + yield* scanner.scan(); + expect(fallbackRuns).toBe(2); + }).pipe(Effect.provide(layer)); +}); + +effectIt.effect("keeps the last Windows fallback snapshot when a retry fails", () => { + let fallbackRuns = 0; + const layer = makeWindowsScannerLayer({ + windowsListeners: Effect.fail( + new NativeTelemetryClient.NativeTelemetryUnavailable({ reason: "test" }), + ), + run: (input) => { + fallbackRuns += 1; + if (fallbackRuns > 1) return processProbeFailure(input); + return Effect.succeed({ + stdout: `127.0.0.1|${LSOF_TEST_PORT}|4242|node\n`, + stderr: "", + code: null, + timedOut: false, + stdoutTruncated: false, + stderrTruncated: false, + stdoutInvalidUtf8: false, + stderrInvalidUtf8: false, + }); + }, + fetch: ((_input: Parameters[0]) => + Promise.resolve( + new Response("app", { headers: { "content-type": "text/html" } }), + )) as typeof globalThis.fetch, + }); + + return Effect.gen(function* () { + const scanner = yield* PortScanner.PortDiscovery; + expect(yield* scanner.scan()).toHaveLength(1); + + yield* TestClock.adjust(Duration.seconds(3)); + expect(yield* scanner.scan()).toHaveLength(1); + expect(yield* scanner.scan()).toHaveLength(1); + expect(fallbackRuns).toBe(2); + }).pipe(Effect.provide(layer)); +}); + effectIt.effect("revalidates a successful HTML probe after its cache entry expires", () => { let responds = true; const requests: string[] = []; diff --git a/apps/server/src/preview/PortScanner.ts b/apps/server/src/preview/PortScanner.ts index 4571aeef4c6b..a3098699db19 100644 --- a/apps/server/src/preview/PortScanner.ts +++ b/apps/server/src/preview/PortScanner.ts @@ -5,8 +5,11 @@ * stable line-prefixed field format; this is the only `lsof` flag set we rely * on). * - * Windows / lsof missing: checks a curated list of common dev ports through - * the shared Net service. + * Windows: asks the persistent resource monitor for the native TCP listener + * table. If the sidecar is unavailable, a backed-off PowerShell probe runs. + * + * lsof / Windows probe missing: checks a curated list of common dev ports + * through the shared Net service. * * Listening ports are published only after a bounded HTTP(S) probe finds a * successful HTML document or a redirect to one. @@ -21,6 +24,7 @@ import { PREVIEW_URL_MAX_LENGTH, ThreadId, type DiscoveredLocalServer, + type ResourceMonitorWindowsListener, } from "@t3tools/contracts"; import { HostProcessPlatform } from "@t3tools/shared/hostProcess"; import * as Net from "@t3tools/shared/Net"; @@ -39,6 +43,7 @@ import * as Semaphore from "effect/Semaphore"; import { FetchHttpClient, HttpClient } from "effect/unstable/http"; import * as ProcessRunner from "../processRunner.ts"; +import * as NativeTelemetryClient from "../resourceTelemetry/NativeTelemetryClient.ts"; export class PortDiscovery extends Context.Service< PortDiscovery, @@ -73,6 +78,9 @@ export const COMMON_DEV_PORTS: ReadonlyArray = Object.freeze([ const POLL_INTERVAL = Duration.seconds(3); const LSOF_TIMEOUT_MS = 5_000; const WINDOWS_LISTENER_TIMEOUT_MS = 5_000; +const WINDOWS_FALLBACK_MAX_RETRY_MS = 60_000; +export const WINDOWS_LISTENER_COMMAND = + '$m = @{}; Get-Process | ForEach-Object { $m[$_.Id] = $_.ProcessName }; Get-NetTCPConnection -State Listen -ErrorAction Stop | ForEach-Object { Write-Output "$($_.LocalAddress)|$($_.LocalPort)|$($_.OwningProcess)|$($m[[int]$_.OwningProcess])" }'; const WEB_PROBE_TIMEOUT = Duration.seconds(1); const WEB_PROBE_CACHE_TTL_MS = Duration.toMillis(Duration.seconds(15)); const WEB_PROBE_CONCURRENCY = 16; @@ -265,6 +273,29 @@ const parseWindowsListenerOutput = ( return [...seen.values()].toSorted((left, right) => left.port - right.port); }; +const windowsListenersToServers = ( + listeners: ReadonlyArray, + terminalByProcessId: ReadonlyMap = new Map(), +): ReadonlyArray => { + const seen = new Map(); + for (const listener of listeners) { + if (seen.has(listener.port)) continue; + seen.set(listener.port, { + host: "localhost", + port: listener.port, + url: `http://localhost:${listener.port}`, + processName: listener.processName?.trim() || null, + pid: listener.pid, + terminal: terminalByProcessId.get(listener.pid) ?? null, + }); + } + return [...seen.values()].toSorted((left, right) => left.port - right.port); +}; + +export function windowsFallbackRetryDelayMs(failureCount: number): number { + return Math.min(3_000 * 2 ** Math.max(0, failureCount - 1), WINDOWS_FALLBACK_MAX_RETRY_MS); +} + const serversEqual = ( left: ReadonlyArray, right: ReadonlyArray, @@ -292,6 +323,7 @@ const serversEqual = ( export const make = Effect.gen(function* PortDiscoveryMake() { const net = yield* Net.NetService; const processRunner = yield* ProcessRunner.ProcessRunner; + const nativeTelemetry = yield* NativeTelemetryClient.NativeTelemetryClient; const hostPlatform = yield* HostProcessPlatform; const httpClient = (yield* HttpClient.HttpClient).pipe(HttpClient.withScope); const stateRef = yield* Ref.make({ @@ -301,6 +333,11 @@ export const make = Effect.gen(function* PortDiscoveryMake() { }); const webProbeCacheRef = yield* Ref.make>(new Map()); const scanSemaphore = yield* Semaphore.make(1); + const windowsFallbackRef = yield* Ref.make({ + failureCount: 0, + nextAttemptAtMillis: 0, + lastSnapshot: null as ReadonlyArray | null, + }); const probeCommonPorts = Effect.fn("PortDiscovery.probeCommonPorts")(function* () { const results = yield* Effect.forEach( @@ -477,6 +514,52 @@ export const make = Effect.gen(function* PortDiscoveryMake() { platform: hostPlatform, }).pipe(Effect.as(null)); + const probeWindowsFallback = Effect.fn("PortDiscovery.probeWindowsFallback")(function* ( + terminalByProcessId: ReadonlyMap, + ) { + const nowMillis = yield* Clock.currentTimeMillis; + const fallback = yield* Ref.get(windowsFallbackRef); + if (nowMillis < fallback.nextAttemptAtMillis) { + return fallback.lastSnapshot ?? (yield* probeCommonPorts()); + } + + const recoverWindowsProbeFailure = recoverProcessProbeFailure("windows-listeners"); + const listeners = yield* processRunner + .run({ + command: "powershell.exe", + args: ["-NoProfile", "-NonInteractive", "-Command", WINDOWS_LISTENER_COMMAND], + timeout: Duration.millis(WINDOWS_LISTENER_TIMEOUT_MS), + maxOutputBytes: 1024 * 1024, + outputMode: "truncate", + }) + .pipe( + Effect.map((result) => parseWindowsListenerOutput(result.stdout, terminalByProcessId)), + Effect.catchTags({ + ProcessSpawnError: recoverWindowsProbeFailure, + ProcessStdinError: recoverWindowsProbeFailure, + ProcessOutputLimitError: recoverWindowsProbeFailure, + ProcessReadError: recoverWindowsProbeFailure, + ProcessTimeoutError: recoverWindowsProbeFailure, + }), + ); + if (listeners !== null) { + yield* Ref.set(windowsFallbackRef, { + failureCount: 0, + nextAttemptAtMillis: 0, + lastSnapshot: listeners, + }); + return listeners; + } + + const failureCount = fallback.failureCount + 1; + yield* Ref.set(windowsFallbackRef, { + failureCount, + nextAttemptAtMillis: nowMillis + windowsFallbackRetryDelayMs(failureCount), + lastSnapshot: fallback.lastSnapshot, + }); + return fallback.lastSnapshot ?? (yield* probeCommonPorts()); + }); + const scanUnlocked = Effect.fn("PortDiscovery.scanUnlocked")(function* ( configuredUrls: ReadonlyArray, ) { @@ -488,29 +571,19 @@ export const make = Effect.gen(function* PortDiscoveryMake() { } } if (hostPlatform === "win32") { - const recoverWindowsProbeFailure = recoverProcessProbeFailure("windows-listeners"); - const command = - 'Get-NetTCPConnection -State Listen -ErrorAction Stop | ForEach-Object { $processName = (Get-Process -Id $_.OwningProcess -ErrorAction SilentlyContinue).ProcessName; Write-Output "$($_.LocalAddress)|$($_.LocalPort)|$($_.OwningProcess)|$processName" }'; - const listeners = yield* processRunner - .run({ - command: "powershell.exe", - args: ["-NoProfile", "-NonInteractive", "-Command", command], - timeout: Duration.millis(WINDOWS_LISTENER_TIMEOUT_MS), - maxOutputBytes: 1024 * 1024, - outputMode: "truncate", - }) - .pipe( - Effect.map((result) => parseWindowsListenerOutput(result.stdout, terminalByProcessId)), - Effect.catchTags({ - ProcessSpawnError: recoverWindowsProbeFailure, - ProcessStdinError: recoverWindowsProbeFailure, - ProcessOutputLimitError: recoverWindowsProbeFailure, - ProcessReadError: recoverWindowsProbeFailure, - ProcessTimeoutError: recoverWindowsProbeFailure, - }), - ); - if (listeners !== null) return yield* probeWebServers(listeners, configuredUrls); - return yield* probeWebServers(yield* probeCommonPorts(), configuredUrls); + const nativeListeners = yield* nativeTelemetry.windowsListeners.pipe( + Effect.map((listeners) => windowsListenersToServers(listeners, terminalByProcessId)), + Effect.catch((cause) => + Effect.logDebug("native Windows listener discovery failed; using fallback", { + cause, + }).pipe(Effect.as(null)), + ), + ); + const listeners = + nativeListeners === null + ? yield* probeWindowsFallback(terminalByProcessId) + : nativeListeners; + return yield* probeWebServers(listeners, configuredUrls); } const recoverLsofProbeFailure = recoverProcessProbeFailure("lsof"); const lsofResult = yield* processRunner diff --git a/apps/server/src/resourceTelemetry/Model.test.ts b/apps/server/src/resourceTelemetry/Model.test.ts index 94690e3967bc..6247215854d6 100644 --- a/apps/server/src/resourceTelemetry/Model.test.ts +++ b/apps/server/src/resourceTelemetry/Model.test.ts @@ -39,7 +39,7 @@ function nativeSnapshot( sequence = 1, ): ResourceMonitorSnapshotEvent { return { - version: 2, + version: 4, type: "snapshot", sequence, sampledAtUnixMs, diff --git a/apps/server/src/resourceTelemetry/NativeTelemetryClient.test.ts b/apps/server/src/resourceTelemetry/NativeTelemetryClient.test.ts index 8a595bc8b480..0fff8aba814c 100644 --- a/apps/server/src/resourceTelemetry/NativeTelemetryClient.test.ts +++ b/apps/server/src/resourceTelemetry/NativeTelemetryClient.test.ts @@ -83,7 +83,7 @@ describe("canCommandNativeTelemetrySidecar", () => { }); describe("NativeTelemetryRequestTimedOut", () => { - it("models history and sample request deadlines without a fabricated cause", () => { + it("models request deadlines without a fabricated cause", () => { const historyTimeout = new NativeTelemetryRequestTimedOut({ operation: "readHistory", timeoutMs: 15_000, @@ -92,6 +92,14 @@ describe("NativeTelemetryRequestTimedOut", () => { operation: "sampleNow", timeoutMs: 5_000, }); + const processTableTimeout = new NativeTelemetryRequestTimedOut({ + operation: "processTable", + timeoutMs: 5_000, + }); + const windowsListenersTimeout = new NativeTelemetryRequestTimedOut({ + operation: "windowsListeners", + timeoutMs: 5_000, + }); expect(historyTimeout.message).toBe( "Resource monitor 'readHistory' request timed out after 15000ms.", @@ -99,8 +107,16 @@ describe("NativeTelemetryRequestTimedOut", () => { expect(sampleTimeout.message).toBe( "Resource monitor 'sampleNow' request timed out after 5000ms.", ); + expect(processTableTimeout.message).toBe( + "Resource monitor 'processTable' request timed out after 5000ms.", + ); + expect(windowsListenersTimeout.message).toBe( + "Resource monitor 'windowsListeners' request timed out after 5000ms.", + ); expect("cause" in historyTimeout).toBe(false); expect("cause" in sampleTimeout).toBe(false); + expect("cause" in processTableTimeout).toBe(false); + expect("cause" in windowsListenersTimeout).toBe(false); }); }); diff --git a/apps/server/src/resourceTelemetry/NativeTelemetryClient.ts b/apps/server/src/resourceTelemetry/NativeTelemetryClient.ts index 232079d9dc9b..a54005672fae 100644 --- a/apps/server/src/resourceTelemetry/NativeTelemetryClient.ts +++ b/apps/server/src/resourceTelemetry/NativeTelemetryClient.ts @@ -5,7 +5,9 @@ import type { ResourceMonitorEvent, ResourceMonitorExternalProcess, ResourceMonitorHelloEvent, + ResourceMonitorProcessTableEntry, ResourceMonitorSnapshotEvent, + ResourceMonitorWindowsListener, ResourceTelemetrySourceStatus, } from "@t3tools/contracts"; import { @@ -44,6 +46,7 @@ const BATTERY_SAMPLE_INTERVAL_MS = 5_000; const CONSTRAINED_SAMPLE_INTERVAL_MS = 15_000; const HANDSHAKE_TIMEOUT = Duration.seconds(5); const SAMPLE_REQUEST_TIMEOUT = Duration.seconds(5); +const PROCESS_TABLE_REQUEST_TIMEOUT = Duration.seconds(5); const HISTORY_REQUEST_TIMEOUT = Duration.seconds(15); const INITIAL_RESTART_DELAY = Duration.millis(500); const MAX_RESTART_DELAY = Duration.seconds(10); @@ -76,7 +79,7 @@ export class NativeTelemetryHandshakeTimedOut extends Schema.TaggedErrorClass()( "NativeTelemetryRequestTimedOut", { - operation: Schema.Literals(["readHistory", "sampleNow"]), + operation: Schema.Literals(["processTable", "readHistory", "sampleNow", "windowsListeners"]), timeoutMs: Schema.Number, }, ) { @@ -192,6 +195,14 @@ export class NativeTelemetryClient extends Context.Service< snapshot: HostPowerSnapshot, ) => Effect.Effect; readonly sampleNow: Effect.Effect; + readonly processTable: Effect.Effect< + ReadonlyArray, + NativeTelemetryClientError + >; + readonly windowsListeners: Effect.Effect< + ReadonlyArray, + NativeTelemetryClientError + >; readonly retry: Effect.Effect; readonly health: Effect.Effect; readonly subscribeHealth: Effect.Effect< @@ -386,6 +397,18 @@ export const make = Effect.fn("resourceTelemetry.nativeTelemetryClient.make")(fu const pendingSamples = yield* Ref.make( new Map>(), ); + const pendingProcessTables = yield* Ref.make( + new Map< + string, + Deferred.Deferred, NativeTelemetryClientError> + >(), + ); + const pendingWindowsListeners = yield* Ref.make( + new Map< + string, + Deferred.Deferred, NativeTelemetryClientError> + >(), + ); const pendingHistories = yield* Ref.make(new Map()); const snapshots = yield* PubSub.sliding(8); const healthChanges = yield* PubSub.sliding(4); @@ -403,10 +426,20 @@ export const make = Effect.fn("resourceTelemetry.nativeTelemetryClient.make")(fu const failPending = (error: NativeTelemetryClientError) => Effect.gen(function* () { const samples = yield* Ref.getAndSet(pendingSamples, new Map()); + const processTables = yield* Ref.getAndSet(pendingProcessTables, new Map()); + const windowsListeners = yield* Ref.getAndSet(pendingWindowsListeners, new Map()); const histories = yield* Ref.getAndSet(pendingHistories, new Map()); yield* Effect.forEach(samples.values(), (deferred) => Deferred.fail(deferred, error), { discard: true, }); + yield* Effect.forEach(processTables.values(), (deferred) => Deferred.fail(deferred, error), { + discard: true, + }); + yield* Effect.forEach( + windowsListeners.values(), + (deferred) => Deferred.fail(deferred, error), + { discard: true }, + ); yield* Effect.forEach( histories.values(), (request) => Deferred.fail(request.deferred, error), @@ -485,6 +518,45 @@ export const make = Effect.fn("resourceTelemetry.nativeTelemetryClient.make")(fu } } }); + case "processTable": + return Ref.modify(pendingProcessTables, (pending) => { + const next = new Map(pending); + const deferred = next.get(event.requestId); + next.delete(event.requestId); + return [Option.fromUndefinedOr(deferred), next] as const; + }).pipe( + Effect.flatMap( + Option.match({ + onNone: () => Effect.void, + onSome: (deferred) => Deferred.succeed(deferred, event.processes), + }), + ), + Effect.asVoid, + ); + case "windowsListeners": + return Ref.modify(pendingWindowsListeners, (pending) => { + const next = new Map(pending); + const deferred = next.get(event.requestId); + next.delete(event.requestId); + return [Option.fromUndefinedOr(deferred), next] as const; + }).pipe( + Effect.flatMap( + Option.match({ + onNone: () => Effect.void, + onSome: (deferred) => + event.error === null + ? Deferred.succeed(deferred, event.listeners) + : Deferred.fail( + deferred, + new NativeTelemetryCommandFailed({ + operation: "windowsListeners", + cause: event.error, + }), + ), + }), + ), + Effect.asVoid, + ); case "historyChunk": return Effect.gen(function* () { const latestSnapshot = event.snapshots.at(-1); @@ -940,6 +1012,116 @@ export const make = Effect.fn("resourceTelemetry.nativeTelemetryClient.make")(fu ); }); + const processTable: NativeTelemetryClient["Service"]["processTable"] = Effect.gen(function* () { + const current = yield* Ref.get(state); + if (!canCommandNativeTelemetrySidecar(current.status, Option.isSome(current.handle))) { + return yield* new NativeTelemetryUnavailable({ + reason: Option.getOrElse(current.lastError, () => "sidecar is not running"), + }); + } + + const requestId = yield* crypto.randomUUIDv4.pipe( + Effect.mapError( + (cause) => new NativeTelemetryCommandFailed({ operation: "createRequestId", cause }), + ), + ); + const deferred = yield* Deferred.make< + ReadonlyArray, + NativeTelemetryClientError + >(); + yield* Ref.update(pendingProcessTables, (pending) => { + const next = new Map(pending); + next.set(requestId, deferred); + return next; + }); + return yield* writeCommand(Option.getOrThrow(current.handle), { + version: RESOURCE_MONITOR_PROTOCOL_VERSION, + type: "processTable", + requestId, + }).pipe( + Effect.andThen( + Deferred.await(deferred).pipe( + Effect.timeoutOption(PROCESS_TABLE_REQUEST_TIMEOUT), + Effect.flatMap( + Option.match({ + onNone: () => + Effect.fail( + new NativeTelemetryRequestTimedOut({ + operation: "processTable", + timeoutMs: Duration.toMillis(PROCESS_TABLE_REQUEST_TIMEOUT), + }), + ), + onSome: Effect.succeed, + }), + ), + ), + ), + Effect.ensuring( + Ref.update(pendingProcessTables, (pending) => { + const next = new Map(pending); + next.delete(requestId); + return next; + }), + ), + ); + }); + + const windowsListeners: NativeTelemetryClient["Service"]["windowsListeners"] = Effect.gen( + function* () { + const current = yield* Ref.get(state); + if (!canCommandNativeTelemetrySidecar(current.status, Option.isSome(current.handle))) { + return yield* new NativeTelemetryUnavailable({ + reason: Option.getOrElse(current.lastError, () => "sidecar is not running"), + }); + } + + const requestId = yield* crypto.randomUUIDv4.pipe( + Effect.mapError( + (cause) => new NativeTelemetryCommandFailed({ operation: "createRequestId", cause }), + ), + ); + const deferred = yield* Deferred.make< + ReadonlyArray, + NativeTelemetryClientError + >(); + yield* Ref.update(pendingWindowsListeners, (pending) => { + const next = new Map(pending); + next.set(requestId, deferred); + return next; + }); + return yield* writeCommand(Option.getOrThrow(current.handle), { + version: RESOURCE_MONITOR_PROTOCOL_VERSION, + type: "windowsListeners", + requestId, + }).pipe( + Effect.andThen( + Deferred.await(deferred).pipe( + Effect.timeoutOption(PROCESS_TABLE_REQUEST_TIMEOUT), + Effect.flatMap( + Option.match({ + onNone: () => + Effect.fail( + new NativeTelemetryRequestTimedOut({ + operation: "windowsListeners", + timeoutMs: Duration.toMillis(PROCESS_TABLE_REQUEST_TIMEOUT), + }), + ), + onSome: Effect.succeed, + }), + ), + ), + ), + Effect.ensuring( + Ref.update(pendingWindowsListeners, (pending) => { + const next = new Map(pending); + next.delete(requestId); + return next; + }), + ), + ); + }, + ); + const health = currentHealth; return NativeTelemetryClient.of({ @@ -961,6 +1143,8 @@ export const make = Effect.fn("resourceTelemetry.nativeTelemetryClient.make")(fu setExternalProcesses, setHostPowerState, sampleNow, + processTable, + windowsListeners, retry: Ref.get(state).pipe( Effect.flatMap((current) => !canRequestNativeTelemetryRetry(current.status, Option.isSome(current.handle)) @@ -1014,6 +1198,16 @@ export const layerTest = ( reason: "No resource monitor sample was configured for this test.", }), ), + processTable: Effect.fail( + new NativeTelemetryUnavailable({ + reason: "No resource monitor process table was configured for this test.", + }), + ), + windowsListeners: Effect.fail( + new NativeTelemetryUnavailable({ + reason: "No Windows listener table was configured for this test.", + }), + ), retry: Effect.succeed(false), health, subscribeHealth: diff --git a/apps/server/src/resourceTelemetry/ResourceTelemetry.test.ts b/apps/server/src/resourceTelemetry/ResourceTelemetry.test.ts index a96423607baf..62e22b31a785 100644 --- a/apps/server/src/resourceTelemetry/ResourceTelemetry.test.ts +++ b/apps/server/src/resourceTelemetry/ResourceTelemetry.test.ts @@ -82,7 +82,7 @@ function nativeSnapshot(input: { }), ]; return { - version: 2, + version: 4, type: "snapshot", sequence: input.sequence, sampledAtUnixMs: input.sampledAtUnixMs, @@ -497,7 +497,7 @@ describe("ResourceTelemetry", () => { const nativeHealth = yield* Ref.make({ status: "healthy", hello: Option.some({ - version: 2, + version: 4, type: "hello", sidecarVersion: "0.1.0", sidecarPid: 9_000, diff --git a/apps/server/src/resourceTelemetry/ResourceTelemetryHistory.test.ts b/apps/server/src/resourceTelemetry/ResourceTelemetryHistory.test.ts index 879c83d86dee..1f7181adac2f 100644 --- a/apps/server/src/resourceTelemetry/ResourceTelemetryHistory.test.ts +++ b/apps/server/src/resourceTelemetry/ResourceTelemetryHistory.test.ts @@ -64,7 +64,7 @@ function snapshot( }), ]; return { - version: 2, + version: 4, type: "snapshot", sequence, sampledAtUnixMs, diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index 3a93adc6d761..249901368efd 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -358,11 +358,15 @@ const CheckpointingLayerLive = Layer.empty.pipe( Layer.provideMerge(CheckpointStore.layer.pipe(Layer.provide(VcsDriverRegistryLayerLive))), ); -const PortScannerLayerLive = PortScanner.layer.pipe(Layer.provide(ProcessRunner.layer)); +const PortScannerLayerLive = PortScanner.layer.pipe( + Layer.provide(ProcessRunner.layer), + Layer.provide(NativeTelemetryLayerLive), +); const TerminalLayerLive = TerminalManager.layer.pipe( Layer.provide(PtyAdapterLive), Layer.provide(PortScannerLayerLive), + Layer.provide(NativeTelemetryLayerLive), ); const PreviewLayerLive = Layer.empty.pipe( diff --git a/apps/server/src/terminal/Manager.test.ts b/apps/server/src/terminal/Manager.test.ts index 47d91e4516ec..0e6a4995bd95 100644 --- a/apps/server/src/terminal/Manager.test.ts +++ b/apps/server/src/terminal/Manager.test.ts @@ -10,6 +10,7 @@ import { } from "@t3tools/contracts"; import { HostProcessPlatform } from "@t3tools/shared/hostProcess"; import * as Data from "effect/Data"; +import * as Clock from "effect/Clock"; import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as Encoding from "effect/Encoding"; @@ -210,6 +211,10 @@ interface CreateManagerOptions { readonly childCommand: string | null; readonly processIds: ReadonlyArray; }>; + processTable?: Effect.Effect< + ReadonlyArray<{ readonly pid: number; readonly ppid: number; readonly name: string }>, + never + >; subprocessPollIntervalMs?: number; processKillGraceMs?: number; maxRetainedInactiveSessions?: number; @@ -248,6 +253,7 @@ const createManager = ( ...(options.subprocessInspector !== undefined ? { subprocessInspector: options.subprocessInspector } : {}), + ...(options.processTable !== undefined ? { processTable: options.processTable } : {}), ...(options.subprocessPollIntervalMs !== undefined ? { subprocessPollIntervalMs: options.subprocessPollIntervalMs } : {}), @@ -1073,6 +1079,96 @@ it.layer( }), ); + it("calculates snapshot failure backoff and success reset delays", () => { + assert.equal(TerminalManager.subprocessSnapshotPollDelayMs(1_000, 0), 1_000); + assert.equal(TerminalManager.subprocessSnapshotPollDelayMs(1_000, 1), 2_000); + assert.equal(TerminalManager.subprocessSnapshotPollDelayMs(1_000, 2), 4_000); + assert.equal(TerminalManager.subprocessSnapshotPollDelayMs(1_000, 30), 60_000); + }); + + it.effect("uses process snapshots from the resource monitor", () => + Effect.gen(function* () { + let snapshotCalls = 0; + const { manager, getEvents } = yield* createManager(5, { + subprocessPollIntervalMs: 20, + processTable: Effect.sync(() => { + snapshotCalls += 1; + return [{ pid: 100, ppid: 9000, name: "ping.exe" }]; + }), + }).pipe(Effect.provide(withHostPlatform("win32"))); + + yield* manager.open(openInput()); + yield* waitFor( + Effect.map(getEvents, (events) => + events.some( + (event) => + event.type === "activity" && event.hasRunningSubprocess && event.label === "ping", + ), + ), + "1200 millis", + ); + expect(snapshotCalls).toBeGreaterThan(0); + }), + ); + + it.effect("backs off the spawned fallback when the resource monitor snapshot fails", () => + Effect.gen(function* () { + const fallbackCalls: Array = []; + const processRunner: ProcessRunner.ProcessRunner["Service"] = { + run: () => + Clock.currentTimeMillis.pipe( + Effect.map((now) => { + fallbackCalls.push(now); + return { + stdout: " 100 9000 vim", + stderr: "", + code: ChildProcessSpawner.ExitCode(0), + timedOut: false, + stdoutTruncated: false, + stderrInvalidUtf8: false, + stdoutInvalidUtf8: false, + stderrTruncated: false, + }; + }), + ), + }; + + const { manager, getEvents } = yield* createManager(5, { + subprocessPollIntervalMs: 20, + processTable: Effect.fail("sidecar unavailable").pipe( + Effect.mapError((cause) => cause as never), + ), + }).pipe( + Effect.provideService(ProcessRunner.ProcessRunner, processRunner), + Effect.provide(withHostPlatform("linux")), + ); + + yield* manager.open(openInput()); + // The fallback data is still applied while the sidecar is down. + yield* waitFor( + Effect.map(getEvents, (events) => + events.some( + (event) => + event.type === "activity" && + event.hasRunningSubprocess === true && + event.label === "vim", + ), + ), + "1200 millis", + ); + + yield* waitFor( + Effect.sync(() => fallbackCalls.length >= 4), + "2000 millis", + ); + // Four snapshots at the 20 ms base cadence would span ~60 ms. Backoff + // (40 + 80 + 160 ms) stretches the same four snapshots past 150 ms, so + // a stalled sidecar no longer hot-loops the spawned fallback. + const spanMs = fallbackCalls[3]! - fallbackCalls[0]!; + expect(spanMs).toBeGreaterThan(150); + }), + ); + it.effect("caps persisted history to configured line limit", () => Effect.gen(function* () { const { manager, ptyAdapter } = yield* createManager(3); diff --git a/apps/server/src/terminal/Manager.ts b/apps/server/src/terminal/Manager.ts index 64c2dbb913fb..a7c84e7793f1 100644 --- a/apps/server/src/terminal/Manager.ts +++ b/apps/server/src/terminal/Manager.ts @@ -26,6 +26,7 @@ import { type TerminalMetadataStreamEvent, type TerminalOpenInput, type TerminalResizeInput, + type ResourceMonitorProcessTableEntry, type TerminalRestartInput, type TerminalSessionSnapshot, type TerminalSessionStatus, @@ -59,6 +60,7 @@ import { } from "../observability/Metrics.ts"; import * as ProcessRunner from "../processRunner.ts"; import * as PortScanner from "../preview/PortScanner.ts"; +import * as NativeTelemetryClient from "../resourceTelemetry/NativeTelemetryClient.ts"; import * as PtyAdapter from "./PtyAdapter.ts"; export { @@ -77,6 +79,7 @@ export { const DEFAULT_HISTORY_LINE_LIMIT = 5_000; const DEFAULT_PERSIST_DEBOUNCE_MS = 40; const DEFAULT_SUBPROCESS_POLL_INTERVAL_MS = 1_000; +const MAX_SUBPROCESS_POLL_INTERVAL_MS = 60_000; const DEFAULT_PROCESS_KILL_GRACE_MS = 1_000; const DEFAULT_MAX_RETAINED_INACTIVE_SESSIONS = 128; const DEFAULT_OPEN_COLS = 120; @@ -89,7 +92,7 @@ class TerminalSubprocessCheckError extends Schema.TaggedErrorClass; } +export function subprocessSnapshotPollDelayMs( + pollIntervalMs: number, + failureCount: number, +): number { + return Math.min(pollIntervalMs * 2 ** failureCount, MAX_SUBPROCESS_POLL_INTERVAL_MS); +} + function parsePosixProcessTable(stdout: string): TerminalProcessTableSnapshot { const childrenByParent = new Map(); const commandById = new Map(); @@ -643,15 +653,15 @@ function parsePosixProcessTable(stdout: string): TerminalProcessTableSnapshot { return { childrenByParent, commandById }; } -function parseWindowsProcessTable(stdout: string): TerminalProcessTableSnapshot { +function processTableSnapshotFromProcesses( + processes: ReadonlyArray, +): TerminalProcessTableSnapshot { const childrenByParent = new Map(); const commandById = new Map(); - for (const line of stdout.split(/\r?\n/g)) { - const [pidRaw, parentPidRaw, nameRaw] = line.trim().split("|", 3); - const pid = Number(pidRaw); - const parentPid = Number(parentPidRaw); + for (const process of processes) { + const { pid, ppid: parentPid, name } = process; if (!Number.isInteger(pid) || !Number.isInteger(parentPid)) continue; - commandById.set(pid, nameRaw?.trim() ?? ""); + commandById.set(pid, name.trim()); const children = childrenByParent.get(parentPid) ?? []; children.push(pid); childrenByParent.set(parentPid, children); @@ -746,14 +756,11 @@ const windowsProcessTableSnapshot = Effect.fn("terminal.windowsProcessTableSnaps TerminalSubprocessCheckError, ProcessRunner.ProcessRunner > { + const processRunner = yield* ProcessRunner.ProcessRunner; const command = 'Get-CimInstance Win32_Process -ErrorAction Stop | ForEach-Object { Write-Output "$($_.ProcessId)|$($_.ParentProcessId)|$($_.Name)" }'; - const processRunner = yield* ProcessRunner.ProcessRunner; const result = yield* processRunner .run({ - // powershell.exe is a real executable — never spawn it through cmd.exe - // shell mode, which would re-tokenize the `-Command` payload (pipes, - // semicolons) before PowerShell ever sees it. command: "powershell.exe", args: ["-NoProfile", "-NonInteractive", "-Command", command], timeout: "1500 millis", @@ -763,16 +770,10 @@ const windowsProcessTableSnapshot = Effect.fn("terminal.windowsProcessTableSnaps }) .pipe( Effect.mapError( - (cause) => - new TerminalSubprocessCheckError({ - cause, - command: "powershell", - }), + (cause) => new TerminalSubprocessCheckError({ cause, command: "powershell" }), ), ); if (result.code !== 0 || result.timedOut || result.stdoutTruncated) { - // Not authoritative: an empty or partial table would mark every terminal - // idle and clear its registered process ids. Failing skips the tick. return yield* new TerminalSubprocessCheckError({ command: "powershell", exitCode: result.code, @@ -780,7 +781,15 @@ const windowsProcessTableSnapshot = Effect.fn("terminal.windowsProcessTableSnaps stdoutTruncated: result.stdoutTruncated, }); } - return parseWindowsProcessTable(result.stdout); + const processes = result.stdout.split(/\r?\n/g).flatMap((line) => { + const [pidRaw, ppidRaw, name = ""] = line.trim().split("|", 3); + const pid = Number(pidRaw); + const ppid = Number(ppidRaw); + return Number.isInteger(pid) && pid > 0 && Number.isInteger(ppid) + ? [{ pid, ppid, name }] + : []; + }); + return processTableSnapshotFromProcesses(processes); }, ); @@ -1115,6 +1124,10 @@ interface TerminalManagerOptions { shellResolver?: () => string; env?: NodeJS.ProcessEnv; subprocessInspector?: TerminalSubprocessInspector; + processTable?: Effect.Effect< + ReadonlyArray, + TerminalSubprocessCheckError + >; subprocessPollIntervalMs?: number; processKillGraceMs?: number; maxRetainedInactiveSessions?: number; @@ -1133,9 +1146,15 @@ export const make = Effect.fn("TerminalManager.make")(function* () { const { terminalLogsDir } = yield* ServerConfig.ServerConfig; const ptyAdapter = yield* PtyAdapter.PtyAdapter; const portDiscovery = yield* PortScanner.PortDiscovery; + const nativeTelemetry = yield* NativeTelemetryClient.NativeTelemetryClient; return yield* makeWithOptions({ logsDir: terminalLogsDir, ptyAdapter, + processTable: nativeTelemetry.processTable.pipe( + Effect.mapError( + (cause) => new TerminalSubprocessCheckError({ cause, command: "resource-monitor" }), + ), + ), registerTerminalProcesses: portDiscovery.registerTerminalProcesses, unregisterTerminal: portDiscovery.unregisterTerminal, }); @@ -1162,23 +1181,60 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func // One process-table snapshot per poll tick, shared across every terminal. // Per-terminal `pgrep`/`ps` calls multiply spawn load by terminal count and // can exhaust the PID space on hosts with many sessions (#6332). - const fetchProcessTableSnapshot = ( + const fallbackProcessTableSnapshot = ( platform === "win32" ? windowsProcessTableSnapshot() : posixProcessTableSnapshot(yield* resolvePosixPsCommand()) ).pipe(Effect.provideService(ProcessRunner.ProcessRunner, processRunner)); + const fetchProcessTableSnapshot: Effect.Effect< + { + readonly snapshot: TerminalProcessTableSnapshot; + /** + * False when the sidecar snapshot failed and this table came from the + * spawned fallback. The data is still applied, but the tick counts as + * a failure so polling backs off instead of hot-looping the fallback. + */ + readonly snapshotSucceeded: boolean; + }, + TerminalSubprocessCheckError + > = options.processTable + ? options.processTable.pipe( + Effect.map((entries) => ({ + snapshot: processTableSnapshotFromProcesses(entries), + snapshotSucceeded: true, + })), + Effect.catch(() => + fallbackProcessTableSnapshot.pipe( + Effect.map((snapshot) => ({ snapshot, snapshotSucceeded: false })), + ), + ), + ) + : fallbackProcessTableSnapshot.pipe( + Effect.map((snapshot) => ({ snapshot, snapshotSucceeded: true })), + ); const customSubprocessInspector = options.subprocessInspector; const acquireSubprocessInspector: Effect.Effect< - TerminalSubprocessInspector, + { + readonly inspector: TerminalSubprocessInspector; + readonly snapshotSucceeded: boolean; + }, TerminalSubprocessCheckError > = customSubprocessInspector !== undefined - ? Effect.succeed(customSubprocessInspector) + ? Effect.succeed({ inspector: customSubprocessInspector, snapshotSucceeded: true }) : Effect.map( fetchProcessTableSnapshot, - (snapshot): TerminalSubprocessInspector => - (terminalPid) => + ({ + snapshot, + snapshotSucceeded, + }): { + readonly inspector: TerminalSubprocessInspector; + readonly snapshotSucceeded: boolean; + } => ({ + inspector: (terminalPid) => Effect.succeed(deriveSubprocessInspectResult(snapshot, terminalPid, platform)), + snapshotSucceeded, + }), ); const subprocessPollIntervalMs = options.subprocessPollIntervalMs ?? DEFAULT_SUBPROCESS_POLL_INTERVAL_MS; @@ -2008,7 +2064,7 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func ); if (runningSessions.length === 0) { - return; + return true; } const inspectorOption = yield* acquireSubprocessInspector.pipe( @@ -2016,15 +2072,22 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func Effect.catch((reason) => Effect.logWarning("failed to snapshot processes for terminal subprocess polling", { reason, - }).pipe(Effect.as(Option.none())), + }).pipe( + Effect.as( + Option.none<{ + readonly inspector: TerminalSubprocessInspector; + readonly snapshotSucceeded: boolean; + }>(), + ), + ), ), ); if (Option.isNone(inspectorOption)) { - return; + return false; } - const subprocessInspector = inspectorOption.value; + const { inspector: subprocessInspector, snapshotSucceeded } = inspectorOption.value; const checkSubprocessActivity = Effect.fn("terminal.checkSubprocessActivity")(function* ( session: TerminalSessionState & { pid: number }, @@ -2093,6 +2156,7 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func concurrency: "unbounded", discard: true, }); + return snapshotSucceeded; }); const hasRunningSessions = readManagerState.pipe( @@ -2101,14 +2165,26 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func ), ); + let subprocessSnapshotFailureCount = 0; yield* Effect.forever( hasRunningSessions.pipe( Effect.flatMap((active) => active ? pollSubprocessActivity().pipe( - Effect.flatMap(() => Effect.sleep(subprocessPollIntervalMs)), + Effect.flatMap((snapshotSucceeded) => { + subprocessSnapshotFailureCount = snapshotSucceeded + ? 0 + : Math.min(subprocessSnapshotFailureCount + 1, 30); + const delayMs = subprocessSnapshotPollDelayMs( + subprocessPollIntervalMs, + subprocessSnapshotFailureCount, + ); + return Effect.sleep(delayMs); + }), ) - : Effect.sleep(subprocessPollIntervalMs), + : Effect.sync(() => { + subprocessSnapshotFailureCount = 0; + }).pipe(Effect.flatMap(() => Effect.sleep(subprocessPollIntervalMs))), ), ), ).pipe(Effect.forkIn(workerScope)); diff --git a/docs/internals/resource-telemetry.md b/docs/internals/resource-telemetry.md index 0f3e6674ff79..049c6c661cc4 100644 --- a/docs/internals/resource-telemetry.md +++ b/docs/internals/resource-telemetry.md @@ -80,10 +80,14 @@ line on stdout: - `setSampleInterval` - `setStreaming` - `sampleNow` +- `processTable` +- `windowsListeners` - `readHistory` - `shutdown` - `hello` - `snapshot` +- `processTable` +- `windowsListeners` - `historyChunk` - `error` diff --git a/native/resource-monitor/src/main.rs b/native/resource-monitor/src/main.rs index 0596aea977cd..cf465c2d01ea 100644 --- a/native/resource-monitor/src/main.rs +++ b/native/resource-monitor/src/main.rs @@ -8,7 +8,7 @@ use sysinfo::{ MINIMUM_CPU_UPDATE_INTERVAL, Pid, ProcessRefreshKind, ProcessesToUpdate, System, UpdateKind, }; -const PROTOCOL_VERSION: u32 = 2; +const PROTOCOL_VERSION: u32 = 4; const MIN_SAMPLE_INTERVAL_MS: u64 = 250; const MAX_SAMPLE_INTERVAL_MS: u64 = 60_000; const PROCESS_START_TIME_PRECISION_MS: u64 = 1_000; @@ -66,6 +66,14 @@ enum Command { version: u32, request_id: String, }, + ProcessTable { + version: u32, + request_id: String, + }, + WindowsListeners { + version: u32, + request_id: String, + }, ReadHistory { version: u32, request_id: String, @@ -84,6 +92,8 @@ impl Command { | Self::SetSampleInterval { version, .. } | Self::SetStreaming { version, .. } | Self::SampleNow { version, .. } + | Self::ProcessTable { version, .. } + | Self::WindowsListeners { version, .. } | Self::ReadHistory { version, .. } | Self::Shutdown { version } => *version, } @@ -146,6 +156,173 @@ struct ProcessSample { io_semantics: IoSemantics, } +#[derive(Debug, Serialize)] +#[serde(rename_all = "camelCase")] +struct ProcessTableEntry { + pid: u32, + ppid: u32, + name: String, +} + +#[derive(Debug, Serialize)] +#[serde(rename_all = "camelCase")] +struct ProcessTableEvent<'a> { + version: u32, + #[serde(rename = "type")] + event_type: &'static str, + request_id: &'a str, + processes: Vec, +} + +#[derive(Debug, Clone, Serialize)] +#[serde(rename_all = "camelCase")] +struct WindowsListener { + port: u16, + pid: u32, + process_name: Option, +} + +#[derive(Debug, Serialize)] +#[serde(rename_all = "camelCase")] +struct WindowsListenersEvent<'a> { + version: u32, + #[serde(rename = "type")] + event_type: &'static str, + request_id: &'a str, + listeners: Vec, + error: Option, +} + +#[cfg(target_os = "windows")] +mod windows_listeners { + use std::ffi::c_void; + use std::io; + use std::mem::size_of; + use std::ptr; + + const AF_INET: u32 = 2; + const AF_INET6: u32 = 23; + const ERROR_INSUFFICIENT_BUFFER: u32 = 122; + const TCP_TABLE_OWNER_PID_LISTENER: u32 = 3; + + #[repr(C)] + #[derive(Clone, Copy)] + struct TcpRow { + state: u32, + local_address: u32, + local_port: u32, + remote_address: u32, + remote_port: u32, + owning_pid: u32, + } + + #[repr(C)] + #[derive(Clone, Copy)] + struct Tcp6Row { + local_address: [u8; 16], + local_scope_id: u32, + local_port: u32, + remote_address: [u8; 16], + remote_scope_id: u32, + remote_port: u32, + state: u32, + owning_pid: u32, + } + + #[link(name = "iphlpapi")] + unsafe extern "system" { + fn GetExtendedTcpTable( + table: *mut c_void, + size: *mut u32, + order: i32, + family: u32, + table_class: u32, + reserved: u32, + ) -> u32; + } + + unsafe fn rows(family: u32) -> io::Result> { + let mut size = 0u32; + let first = unsafe { + GetExtendedTcpTable( + ptr::null_mut(), + &mut size, + 0, + family, + TCP_TABLE_OWNER_PID_LISTENER, + 0, + ) + }; + if first != ERROR_INSUFFICIENT_BUFFER { + return Err(io::Error::from_raw_os_error(first as i32)); + } + loop { + let word_count = (size as usize).div_ceil(size_of::()); + let mut buffer = vec![0u32; word_count]; + let result = unsafe { + GetExtendedTcpTable( + buffer.as_mut_ptr().cast(), + &mut size, + 0, + family, + TCP_TABLE_OWNER_PID_LISTENER, + 0, + ) + }; + if result == ERROR_INSUFFICIENT_BUFFER { + continue; + } + if result != 0 { + return Err(io::Error::from_raw_os_error(result as i32)); + } + let count = unsafe { ptr::read_unaligned(buffer.as_ptr()) } as usize; + let available_bytes = (buffer.len() * size_of::()).saturating_sub(size_of::()); + if count > available_bytes / size_of::() { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "Windows TCP table exceeded its returned buffer", + )); + } + let first_row = unsafe { buffer.as_ptr().cast::().add(size_of::()) }; + return Ok((0..count) + .map(|index| unsafe { + ptr::read_unaligned(first_row.add(index * size_of::()).cast::()) + }) + .collect()); + } + } + + fn port(value: u32) -> u16 { + u16::from_be(value as u16) + } + + pub fn read() -> io::Result> { + let ipv4 = unsafe { rows::(AF_INET)? }; + let ipv6 = unsafe { rows::(AF_INET6)? }; + let mut listeners = ipv4 + .into_iter() + .filter(|row| { + row.local_address == 0 || row.local_address.to_ne_bytes().first() == Some(&127) + }) + .map(|row| (port(row.local_port), row.owning_pid)) + .chain( + ipv6 + .into_iter() + .filter(|row| { + row.local_address.iter().all(|byte| *byte == 0) + || row.local_address + == [0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 1] + }) + .map(|row| (port(row.local_port), row.owning_pid)), + ) + .filter(|(port, pid)| *port > 0 && *pid > 0) + .collect::>(); + listeners.sort_unstable(); + listeners.dedup(); + Ok(listeners) + } +} + impl ProcessSample { fn estimated_history_bytes(&self) -> usize { std::mem::size_of::() @@ -351,6 +528,65 @@ impl Collector { self.cpu_baseline_refreshed_at = Some(Instant::now()); } + fn process_table(&mut self) -> Vec { + self.system.refresh_processes_specifics( + ProcessesToUpdate::All, + true, + ProcessRefreshKind::nothing().without_tasks(), + ); + let mut processes = self + .system + .processes() + .iter() + .filter_map(|(pid, process)| { + let pid = pid.as_u32(); + // Pid 0 is the kernel idle process on some platforms. The + // processTable contract requires positive pids, and one zero + // would fail the whole event decode on the server, so drop it + // here. It can never be a terminal descendant. + if pid == 0 { + return None; + } + Some(ProcessTableEntry { + pid, + ppid: process.parent().map(Pid::as_u32).unwrap_or(0), + name: truncate_utf8( + process.name().to_string_lossy().into_owned(), + MAX_PROCESS_NAME_BYTES, + ), + }) + }) + .collect::>(); + processes.sort_by_key(|process| process.pid); + processes + } + + #[cfg(target_os = "windows")] + fn windows_listeners(&mut self) -> Result, String> { + let process_names = self + .process_table() + .into_iter() + .map(|process| (process.pid, process.name)) + .collect::>(); + windows_listeners::read() + .map(|listeners| { + listeners + .into_iter() + .map(|(port, pid)| WindowsListener { + port, + pid, + process_name: process_names.get(&pid).cloned(), + }) + .collect() + }) + .map_err(|error| error.to_string()) + } + + #[cfg(not(target_os = "windows"))] + fn windows_listeners(&mut self) -> Result, String> { + Err("Windows listener discovery is unavailable on this platform".to_owned()) + } + fn sample(&mut self, config: &CollectorConfig, request_id: Option) -> SnapshotEvent { if let Some(delay) = remaining_cpu_measurement_delay(self.cpu_baseline_refreshed_at.take(), Instant::now()) @@ -815,6 +1051,26 @@ fn main() -> io::Result<()> { )?; } } + Command::ProcessTable { request_id, .. } => { + let event = ProcessTableEvent { + version: PROTOCOL_VERSION, + event_type: "processTable", + request_id: &request_id, + processes: collector.process_table(), + }; + write_event(&mut writer, &event)?; + } + Command::WindowsListeners { request_id, .. } => { + let result = collector.windows_listeners(); + let event = WindowsListenersEvent { + version: PROTOCOL_VERSION, + event_type: "windowsListeners", + request_id: &request_id, + listeners: result.as_ref().cloned().unwrap_or_default(), + error: result.err(), + }; + write_event(&mut writer, &event)?; + } Command::ReadHistory { request_id, window_ms, @@ -845,6 +1101,13 @@ fn main() -> io::Result<()> { mod tests { use super::*; + #[cfg(target_os = "windows")] + #[test] + fn reads_windows_listener_table() { + let listeners = windows_listeners::read().expect("Windows listener table"); + assert!(listeners.iter().all(|(port, pid)| *port > 0 && *pid > 0)); + } + #[test] fn selects_roots_and_all_descendants() { let rows = vec![ @@ -892,7 +1155,7 @@ mod tests { #[test] fn decodes_protocol_commands() { let configure = serde_json::from_str::( - r#"{"version":2,"type":"configure","rootPid":42,"sampleIntervalMs":1000,"externalProcesses":[{"pid":7}]}"#, + r#"{"version":4,"type":"configure","rootPid":42,"sampleIntervalMs":1000,"externalProcesses":[{"pid":7}]}"#, ) .expect("configure command"); @@ -912,7 +1175,7 @@ mod tests { } let read_history = serde_json::from_str::( - r#"{"version":2,"type":"readHistory","requestId":"history-1","windowMs":60000}"#, + r#"{"version":4,"type":"readHistory","requestId":"history-1","windowMs":60000}"#, ) .expect("read history command"); assert!(matches!( @@ -923,6 +1186,24 @@ mod tests { .. } if request_id == "history-1" )); + + let process_table = serde_json::from_str::( + r#"{"version":4,"type":"processTable","requestId":"processes-1"}"#, + ) + .expect("process table command"); + assert!(matches!( + process_table, + Command::ProcessTable { request_id, .. } if request_id == "processes-1" + )); + + let windows_listeners = serde_json::from_str::( + r#"{"version":4,"type":"windowsListeners","requestId":"listeners-1"}"#, + ) + .expect("Windows listeners command"); + assert!(matches!( + windows_listeners, + Command::WindowsListeners { request_id, .. } if request_id == "listeners-1" + )); } #[test] diff --git a/packages/contracts/src/resourceTelemetry.ts b/packages/contracts/src/resourceTelemetry.ts index 3ec1e4de3ef4..2ec1df228b9e 100644 --- a/packages/contracts/src/resourceTelemetry.ts +++ b/packages/contracts/src/resourceTelemetry.ts @@ -4,7 +4,7 @@ import { NonNegativeInt, PositiveInt, TrimmedNonEmptyString } from "./baseSchema import { HostPowerSnapshot } from "./background.ts"; import { DesktopUpdateStateSchema } from "./ipc.ts"; -export const RESOURCE_MONITOR_PROTOCOL_VERSION = 2 as const; +export const RESOURCE_MONITOR_PROTOCOL_VERSION = 4 as const; export const ResourceTelemetryIoSemantics = Schema.Literals([ "storage", @@ -102,6 +102,21 @@ export const ResourceMonitorSampleNowCommand = Schema.Struct({ }); export type ResourceMonitorSampleNowCommand = typeof ResourceMonitorSampleNowCommand.Type; +export const ResourceMonitorProcessTableCommand = Schema.Struct({ + version: Schema.Literal(RESOURCE_MONITOR_PROTOCOL_VERSION), + type: Schema.Literal("processTable"), + requestId: TrimmedNonEmptyString, +}); +export type ResourceMonitorProcessTableCommand = typeof ResourceMonitorProcessTableCommand.Type; + +export const ResourceMonitorWindowsListenersCommand = Schema.Struct({ + version: Schema.Literal(RESOURCE_MONITOR_PROTOCOL_VERSION), + type: Schema.Literal("windowsListeners"), + requestId: TrimmedNonEmptyString, +}); +export type ResourceMonitorWindowsListenersCommand = + typeof ResourceMonitorWindowsListenersCommand.Type; + export const ResourceMonitorSetSampleIntervalCommand = Schema.Struct({ version: Schema.Literal(RESOURCE_MONITOR_PROTOCOL_VERSION), type: Schema.Literal("setSampleInterval"), @@ -137,6 +152,8 @@ export const ResourceMonitorCommand = Schema.Union([ ResourceMonitorSetSampleIntervalCommand, ResourceMonitorSetStreamingCommand, ResourceMonitorSampleNowCommand, + ResourceMonitorProcessTableCommand, + ResourceMonitorWindowsListenersCommand, ResourceMonitorReadHistoryCommand, ResourceMonitorShutdownCommand, ]); @@ -168,6 +185,37 @@ export const ResourceMonitorSnapshotEvent = Schema.Struct({ }); export type ResourceMonitorSnapshotEvent = typeof ResourceMonitorSnapshotEvent.Type; +export const ResourceMonitorProcessTableEntry = Schema.Struct({ + pid: PositiveInt, + ppid: NonNegativeInt, + name: Schema.String, +}); +export type ResourceMonitorProcessTableEntry = typeof ResourceMonitorProcessTableEntry.Type; + +export const ResourceMonitorProcessTableEvent = Schema.Struct({ + version: Schema.Literal(RESOURCE_MONITOR_PROTOCOL_VERSION), + type: Schema.Literal("processTable"), + requestId: TrimmedNonEmptyString, + processes: Schema.Array(ResourceMonitorProcessTableEntry), +}); +export type ResourceMonitorProcessTableEvent = typeof ResourceMonitorProcessTableEvent.Type; + +export const ResourceMonitorWindowsListener = Schema.Struct({ + port: PositiveInt, + pid: PositiveInt, + processName: Schema.NullOr(Schema.String), +}); +export type ResourceMonitorWindowsListener = typeof ResourceMonitorWindowsListener.Type; + +export const ResourceMonitorWindowsListenersEvent = Schema.Struct({ + version: Schema.Literal(RESOURCE_MONITOR_PROTOCOL_VERSION), + type: Schema.Literal("windowsListeners"), + requestId: TrimmedNonEmptyString, + listeners: Schema.Array(ResourceMonitorWindowsListener), + error: Schema.NullOr(Schema.String), +}); +export type ResourceMonitorWindowsListenersEvent = typeof ResourceMonitorWindowsListenersEvent.Type; + export const ResourceMonitorHistoryChunkEvent = Schema.Struct({ version: Schema.Literal(RESOURCE_MONITOR_PROTOCOL_VERSION), type: Schema.Literal("historyChunk"), @@ -189,6 +237,8 @@ export type ResourceMonitorErrorEvent = typeof ResourceMonitorErrorEvent.Type; export const ResourceMonitorEvent = Schema.Union([ ResourceMonitorHelloEvent, ResourceMonitorSnapshotEvent, + ResourceMonitorProcessTableEvent, + ResourceMonitorWindowsListenersEvent, ResourceMonitorHistoryChunkEvent, ResourceMonitorErrorEvent, ]);