diff --git a/.gitignore b/.gitignore index ad38727..c07192d 100644 --- a/.gitignore +++ b/.gitignore @@ -3,3 +3,4 @@ node_modules/ .zig-cache/ zig-out/ zig-cache/ +.dev.vars diff --git a/README.md b/README.md index 0f6873b..980b8f7 100644 --- a/README.md +++ b/README.md @@ -13,11 +13,11 @@ curl "https://typeahead.waow.tech/xrpc/app.bsky.actor.searchActorsTypeahead?q=na ``` jetstream → ingester (zig, fly.io) → worker (cloudflare) ↓ - D1 (sqlite/FTS5) + Turso (libSQL/FTS5) ``` - **ingester**: [zig](https://ziglang.org) on [fly.io](https://fly.io) — streams identity + profile events via [jetstream](https://docs.bsky.app/blog/jetstream), batches to worker -- **worker**: [cloudflare worker](https://workers.cloudflare.com) + D1 + KV + cache API — FTS5 prefix search, edge-cached (60s), rate-limited +- **worker**: [cloudflare worker](https://workers.cloudflare.com) + [Turso](https://turso.tech) + KV + cache API — FTS5 prefix search, edge-cached (60s), rate-limited - **identity**: [slingshot](https://microcosm.blue) for on-demand handle resolution ## dev diff --git a/docs/architecture.md b/docs/architecture.md index 0525997..b27ad63 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -14,8 +14,8 @@ and serves FTS5 prefix search from cloudflare's edge. v worker (cloudflare) | - +---> D1 (actors table + FTS5 index) - +---> KV (cursor, config flags, mod_cursor) + +---> Turso (actors table + FTS5 index) + +---> KV (cursor, mod_cursor) +---> cache API (60s edge cache for search) ``` @@ -57,15 +57,10 @@ the actors table stores one row per DID: plus FTS5 index overhead. roughly ~280 bytes/row total. -D1 has a 10GB hard limit (non-negotiable). at current row size that's ~35M -actors. D1 is designed for per-tenant databases, not single large datasets — -sharding a global search index across multiple D1s degrades FTS5 ranking -(scores aren't comparable across shards) and fans out every query. - -when we outgrow D1, the natural move is **Turso** (hosted libSQL). our sibling -project [leaflet-search](https://tangled.org/zzstoatzz.io/leaflet-search) -already runs Turso + local SQLite read replica in production for FTS5 search. -the schema and queries would port with minimal changes. +storage is [Turso](https://turso.tech) (hosted libSQL). previously used +Cloudflare D1 but migrated to Turso to avoid D1's 10GB hard limit and +per-tenant design constraints that made sharding a global FTS5 index +impractical. ## moderation @@ -86,3 +81,5 @@ by walking the full index over multiple runs. existing actors (run once after adding moderation support) - `scripts/migrate-avatar-cid.sql` — one-shot migration from full avatar URLs to bare CIDs (already applied to production) +- `scripts/migrate-to-turso.py` — one-shot D1-to-Turso data migration + (already applied to production) diff --git a/package-lock.json b/package-lock.json index e00433b..fc11874 100644 --- a/package-lock.json +++ b/package-lock.json @@ -5,6 +5,9 @@ "packages": { "": { "name": "typeahead", + "dependencies": { + "@libsql/client": "^0.17.0" + }, "devDependencies": { "wrangler": "^4" } @@ -1104,6 +1107,173 @@ "@jridgewell/sourcemap-codec": "^1.4.10" } }, + "node_modules/@libsql/client": { + "version": "0.17.0", + "resolved": "https://registry.npmjs.org/@libsql/client/-/client-0.17.0.tgz", + "integrity": "sha512-TLjSU9Otdpq0SpKHl1tD1Nc9MKhrsZbCFGot3EbCxRa8m1E5R1mMwoOjKMMM31IyF7fr+hPNHLpYfwbMKNusmg==", + "license": "MIT", + "dependencies": { + "@libsql/core": "^0.17.0", + "@libsql/hrana-client": "^0.9.0", + "js-base64": "^3.7.5", + "libsql": "^0.5.22", + "promise-limit": "^2.7.0" + } + }, + "node_modules/@libsql/core": { + "version": "0.17.0", + "resolved": "https://registry.npmjs.org/@libsql/core/-/core-0.17.0.tgz", + "integrity": "sha512-hnZRnJHiS+nrhHKLGYPoJbc78FE903MSDrFJTbftxo+e52X+E0Y0fHOCVYsKWcg6XgB7BbJYUrz/xEkVTSaipw==", + "license": "MIT", + "dependencies": { + "js-base64": "^3.7.5" + } + }, + "node_modules/@libsql/darwin-arm64": { + "version": "0.5.22", + "resolved": "https://registry.npmjs.org/@libsql/darwin-arm64/-/darwin-arm64-0.5.22.tgz", + "integrity": "sha512-4B8ZlX3nIDPndfct7GNe0nI3Yw6ibocEicWdC4fvQbSs/jdq/RC2oCsoJxJ4NzXkvktX70C1J4FcmmoBy069UA==", + "cpu": [ + "arm64" + ], + "license": "MIT", + "optional": true, + "os": [ + "darwin" + ] + }, + "node_modules/@libsql/darwin-x64": { + "version": "0.5.22", + "resolved": "https://registry.npmjs.org/@libsql/darwin-x64/-/darwin-x64-0.5.22.tgz", + "integrity": "sha512-ny2HYWt6lFSIdNFzUFIJ04uiW6finXfMNJ7wypkAD8Pqdm6nAByO+Fdqu8t7sD0sqJGeUCiOg480icjyQ2/8VA==", + "cpu": [ + "x64" + ], + "license": "MIT", + "optional": true, + "os": [ + "darwin" + ] + }, + "node_modules/@libsql/hrana-client": { + "version": "0.9.0", + "resolved": "https://registry.npmjs.org/@libsql/hrana-client/-/hrana-client-0.9.0.tgz", + "integrity": "sha512-pxQ1986AuWfPX4oXzBvLwBnfgKDE5OMhAdR/5cZmRaB4Ygz5MecQybvwZupnRz341r2CtFmbk/BhSu7k2Lm+Jw==", + "license": "MIT", + "dependencies": { + "@libsql/isomorphic-ws": "^0.1.5", + "cross-fetch": "^4.0.0", + "js-base64": "^3.7.5", + "node-fetch": "^3.3.2" + } + }, + "node_modules/@libsql/isomorphic-ws": { + "version": "0.1.5", + "resolved": "https://registry.npmjs.org/@libsql/isomorphic-ws/-/isomorphic-ws-0.1.5.tgz", + "integrity": "sha512-DtLWIH29onUYR00i0GlQ3UdcTRC6EP4u9w/h9LxpUZJWRMARk6dQwZ6Jkd+QdwVpuAOrdxt18v0K2uIYR3fwFg==", + "license": "MIT", + "dependencies": { + "@types/ws": "^8.5.4", + "ws": "^8.13.0" + } + }, + "node_modules/@libsql/linux-arm-gnueabihf": { + "version": "0.5.22", + "resolved": "https://registry.npmjs.org/@libsql/linux-arm-gnueabihf/-/linux-arm-gnueabihf-0.5.22.tgz", + "integrity": "sha512-3Uo3SoDPJe/zBnyZKosziRGtszXaEtv57raWrZIahtQDsjxBVjuzYQinCm9LRCJCUT5t2r5Z5nLDPJi2CwZVoA==", + "cpu": [ + "arm" + ], + "license": "MIT", + "optional": true, + "os": [ + "linux" + ] + }, + "node_modules/@libsql/linux-arm-musleabihf": { + "version": "0.5.22", + "resolved": "https://registry.npmjs.org/@libsql/linux-arm-musleabihf/-/linux-arm-musleabihf-0.5.22.tgz", + "integrity": "sha512-LCsXh07jvSojTNJptT9CowOzwITznD+YFGGW+1XxUr7fS+7/ydUrpDfsMX7UqTqjm7xG17eq86VkWJgHJfvpNg==", + "cpu": [ + "arm" + ], + "license": "MIT", + "optional": true, + "os": [ + "linux" + ] + }, + "node_modules/@libsql/linux-arm64-gnu": { + "version": "0.5.22", + "resolved": "https://registry.npmjs.org/@libsql/linux-arm64-gnu/-/linux-arm64-gnu-0.5.22.tgz", + "integrity": "sha512-KSdnOMy88c9mpOFKUEzPskSaF3VLflfSUCBwas/pn1/sV3pEhtMF6H8VUCd2rsedwoukeeCSEONqX7LLnQwRMA==", + "cpu": [ + "arm64" + ], + "license": "MIT", + "optional": true, + "os": [ + "linux" + ] + }, + "node_modules/@libsql/linux-arm64-musl": { + "version": "0.5.22", + "resolved": "https://registry.npmjs.org/@libsql/linux-arm64-musl/-/linux-arm64-musl-0.5.22.tgz", + "integrity": "sha512-mCHSMAsDTLK5YH//lcV3eFEgiR23Ym0U9oEvgZA0667gqRZg/2px+7LshDvErEKv2XZ8ixzw3p1IrBzLQHGSsw==", + "cpu": [ + "arm64" + ], + "license": "MIT", + "optional": true, + "os": [ + "linux" + ] + }, + "node_modules/@libsql/linux-x64-gnu": { + "version": "0.5.22", + "resolved": "https://registry.npmjs.org/@libsql/linux-x64-gnu/-/linux-x64-gnu-0.5.22.tgz", + "integrity": "sha512-kNBHaIkSg78Y4BqAdgjcR2mBilZXs4HYkAmi58J+4GRwDQZh5fIUWbnQvB9f95DkWUIGVeenqLRFY2pcTmlsew==", + "cpu": [ + "x64" + ], + "license": "MIT", + "optional": true, + "os": [ + "linux" + ] + }, + "node_modules/@libsql/linux-x64-musl": { + "version": "0.5.22", + "resolved": "https://registry.npmjs.org/@libsql/linux-x64-musl/-/linux-x64-musl-0.5.22.tgz", + "integrity": "sha512-UZ4Xdxm4pu3pQXjvfJiyCzZop/9j/eA2JjmhMaAhe3EVLH2g11Fy4fwyUp9sT1QJYR1kpc2JLuybPM0kuXv/Tg==", + "cpu": [ + "x64" + ], + "license": "MIT", + "optional": true, + "os": [ + "linux" + ] + }, + "node_modules/@libsql/win32-x64-msvc": { + "version": "0.5.22", + "resolved": "https://registry.npmjs.org/@libsql/win32-x64-msvc/-/win32-x64-msvc-0.5.22.tgz", + "integrity": "sha512-Fj0j8RnBpo43tVZUVoNK6BV/9AtDUM5S7DF3LB4qTYg1LMSZqi3yeCneUTLJD6XomQJlZzbI4mst89yspVSAnA==", + "cpu": [ + "x64" + ], + "license": "MIT", + "optional": true, + "os": [ + "win32" + ] + }, + "node_modules/@neon-rs/load": { + "version": "0.0.4", + "resolved": "https://registry.npmjs.org/@neon-rs/load/-/load-0.0.4.tgz", + "integrity": "sha512-kTPhdZyTQxB+2wpiRcFWrDcejc4JI6tkPuS7UZCG4l6Zvc5kU/gGQ/ozvHTh1XR5tS+UlfAfGuPajjzQjCiHCw==", + "license": "MIT" + }, "node_modules/@poppinss/colors": { "version": "4.1.6", "resolved": "https://registry.npmjs.org/@poppinss/colors/-/colors-4.1.6.tgz", @@ -1153,6 +1323,24 @@ "dev": true, "license": "CC0-1.0" }, + "node_modules/@types/node": { + "version": "25.5.0", + "resolved": "https://registry.npmjs.org/@types/node/-/node-25.5.0.tgz", + "integrity": "sha512-jp2P3tQMSxWugkCUKLRPVUpGaL5MVFwF8RDuSRztfwgN1wmqJeMSbKlnEtQqU8UrhTmzEmZdu2I6v2dpp7XIxw==", + "license": "MIT", + "dependencies": { + "undici-types": "~7.18.0" + } + }, + "node_modules/@types/ws": { + "version": "8.18.1", + "resolved": "https://registry.npmjs.org/@types/ws/-/ws-8.18.1.tgz", + "integrity": "sha512-ThVF6DCVhA8kUGy+aazFQ4kXQ7E1Ty7A3ypFOe0IcJV8O/M511G99AW24irKrW56Wt44yG9+ij8FaqoBGkuBXg==", + "license": "MIT", + "dependencies": { + "@types/node": "*" + } + }, "node_modules/blake3-wasm": { "version": "2.1.5", "resolved": "https://registry.npmjs.org/blake3-wasm/-/blake3-wasm-2.1.5.tgz", @@ -1174,6 +1362,44 @@ "url": "https://opencollective.com/express" } }, + "node_modules/cross-fetch": { + "version": "4.1.0", + "resolved": "https://registry.npmjs.org/cross-fetch/-/cross-fetch-4.1.0.tgz", + "integrity": "sha512-uKm5PU+MHTootlWEY+mZ4vvXoCn4fLQxT9dSc1sXVMSFkINTJVN8cAQROpwcKm8bJ/c7rgZVIBWzH5T78sNZZw==", + "license": "MIT", + "dependencies": { + "node-fetch": "^2.7.0" + } + }, + "node_modules/cross-fetch/node_modules/node-fetch": { + "version": "2.7.0", + "resolved": "https://registry.npmjs.org/node-fetch/-/node-fetch-2.7.0.tgz", + "integrity": "sha512-c4FRfUm/dbcWZ7U+1Wq0AwCyFL+3nt2bEw05wfxSz+DWpWsitgmSgYmy2dQdWyKC1694ELPqMs/YzUSNozLt8A==", + "license": "MIT", + "dependencies": { + "whatwg-url": "^5.0.0" + }, + "engines": { + "node": "4.x || >=6.0.0" + }, + "peerDependencies": { + "encoding": "^0.1.0" + }, + "peerDependenciesMeta": { + "encoding": { + "optional": true + } + } + }, + "node_modules/data-uri-to-buffer": { + "version": "4.0.1", + "resolved": "https://registry.npmjs.org/data-uri-to-buffer/-/data-uri-to-buffer-4.0.1.tgz", + "integrity": "sha512-0R9ikRb668HB7QDxT1vkpuUBtqc53YyAwMwGeUFKRojY/NWKvdZ+9UYtRfGmhqNbRkTSVpMbmyhXipFFv2cb/A==", + "license": "MIT", + "engines": { + "node": ">= 12" + } + }, "node_modules/detect-libc": { "version": "2.1.2", "resolved": "https://registry.npmjs.org/detect-libc/-/detect-libc-2.1.2.tgz", @@ -1236,6 +1462,41 @@ "@esbuild/win32-x64": "0.27.3" } }, + "node_modules/fetch-blob": { + "version": "3.2.0", + "resolved": "https://registry.npmjs.org/fetch-blob/-/fetch-blob-3.2.0.tgz", + "integrity": "sha512-7yAQpD2UMJzLi1Dqv7qFYnPbaPx7ZfFK6PiIxQ4PfkGPyNyl2Ugx+a/umUonmKqjhM4DnfbMvdX6otXq83soQQ==", + "funding": [ + { + "type": "github", + "url": "https://github.com/sponsors/jimmywarting" + }, + { + "type": "paypal", + "url": "https://paypal.me/jimmywarting" + } + ], + "license": "MIT", + "dependencies": { + "node-domexception": "^1.0.0", + "web-streams-polyfill": "^3.0.3" + }, + "engines": { + "node": "^12.20 || >= 14.13" + } + }, + "node_modules/formdata-polyfill": { + "version": "4.0.10", + "resolved": "https://registry.npmjs.org/formdata-polyfill/-/formdata-polyfill-4.0.10.tgz", + "integrity": "sha512-buewHzMvYL29jdeQTVILecSaZKnt/RJWjoZCF5OW60Z67/GmSLBkOFM7qh1PI3zFNtJbaZL5eQu1vLfazOwj4g==", + "license": "MIT", + "dependencies": { + "fetch-blob": "^3.1.2" + }, + "engines": { + "node": ">=12.20.0" + } + }, "node_modules/fsevents": { "version": "2.3.3", "resolved": "https://registry.npmjs.org/fsevents/-/fsevents-2.3.3.tgz", @@ -1251,6 +1512,12 @@ "node": "^8.16.0 || ^10.6.0 || >=11.0.0" } }, + "node_modules/js-base64": { + "version": "3.7.8", + "resolved": "https://registry.npmjs.org/js-base64/-/js-base64-3.7.8.tgz", + "integrity": "sha512-hNngCeKxIUQiEUN3GPJOkz4wF/YvdUdbNL9hsBcMQTkKzboD7T/q3OYOuuPZLUE6dBxSGpwhk5mwuDud7JVAow==", + "license": "BSD-3-Clause" + }, "node_modules/kleur": { "version": "4.1.5", "resolved": "https://registry.npmjs.org/kleur/-/kleur-4.1.5.tgz", @@ -1261,6 +1528,47 @@ "node": ">=6" } }, + "node_modules/libsql": { + "version": "0.5.22", + "resolved": "https://registry.npmjs.org/libsql/-/libsql-0.5.22.tgz", + "integrity": "sha512-NscWthMQt7fpU8lqd7LXMvT9pi+KhhmTHAJWUB/Lj6MWa0MKFv0F2V4C6WKKpjCVZl0VwcDz4nOI3CyaT1DDiA==", + "cpu": [ + "x64", + "arm64", + "wasm32", + "arm" + ], + "license": "MIT", + "os": [ + "darwin", + "linux", + "win32" + ], + "dependencies": { + "@neon-rs/load": "^0.0.4", + "detect-libc": "2.0.2" + }, + "optionalDependencies": { + "@libsql/darwin-arm64": "0.5.22", + "@libsql/darwin-x64": "0.5.22", + "@libsql/linux-arm-gnueabihf": "0.5.22", + "@libsql/linux-arm-musleabihf": "0.5.22", + "@libsql/linux-arm64-gnu": "0.5.22", + "@libsql/linux-arm64-musl": "0.5.22", + "@libsql/linux-x64-gnu": "0.5.22", + "@libsql/linux-x64-musl": "0.5.22", + "@libsql/win32-x64-msvc": "0.5.22" + } + }, + "node_modules/libsql/node_modules/detect-libc": { + "version": "2.0.2", + "resolved": "https://registry.npmjs.org/detect-libc/-/detect-libc-2.0.2.tgz", + "integrity": "sha512-UX6sGumvvqSaXgdKGUsgZWqcUyIXZ/vZTrlRT/iobiKhGL0zL4d3osHj3uqllWJK+i+sixDS/3COVEOFbupFyw==", + "license": "Apache-2.0", + "engines": { + "node": ">=8" + } + }, "node_modules/miniflare": { "version": "4.20260312.1", "resolved": "https://registry.npmjs.org/miniflare/-/miniflare-4.20260312.1.tgz", @@ -1282,6 +1590,44 @@ "node": ">=18.0.0" } }, + "node_modules/node-domexception": { + "version": "1.0.0", + "resolved": "https://registry.npmjs.org/node-domexception/-/node-domexception-1.0.0.tgz", + "integrity": "sha512-/jKZoMpw0F8GRwl4/eLROPA3cfcXtLApP0QzLmUT/HuPCZWyB7IY9ZrMeKw2O/nFIqPQB3PVM9aYm0F312AXDQ==", + "deprecated": "Use your platform's native DOMException instead", + "funding": [ + { + "type": "github", + "url": "https://github.com/sponsors/jimmywarting" + }, + { + "type": "github", + "url": "https://paypal.me/jimmywarting" + } + ], + "license": "MIT", + "engines": { + "node": ">=10.5.0" + } + }, + "node_modules/node-fetch": { + "version": "3.3.2", + "resolved": "https://registry.npmjs.org/node-fetch/-/node-fetch-3.3.2.tgz", + "integrity": "sha512-dRB78srN/l6gqWulah9SrxeYnxeddIG30+GOqK/9OlLVyLg3HPnr6SqOWTWOXKRwC2eGYCkZ59NNuSgvSrpgOA==", + "license": "MIT", + "dependencies": { + "data-uri-to-buffer": "^4.0.0", + "fetch-blob": "^3.1.4", + "formdata-polyfill": "^4.0.10" + }, + "engines": { + "node": "^12.20.0 || ^14.13.1 || >=16.0.0" + }, + "funding": { + "type": "opencollective", + "url": "https://opencollective.com/node-fetch" + } + }, "node_modules/path-to-regexp": { "version": "6.3.0", "resolved": "https://registry.npmjs.org/path-to-regexp/-/path-to-regexp-6.3.0.tgz", @@ -1296,6 +1642,12 @@ "dev": true, "license": "MIT" }, + "node_modules/promise-limit": { + "version": "2.7.0", + "resolved": "https://registry.npmjs.org/promise-limit/-/promise-limit-2.7.0.tgz", + "integrity": "sha512-7nJ6v5lnJsXwGprnGXga4wx6d1POjvi5Qmf1ivTRxTjH4Z/9Czja/UCMLVmB9N93GeWOU93XaFaEt6jbuoagNw==", + "license": "ISC" + }, "node_modules/semver": { "version": "7.7.4", "resolved": "https://registry.npmjs.org/semver/-/semver-7.7.4.tgz", @@ -1367,6 +1719,12 @@ "url": "https://github.com/chalk/supports-color?sponsor=1" } }, + "node_modules/tr46": { + "version": "0.0.3", + "resolved": "https://registry.npmjs.org/tr46/-/tr46-0.0.3.tgz", + "integrity": "sha512-N3WMsuqV66lT30CrXNbEjx4GEwlow3v6rr4mCcv6prnfwhS01rkgyFdjPNBYd9br7LpXV1+Emh01fHnq2Gdgrw==", + "license": "MIT" + }, "node_modules/tslib": { "version": "2.8.1", "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz", @@ -1385,6 +1743,12 @@ "node": ">=20.18.1" } }, + "node_modules/undici-types": { + "version": "7.18.2", + "resolved": "https://registry.npmjs.org/undici-types/-/undici-types-7.18.2.tgz", + "integrity": "sha512-AsuCzffGHJybSaRrmr5eHr81mwJU3kjw6M+uprWvCXiNeN9SOGwQ3Jn8jb8m3Z6izVgknn1R0FTCEAP2QrLY/w==", + "license": "MIT" + }, "node_modules/unenv": { "version": "2.0.0-rc.24", "resolved": "https://registry.npmjs.org/unenv/-/unenv-2.0.0-rc.24.tgz", @@ -1395,6 +1759,31 @@ "pathe": "^2.0.3" } }, + "node_modules/web-streams-polyfill": { + "version": "3.3.3", + "resolved": "https://registry.npmjs.org/web-streams-polyfill/-/web-streams-polyfill-3.3.3.tgz", + "integrity": "sha512-d2JWLCivmZYTSIoge9MsgFCZrt571BikcWGYkjC1khllbTeDlGqZ2D8vD8E/lJa8WGWbb7Plm8/XJYV7IJHZZw==", + "license": "MIT", + "engines": { + "node": ">= 8" + } + }, + "node_modules/webidl-conversions": { + "version": "3.0.1", + "resolved": "https://registry.npmjs.org/webidl-conversions/-/webidl-conversions-3.0.1.tgz", + "integrity": "sha512-2JAn3z8AR6rjK8Sm8orRC0h/bcl/DqL7tRPdGZ4I1CjdF+EaMLmYxBHyXuKL849eucPFhvBoxMsflfOb8kxaeQ==", + "license": "BSD-2-Clause" + }, + "node_modules/whatwg-url": { + "version": "5.0.0", + "resolved": "https://registry.npmjs.org/whatwg-url/-/whatwg-url-5.0.0.tgz", + "integrity": "sha512-saE57nupxk6v3HY35+jzBwYa0rKSy0XR8JSxZPwgLr7ys0IBzhGviA1/TUGJLmSVqs8pb9AnvICXEuOHLprYTw==", + "license": "MIT", + "dependencies": { + "tr46": "~0.0.3", + "webidl-conversions": "^3.0.0" + } + }, "node_modules/workerd": { "version": "1.20260312.1", "resolved": "https://registry.npmjs.org/workerd/-/workerd-1.20260312.1.tgz", @@ -1455,7 +1844,6 @@ "version": "8.18.0", "resolved": "https://registry.npmjs.org/ws/-/ws-8.18.0.tgz", "integrity": "sha512-8VbfWfHLbbwu3+N6OKsOMpBdT4kXPDDB9cJk2bJ6mh9ucxdlnNvH1e+roYkKmN9Nxw2yjz7VzeO9oOz2zJ04Pw==", - "dev": true, "license": "MIT", "engines": { "node": ">=10.0.0" diff --git a/package.json b/package.json index bc1578c..e80a093 100644 --- a/package.json +++ b/package.json @@ -7,5 +7,8 @@ }, "devDependencies": { "wrangler": "^4" + }, + "dependencies": { + "@libsql/client": "^0.17.0" } } diff --git a/scripts/migrate-to-turso.py b/scripts/migrate-to-turso.py new file mode 100755 index 0000000..f8f2cbf --- /dev/null +++ b/scripts/migrate-to-turso.py @@ -0,0 +1,347 @@ +#!/usr/bin/env -S PYTHONUNBUFFERED=1 uv run --script --quiet +# /// script +# requires-python = ">=3.12" +# dependencies = [] +# /// +""" +one-shot: migrate D1 data to Turso. + +reads D1 via Cloudflare REST API (fast), writes to Turso via pipeline API. +uses ON CONFLICT upserts so re-running is safe (idempotent). + +prerequisites: + turso db create typeahead + turso db shell typeahead < schema.sql + +usage: + TURSO_URL=libsql://... TURSO_AUTH_TOKEN=... ./scripts/migrate-to-turso.py + TURSO_URL=libsql://... TURSO_AUTH_TOKEN=... ./scripts/migrate-to-turso.py --verify-only +""" + +import argparse +import json +import os +import re +import subprocess +import sys +import urllib.request + +PAGE_SIZE = 1000 +TURSO_BATCH_SIZE = 200 # rows per Turso pipeline request + +PASS = "\033[32m✓\033[0m" +FAIL = "\033[31m✗\033[0m" +DIM = "\033[2m" +RESET = "\033[0m" + +# D1 config (from wrangler.jsonc) +CF_ACCOUNT_ID = "8feb33b5fb57ce2bc093bc6f4141f40a" +CF_D1_DB_ID = "7e289d5d-dc50-46d1-8084-49aeec2679e5" +D1_API = f"https://api.cloudflare.com/client/v4/accounts/{CF_ACCOUNT_ID}/d1/database/{CF_D1_DB_ID}/query" + +_ANSI_RE = re.compile(r"\x1b\[[0-9;]*m") + + +def get_cf_token() -> str: + """read wrangler's OAuth token from its config file.""" + config_path = os.path.expanduser( + "~/Library/Preferences/.wrangler/config/default.toml" + ) + try: + with open(config_path) as f: + for line in f: + if line.startswith("oauth_token"): + return line.split("=", 1)[1].strip().strip('"') + except FileNotFoundError: + pass + # fallback: try CLOUDFLARE_API_TOKEN env var + token = os.environ.get("CLOUDFLARE_API_TOKEN", "") + if token: + return token + print("error: no Cloudflare API token found", file=sys.stderr) + print(" run `wrangler login` or set CLOUDFLARE_API_TOKEN", file=sys.stderr) + sys.exit(1) + + +def get_turso_url() -> str: + url = os.environ.get("TURSO_URL", "") + if not url: + print("error: TURSO_URL not set", file=sys.stderr) + sys.exit(1) + return url.replace("libsql://", "https://") + + +def get_turso_token() -> str: + token = os.environ.get("TURSO_AUTH_TOKEN", "") + if not token: + print("error: TURSO_AUTH_TOKEN not set", file=sys.stderr) + sys.exit(1) + return token + + +def d1_query(sql: str, cf_token: str, params: list | None = None) -> list[dict]: + """query D1 via Cloudflare REST API.""" + payload: dict = {"sql": sql} + if params: + payload["params"] = params + body = json.dumps(payload).encode() + req = urllib.request.Request( + D1_API, + data=body, + headers={ + "Authorization": f"Bearer {cf_token}", + "Content-Type": "application/json", + }, + ) + try: + with urllib.request.urlopen(req, timeout=30) as resp: + data = json.loads(resp.read()) + if data.get("success"): + return data["result"][0]["results"] + print(f" D1 API error: {data.get('errors')}", file=sys.stderr) + return [] + except urllib.error.HTTPError as e: + body_text = e.read().decode()[:300] + print(f" D1 HTTP {e.code}: {body_text}", file=sys.stderr) + return [] + except Exception as e: + print(f" D1 request failed: {e}", file=sys.stderr) + return [] + + +def d1_query_wrangler(sql: str) -> list[dict]: + """fallback: query D1 via wrangler CLI.""" + result = subprocess.run( + [ + "npx", "wrangler", "d1", "execute", "typeahead-db", + "--remote", "--command", sql, "--json", + ], + capture_output=True, text=True, cwd=".", + ) + stdout = _ANSI_RE.sub("", result.stdout) + bracket = stdout.find("[") + if bracket == -1: + return [] + try: + data = json.loads(stdout[bracket:]) + return data[0]["results"] if data else [] + except (json.JSONDecodeError, IndexError, KeyError): + return [] + + +def turso_batch(stmts: list[dict], turso_url: str, turso_token: str) -> bool: + """execute a batch of statements against Turso via HTTP pipeline API.""" + requests = [{"type": "execute", "stmt": s} for s in stmts] + requests.append({"type": "close"}) + + body = json.dumps({"requests": requests}).encode() + req = urllib.request.Request( + f"{turso_url}/v3/pipeline", + data=body, + headers={ + "Authorization": f"Bearer {turso_token}", + "Content-Type": "application/json", + }, + ) + try: + with urllib.request.urlopen(req, timeout=60) as resp: + result = json.loads(resp.read()) + for r in result.get("results", []): + if r.get("type") == "error": + print(f" Turso error: {r.get('error', {}).get('message', 'unknown')}") + return False + return True + except urllib.error.HTTPError as e: + err_body = e.read().decode()[:300] + print(f" Turso HTTP {e.code}: {err_body}", file=sys.stderr) + return False + except Exception as e: + print(f" Turso request failed: {e}", file=sys.stderr) + return False + + +def turso_count(table: str, turso_url: str, turso_token: str) -> int | str: + """get row count from Turso.""" + body = json.dumps({ + "requests": [ + {"type": "execute", "stmt": {"sql": f"SELECT COUNT(*) AS cnt FROM {table}", "args": []}}, + {"type": "close"}, + ] + }).encode() + req = urllib.request.Request( + f"{turso_url}/v3/pipeline", + data=body, + headers={ + "Authorization": f"Bearer {turso_token}", + "Content-Type": "application/json", + }, + ) + try: + with urllib.request.urlopen(req, timeout=15) as resp: + result = json.loads(resp.read()) + return int(result["results"][0]["response"]["result"]["rows"][0][0]["value"]) + except Exception as e: + return f"error: {e}" + + +def progress(msg: str): + """overwrite current line with progress.""" + sys.stdout.write(f"\r {msg}") + sys.stdout.flush() + + +def migrate_actors(turso_url: str, turso_token: str, cf_token: str) -> int: + print("\n--- actors ---") + + # get total for progress reporting + count_rows = d1_query("SELECT COUNT(*) AS cnt FROM actors", cf_token) + d1_total = count_rows[0]["cnt"] if count_rows else "?" + print(f" D1 has {d1_total} actors") + + cursor = 0 + total = 0 + + while True: + rows = d1_query( + "SELECT rowid, did, handle, display_name, avatar_url, updated_at, hidden " + f"FROM actors WHERE rowid > {cursor} ORDER BY rowid ASC LIMIT {PAGE_SIZE}", + cf_token, + ) + if not rows: + break + + for i in range(0, len(rows), TURSO_BATCH_SIZE): + batch = rows[i : i + TURSO_BATCH_SIZE] + stmts = [] + for r in batch: + stmts.append({ + "sql": ( + "INSERT INTO actors (did, handle, display_name, avatar_url, updated_at, hidden) " + "VALUES (?, ?, ?, ?, ?, ?) " + "ON CONFLICT(did) DO UPDATE SET " + "handle = COALESCE(NULLIF(excluded.handle, ''), actors.handle), " + "display_name = COALESCE(NULLIF(excluded.display_name, ''), actors.display_name), " + "avatar_url = COALESCE(NULLIF(excluded.avatar_url, ''), actors.avatar_url), " + "updated_at = excluded.updated_at, hidden = excluded.hidden" + ), + "args": [ + {"type": "text", "value": r["did"]}, + {"type": "text", "value": r.get("handle") or ""}, + {"type": "text", "value": r.get("display_name") or ""}, + {"type": "text", "value": r.get("avatar_url") or ""}, + {"type": "integer", "value": str(r.get("updated_at") or 0)}, + {"type": "integer", "value": str(r.get("hidden") or 0)}, + ], + }) + + if not turso_batch(stmts, turso_url, turso_token): + print(f"\n batch failed at cursor={cursor}") + return total + total += len(batch) + + cursor = rows[-1]["rowid"] + pct = f" ({total * 100 // d1_total}%)" if isinstance(d1_total, int) else "" + progress(f"{total}/{d1_total} actors{pct} {DIM}cursor={cursor}{RESET}") + + print(f"\r {PASS} actors: {total} rows" + " " * 30) + return total + + +def migrate_table( + table: str, + columns: list[str], + col_types: list[str], + turso_url: str, + turso_token: str, + cf_token: str, +) -> int: + print(f"\n--- {table} ---") + col_list = ", ".join(columns) + rows = d1_query(f"SELECT {col_list} FROM {table}", cf_token) + + if not rows: + print(f" {PASS} {table}: 0 rows (empty)") + return 0 + + total = 0 + for i in range(0, len(rows), TURSO_BATCH_SIZE): + batch = rows[i : i + TURSO_BATCH_SIZE] + stmts = [] + placeholders = ", ".join("?" for _ in columns) + sql = f"INSERT OR REPLACE INTO {table} ({col_list}) VALUES ({placeholders})" + + for r in batch: + args = [] + for col, ctype in zip(columns, col_types): + val = r.get(col) or 0 + if ctype == "text": + args.append({"type": "text", "value": str(val)}) + elif ctype == "float": + args.append({"type": "float", "value": float(val)}) + else: + args.append({"type": "integer", "value": str(int(val))}) + stmts.append({"sql": sql, "args": args}) + + if not turso_batch(stmts, turso_url, turso_token): + print(f" batch failed") + return total + total += len(batch) + progress(f"{total}/{len(rows)} {table}") + + print(f"\r {PASS} {table}: {total} rows" + " " * 20) + return total + + +def verify_counts(turso_url: str, turso_token: str, cf_token: str): + print("\n--- verification ---") + tables = ["actors", "metrics", "snapshots"] + for table in tables: + d1_rows = d1_query(f"SELECT COUNT(*) AS cnt FROM {table}", cf_token) + d1_count = d1_rows[0]["cnt"] if d1_rows else "?" + turso_cnt = turso_count(table, turso_url, turso_token) + match = str(d1_count) == str(turso_cnt) + tag = PASS if match else FAIL + print(f" [{tag}] {table}: D1={d1_count}, Turso={turso_cnt}") + + +def main(): + parser = argparse.ArgumentParser(description="migrate D1 → Turso") + parser.add_argument("--verify-only", action="store_true", help="only compare row counts") + args = parser.parse_args() + + cf_token = get_cf_token() + turso_url = get_turso_url() + turso_token = get_turso_token() + + # quick API check + test = d1_query("SELECT 1 AS ok", cf_token) + if not test: + print("error: D1 API connection failed", file=sys.stderr) + sys.exit(1) + + if args.verify_only: + verify_counts(turso_url, turso_token, cf_token) + return + + print("migrating D1 → Turso") + + migrate_actors(turso_url, turso_token, cf_token) + migrate_table( + "metrics", + ["hour", "searches", "total_ms"], + ["integer", "integer", "float"], + turso_url, turso_token, cf_token, + ) + migrate_table( + "snapshots", + ["hour", "total", "with_handles", "with_avatars"], + ["integer", "integer", "integer", "integer"], + turso_url, turso_token, cf_token, + ) + + verify_counts(turso_url, turso_token, cf_token) + print("\ndone.") + + +if __name__ == "__main__": + main() diff --git a/scripts/smoke.py b/scripts/smoke.py index c81361f..59602f1 100755 --- a/scripts/smoke.py +++ b/scripts/smoke.py @@ -212,19 +212,16 @@ def test_moderation_filtering(base_url: str): extra_keys |= set(a.keys()) - allowed_keys check("actor objects have clean shape", len(extra_keys) == 0, f"extra keys: {extra_keys}" if extra_keys else "") - # verify hidden actors are actually excluded by finding one with !no-unauthenticated - # via bsky API and checking it doesn't appear in our results - print("\n--- hidden actor exclusion ---") - found_hidden = find_hidden_actor(base_url) - if not found_hidden: - check("found a hidden actor to verify", False, "couldn't find one — skipping exclusion check") + # !no-unauthenticated actors should be VISIBLE (it's about content, not identity) + print("\n--- !no-unauthenticated inclusion ---") + found = find_noauth_actor(base_url) + if not found: + check("found a !no-unauthenticated actor to verify", False, "couldn't find one — skipping") -def find_hidden_actor(base_url: str) -> bool: - """find an actor we've indexed that has !no-unauthenticated, verify they're excluded.""" - # search our index for common names and cross-check labels via bsky API +def find_noauth_actor(base_url: str) -> bool: + """find an actor with !no-unauthenticated and verify they ARE included in our results.""" for q in ["alex", "sam", "chris", "jordan"]: - # get actors from bsky that have !no-unauthenticated bsky_data, _ = fetch(f"{BSKY_PUBLIC}/xrpc/app.bsky.actor.searchActors?q={q}&limit=25") if not bsky_data or "_error" in bsky_data: continue @@ -242,17 +239,17 @@ def find_hidden_actor(base_url: str) -> bool: if not handle: continue - # this actor has !no-unauthenticated — check they're NOT in our results + # this actor has !no-unauthenticated — check they ARE in our results our_data, _ = fetch(f"{base_url}{XRPC_PATH}?q={handle}&limit=10") if not our_data or "_error" in our_data: continue our_handles = {a.get("handle") for a in our_data.get("actors", [])} if handle in our_handles: - check(f"hidden actor @{handle} excluded from search", False, "appeared in results") + check(f"!no-unauthenticated actor @{handle} visible in search", True) return True else: - check(f"hidden actor @{handle} excluded from search", True) + check(f"!no-unauthenticated actor @{handle} visible in search", False, "not found in results") return True return False diff --git a/src/index.ts b/src/index.ts index ae03139..63f34d9 100644 --- a/src/index.ts +++ b/src/index.ts @@ -1,9 +1,60 @@ +import { createClient, type Client } from "@libsql/client/web"; + interface Env { - DB: D1Database; KV: KVNamespace; ADMIN_SECRET: string; RATE_LIMITER: RateLimit; RATE_LIMITER_STRICT: RateLimit; + TURSO_URL: string; + TURSO_AUTH_TOKEN: string; +} + +interface Stmt { + bind(...args: unknown[]): Stmt; + all>(): Promise<{ results: T[] }>; + first>(): Promise; + run(): Promise<{ meta: { changes: number } }>; +} + +interface TursoDB { + prepare(sql: string): Stmt; + batch(stmts: Stmt[]): Promise<{ results: unknown[]; meta: { changes: number } }[]>; +} + +function tursoDb(client: Client): TursoDB { + return { + prepare(sql) { + let args: unknown[] = []; + const s: Stmt & { _sql: string; _args: () => unknown[] } = { + _sql: sql, + _args: () => args, + bind(...a) { args = a; return s; }, + async all() { + const r = await client.execute({ sql, args: args as any }); + return { results: r.rows as unknown as T[] }; + }, + async first() { + const r = await client.execute({ sql, args: args as any }); + return (r.rows[0] as unknown as T) ?? null; + }, + async run() { + const r = await client.execute({ sql, args: args as any }); + return { meta: { changes: r.rowsAffected } }; + }, + }; + return s; + }, + async batch(stmts) { + const results = await client.batch( + stmts.map((s) => ({ sql: (s as any)._sql as string, args: (s as any)._args() as any[] })), + "write", + ); + return results.map((r) => ({ + results: r.rows as unknown[], + meta: { changes: r.rowsAffected }, + })); + }, + }; } const CORS_HEADERS = { @@ -69,18 +120,15 @@ const BSKY_TYPEAHEAD_URL = "https://public.api.bsky.app/xrpc/app.bsky.actor.searchActorsTypeahead"; const BSKY_MOD_DID = "did:plc:ar7c4by46qjdydhdevvrndac"; -/** labels from bluesky's moderation service that hide an actor */ +/** labels from bluesky's moderation service that hide an actor from search */ const MOD_HIDE_VALS = new Set(["!hide", "!takedown", "spam"]); -/** labels that hide regardless of issuer (protocol-level, self-labeling respected) */ -const ANY_SRC_HIDE_VALS = new Set(["!no-unauthenticated"]); /** - * true if actor should be hidden from our unauthenticated search. + * returns whether an actor should be hidden from search. * - * two paths: - * 1. bluesky moderation issued !hide or spam → always hide - * 2. anyone (including the actor themselves) issued !no-unauthenticated → hide, - * because our service is unauthenticated and we should respect the user's intent + * only hides actors flagged by bluesky's moderation service (!hide, !takedown, spam). + * !no-unauthenticated is intentionally NOT filtered — it applies to content, not identity. + * bluesky's own public typeahead API returns !no-unauthenticated accounts, and so do we. */ function shouldHide(labels?: any[]): boolean { if (!labels) return false; @@ -88,9 +136,7 @@ function shouldHide(labels?: any[]): boolean { return labels.some((l: any) => { if (l.neg) return false; if (l.exp && new Date(l.exp).getTime() <= now) return false; - if (l.src === BSKY_MOD_DID && MOD_HIDE_VALS.has(l.val)) return true; - if (ANY_SRC_HIDE_VALS.has(l.val)) return true; - return false; + return l.src === BSKY_MOD_DID && MOD_HIDE_VALS.has(l.val); }); } @@ -99,7 +145,7 @@ function shouldHide(labels?: any[]): boolean { async function backfillFromBsky( term: string, limit: number, - env: Env + db: TursoDB, ): Promise { try { const res = await fetch( @@ -114,7 +160,7 @@ async function backfillFromBsky( // upsert all — fills in missing actors AND enriches existing ones // (e.g. actors ingested via Jetstream that lack avatar/displayName) const stmts = actors.map((a) => - env.DB.prepare( + db.prepare( `INSERT INTO actors (did, handle, display_name, avatar_url, hidden, updated_at) VALUES (?1, ?2, ?3, ?4, ?5, unixepoch()) ON CONFLICT(did) DO UPDATE SET @@ -132,14 +178,14 @@ async function backfillFromBsky( ) ); - await env.DB.batch(stmts); + await db.batch(stmts); console.log(JSON.stringify({ event: "backfill", term, upserted: actors.length })); } catch { // best-effort — don't let backfill errors affect anything } } -async function throttledBackfill(term: string, limit: number, env: Env): Promise { +async function throttledBackfill(term: string, limit: number, db: TursoDB, env: Env): Promise { // kill switch — set KV key "backfill" to "off" to disable without redeploying const flag = await env.KV.get("backfill"); if (flag === "off") return; @@ -151,22 +197,22 @@ async function throttledBackfill(term: string, limit: number, env: Env): Promise return; } - return backfillFromBsky(term, limit, env); + return backfillFromBsky(term, limit, db); } // --- end backfill --- /** record an actor-count snapshot for the current hour (idempotent) */ -async function recordSnapshot(env: Env): Promise { +async function recordSnapshot(db: TursoDB): Promise { const hour = Math.floor(Date.now() / 3_600_000); - const row = await env.DB.prepare( + const row = await db.prepare( `SELECT COUNT(*) AS total, SUM(CASE WHEN handle != '' THEN 1 ELSE 0 END) AS with_handles, SUM(CASE WHEN avatar_url != '' THEN 1 ELSE 0 END) AS with_avatars FROM actors WHERE hidden = 0` ).first<{ total: number; with_handles: number; with_avatars: number }>(); if (row) { - await env.DB.prepare( + await db.prepare( `INSERT OR REPLACE INTO snapshots (hour, total, with_handles, with_avatars) VALUES (?1, ?2, ?3, ?4)` ) @@ -176,8 +222,8 @@ async function recordSnapshot(env: Env): Promise { } /** resolve handles for actors missing them via slingshot */ -async function resolveHandles(env: Env): Promise { - const { results } = await env.DB.prepare( +async function resolveHandles(db: TursoDB): Promise { + const { results } = await db.prepare( "SELECT did FROM actors WHERE handle = '' ORDER BY updated_at DESC LIMIT 1000" ).all<{ did: string }>(); if (!results || results.length === 0) return; @@ -191,7 +237,7 @@ async function resolveHandles(env: Env): Promise { if (!res.ok) continue; const identity: SlingshotResponse = await res.json(); if (identity.handle) { - await env.DB.prepare( + await db.prepare( "UPDATE actors SET handle = ?1 WHERE did = ?2 AND handle = ''" ).bind(identity.handle, did).run(); resolved++; @@ -209,12 +255,12 @@ const BSKY_GET_PROFILES_URL = "https://public.api.bsky.app/xrpc/app.bsky.actor.getProfiles"; /** refresh moderation labels, walking the full index over multiple cron runs */ -async function refreshModeration(env: Env): Promise { +async function refreshModeration(db: TursoDB, env: Env): Promise { // resume where we left off (rowid cursor persisted in KV) const cursorStr = await env.KV.get("mod_cursor"); const cursor = cursorStr ? Number(cursorStr) : 0; - const { results } = await env.DB.prepare( + const { results } = await db.prepare( "SELECT rowid, did FROM actors WHERE rowid > ?1 ORDER BY rowid ASC LIMIT 1000" ).bind(cursor).all<{ rowid: number; did: string }>(); @@ -249,17 +295,17 @@ async function refreshModeration(env: Env): Promise { const profiles: any[] = data.profiles || []; checked += profiles.length; - const stmts: D1PreparedStatement[] = []; + const stmts: Stmt[] = []; for (const p of profiles) { const hide = shouldHide(p.labels) ? 1 : 0; stmts.push( - env.DB.prepare( + db.prepare( "UPDATE actors SET hidden = ?1 WHERE did = ?2 AND hidden != ?1" ).bind(hide, p.did) ); } if (stmts.length > 0) { - const batchResults = await env.DB.batch(stmts); + const batchResults = await db.batch(stmts); changed += batchResults.filter((r) => r.meta.changes > 0).length; } } catch { @@ -274,9 +320,9 @@ async function refreshModeration(env: Env): Promise { } /** fire-and-forget: increment hourly search count + accumulate response time */ -async function recordMetric(env: Env, ms: number): Promise { +async function recordMetric(db: TursoDB, ms: number): Promise { const hour = Math.floor(Date.now() / 3_600_000); - await env.DB.prepare( + await db.prepare( `INSERT INTO metrics (hour, searches, total_ms) VALUES (?1, 1, ?2) ON CONFLICT(hour) DO UPDATE SET @@ -289,8 +335,9 @@ async function recordMetric(env: Env, ms: number): Promise { async function handleSearch( request: Request, + db: TursoDB, env: Env, - ctx: ExecutionContext + ctx: ExecutionContext, ): Promise { const url = new URL(request.url); const q = url.searchParams.get("q") || url.searchParams.get("term") || ""; @@ -324,7 +371,7 @@ async function handleSearch( const t0 = Date.now(); const ftsQuery = `"${term}"*`; - const { results } = await env.DB.prepare( + const { results } = await db.prepare( `SELECT a.did, a.handle, a.display_name, a.avatar_url FROM actors_fts JOIN actors a ON a.rowid = actors_fts.rowid @@ -346,11 +393,11 @@ async function handleSearch( // --- backfill: remove this block once at parity with Bluesky --- const hasGaps = actors.length < limit || actors.some((a) => !a.avatar); if (hasGaps) { - ctx.waitUntil(throttledBackfill(term, limit, env)); + ctx.waitUntil(throttledBackfill(term, limit, db, env)); } // --- end backfill --- - ctx.waitUntil(recordMetric(env, Date.now() - t0)); + ctx.waitUntil(recordMetric(db, Date.now() - t0)); const response = json({ actors }); @@ -364,7 +411,8 @@ async function handleSearch( async function handleIngest( request: Request, - env: Env + db: TursoDB, + env: Env, ): Promise { const auth = request.headers.get("Authorization"); if (auth !== `Bearer ${env.ADMIN_SECRET}`) { @@ -393,7 +441,7 @@ async function handleIngest( const stmts = events.map((e) => { const avatarCid = e.avatar_cid || null; const hidden = e.hidden !== undefined ? (e.hidden ? 1 : 0) : null; - return env.DB.prepare( + return db.prepare( `INSERT INTO actors (did, handle, display_name, avatar_url, hidden, updated_at) VALUES (?1, ?2, ?3, ?4, COALESCE(?5, 0), unixepoch()) ON CONFLICT(did) DO UPDATE SET @@ -412,7 +460,7 @@ async function handleIngest( }); try { - await env.DB.batch(stmts); + await db.batch(stmts); } catch (e: any) { console.log(JSON.stringify({ event: "ingest_error", error: e?.message, count: events.length })); return json({ error: e?.message || "db batch failed" }, 500); @@ -431,7 +479,8 @@ async function handleIngest( async function handleDelete( request: Request, - env: Env + db: TursoDB, + env: Env, ): Promise { const auth = request.headers.get("Authorization"); if (auth !== `Bearer ${env.ADMIN_SECRET}`) { @@ -457,9 +506,9 @@ async function handleDelete( } const stmts = dids.map((did) => - env.DB.prepare("DELETE FROM actors WHERE did = ?1").bind(did) + db.prepare("DELETE FROM actors WHERE did = ?1").bind(did) ); - await env.DB.batch(stmts); + await db.batch(stmts); return json({ ok: true, deleted: dids.length }); } @@ -480,10 +529,11 @@ async function handleCursor( return json({ cursor: cursor ? Number(cursor) : null }); } -/** resolve a handle or DID via slingshot, then upsert into D1 */ +/** resolve a handle or DID via slingshot, then upsert into the active DB */ async function handleRequestIndexing( request: Request, - env: Env + db: TursoDB, + env: Env, ): Promise { const url = new URL(request.url); const identifier = @@ -523,7 +573,7 @@ async function handleRequestIndexing( // profile enrichment is best-effort } - await env.DB.prepare( + await db.prepare( `INSERT INTO actors (did, handle, display_name, avatar_url, hidden, updated_at) VALUES (?1, ?2, ?3, ?4, ?5, unixepoch()) ON CONFLICT(did) DO UPDATE SET @@ -536,20 +586,24 @@ async function handleRequestIndexing( .bind(identity.did, identity.handle, displayName, avatarCid, hidden ? 1 : 0) .run(); - return json({ handle: identity.handle, did: identity.did, hidden }); + return json({ + handle: identity.handle, + did: identity.did, + ...(hidden ? { hidden: true, reason: "hidden by moderation" } : { hidden: false }), + }); } -async function handleStats(env: Env): Promise { +async function handleStats(db: TursoDB): Promise { const [totalRes, handlesRes, avatarsRes, hiddenRes, metricsRes, snapshotRes] = - await env.DB.batch([ - env.DB.prepare("SELECT COUNT(*) AS cnt FROM actors WHERE hidden = 0"), - env.DB.prepare("SELECT COUNT(*) AS cnt FROM actors WHERE handle != '' AND hidden = 0"), - env.DB.prepare("SELECT COUNT(*) AS cnt FROM actors WHERE avatar_url != '' AND hidden = 0"), - env.DB.prepare("SELECT COUNT(*) AS cnt FROM actors WHERE hidden = 1"), - env.DB.prepare( + await db.batch([ + db.prepare("SELECT COUNT(*) AS cnt FROM actors WHERE hidden = 0"), + db.prepare("SELECT COUNT(*) AS cnt FROM actors WHERE handle != '' AND hidden = 0"), + db.prepare("SELECT COUNT(*) AS cnt FROM actors WHERE avatar_url != '' AND hidden = 0"), + db.prepare("SELECT COUNT(*) AS cnt FROM actors WHERE hidden != 0"), + db.prepare( "SELECT hour, searches, total_ms FROM metrics ORDER BY hour DESC LIMIT 168" ), - env.DB.prepare( + db.prepare( "SELECT hour, total, with_handles, with_avatars FROM snapshots ORDER BY hour ASC LIMIT 2000" ), ]); @@ -730,7 +784,7 @@ function statsPage(d: StatsData): string {
${d.avatarPct}%
-
hidden by moderation
+
hidden by moderation
${d.hiddenCount.toLocaleString()}
@@ -1103,7 +1157,7 @@ function indexPage(): string { if (data.error) { showMsg(esc(data.error), true); } else { - const hidden = data.hidden ? ' (hidden by moderation)' : ''; + const hidden = data.hidden ? ' (' + esc(data.reason || 'hidden') + ')' : ''; showMsg('indexed @' + esc(data.handle) + '' + hidden, false); handleInput.value = ''; } @@ -1128,9 +1182,10 @@ function html(body: string, status = 200): Response { export default { async scheduled(_event: ScheduledEvent, env: Env, _ctx: ExecutionContext): Promise { - await recordSnapshot(env); - await refreshModeration(env); - await resolveHandles(env); + const db = tursoDb(createClient({ url: env.TURSO_URL, authToken: env.TURSO_AUTH_TOKEN })); + await recordSnapshot(db); + await refreshModeration(db, env); + await resolveHandles(db); }, async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise { @@ -1143,25 +1198,27 @@ export default { if (pathname === "/" && request.method === "GET") { return html(indexPage()); } + if (pathname === "/admin/cursor" && request.method === "GET") { + return handleCursor(request, env); + } + if (pathname === "/request-indexing" && request.method === "GET") { + return new Response(null, { status: 302, headers: { Location: "/" } }); + } + + const db = tursoDb(createClient({ url: env.TURSO_URL, authToken: env.TURSO_AUTH_TOKEN })); if (pathname === "/stats" && request.method === "GET") { - return handleStats(env); + return handleStats(db); } - if (pathname === "/request-indexing") { - if (request.method === "GET") { - // old bookmarks / form fallback — redirect to homepage - return new Response(null, { status: 302, headers: { Location: "/" } }); - } - if (request.method === "POST") { - const ip = clientIP(request); - const { success } = await env.RATE_LIMITER.limit({ key: `index:${ip}` }); - if (!success) { - console.log(JSON.stringify({ event: "rate_limited", endpoint: "/request-indexing", ip })); - return json({ error: "slow down — try again in a minute." }, 429); - } - return handleRequestIndexing(request, env); + if (pathname === "/request-indexing" && request.method === "POST") { + const ip = clientIP(request); + const { success } = await env.RATE_LIMITER.limit({ key: `index:${ip}` }); + if (!success) { + console.log(JSON.stringify({ event: "rate_limited", endpoint: "/request-indexing", ip })); + return json({ error: "slow down — try again in a minute." }, 429); } + return handleRequestIndexing(request, db, env); } if ( @@ -1174,19 +1231,15 @@ export default { console.log(JSON.stringify({ event: "rate_limited", endpoint: "/search", ip })); return json({ error: "rate limited" }, 429); } - return handleSearch(request, env, ctx); + return handleSearch(request, db, env, ctx); } if (pathname === "/admin/ingest" && request.method === "POST") { - return handleIngest(request, env); + return handleIngest(request, db, env); } if (pathname === "/admin/delete" && request.method === "POST") { - return handleDelete(request, env); - } - - if (pathname === "/admin/cursor" && request.method === "GET") { - return handleCursor(request, env); + return handleDelete(request, db, env); } return json({ error: "not found" }, 404); diff --git a/wrangler.jsonc b/wrangler.jsonc index 0128438..f9f5a12 100644 --- a/wrangler.jsonc +++ b/wrangler.jsonc @@ -4,13 +4,6 @@ "compatibility_date": "2024-12-01", "compatibility_flags": ["nodejs_compat"], "triggers": { "crons": ["0 * * * *"] }, - "d1_databases": [ - { - "binding": "DB", - "database_name": "typeahead-db", - "database_id": "7e289d5d-dc50-46d1-8084-49aeec2679e5" - } - ], "kv_namespaces": [ { "binding": "KV",