/**
* Anonymous PostHog telemetry service.
*
* Persists an installation-scoped anonymous identifier, buffers events in
* memory, and flushes batches over Effect's HTTP client. A failed batch is
* retried with backoff and dropped after a few tries. Each event carries a
* uuid, so PostHog can tell a retried copy from a new event.
*
* @module AnalyticsService
*/
import { HostProcessArchitecture, HostProcessPlatform } from "@t3tools/shared/hostProcess";
import type { ClientOs } from "@t3tools/contracts";
import * as Clock from "effect/Clock";
import * as Config from "effect/Config";
import * as Context from "effect/Context";
import * as Crypto from "effect/Crypto";
import * as DateTime from "effect/DateTime";
import * as Effect from "effect/Effect";
import * as Exit from "effect/Exit";
import * as Layer from "effect/Layer";
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 Semaphore from "effect/Semaphore";
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 { getTelemetryIdentifier } from "./Identify.ts";
interface BufferedAnalyticsEvent {
readonly uuid: string;
readonly event: string;
readonly properties?: Readonly<Record<string, unknown>>;
readonly capturedAt: string;
}
interface DeliveryState {
/** Batch that failed last. It is sent again before newer events. */
readonly failedBatch: ReadonlyArray<BufferedAnalyticsEvent>;
/** Failed sends of `failedBatch`. */
readonly batchAttempts: number;
/** Failed sends since the last success. Sets the backoff delay. */
readonly failures: number;
/** The background flush does not send before this time (epoch ms). */
readonly retryAt: number;
}
const FLUSH_INTERVAL_MS = 1_000;
// A hung send would hold the flush lock, and with it the shutdown flush.
const SEND_TIMEOUT = "10 seconds";
const MAX_BATCH_ATTEMPTS = 5;
const RETRY_BASE_DELAY_MS = 2_000;
const RETRY_MAX_DELAY_MS = 300_000;
/**
* Delay before the next send after `failures` consecutive failed sends. The
* ceiling doubles from 2s up to 5 minutes, and the delay is a random point in
* its upper half. `random` is in [0, 1).
*/
export function retryDelayMs(failures: number, random: number): number {
const ceiling = Math.min(RETRY_MAX_DELAY_MS, RETRY_BASE_DELAY_MS * 2 ** (failures - 1));
return Math.round(ceiling / 2 + (ceiling / 2) * random);
}
const TelemetryEnvConfig = Config.all({
posthogKey: Config.String("T3CODE_POSTHOG_KEY").pipe(
Config.withDefault("phc_XOWci4oZP4VvLiEyrFqkFjP4CZn55mjYYBMREK5Wd6m"),
),
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(
Config.withDefault(1_000),
),
wslDistroName: Config.String("WSL_DISTRO_NAME").pipe(Config.option),
});
export class AnalyticsService extends Context.Service<
AnalyticsService,
{
/** Record an anonymous event for best-effort buffered delivery. */
readonly record: (
event: string,
properties?: Readonly<Record<string, unknown>>,
) => Effect.Effect<void>;
/** Flush all currently queued telemetry events. */
readonly flush: Effect.Effect<void>;
}
>()("t3/telemetry/AnalyticsService") {
/** No-op layer for callers that intentionally disable telemetry. */
static readonly layerTest = Layer.succeed(
AnalyticsService,
AnalyticsService.of({
record: () => Effect.void,
flush: Effect.void,
}),
);
}
function serverOsFromNodePlatform(platform: string): ClientOs {
switch (platform) {
case "darwin":
return "macOS";
case "win32":
return "Windows";
case "linux":
return "Linux";
case "android":
return "Android";
default:
return "other";
}
}
/** @public Service construction is part of the canonical Effect module API. */
export const make = Effect.gen(function* () {
const telemetryConfig = yield* TelemetryEnvConfig;
const httpClient = yield* HttpClient.HttpClient;
const serverConfig = yield* ServerConfig.ServerConfig;
const identifier = yield* getTelemetryIdentifier;
const crypto = yield* Crypto.Crypto;
const bufferRef = yield* Ref.make<ReadonlyArray<BufferedAnalyticsEvent>>([]);
const deliveryRef = yield* Ref.make<DeliveryState>({
failedBatch: [],
batchAttempts: 0,
failures: 0,
retryAt: 0,
});
// The background flush and the shutdown flush must not send the same batch at once.
const flushLock = yield* Semaphore.make(1);
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<Record<string, unknown>>,
) =>
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<BufferedAnalyticsEvent>,
) {
if (!telemetryConfig.enabled || !identifier) return;
const payload = {
api_key: telemetryConfig.posthogKey,
batch: events.map((event) => ({
uuid: event.uuid,
event: event.event,
distinct_id: identifier,
properties: {
...event.properties,
$process_person_profile: false,
platform: hostPlatform,
wsl: Option.getOrUndefined(telemetryConfig.wslDistroName),
arch: hostArchitecture,
t3CodeVersion: packageJson.version,
clientType,
serverOs: serverOsFromNodePlatform(hostPlatform),
serverArch: hostArchitecture,
serverWslDistro: Option.getOrUndefined(telemetryConfig.wslDistroName),
serverAppVersion: packageJson.version,
serverMode: serverConfig.mode,
},
timestamp: event.capturedAt,
})),
};
yield* HttpClientRequest.post(`${telemetryConfig.posthogHost}/batch/`).pipe(
HttpClientRequest.bodyJson(payload),
Effect.flatMap(httpClient.execute),
Effect.flatMap(HttpClientResponse.filterStatusOk),
Effect.timeout(SEND_TIMEOUT),
);
});
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* () {
while (true) {
const delivery = yield* Ref.get(deliveryRef);
const batch = delivery.failedBatch.length > 0 ? delivery.failedBatch : yield* takeBatch;
if (batch.length === 0) {
return;
}
const sent = yield* Effect.result(sendBatch(batch));
if (Result.isSuccess(sent)) {
yield* Ref.set(deliveryRef, { 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 });
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 });
yield* Effect.logWarning("Dropped telemetry batch after repeated send failures", {
events: batch.length,
attempts: batchAttempts,
cause: sent.failure,
});
return;
}
}).pipe(flushLock.withPermit);
const flushWhenDue = Effect.gen(function* () {
const { retryAt } = yield* Ref.get(deliveryRef);
if ((yield* Clock.currentTimeMillis) >= retryAt) {
yield* flush;
}
});
const record: AnalyticsService["Service"]["record"] = Effect.fn("AnalyticsService.record")(
function* (event, properties) {
if (!telemetryConfig.enabled || !identifier) 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) {
yield* Effect.logDebug("analytics buffer full; dropping oldest event", {
size: enqueueResult.size,
event,
});
}
},
);
yield* Effect.forever(Effect.sleep(FLUSH_INTERVAL_MS).pipe(Effect.flatMap(() => flushWhenDue)), {
disableYield: true,
}).pipe(Effect.forkScoped);
yield* Effect.addFinalizer(() => flush);
return AnalyticsService.of({ record, flush });
});
export const layer = Layer.effect(AnalyticsService, make);
const layerTest = AnalyticsService.layerTest;