diff --git a/packages/atproto/domain/fetch-missing-post-records.ts b/packages/atproto/domain/fetch-missing-post-records.ts index e0b3dd9..8fd49de 100644 --- a/packages/atproto/domain/fetch-missing-post-records.ts +++ b/packages/atproto/domain/fetch-missing-post-records.ts @@ -34,11 +34,15 @@ export async function fetchMissingPostRecords(ctx: AtContext) { console.warn("[record-fetcher] record does not match feed post"); continue; } + const record = post.record as AppBskyFeedPost.Record; // This update should never be called on a post we dont have await ctx.db.transaction(async (tx) => { await tx .update(postTable) - .set({ flags: postRecordFlags(post.record as AppBskyFeedPost.Record) }) + .set({ + flags: postRecordFlags(record), + created: new Date(record.createdAt), + }) .where(eq(postTable.id, post.uri)); await tx @@ -47,7 +51,7 @@ export async function fetchMissingPostRecords(ctx: AtContext) { cid: post.cid, postId: post.uri, type: "AppBskyFeedPost.Record", - value: post.record as AppBskyFeedPost.Record, + value: record, }) .onConflictDoNothing(); }); diff --git a/packages/atproto/domain/jetstream-subscription.ts b/packages/atproto/domain/jetstream-subscription.ts index 9bdcd9c..f173d24 100644 --- a/packages/atproto/domain/jetstream-subscription.ts +++ b/packages/atproto/domain/jetstream-subscription.ts @@ -10,7 +10,7 @@ import { followTable } from "./user/user-follows.table"; export const LISTEN_NOTIFY_NEW_SUBSCRIBERS = "atproto.subscriber.update"; export async function listenForPosts(ctx: AtContext) { - async function getDids() { + async function getDids(): Promise>> { const wantedDids = await ctx.db .selectDistinct({ did: followTable.follows }) .from(followTable); @@ -20,7 +20,7 @@ export async function listenForPosts(ctx: AtContext) { console.warn("aborting jetstream connection, no dids requested"); return []; } - return dids; + return [...dids]; } const dids = await getDids(); diff --git a/packages/atproto/scripts/classifier.ts b/packages/atproto/scripts/classifier.ts index da79e1e..a88a927 100644 --- a/packages/atproto/scripts/classifier.ts +++ b/packages/atproto/scripts/classifier.ts @@ -102,5 +102,4 @@ export async function classifier() { process.exit(42); } }); - //while (true) {} // listen doesn't wait } diff --git a/packages/jetstream/jetstream.ts b/packages/jetstream/jetstream.ts index 593378b..0b7e7bd 100644 --- a/packages/jetstream/jetstream.ts +++ b/packages/jetstream/jetstream.ts @@ -20,8 +20,8 @@ const JETSTREAM_BASE_URL = "wss://jetstream2.us-east.bsky.network/subscribe"; export class Jetstream { private connections: Array = []; private decoder: TextDecoder = new TextDecoder(); - private zDict: Buffer = fs.readFileSync( - path.join(import.meta.dirname, "./zstd_dictionary.dat") + private zDict = Uint8Array.from( + fs.readFileSync(path.join(import.meta.dirname, "./zstd_dictionary.dat")) ); private args: JetStreamRequest; private listeners: Array = []; @@ -44,7 +44,7 @@ export class Jetstream { const uncompressed = zstd.decompressUsingDict( zstd.createDCtx(), Uint8Array.from(buffer), - Uint8Array.from(this.zDict) + this.zDict ); if (uncompressed.length === 0) throw new Error("Empty message"); const decoded = this.decoder.decode(uncompressed); @@ -54,7 +54,7 @@ export class Jetstream { return data as CommitEvent; } catch (err) { console.log("should be json", decoded); - throw new Error("Failed to pase Jetstream message", { + throw new Error("Failed to parse Jetstream message", { cause: err, }); }