diff --git a/apps/jetstream/package.json b/apps/jetstream/package.json index 04860e9..709773c 100644 --- a/apps/jetstream/package.json +++ b/apps/jetstream/package.json @@ -13,6 +13,7 @@ }, "dependencies": { "@andrioid/atproto": "workspace:*", + "@toondepauw/node-zstd": "^1.2.0", "bufferutil": "^4.0.9", "ws": "^8.18.0" }, diff --git a/bun.lockb b/bun.lockb index 5f8bca809ff75b1304c8537763453f23c3cfac4e..614c5c09b6d473d0fd90eb6c3a2c3056a1f60f2d 100755 GIT binary patch delta 8458 zcmew{S!TjUnF)HDixR^m+7x`6FC_lH!1dC4EqBVKe0i(AV(Mn<3De|ijTt1^89<4l$dHF@D z3=F9ismX~93=9`i85x8b7#a?zGBWToFf?pWWn|!CU}#tkr5C174pLF$lF5K*=gZjK zq%uv}W2fVhjao)4k1km8I@^5i)Xj0XgAQ+seO_-E+^+Tb(c{f>p*r>FMCFeLh%y#g zPV2Y*_51mYGR|kRlI^E;7r8FTS;`usE}N_Lq*=mayMzbh0cp;PI7S8=28ITf$qUQO z8Q)BoOkuQVHHc?qaGQMB&xR#|ks)RBT|aBq&k2kSabQ+pA|pcvh{d@xk&z*dfuVtM z`a%mv3zj5Ch6Io}Ykd+ULlT(vFo}^N2F$WehML5*FnM}Z3Zot8tz9bN9?KrJd85v9% z7#dh7D@K@e=BGk*v4YgFK2Bw1@S8k$oi(Rn8Y6=f149GTBix!N9=8z`#%qr5P9)wp;fxMkD0+EsWVKw*d)2327_X;0NLOP zM*F5$_rH@)iu!CnDa`!Vi57{5lsoI3v-=MhmX@vln)%~L&5O{LjxvWHui;dh=JYjTOq|ArYN z*1RseeTq9|bXNObaoc1Pt8gHBnUd{>*Uz9>rc@Hx02Kr zzoZv#)f2-H-`cnEiPT=Rg$HKx`80E<)yRh$Wpu1%;87NuT4`6Y3}kLSH10}}4PJ4{ zIVHWmZ>HD(pt_Io)}`&)`xFe#?j{}RzvS!`c&Fs^lj+gBAEbxKXBc@p6@NPxx!%~% zW9FpzBL{dNY~7Xd6J#(00|&^+QV;=iaBB39b1otcDmLdMz03EXPB)U^sbBQt)8z*{ znrrl1Cf4kp78Ua9_CG%V1!_eyZ+W-x+|+c4?~T=?d*RQzo%eCGsXz^GU|?Vr!)?F#>XV#yTIKts?v*Z1>`qiapRlsVK&W94J9n6yXX6Btf*Mixbv2yxLkU)eDGw*a? zOJ;Az>gh8rnZu`Z*fR4@{|geRo$hJH96o&qNMNlMvo~Y?^q(Ms3_E7t>9N+#-i(dY zS6VZNPnWP~=9$i8!|cP@JU!BeIehvFkiZ_0KYH%J9thP7qWq_?pT%6UeS!nCN5r&n<}K4FI54wJzcYhbhP5q) zk%4Esqa*Wr#_9XZSq!&@Gpm7FM$-?Jvlwg(XFkC=e2GP1uYx>1@V@LCL~fyPXs}@Q z`LuB6>7XPDPKbtF5{-}?Z8B;4?{XF!*MI*Z091I~Vg}bT3=$v?0|Nt$530`l>` z#S78`m11Clc#MyMf#E1rjFAQ63VsF#h9gih7N|NRv>@47h4Z2w|v{ zECU0BJt&eG7#KvLVsZ=&43ljMXf{wnk7XRR#tI6{wUkR7{P5fngssMVLUv)EO8Uwy-cT zfKrg@^gxgiW>6_j1_p*rP$_e$m=*&A!*Qq)7Em#51_p)%76t}Tg0q~y5M+ckR7w}* zM`%dfK*jVJ7#QY56PhhlOrL>)K@BQqH~k~X2z#iMAp--0EHu6ypkhW03=Gmt44@7i zgCkVTn1O*o6{^l@x?>HCu&pyx%9MeD;WSjr1uAC7z`$^Zi2>C6WN?LwnKLjjoP>(G zO`ixd!W}AQ$-uzy25N){RLqKjf#Ee&%o8eR&A`AgjfsH)l$^b$KLi=!4VAJ5H7=pb zeV}4?3=9kfObiU5lu%$bu7ZRDNrMP z85kHmI2gd$Efp%}$H2hgz`?)(%5G^;F@FXI21lq7>C+2AMr1&x0zshf;U|@L6z`$U@z`(E;>asKj z1_nhY1_n?TTn|;34pj%rf*Ym>f{Xx_5un7D2`c!Q7#Kj=2vnqi(nS^n1A{XY1E?p@ zuoY@VHUk4g6B7diC_`+Uz7S-@cBoV?)Cf?b-T@WMV_;y2Vq#zbCF-3}v3#gH7f|Uw z{UgYT-B77Qs1zts?}3UHF)%Q^0y&3)fnhIHteAm;fd?wKZ@Ob6i?Ho}XndD4Ffgz| zr4B)rmw~E9XaYYB6)R_8U|@!-J2HJD$cST5sY+0d2~`d%FF~rR7#JA-LB&o$)m1Yv zFmQk*K}FB>hae+PL#1jN7#KLAQfHtps{@s%P_gsSSORzcpkf!N8#b{B+g^f7H8Ll zGmEh86R1=#sILUgl~18!eGCi?El{y%P_cfHvFr>CA)xFwy%1!?3ustO1Qjpr3=E(Q z{tBvm5~z5AioJ%4O@?O22vCVX{UFGQw@|663=9m{Sr`~%7#JAdLB*zljDTjT_t5Z} z4o&o+4Do5YU<-?|?PsX+nb6b;O4MJVVzZ#iL5cb+RBSd>tN>K>Oiu(E@f|8Pmw|!d z2Q;#NK*i>PN-{Mp{Z^HsJxu+*v2Ak z3u?WAOyA1Dz~IWvzyL}n%24IoK%Rq&ftqL_b=#rw4N4!X( zsgs}r4OG`KFfh15#ZEzMB#{5zp`JSp6$5$FbGl*&i?FRXRQXv31_mug28II+3=BR{ zv2)Py0Xg0mYQ%YH_<+Rxrx${Z2!MLzBGe;C7#J8rpdPsd^$18T9BS-ks8t}3M@&Bm zG9nTxb(MjEL7S0*0pyw}sMs~Aqd_)DLtSy5fq?-O`4>Rr*wY0&S%j_PpvrGTjR4sa z4;8xw6}tgy6hg%kCV_Gk1H+T$)8jf>Rxt)|XYFFS%;*Xl&Q9Q9U;vGkCvh+^By%t@ zq;N1Wq;fDYq;Wuo&m%Y(7{WLh7@|2C7$P|s7@{~B7{WOi7-BdW7(zK17(%xH?Pi(H zVhS3*F5zHcDCJ;aDC2-M`?5F~81gw77&18+7;-rn81gt67&15*7;-om7z#KT7}B?& zoxrl2nXzbl=MFfdGJVqnN+fefdC#_B-hfS^G}&(8a>Q04izQ zSr{1VSQr>Kf{Gp%28J3I28Ioc3=CB)3=GvQ3=E*gAE;5-%)-F1kCB0)8&u4&Ft{^- zdiV`23=H)w3=9Vu85r7F7#Kj!!44J%22k(#A2S03hz)8Lf-uOTAR3f;EKWf@`D)qp zR|iI@TOoROXZ zC@f&Jc4gs4vb{Barg>>r%unC%__lkA$5AkY*rr`KIlZ=p~Me7UM&XG zm>A;>^o;b37#L(Srhk~t>cb?HG2Lzs>nv`I<&fgsVlC_R*BtCp)Bnz4)!eQzmo?;ow_NE4BQM14Vfoi zue1UrubZA%im3*q1ZqZNQLdQ@b`=mqzzT31gRCgGv^dAX1zYgwrWK{8CKacEoP^IZ zu(wN#bEYdUVU@MT5|T(h1KEZWun?ypjA&p}wZv{2l50>*KoXy>*ubWWBd$Lhvx!K- zqaNxcNH~EU1=i|_-DZdzA!^g}N>OY_lLH%LjoldJ$U`v$SpsZ;Ep`JSc7ehQyDOkF zU?a@28v(Vd!psDl3o6V^z?$u`Yer4wC(f`JFxoRo$4}QvXSCxCiDzVRWMF7uo4hgCoM}`1^i}DMb{sqj zj0{E$3=PvKS}_|;Uf{{WX$KWyp1jfDoT(sTdR7Lb9p}~rMg|iGh6dKj8zam)UqaQe zg4D3aCNeVkO@3Qx&AB>}k->?9p@C`gM{je^w~34l4q)dL`j|88CNWNaU1`gin*?zr z+vGwYbIz4Xj0{#_9U$K0Bt`}^h%J8RoRZ0m3>ILykKX1?fyvXgvKj4Imn1VX*iP0> zv*vsVVlgl@FibAoY|bQ?vi(&yV~~m<8v_GFFb4wzBLf3N2nPcL!*tsTjN+P{3=9l$ zP~muxFarYvHkxz#-U*E2i6GsXP-WO?kbD+Y9z?S+Ffimn`5+o3pAY4OXb``EgMopG zfq|h2L~p;F!x*nj9oNQ!OvmO*P#Oe@gXBT31o1&MNv>_*xs`GI&aF&uTBq-|WA`wv%)HZW zU6{QY>!)|RFo#b+0TOr%5@?)m>B<~Fy~2r^XZjo~W*^38P*M(`uHnqgv)#vynI&@i zg+OMO=^P7~Wmp*!7#VoBC)O~pXWX8!m^qwrdctC6Ll86SHkb(}xEU=fKmo+S5YjpQ zK?9487bw^m{{4pl5a%K@Bz;MMI1CI7Fg_@D6`*35plX;I7#NhHVwa~UHnIrYUV%!n zGB7X%K&)c83Kat--yo>iHD(3|kakdE5&;#vIejC@h}%%*pkyrsW#3_DU=U|uVBlh4 zV2}jqV_;zT43*<%U|`?_Igf#X;mdT!CKloPZ%`>-kQS)acc>U20|UcRsMrsvEBF~0 z7>+>2enZs>FfcIehl>4yiU~3>Fl0f+{zAos7#J9Gp<@4-!P$ThRMs)1K?ND6FKl8F zwq<021gt0n1H(>epfN$k#26SDc0t9Mp<*EK?uLr7O#cWnf|UhqJcA?y1A`4T2eLzz zOEEAo*g?fOpkmSt3=FnVG0y3Z%`C#UTu>=l1_lOuP$V%hFmOY~O6z`(ExDy0V%(_&y?I1V*J zA1bEJz`&5e!oUDZa0b&GK}HxsrF21lgod;+R7{V7fnh#0p_xF%^cff!)SzOf(=URI zFoQ}NGB7a6LgU*UDrUsMz#z@U0IKI0ETCe>3=9maP<58m6p7B?hF;PV_;w?U}9hZrF0jlm^}jnLpfB;b-G|1 zi?FR53j+fv&>a~V7%G_<7(_tzK2*6ADD;^a7(j{N8!G0^z`zg*waRCDBFG3osFW)M z1A`wEB$@d`#oQPe7)+QL7(fLRsAK^Z{O$}447Z>$7dU+*$cQkga!-&0p;F;cF;Hn& z0TlxkIv~rv85kInpm`)}I%7MFux&I{xi13)g9ir#IJ?C_#rzl;7#uhl7(m%A7Aoe? zz`)=LH6m_$AjpV#s8k>*^g&4ufT6Xt39VjD(7vm@e4OB5ZpSD%Hrqz`zKVIt?|piGhKE6)JWHYHTwD z1H(h8*!k&+AR{h7mA8TtI4E0z8Z1!dZJ^W%vXy~>;R;k;JE)}$6}vipBglwrP^nH( zX~n_7;KRVca2+bv#lXPe!ok1*%HTJkV%?x71vInYoX*(8B5ZpbD%H!tz)%Lwm3N?G zeW0ua6}t-+>jxRj&cG1Dz`$^CdLYP%`_QnM2r6FK85lqr{2^5NBvA1J6?+5~n+(m4 z5uk)VeIdw*Cs3)WP!Gf~FfcrYicJF<0nJj+py4wen&?3p;^p*@AR}HumCuBxPEewL z4HcUORSrtjZ=hncp<)FL3=D6lJNB{&+rEQJ&1GO<_yLWq_fWBUppuM{fuV$ff#CzR zvX~Dl$)I9ircdkzHQt~OT*$z{Fa>JFPpA=#Ks6U50|O}K{(_1vW?*2L0af=KDz*ev zw?M`IK*g3aFffEOGBAKr^IxdgG6n{QTBy2z5V1IhX;Rv_AFG4 zi4{`OtpuewMg|5@dS_#WwC+|hFfa%*g1T-D4D3*Ks~H#=(ij;SK#7`rdSf4puq{7S z`C3rJlaYY|lz;`GV(S@WiZgD4{d11Po(rYC}oFos6fQBVt(k%0je zY35L4k3o$9iCIC7IL^Spkio#f0E!vw=^H^t*g&ODf(kT{>p|%qDs~E#Q5hH*K>oLb zhQ(>97|4?j(-|kS2-`YAm7ir`V9;V@U^oB@aj4iikf#_K7(kA9ff{ih8a^N~x9Nc( zBix}Lxd`>h5e5bZZ>UEuK|KNz^Me|D8EO^CS&P7K~Ps*XJB9eMg9fQxWx326Iq0f61<&p zD$8sZQ&8tWn}dNNhl7D3mjlwMOXOf+NaJ8&NZ?>#Na0{$NacWxj3jd~Fr;%ZFvM-Y zIFn^JGh^oV-nlFhT%7ZvjyM2vgqk3zm|xAvz|g|NzyN9^HEma1!J@8MU%R1>UN?8~fHi8N_76yi576yh5j0_BgEDQ`q zEDQ{wW*eySRK>!;u#b^}p&L}Hu`n=z`sd{=3=Cx~3=9WB1r7^noQZ*I1@bsJwpZt+r;S&OIanD zyb`C+TFUAp(*YfqJCyi=$E(F)8WUrjfu50`5d*`tGSg0*rqe`uy2{}@R5~uyS5;^4lCCOsKy`c(-XzmRfXf8vYz_V vrP<@a0Wq+EgB|YtD6R!i(FYtL)oRl%#MvFVHJ}5b8Jq~E?9=CpvwsEvSvC>2 diff --git a/packages/atproto/domain/jetstream-subscription.ts b/packages/atproto/domain/jetstream-subscription.ts index 6b81bbb..dea2aaf 100644 --- a/packages/atproto/domain/jetstream-subscription.ts +++ b/packages/atproto/domain/jetstream-subscription.ts @@ -80,5 +80,5 @@ export async function listenForPosts(ctx: AtContext) { }); let lastStarted = new Date(); - await ctx.db.$client.listen(LISTEN_NOTIFY_NEW_SUBSCRIBERS, handleNotify); + //await ctx.db.$client.listen(LISTEN_NOTIFY_NEW_SUBSCRIBERS, handleNotify); } diff --git a/packages/atproto/worker/tasks/scrape-external-url.ts b/packages/atproto/worker/tasks/scrape-external-url.ts index f77dd4b..49d01dd 100644 --- a/packages/atproto/worker/tasks/scrape-external-url.ts +++ b/packages/atproto/worker/tasks/scrape-external-url.ts @@ -69,6 +69,10 @@ export default async function scrapeExternalUrlTask( // Run papeer against the URL + if (!truncatedScraped || !payload.postId) { + console.warn("scraper received empty text or postId"); + return; + } // Store MAX_LENGTH of the received text in externalTable await ctx.db .insert(postTexts) diff --git a/packages/jetstream/jetstream.ts b/packages/jetstream/jetstream.ts index 49ff10d..1b6495f 100644 --- a/packages/jetstream/jetstream.ts +++ b/packages/jetstream/jetstream.ts @@ -1,5 +1,5 @@ import type { AppBskyFeedPost } from "@atproto/api"; -import zstd from "@bokuweb/zstd-wasm"; +import { Decoder } from "@toondepauw/node-zstd"; import fs from "node:fs"; import path from "node:path"; import { WebSocket, type RawData } from "ws"; @@ -21,8 +21,8 @@ const JETSTREAM_BASE_URL = "wss://jetstream2.us-east.bsky.network/subscribe"; export class Jetstream { private connections: Array = []; private decoder = new TextDecoder(); - private zDict = Uint8Array.from( - fs.readFileSync(path.join(import.meta.dirname, "./zstd_dictionary.dat")) + private zDict = fs.readFileSync( + path.join(import.meta.dirname, "./zstd_dictionary.dat") ); private args: JetStreamRequest; private listeners: Array = []; @@ -33,7 +33,6 @@ export class Jetstream { } static async Create(args: JetStreamRequest) { - await zstd.init(); const instance = new Jetstream(args); return instance; } @@ -42,11 +41,8 @@ export class Jetstream { if (buffer.length === 0) { throw new Error("Received an empty buffer"); } - const uncompressed = zstd.decompressUsingDict( - zstd.createDCtx(), - Uint8Array.from(buffer), - this.zDict - ); + const zDecoder = new Decoder(this.zDict); + const uncompressed = await zDecoder.decode(buffer); if (uncompressed.length === 0) throw new Error("Empty message"); const decoded = this.decoder.decode(uncompressed);