From 9460a682ba58f8a2fa9e0b36794934fd5cab340e Mon Sep 17 00:00:00 2001 From: Takumi Akimoto <81888693+akimon658@users.noreply.github.com> Date: Mon, 4 May 2026 23:36:23 +0900 Subject: [PATCH] Migrate to `@atcute/jetstream` (#8) * Migrate to `@atcute/jetstream` * Fix type errors * fix: Don't throw if `stopResolve` is undefined * Make `getUserSettingByUserId` return `UserSetting` --- database/posts.ts | 2 +- database/userSettings.ts | 3 +- deno.jsonc | 2 +- deno.lock | 45 ++++++--- repository/user.ts | 19 +++- routes/api/settings.ts | 5 +- service/jestream.ts | 206 +++++++++++++++++++++++++-------------- 7 files changed, 186 insertions(+), 96 deletions(-) diff --git a/database/posts.ts b/database/posts.ts index 4e3410a..30ef44a 100644 --- a/database/posts.ts +++ b/database/posts.ts @@ -1,4 +1,4 @@ -import { Generated } from "@kysely/kysely" +import type { Generated } from "@kysely/kysely" export interface PostsTable { id: Generated diff --git a/database/userSettings.ts b/database/userSettings.ts index f1cdb88..ad46505 100644 --- a/database/userSettings.ts +++ b/database/userSettings.ts @@ -1,8 +1,9 @@ +import type { Did } from "@atcute/lexicons" import type { Generated } from "@kysely/kysely" export interface UserSettingsTable { id: Generated - did: string + did: Did target_channel_id: string user_id: string } diff --git a/deno.jsonc b/deno.jsonc index 42ea333..cdfabd0 100644 --- a/deno.jsonc +++ b/deno.jsonc @@ -11,6 +11,7 @@ "imports": { "@atcute/bluesky": "npm:@atcute/bluesky@^3.3.3", "@atcute/client": "npm:@atcute/client@^4.2.1", + "@atcute/jetstream": "npm:@atcute/jetstream@^1.1.2", "@atcute/lexicons": "npm:@atcute/lexicons@^1.3.0", "@badgateway/oauth2-client": "npm:@badgateway/oauth2-client@^3.3.1", "@fresh/plugin-vite": "jsr:@fresh/plugin-vite@^1.1.2", @@ -18,7 +19,6 @@ "@hey-api/vite-plugin": "npm:@hey-api/vite-plugin@^0.3.1", "@kysely/kysely": "jsr:@kysely/kysely@^0.28.16", "@preact/signals": "npm:@preact/signals@^2.9.0", - "@skyware/jetstream": "npm:@skyware/jetstream@^0.2.5", "@std/encoding": "jsr:@std/encoding@^1.0.10", "@std/http": "jsr:@std/http@^1.1.0", "@subframe7536/iconv-lite": "npm:@subframe7536/iconv-lite@^0.8.1", diff --git a/deno.lock b/deno.lock index 8ae6ce6..fcc779e 100644 --- a/deno.lock +++ b/deno.lock @@ -37,6 +37,8 @@ "jsr:@std/uuid@^1.0.9": "1.1.1", "npm:@atcute/bluesky@^3.3.3": "3.3.3", "npm:@atcute/client@^4.2.1": "4.2.1", + "npm:@atcute/jetstream@*": "1.1.2", + "npm:@atcute/jetstream@^1.1.2": "1.1.2", "npm:@atcute/lexicons@^1.3.0": "1.3.0", "npm:@babel/core@^7.28.0": "7.29.0", "npm:@babel/preset-react@^7.27.1": "7.28.5_@babel+core@7.29.0", @@ -48,7 +50,6 @@ "npm:@preact/signals@^2.9.0": "2.9.0_preact@10.29.1", "npm:@prefresh/vite@^2.4.8": "2.4.12_preact@10.29.1_vite@7.3.2__@types+node@25.6.0_@types+node@25.6.0", "npm:@remix-run/node-fetch-server@0.12": "0.12.0", - "npm:@skyware/jetstream@~0.2.5": "0.2.5", "npm:@subframe7536/iconv-lite@~0.8.1": "0.8.1", "npm:@types/babel__core@^7.20.5": "7.20.5", "npm:esbuild-wasm@~0.25.11": "0.25.12", @@ -234,6 +235,18 @@ "@badrap/valita" ] }, + "@atcute/jetstream@1.1.2": { + "integrity": "sha512-u6p/h2xppp7LE6W/9xErAJ6frfN60s8adZuCKtfAaaBBiiYbb1CfpzN8Uc+2qtJZNorqGvuuDb5572Jmh7yHBQ==", + "dependencies": [ + "@atcute/lexicons", + "@badrap/valita", + "@mary-ext/event-iterator", + "@mary-ext/simple-event-emitter", + "partysocket", + "type-fest", + "yocto-queue" + ] + }, "@atcute/lexicons@1.3.0": { "integrity": "sha512-Eq5y+9onnCXNVUlNiMf31beSXHKqptB7lUo/68YbhlmxdaR7ooywHmahya9goP5AsmlYEA1z+dRPXIDAa9O7cg==", "dependencies": [ @@ -925,6 +938,15 @@ "@lukeed/ms@2.0.2": { "integrity": "sha512-9I2Zn6+NJLfaGoz9jN3lpwDgAYvfGeNYdbAIjJOqzs4Tpc+VU3Jqq4IofSUBKajiDS8k9fZIg18/z13mpk1bsA==" }, + "@mary-ext/event-iterator@1.0.0": { + "integrity": "sha512-l6gCPsWJ8aRCe/s7/oCmero70kDHgIK5m4uJvYgwEYTqVxoBOIXbKr5tnkLqUHEg6mNduB4IWvms3h70Hp9ADQ==", + "dependencies": [ + "yocto-queue" + ] + }, + "@mary-ext/simple-event-emitter@1.0.1": { + "integrity": "sha512-9+VvZisxZ/gSg+JJH7hmXaA8Qj42Qjz3O58RSB+INYc8iLA0icATZxHB9vKbj59ojDGZjO3hCKzMXocx3L0H8w==" + }, "@opentelemetry/api@1.9.1": { "integrity": "sha512-gLyJlPHPZYdAk1JENA9LeHejZe1Ti77/pTeFm/nMXmQH/HFZlcS/O2XJB+L8fkbrNSqhdtlvjBVjxwUYanNH5Q==" }, @@ -1097,16 +1119,6 @@ "os": ["win32"], "cpu": ["x64"] }, - "@skyware/jetstream@0.2.5": { - "integrity": "sha512-fM/zs03DLwqRyzZZJFWN20e76KrdqIp97Tlm8Cek+vxn96+tu5d/fx79V6H85L0QN6HvGiX2l9A8hWFqHvYlOA==", - "dependencies": [ - "@atcute/atproto", - "@atcute/bluesky", - "@atcute/lexicons", - "partysocket", - "tiny-emitter" - ] - }, "@standard-schema/spec@1.1.0": { "integrity": "sha512-l2aFy5jALhniG5HgqrD6jXLi/rUWrKvqN/qJx6yoJsgKhblVd+iqqU4RCXavm/jPityDo5TCvKMnpjKnOriy0w==" }, @@ -1648,9 +1660,6 @@ "sql-escaper@1.3.3": { "integrity": "sha512-BsTCV265VpTp8tm1wyIm1xqQCS+Q9NHx2Sr+WcnUrgLrQ6yiDIvHYJV5gHxsj1lMBy2zm5twLaZao8Jd+S8JJw==" }, - "tiny-emitter@2.1.0": { - "integrity": "sha512-NB6Dk1A9xgQPMoGqC5CVXn123gWyte215ONT5Pp5a0yt4nlEoO1ZWeCwpncaekPHXO60i47ihFnZPiRPjRMq4Q==" - }, "tinyglobby@0.2.16": { "integrity": "sha512-pn99VhoACYR8nFHhxqix+uvsbXineAasWm5ojXoN8xEwK5Kd3/TrhNn1wByuD52UxWRLy8pu+kRMniEi6Eq9Zg==", "dependencies": [ @@ -1658,6 +1667,9 @@ "picomatch@4.0.4" ] }, + "type-fest@4.41.0": { + "integrity": "sha512-TeTSQ6H5YHvpqVwBRcnLDCBnDOHWYu7IvGbHT6N8AOymcr9PJGjc1GTtiWZTYg0NCgYwvnYWEkVChQAr9bjfwA==" + }, "typescript@6.0.3": { "integrity": "sha512-y2TvuxSZPDyQakkFRPZHKFm+KKVqIisdg9/CZwm9ftvKXLP8NRWj38/ODjNbr43SsoXqNuAisEf1GdCxqWcdBw==", "bin": true @@ -1716,6 +1728,9 @@ "yaml@2.8.3": { "integrity": "sha512-AvbaCLOO2Otw/lW5bmh9d/WEdcDFdQp2Z2ZUH3pX9U2ihyUY0nvLv7J6TrWowklRGPYbB/IuIMfYgxaCPg5Bpg==", "bin": true + }, + "yocto-queue@1.2.2": { + "integrity": "sha512-4LCcse/U2MHZ63HAJVE+v71o7yOdIe4cZ70Wpf8D/IyjDKYQLV5GD46B+hSTjJsvV5PztjvHoU580EftxjDZFQ==" } }, "workspace": { @@ -1727,12 +1742,12 @@ "jsr:@std/http@^1.1.0", "npm:@atcute/bluesky@^3.3.3", "npm:@atcute/client@^4.2.1", + "npm:@atcute/jetstream@^1.1.2", "npm:@atcute/lexicons@^1.3.0", "npm:@badgateway/oauth2-client@^3.3.1", "npm:@hey-api/openapi-ts@0.97", "npm:@hey-api/vite-plugin@~0.3.1", "npm:@preact/signals@^2.9.0", - "npm:@skyware/jetstream@~0.2.5", "npm:@subframe7536/iconv-lite@~0.8.1", "npm:jose@^6.2.3", "npm:mysql2@^3.22.3", diff --git a/repository/user.ts b/repository/user.ts index d154a2e..73af5ab 100644 --- a/repository/user.ts +++ b/repository/user.ts @@ -1,3 +1,4 @@ +import type { Did } from "@atcute/lexicons" import { db } from "../database/db.ts" export const getAllDids = async () => { @@ -13,7 +14,7 @@ interface UserSetting { } export const getUserSettingByDid = async ( - did: string, + did: Did, ): Promise => { const result = await db .selectFrom("user_settings") @@ -28,17 +29,27 @@ export const getUserSettingByDid = async ( } } -export const getUserSettingByUserId = async (userId: string) => { - return await db +export const getUserSettingByUserId = async ( + userId: string, +): Promise => { + const result = await db .selectFrom("user_settings") .select(["did", "target_channel_id"]) .where("user_id", "=", userId) .executeTakeFirst() + + return result + ? { + userId, + did: result.did, + targetChannelId: result.target_channel_id, + } + : undefined } export const saveUserSettings = async ( userId: string, - did: string, + did: Did, targetChannelId: string, ) => { await db diff --git a/routes/api/settings.ts b/routes/api/settings.ts index e925770..19edb7b 100644 --- a/routes/api/settings.ts +++ b/routes/api/settings.ts @@ -1,3 +1,4 @@ +import { isDid } from "@atcute/lexicons/syntax" import { define } from "../../lib/define.ts" import { getUserSettingByUserId, @@ -14,9 +15,9 @@ export const handler = define.handlers({ PUT: async (ctx) => { const { did, targetChannelId } = await ctx.req.json() - if (typeof did !== "string" || typeof targetChannelId !== "string") { + if (!isDid(did) || typeof targetChannelId !== "string") { return Response.json( - { error: "did and targetChannelId are required" }, + { error: "Invalid request body" }, { status: 400 }, ) } diff --git a/service/jestream.ts b/service/jestream.ts index 616fb56..77a9c29 100644 --- a/service/jestream.ts +++ b/service/jestream.ts @@ -1,5 +1,6 @@ -import type { AppBskyFeedPost } from "@atcute/bluesky" -import { Jetstream } from "@skyware/jetstream" +import { AppBskyFeedPost } from "@atcute/bluesky" +import { JetstreamSubscription } from "@atcute/jetstream" +import { type Did, is } from "@atcute/lexicons" import { buildAtProtoUri } from "../lib/atProto.ts" import { buildMessageContent } from "../lib/buildMessageContent.ts" import { config } from "../lib/config.ts" @@ -25,85 +26,145 @@ client.setConfig({ }) export class JetstreamService { - private jetstream: Jetstream + private subscription: JetstreamSubscription private cursor?: number - private readonly subscribingDids: string[] + private readonly subscribingDids: Did[] + private stopResolve?: () => void + private loopPromise?: Promise - constructor(opts: { wantedDids: string[]; cursor?: number }) { + constructor(opts: { wantedDids: Did[]; cursor?: number }) { this.cursor = opts.cursor this.subscribingDids = opts.wantedDids - this.jetstream = new Jetstream({ + this.subscription = new JetstreamSubscription({ + url: [ + "wss://jetstream1.us-east.bsky.network", + "wss://jetstream2.us-east.bsky.network", + "wss://jetstream1.us-west.bsky.network", + "wss://jetstream2.us-west.bsky.network", + ], wantedCollections: ["app.bsky.feed.post"], wantedDids: opts.wantedDids, cursor: opts.cursor, }) - this.jetstream.onCreate("app.bsky.feed.post", async (event) => { - const atProtoUri = buildAtProtoUri({ - userDid: event.did, - recordKey: event.commit.rkey, - }) - const hasNonSupportedEmbed = event.commit.record.embed && - !SUPPORTED_EMBED_TYPES.includes(event.commit.record.embed.$type) - const isReply = !!event.commit.record.reply?.parent - - if (hasNonSupportedEmbed || isReply) { - console.warn( - `Skipping post ${atProtoUri} because it has unsupported embed or is a reply.`, - ) - - this.cursor = event.time_us - - return - } - - const traqMessageId = await getTraqMessageIdByAtProtoUri(atProtoUri) - - if (traqMessageId) { - // This message is already posted to traQ - return - } - - const userSetting = await getUserSettingByDid(event.did) - const accessToken = await getUserAccessToken(userSetting.userId) - let imageIds: string[] | undefined - - if (event.commit.record.embed?.$type === "app.bsky.embed.images") { - imageIds = await uploadImages({ - accessToken, - did: event.did, - images: event.commit.record.embed.images, - targetChannelId: userSetting.targetChannelId, - }) + } - console.debug(`Uploaded images for post ${atProtoUri}: ${imageIds}`) - } + private async runLoop(): Promise { + const stopPromise = new Promise((resolve) => { + this.stopResolve = resolve + }) - const { data, error } = await postMessage({ - headers: { - Authorization: `Bearer ${accessToken}`, - }, - path: { - channelId: userSetting.targetChannelId, - }, - body: { - content: buildMessageContent({ - imageIds, - post: event.commit.record, - }), - }, - }) - - if (!data) { - throw new Error("Failed to post message to traQ", { cause: error }) + const iterator = this.subscription[Symbol.asyncIterator]() + + try { + while (true) { + const result = await Promise.race([ + iterator.next(), + stopPromise.then( + (): IteratorReturnResult => ({ + done: true, + value: undefined, + }), + ), + ]) + + if (result.done) { + break + } + + const event = result.value + + if (event.kind !== "commit" || event.commit.operation !== "create") { + this.cursor = event.time_us + + continue + } + + try { + if (!is(AppBskyFeedPost.mainSchema, event.commit.record)) { + console.warn("Invalid record", event.commit.record) + + this.cursor = event.time_us + + continue + } + + const atProtoUri = buildAtProtoUri({ + userDid: event.did, + recordKey: event.commit.rkey, + }) + const hasNonSupportedEmbed = event.commit.record.embed && + !SUPPORTED_EMBED_TYPES.includes(event.commit.record.embed.$type) + const isReply = !!event.commit.record.reply?.parent + + if (hasNonSupportedEmbed || isReply) { + console.warn( + `Skipping post ${atProtoUri} because it has unsupported embed or is a reply.`, + ) + + this.cursor = event.time_us + + continue + } + + const traqMessageId = await getTraqMessageIdByAtProtoUri(atProtoUri) + + if (traqMessageId) { + // This message is already posted to traQ + + this.cursor = event.time_us + + continue + } + + const userSetting = await getUserSettingByDid(event.did) + const accessToken = await getUserAccessToken(userSetting.userId) + let imageIds: string[] | undefined + + if (event.commit.record.embed?.$type === "app.bsky.embed.images") { + imageIds = await uploadImages({ + accessToken, + did: event.did, + images: event.commit.record.embed.images, + targetChannelId: userSetting.targetChannelId, + }) + + console.debug( + `Uploaded images for post ${atProtoUri}: ${imageIds}`, + ) + } + + const { data, error } = await postMessage({ + headers: { + Authorization: `Bearer ${accessToken}`, + }, + path: { + channelId: userSetting.targetChannelId, + }, + body: { + content: buildMessageContent({ + imageIds, + post: event.commit.record, + }), + }, + }) + + if (!data) { + throw new Error("Failed to post message to traQ", { cause: error }) + } + + await savePostMetadata({ + atProtoUri, + traqMessageId: data.id, + }) + + this.cursor = event.time_us + } catch (err) { + console.error("Error handling Jetstream event:", err) + } } - - await savePostMetadata({ - atProtoUri, - traqMessageId: data.id, - }) - - this.cursor = event.time_us - }) + } finally { + await iterator.return() + } } start() { @@ -112,11 +173,12 @@ export class JetstreamService { return } - this.jetstream.start() + this.loopPromise = this.runLoop() } async close() { - this.jetstream.close() + this.stopResolve?.() + await this.loopPromise if (this.cursor) { await saveJetstreamCursor(this.cursor) -- 2.51.2