diff --git a/deno.jsonc b/deno.jsonc index 3f7b192..8e62375 100644 --- a/deno.jsonc +++ b/deno.jsonc @@ -9,6 +9,7 @@ "imports": { "@fresh/plugin-vite": "jsr:@fresh/plugin-vite@^1.1.2", "@kysely/kysely": "jsr:@kysely/kysely@^0.28.16", + "@skyware/jetstream": "npm:@skyware/jetstream@^0.2.5", "fresh": "jsr:@fresh/core@^2.3.3", "mysql2": "npm:mysql2@^3.22.3", "preact": "npm:preact@^10.29.1", diff --git a/deno.lock b/deno.lock index 6079e87..1db8a25 100644 --- a/deno.lock +++ b/deno.lock @@ -35,6 +35,7 @@ "npm:@preact/signals@^2.5.1": "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:@types/babel__core@^7.20.5": "7.20.5", "npm:esbuild-wasm@~0.25.11": "0.25.12", "npm:esbuild@0.25.7": "0.25.7", @@ -177,6 +178,37 @@ } }, "npm": { + "@atcute/atproto@3.1.11": { + "integrity": "sha512-yh+ASvA+iHHQij6UeHEKp2+rwvFvQR8A6/5Dk/xvqDslIikWEFx9VlprNwm/clQIPl2bLuQg+LHS8uY9o5nFTA==", + "dependencies": [ + "@atcute/lexicons" + ] + }, + "@atcute/bluesky@3.3.3": { + "integrity": "sha512-R1VyEmbvIUhub6ONnt7sm6FcRoKJISWf+I0qraXtk59KLmNCsJYMX8fLdIcdv5o1i9XtDg8/Um/wqLrl84IuFg==", + "dependencies": [ + "@atcute/atproto", + "@atcute/lexicons" + ] + }, + "@atcute/lexicons@1.3.0": { + "integrity": "sha512-Eq5y+9onnCXNVUlNiMf31beSXHKqptB7lUo/68YbhlmxdaR7ooywHmahya9goP5AsmlYEA1z+dRPXIDAa9O7cg==", + "dependencies": [ + "@atcute/uint8array", + "@atcute/util-text", + "@standard-schema/spec", + "esm-env" + ] + }, + "@atcute/uint8array@1.1.1": { + "integrity": "sha512-3LsC8XB8TKe9q/5hOA5sFuzGaIFdJZJNewC5OKa3o/eU6+K7JR6see9Zy2JbQERNVnRl11EzbNov1efgLMAs4g==" + }, + "@atcute/util-text@1.3.1": { + "integrity": "sha512-MRgJXkx67znuBXuoAYCJkBZyd3OApL7zZlNf5kXhuoCXcdiu1nblRDycYTADSkym4epBSQWxh26kmI9sewaq6A==", + "dependencies": [ + "unicode-segmenter" + ] + }, "@babel/code-frame@7.29.0": { "integrity": "sha512-9NhCeYjq9+3uxgdtp20LSiJXJvN0FeCtNGpJxuMFZ1Kv3cWUNb6DOhJwUvcVCzKGR66cw4njwM6hrJLqgOwbcw==", "dependencies": [ @@ -947,6 +979,19 @@ "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==" + }, "@types/babel__core@7.20.5": { "integrity": "sha512-qoQprZvz5wQFJwMDqeseRXWv3rqMvhgpbXFfVyWhbx9X47POIA6i/+dXefEmZKoAgOaTdaIgNSMqMIU61yRyzA==", "dependencies": [ @@ -1127,9 +1172,15 @@ "escalade@3.2.0": { "integrity": "sha512-WUj2qlxaQtO4g6Pq5c29GTcWGDyd8itL8zTlipgECz3JesAiiOKotd8JU6otB3PACgG6xkJUyVhboMS+bje/jA==" }, + "esm-env@1.2.2": { + "integrity": "sha512-Epxrv+Nr/CaL4ZcFGPJIYLWFom+YeV1DqMLHJoEd9SYRxNbaFruBwfEX/kkHUJf55j2+TUbmDcmuilbP1TmXHA==" + }, "estree-walker@2.0.2": { "integrity": "sha512-Rfkk/Mp/DL7JVje3u18FxFujQlTNR2q6QfMSMB7AvCBx91NGj/ba3kCfza0f6dVDbw7YlRf/nDrn7pQrCCyQ/w==" }, + "event-target-polyfill@0.0.4": { + "integrity": "sha512-Gs6RLjzlLRdT8X9ZipJdIZI/Y6/HhRLyq9RdDlCsnpxr/+Nn6bU2EFGuC94GjxqhM+Nmij2Vcq98yoHrU8uNFQ==" + }, "fdir@6.5.0_picomatch@4.0.4": { "integrity": "sha512-tIbYtZbucOs0BRGqPJkshJUYdL+SDH7dVM8gjy+ERp3WAUjLEFJE+02kanyHtwjWOnwrKYBiwAmM0p4kLJAnXg==", "dependencies": [ @@ -1215,6 +1266,12 @@ "node-releases@2.0.38": { "integrity": "sha512-3qT/88Y3FbH/Kx4szpQQ4HzUbVrHPKTLVpVocKiLfoYvw9XSGOX2FmD2d6DrXbVYyAQTF2HeF6My8jmzx7/CRw==" }, + "partysocket@1.1.18": { + "integrity": "sha512-SyuvH9VavWOSa14v6dYdp3yfSUDII4BQB1+TkGOFBkjfZKjnDBiba4fhdhwBlqGBkqw4ea3gTA1DYhSffX24Wg==", + "dependencies": [ + "event-target-polyfill" + ] + }, "picocolors@1.1.1": { "integrity": "sha512-xceH2snhtb5M9liqDsmEw56le376mTZkEX/jEb/RxNFyegNul7eNslCXP9FDj/Lcu0X8KEyMceP2ntpaHrDEVA==" }, @@ -1289,6 +1346,9 @@ "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": [ @@ -1299,6 +1359,9 @@ "undici-types@7.19.2": { "integrity": "sha512-qYVnV5OEm2AW8cJMCpdV20CDyaN3g0AjDlOGf1OW4iaDEx8MwdtChUp4zu4H0VP3nDRF/8RKWH+IPp9uW0YGZg==" }, + "unicode-segmenter@0.14.5": { + "integrity": "sha512-jHGmj2LUuqDcX3hqY12Ql+uhUTn8huuxNZGq7GvtF6bSybzH3aFgedYu/KTzQStEgt1Ra2F3HxadNXsNjb3m3g==" + }, "update-browserslist-db@1.2.3_browserslist@4.28.2": { "integrity": "sha512-Js0m9cx+qOgDxo0eMiFGEueWztz+d4+M3rGlmKPT+T4IS/jP4ylw3Nwpu6cpTTP8R1MAC1kF4VbdLt3ARf209w==", "dependencies": [ @@ -1336,6 +1399,7 @@ "jsr:@fresh/core@^2.3.3", "jsr:@fresh/plugin-vite@^1.1.2", "jsr:@kysely/kysely@~0.28.16", + "npm:@skyware/jetstream@~0.2.5", "npm:mysql2@^3.22.3", "npm:preact@^10.29.1", "npm:vite@^7.3.2" diff --git a/lib/atProto.ts b/lib/atProto.ts new file mode 100644 index 0000000..ba3988f --- /dev/null +++ b/lib/atProto.ts @@ -0,0 +1,6 @@ +export const buildAtProtoUri = (params: { + userDid: string + recordKey: string +}): string => { + return `at://${params.userDid}/app.bsky.feed.post/${params.recordKey}` +} diff --git a/main.ts b/main.ts index 911a422..fe1db6e 100644 --- a/main.ts +++ b/main.ts @@ -2,6 +2,9 @@ import { Migrator } from "@kysely/kysely" import { App, staticFiles } from "fresh" import { db } from "./database/db.ts" import { migrationProvider } from "./database/migrator.ts" +import { getJetstreamCursor } from "./repository/systemState.ts" +import { getAllDids } from "./repository/user.ts" +import { JetstreamService } from "./service/jestream.ts" const migrator = new Migrator({ db, provider: migrationProvider }) const { error } = await migrator.migrateToLatest() @@ -13,3 +16,17 @@ if (error) { console.log("Database migrated") export const app = new App().use(staticFiles()).fsRoutes() +const jetstreamService = new JetstreamService({ + wantedDids: await getAllDids(), + cursor: await getJetstreamCursor(), +}) + +jetstreamService.start() + +Deno.addSignalListener("SIGTERM", async () => { + console.log("SIGTERM received, shutting down...") + + await jetstreamService.close() + + console.log("Shutdown complete") +}) diff --git a/repository/post.ts b/repository/post.ts new file mode 100644 index 0000000..2b4365f --- /dev/null +++ b/repository/post.ts @@ -0,0 +1,14 @@ +export const getTraqMessageIdByAtProtoUri = ( + atProtoUri: string, +): Promise => { + // TODO + return Promise.resolve(undefined) +} + +export const savePostMetadata = (data: { + atProtoUri: string + traqMessageId: string +}): Promise => { + // TODO + return Promise.resolve() +}) diff --git a/repository/systemState.ts b/repository/systemState.ts new file mode 100644 index 0000000..c0ebe91 --- /dev/null +++ b/repository/systemState.ts @@ -0,0 +1,9 @@ +export const getJetstreamCursor = (): Promise => { + // TODO + return Promise.resolve(undefined) +} + +export const saveJetstreamCursor = (cursor: number): Promise => { + // TODO + return Promise.resolve() +} diff --git a/repository/user.ts b/repository/user.ts new file mode 100644 index 0000000..e18fc9d --- /dev/null +++ b/repository/user.ts @@ -0,0 +1,22 @@ +import { UserSettingsTable } from "../database/userSettings.ts" + +export const getAllDids = (): Promise => { + // TODO + return Promise.resolve([]) +} + +export const getUserSettingByDid = ( + did: string, +): Promise => { + // TODO + return Promise.resolve({ + did, + targetChannelId: "", + userId: "", + }) +} + +export const getUserAccessToken = (userId: string) => { + // TODO + return Promise.resolve("") +} diff --git a/service/jestream.ts b/service/jestream.ts new file mode 100644 index 0000000..f2d192e --- /dev/null +++ b/service/jestream.ts @@ -0,0 +1,60 @@ +import { Jetstream } from "@skyware/jetstream" +import { buildAtProtoUri } from "../lib/atProto.ts" +import { + getTraqMessageIdByAtProtoUri, + savePostMetadata, +} from "../repository/post.ts" +import { saveJetstreamCursor } from "../repository/systemState.ts" +import { getUserAccessToken, getUserSettingByDid } from "../repository/user.ts" +import { postMessage } from "./traq.ts" + +export class JetstreamService { + private jetstream: Jetstream + private cursor?: number + + constructor(opts: { wantedDids: string[]; cursor?: number }) { + this.cursor = opts.cursor + this.jetstream = new Jetstream({ + 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 traqMessageId = await getTraqMessageIdByAtProtoUri(atProtoUri) + + if (traqMessageId) { + // This message is already posted to traQ + return + } + + const userSetting = await getUserSettingByDid(event.did) + const messageId = await postMessage({ + token: await getUserAccessToken(userSetting.userId), + channelId: userSetting.targetChannelId, + content: event.commit.record.text, + }) + + await savePostMetadata({ + atProtoUri, + traqMessageId: messageId, + }) + + this.cursor = event.time_us + }) + } + + start() { + this.jetstream.start() + } + + async close() { + this.jetstream.close() + + if (this.cursor) { + await saveJetstreamCursor(this.cursor) + } + } +} diff --git a/service/traq.ts b/service/traq.ts new file mode 100644 index 0000000..dcd6bf6 --- /dev/null +++ b/service/traq.ts @@ -0,0 +1,8 @@ +export const postMessage = async (params: { + token: string + channelId: string + content: string +}): Promise => { + // TODO + return Promise.resolve("") +})