diff --git a/deno.jsonc b/deno.jsonc index a579f65..d646c9d 100644 --- a/deno.jsonc +++ b/deno.jsonc @@ -22,6 +22,7 @@ "@imagemagick/magick-wasm": "npm:@imagemagick/magick-wasm@^0.0.40", "@kysely/kysely": "jsr:@kysely/kysely@^0.28.16", "@preact/signals": "npm:@preact/signals@^2.9.0", + "@std/async": "jsr:@std/async@^1.3.0", "@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 0b33140..f201030 100644 --- a/deno.lock +++ b/deno.lock @@ -9,8 +9,10 @@ "jsr:@fresh/core@^2.3.3": "2.3.3", "jsr:@fresh/plugin-vite@^1.1.2": "1.1.2", "jsr:@kysely/kysely@~0.28.16": "0.28.16", + "jsr:@std/async@^1.3.0": "1.3.0", "jsr:@std/bytes@^1.0.6": "1.0.6", "jsr:@std/cli@^1.0.29": "1.0.29", + "jsr:@std/data-structures@^1.0.11": "1.0.11", "jsr:@std/dotenv@~0.225.5": "0.225.6", "jsr:@std/encoding@^1.0.10": "1.0.10", "jsr:@std/fmt@^1.0.10": "1.0.10", @@ -138,12 +140,21 @@ "@kysely/kysely@0.28.16": { "integrity": "b620bf5d8dc892cb30419aaaef62f0f1ba1846a30eaa45f882ec20183b88a18b" }, + "@std/async@1.3.0": { + "integrity": "80485538a4f7baaa46bfe2246168069e02ed142b9f9079cd164f43bb060ad9e9", + "dependencies": [ + "jsr:@std/data-structures" + ] + }, "@std/bytes@1.0.6": { "integrity": "f6ac6adbd8ccd99314045f5703e23af0a68d7f7e58364b47d2c7f408aeb5820a" }, "@std/cli@1.0.29": { "integrity": "fa4ef29130baa834d8a13b7d138240c3a2fcfba740bfb7afa646a360a15ec84f" }, + "@std/data-structures@1.0.11": { + "integrity": "53b98ed7efa61f107dfc14244bd2ec5557f7f7ee0bbaef6d449d7937facacb89" + }, "@std/dotenv@0.225.6": { "integrity": "1d6f9db72f565bd26790fa034c26e45ecb260b5245417be76c2279e5734c421b" }, @@ -1757,6 +1768,7 @@ "jsr:@fresh/core@^2.3.3", "jsr:@fresh/plugin-vite@^1.1.2", "jsr:@kysely/kysely@~0.28.16", + "jsr:@std/async@^1.3.0", "jsr:@std/encoding@^1.0.10", "jsr:@std/http@^1.1.0", "npm:@atcute/atproto@^3.1.11", diff --git a/src/service/jestream.ts b/src/service/jestream.ts index 0168aea..5daec74 100644 --- a/src/service/jestream.ts +++ b/src/service/jestream.ts @@ -6,6 +6,7 @@ import { } from "@atcute/bluesky" import { JetstreamSubscription } from "@atcute/jetstream" import { type Did, is } from "@atcute/lexicons" +import { debounce } from "@std/async" import { buildAtProtoUri } from "../lib/atProto.ts" import { config } from "../lib/config.ts" import { uploadImages } from "../lib/image.ts" @@ -20,19 +21,19 @@ import { client } from "../traq/client.gen.ts" import { postMessage } from "../traq/index.ts" import { MessageBuilder } from "./messageBuilder.ts" +const UPDATE_CURSOR_DEBOUNCE_MS = 30000 + client.setConfig({ baseUrl: `${config.traqBaseUrl}/api/v3`, }) export class JetstreamService { private subscription: JetstreamSubscription - private cursor?: number private readonly subscribingDids: Did[] private stopResolve?: () => void private loopPromise?: Promise constructor(opts: { wantedDids: Did[]; cursor?: number }) { - this.cursor = opts.cursor this.subscribingDids = opts.wantedDids this.subscription = new JetstreamSubscription({ url: [ @@ -47,6 +48,11 @@ export class JetstreamService { }) } + private updateCursor = debounce( + saveJetstreamCursor, + UPDATE_CURSOR_DEBOUNCE_MS, + ) + private async runLoop(): Promise { const stopPromise = new Promise((resolve) => { this.stopResolve = resolve @@ -73,7 +79,8 @@ export class JetstreamService { const event = result.value if (event.kind !== "commit" || event.commit.operation !== "create") { - this.cursor = event.time_us + console.debug("Skipping non-create commit event", event.did) + this.updateCursor(event.time_us) continue } @@ -81,8 +88,7 @@ export class JetstreamService { try { if (!is(AppBskyFeedPost.mainSchema, event.commit.record)) { console.warn("Invalid record", event.commit.record) - - this.cursor = event.time_us + this.updateCursor(event.time_us) continue } @@ -102,8 +108,7 @@ export class JetstreamService { console.warn( `Skipping post ${atProtoUri} because it has video or is not a self thread`, ) - - this.cursor = event.time_us + this.updateCursor(event.time_us) continue } @@ -112,8 +117,10 @@ export class JetstreamService { if (traqMessageId) { // This message is already posted to traQ - - this.cursor = event.time_us + console.warn( + `Skipping post ${atProtoUri} because it is already posted to traQ with message ID ${traqMessageId}`, + ) + this.updateCursor(event.time_us) continue } @@ -175,8 +182,8 @@ export class JetstreamService { atProtoUri, traqMessageId: data.id, }) - - this.cursor = event.time_us + console.info(`Posted message for post ${atProtoUri} to traQ`) + this.updateCursor(event.time_us) } catch (err) { console.error("Error handling Jetstream event:", err) } @@ -198,9 +205,6 @@ export class JetstreamService { async close() { this.stopResolve?.() await this.loopPromise - - if (this.cursor) { - await saveJetstreamCursor(this.cursor) - } + this.updateCursor.flush() } }