diff --git a/apps/mobile/src/components/AppSymbol.tsx b/apps/mobile/src/components/AppSymbol.tsx index 9e9a6e1173c7..e1037c36e24f 100644 --- a/apps/mobile/src/components/AppSymbol.tsx +++ b/apps/mobile/src/components/AppSymbol.tsx @@ -144,6 +144,7 @@ const ANDROID_ICON_BY_SF_SYMBOL = { "bolt.horizontal.circle": IconBolt, brain: IconBrain, camera: IconCamera, + "chart.bar": IconChartBar, "chart.bar.xaxis": IconChartBar, checkmark: IconCheck, "checkmark.circle": IconCircleCheck, diff --git a/apps/mobile/src/features/settings/SettingsServerControlsRouteScreen.tsx b/apps/mobile/src/features/settings/SettingsServerControlsRouteScreen.tsx index ae4f43b5fb37..e07cafe89eb6 100644 --- a/apps/mobile/src/features/settings/SettingsServerControlsRouteScreen.tsx +++ b/apps/mobile/src/features/settings/SettingsServerControlsRouteScreen.tsx @@ -230,6 +230,27 @@ function ServerSettingsDetail(props: { readonly page: SettingsPage }) { (target) => target.environment.serverConfig.environment.capabilities.threadRestartContinuation === true, ); + + const telemetryOverridden = targets.some( + (target) => target.environment.serverConfig.telemetryDisabledByEnvironment === true, + ); + + const telemetryPartiallyOverridden = + telemetryOverridden && + targets.some( + (target) => target.environment.serverConfig.telemetryDisabledByEnvironment !== true, + ); + + const telemetryValue = telemetryPartiallyOverridden + ? null + : !telemetryOverridden && uniform("telemetryEnabled"); + + const telemetryDescription = telemetryPartiallyOverridden + ? "Disabled by some selected servers' environment configuration." + : telemetryOverridden + ? "Disabled by the server's environment configuration." + : "Share anonymous usage data to help improve T3 Code."; + const disabledFor = (key: string) => disabled || (projectSelected && @@ -449,6 +470,23 @@ function ServerSettingsDetail(props: { readonly page: SettingsPage }) { ))} ) : null} + + { + write({ telemetryEnabled: value }); + }} + /> + , "children"> & { readonly value: boolean | null; + readonly mixedValue?: boolean; readonly onValueChange: (value: boolean) => void; }, ) { + const mixedValue = props.mixedValue ?? true; + const mixedLabel = mixedValue ? "on" : "off"; return ( {props.value === null ? ( props.onValueChange(true)} + onPress={() => props.onValueChange(mixedValue)} > - Mixed · Set on + Mixed · Set {mixedLabel} ) : ( Effect.die("Unexpected provider mutation"), withSettingsSnapshot: (use) => Ref.get(settings).pipe(Effect.flatMap(use)), streamChanges: Stream.fromPubSub(settingsChanges), + subscribePersistedChanges: PubSub.subscribe(settingsChanges).pipe( + Effect.map(Stream.fromSubscription), + ), subscribeChanges: PubSub.subscribe(settingsChanges).pipe( Effect.map((subscription) => Stream.fromSubscription(subscription)), ), diff --git a/apps/server/src/provider/ProviderRegistry.test.ts b/apps/server/src/provider/ProviderRegistry.test.ts index 155eed099f25..80dd36af9cdb 100644 --- a/apps/server/src/provider/ProviderRegistry.test.ts +++ b/apps/server/src/provider/ProviderRegistry.test.ts @@ -359,6 +359,9 @@ function makeMutableServerSettingsService( get streamChanges() { return Stream.fromPubSub(changes); }, + get subscribePersistedChanges() { + return PubSub.subscribe(changes).pipe(Effect.map(Stream.fromSubscription)); + }, get subscribeChanges() { return PubSub.subscribe(changes).pipe( Effect.map((subscription) => Stream.fromSubscription(subscription)), diff --git a/apps/server/src/provider/makeManagedServerProvider.test.ts b/apps/server/src/provider/makeManagedServerProvider.test.ts index 17c1e442adb7..8af54cbaeb2d 100644 --- a/apps/server/src/provider/makeManagedServerProvider.test.ts +++ b/apps/server/src/provider/makeManagedServerProvider.test.ts @@ -303,6 +303,9 @@ describe("makeManagedServerProvider", () => { updateProviderInstance: () => Effect.die(new Error("unused in this test")), withSettingsSnapshot: (use) => Ref.get(serverSettingsRef).pipe(Effect.flatMap(use)), streamChanges: Stream.empty, + subscribePersistedChanges: PubSub.subscribe(serverSettingsChanges).pipe( + Effect.map(Stream.fromSubscription), + ), subscribeChanges: PubSub.subscribe(serverSettingsChanges).pipe( Effect.map((subscription) => Stream.fromSubscription(subscription)), ), diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index 192b17204aea..4e8efd567a83 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -636,7 +636,7 @@ const layerRuntimeDependencies = layerRuntimeCoreDependencies.pipe( Layer.provideMerge(layerResourceDiagnostics), Layer.provideMerge(layerUsage), Layer.provideMerge(TraceDiagnostics.layer), - Layer.provideMerge(AnalyticsService.layer), + Layer.provideMerge(AnalyticsService.layer.pipe(Layer.provide(layerServerSettings))), Layer.provideMerge(ExternalLauncher.layer), Layer.provideMerge(RemoteOpenTargets.layer), Layer.provideMerge(DirectEndpoints.layer), diff --git a/apps/server/src/serverSettings.test.ts b/apps/server/src/serverSettings.test.ts index 051f05331613..4f5c7e849504 100644 --- a/apps/server/src/serverSettings.test.ts +++ b/apps/server/src/serverSettings.test.ts @@ -97,6 +97,26 @@ const recordProviderUsage = (provider: string, instanceId: string | null = provi }); it.layer(NodeServices.layer)("server settings", (it) => { + it.effect("persists the analytics opt-out and allows opting back in", () => + Effect.gen(function* () { + const service = yield* ServerSettingsModule.ServerSettingsService; + const config = yield* ServerConfig.ServerConfig; + const fs = yield* FileSystem.FileSystem; + assert.isTrue((yield* service.getSettings).telemetryEnabled); + + for (const telemetryEnabled of [false, true]) { + yield* service.updateSettings({ telemetryEnabled }); + + const persisted = yield* decodeServerSettingsJson( + yield* fs.readFileString(config.settingsPath), + ); + + assert.equal(persisted.telemetryEnabled, telemetryEnabled); + assert.equal((yield* service.getSettings).telemetryEnabled, telemetryEnabled); + } + }).pipe(Effect.provide(layerServerSettings())), + ); + it.effect("migrates saved token delivery to paragraph buffering without resetting settings", () => Effect.gen(function* () { const config = yield* ServerConfig.ServerConfig; diff --git a/apps/server/src/serverSettings.ts b/apps/server/src/serverSettings.ts index 7cbcd895ac14..1e9204fb44db 100644 --- a/apps/server/src/serverSettings.ts +++ b/apps/server/src/serverSettings.ts @@ -288,6 +288,13 @@ export class ServerSettingsService extends Context.Service< * snapshot and a lazily started stream must not be lost. */ readonly subscribeChanges: Effect.Effect, never, Scope.Scope>; + + /** Subscribe without loading secrets, for settings that must react before secret-store I/O. */ + readonly subscribePersistedChanges: Effect.Effect< + Stream.Stream, + never, + Scope.Scope + >; } >()("t3/serverSettings/ServerSettingsService") { /** @deprecated Import and use `layerTest` from this module. */ @@ -309,6 +316,7 @@ const makeTest = (overrides: DeepPartial = {}) => : {}), }); const currentSettingsRef = yield* Ref.make(initialSettings); + const changesPubSub = yield* PubSub.unbounded(); const writeSemaphore = yield* Semaphore.make(1); const getSettings = Ref.get(currentSettingsRef).pipe(Effect.map(resolveTextGenerationProvider)); @@ -320,6 +328,7 @@ const makeTest = (overrides: DeepPartial = {}) => Effect.flatMap(update), Effect.flatMap(normalizeServerSettings), Effect.tap((nextSettings) => Ref.set(currentSettingsRef, nextSettings)), + Effect.tap((nextSettings) => PubSub.publish(changesPubSub, nextSettings)), Effect.map(resolveTextGenerationProvider), ), ); @@ -346,8 +355,11 @@ const makeTest = (overrides: DeepPartial = {}) => ), withSettingsSnapshot: (use) => writeSemaphore.withPermits(1)(getSettings.pipe(Effect.flatMap(use))), - streamChanges: Stream.empty, - subscribeChanges: Effect.succeed(Stream.empty), + streamChanges: Stream.fromPubSub(changesPubSub), + subscribeChanges: PubSub.subscribe(changesPubSub).pipe(Effect.map(Stream.fromSubscription)), + subscribePersistedChanges: PubSub.subscribe(changesPubSub).pipe( + Effect.map(Stream.fromSubscription), + ), } satisfies ServerSettingsService["Service"]; }); @@ -1390,6 +1402,9 @@ const make = Effect.gen(function* () { Effect.map((subscription) => materializeChanges(Stream.fromSubscription(subscription))), ); }, + get subscribePersistedChanges() { + return PubSub.subscribe(changesPubSub).pipe(Effect.map(Stream.fromSubscription)); + }, } satisfies ServerSettingsService["Service"]; }); diff --git a/apps/server/src/telemetry/AnalyticsService.test.ts b/apps/server/src/telemetry/AnalyticsService.test.ts index e0a4786ece6f..06163476564d 100644 --- a/apps/server/src/telemetry/AnalyticsService.test.ts +++ b/apps/server/src/telemetry/AnalyticsService.test.ts @@ -2,18 +2,23 @@ import * as NodeHttpServer from "@effect/platform-node/NodeHttpServer"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { assert, it } from "@effect/vitest"; import * as ConfigProvider from "effect/ConfigProvider"; +import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; +import * as Fiber from "effect/Fiber"; import * as Layer from "effect/Layer"; +import * as Predicate from "effect/Predicate"; import * as Schema from "effect/Schema"; import * as TestClock from "effect/testing/TestClock"; import * as HttpClient from "effect/http/HttpClient"; import * as HttpClientError from "effect/http/HttpClientError"; +import * as HttpClientResponse from "effect/http/HttpClientResponse"; import * as HttpServer from "effect/http/HttpServer"; import * as HttpServerRequest from "effect/http/HttpServerRequest"; import * as HttpServerResponse from "effect/http/HttpServerResponse"; import { HostProcessArchitecture, HostProcessPlatform } from "@t3tools/shared/hostProcess"; import * as ServerConfig from "../config.ts"; +import * as ServerSettings from "../serverSettings.ts"; import { getTelemetryIdentifier } from "./Identify.ts"; import * as AnalyticsService from "./AnalyticsService.ts"; @@ -52,10 +57,12 @@ interface RecordedBatchBody { const SentBatch = Schema.fromJsonString( Schema.Struct({ - batch: Schema.Array(Schema.Struct({ uuid: Schema.String })), + batch: Schema.Array(Schema.Struct({ uuid: Schema.String, event: Schema.String })), }), ); +const decodeSentBatch = Schema.decodeEffect(SentBatch); + /** * HTTP client that reads each batch, then fails as if the connection dropped * before the response arrived. PostHog stores these batches, so the server @@ -67,9 +74,9 @@ const layerAcceptThenFailClient = (batches: Array Effect.gen(function* () { if (request.body._tag === "Uint8Array") { - const body = yield* Schema.decodeEffect(SentBatch)( - new TextDecoder().decode(request.body.body), - ).pipe(Effect.orDie); + const body = yield* decodeSentBatch(new TextDecoder().decode(request.body.body)).pipe( + Effect.orDie, + ); batches.push(body.batch); } return yield* new HttpClientError.HttpClientError({ @@ -79,6 +86,33 @@ const layerAcceptThenFailClient = (batches: Array, + telemetryEnabled = true, + settingsLayer = ServerSettings.layerTest({ telemetryEnabled }), +) => + AnalyticsService.layer.pipe( + Layer.provideMerge(settingsLayer), + Layer.provide( + ServerConfig.ServerConfig.layerTest(process.cwd(), { prefix: "t3-telemetry-toggle-" }), + ), + Layer.provide( + ConfigProvider.layer( + ConfigProvider.fromUnknown({ + T3CODE_TELEMETRY_ENABLED: true, + T3CODE_TELEMETRY_FLUSH_BATCH_SIZE: 20, + }), + ), + ), + Layer.provide( + Layer.mergeAll( + Layer.succeed(HostProcessPlatform, "linux"), + Layer.succeed(HostProcessArchitecture, "arm64"), + clientLayer, + ), + ), + ); + it("retryDelayMs doubles from 2s and stays under the 5 minute cap", () => { assert.equal(AnalyticsService.retryDelayMs(1, 0), 1_000); assert.equal(AnalyticsService.retryDelayMs(2, 0.999_999), 4_000); @@ -86,7 +120,9 @@ it("retryDelayMs doubles from 2s and stays under the 5 minute cap", () => { assert.equal(AnalyticsService.retryDelayMs(30, 0.999_999), 300_000); }); -it.layer(NodeServices.layer)("AnalyticsService test", (it) => { +const layerTest = Layer.mergeAll(NodeServices.layer, ServerSettings.layerTest()); + +it.layer(layerTest)("AnalyticsService test", (it) => { it.effect("a batch that keeps failing is retried with backoff, then dropped", () => Effect.gen(function* () { const batches: Array> = []; @@ -267,11 +303,156 @@ it.layer(NodeServices.layer)("AnalyticsService test", (it) => { yield* Effect.gen(function* () { yield* Layer.launch(layerBatchServer).pipe(Effect.forkScoped); const analytics = yield* AnalyticsService.AnalyticsService; + const settings = yield* ServerSettings.ServerSettingsService; yield* analytics.record("test.disabled", { index: 1 }); yield* analytics.flush; + yield* settings.updateSettings({ telemetryEnabled: false }); + yield* settings.updateSettings({ telemetryEnabled: true }); + yield* analytics.record("test.environment-still-disabled"); + yield* analytics.flush; }).pipe(Effect.provide(layerRuntime)); assert.deepEqual(capturedPaths, []); }), ); + + it.effect("records and flushes without rereading settings after startup", () => + Effect.gen(function* () { + const settings = yield* ServerSettings.ServerSettingsService; + let settingsReads = 0; + let batchesSent = 0; + const settingsLayer = Layer.succeed(ServerSettings.ServerSettingsService, { + ...settings, + getSettings: Effect.sync(() => settingsReads++).pipe(Effect.andThen(settings.getSettings)), + }); + const clientLayer = Layer.succeed( + HttpClient.HttpClient, + HttpClient.make((request) => + Effect.sync(() => { + batchesSent++; + return HttpClientResponse.fromWeb(request, new Response("{}")); + }), + ), + ); + yield* Effect.gen(function* () { + const analytics = yield* AnalyticsService.AnalyticsService; + assert.equal(settingsReads, 1); + for (let index = 0; index < 25; index++) { + yield* analytics.record("test.cached-gate", { index }); + } + yield* analytics.flush; + assert.equal(batchesSent, 2); + yield* settings.updateSettings({ telemetryEnabled: false }); + yield* analytics.record("test.disabled"); + yield* analytics.flush; + assert.equal(batchesSent, 2); + yield* settings.updateSettings({ telemetryEnabled: true }); + yield* analytics.record("test.reenabled"); + yield* analytics.flush; + assert.equal(batchesSent, 3); + assert.equal(settingsReads, 1); + }).pipe(Effect.provide(layerToggleTest(clientLayer, true, settingsLayer))); + }), + ); + + it.effect("rapid opt-out and opt-in discard queued events and retries", () => + Effect.gen(function* () { + const batches: Array = []; + const layerRuntime = layerToggleTest(layerAcceptThenFailClient(batches), false); + yield* Effect.gen(function* () { + const analytics = yield* AnalyticsService.AnalyticsService; + const settings = yield* ServerSettings.ServerSettingsService; + yield* analytics.record("test.disabled-at-startup"); + yield* analytics.flush; + assert.deepEqual(batches, []); + yield* settings.updateSettings({ telemetryEnabled: true }); + yield* analytics.record("server.boot.heartbeat"); + yield* analytics.flush; + assert.equal(batches.length, 1); + yield* analytics.record("client.connected"); + yield* settings.updateSettings({ telemetryEnabled: false }); + yield* analytics.record("test.disabled"); + yield* settings.updateSettings({ telemetryEnabled: true }); + yield* analytics.flush; + assert.equal(batches.length, 1, "opt-out discarded queued events and retries"); + yield* analytics.record("client.connected"); + yield* analytics.flush; + assert.equal(batches.length, 2); + assert.deepEqual( + batches[1]?.map((event) => event.event), + ["client.connected"], + ); + assert.notEqual(batches[0]?.[0]?.uuid, batches[1]?.[0]?.uuid); + // Prevent the shutdown flush from retrying the new failed batch. + yield* settings.updateSettings({ telemetryEnabled: false }); + }).pipe(Effect.provide(layerRuntime)); + }), + ); + + it.effect.each(["successful", "failed"] as const)( + "opt-out during a %s send discards pending batches", + (firstRequestOutcome) => + Effect.gen(function* () { + const started = yield* Deferred.make(); + const release = yield* Deferred.make(); + const batches: Array = []; + + const layerClient = Layer.succeed( + HttpClient.HttpClient, + HttpClient.make((request) => + Effect.gen(function* () { + if (Predicate.isTagged(request.body, "Uint8Array")) { + const body = yield* decodeSentBatch( + new TextDecoder().decode(request.body.body), + ).pipe(Effect.orDie); + + batches.push(body.batch); + } + + if (batches.length === 1) { + yield* Deferred.succeed(started, undefined); + yield* Deferred.await(release); + + if (firstRequestOutcome === "failed") { + return yield* new HttpClientError.HttpClientError({ + reason: new HttpClientError.TransportError({ + request, + cause: "connection reset", + }), + }); + } + } + + return HttpClientResponse.fromWeb(request, new Response("{}")); + }), + ), + ); + + const layerRuntime = layerToggleTest(layerClient); + yield* Effect.gen(function* () { + const analytics = yield* AnalyticsService.AnalyticsService; + const settings = yield* ServerSettings.ServerSettingsService; + + for (let index = 0; index < 21; index++) { + yield* analytics.record("test.before-opt-out", { index }); + } + + const flushing = yield* analytics.flush.pipe(Effect.forkScoped); + yield* Deferred.await(started); + yield* settings.updateSettings({ telemetryEnabled: false }); + yield* settings.updateSettings({ telemetryEnabled: true }); + yield* analytics.record("test.after-opt-in"); + yield* Deferred.succeed(release, undefined); + yield* Fiber.join(flushing); + assert.equal(batches.length, 1, "the old flush stops after its in-flight request"); + yield* analytics.flush; + assert.equal(batches.length, 2); + assert.deepEqual( + batches[1]?.map((event) => event.event), + ["test.after-opt-in"], + "only the new event is sent without old retries", + ); + }).pipe(Effect.provide(layerRuntime)); + }), + ); }); diff --git a/apps/server/src/telemetry/AnalyticsService.ts b/apps/server/src/telemetry/AnalyticsService.ts index 3d10e93a26d6..3415b51517c8 100644 --- a/apps/server/src/telemetry/AnalyticsService.ts +++ b/apps/server/src/telemetry/AnalyticsService.ts @@ -22,19 +22,22 @@ import * as Option from "effect/Option"; import * as Random from "effect/Random"; import * as Ref from "effect/Ref"; import * as Result from "effect/Result"; +import * as Schema from "effect/Schema"; import * as Semaphore from "effect/Semaphore"; +import * as Stream from "effect/Stream"; import * as HttpClient from "effect/http/HttpClient"; import * as HttpClientRequest from "effect/http/HttpClientRequest"; import * as HttpClientResponse from "effect/http/HttpClientResponse"; import packageJson from "../../package.json" with { type: "json" }; import * as ServerConfig from "../config.ts"; +import * as ServerSettings from "../serverSettings.ts"; import { getTelemetryIdentifier } from "./Identify.ts"; interface BufferedAnalyticsEvent { readonly uuid: string; readonly event: string; - readonly properties?: Readonly>; + readonly properties: Parameters[1]; readonly capturedAt: string; } @@ -66,6 +69,10 @@ export function retryDelayMs(failures: number, random: number): number { return Math.round(ceiling / 2 + (ceiling / 2) * random); } +export const TelemetryEnabledConfig = Config.Boolean("T3CODE_TELEMETRY_ENABLED").pipe( + Config.withDefault(true), +); + const TelemetryEnvConfig = Config.all({ posthogKey: Config.String("T3CODE_POSTHOG_KEY").pipe( Config.withDefault("phc_XOWci4oZP4VvLiEyrFqkFjP4CZn55mjYYBMREK5Wd6m"), @@ -73,9 +80,12 @@ const TelemetryEnvConfig = Config.all({ posthogHost: Config.String("T3CODE_POSTHOG_HOST").pipe( Config.withDefault("https://us.i.posthog.com"), ), - enabled: Config.Boolean("T3CODE_TELEMETRY_ENABLED").pipe(Config.withDefault(true)), - flushBatchSize: Config.Number("T3CODE_TELEMETRY_FLUSH_BATCH_SIZE").pipe(Config.withDefault(20)), - maxBufferedEvents: Config.Number("T3CODE_TELEMETRY_MAX_BUFFERED_EVENTS").pipe( + enabled: TelemetryEnabledConfig, + flushBatchSize: Config.schema( + Schema.Int.check(Schema.isGreaterThan(0)), + "T3CODE_TELEMETRY_FLUSH_BATCH_SIZE", + ).pipe(Config.withDefault(20)), + maxBufferedEvents: Config.schema(Schema.Natural, "T3CODE_TELEMETRY_MAX_BUFFERED_EVENTS").pipe( Config.withDefault(1_000), ), wslDistroName: Config.String("WSL_DISTRO_NAME").pipe(Config.option), @@ -124,58 +134,51 @@ export const make = Effect.gen(function* () { const telemetryConfig = yield* TelemetryEnvConfig; const httpClient = yield* HttpClient.HttpClient; const serverConfig = yield* ServerConfig.ServerConfig; + const serverSettings = yield* ServerSettings.ServerSettingsService; + const settingsChanges = yield* serverSettings.subscribePersistedChanges; + + const initialEnabled = yield* telemetryConfig.enabled + ? serverSettings.getSettings.pipe( + Effect.map((settings) => settings.telemetryEnabled), + Effect.orElseSucceed(() => false), + ) + : Effect.succeed(false); + const enabledRef = yield* Ref.make(initialEnabled); + const enabled = Ref.get(enabledRef); const identifier = yield* getTelemetryIdentifier; const crypto = yield* Crypto.Crypto; - const bufferRef = yield* Ref.make>([]); - const deliveryRef = yield* Ref.make({ - failedBatch: [], - batchAttempts: 0, - failures: 0, - retryAt: 0, - }); + + const makeQueue = () => { + const events: BufferedAnalyticsEvent[] = []; + const delivery: DeliveryState = { failedBatch: [], batchAttempts: 0, failures: 0, retryAt: 0 }; + + return { events, delivery }; + }; + + let queue = makeQueue(); // The background flush and the shutdown flush must not send the same batch at once. const flushLock = yield* Semaphore.make(1); + + const discardQueuedEvents = Effect.sync(() => { + queue = makeQueue(); + }); + + yield* settingsChanges.pipe( + Stream.runForEach((settings) => + Ref.set(enabledRef, telemetryConfig.enabled && settings.telemetryEnabled).pipe( + Effect.andThen(settings.telemetryEnabled ? Effect.void : discardQueuedEvents), + ), + ), + Effect.forkScoped({ startImmediately: true }), + ); + const clientType = serverConfig.mode === "desktop" ? "desktop-app" : "cli-web-client"; const hostPlatform = yield* HostProcessPlatform; const hostArchitecture = yield* HostProcessArchitecture; - const enqueueBufferedEvent = ( - uuid: string, - event: string, - properties?: Readonly>, - ) => - Effect.flatMap(DateTime.now, (now) => - Ref.modify(bufferRef, (current) => { - const appended = [ - ...current, - { - uuid, - event, - ...(properties ? { properties } : {}), - capturedAt: DateTime.formatIso(now), - } satisfies BufferedAnalyticsEvent, - ]; - - const next = - appended.length > telemetryConfig.maxBufferedEvents - ? appended.slice(appended.length - telemetryConfig.maxBufferedEvents) - : appended; - - return [ - { - size: next.length, - dropped: next.length !== appended.length, - } as const, - next, - ] as const; - }), - ); - const sendBatch = Effect.fn("AnalyticsService.sendBatch")(function* ( events: ReadonlyArray, ) { - if (!telemetryConfig.enabled || !identifier) return; - const payload = { api_key: telemetryConfig.posthogKey, batch: events.map((event) => ({ @@ -208,39 +211,57 @@ export const make = Effect.gen(function* () { ); }); - const takeBatch = Ref.modify(bufferRef, (current) => { - const nextBatch = current.slice(0, telemetryConfig.flushBatchSize); - return [nextBatch, current.slice(nextBatch.length)] as const; - }); - // Sends batches until the buffer is empty or a send fails. A failed batch is // kept for the next flush, and dropped after MAX_BATCH_ATTEMPTS failed sends. - const flush: AnalyticsService["Service"]["flush"] = Effect.gen(function* () { + const flushBatches = Effect.gen(function* () { + const pending = queue; + while (true) { - const delivery = yield* Ref.get(deliveryRef); - const batch = delivery.failedBatch.length > 0 ? delivery.failedBatch : yield* takeBatch; + const delivery = pending.delivery; + + const batch = + delivery.failedBatch.length > 0 + ? delivery.failedBatch + : pending.events.splice(0, telemetryConfig.flushBatchSize); + if (batch.length === 0) { return; } + const telemetryEnabled = yield* enabled; + + if (pending !== queue) return; + + if (!identifier || !telemetryEnabled) { + yield* discardQueuedEvents; + + return; + } + const sent = yield* Effect.result(sendBatch(batch)); + + // Opt-out replaces the queue, including retries and unfinished recordings. + if (pending !== queue) return; + if (Result.isSuccess(sent)) { - yield* Ref.set(deliveryRef, { failedBatch: [], batchAttempts: 0, failures: 0, retryAt: 0 }); + pending.delivery = { failedBatch: [], batchAttempts: 0, failures: 0, retryAt: 0 }; continue; } const failures = delivery.failures + 1; const batchAttempts = delivery.batchAttempts + 1; const retryAt = (yield* Clock.currentTimeMillis) + retryDelayMs(failures, yield* Random.next); + if (batchAttempts < MAX_BATCH_ATTEMPTS) { - yield* Ref.set(deliveryRef, { failedBatch: batch, batchAttempts, failures, retryAt }); + pending.delivery = { failedBatch: batch, batchAttempts, failures, retryAt }; yield* Effect.logDebug("Failed to send telemetry batch; will retry", { attempt: batchAttempts, cause: sent.failure, }); return; } - yield* Ref.set(deliveryRef, { failedBatch: [], batchAttempts: 0, failures, retryAt }); + + pending.delivery = { failedBatch: [], batchAttempts: 0, failures, retryAt }; yield* Effect.logWarning("Dropped telemetry batch after repeated send failures", { events: batch.length, attempts: batchAttempts, @@ -248,28 +269,39 @@ export const make = Effect.gen(function* () { }); return; } - }).pipe(flushLock.withPermit); + }); + + const flush = flushBatches.pipe(flushLock.withPermit); const flushWhenDue = Effect.gen(function* () { - const { retryAt } = yield* Ref.get(deliveryRef); - if ((yield* Clock.currentTimeMillis) >= retryAt) { - yield* flush; + if ((yield* Clock.currentTimeMillis) >= queue.delivery.retryAt) { + yield* flushBatches; } - }); + }).pipe(flushLock.withPermit); const record: AnalyticsService["Service"]["record"] = Effect.fn("AnalyticsService.record")( function* (event, properties) { - if (!telemetryConfig.enabled || !identifier) return; + const pending = queue; + + if (!identifier || !(yield* enabled)) return; // Telemetry is best effort: an event without a uuid is not sent. The // Node implementation throws (a defect) rather than failing, so catch both. const uuid = yield* Effect.exit(crypto.randomUUIDv7); if (Exit.isFailure(uuid)) return; - const enqueueResult = yield* enqueueBufferedEvent(uuid.value, event, properties); - if (enqueueResult.dropped) { + const now = yield* DateTime.now; + pending.events.push({ + uuid: uuid.value, + event, + properties, + capturedAt: DateTime.formatIso(now), + }); + + if (pending.events.length > telemetryConfig.maxBufferedEvents) { + pending.events.splice(0, pending.events.length - telemetryConfig.maxBufferedEvents); yield* Effect.logDebug("analytics buffer full; dropping oldest event", { - size: enqueueResult.size, + size: pending.events.length, event, }); } @@ -286,5 +318,3 @@ export const make = Effect.gen(function* () { }); export const layer = Layer.effect(AnalyticsService, make); - -const layerTest = AnalyticsService.layerTest; diff --git a/apps/server/src/terminal/Manager.test.ts b/apps/server/src/terminal/Manager.test.ts index 5a493cc28053..7332d1e3c71d 100644 --- a/apps/server/src/terminal/Manager.test.ts +++ b/apps/server/src/terminal/Manager.test.ts @@ -2204,6 +2204,7 @@ it.layer( withSettingsSnapshot: () => Effect.fail(settingsError), streamChanges: Stream.empty, subscribeChanges: Effect.succeed(Stream.empty), + subscribePersistedChanges: Effect.succeed(Stream.empty), }); const error = yield* TerminalManager.resolveProviderInstanceTerminalEnvironment({ diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index 53bf27800583..0e66152bd3d5 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -1210,6 +1210,7 @@ const layerWsRpc = ( const providerSessionsV2 = yield* ProviderSessionManager.ProviderSessionManagerV2; const mcpAppRequests = yield* McpAppRequests.McpAppRequests; const analytics = yield* AnalyticsService.AnalyticsService; + const telemetryDisabledByEnvironment = !(yield* AnalyticsService.TelemetryEnabledConfig); // Client-origin attribution (#7774): every thread/turn the connecting // client starts is credited to its surface + app version. Best-effort: // attribution must never fail the user's command. @@ -1733,6 +1734,7 @@ const layerWsRpc = ( onSome: (root) => ({ scratchWorkspaceRoot: root }), }), newProjectsRoot: managedFolders.namedProjectsRoot, + telemetryDisabledByEnvironment, }; }); diff --git a/apps/web/src/components/settings/SettingsPanels.tsx b/apps/web/src/components/settings/SettingsPanels.tsx index bee57d21b253..f5fa298d5997 100644 --- a/apps/web/src/components/settings/SettingsPanels.tsx +++ b/apps/web/src/components/settings/SettingsPanels.tsx @@ -2207,6 +2207,26 @@ export function GeneralSettingsPanel() { (target) => target.serverConfig?.environment.capabilities.threadRestartContinuation === true, ); + const mixedTelemetryEnabled = useScopedSettingsMixed(["telemetryEnabled"]); + + const telemetryOverridden = connectedEnvironments.some( + (target) => target.serverConfig?.telemetryDisabledByEnvironment === true, + ); + + const telemetryPartiallyOverridden = + telemetryOverridden && + connectedEnvironments.some( + (target) => target.serverConfig?.telemetryDisabledByEnvironment !== true, + ); + + const telemetryMixed = telemetryOverridden ? telemetryPartiallyOverridden : mixedTelemetryEnabled; + + const telemetryDescription = telemetryPartiallyOverridden + ? "Disabled by some selected servers' environment configuration." + : telemetryOverridden + ? "Disabled by the server's environment configuration." + : "Share anonymous usage data to help improve T3 Code."; + const textGenerationProviders = serverProviders.filter( (provider) => provider.supportsTextGeneration !== false, ); @@ -3339,6 +3359,26 @@ export function GeneralSettingsPanel() { /> + + { + updateSettings({ telemetryEnabled: telemetryMixed ? false : Boolean(checked) }); + }} + aria-label="Anonymous analytics" + /> + } + /> + + {isElectron || HOSTED_APP_CHANNEL ? ( diff --git a/apps/web/src/components/settings/settingsSearch.ts b/apps/web/src/components/settings/settingsSearch.ts index 853886d81fcd..1ed5f7f088b2 100644 --- a/apps/web/src/components/settings/settingsSearch.ts +++ b/apps/web/src/components/settings/settingsSearch.ts @@ -510,6 +510,13 @@ export const SETTINGS_SEARCH_ITEMS = [ searchTerms: ["cli terminal shell path install command line"], desktopOnly: true, }, + { + id: "anonymous-analytics", + title: "Anonymous analytics", + to: "/settings/general", + scope: "environment-defaults", + searchTerms: ["privacy telemetry usage data tracking opt out posthog"], + }, { id: "privacy-policy", title: "Privacy policy", diff --git a/docs/user/telemetry.md b/docs/user/telemetry.md index 2f7e44825dc6..6c8e61a42d6a 100644 --- a/docs/user/telemetry.md +++ b/docs/user/telemetry.md @@ -7,7 +7,11 @@ turn result, duration, and main-agent token totals when available. Events do not include prompts, responses, file contents, authentication tokens, conversation IDs, raw provider events, or child-agent output. Child-agent token use is excluded from the totals. -To disable collection, set `T3CODE_TELEMETRY_ENABLED=false` in the server's environment before +To disable collection without restarting, turn off **Anonymous analytics** in **Settings → General** +on desktop or web, or **Settings → Maintenance** on mobile. This stops both server and client +usage events on the selected environments. + +You can also set `T3CODE_TELEMETRY_ENABLED=false` in the server's environment before starting it. This stops product events from being recorded or sent. The desktop app reads the variable from your shell profile (for example `~/.zshrc`) on macOS and diff --git a/packages/client-runtime/src/state/sharedSettings.test.ts b/packages/client-runtime/src/state/sharedSettings.test.ts index 6f94fe0f9037..e7193a5efb99 100644 --- a/packages/client-runtime/src/state/sharedSettings.test.ts +++ b/packages/client-runtime/src/state/sharedSettings.test.ts @@ -44,6 +44,25 @@ describe("supportsSharedSettingsSync", () => { }); describe("splitSharedServerPatch", () => { + it("shares analytics opt-out and detects environments still collecting", () => { + const patch = { telemetryEnabled: false }; + expect(splitSharedServerPatch(patch)).toEqual({ sharedPatch: patch, localPatch: {} }); + expect( + findSharedSettingsMismatches({ + primaryEnvironmentId: primaryId, + primarySettings: { ...DEFAULT_SERVER_SETTINGS, ...patch }, + environments: [ + { + environmentId: boxId, + label: "Remote Box", + syncEligible: true, + settings: DEFAULT_SERVER_SETTINGS, + }, + ], + }), + ).toEqual([{ environmentId: boxId, label: "Remote Box" }]); + }); + it("keeps project overrides local: project ids belong to one environment", () => { const patch = { projectSettingsOverrides: { [ProjectId.make("project")]: { defaultAutoPull: true } }, @@ -130,6 +149,7 @@ describe("pickSharedServerSettings", () => { "sidebarAutoSettleOnMerge", "snoozeLimitedThreads", "sourceControlWritingStyle", + "telemetryEnabled", "textGenerationModelSelection", ]); }); diff --git a/packages/client-runtime/src/state/sharedSettings.ts b/packages/client-runtime/src/state/sharedSettings.ts index 1d6fb5e14ca7..6ef986576c1b 100644 --- a/packages/client-runtime/src/state/sharedSettings.ts +++ b/packages/client-runtime/src/state/sharedSettings.ts @@ -30,6 +30,7 @@ const SHARED_SERVER_SETTING_KEYS = [ "newWorktreesStartFromOrigin", "sourceControlWritingStyle", "textGenerationModelSelection", + "telemetryEnabled", ] as const satisfies ReadonlyArray; export type SharedServerSettingKey = (typeof SHARED_SERVER_SETTING_KEYS)[number]; diff --git a/packages/contracts/src/server.ts b/packages/contracts/src/server.ts index c46c8b732a6c..dac461157293 100644 --- a/packages/contracts/src/server.ts +++ b/packages/contracts/src/server.ts @@ -694,6 +694,8 @@ export const ServerConfig = Schema.Struct({ * stays absent for subscribers that did not opt in. */ usageLimitSources: Schema.optional(UsageLimitSourceSnapshots), + /** Whether server environment configuration forces anonymous analytics off. */ + telemetryDisabledByEnvironment: Schema.optionalKey(Schema.Boolean), }); export type ServerConfig = typeof ServerConfig.Type; diff --git a/packages/contracts/src/settings.ts b/packages/contracts/src/settings.ts index 942fbb915ffe..87f684b83f2b 100644 --- a/packages/contracts/src/settings.ts +++ b/packages/contracts/src/settings.ts @@ -1480,6 +1480,7 @@ export const ServerSettings = Schema.Struct({ usageModelAliases: Schema.Record(TrimmedNonEmptyString, TrimmedNonEmptyString).pipe( Schema.withDecodingDefault(Effect.succeed({})), ), + telemetryEnabled: Schema.Boolean.pipe(Schema.withDecodingDefault(Effect.succeed(true))), }); export type ServerSettings = typeof ServerSettings.Type; @@ -1792,6 +1793,7 @@ export const ServerSettingsPatch = Schema.Struct({ usageModelAliases: Schema.optionalKey( Schema.Record(TrimmedNonEmptyString, Schema.NullOr(TrimmedNonEmptyString)), ), + telemetryEnabled: Schema.optionalKey(Schema.Boolean), }); export type ServerSettingsPatch = typeof ServerSettingsPatch.Type;