diff --git a/bun.lock b/bun.lock index a216ef036f..7935b9b459 100644 --- a/bun.lock +++ b/bun.lock @@ -627,6 +627,8 @@ "name": "@opencode-ai/stats-core", "version": "1.14.50", "dependencies": { + "@aws-sdk/client-athena": "3.933.0", + "@aws-sdk/client-firehose": "3.933.0", "@planetscale/database": "1.19.0", "drizzle-orm": "catalog:", "effect": "catalog:", @@ -765,7 +767,6 @@ "@opentui/core": "catalog:", "@opentui/keymap": "catalog:", "@opentui/solid": "catalog:", - "@pulumi/cloudflare": "6.10.0", "@types/bun": "catalog:", "@types/node": "catalog:", }, @@ -945,8 +946,12 @@ "@aws-crypto/util": ["@aws-crypto/util@5.2.0", "", { "dependencies": { "@aws-sdk/types": "^3.222.0", "@smithy/util-utf8": "^2.0.0", "tslib": "^2.6.2" } }, "sha512-4RkU9EsI6ZpBve5fseQlGNUWKMa1RLPQ1dnjnQoe07ldfIzcsGb5hC5W0Dm7u423KWzawlrpbjXBrXCEv9zazQ=="], + "@aws-sdk/client-athena": ["@aws-sdk/client-athena@3.933.0", "", { "dependencies": { "@aws-crypto/sha256-browser": "5.2.0", "@aws-crypto/sha256-js": "5.2.0", "@aws-sdk/core": "3.932.0", "@aws-sdk/credential-provider-node": "3.933.0", "@aws-sdk/middleware-host-header": "3.930.0", "@aws-sdk/middleware-logger": "3.930.0", "@aws-sdk/middleware-recursion-detection": "3.933.0", "@aws-sdk/middleware-user-agent": "3.932.0", "@aws-sdk/region-config-resolver": "3.930.0", "@aws-sdk/types": "3.930.0", "@aws-sdk/util-endpoints": "3.930.0", "@aws-sdk/util-user-agent-browser": "3.930.0", "@aws-sdk/util-user-agent-node": "3.932.0", "@smithy/config-resolver": "^4.4.3", "@smithy/core": "^3.18.2", "@smithy/fetch-http-handler": "^5.3.6", "@smithy/hash-node": "^4.2.5", "@smithy/invalid-dependency": "^4.2.5", "@smithy/middleware-content-length": "^4.2.5", "@smithy/middleware-endpoint": "^4.3.9", "@smithy/middleware-retry": "^4.4.9", "@smithy/middleware-serde": "^4.2.5", "@smithy/middleware-stack": "^4.2.5", "@smithy/node-config-provider": "^4.3.5", "@smithy/node-http-handler": "^4.4.5", "@smithy/protocol-http": "^5.3.5", "@smithy/smithy-client": "^4.9.5", "@smithy/types": "^4.9.0", "@smithy/url-parser": "^4.2.5", "@smithy/util-base64": "^4.3.0", "@smithy/util-body-length-browser": "^4.2.0", "@smithy/util-body-length-node": "^4.2.1", "@smithy/util-defaults-mode-browser": "^4.3.8", "@smithy/util-defaults-mode-node": "^4.2.11", "@smithy/util-endpoints": "^3.2.5", "@smithy/util-middleware": "^4.2.5", "@smithy/util-retry": "^4.2.5", "@smithy/util-utf8": "^4.2.0", "tslib": "^2.6.2" } }, "sha512-9eMUCu1Ay3C9ojo+dJcynSdpbxuwDVtZUt/Xhce+c2+mgDsmvRzjww+wfLpZwRNWxBWmeauQQAZk52tCwQgXsQ=="], + "@aws-sdk/client-cognito-identity": ["@aws-sdk/client-cognito-identity@3.993.0", "", { "dependencies": { "@aws-crypto/sha256-browser": "5.2.0", "@aws-crypto/sha256-js": "5.2.0", "@aws-sdk/core": "^3.973.11", "@aws-sdk/credential-provider-node": "^3.972.10", "@aws-sdk/middleware-host-header": "^3.972.3", "@aws-sdk/middleware-logger": "^3.972.3", "@aws-sdk/middleware-recursion-detection": "^3.972.3", "@aws-sdk/middleware-user-agent": "^3.972.11", "@aws-sdk/region-config-resolver": "^3.972.3", "@aws-sdk/types": "^3.973.1", "@aws-sdk/util-endpoints": "3.993.0", "@aws-sdk/util-user-agent-browser": "^3.972.3", "@aws-sdk/util-user-agent-node": "^3.972.9", "@smithy/config-resolver": "^4.4.6", "@smithy/core": "^3.23.2", "@smithy/fetch-http-handler": "^5.3.9", "@smithy/hash-node": "^4.2.8", "@smithy/invalid-dependency": "^4.2.8", "@smithy/middleware-content-length": "^4.2.8", "@smithy/middleware-endpoint": "^4.4.16", "@smithy/middleware-retry": "^4.4.33", "@smithy/middleware-serde": "^4.2.9", "@smithy/middleware-stack": "^4.2.8", "@smithy/node-config-provider": "^4.3.8", "@smithy/node-http-handler": "^4.4.10", "@smithy/protocol-http": "^5.3.8", "@smithy/smithy-client": "^4.11.5", "@smithy/types": "^4.12.0", "@smithy/url-parser": "^4.2.8", "@smithy/util-base64": "^4.3.0", "@smithy/util-body-length-browser": "^4.2.0", "@smithy/util-body-length-node": "^4.2.1", "@smithy/util-defaults-mode-browser": "^4.3.32", "@smithy/util-defaults-mode-node": "^4.2.35", "@smithy/util-endpoints": "^3.2.8", "@smithy/util-middleware": "^4.2.8", "@smithy/util-retry": "^4.2.8", "@smithy/util-utf8": "^4.2.0", "tslib": "^2.6.2" } }, "sha512-7Ne3Yk/bgQPVebAkv7W+RfhiwTRSbfER9BtbhOa2w/+dIr902LrJf6vrZlxiqaJbGj2ALx8M+ZK1YIHVxSwu9A=="], + "@aws-sdk/client-firehose": ["@aws-sdk/client-firehose@3.933.0", "", { "dependencies": { "@aws-crypto/sha256-browser": "5.2.0", "@aws-crypto/sha256-js": "5.2.0", "@aws-sdk/core": "3.932.0", "@aws-sdk/credential-provider-node": "3.933.0", "@aws-sdk/middleware-host-header": "3.930.0", "@aws-sdk/middleware-logger": "3.930.0", "@aws-sdk/middleware-recursion-detection": "3.933.0", "@aws-sdk/middleware-user-agent": "3.932.0", "@aws-sdk/region-config-resolver": "3.930.0", "@aws-sdk/types": "3.930.0", "@aws-sdk/util-endpoints": "3.930.0", "@aws-sdk/util-user-agent-browser": "3.930.0", "@aws-sdk/util-user-agent-node": "3.932.0", "@smithy/config-resolver": "^4.4.3", "@smithy/core": "^3.18.2", "@smithy/fetch-http-handler": "^5.3.6", "@smithy/hash-node": "^4.2.5", "@smithy/invalid-dependency": "^4.2.5", "@smithy/middleware-content-length": "^4.2.5", "@smithy/middleware-endpoint": "^4.3.9", "@smithy/middleware-retry": "^4.4.9", "@smithy/middleware-serde": "^4.2.5", "@smithy/middleware-stack": "^4.2.5", "@smithy/node-config-provider": "^4.3.5", "@smithy/node-http-handler": "^4.4.5", "@smithy/protocol-http": "^5.3.5", "@smithy/smithy-client": "^4.9.5", "@smithy/types": "^4.9.0", "@smithy/url-parser": "^4.2.5", "@smithy/util-base64": "^4.3.0", "@smithy/util-body-length-browser": "^4.2.0", "@smithy/util-body-length-node": "^4.2.1", "@smithy/util-defaults-mode-browser": "^4.3.8", "@smithy/util-defaults-mode-node": "^4.2.11", "@smithy/util-endpoints": "^3.2.5", "@smithy/util-middleware": "^4.2.5", "@smithy/util-retry": "^4.2.5", "@smithy/util-utf8": "^4.2.0", "tslib": "^2.6.2" } }, "sha512-tDrtgczN2lQsflLDPYu/wdOoyCZLVYtgzmWnYzSEOBWd/cp2AbuQ7D+FemSwUTzyoMTuhhIevyEJKzqsF+QYxA=="], + "@aws-sdk/client-lambda": ["@aws-sdk/client-lambda@3.1048.0", "", { "dependencies": { "@aws-crypto/sha256-browser": "5.2.0", "@aws-crypto/sha256-js": "5.2.0", "@aws-sdk/core": "^3.974.11", "@aws-sdk/credential-provider-node": "^3.972.42", "@aws-sdk/types": "^3.973.8", "@smithy/core": "^3.24.2", "@smithy/fetch-http-handler": "^5.4.2", "@smithy/node-http-handler": "^4.7.2", "@smithy/types": "^4.14.1", "tslib": "^2.6.2" } }, "sha512-ryEYNVdilyWkKsOs/7Xy/l7+qjtSz4sll8NpcWD6AtONxjG/5OMaAhxxDkQb4iBoNMKnISxsARzQAp/Wa8pXIg=="], "@aws-sdk/client-s3": ["@aws-sdk/client-s3@3.933.0", "", { "dependencies": { "@aws-crypto/sha1-browser": "5.2.0", "@aws-crypto/sha256-browser": "5.2.0", "@aws-crypto/sha256-js": "5.2.0", "@aws-sdk/core": "3.932.0", "@aws-sdk/credential-provider-node": "3.933.0", "@aws-sdk/middleware-bucket-endpoint": "3.930.0", "@aws-sdk/middleware-expect-continue": "3.930.0", "@aws-sdk/middleware-flexible-checksums": "3.932.0", "@aws-sdk/middleware-host-header": "3.930.0", "@aws-sdk/middleware-location-constraint": "3.930.0", "@aws-sdk/middleware-logger": "3.930.0", "@aws-sdk/middleware-recursion-detection": "3.933.0", "@aws-sdk/middleware-sdk-s3": "3.932.0", "@aws-sdk/middleware-ssec": "3.930.0", "@aws-sdk/middleware-user-agent": "3.932.0", "@aws-sdk/region-config-resolver": "3.930.0", "@aws-sdk/signature-v4-multi-region": "3.932.0", "@aws-sdk/types": "3.930.0", "@aws-sdk/util-endpoints": "3.930.0", "@aws-sdk/util-user-agent-browser": "3.930.0", "@aws-sdk/util-user-agent-node": "3.932.0", "@smithy/config-resolver": "^4.4.3", "@smithy/core": "^3.18.2", "@smithy/eventstream-serde-browser": "^4.2.5", "@smithy/eventstream-serde-config-resolver": "^4.3.5", "@smithy/eventstream-serde-node": "^4.2.5", "@smithy/fetch-http-handler": "^5.3.6", "@smithy/hash-blob-browser": "^4.2.6", "@smithy/hash-node": "^4.2.5", "@smithy/hash-stream-node": "^4.2.5", "@smithy/invalid-dependency": "^4.2.5", "@smithy/md5-js": "^4.2.5", "@smithy/middleware-content-length": "^4.2.5", "@smithy/middleware-endpoint": "^4.3.9", "@smithy/middleware-retry": "^4.4.9", "@smithy/middleware-serde": "^4.2.5", "@smithy/middleware-stack": "^4.2.5", "@smithy/node-config-provider": "^4.3.5", "@smithy/node-http-handler": "^4.4.5", "@smithy/protocol-http": "^5.3.5", "@smithy/smithy-client": "^4.9.5", "@smithy/types": "^4.9.0", "@smithy/url-parser": "^4.2.5", "@smithy/util-base64": "^4.3.0", "@smithy/util-body-length-browser": "^4.2.0", "@smithy/util-body-length-node": "^4.2.1", "@smithy/util-defaults-mode-browser": "^4.3.8", "@smithy/util-defaults-mode-node": "^4.2.11", "@smithy/util-endpoints": "^3.2.5", "@smithy/util-middleware": "^4.2.5", "@smithy/util-retry": "^4.2.5", "@smithy/util-stream": "^4.5.6", "@smithy/util-utf8": "^4.2.0", "@smithy/util-waiter": "^4.2.5", "tslib": "^2.6.2" } }, "sha512-KxwZvdxdCeWK6o8mpnb+kk7Kgb8V+8AjTwSXUWH1UAD85B0tjdo1cSfE5zoR5fWGol4Ml5RLez12a6LPhsoTqA=="], @@ -5229,6 +5234,16 @@ "@aws-crypto/util/@smithy/util-utf8": ["@smithy/util-utf8@2.3.0", "", { "dependencies": { "@smithy/util-buffer-from": "^2.2.0", "tslib": "^2.6.2" } }, "sha512-R8Rdn8Hy72KKcebgLiv8jQcQkXoLMOGGv5uI1/k0l+snqkOzQ1R0ChUBCxWMlBsFMekWjq0wRudIweFs7sKT5A=="], + "@aws-sdk/client-athena/@smithy/core": ["@smithy/core@3.24.3", "", { "dependencies": { "@aws-crypto/crc32": "5.2.0", "@smithy/types": "^4.14.2", "tslib": "^2.6.2" } }, "sha512-Ep/7tPamGY8mgESE3LyLKtxJyy6U52WWAqr/3wial47Sj4u3PiIF73AOGI27UyLy9duTkhZbgzodOfLV4TduZg=="], + + "@aws-sdk/client-athena/@smithy/fetch-http-handler": ["@smithy/fetch-http-handler@5.4.3", "", { "dependencies": { "@smithy/core": "^3.24.3", "@smithy/types": "^4.14.2", "tslib": "^2.6.2" } }, "sha512-F+DRf8IJazRJgYog2A/yJK7eYVc0rqTlRzO+5ZxjJd4WkZoKz0IJRncf7G6t1pdVT3kryJcwuTFhN1c5m6N47A=="], + + "@aws-sdk/client-athena/@smithy/node-http-handler": ["@smithy/node-http-handler@4.7.3", "", { "dependencies": { "@smithy/core": "^3.24.3", "@smithy/types": "^4.14.2", "tslib": "^2.6.2" } }, "sha512-/jPhevcTFPMVl6KNjbaI47iOg1zxC7IsnX4PQDGVZKMFceOXtB8IEYaB7a9VvkP/3oC60WzTeKocvSI7vLT0vA=="], + + "@aws-sdk/client-athena/@smithy/types": ["@smithy/types@4.14.2", "", { "dependencies": { "tslib": "^2.6.2" } }, "sha512-P+otAxbV4CqBybp7EkcJCrig63yE2E7PuNVOmilVMRcx/O+QDzGULTrKsq4DV13gSfak9ObPrWaHl/9bL5YcWw=="], + + "@aws-sdk/client-athena/@smithy/util-utf8": ["@smithy/util-utf8@4.2.2", "", { "dependencies": { "@smithy/util-buffer-from": "^4.2.2", "tslib": "^2.6.2" } }, "sha512-75MeYpjdWRe8M5E3AW0O4Cx3UadweS+cwdXjwYGBW5h/gxxnbeZ877sLPX/ZJA9GVTlL/qG0dXP29JWFCD1Ayw=="], + "@aws-sdk/client-cognito-identity/@aws-sdk/core": ["@aws-sdk/core@3.973.27", "", { "dependencies": { "@aws-sdk/types": "^3.973.7", "@aws-sdk/xml-builder": "^3.972.17", "@smithy/core": "^3.23.14", "@smithy/node-config-provider": "^4.3.13", "@smithy/property-provider": "^4.2.13", "@smithy/protocol-http": "^5.3.13", "@smithy/signature-v4": "^5.3.13", "@smithy/smithy-client": "^4.12.9", "@smithy/types": "^4.14.0", "@smithy/util-base64": "^4.3.2", "@smithy/util-middleware": "^4.2.13", "@smithy/util-utf8": "^4.2.2", "tslib": "^2.6.2" } }, "sha512-CUZ5m8hwMCH6OYI4Li/WgMfIEx10Q2PLI9Y3XOUTPGZJ53aZ0007jCv+X/ywsaERyKPdw5MRZWk877roQksQ4A=="], "@aws-sdk/client-cognito-identity/@aws-sdk/credential-provider-node": ["@aws-sdk/credential-provider-node@3.972.30", "", { "dependencies": { "@aws-sdk/credential-provider-env": "^3.972.25", "@aws-sdk/credential-provider-http": "^3.972.27", "@aws-sdk/credential-provider-ini": "^3.972.29", "@aws-sdk/credential-provider-process": "^3.972.25", "@aws-sdk/credential-provider-sso": "^3.972.29", "@aws-sdk/credential-provider-web-identity": "^3.972.29", "@aws-sdk/types": "^3.973.7", "@smithy/credential-provider-imds": "^4.2.13", "@smithy/property-provider": "^4.2.13", "@smithy/shared-ini-file-loader": "^4.4.8", "@smithy/types": "^4.14.0", "tslib": "^2.6.2" } }, "sha512-FMnAnWxc8PG+ZrZ2OBKzY4luCUJhe9CG0B9YwYr4pzrYGLXBS2rl+UoUvjGbAwiptxRL6hyA3lFn03Bv1TLqTw=="], @@ -5253,6 +5268,16 @@ "@aws-sdk/client-cognito-identity/@smithy/util-utf8": ["@smithy/util-utf8@4.2.2", "", { "dependencies": { "@smithy/util-buffer-from": "^4.2.2", "tslib": "^2.6.2" } }, "sha512-75MeYpjdWRe8M5E3AW0O4Cx3UadweS+cwdXjwYGBW5h/gxxnbeZ877sLPX/ZJA9GVTlL/qG0dXP29JWFCD1Ayw=="], + "@aws-sdk/client-firehose/@smithy/core": ["@smithy/core@3.24.3", "", { "dependencies": { "@aws-crypto/crc32": "5.2.0", "@smithy/types": "^4.14.2", "tslib": "^2.6.2" } }, "sha512-Ep/7tPamGY8mgESE3LyLKtxJyy6U52WWAqr/3wial47Sj4u3PiIF73AOGI27UyLy9duTkhZbgzodOfLV4TduZg=="], + + "@aws-sdk/client-firehose/@smithy/fetch-http-handler": ["@smithy/fetch-http-handler@5.4.3", "", { "dependencies": { "@smithy/core": "^3.24.3", "@smithy/types": "^4.14.2", "tslib": "^2.6.2" } }, "sha512-F+DRf8IJazRJgYog2A/yJK7eYVc0rqTlRzO+5ZxjJd4WkZoKz0IJRncf7G6t1pdVT3kryJcwuTFhN1c5m6N47A=="], + + "@aws-sdk/client-firehose/@smithy/node-http-handler": ["@smithy/node-http-handler@4.7.3", "", { "dependencies": { "@smithy/core": "^3.24.3", "@smithy/types": "^4.14.2", "tslib": "^2.6.2" } }, "sha512-/jPhevcTFPMVl6KNjbaI47iOg1zxC7IsnX4PQDGVZKMFceOXtB8IEYaB7a9VvkP/3oC60WzTeKocvSI7vLT0vA=="], + + "@aws-sdk/client-firehose/@smithy/types": ["@smithy/types@4.14.2", "", { "dependencies": { "tslib": "^2.6.2" } }, "sha512-P+otAxbV4CqBybp7EkcJCrig63yE2E7PuNVOmilVMRcx/O+QDzGULTrKsq4DV13gSfak9ObPrWaHl/9bL5YcWw=="], + + "@aws-sdk/client-firehose/@smithy/util-utf8": ["@smithy/util-utf8@4.2.2", "", { "dependencies": { "@smithy/util-buffer-from": "^4.2.2", "tslib": "^2.6.2" } }, "sha512-75MeYpjdWRe8M5E3AW0O4Cx3UadweS+cwdXjwYGBW5h/gxxnbeZ877sLPX/ZJA9GVTlL/qG0dXP29JWFCD1Ayw=="], + "@aws-sdk/client-lambda/@aws-sdk/core": ["@aws-sdk/core@3.974.11", "", { "dependencies": { "@aws-sdk/types": "^3.973.8", "@aws-sdk/xml-builder": "^3.972.24", "@aws/lambda-invoke-store": "^0.2.2", "@smithy/core": "^3.24.2", "@smithy/signature-v4": "^5.4.2", "@smithy/types": "^4.14.1", "bowser": "^2.11.0", "tslib": "^2.6.2" } }, "sha512-QpnINq5FZH6EOaDEkmHdT7eUunbvD27pDNQypaWjFyYz7Zl1q3UCMQErBZxpmfGfI7MvI2TlK8KTkgNpv8b1ug=="], "@aws-sdk/client-lambda/@aws-sdk/credential-provider-node": ["@aws-sdk/credential-provider-node@3.972.42", "", { "dependencies": { "@aws-sdk/credential-provider-env": "^3.972.37", "@aws-sdk/credential-provider-http": "^3.972.39", "@aws-sdk/credential-provider-ini": "^3.972.41", "@aws-sdk/credential-provider-process": "^3.972.37", "@aws-sdk/credential-provider-sso": "^3.972.41", "@aws-sdk/credential-provider-web-identity": "^3.972.41", "@aws-sdk/types": "^3.973.8", "@smithy/core": "^3.24.2", "@smithy/credential-provider-imds": "^4.3.2", "@smithy/types": "^4.14.1", "tslib": "^2.6.2" } }, "sha512-D4oon2zbqqsWOJUM99Gm3/ZyJ0IJvTXVN3PyloGb3kQEyI36fjCZheZj422lAgTWWd6TSHgiImLt3RIaLdv3dQ=="], diff --git a/infra/console.ts b/infra/console.ts index 3b19c697de..11d6c5b54f 100644 --- a/infra/console.ts +++ b/infra/console.ts @@ -1,6 +1,7 @@ import { domain } from "./stage" import { EMAILOCTOPUS_API_KEY } from "./app" import { SECRET } from "./secret" +import { lakeIngest } from "./stats" //////////////// // DATABASE @@ -240,7 +241,7 @@ const SALESFORCE_INSTANCE_URL = new sst.Secret("SALESFORCE_INSTANCE_URL") const logProcessor = new sst.cloudflare.Worker("LogProcessor", { handler: "packages/console/function/src/log-processor.ts", - link: [SECRET.HoneycombApiKey], + link: [SECRET.HoneycombApiKey, lakeIngest], }) new sst.cloudflare.x.SolidStart("Console", { diff --git a/infra/stats.ts b/infra/stats.ts index bafe85dce5..8113791ac5 100644 --- a/infra/stats.ts +++ b/infra/stats.ts @@ -1,11 +1,69 @@ -import { SECRET } from "./secret" - const domain = (() => { if ($app.stage === "production") return "stats.opencode.ai" if ($app.stage === "dev") return "stats.dev.opencode.ai" - return `stats-${$app.stage}.dev.opencode.ai` + return `stats.${$app.stage}.dev.opencode.ai` })() +const current = aws.getCallerIdentityOutput({}) +const partition = aws.getPartitionOutput({}) +const region = aws.getRegionOutput({}) + +const tableBucketName = `opencode-${$app.stage}-datalake` +const tableNamespaceName = "inference" +const tableName = "event" +const glueCatalogName = "s3tablescatalog" +const glueCatalogArn = $interpolate`arn:${partition.partition}:glue:${region.region}:${current.accountId}:catalog` +const glueS3TablesCatalogArn = $interpolate`${glueCatalogArn}/${glueCatalogName}` +const glueS3TablesChildCatalogArn = $interpolate`${glueS3TablesCatalogArn}/${tableBucketName}` +const glueS3TablesDatabaseArn = $interpolate`arn:${partition.partition}:glue:${region.region}:${current.accountId}:database/${glueCatalogName}/${tableBucketName}/${tableNamespaceName}` +const glueS3TablesTableArn = $interpolate`arn:${partition.partition}:glue:${region.region}:${current.accountId}:table/${glueCatalogName}/${tableBucketName}/${tableNamespaceName}/${tableName}` +const s3TablesBucketWildcardArn = $interpolate`arn:${partition.partition}:s3tables:${region.region}:${current.accountId}:bucket/*` + +const eventSchema = [ + { name: "event_timestamp", type: "string", required: false }, + { name: "event_date", type: "string", required: false }, + { name: "event_type", type: "string", required: false }, + { name: "dataset", type: "string", required: false }, + { name: "client", type: "string", required: false }, + { name: "source", type: "string", required: false }, + { name: "tier", type: "string", required: false }, + { name: "provider", type: "string", required: false }, + { name: "provider_model", type: "string", required: false }, + { name: "model", type: "string", required: false }, + { name: "session", type: "string", required: false }, + { name: "request", type: "string", required: false }, + { name: "user_agent", type: "string", required: false }, + { name: "ip", type: "string", required: false }, + { name: "status", type: "int", required: false }, + { name: "is_stream", type: "boolean", required: false }, + { name: "duration_ms", type: "long", required: false }, + { name: "ttfb_ms", type: "long", required: false }, + { name: "request_length", type: "long", required: false }, + { name: "response_length", type: "long", required: false }, + { name: "timestamp_first_byte", type: "long", required: false }, + { name: "timestamp_last_byte", type: "long", required: false }, + { name: "tokens_input", type: "long", required: false }, + { name: "tokens_output", type: "long", required: false }, + { name: "tokens_reasoning", type: "long", required: false }, + { name: "tokens_cache_read", type: "long", required: false }, + { name: "tokens_cache_write_5m", type: "long", required: false }, + { name: "tokens_cache_write_1h", type: "long", required: false }, + { name: "tokens_total", type: "long", required: false }, + { name: "cost_input_microcents", type: "long", required: false }, + { name: "cost_output_microcents", type: "long", required: false }, + { name: "cost_cache_read_microcents", type: "long", required: false }, + { name: "cost_cache_write_microcents", type: "long", required: false }, + { name: "cost_total_microcents", type: "long", required: false }, + { name: "output_tps", type: "double", required: false }, + { name: "cf_continent", type: "string", required: false }, + { name: "cf_country", type: "string", required: false }, + { name: "cf_city", type: "string", required: false }, + { name: "cf_region", type: "string", required: false }, + { name: "cf_latitude", type: "double", required: false }, + { name: "cf_longitude", type: "double", required: false }, + { name: "cf_timezone", type: "string", required: false }, +] + //////////////// // DATABASE //////////////// @@ -63,6 +121,268 @@ new sst.x.DevCommand("StatsStudio", { }, }) +//////////////// +// DATA LAKE +//////////////// + +const tableBucket = new aws.s3tables.TableBucket("StatsLakeTableBucket", { + name: tableBucketName, + forceDestroy: $app.stage !== "production", + tags: { + app: $app.name, + stage: $app.stage, + }, +}) + +const namespace = new aws.s3tables.Namespace("StatsLakeNamespace", { + namespace: tableNamespaceName, + tableBucketArn: tableBucket.arn, +}) + +const eventsTable = new aws.s3tables.Table("StatsLakeEventsTable", { + name: tableName, + namespace: namespace.namespace, + tableBucketArn: namespace.tableBucketArn, + format: "ICEBERG", + metadata: { + iceberg: { + schema: { + fields: eventSchema, + }, + }, + }, + tags: { + app: $app.name, + stage: $app.stage, + }, +}) + +const s3TablesCatalog = new aws.cloudcontrol.Resource( + "StatsLakeS3TablesCatalog", + { + typeName: "AWS::Glue::Catalog", + desiredState: $jsonStringify({ + Name: glueCatalogName, + Description: "Federated catalog for S3 Tables", + FederatedCatalog: { + Identifier: s3TablesBucketWildcardArn, + ConnectionName: "aws:s3tables", + }, + CreateDatabaseDefaultPermissions: [ + { + Principal: { + DataLakePrincipalIdentifier: "IAM_ALLOWED_PRINCIPALS", + }, + Permissions: ["ALL"], + }, + ], + CreateTableDefaultPermissions: [ + { + Principal: { + DataLakePrincipalIdentifier: "IAM_ALLOWED_PRINCIPALS", + }, + Permissions: ["ALL"], + }, + ], + AllowFullTableExternalDataAccess: "True", + }), + }, + { dependsOn: [tableBucket] }, +) + +const athenaResultsBucket = new aws.s3.Bucket("StatsLakeAthenaResults", { + bucket: `opencode-${$app.stage}-stats-athena-results`, + forceDestroy: $app.stage !== "production", + tags: { + app: $app.name, + stage: $app.stage, + }, +}) + +const firehoseErrorBucket = new aws.s3.Bucket("StatsLakeFirehoseErrors", { + bucket: `opencode-${$app.stage}-stats-firehose-errors`, + forceDestroy: $app.stage !== "production", + tags: { + app: $app.name, + stage: $app.stage, + }, +}) + +const athenaWorkgroup = new aws.athena.Workgroup("StatsLakeAthenaWorkgroup", { + name: `opencode-${$app.stage}-stats`, + forceDestroy: $app.stage !== "production", + configuration: { + enforceWorkgroupConfiguration: true, + publishCloudwatchMetricsEnabled: true, + resultConfiguration: { + outputLocation: $interpolate`s3://${athenaResultsBucket.bucket}/`, + }, + }, + tags: { + app: $app.name, + stage: $app.stage, + }, +}) + +const firehoseRole = new aws.iam.Role("StatsLakeFirehoseRole", { + assumeRolePolicy: aws.iam.getPolicyDocumentOutput({ + statements: [ + { + effect: "Allow", + actions: ["sts:AssumeRole"], + principals: [ + { + type: "Service", + identifiers: ["firehose.amazonaws.com"], + }, + ], + }, + ], + }).json, + tags: { + app: $app.name, + stage: $app.stage, + }, +}) + +const firehosePolicy = new aws.iam.RolePolicy("StatsLakeFirehosePolicy", { + role: firehoseRole.id, + policy: aws.iam.getPolicyDocumentOutput({ + statements: [ + { + effect: "Allow", + actions: [ + "s3tables:ListTableBuckets", + "s3tables:GetTableBucket", + "s3tables:GetNamespace", + "s3tables:GetTable", + "s3tables:GetTableData", + "s3tables:GetTableMetadataLocation", + "s3tables:ListNamespaces", + "s3tables:ListTables", + "s3tables:PutTableData", + "s3tables:UpdateTableMetadataLocation", + ], + resources: ["*"], + }, + { + effect: "Allow", + actions: [ + "glue:GetCatalog", + "glue:GetCatalogs", + "glue:GetDatabase", + "glue:GetDatabases", + "glue:GetTable", + "glue:GetTables", + "glue:UpdateTable", + ], + resources: [ + glueCatalogArn, + glueS3TablesCatalogArn, + $interpolate`${glueS3TablesCatalogArn}/*`, + glueS3TablesDatabaseArn, + glueS3TablesTableArn, + $interpolate`arn:${partition.partition}:glue:${region.region}:${current.accountId}:database/*`, + $interpolate`arn:${partition.partition}:glue:${region.region}:${current.accountId}:table/*/*`, + $interpolate`arn:${partition.partition}:glue:${region.region}:${current.accountId}:table/${glueCatalogName}/*`, + ], + }, + { + effect: "Allow", + actions: [ + "s3:AbortMultipartUpload", + "s3:GetBucketLocation", + "s3:GetObject", + "s3:ListBucket", + "s3:ListBucketMultipartUploads", + "s3:PutObject", + ], + resources: [firehoseErrorBucket.arn, $interpolate`${firehoseErrorBucket.arn}/*`], + }, + { + effect: "Allow", + actions: ["lakeformation:GetDataAccess"], + resources: ["*"], + }, + ], + }).json, +}) + +const firehose = new aws.kinesis.FirehoseDeliveryStream( + "StatsLakeFirehose", + { + name: `opencode-${$app.stage}-datalake-ingest`, + destination: "iceberg", + icebergConfiguration: { + appendOnly: true, + bufferingInterval: 60, + bufferingSize: 1, + catalogArn: glueS3TablesChildCatalogArn, + destinationTableConfigurations: [ + { + databaseName: namespace.namespace, + tableName: eventsTable.name, + s3ErrorOutputPrefix: "events/", + }, + ], + roleArn: firehoseRole.arn, + s3BackupMode: "FailedDataOnly", + s3Configuration: { + roleArn: firehoseRole.arn, + bucketArn: firehoseErrorBucket.arn, + errorOutputPrefix: "errors/!{firehose:error-output-type}/", + }, + }, + tags: { + app: $app.name, + stage: $app.stage, + }, + }, + { dependsOn: [s3TablesCatalog, eventsTable, firehosePolicy] }, +) + +export const lake = new sst.Linkable("StatsLake", { + properties: { + region: region.region, + catalog: $interpolate`${glueCatalogName}/${tableBucket.name}`, + database: namespace.namespace, + table: eventsTable.name, + tableBucket: tableBucket.name, + workgroup: athenaWorkgroup.name, + dataset: "zen", + }, +}) + +const ingestSecret = new random.RandomPassword("StatsLakeIngestSecret", { length: 32 }) + +const ingestConfig = new sst.Linkable("StatsLakeIngestConfig", { + properties: { + streamName: firehose.name, + secret: ingestSecret.result, + }, +}) + +const ingestFunction = new sst.aws.Function("StatsLakeIngestFunction", { + handler: "packages/stats/core/src/ingest.handler", + runtime: "nodejs22.x", + timeout: "30 seconds", + url: true, + link: [ingestConfig], + permissions: [ + { + actions: ["firehose:PutRecord", "firehose:PutRecordBatch"], + resources: [firehose.arn], + }, + ], +}) + +export const lakeIngest = new sst.Linkable("StatsLakeIngest", { + properties: { + url: ingestFunction.url, + secret: ingestSecret.result, + }, +}) + //////////////// // APP //////////////// @@ -92,6 +412,57 @@ export const statSync = new sst.aws.Cron("StatsSync", { handler: "packages/stats/core/src/cron/stat.handler", runtime: "nodejs22.x", timeout: "5 minutes", - link: [database, SECRET.HoneycombApiKey], + link: [database, lake], + permissions: [ + { + actions: ["athena:StartQueryExecution", "athena:GetQueryExecution", "athena:GetQueryResults"], + resources: [athenaWorkgroup.arn], + }, + { + actions: [ + "glue:GetCatalog", + "glue:GetCatalogs", + "glue:GetDatabase", + "glue:GetDatabases", + "glue:GetTable", + "glue:GetTables", + "glue:GetPartitions", + ], + resources: [ + glueCatalogArn, + glueS3TablesCatalogArn, + $interpolate`${glueS3TablesCatalogArn}/*`, + glueS3TablesDatabaseArn, + glueS3TablesTableArn, + $interpolate`arn:${partition.partition}:glue:${region.region}:${current.accountId}:database/*`, + $interpolate`arn:${partition.partition}:glue:${region.region}:${current.accountId}:table/*/*`, + $interpolate`arn:${partition.partition}:glue:${region.region}:${current.accountId}:table/${glueCatalogName}/*`, + ], + }, + { + actions: ["s3:GetBucketLocation", "s3:ListBucket"], + resources: [athenaResultsBucket.arn], + }, + { + actions: ["s3:GetObject", "s3:PutObject", "s3:AbortMultipartUpload", "s3:ListBucketMultipartUploads"], + resources: [$interpolate`${athenaResultsBucket.arn}/*`], + }, + { + actions: [ + "s3tables:GetTableBucket", + "s3tables:GetNamespace", + "s3tables:GetTable", + "s3tables:GetTableData", + "s3tables:GetTableMetadataLocation", + "s3tables:ListNamespaces", + "s3tables:ListTables", + ], + resources: ["*"], + }, + { + actions: ["lakeformation:GetDataAccess"], + resources: ["*"], + }, + ], }, }) diff --git a/packages/console/core/script/create-api-key.ts b/packages/console/core/script/create-api-key.ts new file mode 100644 index 0000000000..dba2ee946c --- /dev/null +++ b/packages/console/core/script/create-api-key.ts @@ -0,0 +1,146 @@ +import { Resource } from "@opencode-ai/console-resource" +import { and, Database, eq, isNull } from "../src/drizzle/index.js" +import { Identifier } from "../src/identifier.js" +import { AccountTable } from "../src/schema/account.sql.js" +import { AuthTable } from "../src/schema/auth.sql.js" +import { BillingTable } from "../src/schema/billing.sql.js" +import { KeyTable } from "../src/schema/key.sql.js" +import { UserTable } from "../src/schema/user.sql.js" +import { WorkspaceTable } from "../src/schema/workspace.sql.js" +import { centsToMicroCents } from "../src/util/price.js" + +const args = parseArgs(process.argv.slice(2)) +if (!args.email) { + console.error( + "Usage: bun script/create-api-key.ts --email [--workspace-id ] [--workspace-name ] [--key-name ] [--balance-dollars ] [--allow-production]", + ) + process.exit(1) +} +if (Resource.App.stage === "production" && !args.allowProduction) { + throw new Error("Refusing to create a production API key without --allow-production") +} + +const result = await Database.transaction(async (tx) => { + const auth = await tx + .select() + .from(AuthTable) + .where(and(eq(AuthTable.provider, "email"), eq(AuthTable.subject, args.email))) + .then((rows) => rows[0]) + const accountID = auth?.accountID ?? Identifier.create("account") + if (!auth) { + await tx.insert(AccountTable).values({ id: accountID }) + await tx.insert(AuthTable).values({ + id: Identifier.create("auth"), + provider: "email", + subject: args.email, + accountID, + }) + } + + const workspace = args.workspaceID + ? await tx + .select() + .from(WorkspaceTable) + .where(eq(WorkspaceTable.id, args.workspaceID)) + .then((rows) => rows[0]) + : await tx + .select({ workspace: WorkspaceTable }) + .from(UserTable) + .innerJoin(WorkspaceTable, eq(WorkspaceTable.id, UserTable.workspaceID)) + .where(and(eq(UserTable.accountID, accountID), isNull(UserTable.timeDeleted))) + .then((rows) => rows[0]?.workspace) + if (args.workspaceID && !workspace) throw new Error(`Workspace not found: ${args.workspaceID}`) + const workspaceID = workspace?.id ?? Identifier.create("workspace") + if (!workspace) { + await tx.insert(WorkspaceTable).values({ + id: workspaceID, + slug: null, + name: args.workspaceName ?? `${args.email} manual`, + }) + } + + const user = await tx + .select() + .from(UserTable) + .where( + and(eq(UserTable.workspaceID, workspaceID), eq(UserTable.accountID, accountID), isNull(UserTable.timeDeleted)), + ) + .then((rows) => rows[0]) + const userID = user?.id ?? Identifier.create("user") + if (!user) { + await tx.insert(UserTable).values({ + id: userID, + workspaceID, + accountID, + email: args.email, + name: args.email, + role: "admin", + }) + } + + const balance = centsToMicroCents(args.balanceDollars * 100) + const billing = await tx + .select() + .from(BillingTable) + .where(eq(BillingTable.workspaceID, workspaceID)) + .then((rows) => rows[0]) + if (!billing) { + await tx.insert(BillingTable).values({ + id: Identifier.create("billing"), + workspaceID, + balance, + }) + } else if (billing.balance < balance) { + await tx.update(BillingTable).set({ balance }).where(eq(BillingTable.workspaceID, workspaceID)) + } + + const secretKey = createSecretKey() + const keyID = Identifier.create("key") + await tx.insert(KeyTable).values({ + id: keyID, + workspaceID, + userID, + name: args.keyName ?? "Manual API Key", + key: secretKey, + timeUsed: null, + }) + + return { accountID, workspaceID, userID, keyID, secretKey } +}) + +console.log(JSON.stringify({ stage: Resource.App.stage, ...result }, null, 2)) + +function createSecretKey() { + const chars = "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789" + const values = new Uint32Array(64) + crypto.getRandomValues(values) + return `sk-${Array.from(values, (value) => chars[value % chars.length]).join("")}` +} + +function parseArgs(argv: string[]) { + const parsed = { + email: "", + workspaceID: "", + workspaceName: "", + keyName: "", + balanceDollars: 100, + allowProduction: false, + } + for (let index = 0; index < argv.length; index++) { + const arg = argv[index] + if (arg === "--email") parsed.email = requiredValue(argv, ++index, arg) + if (arg === "--workspace-id") parsed.workspaceID = requiredValue(argv, ++index, arg) + if (arg === "--workspace-name") parsed.workspaceName = requiredValue(argv, ++index, arg) + if (arg === "--key-name") parsed.keyName = requiredValue(argv, ++index, arg) + if (arg === "--balance-dollars") parsed.balanceDollars = Number(requiredValue(argv, ++index, arg)) + if (arg === "--allow-production") parsed.allowProduction = true + } + if (!Number.isFinite(parsed.balanceDollars) || parsed.balanceDollars < 0) throw new Error("Invalid --balance-dollars") + return parsed +} + +function requiredValue(argv: string[], index: number, arg: string) { + const value = argv[index] + if (!value || value.startsWith("--")) throw new Error(`Missing value for ${arg}`) + return value +} diff --git a/packages/console/function/src/log-processor.ts b/packages/console/function/src/log-processor.ts index 2bb741b7aa..6de344dd0a 100644 --- a/packages/console/function/src/log-processor.ts +++ b/packages/console/function/src/log-processor.ts @@ -1,6 +1,9 @@ import { Resource } from "@opencode-ai/console-resource" import type { TraceItem } from "@cloudflare/workers-types" +type MetricData = Record +type MetricEvent = { time: string; data: MetricData } + export default { async tail(events: TraceItem[]) { for (const event of events) { @@ -21,7 +24,7 @@ export default { ) continue - let data = { + let data: MetricData = { "cf.continent": event.event.request.cf?.continent, "cf.country": event.event.request.cf?.country, "cf.city": event.event.request.cf?.city, @@ -35,11 +38,11 @@ export default { ip: event.event.request.headers["x-real-ip"], } const time = new Date(event.eventTimestamp ?? Date.now()).toISOString() - const events = [] + const events: MetricEvent[] = [] for (const log of event.logs) { for (const message of log.message) { if (!message.startsWith("_metric:")) continue - const json = JSON.parse(message.slice(8)) + const json = JSON.parse(message.slice(8)) as MetricData data = { ...data, ...json } if ("llm.error.code" in json) { events.push({ time, data: { ...data, event_type: "llm.error" } }) @@ -49,16 +52,135 @@ export default { events.push({ time, data: { ...data, event_type: "completions" } }) console.log(JSON.stringify(data, null, 2)) - const ret = await fetch("https://api.honeycomb.io/1/batch/zen", { - method: "POST", - headers: { - "Content-Type": "application/json", - "X-Honeycomb-Team": Resource.HONEYCOMB_API_KEY.value, - }, - body: JSON.stringify(events), - }) - console.log(ret.status) - console.log(await ret.text()) + const [honeycomb, lake] = await Promise.all([ + fetch("https://api.honeycomb.io/1/batch/zen", { + method: "POST", + headers: { + "Content-Type": "application/json", + "X-Honeycomb-Team": Resource.HONEYCOMB_API_KEY.value, + }, + body: JSON.stringify(events), + }), + fetch(Resource.StatsLakeIngest.url, { + method: "POST", + headers: { + "Content-Type": "application/json", + Authorization: `Bearer ${Resource.StatsLakeIngest.secret}`, + }, + body: JSON.stringify({ events: events.map((event) => toLakeEvent(event.time, event.data)) }), + }), + ]) + console.log(honeycomb.status) + console.log(await honeycomb.text()) + console.log(lake.status) + console.log(await lake.text()) } }, } + +function toLakeEvent(time: string, data: MetricData) { + const tokensInput = integer(data, "tokens.input") + const tokensOutput = integer(data, "tokens.output") + const tokensReasoning = integer(data, "tokens.reasoning") + const tokensCacheRead = integer(data, "tokens.cache_read") + const tokensCacheWrite5m = integer(data, "tokens.cache_write_5m") + const tokensCacheWrite1h = integer(data, "tokens.cache_write_1h") + const timestampFirstByte = integer(data, "timestamp.first_byte") + const timestampLastByte = integer(data, "timestamp.last_byte") + const source = string(data, "source") + + return { + event_timestamp: time, + event_date: time.slice(0, 10), + event_type: string(data, "event_type"), + dataset: "zen", + client: string(data, "client"), + source, + tier: tier(source), + provider: string(data, "provider"), + provider_model: string(data, "provider.model"), + model: string(data, "model"), + session: string(data, "session"), + request: string(data, "request"), + user_agent: string(data, "user_agent"), + ip: string(data, "ip"), + status: integer(data, "status"), + is_stream: boolean(data, "is_stream"), + duration_ms: integer(data, "duration"), + ttfb_ms: integer(data, "time_to_first_byte"), + request_length: integer(data, "request_length"), + response_length: integer(data, "response_length"), + timestamp_first_byte: timestampFirstByte, + timestamp_last_byte: timestampLastByte, + tokens_input: tokensInput, + tokens_output: tokensOutput, + tokens_reasoning: tokensReasoning, + tokens_cache_read: tokensCacheRead, + tokens_cache_write_5m: tokensCacheWrite5m, + tokens_cache_write_1h: tokensCacheWrite1h, + tokens_total: + integer(data, "tokens") ?? + (tokensInput ?? 0) + + (tokensOutput ?? 0) + + (tokensReasoning ?? 0) + + (tokensCacheRead ?? 0) + + (tokensCacheWrite5m ?? 0) + + (tokensCacheWrite1h ?? 0), + cost_input_microcents: integer(data, "cost.input.microcents"), + cost_output_microcents: integer(data, "cost.output.microcents"), + cost_cache_read_microcents: integer(data, "cost.cache_read.microcents"), + cost_cache_write_microcents: integer(data, "cost.cache_write.microcents"), + cost_total_microcents: integer(data, "cost.total.microcents"), + output_tps: number(data, "tps.output") ?? outputTps(tokensOutput, timestampFirstByte, timestampLastByte), + cf_continent: string(data, "cf.continent"), + cf_country: string(data, "cf.country"), + cf_city: string(data, "cf.city"), + cf_region: string(data, "cf.region"), + cf_latitude: number(data, "cf.latitude"), + cf_longitude: number(data, "cf.longitude"), + cf_timezone: string(data, "cf.timezone"), + } +} + +function tier(source: string | undefined) { + if (source === "anonymous" || source === "free") return "Free" + if (source === "lite") return "Go" + if (source === "subscription" || source === "balance") return "Zen" + if (source === "byok") return "BYOK" + return undefined +} + +function outputTps(tokens: number | undefined, firstByte: number | undefined, lastByte: number | undefined) { + if (!tokens || !firstByte || !lastByte || lastByte <= firstByte) return undefined + return Number(((tokens / (lastByte - firstByte)) * 1000).toFixed(6)) +} + +function string(data: MetricData, key: string) { + const value = data[key] + if (typeof value === "string") return value + if (typeof value === "number" || typeof value === "boolean") return String(value) + return undefined +} + +function boolean(data: MetricData, key: string) { + const value = data[key] + if (typeof value === "boolean") return value + if (typeof value === "string") return value === "true" ? true : value === "false" ? false : undefined + return undefined +} + +function integer(data: MetricData, key: string) { + const value = number(data, key) + if (value === undefined) return undefined + return Math.round(value) +} + +function number(data: MetricData, key: string) { + const value = data[key] + if (typeof value === "number") return Number.isFinite(value) ? value : undefined + if (typeof value === "string") { + const parsed = Number(value) + return Number.isFinite(parsed) ? parsed : undefined + } + return undefined +} diff --git a/packages/console/resource/resource.cloudflare.ts b/packages/console/resource/resource.cloudflare.ts index 745212ca9c..5108042a70 100644 --- a/packages/console/resource/resource.cloudflare.ts +++ b/packages/console/resource/resource.cloudflare.ts @@ -4,11 +4,15 @@ export { waitUntil } from "cloudflare:workers" export const Resource = new Proxy( {}, { - get(_target, prop: string) { + get(_target, prop: string | symbol) { + if (typeof prop !== "string") return undefined if (prop in env) { // @ts-expect-error const value = env[prop] return typeof value === "string" ? JSON.parse(value) : value + } else if (`SST_RESOURCE_${prop}` in env) { + // @ts-expect-error + return JSON.parse(env[`SST_RESOURCE_${prop}`]) } else if (prop === "App") { // @ts-expect-error return JSON.parse(env.SST_RESOURCE_App) diff --git a/packages/console/resource/resource.node.ts b/packages/console/resource/resource.node.ts index 1470bacf26..ce11abcc44 100644 --- a/packages/console/resource/resource.node.ts +++ b/packages/console/resource/resource.node.ts @@ -11,6 +11,7 @@ export const Resource = new Proxy( { get(_target, prop: keyof typeof ResourceBase) { const value = ResourceBase[prop] + const secrets = ResourceBase as unknown as Record if ("type" in value) { // @ts-ignore if (value.type === "sst.cloudflare.Bucket") { @@ -21,11 +22,11 @@ export const Resource = new Proxy( // @ts-ignore if (value.type === "sst.cloudflare.Kv") { const client = new Cloudflare({ - apiToken: ResourceBase.CLOUDFLARE_API_TOKEN.value, + apiToken: secrets.CLOUDFLARE_API_TOKEN.value, }) // @ts-ignore const namespaceId = value.namespaceId - const accountId = ResourceBase.CLOUDFLARE_DEFAULT_ACCOUNT_ID.value + const accountId = secrets.CLOUDFLARE_DEFAULT_ACCOUNT_ID.value return { get: (k: string | string[]) => { const isMulti = Array.isArray(k) diff --git a/packages/stats/app/app.config.ts b/packages/stats/app/app.config.ts new file mode 100644 index 0000000000..17ed6bf999 --- /dev/null +++ b/packages/stats/app/app.config.ts @@ -0,0 +1,5 @@ +export default { + server: { + preset: "aws-lambda", + }, +} diff --git a/packages/stats/app/vite.config.ts b/packages/stats/app/vite.config.ts index b85db17590..afcf20617c 100644 --- a/packages/stats/app/vite.config.ts +++ b/packages/stats/app/vite.config.ts @@ -3,7 +3,7 @@ import { nitro } from "nitro/vite" import { defineConfig, type PluginOption } from "vite" export default defineConfig({ - plugins: [solidStart() as PluginOption, nitro()], + plugins: [solidStart() as PluginOption, nitro({ preset: "aws-lambda" })], server: { allowedHosts: true, }, diff --git a/packages/stats/core/package.json b/packages/stats/core/package.json index 923f682864..3b2b3053ba 100644 --- a/packages/stats/core/package.json +++ b/packages/stats/core/package.json @@ -21,6 +21,8 @@ "typecheck": "tsgo --noEmit" }, "dependencies": { + "@aws-sdk/client-athena": "3.933.0", + "@aws-sdk/client-firehose": "3.933.0", "@planetscale/database": "1.19.0", "drizzle-orm": "catalog:", "effect": "catalog:" diff --git a/packages/stats/core/src/cron/stat.ts b/packages/stats/core/src/cron/stat.ts index 5751e734ad..ca56211423 100644 --- a/packages/stats/core/src/cron/stat.ts +++ b/packages/stats/core/src/cron/stat.ts @@ -1,121 +1,89 @@ -import { createHash } from "node:crypto" +import { + AthenaClient, + GetQueryExecutionCommand, + GetQueryResultsCommand, + StartQueryExecutionCommand, + type Row, +} from "@aws-sdk/client-athena" import { Client } from "@planetscale/database" import { sql } from "drizzle-orm" import { drizzle } from "drizzle-orm/planetscale-serverless" -import { DateTime, Effect, Option, Schema } from "effect" +import { DateTime, Effect, Schema } from "effect" import { Resource } from "sst" import { stat } from "../database/schema" -const HONEYCOMB_API_URL = "https://api.honeycomb.io" -const HONEYCOMB_DATASET = "zen" -const DAY_SECONDS = 86_400 -const MAX_POLL_ATTEMPTS = 15 +const ATHENA_MAX_POLL_ATTEMPTS = 60 +const ATHENA_PAGE_SIZE = 1000 +const DATALAKE_INGESTION_LAG_MS = 5 * 60_000 const UPSERT_CHUNK_SIZE = 500 -type HoneycombScalar = string | number | boolean | null -type HoneycombData = Record -type HoneycombQueryResult = { results: HoneycombData[]; series: { time: Date; data: HoneycombData }[] } +type AthenaData = Record type StatRow = typeof stat.$inferInsert +type StatAggregate = { + grain: "day" | "week" + period_start: Date + period_end: Date + dataset: string + tier: string + provider: string + model: string + sessions: number + requests: number + input_tokens: number + output_tokens: number + reasoning_tokens: number + cache_read_tokens: number + total_tokens: number + input_cost_microcents: number + output_cost_microcents: number + total_cost_microcents: number + avg_duration_ms: number | null + p50_duration_ms: number | null + p95_duration_ms: number | null + avg_ttfb_ms: number | null + p50_ttfb_ms: number | null + p95_ttfb_ms: number | null + avg_output_tps: number | null + success_count: number + error_count: number + sample_count: number +} type SyncResult = { ok: true; rows: number; startedAt: string; periodStart: string; periodEnd: string } -type SyncError = HoneycombApiError | HoneycombQueryTimeoutError | StatDatabaseError +type SyncError = AthenaQueryError | AthenaQueryTimeoutError | StatDatabaseError -class HoneycombApiError extends Schema.TaggedErrorClass()("HoneycombApiError", { +class AthenaQueryError extends Schema.TaggedErrorClass()("AthenaQueryError", { message: Schema.String, - status: Schema.optional(Schema.Number), + queryExecutionId: Schema.optional(Schema.String), cause: Schema.optional(Schema.Defect), }) {} -class HoneycombQueryTimeoutError extends Schema.TaggedErrorClass()( - "HoneycombQueryTimeoutError", - { - message: Schema.String, - resultId: Schema.String, - }, -) {} +class AthenaQueryTimeoutError extends Schema.TaggedErrorClass()("AthenaQueryTimeoutError", { + message: Schema.String, + queryExecutionId: Schema.String, +}) {} class StatDatabaseError extends Schema.TaggedErrorClass()("StatDatabaseError", { message: Schema.String, cause: Schema.optional(Schema.Defect), }) {} -const decodeJson = Schema.decodeUnknownOption(Schema.UnknownFromJsonString) - -const calculations = [ - { op: "COUNT" }, - { op: "COUNT_DISTINCT", column: "session" }, - { op: "SUM", column: "tokens.input" }, - { op: "SUM", column: "tokens.output" }, - { op: "SUM", column: "tokens.reasoning" }, - { op: "SUM", column: "tokens.cache_read" }, - { op: "SUM", column: "tokens" }, - { op: "SUM", column: "cost.input.microcents" }, - { op: "SUM", column: "cost.output.microcents" }, - { op: "SUM", column: "cost.total.microcents" }, - { op: "AVG", column: "duration" }, - { op: "P50", column: "duration" }, - { op: "P95", column: "duration" }, - { op: "AVG", column: "time_to_first_byte" }, - { op: "P50", column: "time_to_first_byte" }, - { op: "P95", column: "time_to_first_byte" }, - { op: "AVG", column: "tps.output" }, - { op: "SUM", column: "stats_success" }, - { op: "SUM", column: "stats_error" }, -] as const - export function handler(): Promise { return Effect.runPromise(syncStats()) } const syncStats: () => Effect.Effect = Effect.fn("StatsCron.sync")(function* () { const startedAt = yield* DateTime.nowAsDate - const periodEnd = new Date(Math.floor(startedAt.getTime() / 3_600_000) * 3_600_000) + const periodEnd = new Date(Math.floor((startedAt.getTime() - DATALAKE_INGESTION_LAG_MS) / 60_000) * 60_000) const periodStart = new Date( Date.UTC(periodEnd.getUTCFullYear(), periodEnd.getUTCMonth(), periodEnd.getUTCDate() - 6), ) - yield* logHoneycombRuntimeCheck() + yield* logAthenaRuntimeCheck() - const result = yield* runHoneycombQuery({ - start_time: Math.floor(periodStart.getTime() / 1000), - end_time: Math.floor(periodEnd.getTime() / 1000), - granularity: DAY_SECONDS, - breakdowns: ["tier", "provider", "model"], - calculations, - calculated_fields: [ - { - name: "stats_success", - expression: `IF(AND(GTE($status, "200"), LT($status, "400")), 1, 0)`, - }, - { - name: "stats_error", - expression: `IF(GTE($status, "400"), 1, 0)`, - }, - ], - filters: [ - { column: "event_type", op: "=", value: "completions" }, - { column: "model", op: "exists" }, - { column: "user_agent", op: "contains", value: "opencode" }, - ], - filter_combination: "AND", - orders: [{ column: "tokens", op: "SUM", order: "descending" }], - limit: 1000, - }) + const aggregates = (yield* runAthenaQuery(buildStatsQuery(periodStart, periodEnd))).flatMap(toStatAggregate) const rows = rankRows([ - ...synthesizeAllTierRows( - collapseRows(result.results.map((item) => toStatRow("week", periodStart, periodEnd, item))), - ), - ...synthesizeAllTierRows( - collapseRows( - result.series.map((item) => - toStatRow( - "day", - item.time, - new Date(Math.min(item.time.getTime() + DAY_SECONDS * 1000, periodEnd.getTime())), - item.data, - ), - ), - ), - ), + ...synthesizeAllTierRows(collapseRows(aggregates.filter((item) => item.grain === "week").map(toStatRow))), + ...synthesizeAllTierRows(collapseRows(aggregates.filter((item) => item.grain === "day").map(toStatRow))), ]) yield* saveRows(rows) @@ -139,97 +107,89 @@ const syncStats: () => Effect.Effect = Effect.fn(" } }) -const runHoneycombQuery: ( - query: Record, -) => Effect.Effect = Effect.fn( - "StatsCron.runHoneycombQuery", -)(function* (query: Record) { - const created = asRecord(yield* honeycombRequest(`/1/queries/${HONEYCOMB_DATASET}`, "POST", query)) - const queryId = asString(created.id) - if (!queryId) return yield* new HoneycombApiError({ message: "Honeycomb did not return a query id" }) - - const queued = asRecord( - yield* honeycombRequest(`/1/query_results/${HONEYCOMB_DATASET}`, "POST", { - query_id: queryId, - disable_series: false, - disable_total_by_aggregate: true, - disable_other_by_aggregate: true, - limit: 1000, - }), - ) - const resultId = asString(queued.id) - if (!resultId) return yield* new HoneycombApiError({ message: "Honeycomb did not return a query result id" }) - - return yield* pollHoneycombResult(resultId) -}) - -const pollHoneycombResult: ( - resultId: string, - attempt?: number, -) => Effect.Effect = Effect.fn( - "StatsCron.pollHoneycombResult", -)(function* (resultId: string, attempt = 0) { - if (attempt > 0) yield* Effect.sleep("1000 millis") - const result = asRecord(yield* honeycombRequest(`/1/query_results/${HONEYCOMB_DATASET}/${resultId}`, "GET")) - - if (result.complete === true) { - const data = asRecord(result.data) - return { - results: asArray(data.results).map((item) => asData(asRecord(item).data)), - series: asArray(data.series).flatMap((item) => { - const time = new Date(String(asRecord(item).time ?? "")) - if (Number.isNaN(time.getTime())) return [] - return [{ time, data: asData(asRecord(item).data) }] - }), - } - } - - if (attempt >= MAX_POLL_ATTEMPTS - 1) - return yield* new HoneycombQueryTimeoutError({ - message: `Honeycomb query result ${resultId} did not complete`, - resultId, - }) - - return yield* pollHoneycombResult(resultId, attempt + 1) -}) - -const honeycombRequest: ( - path: string, - method: "GET" | "POST", - body?: Record, -) => Effect.Effect = Effect.fn("StatsCron.honeycombRequest")(function* ( - path: string, - method: "GET" | "POST", - body?: Record, -) { - const response = yield* Effect.tryPromise({ +const runAthenaQuery: ( + query: string, +) => Effect.Effect = Effect.fn( + "StatsCron.runAthenaQuery", +)(function* (query: string) { + const client = new AthenaClient({ region: Resource.StatsLake.region }) + const started = yield* Effect.tryPromise({ try: () => - fetch(`${HONEYCOMB_API_URL}${path}`, { - method, - headers: { - "Content-Type": "application/json", - "X-Honeycomb-Team": Resource.HONEYCOMB_API_KEY.value, - }, - body: body ? JSON.stringify(body) : undefined, - }), - catch: (cause) => new HoneycombApiError({ message: `Honeycomb ${method} ${path} request failed`, cause }), - }) - const text = yield* Effect.tryPromise({ - try: () => response.text(), - catch: (cause) => new HoneycombApiError({ message: `Honeycomb ${method} ${path} response read failed`, cause }), + client.send( + new StartQueryExecutionCommand({ + QueryString: query, + WorkGroup: Resource.StatsLake.workgroup, + QueryExecutionContext: { + Catalog: Resource.StatsLake.catalog, + Database: Resource.StatsLake.database, + }, + }), + ), + catch: (cause) => new AthenaQueryError({ message: "Failed to start Athena stats query", cause }), }) + const queryExecutionId = started.QueryExecutionId + if (!queryExecutionId) return yield* new AthenaQueryError({ message: "Athena did not return a query execution id" }) - if (!response.ok) - return yield* new HoneycombApiError({ - message: `Honeycomb ${method} ${path} failed: ${response.status} ${text.slice(0, 500)}`, - status: response.status, + yield* pollAthenaQuery(client, queryExecutionId) + return yield* getAthenaResults(client, queryExecutionId) +}) + +const pollAthenaQuery: ( + client: AthenaClient, + queryExecutionId: string, + attempt?: number, +) => Effect.Effect = Effect.fn("StatsCron.pollAthenaQuery")( + function* (client: AthenaClient, queryExecutionId: string, attempt = 0) { + if (attempt > 0) yield* Effect.sleep("2 seconds") + + const result = yield* Effect.tryPromise({ + try: () => client.send(new GetQueryExecutionCommand({ QueryExecutionId: queryExecutionId })), + catch: (cause) => new AthenaQueryError({ message: "Failed to poll Athena stats query", queryExecutionId, cause }), }) - if (!text) return {} + const status = result.QueryExecution?.Status - const parsed = decodeJson(text) - if (Option.isNone(parsed)) - return yield* new HoneycombApiError({ message: `Honeycomb ${method} ${path} returned invalid JSON` }) - return parsed.value + if (status?.State === "SUCCEEDED") return + if (status?.State === "FAILED" || status?.State === "CANCELLED") + return yield* new AthenaQueryError({ + message: `Athena stats query ${status.State.toLowerCase()}: ${status.StateChangeReason ?? "unknown reason"}`, + queryExecutionId, + }) + + if (attempt >= ATHENA_MAX_POLL_ATTEMPTS - 1) + return yield* new AthenaQueryTimeoutError({ + message: `Athena stats query ${queryExecutionId} did not complete`, + queryExecutionId, + }) + + return yield* pollAthenaQuery(client, queryExecutionId, attempt + 1) + }, +) + +const getAthenaResults: ( + client: AthenaClient, + queryExecutionId: string, + nextToken?: string, +) => Effect.Effect = Effect.fn("StatsCron.getAthenaResults")(function* ( + client: AthenaClient, + queryExecutionId: string, + nextToken?: string, +) { + const result = yield* Effect.tryPromise({ + try: () => + client.send( + new GetQueryResultsCommand({ + QueryExecutionId: queryExecutionId, + NextToken: nextToken, + MaxResults: ATHENA_PAGE_SIZE, + }), + ), + catch: (cause) => new AthenaQueryError({ message: "Failed to read Athena stats results", queryExecutionId, cause }), + }) + const columns = result.ResultSet?.ResultSetMetadata?.ColumnInfo?.map((item) => item.Name ?? "") ?? [] + const rows = (result.ResultSet?.Rows ?? []).slice(nextToken ? 0 : 1).map((row) => rowData(columns, row)) + + if (!result.NextToken) return rows + return [...rows, ...(yield* getAthenaResults(client, queryExecutionId, result.NextToken))] }) const saveRows: (rows: StatRow[]) => Effect.Effect = Effect.fn("StatsCron.saveRows")( @@ -286,13 +246,101 @@ const saveRows: (rows: StatRow[]) => Effect.Effect= 200 AND status < 400 THEN 1 ELSE 0 END) AS success_count, + SUM(CASE WHEN status >= 400 THEN 1 ELSE 0 END) AS error_count, + COUNT(*) AS sample_count` + + return ` +WITH filtered AS ( + SELECT + from_iso8601_timestamp(event_timestamp) AS event_time, + COALESCE(NULLIF(tier, ''), 'unknown') AS tier, + COALESCE(NULLIF(provider, ''), 'unknown') AS provider, + COALESCE(NULLIF(model, ''), 'unknown') AS model, + session, + status, + duration_ms, + ttfb_ms, + output_tps, + tokens_input, + tokens_output, + tokens_reasoning, + tokens_cache_read, + tokens_total, + cost_input_microcents, + cost_output_microcents, + cost_total_microcents + FROM ${sourceTable} + WHERE event_type = 'completions' + AND model IS NOT NULL + AND model <> '' + AND user_agent LIKE '%opencode%' + AND event_timestamp >= ${periodStartValue} + AND event_timestamp < ${periodEndValue} +), daily AS ( + SELECT date_trunc('day', event_time) AS day, * + FROM filtered +) +SELECT + 'week' AS grain, + ${periodStartValue} AS period_start, + ${periodEndValue} AS period_end, + ${sqlString(Resource.StatsLake.dataset)} AS dataset, + tier, + provider, + model, + ${aggregateColumns} +FROM filtered +GROUP BY tier, provider, model +UNION ALL +SELECT + 'day' AS grain, + to_iso8601(day) AS period_start, + to_iso8601(least(day + INTERVAL '1' DAY, from_iso8601_timestamp(${periodEndValue}))) AS period_end, + ${sqlString(Resource.StatsLake.dataset)} AS dataset, + tier, + provider, + model, + ${aggregateColumns} +FROM daily +GROUP BY day, tier, provider, model +ORDER BY grain, period_start, total_tokens DESC +` +} + +function logAthenaRuntimeCheck() { + return Effect.logInfo("athena stats runtime check").pipe( Effect.annotateLogs({ - hasHoneycombApiKey: Boolean(Resource.HONEYCOMB_API_KEY.value), - honeycombApiKeyLength: Resource.HONEYCOMB_API_KEY.value.length, - honeycombApiKeySha256: createHash("sha256").update(Resource.HONEYCOMB_API_KEY.value).digest("hex").slice(0, 12), - honeycombApiUrl: HONEYCOMB_API_URL, + catalog: Resource.StatsLake.catalog, + database: Resource.StatsLake.database, + table: Resource.StatsLake.table, + workgroup: Resource.StatsLake.workgroup, + region: Resource.StatsLake.region, + stage: Resource.App.stage, }), ) } @@ -301,38 +349,77 @@ function inserted(column: string) { return sql.raw(`values(\`${column}\`)`) } -function toStatRow(grain: "day" | "week", periodStart: Date, periodEnd: Date, data: HoneycombData): StatRow { +function toStatAggregate(data: AthenaData): StatAggregate[] { + const grain = data.grain === "day" || data.grain === "week" ? data.grain : undefined + const periodStart = new Date(data.period_start ?? "") + const periodEnd = new Date(data.period_end ?? "") + if (!grain || Number.isNaN(periodStart.getTime()) || Number.isNaN(periodEnd.getTime())) return [] + + return [ + { + grain, + period_start: periodStart, + period_end: periodEnd, + dataset: data.dataset || Resource.StatsLake.dataset, + tier: normalizeTier(data.tier || "unknown"), + provider: data.provider || "unknown", + model: data.model || "unknown", + sessions: integer(data, "sessions"), + requests: integer(data, "requests"), + input_tokens: integer(data, "input_tokens"), + output_tokens: integer(data, "output_tokens"), + reasoning_tokens: integer(data, "reasoning_tokens"), + cache_read_tokens: integer(data, "cache_read_tokens"), + total_tokens: integer(data, "total_tokens"), + input_cost_microcents: integer(data, "input_cost_microcents"), + output_cost_microcents: integer(data, "output_cost_microcents"), + total_cost_microcents: integer(data, "total_cost_microcents"), + avg_duration_ms: nullableNumber(data, "avg_duration_ms"), + p50_duration_ms: nullableInteger(data, "p50_duration_ms"), + p95_duration_ms: nullableInteger(data, "p95_duration_ms"), + avg_ttfb_ms: nullableNumber(data, "avg_ttfb_ms"), + p50_ttfb_ms: nullableInteger(data, "p50_ttfb_ms"), + p95_ttfb_ms: nullableInteger(data, "p95_ttfb_ms"), + avg_output_tps: nullableNumber(data, "avg_output_tps"), + success_count: integer(data, "success_count"), + error_count: integer(data, "error_count"), + sample_count: integer(data, "sample_count"), + }, + ] +} + +function toStatRow(data: StatAggregate): StatRow { return { - grain, - period_start: periodStart, - period_end: periodEnd, - dataset: HONEYCOMB_DATASET, - tier: normalizeTier(asString(data.tier) || "unknown"), + grain: data.grain, + period_start: data.period_start, + period_end: data.period_end, + dataset: data.dataset, + tier: data.tier, client: "all", source: "all", - provider: asString(data.provider) || "unknown", - model: asString(data.model) || "unknown", + provider: data.provider, + model: data.model, provider_model: "", - sessions: Math.round(number(data, "COUNT_DISTINCT(session)")), - requests: Math.round(number(data, "COUNT")), - input_tokens: Math.round(number(data, "SUM(tokens.input)")), - output_tokens: Math.round(number(data, "SUM(tokens.output)")), - reasoning_tokens: Math.round(number(data, "SUM(tokens.reasoning)")), - cache_read_tokens: Math.round(number(data, "SUM(tokens.cache_read)")), - total_tokens: Math.round(number(data, "SUM(tokens)")), - input_cost_microcents: Math.round(number(data, "SUM(cost.input.microcents)")), - output_cost_microcents: Math.round(number(data, "SUM(cost.output.microcents)")), - total_cost_microcents: Math.round(number(data, "SUM(cost.total.microcents)")), - avg_duration_ms: nullableNumber(data, "AVG(duration)"), - p50_duration_ms: nullableInteger(data, "P50(duration)"), - p95_duration_ms: nullableInteger(data, "P95(duration)"), - avg_ttfb_ms: nullableNumber(data, "AVG(time_to_first_byte)"), - p50_ttfb_ms: nullableInteger(data, "P50(time_to_first_byte)"), - p95_ttfb_ms: nullableInteger(data, "P95(time_to_first_byte)"), - avg_output_tps: nullableNumber(data, "AVG(tps.output)"), - success_count: Math.round(number(data, "SUM(stats_success)")), - error_count: Math.round(number(data, "SUM(stats_error)")), - sample_count: Math.round(number(data, "COUNT")), + sessions: data.sessions, + requests: data.requests, + input_tokens: data.input_tokens, + output_tokens: data.output_tokens, + reasoning_tokens: data.reasoning_tokens, + cache_read_tokens: data.cache_read_tokens, + total_tokens: data.total_tokens, + input_cost_microcents: data.input_cost_microcents, + output_cost_microcents: data.output_cost_microcents, + total_cost_microcents: data.total_cost_microcents, + avg_duration_ms: data.avg_duration_ms, + p50_duration_ms: data.p50_duration_ms, + p95_duration_ms: data.p95_duration_ms, + avg_ttfb_ms: data.avg_ttfb_ms, + p50_ttfb_ms: data.p50_ttfb_ms, + p95_ttfb_ms: data.p95_ttfb_ms, + avg_output_tps: data.avg_output_tps, + success_count: data.success_count, + error_count: data.error_count, + sample_count: data.sample_count, } } @@ -452,49 +539,39 @@ function normalizeTier(value: string) { return value } -function number(data: HoneycombData, key: string) { - const value = data[key] - if (typeof value === "number") return Number.isFinite(value) ? value : 0 - if (typeof value === "string") { - const parsed = Number(value) - return Number.isFinite(parsed) ? parsed : 0 - } - return 0 -} - -function nullableNumber(data: HoneycombData, key: string) { - const value = number(data, key) - if (value === 0 && data[key] === undefined) return null - return Number(value.toFixed(2)) -} - -function nullableInteger(data: HoneycombData, key: string) { - if (data[key] === undefined) return null +function integer(data: AthenaData, key: string) { return Math.round(number(data, key)) } -function asData(value: unknown): HoneycombData { +function nullableNumber(data: AthenaData, key: string) { + if (data[key] === undefined || data[key] === "") return null + return Number(number(data, key).toFixed(2)) +} + +function nullableInteger(data: AthenaData, key: string) { + if (data[key] === undefined || data[key] === "") return null + return Math.round(number(data, key)) +} + +function number(data: AthenaData, key: string) { + const value = Number(data[key]) + return Number.isFinite(value) ? value : 0 +} + +function rowData(columns: string[], row: Row): AthenaData { return Object.fromEntries( - Object.entries(asRecord(value)).flatMap(([key, item]) => { - if (typeof item === "string" || typeof item === "number" || typeof item === "boolean" || item === null) - return [[key, item]] - return [] + columns.flatMap((column, index) => { + const value = row.Data?.[index]?.VarCharValue + if (!column || value === undefined) return [] + return [[column, value]] }), ) } -function asRecord(value: unknown): Record { - if (!value || typeof value !== "object" || Array.isArray(value)) return {} - return value as Record +function sqlIdentifier(value: string) { + return `"${value.replace(/"/g, '""')}"` } -function asArray(value: unknown) { - if (!Array.isArray(value)) return [] - return value -} - -function asString(value: unknown) { - if (typeof value === "string") return value - if (typeof value === "number" || typeof value === "boolean") return String(value) - return "" +function sqlString(value: string) { + return `'${value.replace(/'/g, "''")}'` } diff --git a/packages/stats/core/src/ingest.ts b/packages/stats/core/src/ingest.ts new file mode 100644 index 0000000000..c3548e6bd8 --- /dev/null +++ b/packages/stats/core/src/ingest.ts @@ -0,0 +1,98 @@ +import { Buffer } from "node:buffer" +import { timingSafeEqual } from "node:crypto" +import { FirehoseClient, PutRecordBatchCommand } from "@aws-sdk/client-firehose" +import { Resource } from "sst" + +const MAX_FIREHOSE_BATCH_SIZE = 500 +const MAX_FIREHOSE_ATTEMPTS = 3 +const client = new FirehoseClient({}) + +type FirehoseRecord = { Data: Uint8Array } + +type FunctionUrlEvent = { + body?: string | null + headers?: Record + isBase64Encoded?: boolean + requestContext?: { + http?: { + method?: string + } + } +} + +type IngestPayload = { + events?: unknown +} + +export async function handler(event: FunctionUrlEvent) { + if (event.requestContext?.http?.method !== "POST") return response(405, { ok: false, error: "Method Not Allowed" }) + if (!isAuthorized(event.headers ?? {})) return response(401, { ok: false, error: "Unauthorized" }) + + const payload = parsePayload(event) + if (!payload) return response(400, { ok: false, error: "Invalid JSON body" }) + + const events = Array.isArray(payload.events) + ? payload.events.filter( + (item): item is Record => Boolean(item) && typeof item === "object" && !Array.isArray(item), + ) + : [] + if (events.length === 0) return response(202, { ok: true, records: 0 }) + + const failed = ( + await Promise.all( + chunks( + events.map((item) => ({ Data: Buffer.from(JSON.stringify(item)) })), + MAX_FIREHOSE_BATCH_SIZE, + ).map((batch) => putRecords(Resource.StatsLakeIngestConfig.streamName, batch)), + ) + ).reduce((sum, item) => sum + item, 0) + if (failed > 0) return response(502, { ok: false, records: events.length, failed }) + + return response(202, { ok: true, records: events.length }) +} + +async function putRecords(streamName: string, records: FirehoseRecord[], attempt = 1): Promise { + const result = await client.send(new PutRecordBatchCommand({ DeliveryStreamName: streamName, Records: records })) + const failed = + result.RequestResponses?.flatMap((item, index) => { + const record = records[index] + if (!item.ErrorCode || !record) return [] + return [record] + }) ?? [] + if (failed.length === 0) return 0 + if (attempt >= MAX_FIREHOSE_ATTEMPTS) return failed.length + + await new Promise((resolve) => setTimeout(resolve, 250 * 2 ** (attempt - 1))) + return putRecords(streamName, failed, attempt + 1) +} + +function parsePayload(event: FunctionUrlEvent): IngestPayload | undefined { + try { + return JSON.parse( + event.isBase64Encoded ? Buffer.from(event.body ?? "", "base64").toString("utf8") : (event.body ?? ""), + ) as IngestPayload + } catch { + return undefined + } +} + +function isAuthorized(headers: Record) { + const actual = Buffer.from(headers.authorization ?? headers.Authorization ?? "") + const expected = Buffer.from(`Bearer ${Resource.StatsLakeIngestConfig.secret}`) + if (actual.length !== expected.length) return false + return timingSafeEqual(actual, expected) +} + +function chunks(items: T[], size: number) { + return Array.from({ length: Math.ceil(items.length / size) }, (_, index) => + items.slice(index * size, (index + 1) * size), + ) +} + +function response(statusCode: number, body: Record) { + return { + statusCode, + headers: { "content-type": "application/json" }, + body: JSON.stringify(body), + } +} diff --git a/packages/stats/core/src/resource.d.ts b/packages/stats/core/src/resource.d.ts index 83ad689adf..9174bee69e 100644 --- a/packages/stats/core/src/resource.d.ts +++ b/packages/stats/core/src/resource.d.ts @@ -2,9 +2,20 @@ import "sst" declare module "sst" { export interface Resource { - HONEYCOMB_API_KEY: { - type: "sst.sst.Secret" - value: string + StatsLake: { + catalog: string + database: string + dataset: string + region: string + table: string + tableBucket: string + type: "sst.sst.Linkable" + workgroup: string + } + StatsLakeIngestConfig: { + secret: string + streamName: string + type: "sst.sst.Linkable" } StatsDatabase: { database: string diff --git a/sst-env.d.ts b/sst-env.d.ts index bbb62c5a58..abaf85da0e 100644 --- a/sst-env.d.ts +++ b/sst-env.d.ts @@ -146,6 +146,35 @@ declare module "sst" { "url": string "username": string } + "StatsLake": { + "catalog": string + "database": string + "dataset": string + "region": string + "table": string + "tableBucket": string + "type": "sst.sst.Linkable" + "workgroup": string + } + "StatsLakeIngest": { + "secret": string + "type": "sst.sst.Linkable" + "url": string + } + "StatsLakeIngestConfig": { + "secret": string + "streamName": string + "type": "sst.sst.Linkable" + } + "StatsLakeIngestFunction": { + "name": string + "type": "sst.aws.Function" + "url": string + } + "StatsLakeIngestSecret": { + "type": "random.index/randomPassword.RandomPassword" + "value": string + } "Teams": { "type": "sst.cloudflare.SolidStart" "url": string diff --git a/sst.config.ts b/sst.config.ts index c07d2eba69..0d4fdb992b 100644 --- a/sst.config.ts +++ b/sst.config.ts @@ -9,6 +9,7 @@ export default $config({ home: "cloudflare", providers: { aws: { + version: "7.30.0", region: "us-east-1", profile: process.env.GITHUB_ACTIONS ? undefined