diff --git a/.dockerignore b/.dockerignore new file mode 100644 index 000000000..b6d8d7351 --- /dev/null +++ b/.dockerignore @@ -0,0 +1,12 @@ +.git +.opencode +.sst +.turbo +.wrangler +node_modules +**/node_modules +**/.output +**/dist +**/.turbo +**/.vite +**/coverage diff --git a/.github/workflows/deploy.yml b/.github/workflows/deploy.yml index 68f00dd4a..18e6cf7ac 100644 --- a/.github/workflows/deploy.yml +++ b/.github/workflows/deploy.yml @@ -9,9 +9,15 @@ on: concurrency: ${{ github.workflow }}-${{ github.ref }} +permissions: + contents: read + id-token: write + jobs: deploy: + if: github.repository == 'anomalyco/opencode' && (github.ref_name == 'dev' || github.ref_name == 'production') runs-on: ubuntu-latest + environment: ${{ github.ref_name }} steps: - uses: actions/checkout@f43a0e5ff2bd294095638e18286ca9a3d1956744 # v3.6.0 @@ -21,6 +27,12 @@ jobs: with: node-version: "24" + - uses: aws-actions/configure-aws-credentials@7474bc4690e29a8392af63c5b98e7449536d5c3a # v4.3.1 + with: + role-to-assume: ${{ vars.AWS_DEPLOY_ROLE_ARN }} + role-session-name: opencode-${{ github.run_id }} + aws-region: us-east-1 + - run: bun sst deploy --stage=${{ github.ref_name }} env: CLOUDFLARE_API_TOKEN: ${{ secrets.CLOUDFLARE_API_TOKEN }} diff --git a/bun.lock b/bun.lock index 514551950..d85e40b08 100644 --- a/bun.lock +++ b/bun.lock @@ -23,7 +23,7 @@ "oxlint-tsgolint": "0.21.0", "prettier": "3.6.2", "semver": "^7.6.0", - "sst": "4.13.1", + "sst": "catalog:", "turbo": "2.8.13", }, }, @@ -102,6 +102,7 @@ "@solidjs/router": "catalog:", "@solidjs/start": "catalog:", "@stripe/stripe-js": "8.6.1", + "@upstash/redis": "1.38.0", "chart.js": "4.5.1", "nitro": "3.0.1-alpha.1", "solid-js": "catalog:", @@ -618,6 +619,66 @@ "typescript": "catalog:", }, }, + "packages/stats/app": { + "name": "@opencode-ai/stats-app", + "version": "1.14.50", + "dependencies": { + "@opencode-ai/stats-core": "workspace:*", + "@opencode-ai/ui": "workspace:*", + "@solidjs/meta": "catalog:", + "@solidjs/router": "catalog:", + "@solidjs/start": "catalog:", + "d3-scale": "4.0.2", + "effect": "catalog:", + "nitro": "3.0.1-alpha.1", + "solid-js": "catalog:", + "vite": "catalog:", + }, + "devDependencies": { + "@cloudflare/workers-types": "catalog:", + "@types/bun": "catalog:", + "@types/d3-scale": "4.0.9", + "@typescript/native-preview": "catalog:", + "typescript": "catalog:", + }, + }, + "packages/stats/core": { + "name": "@opencode-ai/stats-core", + "version": "1.14.50", + "dependencies": { + "@aws-sdk/client-athena": "3.933.0", + "@planetscale/database": "1.19.0", + "drizzle-orm": "catalog:", + "effect": "catalog:", + "sst": "catalog:", + }, + "devDependencies": { + "@tsconfig/node22": "catalog:", + "@types/bun": "catalog:", + "@types/node": "catalog:", + "@typescript/native-preview": "catalog:", + "drizzle-kit": "catalog:", + "typescript": "catalog:", + }, + }, + "packages/stats/server": { + "name": "@opencode-ai/stats-server", + "version": "1.14.50", + "dependencies": { + "@aws-sdk/client-firehose": "3.933.0", + "@effect/platform-node": "catalog:", + "@opencode-ai/stats-core": "workspace:*", + "effect": "catalog:", + "sst": "catalog:", + }, + "devDependencies": { + "@tsconfig/node22": "catalog:", + "@types/bun": "catalog:", + "@types/node": "catalog:", + "@typescript/native-preview": "catalog:", + "typescript": "catalog:", + }, + }, "packages/storybook": { "name": "@opencode-ai/storybook", "devDependencies": { @@ -798,6 +859,7 @@ "shiki": "3.20.0", "solid-js": "1.9.10", "solid-list": "0.3.0", + "sst": "4.13.1", "tailwindcss": "4.1.11", "typescript": "5.8.2", "ulid": "3.0.1", @@ -923,8 +985,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=="], @@ -1605,6 +1671,12 @@ "@opencode-ai/slack": ["@opencode-ai/slack@workspace:packages/slack"], + "@opencode-ai/stats-app": ["@opencode-ai/stats-app@workspace:packages/stats/app"], + + "@opencode-ai/stats-core": ["@opencode-ai/stats-core@workspace:packages/stats/core"], + + "@opencode-ai/stats-server": ["@opencode-ai/stats-server@workspace:packages/stats/server"], + "@opencode-ai/storybook": ["@opencode-ai/storybook@workspace:packages/storybook"], "@opencode-ai/ui": ["@opencode-ai/ui@workspace:packages/ui"], @@ -2343,6 +2415,10 @@ "@types/cross-spawn": ["@types/cross-spawn@6.0.6", "", { "dependencies": { "@types/node": "*" } }, "sha512-fXRhhUkG4H3TQk5dBhQ7m/JDdSNHKwR2BBia62lhwEIq9xGiQKLxd6LymNhn47SjXhsUEPmxi+PKw2OkW4LLjA=="], + "@types/d3-scale": ["@types/d3-scale@4.0.9", "", { "dependencies": { "@types/d3-time": "*" } }, "sha512-dLmtwB8zkAeO/juAMfnV+sItKjlsw2lKdZVVy6LRr0cBmegxSABiLEpGVmSJJ8O08i4+sGR6qQtb6WtuwJdvVw=="], + + "@types/d3-time": ["@types/d3-time@3.0.4", "", {}, "sha512-yuzZug1nkAAaBlBBikKZTgzCeA+k1uy4ZFwWANOfKw5z5LRhV0gNA7gNkKm7HoK+HRN0wX3EkxGk0fpbWhmB7g=="], + "@types/debug": ["@types/debug@4.1.13", "", { "dependencies": { "@types/ms": "*" } }, "sha512-KSVgmQmzMwPlmtljOomayoR89W4FynCAi3E8PPs7vmDVPe84hT+vGPKkJfThkmXs0x0jAaa9U8uW8bbfyS2fWw=="], "@types/deep-eql": ["@types/deep-eql@4.0.2", "", {}, "sha512-c9h9dVVMigMPc4bwTvC5dxqtqJZwQPePsWjPlpSOnojbor6pGqdk541lfA7AqFQr5pB1BRdq0juY9db81BwyFw=="], @@ -2487,6 +2563,8 @@ "@ungap/structured-clone": ["@ungap/structured-clone@1.3.0", "", {}, "sha512-WmoN8qaIAo7WTYWbAZuG8PYEhn5fkz7dZrqTBZ7dtt//lL2Gwms1IcnQ5yHqjDfX8Ft5j4YzDM23f87zBfDe9g=="], + "@upstash/redis": ["@upstash/redis@1.38.0", "", { "dependencies": { "uncrypto": "^0.1.3" } }, "sha512-wu+dZBptlLy0+MCUEoHmzrY/TnmgDey3+c7EbIGwrLqAvkP8yi5MWZHYGIFtAygmL4Bkz2TdFu+eU0vFPncIcg=="], + "@valibot/to-json-schema": ["@valibot/to-json-schema@1.6.0", "", { "peerDependencies": { "valibot": "^1.3.0" } }, "sha512-d6rYyK5KVa2XdqamWgZ4/Nr+cXhxjy7lmpe6Iajw15J/jmU+gyxl2IEd1Otg1d7Rl3gOQL5reulnSypzBtYy1A=="], "@vercel/oidc": ["@vercel/oidc@3.2.0", "", {}, "sha512-UycprH3T6n3jH0k44NHMa7pnFHGu/N05MjojYr+Mc6I7obkoLIJujSWwin1pCvdy/eOxrI/l3uDLQsmcrOb4ug=="], @@ -2913,6 +2991,20 @@ "csstype": ["csstype@3.2.3", "", {}, "sha512-z1HGKcYy2xA8AGQfwrn0PAy+PB7X/GSj3UVJW9qKyn43xWa+gl5nXmU4qqLMRzWVLFC8KusUX8T/0kCiOYpAIQ=="], + "d3-array": ["d3-array@3.2.4", "", { "dependencies": { "internmap": "1 - 2" } }, "sha512-tdQAmyA18i4J7wprpYq8ClcxZy3SC31QMeByyCFyRt7BVHdREQZ5lpzoe5mFEYZUWe+oq8HBvk9JjpibyEV4Jg=="], + + "d3-color": ["d3-color@3.1.0", "", {}, "sha512-zg/chbXyeBtMQ1LbD/WSoW2DpC3I0mpmPdW+ynRTj/x2DAWYrIY7qeZIHidozwV24m4iavr15lNwIwLxRmOxhA=="], + + "d3-format": ["d3-format@3.1.2", "", {}, "sha512-AJDdYOdnyRDV5b6ArilzCPPwc1ejkHcoyFarqlPqT7zRYjhavcT3uSrqcMvsgh2CgoPbK3RCwyHaVyxYcP2Arg=="], + + "d3-interpolate": ["d3-interpolate@3.0.1", "", { "dependencies": { "d3-color": "1 - 3" } }, "sha512-3bYs1rOD33uo8aqJfKP3JWPAibgw8Zm2+L9vBKEHJ2Rg+viTR7o5Mmv5mZcieN+FRYaAOWX5SJATX6k1PWz72g=="], + + "d3-scale": ["d3-scale@4.0.2", "", { "dependencies": { "d3-array": "2.10.0 - 3", "d3-format": "1 - 3", "d3-interpolate": "1.2.0 - 3", "d3-time": "2.1.1 - 3", "d3-time-format": "2 - 4" } }, "sha512-GZW464g1SH7ag3Y7hXjf8RoUuAFIqklOAq3MRl4OaWabTFJY9PN/E1YklhXLh+OQ3fM9yS2nOkCoS+WLZ6kvxQ=="], + + "d3-time": ["d3-time@3.1.0", "", { "dependencies": { "d3-array": "2 - 3" } }, "sha512-VqKjzBLejbSMT4IgbmVgDjpkYrNWUYJnbCGo874u7MMKIWsILRX+OpX/gTk8MqjpT1A/c6HY2dCA77ZN0lkQ2Q=="], + + "d3-time-format": ["d3-time-format@4.1.0", "", { "dependencies": { "d3-time": "1 - 3" } }, "sha512-dJxPBlzC7NugB2PDLwo9Q8JiTR3M3e4/XANkreKSUxF8vvXKqm1Yfq4Q5dl8budlunRVlUUaDUgFt7eA8D6NLg=="], + "data-uri-to-buffer": ["data-uri-to-buffer@4.0.1", "", {}, "sha512-0R9ikRb668HB7QDxT1vkpuUBtqc53YyAwMwGeUFKRojY/NWKvdZ+9UYtRfGmhqNbRkTSVpMbmyhXipFFv2cb/A=="], "data-view-buffer": ["data-view-buffer@1.0.2", "", { "dependencies": { "call-bound": "^1.0.3", "es-errors": "^1.3.0", "is-data-view": "^1.0.2" } }, "sha512-EmKO5V3OLXh1rtK2wgXRansaK1/mtVdTUEiEI0W8RkvgT05kfxaH29PliLnpLP73yYO6142Q72QNa8Wx/A5CqQ=="], @@ -3487,6 +3579,8 @@ "internal-slot": ["internal-slot@1.1.0", "", { "dependencies": { "es-errors": "^1.3.0", "hasown": "^2.0.2", "side-channel": "^1.1.0" } }, "sha512-4gd7VpWNQNB4UKKCFFVcp1AVv+FMOgs9NKzjHKusc8jTMhd5eL1NqQqOpE0KzMds804/yHlglp3uxgluOqAPLw=="], + "internmap": ["internmap@2.0.3", "", {}, "sha512-5Hh7Y1wQbvY5ooGgPbDaL5iYLAPzMTUrjMulskHLH6wnv/A+1q5rgEaiuqEjB+oxGXIVZs1FF+R/KPN3ZSQYYg=="], + "ioredis": ["ioredis@5.10.1", "", { "dependencies": { "@ioredis/commands": "1.5.1", "cluster-key-slot": "^1.1.0", "debug": "^4.3.4", "denque": "^2.1.0", "lodash.defaults": "^4.2.0", "lodash.isarguments": "^3.1.0", "redis-errors": "^1.2.0", "redis-parser": "^3.0.0", "standard-as-callback": "^2.1.0" } }, "sha512-HuEDBTI70aYdx1v6U97SbNx9F1+svQKBDo30o0b9fw055LMepzpOOd0Ccg9Q6tbqmBSJaMuY0fB7yw9/vjBYCA=="], "ip-address": ["ip-address@10.1.0", "", {}, "sha512-XXADHxXmvT9+CRxhXg56LJovE+bmWnEWB78LB83VZTprKTmaC5QfruXocxzTZ2Kl0DNwKuBdlIhjL8LeY8Sf8Q=="], @@ -5195,6 +5289,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=="], @@ -5219,6 +5323,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/app.ts b/infra/app.ts index 2ede5a1f4..7b532bcb2 100644 --- a/infra/app.ts +++ b/infra/app.ts @@ -30,7 +30,7 @@ export const api = new sst.cloudflare.Worker("Api", { transform: { worker: (args) => { args.logpush = true - if ($app.stage === "vimtor") return + if ($app.stage === "vimtor" || $app.stage === "adam") return args.bindings = $resolve(args.bindings).apply((bindings) => [ ...bindings, { diff --git a/infra/console.ts b/infra/console.ts index 29e473de3..0a304a7be 100644 --- a/infra/console.ts +++ b/infra/console.ts @@ -1,7 +1,9 @@ -import { domain } from "./stage" +import { deployAws, domain } from "./stage" import { EMAILOCTOPUS_API_KEY } from "./app" import { SECRET } from "./secret" +const lake = deployAws ? await import("./lake") : undefined + //////////////// // DATABASE //////////////// @@ -240,7 +242,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: [new sst.Secret("HONEYCOMB_API_KEY")], + link: [SECRET.HoneycombApiKey, ...(lake?.lakeIngest ? [lake.lakeIngest] : [])], }) new sst.cloudflare.x.SolidStart("Console", { @@ -250,6 +252,8 @@ new sst.cloudflare.x.SolidStart("Console", { bucket, bucketNew, database, + SECRET.UpstashRedisRestUrl, + SECRET.UpstashRedisRestToken, AUTH_API_URL, STRIPE_WEBHOOK_SECRET, DISCORD_INCIDENT_WEBHOOK_URL, @@ -281,7 +285,7 @@ new sst.cloudflare.x.SolidStart("Console", { }, transform: { server: { - placement: { region: "aws:us-east-1" }, + placement: { region: "aws:us-east-2" }, transform: { worker: { tailConsumers: [{ service: logProcessor.nodes.worker.scriptName }], diff --git a/infra/lake.ts b/infra/lake.ts new file mode 100644 index 000000000..04a3c46ba --- /dev/null +++ b/infra/lake.ts @@ -0,0 +1,322 @@ +import { domain } from "./stage" + +const current = aws.getCallerIdentityOutput({}) +const partition = aws.getPartitionOutput({}) +const region = aws.getRegionOutput({}) + +const tableBucketName = `opencode-${$app.stage}-lake` +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 glueS3TablesDatabaseWildcardArn = $interpolate`arn:${partition.partition}:glue:${region.region}:${current.accountId}:database/${glueCatalogName}/${tableBucketName}/*` +const glueS3TablesTableWildcardArn = $interpolate`arn:${partition.partition}:glue:${region.region}:${current.accountId}:table/${glueCatalogName}/${tableBucketName}/*/*` +const s3TablesBucketWildcardArn = $interpolate`arn:${partition.partition}:s3tables:${region.region}:${current.accountId}:bucket/*` + +export const tableBucket = new aws.s3tables.TableBucket("LakeTableBucket", { + name: tableBucketName, + forceDestroy: $app.stage !== "production", +}) + +const s3TablesCatalog = new aws.cloudcontrol.Resource( + "LakeS3TablesCatalog", + { + 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("LakeAthenaResults", { + bucket: `opencode-${$app.stage}-lake-athena-results`, + forceDestroy: $app.stage !== "production", +}) + +const firehoseErrorBucket = new aws.s3.Bucket("LakeFirehoseErrors", { + bucket: `opencode-${$app.stage}-lake-firehose-errors`, + forceDestroy: $app.stage !== "production", +}) + +const athenaWorkgroup = new aws.athena.Workgroup("LakeAthenaWorkgroup", { + name: `opencode-${$app.stage}-lake-workgroup`, + forceDestroy: $app.stage !== "production", + configuration: { + enforceWorkgroupConfiguration: true, + publishCloudwatchMetricsEnabled: true, + resultConfiguration: { + outputLocation: $interpolate`s3://${athenaResultsBucket.bucket}/`, + }, + }, +}) + +const firehoseRole = new aws.iam.Role("LakeFirehoseRole", { + assumeRolePolicy: aws.iam.getPolicyDocumentOutput({ + statements: [ + { + effect: "Allow", + actions: ["sts:AssumeRole"], + principals: [ + { + type: "Service", + identifiers: ["firehose.amazonaws.com"], + }, + ], + }, + ], + }).json, +}) + +const firehosePolicy = new aws.iam.RolePolicy("LakeFirehosePolicy", { + 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}/*`, + glueS3TablesDatabaseWildcardArn, + glueS3TablesTableWildcardArn, + $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( + "LakeFirehose", + { + name: `opencode-${$app.stage}-lake-ingest`, + destination: "iceberg", + icebergConfiguration: { + appendOnly: true, + bufferingInterval: 60, + bufferingSize: 1, + catalogArn: glueS3TablesChildCatalogArn, + processingConfiguration: { + enabled: true, + processors: [ + { + type: "MetadataExtraction", + parameters: [ + { parameterName: "JsonParsingEngine", parameterValue: "JQ-1.6" }, + { + parameterName: "MetadataExtractionQuery", + parameterValue: + '{destinationDatabaseName:._lake_database,destinationTableName:._lake_table,operation:(._lake_operation // "insert")}', + }, + ], + }, + ], + }, + roleArn: firehoseRole.arn, + s3BackupMode: "FailedDataOnly", + s3Configuration: { + roleArn: firehoseRole.arn, + bucketArn: firehoseErrorBucket.arn, + errorOutputPrefix: "errors/!{firehose:error-output-type}/", + }, + }, + }, + { dependsOn: [s3TablesCatalog, firehosePolicy] }, +) + +export const lakeVpc = new sst.aws.Vpc("LakeVpc") +export const lakeCluster = new sst.aws.Cluster("LakeCluster", { vpc: lakeVpc }) +export const lakeRegion = region.region +export const lakeCatalog = $interpolate`${glueCatalogName}/${tableBucket.name}` +export const lakeAthenaWorkgroup = athenaWorkgroup + +const ingestSecret = new random.RandomPassword("LakeIngestSecret", { length: 32 }) + +const ingestConfig = new sst.Linkable("LakeIngestConfig", { + properties: { + streamName: firehose.name, + secret: ingestSecret.result, + }, +}) + +const ingestService = new sst.aws.Service("LakeIngestService", { + cluster: lakeCluster, + architecture: "arm64", + cpu: "0.5 vCPU", + memory: "1 GB", + image: { + context: ".", + dockerfile: "packages/stats/server/Dockerfile", + }, + link: [ingestConfig], + permissions: [ + { + actions: ["firehose:PutRecord", "firehose:PutRecordBatch"], + resources: [firehose.arn], + }, + ], + scaling: { + min: $app.stage === "production" ? 2 : 1, + max: $app.stage === "production" ? 32 : 4, + cpuUtilization: 60, + memoryUtilization: 70, + }, + loadBalancer: { + domain: { + name: `lake.${domain}`, + dns: sst.cloudflare.dns(), + }, + rules: [ + { listen: "80/http", redirect: "443/https" }, + { listen: "443/https", forward: "3000/http" }, + ], + health: { + "3000/http": { + path: "/ready", + successCodes: "200-299", + }, + }, + }, + health: { + command: [ + "CMD-SHELL", + "bun --eval \"fetch('http://localhost:3000/health').then((r) => process.exit(r.ok ? 0 : 1)).catch(() => process.exit(1))\"", + ], + interval: "30 seconds", + retries: 3, + startPeriod: "30 seconds", + timeout: "5 seconds", + }, + dev: { + command: "bun run start", + directory: "packages/stats/server", + url: "http://localhost:3000", + }, + wait: $app.stage === "production", +}) + +export const lakeIngest = new sst.Linkable("LakeIngest", { + properties: { + url: ingestService.url, + secret: ingestSecret.result, + }, +}) + +export const lakeQueryPermissions = [ + { + 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}/*`, + glueS3TablesDatabaseWildcardArn, + glueS3TablesTableWildcardArn, + $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/infra/secret.ts b/infra/secret.ts index d4e8b148f..65ada2f1f 100644 --- a/infra/secret.ts +++ b/infra/secret.ts @@ -7,5 +7,8 @@ sst.Linkable.wrap(random.RandomPassword, (resource) => ({ export const SECRET = { R2AccessKey: new sst.Secret("R2AccessKey", "unknown"), R2SecretKey: new sst.Secret("R2SecretKey", "unknown"), + HoneycombApiKey: new sst.Secret("HONEYCOMB_API_KEY"), HoneycombWebhookSecret: new random.RandomPassword("HoneycombWebhookSecret", { length: 24 }), + UpstashRedisRestUrl: new sst.Secret("UpstashRedisRestUrl"), + UpstashRedisRestToken: new sst.Secret("UpstashRedisRestToken"), } diff --git a/infra/stage.ts b/infra/stage.ts index f9a6fd755..f98867238 100644 --- a/infra/stage.ts +++ b/infra/stage.ts @@ -5,6 +5,52 @@ export const domain = (() => { })() export const zoneID = "430ba34c138cfb5360826c4909f99be8" +// Dev owns the shared AWS lake/stats infra for all non-production stages. +export const awsStage = $app.stage === "production" ? "production" : "dev" +export const deployAws = $app.stage === awsStage + +const githubActionsDeployRole = (() => { + if ($app.stage !== "dev" && $app.stage !== "production") return + + const provider = new aws.iam.OpenIdConnectProvider("GithubActionsOidcProvider", { + url: "https://token.actions.githubusercontent.com", + clientIdLists: ["sts.amazonaws.com"], + }) + const role = new aws.iam.Role("GithubActionsDeployRole", { + name: `opencode-${$app.stage}-github-actions-deploy`, + maxSessionDuration: 3600, + assumeRolePolicy: aws.iam.getPolicyDocumentOutput({ + statements: [ + { + effect: "Allow", + actions: ["sts:AssumeRoleWithWebIdentity"], + principals: [{ type: "Federated", identifiers: [provider.arn] }], + conditions: [ + { + test: "StringEquals", + variable: "token.actions.githubusercontent.com:aud", + values: ["sts.amazonaws.com"], + }, + { + test: "StringEquals", + variable: "token.actions.githubusercontent.com:sub", + values: [`repo:anomalyco/opencode:environment:${$app.stage}`], + }, + ], + }, + ], + }).json, + }) + + new aws.iam.RolePolicyAttachment("GithubActionsDeployRoleAdmin", { + role: role.name, + policyArn: "arn:aws:iam::aws:policy/AdministratorAccess", + }) + + return role +})() + +export const githubActionsDeployRoleArn = githubActionsDeployRole?.arn new cloudflare.RegionalHostname("RegionalHostname", { hostname: domain, diff --git a/infra/stats.ts b/infra/stats.ts new file mode 100644 index 000000000..28d50fe0e --- /dev/null +++ b/infra/stats.ts @@ -0,0 +1,207 @@ +import { lakeAthenaWorkgroup, lakeCatalog, lakeCluster, lakeQueryPermissions, lakeRegion, tableBucket } from "./lake" + +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` +})() + +//////////////// +// LAKE +//////////////// + +const inferenceNamespace = new aws.s3tables.Namespace("LakeInferenceNamespace", { + namespace: "inference", + tableBucketArn: tableBucket.arn, +}) + +const inferenceEventTable = new aws.s3tables.Table( + "LakeInferenceEventTable", + { + name: "event", + namespace: inferenceNamespace.namespace, + tableBucketArn: inferenceNamespace.tableBucketArn, + format: "ICEBERG", + metadata: { + iceberg: { + schema: { + fields: [ + { 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: "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 }, + { name: "duration", type: "double", required: false }, + { name: "request_length", type: "long", required: false }, + { name: "status", type: "int", required: false }, + { name: "ip", type: "string", required: false }, + { name: "is_stream", type: "boolean", required: false }, + { name: "session", type: "string", required: false }, + { name: "request", type: "string", required: false }, + { name: "client", type: "string", required: false }, + { name: "user_agent", type: "string", required: false }, + { name: "model_variant", type: "string", required: false }, + { name: "source", type: "string", required: false }, + { name: "provider", type: "string", required: false }, + { name: "provider_model", type: "string", required: false }, + { name: "model", type: "string", required: false }, + { name: "llm_error_code", type: "int", required: false }, + { name: "llm_error_message", type: "string", required: false }, + { name: "error_response", type: "string", required: false }, + { name: "error_type", type: "string", required: false }, + { name: "error_message", type: "string", required: false }, + { name: "error_cause", type: "string", required: false }, + { name: "error_cause2", type: "string", required: false }, + { name: "api_key", type: "string", required: false }, + { name: "workspace", type: "string", required: false }, + { name: "is_subscription", type: "boolean", required: false }, + { name: "subscription", type: "string", required: false }, + { name: "response_length", type: "long", required: false }, + { name: "time_to_first_byte", 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: "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: "cost_input", type: "long", required: false }, + { name: "cost_output", type: "long", required: false }, + { name: "cost_cache_read", type: "long", required: false }, + { name: "cost_cache_write_5m", type: "long", required: false }, + { name: "cost_cache_write_1h", type: "long", required: false }, + { name: "cost_total", type: "long", required: false }, + ], + }, + }, + }, + }, + { deleteBeforeReplace: $app.stage !== "production" }, +) + +export const inferenceEvent = new sst.Linkable("InferenceEvent", { + properties: { + region: lakeRegion, + catalog: lakeCatalog, + database: inferenceNamespace.namespace, + table: inferenceEventTable.name, + tableBucket: tableBucket.name, + workgroup: lakeAthenaWorkgroup.name, + }, +}) + +//////////////// +// DATABASE +//////////////// + +const cluster = planetscale.getDatabaseOutput({ + name: "opencode-stats", + organization: "anomalyco", +}) + +const branch = + $app.stage === "production" + ? planetscale.getBranchOutput({ + name: "production", + organization: cluster.organization, + database: cluster.name, + }) + : new planetscale.Branch("StatsDatabaseBranch", { + database: cluster.name, + organization: cluster.organization, + name: $app.stage, + parentBranch: "production", + }) + +const password = new planetscale.Password("StatsDatabasePassword", { + name: $app.stage, + database: cluster.name, + organization: cluster.organization, + branch: branch.name, +}) + +const databaseUrl = $interpolate`mysql://${password.username.apply(encodeURIComponent)}:${password.plaintext.apply( + encodeURIComponent, +)}@${password.accessHostUrl}/${cluster.name}` + +export const database = new sst.Linkable("StatsDatabase", { + properties: { + host: password.accessHostUrl, + database: cluster.name, + username: password.username, + password: password.plaintext, + port: 3306, + url: databaseUrl, + }, +}) + +new sst.x.DevCommand("StatsStudio", { + link: [database], + environment: { + DATABASE_URL: databaseUrl, + }, + dev: { + command: "bun db:studio", + directory: "packages/stats/core", + autostart: false, + }, +}) + +//////////////// +// APP +//////////////// + +// export const app = new sst.cloudflare.x.SolidStart("Stats", { +// path: "packages/stats/app", +// buildCommand: "bun run build", +// domain, +// link: [database], +// environment: { +// PUBLIC_URL: `https://${domain}`, +// }, +// }) + +//////////////// +// SERVICES +//////////////// + +const statsSyncConfig = new sst.Linkable("StatsSyncConfig", { + properties: { + dataset: "zen", + }, +}) + +export const statSync = new sst.aws.Service("StatsSyncService", { + cluster: lakeCluster, + architecture: "arm64", + cpu: "0.25 vCPU", + memory: "0.5 GB", + image: { + context: ".", + dockerfile: "packages/stats/server/Dockerfile", + }, + command: ["bun", "src/stat-sync.ts"], + link: [database, inferenceEvent, statsSyncConfig], + permissions: lakeQueryPermissions, + scaling: { + min: 1, + max: 1, + }, + dev: { + command: "bun src/stat-sync.ts", + directory: "packages/stats/server", + autostart: false, + }, +}) diff --git a/nix/hashes.json b/nix/hashes.json index 75a54eac6..a4353299a 100644 --- a/nix/hashes.json +++ b/nix/hashes.json @@ -1,8 +1,8 @@ { "nodeModules": { - "x86_64-linux": "sha256-Nol27sqG6za6oR32MqZrmRxch0/XFeQ61C9Yqtxo7Hk=", - "aarch64-linux": "sha256-+jfkGr4tuv/s/v890XaFTRqUxwA8H8KVMiys6sEXi/4=", - "aarch64-darwin": "sha256-FYtKluNazzuQcrr2iHm7B2pJfYrAHkbcj4kWs2FtBgc=", - "x86_64-darwin": "sha256-PpyWDkzb0p50US7Bs3m3A4PZ0JM3dccGgTKYN9JFYR0=" + "x86_64-linux": "sha256-6s5msV+dYMHHt9Gc/CqvCrUj8K7ELjxoAe6ejHSJo4I=", + "aarch64-linux": "sha256-SP94UPy5LePd7ZdC3eENIXiozc67blpg1SN9Ug5yiv8=", + "aarch64-darwin": "sha256-jv3lffiAQ5kxDAbXNavsHD+tjjdgT0dT0JxN0bWMYTE=", + "x86_64-darwin": "sha256-aYNtcarg516ZmqaO62mnUZYiSyWt8rJjUHQslhrhGHM=" } } diff --git a/package.json b/package.json index 86f5a1eac..fcf3873d3 100644 --- a/package.json +++ b/package.json @@ -10,6 +10,7 @@ "dev:desktop": "bun --cwd packages/desktop dev", "dev:web": "bun --cwd packages/app dev", "dev:console": "ulimit -n 10240 2>/dev/null; bun run --cwd packages/console/app dev", + "dev:stats": "bun run --cwd packages/stats/app dev", "dev:storybook": "bun --cwd packages/storybook storybook", "lint": "oxlint", "typecheck": "bun turbo typecheck", @@ -17,13 +18,14 @@ "postinstall": "bun run --cwd packages/opencode fix-node-pty", "prepare": "husky", "random": "echo 'Random script'", - "hello": "echo 'Hello World!'", + "sso": "aws sso login --sso-session=opencode --no-browser", "test": "echo 'do not run tests from root' && exit 1" }, "workspaces": { "packages": [ "packages/*", "packages/console/*", + "packages/stats/*", "packages/sdk/js", "packages/slack" ], @@ -72,6 +74,7 @@ "@typescript/native-preview": "7.0.0-dev.20251207.1", "zod": "4.1.8", "remeda": "2.26.0", + "sst": "4.13.1", "shiki": "3.20.0", "solid-list": "0.3.0", "tailwindcss": "4.1.11", @@ -98,7 +101,7 @@ "oxlint-tsgolint": "0.21.0", "prettier": "3.6.2", "semver": "^7.6.0", - "sst": "4.13.1", + "sst": "catalog:", "turbo": "2.8.13" }, "dependencies": { diff --git a/packages/console/app/package.json b/packages/console/app/package.json index 86e7da253..b93c15cdc 100644 --- a/packages/console/app/package.json +++ b/packages/console/app/package.json @@ -26,6 +26,7 @@ "@solidjs/router": "catalog:", "@solidjs/start": "catalog:", "@stripe/stripe-js": "8.6.1", + "@upstash/redis": "1.38.0", "chart.js": "4.5.1", "nitro": "3.0.1-alpha.1", "solid-js": "catalog:", diff --git a/packages/console/app/src/routes/zen/util/handler.ts b/packages/console/app/src/routes/zen/util/handler.ts index e4b42d741..0434f9100 100644 --- a/packages/console/app/src/routes/zen/util/handler.ts +++ b/packages/console/app/src/routes/zen/util/handler.ts @@ -249,8 +249,9 @@ export async function handler( if (!isStream || [400, 404, 429].includes(res.status)) { const json = await res.json() await rateLimiter?.track() - if (json.usage) { - const usageInfo = providerInfo.normalizeUsage(json.usage) + const usage = providerInfo.extractUsage(json) + if (usage) { + const usageInfo = providerInfo.normalizeUsage(usage) const costInfo = calculateCost(modelInfo, usageInfo) await trialLimiter?.track(usageInfo) await modelTpmLimiter?.track(providerInfo.id, providerInfo.model, usageInfo) diff --git a/packages/console/app/src/routes/zen/util/ipRateLimiter.ts b/packages/console/app/src/routes/zen/util/ipRateLimiter.ts index 177d78514..81f73a4e5 100644 --- a/packages/console/app/src/routes/zen/util/ipRateLimiter.ts +++ b/packages/console/app/src/routes/zen/util/ipRateLimiter.ts @@ -2,6 +2,7 @@ import { Database, eq, and, sql, inArray } from "@opencode-ai/console-core/drizz import { IpRateLimitTable } from "@opencode-ai/console-core/schema/ip.sql.js" import { FreeUsageLimitError } from "./error" import { logger } from "./logger" +import { buildRateLimitKey, getRedis } from "./redis" import { i18n } from "~/i18n" import { localeFromRequest } from "~/lib/language" import { Subscription } from "@opencode-ai/console-core/subscription.js" @@ -23,43 +24,62 @@ export function createRateLimiter(modelId: string, rateLimit: number | undefined const now = Date.now() const lifetimeInterval = "" const dailyInterval = rateLimit ? `${buildYYYYMMDD(now)}${modelId.substring(0, 2)}` : buildYYYYMMDD(now) - - let _isNew: boolean + const retryAfter = getRetryAfterDay(now) + const redis = getRedis() + const lifetimeKey = buildRateLimitKey("ip", ip) + const dailyKey = buildRateLimitKey("ip", ip, dailyInterval) + let isNew = false return { check: async () => { - const rows = await Database.use((tx) => - tx - .select({ interval: IpRateLimitTable.interval, count: IpRateLimitTable.count }) - .from(IpRateLimitTable) - .where( - and( - eq(IpRateLimitTable.ip, ip), - isDefaultModel - ? inArray(IpRateLimitTable.interval, [lifetimeInterval, dailyInterval]) - : inArray(IpRateLimitTable.interval, [dailyInterval]), + const [counts, rows] = await Promise.all([ + redis.mget<(string | number | null)[]>(isDefaultModel ? [lifetimeKey, dailyKey] : [dailyKey]).catch(() => []), + Database.use((tx) => + tx + .select({ interval: IpRateLimitTable.interval, count: IpRateLimitTable.count }) + .from(IpRateLimitTable) + .where( + and( + eq(IpRateLimitTable.ip, ip), + isDefaultModel + ? inArray(IpRateLimitTable.interval, [lifetimeInterval, dailyInterval]) + : inArray(IpRateLimitTable.interval, [dailyInterval]), + ), ), - ), - ) - const lifetimeCount = rows.find((r) => r.interval === lifetimeInterval)?.count ?? 0 - const dailyCount = rows.find((r) => r.interval === dailyInterval)?.count ?? 0 + ), + ]) + const redisLifetimeCount = isDefaultModel ? Number(counts[0] ?? 0) : 0 + const redisDailyCount = Number(counts[isDefaultModel ? 1 : 0] ?? 0) + const databaseLifetimeCount = rows.find((r) => r.interval === lifetimeInterval)?.count ?? 0 + const databaseDailyCount = rows.find((r) => r.interval === dailyInterval)?.count ?? 0 + const lifetimeCount = Math.max(redisLifetimeCount, databaseLifetimeCount) + const dailyCount = Math.max(redisDailyCount, databaseDailyCount) logger.debug(`rate limit lifetime: ${lifetimeCount}, daily: ${dailyCount}`) - _isNew = isDefaultModel && lifetimeCount < dailyLimit * 7 + isNew = isDefaultModel && lifetimeCount < dailyLimit * 7 + if (isDefaultModel && databaseLifetimeCount > redisLifetimeCount) + await redis.set(lifetimeKey, databaseLifetimeCount).catch(() => {}) - if ((_isNew && dailyCount >= dailyLimit * 2) || (!_isNew && dailyCount >= dailyLimit)) - throw new FreeUsageLimitError(dict["zen.api.error.rateLimitExceeded"], getRetryAfterDay(now)) + if ((isNew && dailyCount >= dailyLimit * 2) || (!isNew && dailyCount >= dailyLimit)) + throw new FreeUsageLimitError(dict["zen.api.error.rateLimitExceeded"], retryAfter) }, track: async () => { - await Database.use((tx) => - tx - .insert(IpRateLimitTable) - .values([ - { ip, interval: dailyInterval, count: 1 }, - ...(_isNew ? [{ ip, interval: lifetimeInterval, count: 1 }] : []), - ]) - .onDuplicateKeyUpdate({ set: { count: sql`${IpRateLimitTable.count} + 1` } }), - ) + const pipeline = redis.pipeline() + pipeline.incr(dailyKey) + pipeline.expire(dailyKey, retryAfter) + if (isNew) pipeline.incr(lifetimeKey) + await Promise.all([ + pipeline.exec().catch(() => {}), + Database.use((tx) => + tx + .insert(IpRateLimitTable) + .values([ + { ip, interval: dailyInterval, count: 1 }, + ...(isNew ? [{ ip, interval: lifetimeInterval, count: 1 }] : []), + ]) + .onDuplicateKeyUpdate({ set: { count: sql`${IpRateLimitTable.count} + 1` } }), + ), + ]) }, } } diff --git a/packages/console/app/src/routes/zen/util/provider/anthropic.ts b/packages/console/app/src/routes/zen/util/provider/anthropic.ts index 8c394ee3e..64053fd73 100644 --- a/packages/console/app/src/routes/zen/util/provider/anthropic.ts +++ b/packages/console/app/src/routes/zen/util/provider/anthropic.ts @@ -175,6 +175,7 @@ export const anthropicHelper: ProviderHelper = ({ reqModel, providerModel }) => retrieve: () => usage, } }, + extractUsage: (response: any) => response.usage, normalizeUsage: (usage: Usage) => ({ inputTokens: usage.input_tokens ?? 0, outputTokens: usage.output_tokens ?? 0, diff --git a/packages/console/app/src/routes/zen/util/provider/google.ts b/packages/console/app/src/routes/zen/util/provider/google.ts index 2954024e2..eead927c8 100644 --- a/packages/console/app/src/routes/zen/util/provider/google.ts +++ b/packages/console/app/src/routes/zen/util/provider/google.ts @@ -58,6 +58,7 @@ export const googleHelper: ProviderHelper = ({ providerModel }) => ({ retrieve: () => usage, } }, + extractUsage: (response: any) => response.usageMetadata, normalizeUsage: (usage: Usage) => { const inputTokens = usage.promptTokenCount ?? 0 const outputTokens = usage.candidatesTokenCount ?? 0 diff --git a/packages/console/app/src/routes/zen/util/provider/openai-compatible.ts b/packages/console/app/src/routes/zen/util/provider/openai-compatible.ts index 912c89092..9e5e15d94 100644 --- a/packages/console/app/src/routes/zen/util/provider/openai-compatible.ts +++ b/packages/console/app/src/routes/zen/util/provider/openai-compatible.ts @@ -58,6 +58,7 @@ export const oaCompatHelper: ProviderHelper = ({ adjustCacheUsage }) => ({ retrieve: () => usage, } }, + extractUsage: (response: any) => response.usage, normalizeUsage: (usage: Usage) => { let inputTokens = usage.prompt_tokens ?? 0 const outputTokens = usage.completion_tokens ?? 0 diff --git a/packages/console/app/src/routes/zen/util/provider/openai.ts b/packages/console/app/src/routes/zen/util/provider/openai.ts index 4b39407d4..e55dcd878 100644 --- a/packages/console/app/src/routes/zen/util/provider/openai.ts +++ b/packages/console/app/src/routes/zen/util/provider/openai.ts @@ -43,6 +43,7 @@ export const openaiHelper: ProviderHelper = ({ workspaceID }) => ({ retrieve: () => usage, } }, + extractUsage: (response: any) => response.usage ?? response.response?.usage, normalizeUsage: (usage: Usage) => { const inputTokens = usage.input_tokens ?? 0 const outputTokens = usage.output_tokens ?? 0 diff --git a/packages/console/app/src/routes/zen/util/provider/provider.ts b/packages/console/app/src/routes/zen/util/provider/provider.ts index 319f8fdca..d9fe55681 100644 --- a/packages/console/app/src/routes/zen/util/provider/provider.ts +++ b/packages/console/app/src/routes/zen/util/provider/provider.ts @@ -49,6 +49,7 @@ export type ProviderHelper = (input: { parse: (chunk: string) => void retrieve: () => any } + extractUsage: (response: any) => any normalizeUsage: (usage: any) => UsageInfo } diff --git a/packages/console/app/src/routes/zen/util/redis.ts b/packages/console/app/src/routes/zen/util/redis.ts new file mode 100644 index 000000000..512523298 --- /dev/null +++ b/packages/console/app/src/routes/zen/util/redis.ts @@ -0,0 +1,18 @@ +import { Resource } from "@opencode-ai/console-resource" +import { Redis } from "@upstash/redis/cloudflare" + +let redis: Redis | undefined + +export function getRedis() { + if (redis) return redis + redis = new Redis({ + url: Resource.UpstashRedisRestUrl.value, + token: Resource.UpstashRedisRestToken.value, + enableTelemetry: false, + }) + return redis +} + +export function buildRateLimitKey(kind: string, identifier: string, interval?: string) { + return `${Resource.App.stage}:ratelimit:${kind}:${identifier}${interval ? `:${interval}` : ""}` +} diff --git a/packages/console/app/test/providerUsage.test.ts b/packages/console/app/test/providerUsage.test.ts new file mode 100644 index 000000000..d39c9fa86 --- /dev/null +++ b/packages/console/app/test/providerUsage.test.ts @@ -0,0 +1,68 @@ +import { describe, expect, test } from "bun:test" +import type { ZenData } from "@opencode-ai/console-core/model.js" +import type { ProviderHelper } from "../src/routes/zen/util/provider/provider" +import { anthropicHelper } from "../src/routes/zen/util/provider/anthropic" +import { googleHelper } from "../src/routes/zen/util/provider/google" +import { oaCompatHelper } from "../src/routes/zen/util/provider/openai-compatible" +import { openaiHelper } from "../src/routes/zen/util/provider/openai" + +const providers = { + anthropic: anthropicHelper({ reqModel: "claude-haiku-4-5", providerModel: "claude-haiku-4-5" }), + google: googleHelper({ reqModel: "gemini-3-flash", providerModel: "gemini-3-flash" }), + openai: openaiHelper({ reqModel: "gpt-5", providerModel: "gpt-5" }), + "oa-compat": oaCompatHelper({ reqModel: "gpt-5-nano", providerModel: "gpt-5-nano" }), +} satisfies Record> + +describe("provider usage extraction", () => { + test("extracts Google non-stream usage metadata", () => { + const usage = providers.google.extractUsage({ + usageMetadata: { + promptTokenCount: 10, + candidatesTokenCount: 3, + thoughtsTokenCount: 2, + cachedContentTokenCount: 4, + }, + }) + + expect(providers.google.normalizeUsage(usage)).toEqual({ + inputTokens: 6, + outputTokens: 3, + reasoningTokens: 2, + cacheReadTokens: 4, + cacheWrite5mTokens: undefined, + cacheWrite1hTokens: undefined, + }) + }) + + test("parses Google stream usage metadata", () => { + const usageParser = providers.google.createUsageParser() + usageParser.parse( + 'data: {"usageMetadata":{"promptTokenCount":10,"candidatesTokenCount":3,"thoughtsTokenCount":2,"cachedContentTokenCount":4}}', + ) + + expect(providers.google.normalizeUsage(usageParser.retrieve())).toEqual({ + inputTokens: 6, + outputTokens: 3, + reasoningTokens: 2, + cacheReadTokens: 4, + cacheWrite5mTokens: undefined, + cacheWrite1hTokens: undefined, + }) + }) + + test("extracts nested OpenAI Responses usage", () => { + expect( + providers.openai.extractUsage({ + response: { + usage: { + input_tokens: 5, + output_tokens: 7, + }, + }, + }), + ).toEqual({ + input_tokens: 5, + output_tokens: 7, + }) + }) +}) diff --git a/packages/console/app/vite.config.ts b/packages/console/app/vite.config.ts index 951c9a427..fb753b6be 100644 --- a/packages/console/app/vite.config.ts +++ b/packages/console/app/vite.config.ts @@ -9,7 +9,7 @@ export default defineConfig({ }) as PluginOption, nitro({ compatibilityDate: "2024-09-19", - preset: "cloudflare_module", + preset: "cloudflare-module", cloudflare: { nodeCompat: true, }, 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 000000000..dba2ee946 --- /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 2bb741b7a..25e1838ed 100644 --- a/packages/console/function/src/log-processor.ts +++ b/packages/console/function/src/log-processor.ts @@ -21,7 +21,7 @@ export default { ) continue - let data = { + let data: Record = { "cf.continent": event.event.request.cf?.continent, "cf.country": event.event.request.cf?.country, "cf.city": event.event.request.cf?.city, @@ -35,30 +35,152 @@ export default { ip: event.event.request.headers["x-real-ip"], } const time = new Date(event.eventTimestamp ?? Date.now()).toISOString() - const events = [] - for (const log of event.logs) { - for (const message of log.message) { - if (!message.startsWith("_metric:")) continue - const json = JSON.parse(message.slice(8)) - data = { ...data, ...json } - if ("llm.error.code" in json) { - events.push({ time, data: { ...data, event_type: "llm.error" } }) - } - } - } - events.push({ time, data: { ...data, event_type: "completions" } }) + const events = [ + ...event.logs.flatMap((log) => + log.message.flatMap((message: string) => { + if (!message.startsWith("_metric:")) return [] + const json = JSON.parse(message.slice(8)) as Record + data = { ...data, ...json } + if ("llm.error.code" in json) { + return [{ time, data: { ...data, event_type: "llm.error" } }] + } + return [] + }), + ), + { 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 lakeIngest = getLakeIngest() + 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), + }), + ...(lakeIngest + ? [ + fetch(lakeIngest.url, { + method: "POST", + headers: { + "Content-Type": "application/json", + Authorization: `Bearer ${lakeIngest.secret}`, + }, + body: JSON.stringify({ events: events.map((event) => toLakeEvent(event.time, event.data)) }), + }), + ] + : []), + ]) + console.log(honeycomb.status) + console.log(await honeycomb.text()) + if (lake) { + console.log(lake.status) + console.log(await lake.text()) + } } }, } + +function getLakeIngest(): { url: string; secret: string } | undefined { + try { + return Resource.LakeIngest + } catch { + return undefined + } +} + +function toLakeEvent(time: string, data: Record) { + return { + _datalake_key: "inference.event", + event_timestamp: time, + event_date: time.slice(0, 10), + event_type: string(data, "event_type"), + dataset: "zen", + 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"), + duration: number(data, "duration"), + request_length: integer(data, "request_length"), + status: integer(data, "status"), + ip: string(data, "ip"), + is_stream: boolean(data, "is_stream"), + session: string(data, "session"), + request: string(data, "request"), + client: string(data, "client"), + user_agent: string(data, "user_agent"), + model_variant: string(data, "model.variant"), + source: string(data, "source"), + provider: string(data, "provider"), + provider_model: string(data, "provider.model"), + model: string(data, "model"), + llm_error_code: integer(data, "llm.error.code"), + llm_error_message: string(data, "llm.error.message"), + error_response: string(data, "error.response"), + error_type: string(data, "error.type"), + error_message: string(data, "error.message"), + error_cause: string(data, "error.cause"), + error_cause2: string(data, "error.cause2"), + api_key: string(data, "api_key"), + workspace: string(data, "workspace"), + is_subscription: boolean(data, "isSubscription"), + subscription: string(data, "subscription"), + response_length: integer(data, "response_length"), + time_to_first_byte: integer(data, "time_to_first_byte"), + timestamp_first_byte: integer(data, "timestamp.first_byte"), + timestamp_last_byte: integer(data, "timestamp.last_byte"), + tokens_input: integer(data, "tokens.input"), + tokens_output: integer(data, "tokens.output"), + tokens_reasoning: integer(data, "tokens.reasoning"), + tokens_cache_read: integer(data, "tokens.cache_read"), + tokens_cache_write_5m: integer(data, "tokens.cache_write_5m"), + tokens_cache_write_1h: integer(data, "tokens.cache_write_1h"), + 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"), + cost_input: integer(data, "cost.input"), + cost_output: integer(data, "cost.output"), + cost_cache_read: integer(data, "cost.cache_read"), + cost_cache_write_5m: integer(data, "cost.cache_write_5m"), + cost_cache_write_1h: integer(data, "cost.cache_write_1h"), + cost_total: integer(data, "cost.total"), + } +} + +function string(data: Record, 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: Record, 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: Record, key: string) { + const value = number(data, key) + if (value === undefined) return undefined + return Math.round(value) +} + +function number(data: Record, 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.node.ts b/packages/console/resource/resource.node.ts index 1470bacf2..ce11abcc4 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/core/src/flag/flag.ts b/packages/core/src/flag/flag.ts index 54f2445e0..504e156ac 100644 --- a/packages/core/src/flag/flag.ts +++ b/packages/core/src/flag/flag.ts @@ -8,6 +8,10 @@ function truthy(key: string) { const OPENCODE_EXPERIMENTAL = truthy("OPENCODE_EXPERIMENTAL") const copy = process.env["OPENCODE_EXPERIMENTAL_DISABLE_COPY_ON_SELECT"] +function enabledByExperimental(key: string) { + return process.env[key] === undefined ? OPENCODE_EXPERIMENTAL : truthy(key) +} + export const Flag = { OTEL_EXPORTER_OTLP_ENDPOINT: process.env["OTEL_EXPORTER_OTLP_ENDPOINT"], OTEL_EXPORTER_OTLP_HEADERS: process.env["OTEL_EXPORTER_OTLP_HEADERS"], @@ -42,7 +46,7 @@ export const Flag = { OPENCODE_DB: process.env["OPENCODE_DB"], OPENCODE_WORKSPACE_ID: process.env["OPENCODE_WORKSPACE_ID"], - OPENCODE_EXPERIMENTAL_WORKSPACES: OPENCODE_EXPERIMENTAL || truthy("OPENCODE_EXPERIMENTAL_WORKSPACES"), + OPENCODE_EXPERIMENTAL_WORKSPACES: enabledByExperimental("OPENCODE_EXPERIMENTAL_WORKSPACES"), // Evaluated at access time (not module load) because tests, the CLI, and // external tooling set these env vars at runtime. diff --git a/packages/enterprise/vite.config.ts b/packages/enterprise/vite.config.ts index 11ca1729d..531732c2b 100644 --- a/packages/enterprise/vite.config.ts +++ b/packages/enterprise/vite.config.ts @@ -8,7 +8,7 @@ const nitroConfig: any = (() => { if (target === "cloudflare") { return { compatibilityDate: "2024-09-19", - preset: "cloudflare_module", + preset: "cloudflare-module", cloudflare: { nodeCompat: true, }, diff --git a/packages/opencode/src/acp-next/agent.ts b/packages/opencode/src/acp-next/agent.ts new file mode 100644 index 000000000..a4ae00695 --- /dev/null +++ b/packages/opencode/src/acp-next/agent.ts @@ -0,0 +1,75 @@ +import { + RequestError, + type Agent as ACPAgent, + type AgentSideConnection, + type AuthenticateRequest, + type CancelNotification, + type InitializeRequest, + type LoadSessionRequest, + type NewSessionRequest, + type PromptRequest, + type SetSessionConfigOptionRequest, + type SetSessionModelRequest, + type SetSessionModeRequest, +} from "@agentclientprotocol/sdk" +import { Effect } from "effect" +import type { OpencodeClient } from "@opencode-ai/sdk/v2" +import * as ACPNextError from "./error" +import * as ACPNextService from "./service" + +export function init({ sdk: _sdk }: { sdk: OpencodeClient }) { + return { + create: (connection: AgentSideConnection) => { + return new Agent(ACPNextService.make({ sdk: _sdk, connection })) + }, + } +} + +export class Agent implements ACPAgent { + constructor(private readonly service: ACPNextService.Interface) {} + + initialize(params: InitializeRequest) { + return run(this.service.initialize(params)) + } + + authenticate(params: AuthenticateRequest) { + return run(this.service.authenticate(params)) + } + + newSession(params: NewSessionRequest) { + return run(this.service.newSession(params)) + } + + loadSession(params: LoadSessionRequest) { + return run(this.service.loadSession(params)) + } + + setSessionConfigOption(params: SetSessionConfigOptionRequest) { + return run(this.service.setSessionConfigOption(params)) + } + + setSessionMode(params: SetSessionModeRequest) { + return run(this.service.setSessionMode(params)) + } + + unstable_setSessionModel(params: SetSessionModelRequest) { + return run(this.service.setSessionModel(params)) + } + + prompt(params: PromptRequest) { + return run(this.service.prompt(params)) + } + + cancel(params: CancelNotification) { + return run(this.service.cancel(params)) + } +} + +function run(effect: Effect.Effect) { + return Effect.runPromise(effect.pipe(Effect.mapError(ACPNextError.toRequestError))).catch((defect: unknown) => { + if (defect instanceof RequestError) throw defect + throw ACPNextError.toRequestError(ACPNextError.fromUnknownDefect(defect)) + }) +} + +export * as ACPNext from "./agent" diff --git a/packages/opencode/src/acp-next/config-option.ts b/packages/opencode/src/acp-next/config-option.ts new file mode 100644 index 000000000..b730ae075 --- /dev/null +++ b/packages/opencode/src/acp-next/config-option.ts @@ -0,0 +1,203 @@ +import type { SessionConfigOption } from "@agentclientprotocol/sdk" + +export const DEFAULT_VARIANT_VALUE = "default" + +export type ConfigOptionModel = { + id: string + name: string + variants?: Record> +} + +export type ConfigOptionProvider = { + id: string + name: string + models: Record +} + +export type ConfigOptionMode = { + id: string + name: string + description?: string +} + +export type ModelSelection = { + model: { + providerID: string + modelID: string + } + variant?: string +} + +export function buildModelSelectOption(input: { + providers: readonly ConfigOptionProvider[] + currentModel: ModelSelection["model"] + currentVariant?: string + includeVariants?: boolean +}): SessionConfigOption { + return { + id: "model", + name: "Model", + category: "model", + type: "select", + currentValue: formatCurrentModelId({ + model: input.currentModel, + variant: input.currentVariant, + variants: variantsForModel(input.providers, input.currentModel), + includeVariant: input.includeVariants ?? false, + }), + options: buildModelSelectOptions(input.providers, { includeVariants: input.includeVariants ?? false }), + } +} + +export function buildEffortSelectOption(input: { + variants: readonly string[] + currentVariant?: string +}): SessionConfigOption | undefined { + if (input.variants.length === 0) return undefined + + return { + id: "effort", + name: "Effort", + description: "Available effort levels for this model", + category: "thought_level", + type: "select", + currentValue: selectVariant(input.currentVariant, input.variants), + options: input.variants.map((variant) => ({ + value: variant, + name: formatVariantName(variant), + })), + } +} + +export function buildModeSelectOption(input: { + modes: readonly ConfigOptionMode[] + currentModeId: string +}): SessionConfigOption { + return { + id: "mode", + name: "Session Mode", + category: "mode", + type: "select", + currentValue: input.currentModeId, + options: input.modes.map((mode) => ({ + value: mode.id, + name: mode.name, + ...(mode.description ? { description: mode.description } : {}), + })), + } +} + +export function buildConfigOptions(input: { + providers: readonly ConfigOptionProvider[] + currentModel: ModelSelection["model"] + currentVariant?: string + includeModelVariants?: boolean + modes?: readonly ConfigOptionMode[] + currentModeId?: string +}): SessionConfigOption[] { + const variants = variantsForModel(input.providers, input.currentModel) + const effort = buildEffortSelectOption({ variants, currentVariant: input.currentVariant }) + + return [ + buildModelSelectOption({ + providers: input.providers, + currentModel: input.currentModel, + currentVariant: input.currentVariant, + includeVariants: input.includeModelVariants ?? false, + }), + ...(effort ? [effort] : []), + ...(input.modes && input.currentModeId + ? [buildModeSelectOption({ modes: input.modes, currentModeId: input.currentModeId })] + : []), + ] +} + +export function parseModelSelection(modelId: string, providers: readonly ConfigOptionProvider[]): ModelSelection { + const provider = providers.find((item) => modelId.startsWith(`${item.id}/`)) + if (provider) { + const modelID = modelId.slice(provider.id.length + 1) + if (provider.models[modelID]) { + return { model: { providerID: provider.id, modelID } } + } + + const separator = modelID.lastIndexOf("/") + if (separator > -1) { + const baseModelID = modelID.slice(0, separator) + const variant = modelID.slice(separator + 1) + if (provider.models[baseModelID]?.variants?.[variant]) { + return { model: { providerID: provider.id, modelID: baseModelID }, variant } + } + } + + return { model: { providerID: provider.id, modelID } } + } + + const separator = modelId.indexOf("/") + if (separator === -1) { + return { model: { providerID: modelId, modelID: "" } } + } + + return { + model: { + providerID: modelId.slice(0, separator), + modelID: modelId.slice(separator + 1), + }, + } +} + +export function formatCurrentModelId(input: { + model: ModelSelection["model"] + variant?: string + variants?: readonly string[] + includeVariant?: boolean +}) { + const base = `${input.model.providerID}/${input.model.modelID}` + if (!input.includeVariant || !input.variants?.length) return base + return `${base}/${selectVariant(input.variant, input.variants)}` +} + +export function formatVariantName(variant: string) { + return variant + .split(/[_-]/) + .map((part) => (part ? part.charAt(0).toUpperCase() + part.slice(1) : part)) + .join(" ") +} + +function buildModelSelectOptions( + providers: readonly ConfigOptionProvider[], + options: { includeVariants: boolean }, +): Array<{ value: string; name: string }> { + return providers.flatMap((provider) => + Object.values(provider.models) + .sort((a, b) => a.name.localeCompare(b.name)) + .flatMap((model) => { + const base = { + value: `${provider.id}/${model.id}`, + name: `${provider.name}/${model.name}`, + } + if (!options.includeVariants || !model.variants) return [base] + + return [ + base, + ...Object.keys(model.variants) + .filter((variant) => variant !== DEFAULT_VARIANT_VALUE) + .map((variant) => ({ + value: `${provider.id}/${model.id}/${variant}`, + name: `${provider.name}/${model.name} (${formatVariantName(variant)})`, + })), + ] + }), + ) +} + +function variantsForModel(providers: readonly ConfigOptionProvider[], model: ModelSelection["model"]) { + return Object.keys( + providers.find((provider) => provider.id === model.providerID)?.models[model.modelID]?.variants ?? {}, + ) +} + +function selectVariant(variant: string | undefined, variants: readonly string[]) { + if (variant && variants.includes(variant)) return variant + if (variants.includes(DEFAULT_VARIANT_VALUE)) return DEFAULT_VARIANT_VALUE + return variants[0] +} diff --git a/packages/opencode/src/acp-next/content.ts b/packages/opencode/src/acp-next/content.ts new file mode 100644 index 000000000..32630a620 --- /dev/null +++ b/packages/opencode/src/acp-next/content.ts @@ -0,0 +1,250 @@ +import type { ContentBlock, ContentChunk, ResourceLink, Role } from "@agentclientprotocol/sdk" +import path from "node:path" +import { pathToFileURL } from "node:url" +import { SessionLegacy } from "@opencode-ai/core/session/legacy" + +export type PromptPart = SessionLegacy.TextPartInput | SessionLegacy.FilePartInput + +export type ReplayPart = + | { + type: "text" + text: string + synthetic?: boolean + ignored?: boolean + } + | { + type: "file" + url: string + mime: string + filename?: string + } + | { + type: "reasoning" + text: string + } + +export function promptContentToParts(content: readonly ContentBlock[]): PromptPart[] { + return content.flatMap(contentBlockToParts) +} + +export function contentBlockToParts(block: ContentBlock): PromptPart[] { + switch (block.type) { + case "text": + return [ + { + type: "text", + text: block.text, + ...audienceFlags(block.annotations?.audience ?? undefined), + }, + ] + + case "image": + if (block.data) { + return [ + { + type: "file", + url: `data:${block.mimeType};base64,${block.data}`, + filename: filenameFromUri(block.uri ?? undefined) ?? "image", + mime: block.mimeType, + }, + ] + } + if (block.uri?.startsWith("data:")) { + return [ + { + type: "file", + url: block.uri, + filename: filenameFromUri(block.uri) ?? "image", + mime: block.mimeType, + }, + ] + } + if (block.uri?.startsWith("http://") || block.uri?.startsWith("https://")) { + return [ + { + type: "file", + url: block.uri, + filename: filenameFromUri(block.uri) ?? "image", + mime: block.mimeType, + }, + ] + } + return [] + + case "resource_link": + return [resourceLinkToPart(block)] + + case "resource": + if ("text" in block.resource) { + return [{ type: "text", text: block.resource.text }] + } + if (block.resource.mimeType) { + return [ + { + type: "file", + url: block.resource.uri.startsWith("data:") + ? block.resource.uri + : `data:${block.resource.mimeType};base64,${block.resource.blob}`, + filename: filenameFromUri(block.resource.uri) ?? "file", + mime: block.resource.mimeType, + }, + ] + } + return [] + + default: + return [] + } +} + +export function partsToContentChunks(parts: readonly ReplayPart[]): ContentChunk[] { + return parts.flatMap(partToContentChunks) +} + +export function partToContentChunks(part: ReplayPart): ContentChunk[] { + switch (part.type) { + case "text": + if (!part.text) return [] + return [ + { + content: { + type: "text", + text: part.text, + ...partAudience(part), + }, + }, + ] + + case "file": + return filePartToContentChunks(part) + + case "reasoning": + if (!part.text) return [] + return [ + { + content: { + type: "text", + text: part.text, + }, + }, + ] + } +} + +function resourceLinkToPart(link: ResourceLink): PromptPart { + const parsed = uriToFilePart(link.uri, link.mimeType ?? "text/plain", link.name) + if (parsed.type === "file") return parsed + return { type: "text", text: parsed.text } +} + +function uriToFilePart( + uri: string, + mime: string, + filename?: string, +): SessionLegacy.FilePartInput | SessionLegacy.TextPartInput { + try { + if (uri.startsWith("file://")) { + return { + type: "file", + url: uri, + filename: filename ?? filenameFromUri(uri) ?? "file", + mime, + } + } + if (uri.startsWith("zed://")) { + const pathname = new URL(uri).searchParams.get("path") + if (pathname) { + return { + type: "file", + url: pathToFileURL(pathname).href, + filename: filename ?? (path.basename(pathname) || "file"), + mime, + } + } + } + return { type: "text", text: uri } + } catch { + return { type: "text", text: uri } + } +} + +function filePartToContentChunks(part: Extract): ContentChunk[] { + if (part.url.startsWith("file://")) { + return [ + { + content: { + type: "resource_link", + uri: part.url, + name: part.filename ?? "file", + mimeType: part.mime, + }, + }, + ] + } + if (!part.url.startsWith("data:")) return [] + + const data = decodeDataUrl(part.url) + if (!data) return [] + if (data.mime.startsWith("image/")) { + return [ + { + content: { + type: "image", + mimeType: data.mime, + data: data.base64, + uri: pathToFileURL(part.filename ?? "image").href, + }, + }, + ] + } + + return [ + { + content: { + type: "resource", + resource: + data.mime.startsWith("text/") || data.mime === "application/json" + ? { + uri: pathToFileURL(part.filename ?? "file").href, + mimeType: data.mime, + text: Buffer.from(data.base64, "base64").toString("utf8"), + } + : { + uri: pathToFileURL(part.filename ?? "file").href, + mimeType: data.mime, + blob: data.base64, + }, + }, + }, + ] +} + +function decodeDataUrl(url: string) { + const match = /^data:([^;]+);base64,(.*)$/.exec(url) + if (!match) return + return { mime: match[1], base64: match[2] } +} + +function audienceFlags(audience: readonly Role[] | null | undefined) { + if (audience?.length === 1 && audience[0] === "assistant") return { synthetic: true } + if (audience?.length === 1 && audience[0] === "user") return { ignored: true } + return {} +} + +function partAudience(part: Extract) { + const audience: Role[] | undefined = part.synthetic ? ["assistant"] : part.ignored ? ["user"] : undefined + if (!audience) return {} + return { annotations: { audience } } +} + +function filenameFromUri(uri: string | undefined) { + if (!uri) return + if (uri.startsWith("data:")) return + try { + const parsed = new URL(uri) + const name = path.basename(parsed.pathname) + return name || undefined + } catch { + return path.basename(uri) || undefined + } +} diff --git a/packages/opencode/src/acp-next/directory.ts b/packages/opencode/src/acp-next/directory.ts new file mode 100644 index 000000000..dabe498b8 --- /dev/null +++ b/packages/opencode/src/acp-next/directory.ts @@ -0,0 +1,209 @@ +import { Agent } from "@/agent/agent" +import { Command } from "@/command" +import { InstanceRef } from "@/effect/instance-ref" +import { InstanceStore } from "@/project/instance-store" +import { ProviderV2 } from "@opencode-ai/core/provider" +import { Provider } from "@/provider/provider" +import { Context, Effect, Layer, SynchronizedRef } from "effect" +import type * as ACPNextError from "./error" + +export type ModelOption = { + readonly providerID: ProviderV2.ID + readonly providerName: string + readonly modelID: ProviderV2.ModelID + readonly modelName: string +} + +export type ModeOption = { + readonly id: string + readonly name: string + readonly description?: string +} + +export type ModelVariants = NonNullable + +export type DefaultModel = { + readonly providerID: ProviderV2.ID + readonly modelID: ProviderV2.ModelID +} + +export type Snapshot = { + readonly directory: string + readonly providers: Record + readonly modelOptions: readonly ModelOption[] + readonly variantsByModel: Readonly> + readonly availableModes: readonly ModeOption[] + readonly defaultModeID: string + readonly availableCommands: readonly Command.Info[] + readonly defaultModel?: DefaultModel +} + +export interface LoaderInterface { + readonly load: (directory: string) => Effect.Effect +} + +export interface Interface { + readonly get: (directory: string) => Effect.Effect + readonly refresh: (directory: string) => Effect.Effect + readonly variants: (snapshot: Snapshot, model: DefaultModel) => ModelVariants | undefined +} + +export class Loader extends Context.Service()("@opencode/ACPNextDirectoryLoader") {} + +export class Service extends Context.Service()("@opencode/ACPNextDirectory") {} + +export const modelKey = (model: DefaultModel) => `${model.providerID}/${model.modelID}` + +export const variants = (snapshot: Snapshot, model: DefaultModel) => snapshot.variantsByModel[modelKey(model)] + +export const build = (input: { + readonly directory: string + readonly providers: Record + readonly modes: readonly ModeOption[] + readonly defaultModeID: string + readonly commands: readonly Command.Info[] + readonly defaultModel?: DefaultModel +}): Snapshot => { + const modelOptions = Provider.sort( + Object.values(input.providers).flatMap((provider) => + Object.values(provider.models).map((model) => ({ + id: model.id, + providerID: provider.id, + providerName: provider.name, + modelID: model.id, + modelName: model.name, + })), + ), + ).map((model) => ({ + providerID: model.providerID, + providerName: model.providerName, + modelID: model.modelID, + modelName: model.modelName, + })) + + return { + directory: input.directory, + providers: input.providers, + modelOptions, + variantsByModel: Object.fromEntries( + Object.values(input.providers).flatMap((provider) => + Object.values(provider.models).flatMap((model) => + model.variants ? [[modelKey({ providerID: provider.id, modelID: model.id }), model.variants]] : [], + ), + ), + ), + availableModes: input.modes, + defaultModeID: input.modes.some((mode) => mode.id === input.defaultModeID) + ? input.defaultModeID + : (input.modes[0]?.id ?? input.defaultModeID), + availableCommands: input.commands, + ...(input.defaultModel ? { defaultModel: input.defaultModel } : {}), + } +} + +export const loaderLayer = Layer.effect( + Loader, + Effect.gen(function* () { + const store = yield* InstanceStore.Service + const provider = yield* Provider.Service + const agent = yield* Agent.Service + const command = yield* Command.Service + + return Loader.of({ + load: Effect.fn("ACPNextDirectoryLoader.load")(function* (directory) { + const ctx = yield* store.load({ directory }) + return yield* Effect.gen(function* () { + const providers = yield* provider.list() + const [agents, defaultAgent, commands, defaultModel] = yield* Effect.all( + [agent.list(), agent.defaultInfo(), command.list(), provider.defaultModel().pipe(Effect.option)], + { concurrency: "unbounded" }, + ) + return build({ + directory, + providers, + modes: agents + .filter((item) => item.mode !== "subagent" && item.hidden !== true) + .map((item) => ({ + id: item.name, + name: item.name, + ...(item.description ? { description: item.description } : {}), + })), + defaultModeID: defaultAgent.name, + commands: commands.toSorted((a, b) => a.name.localeCompare(b.name)), + ...(defaultModel._tag === "Some" ? { defaultModel: defaultModel.value } : {}), + }) + }).pipe(Effect.provideService(InstanceRef, ctx)) + }), + }) + }), +) + +export const layer = Layer.effect( + Service, + Effect.gen(function* () { + const loader = yield* Loader + const snapshots = yield* SynchronizedRef.make(new Map>()) + + const cached = Effect.fnUntraced(function* (directory: string) { + return yield* SynchronizedRef.modifyEffect( + snapshots, + Effect.fnUntraced(function* (items) { + const current = items.get(directory) + if (current) return [current, items] as const + const next = yield* Effect.cached( + loader.load(directory).pipe( + Effect.tapError(() => + SynchronizedRef.update(snapshots, (state) => { + const next = new Map(state) + next.delete(directory) + return next + }), + ), + ), + ) + return [next, new Map(items).set(directory, next)] as const + }), + ) + }) + + const get = Effect.fn("ACPNextDirectory.get")(function* (directory: string) { + return yield* yield* cached(directory) + }) + + const refresh = Effect.fn("ACPNextDirectory.refresh")(function* (directory: string) { + return yield* SynchronizedRef.modifyEffect( + snapshots, + Effect.fnUntraced(function* (items) { + const next = yield* Effect.cached( + loader.load(directory).pipe( + Effect.tapError(() => + SynchronizedRef.update(snapshots, (state) => { + const next = new Map(state) + next.delete(directory) + return next + }), + ), + ), + ) + return [next, new Map(items).set(directory, next)] as const + }), + ).pipe(Effect.flatten) + }) + + return Service.of({ + get, + refresh, + variants, + }) + }), +) + +export const defaultLayer = layer.pipe( + Layer.provide(loaderLayer), + Layer.provide(Provider.defaultLayer), + Layer.provide(Agent.defaultLayer), + Layer.provide(Command.defaultLayer), + Layer.provide(InstanceStore.defaultLayer), +) + +export * as Directory from "./directory" diff --git a/packages/opencode/src/acp-next/error.ts b/packages/opencode/src/acp-next/error.ts new file mode 100644 index 000000000..1d4af53b5 --- /dev/null +++ b/packages/opencode/src/acp-next/error.ts @@ -0,0 +1,93 @@ +import { RequestError } from "@agentclientprotocol/sdk" +import { Schema } from "effect" + +export class SessionNotFoundError extends Schema.TaggedErrorClass()( + "ACPNextSessionNotFoundError", + { + sessionId: Schema.String, + }, +) {} + +export class InvalidConfigOptionError extends Schema.TaggedErrorClass()( + "ACPNextInvalidConfigOptionError", + { + configId: Schema.String, + }, +) {} + +export class InvalidModelError extends Schema.TaggedErrorClass()("ACPNextInvalidModelError", { + modelId: Schema.String, + providerId: Schema.optional(Schema.String), +}) {} + +export class InvalidEffortError extends Schema.TaggedErrorClass()("ACPNextInvalidEffortError", { + effort: Schema.String, +}) {} + +export class InvalidModeError extends Schema.TaggedErrorClass()("ACPNextInvalidModeError", { + mode: Schema.String, +}) {} + +export class AuthRequiredError extends Schema.TaggedErrorClass()("ACPNextAuthRequiredError", { + providerId: Schema.optional(Schema.String), +}) {} + +export class UnknownAuthMethodError extends Schema.TaggedErrorClass()( + "ACPNextUnknownAuthMethodError", + { + methodId: Schema.String, + }, +) {} + +export class UnsupportedOperationError extends Schema.TaggedErrorClass()( + "ACPNextUnsupportedOperationError", + { + method: Schema.String, + }, +) {} + +export class ServiceFailureError extends Schema.TaggedErrorClass()("ACPNextServiceFailureError", { + safeMessage: Schema.String, + service: Schema.optional(Schema.String), +}) {} + +export type Error = + | SessionNotFoundError + | InvalidConfigOptionError + | InvalidModelError + | InvalidEffortError + | InvalidModeError + | AuthRequiredError + | UnknownAuthMethodError + | UnsupportedOperationError + | ServiceFailureError + +export function toRequestError(error: Error) { + switch (error._tag) { + case "ACPNextSessionNotFoundError": + return RequestError.invalidParams({ sessionId: error.sessionId }, `session not found: ${error.sessionId}`) + case "ACPNextInvalidConfigOptionError": + return RequestError.invalidParams({ configId: error.configId }, `unknown config option: ${error.configId}`) + case "ACPNextInvalidModelError": + return RequestError.invalidParams( + { providerId: error.providerId, modelId: error.modelId }, + `model not found: ${error.modelId}`, + ) + case "ACPNextInvalidEffortError": + return RequestError.invalidParams({ effort: error.effort }, `effort not found: ${error.effort}`) + case "ACPNextInvalidModeError": + return RequestError.invalidParams({ mode: error.mode }, `mode not found: ${error.mode}`) + case "ACPNextAuthRequiredError": + return RequestError.authRequired({ providerId: error.providerId }, "provider authentication required") + case "ACPNextUnknownAuthMethodError": + return RequestError.invalidParams({ methodId: error.methodId }, `unknown auth method: ${error.methodId}`) + case "ACPNextUnsupportedOperationError": + return RequestError.methodNotFound(error.method) + case "ACPNextServiceFailureError": + return RequestError.internalError({ service: error.service }, error.safeMessage) + } +} + +export function fromUnknownDefect(_defect: unknown, safeMessage = "Internal service failure") { + return new ServiceFailureError({ safeMessage }) +} diff --git a/packages/opencode/src/acp-next/service.ts b/packages/opencode/src/acp-next/service.ts new file mode 100644 index 000000000..4516a3047 --- /dev/null +++ b/packages/opencode/src/acp-next/service.ts @@ -0,0 +1,654 @@ +import { + type AgentSideConnection, + type AuthenticateRequest, + type AuthenticateResponse, + type AuthMethod, + type CancelNotification, + type InitializeRequest, + type InitializeResponse, + type LoadSessionRequest, + type LoadSessionResponse, + type McpServer, + type NewSessionRequest, + type NewSessionResponse, + type PromptRequest, + type PromptResponse, + type SetSessionConfigOptionRequest, + type SetSessionConfigOptionResponse, + type SetSessionModelRequest, + type SetSessionModelResponse, + type SetSessionModeRequest, + type SetSessionModeResponse, +} from "@agentclientprotocol/sdk" +import { InstallationVersion } from "@opencode-ai/core/installation/version" +import type { OpencodeClient } from "@opencode-ai/sdk/v2" +import { Context, Effect, Layer, ManagedRuntime } from "effect" +import * as ACPNextError from "./error" +import { buildConfigOptions, parseModelSelection } from "./config-option" +import { Directory } from "./directory" +import { ACPNextSession } from "./session" +import { ProviderV2 } from "@opencode-ai/core/provider" +import { Provider } from "@/provider/provider" +import type { Command } from "@/command" + + +export const AuthMethodID = "opencode-login" + +export type Error = ACPNextError.Error + +export type Interface = { + readonly initialize: (input: InitializeRequest) => Effect.Effect + readonly authenticate: (input: AuthenticateRequest) => Effect.Effect + readonly newSession: (input: NewSessionRequest) => Effect.Effect + readonly loadSession: (input: LoadSessionRequest) => Effect.Effect + readonly setSessionConfigOption: ( + input: SetSessionConfigOptionRequest, + ) => Effect.Effect + readonly setSessionMode: (input: SetSessionModeRequest) => Effect.Effect + readonly setSessionModel: (input: SetSessionModelRequest) => Effect.Effect + readonly prompt: (input: PromptRequest) => Effect.Effect + readonly cancel: (input: CancelNotification) => Effect.Effect +} + +export class Service extends Context.Service()("@opencode/ACPNext/Service") {} + +export function make(input: { + sdk: OpencodeClient + connection?: Pick + directory?: Directory.Interface + session?: ACPNextSession.Interface +}): Interface { + const session = input.session ?? makeSessionService() + const directoryService = input.directory ?? makeDirectoryService(input.sdk) + const registeredMcp = new Map>() + + const initialize = Effect.fn("ACPNext.initialize")(function* (params: InitializeRequest) { + const authMethod: AuthMethod = { + description: "Run `opencode auth login` in the terminal", + name: "Login with opencode", + id: AuthMethodID, + } + + if (params.clientCapabilities?._meta?.["terminal-auth"] === true) { + authMethod._meta = { + "terminal-auth": { + command: "opencode", + args: ["auth", "login"], + label: "OpenCode Login", + }, + } + } + + return { + protocolVersion: 1, + agentCapabilities: { + loadSession: true, + mcpCapabilities: { + http: true, + sse: true, + }, + promptCapabilities: { + embeddedContext: true, + image: true, + }, + }, + authMethods: [authMethod], + agentInfo: { + name: "OpenCode", + version: InstallationVersion, + }, + } + }) + + const authenticate = Effect.fn("ACPNext.authenticate")(function* (params: AuthenticateRequest) { + if (params.methodId !== AuthMethodID) { + return yield* new ACPNextError.UnknownAuthMethodError({ methodId: params.methodId }) + } + return {} + }) + + const directorySnapshot = Effect.fn("ACPNext.directorySnapshot")(function* (cwd: string) { + return yield* directoryService.get(cwd) + }) + + const newSession = Effect.fn("ACPNext.newSession")(function* (params: NewSessionRequest) { + const snapshot = yield* directorySnapshot(params.cwd) + const selected = selectDefaultModel(snapshot) + const variant = selectVariant(snapshot, selected) + const modeId = snapshot.availableModes.length > 0 ? snapshot.defaultModeID : undefined + const created = yield* request( + () => + input.sdk.session.create( + { + directory: params.cwd, + ...(modeId ? { agent: modeId } : {}), + model: { + providerID: selected.providerID, + id: selected.modelID, + ...(variant ? { variant } : {}), + }, + }, + { throwOnError: true }, + ), + "session", + ) + const state = yield* session.create({ + id: created.id, + cwd: params.cwd, + mcpServers: params.mcpServers, + model: selected, + variant, + modeId, + }) + + yield* registerMcpServers(input.sdk, registeredMcp, params.cwd, state.id, params.mcpServers) + yield* sendAvailableCommands(input.connection, state.id, snapshot) + + return { + sessionId: state.id, + configOptions: configOptions(snapshot, { + model: state.model ?? selected, + variant: state.variant, + modeId: state.modeId, + }), + } + }) + + const loadSession = Effect.fn("ACPNext.loadSession")(function* (params: LoadSessionRequest) { + const snapshot = yield* directorySnapshot(params.cwd) + yield* request( + () => input.sdk.session.get({ directory: params.cwd, sessionID: params.sessionId }, { throwOnError: true }), + "session", + ) + const messages = yield* request( + () => + input.sdk.session.messages( + { directory: params.cwd, sessionID: params.sessionId, limit: 100 }, + { throwOnError: true }, + ), + "session", + ) + const restored = restoreFromMessages(messages.map((item) => item.info)) + const model = restored.model ?? selectDefaultModel(snapshot) + const state = yield* session.load({ + id: params.sessionId, + cwd: params.cwd, + mcpServers: params.mcpServers, + model, + variant: restored.variant ?? selectVariant(snapshot, model), + modeId: restored.modeId ?? (snapshot.availableModes.length > 0 ? snapshot.defaultModeID : undefined), + }) + + yield* registerMcpServers(input.sdk, registeredMcp, params.cwd, state.id, params.mcpServers) + yield* sendAvailableCommands(input.connection, state.id, snapshot) + + return { + sessionId: state.id, + configOptions: configOptions(snapshot, { + model: state.model ?? model, + variant: state.variant, + modeId: state.modeId, + }), + } + }) + + const setSessionConfigOption = Effect.fn("ACPNext.setSessionConfigOption")(function* ( + params: SetSessionConfigOptionRequest, + ) { + const current = yield* session.get(params.sessionId) + const snapshot = yield* directorySnapshot(current.cwd) + if (typeof params.value !== "string") { + return yield* new ACPNextError.InvalidConfigOptionError({ configId: params.configId }) + } + + if (params.configId === "model") { + const selected = yield* parseSelectedModel(snapshot, params.value) + const variant = selected.variant ?? selectVariant(snapshot, selected.model) + const state = yield* session + .setVariant(params.sessionId, Directory.variants(snapshot, selected.model) ? variant : undefined) + .pipe(Effect.andThen(session.setModel(params.sessionId, selected.model))) + return { + configOptions: configOptions(snapshot, { + model: state.model ?? selected.model, + variant: state.variant, + modeId: state.modeId, + }), + } + } + + if (params.configId === "effort") { + const model = current.model ?? selectDefaultModel(snapshot) + const variants = Directory.variants(snapshot, model) + if (!variants || !Object.keys(variants).includes(params.value)) { + return yield* new ACPNextError.InvalidEffortError({ effort: params.value }) + } + const state = yield* session.setVariant(params.sessionId, params.value) + return { + configOptions: configOptions(snapshot, { + model: state.model ?? model, + variant: state.variant, + modeId: state.modeId, + }), + } + } + + if (params.configId === "mode") { + if (!snapshot.availableModes.some((mode) => mode.id === params.value)) { + return yield* new ACPNextError.InvalidModeError({ mode: params.value }) + } + const state = yield* session.setMode(params.sessionId, params.value) + return { + configOptions: configOptions(snapshot, { + model: state.model ?? selectDefaultModel(snapshot), + variant: state.variant, + modeId: state.modeId, + }), + } + } + + return yield* new ACPNextError.InvalidConfigOptionError({ configId: params.configId }) + }) + + const setSessionMode = Effect.fn("ACPNext.setSessionMode")(function* (params: SetSessionModeRequest) { + const current = yield* session.get(params.sessionId) + const snapshot = yield* directorySnapshot(current.cwd) + if (!snapshot.availableModes.some((mode) => mode.id === params.modeId)) { + return yield* new ACPNextError.InvalidModeError({ mode: params.modeId }) + } + yield* session.setMode(params.sessionId, params.modeId) + return {} + }) + + const setSessionModel = Effect.fn("ACPNext.setSessionModel")(function* (params: SetSessionModelRequest) { + const current = yield* session.get(params.sessionId) + const snapshot = yield* directorySnapshot(current.cwd) + const selected = yield* parseSelectedModel(snapshot, params.modelId) + yield* session + .setVariant( + params.sessionId, + Directory.variants(snapshot, selected.model) + ? (selected.variant ?? selectVariant(snapshot, selected.model)) + : undefined, + ) + .pipe(Effect.andThen(session.setModel(params.sessionId, selected.model))) + return {} + }) + + return { + initialize, + authenticate, + newSession, + loadSession, + setSessionConfigOption, + setSessionMode, + setSessionModel, + prompt: Effect.fn("ACPNext.prompt")(function* (_input: PromptRequest) { + return yield* new ACPNextError.UnsupportedOperationError({ method: "session/prompt" }) + }), + cancel: Effect.fn("ACPNext.cancel")(function* (_input: CancelNotification) { + return yield* new ACPNextError.UnsupportedOperationError({ method: "session/cancel" }) + }), + } +} + +function makeSessionService() { + return ManagedRuntime.make(ACPNextSession.defaultLayer).runSync( + ACPNextSession.Service.use((service) => Effect.succeed(service)), + ) +} + +function makeDirectoryService(sdk: OpencodeClient) { + return ManagedRuntime.make( + Directory.layer.pipe( + Layer.provide( + Layer.succeed( + Directory.Loader, + Directory.Loader.of({ + load: (directory) => request(() => loadDirectorySnapshot(sdk, directory), "directory"), + }), + ), + ), + ), + ).runSync(Directory.Service.use((service) => Effect.succeed(service))) +} + +type ConfigState = { + readonly model: Directory.DefaultModel + readonly variant?: string + readonly modeId?: string +} + +type SdkResponse = { + readonly data?: T + readonly error?: unknown +} + +type MessageInfo = { + readonly role?: string + readonly model?: { + readonly providerID?: string + readonly modelID?: string + readonly variant?: string + } + readonly providerID?: string + readonly modelID?: string + readonly variant?: string + readonly mode?: string + readonly agent?: string +} + +function request(fn: () => Promise>, service?: string) { + return Effect.tryPromise({ + try: async () => { + const result = await fn() + if (isSdkResponse(result)) { + if (result.error) throw result.error + if (result.data !== undefined) return result.data + } + return result as T + }, + catch: (error) => fromUnknownError(error, service), + }) +} + +async function loadDirectorySnapshot(sdk: OpencodeClient, directory: string) { + const [providersResponse, agentsResponse, commandsResponse, skillsResponse] = await Promise.all([ + sdk.config.providers({ directory }, { throwOnError: true }), + sdk.app.agents({ directory }, { throwOnError: true }), + sdk.command.list({ directory }, { throwOnError: true }), + sdk.app.skills({ directory }, { throwOnError: true }), + ]) + const providersData = providersResponse.data! + const agents = agentsResponse.data! + const commandsData = commandsResponse.data! + const skills = skillsResponse.data! + const providers = Object.fromEntries(providersData.providers.map((provider) => [provider.id, provider])) as Record< + ProviderV2.ID, + Provider.Info + > + const defaultModel = await defaultModelFromSdk(sdk, directory, providers) + const modes = agents + .filter((agent) => agent.mode !== "subagent" && agent.hidden !== true) + .map((agent) => ({ + id: agent.name, + name: agent.name, + ...(agent.description ? { description: agent.description } : {}), + })) + const commands = [ + ...commandsData, + ...skills + .filter((skill) => !commandsData.some((command) => command.name === skill.name)) + .map((skill) => ({ + name: skill.name, + description: skill.description, + source: "skill" as const, + template: skill.content, + hints: [], + })), + ] as Command.Info[] + + return Directory.build({ + directory, + providers, + modes, + defaultModeID: agents.find((agent) => agent.mode === "primary" && agent.hidden !== true)?.name ?? "build", + commands: commands.toSorted((a, b) => a.name.localeCompare(b.name)), + ...(defaultModel ? { defaultModel } : {}), + }) +} + +async function defaultModelFromSdk( + sdk: OpencodeClient, + directory: string, + providers: Record, +): Promise { + const configured = await sdk.config + .get({ directory }, { throwOnError: true }) + .then((response) => (response.data?.model ? Provider.parseModel(response.data.model) : undefined)) + .catch(() => undefined) + if (configured && providers[configured.providerID]?.models[configured.modelID]) return configured + + const lastUsed = await lastUsedModel(sdk, directory, providers) + if (lastUsed) return lastUsed + + const opencodeProvider = providers[ProviderV2.ID.make("opencode")] + const opencodeModel = opencodeProvider ? Provider.sort(Object.values(opencodeProvider.models))[0] : undefined + if (opencodeProvider && opencodeModel) return { providerID: opencodeProvider.id, modelID: opencodeModel.id } + + const best = Provider.sort(Object.values(providers).flatMap((provider) => Object.values(provider.models)))[0] + if (best) return { providerID: best.providerID, modelID: best.id } + if (configured) return configured +} + +async function lastUsedModel( + sdk: OpencodeClient, + directory: string, + providers: Record, +): Promise { + const session = await sdk.session + .list({ directory, roots: true, limit: 1 }, { throwOnError: true }) + .then((response) => response.data?.[0]) + .catch(() => undefined) + if (!session) return + + const lastUser = await sdk.session + .messages({ directory, sessionID: session.id, limit: 20 }, { throwOnError: true }) + .then((response) => response.data?.findLast((message) => message.info.role === "user")?.info) + .catch(() => undefined) + if (lastUser?.role !== "user") return + if (!providers[ProviderV2.ID.make(lastUser.model.providerID)]?.models[ProviderV2.ModelID.make(lastUser.model.modelID)]) return + + return { + providerID: ProviderV2.ID.make(lastUser.model.providerID), + modelID: ProviderV2.ModelID.make(lastUser.model.modelID), + } +} + +function selectDefaultModel(snapshot: Directory.Snapshot) { + if (snapshot.defaultModel) return snapshot.defaultModel + const model = snapshot.modelOptions[0] + if (model) return { providerID: model.providerID, modelID: model.modelID } + return { providerID: "unknown" as ProviderV2.ID, modelID: "unknown" as ProviderV2.ModelID } +} + +function selectVariant(snapshot: Directory.Snapshot, model: Directory.DefaultModel) { + const variants = Directory.variants(snapshot, model) + if (!variants) return + if (variants.default) return "default" + return Object.keys(variants)[0] +} + +function configOptions(snapshot: Directory.Snapshot, session: ConfigState) { + return buildConfigOptions({ + providers: Object.values(snapshot.providers), + currentModel: session.model, + currentVariant: session.variant, + modes: snapshot.availableModes, + currentModeId: session.modeId, + }) +} + +function parseSelectedModel(snapshot: Directory.Snapshot, modelId: string) { + const selected = parseModelSelection(modelId, Object.values(snapshot.providers)) + const provider = snapshot.providers[ProviderV2.ID.make(selected.model.providerID)] + const model = provider?.models[ProviderV2.ModelID.make(selected.model.modelID)] + if (!model) { + return Effect.fail( + new ACPNextError.InvalidModelError({ + providerId: selected.model.providerID, + modelId, + }), + ) + } + if (selected.variant && !model.variants?.[selected.variant]) { + return Effect.fail(new ACPNextError.InvalidEffortError({ effort: selected.variant })) + } + return Effect.succeed({ + model: { + providerID: provider.id, + modelID: model.id, + }, + variant: selected.variant, + }) +} + +function sendAvailableCommands( + connection: Pick | undefined, + sessionId: string, + snapshot: Directory.Snapshot, +) { + if (!connection) return Effect.void + return Effect.sync(() => { + setTimeout(() => { + void connection.sessionUpdate({ + sessionId, + update: { + sessionUpdate: "available_commands_update", + availableCommands: snapshot.availableCommands.map((command) => ({ + name: command.name, + description: command.description ?? "", + })), + }, + }) + }, 0) + }) +} + +function registerMcpServers( + sdk: OpencodeClient, + registered: Map>, + directory: string, + sessionId: string, + servers: readonly McpServer[], +) { + const current = registered.get(sessionId) ?? new Set() + registered.set(sessionId, current) + const pending = new Set() + + return Effect.all( + servers + .map((server) => ({ server, config: mcpConfig(server) })) + .filter((entry) => { + const key = mcpRegistrationKey(entry.server.name, entry.config) + if (current.has(key) || pending.has(key)) return false + pending.add(key) + return true + }) + .map((entry) => + request( + () => + sdk.mcp.add( + { + directory, + name: entry.server.name, + config: entry.config, + }, + { throwOnError: true }, + ), + "mcp", + ).pipe( + Effect.tap(() => Effect.sync(() => current.add(mcpRegistrationKey(entry.server.name, entry.config)))), + Effect.ignore, + ), + ), + { concurrency: "unbounded" }, + ).pipe(Effect.asVoid) +} + +function mcpRegistrationKey(name: string, config: ReturnType) { + return `${name}:${stableStringify(config)}` +} + +function mcpConfig(server: McpServer) { + if ("type" in server) { + return { + type: "remote" as const, + url: server.url, + headers: Object.fromEntries(server.headers.map((header) => [header.name, header.value])), + } + } + return { + type: "local" as const, + command: [server.command, ...server.args], + environment: Object.fromEntries(server.env.map((entry) => [entry.name, entry.value])), + } +} + +function stableStringify(value: unknown): string { + if (Array.isArray(value)) return `[${value.map(stableStringify).join(",")}]` + if (!value || typeof value !== "object") return JSON.stringify(value) + return `{${Object.entries(value) + .toSorted(([a], [b]) => a.localeCompare(b)) + .map(([key, item]) => `${JSON.stringify(key)}:${stableStringify(item)}`) + .join(",")}}` +} + +function restoreFromMessages(messages: readonly MessageInfo[]) { + const user = messages.findLast( + (message) => message.role === "user" && message.model?.providerID && message.model.modelID, + ) + if (user?.model?.providerID && user.model.modelID) { + return { + model: { providerID: user.model.providerID as ProviderV2.ID, modelID: user.model.modelID as ProviderV2.ModelID }, + variant: user.model.variant, + modeId: user.agent, + } + } + + const assistant = messages.findLast((message) => message.providerID && message.modelID) + if (assistant?.providerID && assistant.modelID) { + return { + model: { providerID: assistant.providerID as ProviderV2.ID, modelID: assistant.modelID as ProviderV2.ModelID }, + variant: assistant.variant, + modeId: assistant.mode ?? assistant.agent, + } + } + + return {} +} + +function isSdkResponse(value: T | SdkResponse): value is SdkResponse { + return typeof value === "object" && value !== null && ("data" in value || "error" in value) +} + +function fromUnknownError(error: unknown, service?: string): Error { + if (isACPNextError(error)) return error + if (isAuthRequired(error)) { + return new ACPNextError.AuthRequiredError({ providerId: findProviderID(error) }) + } + return new ACPNextError.ServiceFailureError({ safeMessage: "OpenCode service failure", service }) +} + +function isACPNextError(error: unknown): error is Error { + return ( + typeof error === "object" && + error !== null && + "_tag" in error && + typeof error._tag === "string" && + error._tag.startsWith("ACPNext") + ) +} + +function isAuthRequired(value: unknown): boolean { + if (typeof value !== "object" || value === null) return false + if (value instanceof Error && (value.name === "ProviderAuthError" || value.name === "LoadAPIKeyError")) return true + if ( + value instanceof Error && + (value.message.includes("ProviderAuthError") || value.message.includes("LoadAPIKeyError")) + ) { + return true + } + if ("name" in value && (value.name === "ProviderAuthError" || value.name === "LoadAPIKeyError")) return true + if ("_tag" in value && (value._tag === "ProviderAuthError" || value._tag === "LoadAPIKeyError")) return true + if ("error" in value && isAuthRequired(value.error)) return true + if ("data" in value && isAuthRequired(value.data)) return true + return false +} + +function findProviderID(value: unknown): string | undefined { + if (typeof value !== "object" || value === null) return + if ("providerID" in value && typeof value.providerID === "string") return value.providerID + if ("providerId" in value && typeof value.providerId === "string") return value.providerId + if ("data" in value) return findProviderID(value.data) + if ("error" in value) return findProviderID(value.error) +} diff --git a/packages/opencode/src/acp-next/session.ts b/packages/opencode/src/acp-next/session.ts new file mode 100644 index 000000000..99c8542e6 --- /dev/null +++ b/packages/opencode/src/acp-next/session.ts @@ -0,0 +1,216 @@ +import type { McpServer } from "@agentclientprotocol/sdk" +import { Context, Effect, Layer, Ref } from "effect" +import type { ProviderV2 } from "@opencode-ai/core/provider" +import * as ACPNextError from "./error" + + +export type SelectedModel = { + providerID: ProviderV2.ID + modelID: ProviderV2.ModelID +} + +export type KnownMessagePartMetadata = { + messageId: string + partId: string + toolCallId?: string + metadata?: unknown +} + +export type Info = { + id: string + cwd: string + mcpServers: readonly McpServer[] + createdAt: Date + model?: SelectedModel + variant?: string + modeId?: string + knownParts: ReadonlyMap +} + +export type StoreInput = { + id: string + cwd: string + mcpServers?: readonly McpServer[] + createdAt?: Date + model?: SelectedModel + variant?: string + modeId?: string +} + +export type RecordPartMetadataInput = { + sessionId: string + messageId: string + partId: string + toolCallId?: string + metadata?: unknown +} + +export type PartMetadataLookupInput = { + sessionId: string + messageId: string + partId: string +} + +export type Interface = { + readonly create: (input: StoreInput) => Effect.Effect + readonly load: (input: StoreInput) => Effect.Effect + readonly get: (sessionId: string) => Effect.Effect + readonly tryGet: (sessionId: string) => Effect.Effect + readonly remove: (sessionId: string) => Effect.Effect + readonly setModel: ( + sessionId: string, + model: SelectedModel | undefined, + ) => Effect.Effect + readonly getModel: (sessionId: string) => Effect.Effect + readonly setVariant: ( + sessionId: string, + variant: string | undefined, + ) => Effect.Effect + readonly getVariant: (sessionId: string) => Effect.Effect + readonly setMode: ( + sessionId: string, + modeId: string | undefined, + ) => Effect.Effect + readonly getMode: (sessionId: string) => Effect.Effect + readonly recordPartMetadata: ( + input: RecordPartMetadataInput, + ) => Effect.Effect + readonly getPartMetadata: ( + input: PartMetadataLookupInput, + ) => Effect.Effect + readonly tryGetPartMetadata: (input: PartMetadataLookupInput) => Effect.Effect +} + +export class Service extends Context.Service()("@opencode/ACPNext/Session") {} + +type State = Map + +export const layer = Layer.effect( + Service, + Effect.gen(function* () { + const sessions = yield* Ref.make(new Map()) + + const store = Effect.fn("ACPNext.Session.store")(function* (input: StoreInput) { + const session = makeSession(input) + yield* Ref.update(sessions, (state) => new Map(state).set(session.id, session)) + return snapshot(session) + }) + + const tryGet = Effect.fn("ACPNext.Session.tryGet")(function* (sessionId: string) { + const session = (yield* Ref.get(sessions)).get(sessionId) + if (!session) return + return snapshot(session) + }) + + const get = Effect.fn("ACPNext.Session.get")(function* (sessionId: string) { + const session = yield* tryGet(sessionId) + if (session) return session + return yield* new ACPNextError.SessionNotFoundError({ sessionId }) + }) + + const update = Effect.fn("ACPNext.Session.update")(function* (sessionId: string, fn: (session: Info) => Info) { + const result = yield* Ref.modify(sessions, (state) => { + const session = state.get(sessionId) + if (!session) return [undefined, state] as const + const next = fn(session) + return [snapshot(next), new Map(state).set(sessionId, next)] as const + }) + if (result) return result + return yield* new ACPNextError.SessionNotFoundError({ sessionId }) + }) + + const remove = Effect.fn("ACPNext.Session.remove")(function* (sessionId: string) { + return yield* Ref.modify(sessions, (state) => { + const session = state.get(sessionId) + if (!session) return [undefined, state] as const + const next = new Map(state) + next.delete(sessionId) + return [snapshot(session), next] as const + }) + }) + + const setModel: Interface["setModel"] = Effect.fn("ACPNext.Session.setModel")((sessionId, model) => + update(sessionId, (session) => ({ ...session, model })), + ) + + const setVariant: Interface["setVariant"] = Effect.fn("ACPNext.Session.setVariant")((sessionId, variant) => + update(sessionId, (session) => ({ ...session, variant })), + ) + + const setMode: Interface["setMode"] = Effect.fn("ACPNext.Session.setMode")((sessionId, modeId) => + update(sessionId, (session) => ({ ...session, modeId })), + ) + + const recordPartMetadata: Interface["recordPartMetadata"] = Effect.fn("ACPNext.Session.recordPartMetadata")(( + input, + ) => { + const metadata = { + messageId: input.messageId, + partId: input.partId, + toolCallId: input.toolCallId, + metadata: input.metadata, + } + return update(input.sessionId, (session) => ({ + ...session, + knownParts: new Map(session.knownParts).set(partMetadataKey(input), metadata), + })).pipe(Effect.as(metadata)) + }) + + return Service.of({ + create: store, + load: store, + get, + tryGet, + remove, + setModel, + getModel: Effect.fn("ACPNext.Session.getModel")(function* (sessionId) { + return (yield* get(sessionId)).model + }), + setVariant, + getVariant: Effect.fn("ACPNext.Session.getVariant")(function* (sessionId) { + return (yield* get(sessionId)).variant + }), + setMode, + getMode: Effect.fn("ACPNext.Session.getMode")(function* (sessionId) { + return (yield* get(sessionId)).modeId + }), + recordPartMetadata, + getPartMetadata: Effect.fn("ACPNext.Session.getPartMetadata")(function* (input) { + return (yield* get(input.sessionId)).knownParts.get(partMetadataKey(input)) + }), + tryGetPartMetadata: Effect.fn("ACPNext.Session.tryGetPartMetadata")(function* (input) { + return (yield* tryGet(input.sessionId))?.knownParts.get(partMetadataKey(input)) + }), + }) + }), +) + +export const defaultLayer = layer + +function makeSession(input: StoreInput): Info { + return { + id: input.id, + cwd: input.cwd, + mcpServers: [...(input.mcpServers ?? [])], + createdAt: input.createdAt ? new Date(input.createdAt) : new Date(), + model: input.model, + variant: input.variant, + modeId: input.modeId, + knownParts: new Map(), + } +} + +function snapshot(session: Info): Info { + return { + ...session, + mcpServers: [...session.mcpServers], + createdAt: new Date(session.createdAt), + knownParts: new Map(session.knownParts), + } +} + +function partMetadataKey(input: { messageId: string; partId: string }) { + return `${input.messageId}:${input.partId}` +} + +export * as ACPNextSession from "./session" diff --git a/packages/opencode/src/acp-next/tool.ts b/packages/opencode/src/acp-next/tool.ts new file mode 100644 index 000000000..128c4c9c8 --- /dev/null +++ b/packages/opencode/src/acp-next/tool.ts @@ -0,0 +1,174 @@ +import type { ToolCallContent, ToolCallLocation, ToolKind } from "@agentclientprotocol/sdk" + +export type ToolInput = Record + +export type ToolAttachment = { + readonly mime?: string + readonly url?: string + readonly [key: string]: unknown +} + +export type CompletedToolState = { + readonly status: "completed" + readonly input: ToolInput + readonly output: string + readonly metadata?: unknown + readonly attachments?: ReadonlyArray +} + +export type ImageAttachment = { + readonly mimeType: string + readonly data: string +} + +export function toToolKind(toolName: string): ToolKind { + const tool = toolName.toLocaleLowerCase() + + switch (tool) { + case "bash": + case "shell": + return "execute" + + case "webfetch": + return "fetch" + + case "edit": + case "patch": + case "write": + return "edit" + + case "grep": + case "glob": + case "repo_clone": + case "repo_overview": + case "context": + case "context7_resolve_library_id": + case "context7_get_library_docs": + return "search" + + case "read": + return "read" + + default: + return "other" + } +} + +export function toLocations(toolName: string, input: ToolInput): ToolCallLocation[] { + const tool = toolName.toLocaleLowerCase() + + switch (tool) { + case "read": + case "edit": + case "write": + return locationFrom(input.filePath) + + case "grep": + case "glob": + case "repo_clone": + case "repo_overview": + case "context": + case "context7_resolve_library_id": + case "context7_get_library_docs": + return locationFrom(input.path) + + case "bash": + case "shell": + return [] + + default: + return [] + } +} + +export function completedToolContent(toolName: string, state: CompletedToolState): ToolCallContent[] { + const content: ToolCallContent[] = [ + { + type: "content", + content: { + type: "text", + text: state.output, + }, + }, + ] + + if (toToolKind(toolName) === "edit") { + content.push(...diffContent(state.input)) + } + + content.push(...imageContents(state.attachments ?? [])) + return content +} + +export function completedToolRawOutput(state: CompletedToolState) { + return { + output: state.output, + ...(state.metadata !== undefined ? { metadata: state.metadata } : {}), + ...(state.attachments?.length ? { attachments: state.attachments } : {}), + } +} + +export function imageContents(attachments: ReadonlyArray): ToolCallContent[] { + return extractImageAttachments(attachments).map((attachment): ToolCallContent => { + return { + type: "content", + content: { + type: "image", + mimeType: attachment.mimeType, + data: attachment.data, + }, + } + }) +} + +export function extractImageAttachments(attachments: ReadonlyArray): ImageAttachment[] { + return attachments.flatMap((attachment): ImageAttachment[] => { + const data = dataUrlImage(attachment) + return data ? [data] : [] + }) +} + +export function shellOutputSnapshot(state: { readonly metadata?: unknown }) { + if (!state.metadata || typeof state.metadata !== "object") return undefined + return stringValue((state.metadata as Record).output) +} + +export const mapToolKind = toToolKind +export const extractLocations = toLocations +export const buildCompletedToolContent = completedToolContent +export const buildCompletedRawOutput = completedToolRawOutput +export const extractShellOutputSnapshot = shellOutputSnapshot + +function locationFrom(value: unknown): ToolCallLocation[] { + const path = stringValue(value) + return path ? [{ path }] : [] +} + +function diffContent(input: ToolInput): ToolCallContent[] { + const oldText = stringValue(input.oldString) + const newText = stringValue(input.newString) ?? stringValue(input.content) + if (oldText === undefined || newText === undefined) return [] + + return [ + { + type: "diff", + path: stringValue(input.filePath) ?? "", + oldText, + newText, + }, + ] +} + +function dataUrlImage(attachment: ToolAttachment) { + const match = stringValue(attachment.url)?.match(/^data:([^;,]+)(?:;[^,]*)*;base64,(.*)$/) + const mime = match?.[1] ?? stringValue(attachment.mime) + if (!mime?.startsWith("image/")) return undefined + + const data = match?.[2] + if (data === undefined) return undefined + return { mimeType: mime, data } +} + +function stringValue(value: unknown) { + return typeof value === "string" ? value : undefined +} diff --git a/packages/opencode/src/acp-next/usage.ts b/packages/opencode/src/acp-next/usage.ts new file mode 100644 index 000000000..54e3b4f9c --- /dev/null +++ b/packages/opencode/src/acp-next/usage.ts @@ -0,0 +1,250 @@ +import type { AgentSideConnection, Usage } from "@agentclientprotocol/sdk" +import * as Log from "@opencode-ai/core/util/log" +import { InstanceRef } from "@/effect/instance-ref" +import { InstanceStore } from "@/project/instance-store" +import { ProviderV2 } from "@opencode-ai/core/provider" +import { Provider } from "@/provider/provider" +import { Context, Effect, Layer, SynchronizedRef } from "effect" + +const log = Log.create({ service: "acp-next-usage" }) + +export type AssistantTokenCost = { + readonly cost: number + readonly tokens: { + readonly input: number + readonly output: number + readonly reasoning: number + readonly cache: { + readonly read: number + readonly write: number + } + } +} + +export type AssistantMessage = AssistantTokenCost & { + readonly role: "assistant" + readonly providerID?: string + readonly modelID?: string +} + +export type SessionMessage = { + readonly info: { readonly role: string } | AssistantMessage +} + +export type MessagesInput = { + readonly sessionID: string + readonly directory: string +} + +export type SDK = { + readonly session: { + readonly messages: ( + parameters: { readonly sessionID: string; readonly directory: string }, + options: { readonly throwOnError: true }, + ) => Promise<{ readonly data?: readonly SessionMessage[] | null }> + } +} + +export interface MessageLoaderInterface { + readonly messages: (input: MessagesInput) => Effect.Effect +} + +export interface ContextLimitLoaderInterface { + readonly providers: (directory: string) => Effect.Effect, unknown> +} + +export type UsageConnection = Pick + +export interface Interface { + readonly buildUsage: (message: AssistantTokenCost) => Usage + readonly latestAssistantMessage: (messages: readonly SessionMessage[]) => AssistantMessage | undefined + readonly totalSessionCost: (messages: readonly SessionMessage[]) => number + readonly contextLimit: (input: { + readonly directory: string + readonly providerID: ProviderV2.ID + readonly modelID: ProviderV2.ModelID + }) => Effect.Effect + readonly sendUpdate: (input: { + readonly connection: UsageConnection + readonly sessionID: string + readonly directory: string + }) => Effect.Effect +} + +export class MessageLoader extends Context.Service()( + "@opencode/ACPNextUsageMessageLoader", +) {} + +export class ContextLimitLoader extends Context.Service()( + "@opencode/ACPNextUsageContextLimitLoader", +) {} + +export class Service extends Context.Service()("@opencode/ACPNextUsage") {} + +export function messageLoaderFromSDK(sdk: SDK): MessageLoaderInterface { + return MessageLoader.of({ + messages: (input) => + Effect.promise(() => + sdk.session + .messages({ sessionID: input.sessionID, directory: input.directory }, { throwOnError: true }) + .then((response) => response.data ?? []), + ), + }) +} + +export const messageLoaderLayer = (sdk: SDK) => Layer.succeed(MessageLoader, messageLoaderFromSDK(sdk)) + +export function buildUsage(message: AssistantTokenCost): Usage { + const cachedReadTokens = message.tokens.cache.read + const cachedWriteTokens = message.tokens.cache.write + const thoughtTokens = message.tokens.reasoning + + return { + inputTokens: message.tokens.input, + outputTokens: message.tokens.output, + totalTokens: message.tokens.input + message.tokens.output + thoughtTokens + cachedReadTokens + cachedWriteTokens, + ...(thoughtTokens > 0 ? { thoughtTokens } : {}), + ...(cachedReadTokens > 0 ? { cachedReadTokens } : {}), + ...(cachedWriteTokens > 0 ? { cachedWriteTokens } : {}), + } +} + +export function latestAssistantMessage(messages: readonly SessionMessage[]): AssistantMessage | undefined { + return messages + .filter((message): message is { readonly info: AssistantMessage } => message.info.role === "assistant") + .at(-1)?.info +} + +export function totalSessionCost(messages: readonly SessionMessage[]): number { + return messages + .filter((message): message is { readonly info: AssistantMessage } => message.info.role === "assistant") + .reduce((sum, message) => sum + message.info.cost, 0) +} + +export function findContextLimit( + providers: Record, + providerID: ProviderV2.ID, + modelID: ProviderV2.ModelID, +): number | undefined { + return providers[providerID]?.models[modelID]?.limit.context +} + +export const contextLimitLoaderLayer = Layer.effect( + ContextLimitLoader, + Effect.gen(function* () { + const store = yield* InstanceStore.Service + const provider = yield* Provider.Service + + return ContextLimitLoader.of({ + providers: Effect.fn("ACPNextUsageContextLimitLoader.providers")(function* (directory) { + const ctx = yield* store.load({ directory }) + return yield* Effect.gen(function* () { + return yield* provider.list() + }).pipe(Effect.provideService(InstanceRef, ctx)) + }), + }) + }), +) + +export const layer = Layer.effect( + Service, + Effect.gen(function* () { + const messageLoader = yield* MessageLoader + const contextLimitLoader = yield* ContextLimitLoader + const limits = yield* SynchronizedRef.make(new Map>()) + + const cachedLimit = Effect.fnUntraced(function* (input: { + readonly directory: string + readonly providerID: ProviderV2.ID + readonly modelID: ProviderV2.ModelID + }) { + return yield* SynchronizedRef.modifyEffect( + limits, + Effect.fnUntraced(function* (items) { + const key = `${input.directory}\u0000${input.providerID}\u0000${input.modelID}` + const current = items.get(key) + if (current) return [current, items] as const + const next = yield* Effect.cached( + contextLimitLoader.providers(input.directory).pipe( + Effect.map((providers) => findContextLimit(providers, input.providerID, input.modelID)), + Effect.catch((error) => + Effect.sync(() => { + log.error("failed to get providers for usage context limit", { error }) + return undefined + }), + ), + ), + ) + return [next, new Map(items).set(key, next)] as const + }), + ) + }) + + const contextLimit = Effect.fn("ACPNextUsage.contextLimit")(function* (input: { + readonly directory: string + readonly providerID: ProviderV2.ID + readonly modelID: ProviderV2.ModelID + }) { + return yield* yield* cachedLimit(input) + }) + + const sendUpdate = Effect.fn("ACPNextUsage.sendUpdate")(function* (input: { + readonly connection: UsageConnection + readonly sessionID: string + readonly directory: string + }) { + const messages = yield* messageLoader.messages({ sessionID: input.sessionID, directory: input.directory }).pipe( + Effect.catch((error) => + Effect.sync(() => { + log.error("failed to fetch messages for usage update", { error }) + return undefined + }), + ), + ) + if (!messages) return + + const message = latestAssistantMessage(messages) + if (!message) return + if (!message.providerID || !message.modelID) return + + const size = yield* contextLimit({ + directory: input.directory, + providerID: ProviderV2.ID.make(message.providerID), + modelID: ProviderV2.ModelID.make(message.modelID), + }) + if (!size) return + + yield* Effect.promise(() => + input.connection + .sessionUpdate({ + sessionId: input.sessionID, + update: { + sessionUpdate: "usage_update", + used: message.tokens.input + message.tokens.cache.read, + size, + cost: { amount: totalSessionCost(messages), currency: "USD" }, + }, + }) + .catch((error) => { + log.error("failed to send usage update", { error }) + }), + ) + }) + + return Service.of({ + buildUsage, + latestAssistantMessage, + totalSessionCost, + contextLimit, + sendUpdate, + }) + }), +) + +export const defaultLayer = layer.pipe( + Layer.provide(contextLimitLoaderLayer), + Layer.provide(Provider.defaultLayer), + Layer.provide(InstanceStore.defaultLayer), +) + +export * as UsageService from "./usage" diff --git a/packages/opencode/src/cli/cmd/acp.ts b/packages/opencode/src/cli/cmd/acp.ts index b3b7df486..b113a278f 100644 --- a/packages/opencode/src/cli/cmd/acp.ts +++ b/packages/opencode/src/cli/cmd/acp.ts @@ -3,10 +3,12 @@ import { Effect } from "effect" import { effectCmd } from "../effect-cmd" import { AgentSideConnection, ndJsonStream } from "@agentclientprotocol/sdk" import { ACP } from "@/acp/agent" +import { ACPNext } from "@/acp-next/agent" import { Server } from "@/server/server" import { ServerAuth } from "@/server/auth" import { createOpencodeClient } from "@opencode-ai/sdk/v2" import { withNetworkOptions, resolveNetworkOptions } from "../network" +import { RuntimeFlags } from "@/effect/runtime-flags" const log = Log.create({ service: "acp-command" }) @@ -22,6 +24,7 @@ export const AcpCommand = effectCmd({ }, handler: Effect.fn("Cli.acp")(function* (args) { process.env.OPENCODE_CLIENT = "acp" + const flags = yield* RuntimeFlags.Service const opts = yield* resolveNetworkOptions(args) const server = yield* Effect.promise(() => Server.listen(opts)) @@ -54,7 +57,7 @@ export const AcpCommand = effectCmd({ }) const stream = ndJsonStream(input, output) - const agent = ACP.init({ sdk }) + const agent = flags.acpNext ? ACPNext.init({ sdk }) : ACP.init({ sdk }) new AgentSideConnection((conn) => { return agent.create(conn, { sdk }) diff --git a/packages/opencode/src/cli/cmd/tui/component/prompt/index.tsx b/packages/opencode/src/cli/cmd/tui/component/prompt/index.tsx index 0566e07b3..64b7181e6 100644 --- a/packages/opencode/src/cli/cmd/tui/component/prompt/index.tsx +++ b/packages/opencode/src/cli/cmd/tui/component/prompt/index.tsx @@ -1464,11 +1464,13 @@ export function Prompt(props: PromptProps) { }), } }) + const maxHeight = createMemo(() => tuiConfig.prompt?.max_height ?? Math.max(6, Math.floor(dimensions().height / 3))) return ( <> - (anchor = r)} visible={props.visible !== false}> + (anchor = r)} visible={props.visible !== false} width="100%">