From 4f0f00f465f4cdaff36242b801a8675ae6b31ae2 Mon Sep 17 00:00:00 2001 From: Kristina Date: Wed, 25 Dec 2024 17:15:39 +0400 Subject: [PATCH] Init query (#9) --- bun.lockb | Bin 102896 -> 104400 bytes packages/client-query/package.json | 19 ++ packages/client-query/src/index.ts | 26 +++ packages/client-query/src/query.ts | 65 +++++++ packages/client-query/tsconfig.json | 9 + packages/client-sqlite/src/client.ts | 2 +- packages/query/package.json | 19 ++ packages/query/src/index.ts | 1 + packages/query/src/lq.ts | 150 ++++++++++++++++ packages/query/src/messages/query.ts | 202 ++++++++++++++++++++++ packages/query/src/notifications/query.ts | 129 ++++++++++++++ packages/query/src/query.ts | 198 +++++++++++++++++++++ packages/query/src/result.ts | 92 ++++++++++ packages/query/src/types.ts | 33 ++++ packages/query/src/window.ts | 34 ++++ packages/query/tsconfig.json | 9 + packages/sdk-types/src/client.ts | 1 - packages/sdk-types/src/index.ts | 3 +- packages/sdk-types/src/query.ts | 6 + 19 files changed, 995 insertions(+), 3 deletions(-) create mode 100644 packages/client-query/package.json create mode 100644 packages/client-query/src/index.ts create mode 100644 packages/client-query/src/query.ts create mode 100644 packages/client-query/tsconfig.json create mode 100644 packages/query/package.json create mode 100644 packages/query/src/index.ts create mode 100644 packages/query/src/lq.ts create mode 100644 packages/query/src/messages/query.ts create mode 100644 packages/query/src/notifications/query.ts create mode 100644 packages/query/src/query.ts create mode 100644 packages/query/src/result.ts create mode 100644 packages/query/src/types.ts create mode 100644 packages/query/src/window.ts create mode 100644 packages/query/tsconfig.json create mode 100644 packages/sdk-types/src/query.ts diff --git a/bun.lockb b/bun.lockb index ecf187af96e4354f27cb4a87c734623b2bd2eedf..1a93241bbc237c225b0dc6bb8706a74378d7120f 100755 GIT binary patch delta 15445 zcmeHOd3;UR_CNc`P3}n~k;r^Q2D8XaZb;mdM3KZ8^HfA!M3aczm}PT|p_!-xJ_>Q_%+i}JLu^jG}8>kJ`&yx-4zfBin6p7v{j!r>ms6tWuIjcu1OjKWHP!OPo{7NPCwd-(b*JK^2sb zbvn!Ai(OJL@c!T*G347kb-uBeZqEYh*AoiqP-qC+7_$-L8Ry$Gi5v!N?p{PB&~y-8hi(>QG@=CxGcBqSb$igs(V7AD;_?M z?*^XK!UH7f79-oJxVd8xBvyffua zRteOjNQ^6VRZhmRRXC@(ipR~6uE0wg|LrllBR4_)Ga-1=sNiyzPIXpD(j0h7J<0)v zv+fq?8TB9@JT-`Ps;n$_R_x@-9-+f7qqF4Ja!|z7{U&ArX}!QC2{Css2PJ-WOFb>l zf>Iow0VVk&lxy`3y%JE0`y-5R@JOgS+e$b4GkQW!*icNvG0jNIQWwUz=&_2HZS`Ki z5hqFQP`=Qh!$1?jCm1{fjR$`M`iY=BK&i96!T5kWr=xo=m-IrC7OgREU))f^Q=rIf z_bkvxpk<&m45JO&ACwZfJ1F(f8Me3)VDTccLHbUlVuyudT6YHCNlr4mqzY-w?+Gk)wC zX|fTr(VcV;zXPR6e+JqJbQdV~s9%QOlSN1lfAE;tlctp+k-$?At3atI<0qACJuEWx zUPXO)-OpVB0nN>>pcIM(P)dzl7$%1fL6JiFCEO?eimG?_fd|~y-A(V#d{D};&HT7mX!hsb_0|r8Qjp#RrB0WRkFT6j?sP@w>3x`9sgJWi zcnS^!r66A8ZLF=Ts^BYeJ9_9@=IE)nnF>nQJ~ikQC|AI*1tp)B_R{;Z1UyADuD5Qd z685OxF(ejwwGDC_x;fy<_lA9RzAt$6)zghNLbC1LYmKue`F-$Ih3n#J&)1$%+V^ZJ zCKrBFF)^)Kvw`zIJk-BkvWIjbdw}0}Gh3H+k6&{nAZzhFp1n7VX9ZZVPH`UJzhs<9Ka=%0A&rpel!2cz&Qm?r!0W0v&8VuMJe$ey%iD*$tl8SXB}|B?+;0 zGhP#DXH$4BO7?LjNR^v-@%$hM8_KJLRQ3X|4N{e_y^v^{azT(??qlWo!4CE`uMSo% zNANgx@Lr8GWJTfLAr7T8Lgs+3C$9;zD>K2RfwS=AA$Il=uMJTZ3(`y1fK2HE4l~8j zn8~X{RqGx@Lof?VqV6icL6QqCE7?(cBHTz{w+Dwm%0h6|MgtyzHumw_Fjct$8MV=X zG-P{2o*(W|W+U%OQRW4~cD9q(hO6uvS0WG%WC_L5o7ZHe^XdrI@-8HuxNk^?9OlQp zZ4RXeQYH&3R$dcfSC)grgn-v60{?7!7z$uw8i_T#BYip*&Z|k!BF)DwV;3xO6KROe{12 z-z*!!rIRu*&yC4Aj9$|*1ufst`lbYNbPPt?gnUg zzu>ML-6moji<3uyrEEpU_+7HzmK(YhTkEzeYNR1X0`M!LW`G|l7aE82)t-(AsV}0Zxxm z>t}6h!e=qXdBvp=Y#l4dqtdQp=t4#3f zWYsc~VO8go52o@1$yQ4exFC(ARmj>}l9JJy4~1Lq*qSd&aj=QJHbqrdqsT}^tSu*b zUaBg`w&9CX9hP$vo)?D2hRolz#y0c3_A0x?tJ|x}#J2kSrVr%VI7uo)J37}JE6E*h z54eiDvV{1%%bo{UURTyIL6Ry>j{TKar>V-&M183BB~Wn1t}3r5>GMnC$HVN(EpSFM zAT=`4U77|+6r|Ljz|ktKO?KrNI7%=KT!dW?OyP^t9m?PobE;u@*mGRzs9H`#mduy5 z$gudNYEu$RN)OyoGjIgWRD+}F!4#S~SXV|XoevTzODm)K&;y(iXK?fDxZ`nlu~tS*={K^aOB?kIsYx8Y%4rpk`+>P%I+i4uL)Mtw^gC?)Zu;Tg(M zy3=X|B*Is4x~~NZ@D&{8o8IPsfzy4((q_%T+5#u^v4S^BF-Nw6GyDs(BMlT) zmDAFBzUojG!#*j(Xk)u_1f0?3OuMB4HjpX&C?2>SI`aH1hcXuevLn%>2+#Q1EES8Y zlC8?Aop^q>Ln(oOMzdJeuY>EW1!5G|ZGAKmShk7h<a~vqg^2<%@b^bAzEI z8xZ!yQ~;;v422>-&z!}?O#_D+Nxed-w}8`)9q($l{s;~|xHD#nh&2iH{C0uUV@jj6 z#o+LmMw4g2QJTO3WOw5pI)|-CX1Wy&l0o0&DJQ@cY8={^3wmim z9N3E=?(2}d_2S+I4rLvJiprxqu7E0_!`gK}O65ox%^-$ZyNHq;!>3Vu-?IKF$ zm~dJ>1WCJylDvv%M9I!@ zfaD{Z;~%<4=tQGb@TkEPmHFNQp~)mI1c;6S=psr3H;x!wcT>ij3=HM%23k#2k%PsW zUVTayCmM313g9w$qEudK@I=YZRD&l<<ke6@cVh0lMx-N$<}9IkbcNPdB>(y6&dbHLSkPz1afcS#| z=^rxq_d)3*O64CI^g~e6tNPfuISh(dTj@AJ7g16;2~fq)0J`o+X)M14{DAWS)w=*t zz3&0C`x8LdJ+$gxK|P)fT+yoOGz$0r5R$wR{qvC21NuRz=L;I6)kpm4hVElEflg^i zS{eMmL&-#I!_Ivv+OMZTp-MC=+>bUyc?QZ!CDYI&NuyRpQV2OYR%B2ollbbP3g0s{ zntwO65ye7!%&2=er8=dCoG95UH+Z6SI~o6w?^8gjd^#xAd)%P2K=Hpc+n{qmNpCI* zb=pZ&m#rJ+ZZ+kK(ja({!t z-%ZKwmkc>k(%Zzal#y&R{vo;r|Ih-mjTnQ{7{6igM5+88gObj__t@9^fB$^|wV4Bu zr;h-1{WD7XT{W%yQ!4KUkiA?(FPGYYSHAt@ z;l4g?^gly}fHvd-9}V_RYpefJwg5C7ZOq(xy*5mN#mt?JP7X=oUtcQZjZP-}u+3Il zFY(@@hzbkKsgiEp)-O$}^;h#&V=45cPD|=JT_Q0oGx;lZs|?^1PO2{hKzRcyEF-$H)QmL z4Z}S-;A6<>^&$oL)K){1(PoO?mRf_M#(WJ$devJ7(B)^y{Bd6cP-Es>S{Sv4Ylxp1 z>BC}0q!-I_k|M@h+k5Uu`A*~m2U@LG6<)?wKfHS~Z;9G#g{GY&gz%k$x;5~ps z9ReU(=r7GuxY+B(9&>iUZ5lxD)o%f`#SH`+1GKZJ(;bU}r+_8EQh@g4&jG7|)d0P6 zlmU|fdM%+7AXiZTD)4vUnpk3GqpB!06p9NF9S3L(=y^O1CkW0aydf zgyJ-yH_#WzhwLTbWne9^4k&`WKQI6o1Plg-07HRcz;Iv-^fm$OfmeVHz((K$;1?RS z$8j?rp!pXGv;t_}(KMqkE;RtjOgs;MC9seZjA*bnohae8`AqZN3-|%Nh3#JC!&Z65 zctWrkg89Jrz@xxiz?4}Td-@fGfR==oiRZRz>07CgA{)asDFOlO}pF zB$Sz0*4*X^jc89wLpWD=?!!FG5$#!ml0cMf1`PxKYiE9JXYb{>%nf4UveZC44~voJ`M_|`!Od;0r!I$LQbI>kjBX%4g+`j&=vIhjXp`Pe z=GV#$PLofPd1%qq^Ra=WyUbgl)lErA#WZRlng+vP^Vnmj&4cIlfj6Sd17diw2w~%G=26k` zk86~#7red|dZ~yf+P@_BhOkt2Q(Og+tsbIfC`(B*53~-7$UPo)btigadSf1WExXb2 zdasH7Y+8$H2`N}4Jj8(}EKKg@A=*W;Fg8LQ429(lq9%-m3%@WHVKWcf9`L)crTfmY z6zgQfx&xL^(IpK1F;D7FfBP4^rBfdHiJ?eH!%WxAuw7y+tQkhxHgN`=Z1EDG#vq2~ z`P|wiv;0Ek&!@`_p`Z{-3q*W4eAzAv!%=&Wm>JG`MVjYr6E`^y{B^;Y*=P*f6ucZC z@m)CU$u5iJ2XM|@Fb5wHZB%N>-3Fk<0H*e zvsuAsR(9HNrDmv4XvV*hc*4d~BD*1zv1!0j#P>FYA3hWWjrO`D1p(p&G;HRP+A))x zch2(hPd99!YZF9NBm(z<0fxfYqBIf?UlmV6k7>yiBtD76si#Wuc_iC~kS=b5p*Mz) z9Tca)vTuZMQ_NKJ)a=FUW4<2M*xFumO>=&R$c6%p3~zc@?IU7KQx<0}Lo{ef6Xi`< za0XHhXD`7k5WPP%=~wTk-h3LCk`vlf7|a8*$AWXSU;KIE_gYCRl^Cfc>vNiE6EZl8 znm3QA*6coXs5mxx;hpI{N+d?ZuQ}q+Q7q18&hmSggGOwPxE94?*ex+P2I_a$3YSA7 z#pdR$Et@Kyj%EWR%>%m0>WeLA%@2rxgQ*Frdd{z)#CvcL>YLZI#%yz@SdFG5&11kr z_g}U=|IMY=uul1bkQsU1!91tCc zF^ro;0f@~!Qao_<#oNCvEBuFMK^vB_Vr?wO**sh9F(hTcQO{RCfC3gSa^iP!k_=rD zUd>o9+y8Whe8mKowu?#4VCZIy{&dKQ${qRmGx3&@v53eA$4wt`b08FR z7xqqE@{JFH-L5ZmW!{NqZ3<&e+vxu3!+U6tf801v=$J)DYg=MMijP0*=jKp z#AY5=RwiYB_(Q>}cV#xGLqbYgLVIboSepq2+Q(5k&%8MP=)%AsV`WB9uv9yCdUqRQ zGsV@;FuO=>BkT3J?)N@*W9y!}_dY>m>z<(IvFNgwC$(L+JEtU;oe*7GvFiG!Z4Wk3 z&jn21KCM}J-Er0Zhc*xYc*p5)TI<7)DQKS7-gso?jI|R_6r&_Pp+iD)s`jV3yEhuK zbsLR`?w1DwxkUWfnyt2d-bUXVOmMw+`#1ls#X+oZ)8kQtpYlq!JGej4!g>KvEYc)>{59x~(qbGiQ{GG}u znK(Rg*xr*K%_+CE$%-@O;&;@wnP<&c&uV8mf9Th5wd~WL!u>@e+P0a8(|x9nuloGV ztvGrVB0&+JGI0rBu~}j^>6wSyd!L^ATUq(Owq!^T(hFic6ilzpj$Yi2t~H^4Hz6zc z@2TOXSd++lGp)T;c3Z?GVa^%x((9R_*}W%;HDz+Ubnzd_*f`XWTl)vZ32RPOx*owB zzOlQ4v({aanahv}M>2~LQ&U*PJyuzX_$iHrc%-zKriwEuSQBu~CAY zD##Sa+Sj?piUjL`9#x1JKt%fvEP_>uJ{=Gf^9KMcT72xB64)PeON$9i=*nj29_)nP zd*P7=qedL{sNXgiHqqmYapyBM4Q0!R6 zzZaeDP&beGPn$mdvvsEy`Jw|DQY_E*cf@2Uu;XGek@?~!)Uut!3oyM^jr?MFQkwzS zj*;~g&H8nbosL-$_=xaJXDLi)*$Oc{9gXRn*qn||jueN|k>B%0Ku6fsTd|qHFc`7q z{Z|`Azf%b}jHjyMM!m&2{-M@QajGNRh=@*TijbY36Sg@;;4JM$G(Tjk%ZWAL;scwxC`IOp-WhBh zdqOn9^N%eNpJZT>HGjA8R$S%0hTEb?$*e8Lgf8=!4{_a>E|2VyMtgngJK8WT=nv9J z^A{2S$m_H@=V-=LcPu!1ijqt$cIIy-R{i8Yw5HWdJ(=iGh*v9qh3n% z`B!Y4zXh51q5k`hSVn(pr8Tae?cxX9zW3EW_JObVcQz*f#~xMP`xPSA^8iMYV)9^r z9=w}*aH|XHYZ|WB)EQjfnXZbWgktgQQsyVNJ6J2c|BEjiEG1w;6FoXlq<8oWo9=gO zhrn+(qVnAdl7!8OyZJwYf{W6=ySxQFy`|NB_4caB=c|L>0N(=q$86iGCeLhGaSA-% z1Kk(;9sX$ISIV>@M#FD&F;8V#(b|tQvx+8_mQE=vE-I`no>Z1tR8s6LtBlu*#154` z9ezJ8O}0BdiyiIdpQZf>DY2+xMtS9=_=zHW9E%id#xXCODRL#o*A=_swHy^}DVySb zoxVTf%33bmOIeH7=JzJ^>yr5`$^3F;ekT$Kma=H?g|y?vmA+#44@+T>Ck1-o%YG8j z&U=D!EaUYRWV+=;n!aw-&JyB>Wh^I~zBid4y6|;M`wnG(#5w@p59Q|fsfO^1zDAkf Mp0qDb_&WsuZ+v7s=Kufz delta 14649 zcmeHO2~<=^)_(P*u^%c5BD;bp3Mvf^Leqf#02MdDC9Whcpshj#5pcmJw3OANSo` zx2kU4y0yI5<=Z>GFK_a$wZ(4w{tvJEjh=P&hnz8sk8U66%9y<4(2VmLW8!}<3rjdL zGJC=bCQ-7k7!^@7)uJ<+MDmd&&uj?%K;4zoip#1b$;S%-Ukcg=beels8EH>6c{s1V|R|o2wi-%93&=T|jXf$YXW!#i$(~G1e ztFHJW@;liWo!r;tJ=H77~=kW+(G(Hb?luQiwDF3#&P5=A{z z3oGMb6>%oUYzLXcrK{Xku6ar5?C~1cojbD=`sRQ`S zpd^oKcWeF`$f=wSgH64B6FhlqdZ6CRlFGR1!qRC{Va1g4!iq}wM98U^tAlj93>5X& zc)AAb?w}DXg!>MnB2xSlD0v|&M3T_9=P~ej^vnh&!%TO2fsX`#5r$Cy;V?<+4B84b z0`yr>G~qdqK2yE1?RA6osH(XURypCiAtO*b^=K_9gko|s%%Sv3R0R^gsmSv{Kt3H4cP(OIs<}aqd=v*be6k9 zl6t^m>d|Fb3}Zd-p=Z>Cli;aAgi}>jvAd#>+r5Gco+tlsOQ?&tvWQw(YawShm?O_%QkjRWrijR&0oN}c@)N))y3Ip|(xr8Fo(^Hz~( zFSIBl4FXzKx`C2=!$D~nd=1J#DNw(MKJ~BCT{a1}N}n72E~6f|J8p7uMP=1ygMW_i zwnQjv3~oOE(4so$gCpOZ2Bo@YmX#NmV+y}x@J9^oM-2LrL-%1EAL<=Z)7hz8917}# z(xt_v?zl-s(hS38zG-@`kAjk?7lBfbCW2Crq;$O}sR)kN;4!hM&n|@1gzCGV!%QJuOKc)A<%BTMFzM5P2sWrsP+xP|Z!k04n zu3&{%`>9H{EJ?{Ir0|7)4i>;0{ZzJqJN;Gp6)!&5-z9(V#T)!x%+4GARaVNKZB@3F zSGQG_(-ujBmpqIg@OQ9Iyb(FmxHCYt9`}|cgn>sCHfIUE+|SCI^ZEdl4dRUfs zmLaX&!T^VStvR3D&c%B1`gW?N8hfORFKC-C@3HcCcF9=Fk zX5xwDUVLGYL#Z=3Z{p-{6h1e|#bSAVkZLVKNRc8h+#k=U_w=!X$DnwuTEVN}$WT-h z;80TF%oK3VcvheTU1BF-^9W0F3hpLuh#6cSG6__p`X~omRTuLJZnNYIw1ECHj8v}_601Hq=nGbG=#!*LJ z_vLBfF6Czkdq9ZRLL5prOhPiKInVNUD6_%Q7$6URI0g=JOvZw1iC9S0xC8AR$~bU| znkGD{)HZRXdC=gzI7415%tp)>v_J+91xK~$o?Ze@_dA(;0vxr6)=@4Ju}RKkd|_J$ zyNB0Ds7gI#_h=dif*h81i1s0Tfq%N<##2wdR>aO~ULUE-@3iGtBV9@mCK#3oRD%v$ zhJnk{if}xQ(5uHtd=|j3c670L?u=5EV(bA}4?Ge-5b0pg@OsEDL6!`e7S~EA%yDvh zGp$dHz+tWs2hW@VN2BS*VWSd&)r0!0_bC^gUXJ>^+2C|vTmVPCM}HBU37COtd>>-7 z1WyjgeQ6XB8P(CM{9!O}h;}IsoCyZ%mBU2~xwDI^9D{5aWT+IbaPayrs^tMJQ6u>R zzjVtaEN=a^r_b;-PJ0?0E=j|+r%iaup(o{YJdvULWK?-|jLPQm`WRJt0lT&lA{e9} z!BM>FF+QGG$Eud6J7}gyrYmRgM8kprf@8vQrpN$?Kma!x9EF8N8|*r8R8Mm))V=~o z<=_egWyeT8NGv?7Ln^QCuFAV3c|&)Xd_Iz2?e1cs+}T4_?(NuQFTQLrc2b}Jk|H2c(30~^b$@IfCJds_RALJK1@wqmaatXrW5L!5c z+b8ozo65FuXS}L>&{-dQL{2-0)fcCgkI9Vy z=k@rzhBqdt^2shd&F*5Yd9__-eR;iIwX|edsriQeN&K?iYB?ie_t3b6u8Os*B-v40 z3v#vmX;qnLhbh4~J3=4)#*#P;NCiszo#1@;*3B_keW85339daE*b51@t%cX?EJfl(cY74j_nipDO$<~Jw ziy<4qtFviZ)Mu+oKP(}}%!U^q21mWX%!c<5fXfEQG|QUxGP^(vaTYk;PY4I=3~*G7 zHZ`gDdmtcJ!Ht-k-+;^2INFJNp#!>+#4P|fvMKL9aI~mt!)R@dRfghGpYZ_-Z|qCQ z8)uHHoI@5hfTaNqgrPnRfNm)p${V4v0y3%!u8K)znS5?Pm%^|Q=~grN0f(l3Di$1P ze^uFp9O^5!8BF*KS-hdYi!I@e{Z-2=kWb_L`llNrfO9Thnr{swRM_3%Bu&dmdTR2Y`}H$^&Eb6 zkPBOlGgnns<0wmJX;Vc$*^f8m;(&vpq>X#XJC~M+|x5BshHv6Bj)|KUSbuaP=r~ zdSxufVO1Zl2 zEfJ-9#sDNQFkM^#h7wWAxSJRxqEzu%faK!hWHH3hVxfO&Y_oY9-ka}*+u(0UDeE31|7J=tG9Gf7mlj0$fsv*dS{mi2i@Nfk zih^rMugu6HN{`d=hYYI(rTkf-)WJChoePS8(maFC2PHji5i=;|d+>)m{2(y~C4RZV zgVF|Zi$L=bhXU!WG!%agW&DkaK{aId8k8cs&L~Zk%zw(@Z>D7SCPPk?^frT%Xp2F& zg3`D>OG1OvT0!X-yLioidK=e_$i`o?t+%<1_SVF=ZtJAb2cRD282Lm=-p`=@4NBWL zC8Csn=WU#{v;(G2X&On<5(8{ z|7_##|Hw98w|Y{K0zUKQLHIGL@TI}x<(2pdQoFw;ORjrz|378RJ{(&WysPiQ4^Cc6 zpEmmMf%$b~5C5IXZwJ=7j{G27)F>QpdDih!y!NQovZWJ#Fhh}fRHT*p8}-&`6oqsX zN%`Bytd`AP@$(*)@~Mx}e{wqCnlD zzDrQL4v>t7B^1yS1SF%A(Kot8Vh~VWbUOMDz;lh%%uuA0A2mqH%aD;F^h}BDrbh*M zn*>PShKx?db(#bxN<&63s)2YW1FVLO=6M31sVzm5;YdSoAM`Rnjrkag^y;)6pwz;U z(OK>hfEuIk?D!`|Qfo-VTZnHJRuoHTz}>(O;C|o%;19t2(D?xP5I7B-5et1-+mP4s z_y({aAZG>u?SMeBp*3qO_WQ6&Zh8lO5uoM~WAxLtc4B@A)h+@42rLDb0jmH3JPxb| z?ggd+r9d%I0?^9;HE0E&Q# z0DTO(2cWMgqkz%C2w)g+7myDO2I#yC@7 zhSO<09_&C*AQ|WgbOz|H^b<7rDexKaIq(H=5%?0Q04jkQz$Ab^1C0gfyh-nFbV42k z?24jv<#WI(8v>oeG6A4*3TKF_;xbOLSF=dh7 zG-U%Jp+F1g`u86%b520?H3=lvT;tEJ=vBZrHXxlEQDo?_XAlH8zfo; zfh!czK`bfUyl!|WD*lejS%zo~W_c`M^b3Jk#)x?#EF{dl zPPle;=a=&M)+;ELs+&|Tc80JiVdj0skk?iyrytn00}4qtJDIdr*h5j#4&eroUu-7U zgtA2WKr``XC<{(8Z&{W#I{%edGAs<`dfHNKiIQYq(2S7(F>dIkOXWt5$oa9Ecq$5> z_YuiqsCkGeCQ=}t41@U-#UA2H#YJ#o<}J{=qcx9aPkA;PmDpi7dOJ}>gu|70RxmGy zj(Pc#<%u(2#b`zCHY}=j-l8lVt~9TV9({PeZ;<@AIWkMIB~aZ`lGqT=lETJgQ*D9K z7&-KKeqozGzlthSpn-ub5nqSH$>T+@_AErc>?20EXSreK<A?JzO^5)RB~Fpm0Uif1bt;-c_)*US1&@FocNlFb zEEVS>(B^vaBlK99h>rvpB?d&Y7cex}BT@Ly13O*hbwsIUVqr&=S}oRsgqc@4FSd*f zX|X6e5yp^L;KV5L7S#Vuebu6U6zk5y#egW*Q9co|q8Ag1ma-8+8XrY$;lW{wo|)~H>5I(q(Mzfi5u+Bo0Fc~~=SXv63hyR|Jn zE=*ypXdR8k={BKljIR;eELn_&VwidB^yMCFKT01FJP>x!^ues&Ay!0Vs78oaK$6W{ zts^FUdHrf_;SXBCVE$vNkI+}&Uf=ECeXLFLROC3YE+A*Li0cBMJR|ybK}GAtye=#+ z*}S#7bkfS*fBop_-(^OIC1XmB@2D@A4>pfom$&J8S`<-@4M&?D;u=-d)NgxNtTW~y zo2#|q?51T}*A6r9{Z{^U#Ew77t-nF7J#9THxU@lT8l2kfUi`}rLP7IXl6E4x`6v-}$D_or=~ucK--E4Xi^Y#UFd-g})z1oVCkG54bLvWS z8cv)Lu&|zp6$x?Jr_3A0S6B5<_dC4hkjyYXwD98Yuzu&bZtv8yFYMj$7cXtm#uSlu zb`$rZY?ygVxud;y$q{+@eK$0$Vs{*dDMcIy2{&&!hwKV??(p+JFS?=klkG(tD`9ND zxHleAU|z8;J-q2Y^~?uD5iK-Lur)m*)axr_ilfOa{TJoJ%zN7L zjEix{3vpt?7-8}3DRmZirC?1kFKo98{8z!{WoH*?3t9qBYEpsNfMLUSIm3a$Ht$u> zdaHIr?%)&e8G2Zv7KyK+mu%kn{=C)1P8Aakw8aYN2 zk~rsNfnHc&^~RW9W!P^nD4$9fzNv@{^LF_2JH7b{d(P4u)|&UngV)bn{qVNow2vEM zryCPy-Z-~q9CCm5$Gb+}C>kmJoH*QUOJ$)fP8>>w_snbTtD;|X&-A|w^HHk}UeT3t z93&^2b{VXg<3trk=XJ#<4b2xJW~Skf$3%P@Tx{ND|8&QNyULzl9Y8USmk31X9fmon zr}Us$kcJ&#mDr6Muu&c-GF7b40I3nziL;B4bg0ydymU+f^GbX3hoU!+oKcsdO&mMT zf&IP3qI8_g?yS*AyhSx_5zR6{S_(%7%L_B_z8_iBFt*_>Uu-IP!$`3uIUKqfVdf?I zdmcXgY_rIQDk#vZhkhA$-m-Yd{igjx~ku>XD23mu&I1;FP7t8Ww?1G|3dzy?eDHTCqaXv z2u?cIUmSF?$t+Q*-hb7)<~-a3j!y(we!)=Ai%YswIX zv|Sr|=Kl-WUNCft--R7BZnStyeYbAGgMZWL81KgB|1?;1|Rn78;f^&dMmwU z8zmZ4X@s0e(nPPy5=#4KDAGlxz(i~L`T!~kOvAtPW@lbD;;Fr{1H{$-( zyc5>Ny_8vJW)=bHA9 zZd~;sczPNBbl9<1OHL}Yx6wB#q{+)gVjq?%HuhmPdwTa}hx5c}&Sv;*qi1F14<8T@&Et; diff --git a/packages/client-query/package.json b/packages/client-query/package.json new file mode 100644 index 0000000000..b5582f16cd --- /dev/null +++ b/packages/client-query/package.json @@ -0,0 +1,19 @@ +{ + "name": "@communication/client-query", + "version": "0.1.0", + "main": "src/index.ts", + "module": "src/index.ts", + "type": "module", + "devDependencies": { + "@types/bun": "^1.1.14" + }, + "dependencies": { + "@communication/types": "workspace:*", + "@communication/sdk-types": "workspace:*", + "@communication/query": "workspace:*", + "fast-equals": "^5.0.1" + }, + "peerDependencies": { + "typescript": "^5.6.3" + } +} diff --git a/packages/client-query/src/index.ts b/packages/client-query/src/index.ts new file mode 100644 index 0000000000..4a1c969701 --- /dev/null +++ b/packages/client-query/src/index.ts @@ -0,0 +1,26 @@ +import { LiveQueries } from '@communication/query' +import type { Client } from '@communication/sdk-types' + +import { MessagesQuery, NotificationsQuery } from './query' + +let lq: LiveQueries + +export function createMessagesQuery(): MessagesQuery { + return new MessagesQuery(lq) +} + +export function createNotificationsQuery(): NotificationsQuery { + return new NotificationsQuery(lq) +} + +export function initLiveQueries(client: Client) { + if (lq != null) { + lq.close() + } + + lq = new LiveQueries(client) + + client.onEvent = (event) => { + void lq.onEvent(event) + } +} diff --git a/packages/client-query/src/query.ts b/packages/client-query/src/query.ts new file mode 100644 index 0000000000..bbaad06b38 --- /dev/null +++ b/packages/client-query/src/query.ts @@ -0,0 +1,65 @@ +import { type LiveQueries } from '@communication/query' +import type { MessagesQueryCallback, NotificationsQueryCallback, QueryCallback } from '@communication/sdk-types' +import { type FindMessagesParams, type FindNotificationsParams } from '@communication/types' +import { deepEqual } from 'fast-equals' + +class BaseQuery

, C extends QueryCallback> { + private oldQuery: P | undefined + private oldCallback: QueryCallback | undefined + + constructor(protected readonly lq: LiveQueries) {} + + unsubscribe: () => void = () => {} + + query(params: P, callback: C): boolean { + if (!this.needUpdate(params, callback)) { + return false + } + this.doQuery(params, callback) + return true + } + + private doQuery(query: P, callback: C): void { + this.unsubscribe() + this.oldCallback = callback + this.oldQuery = query + + const { unsubscribe } = this.createQuery(query, callback) + this.unsubscribe = () => { + unsubscribe() + this.oldCallback = undefined + this.oldQuery = undefined + this.unsubscribe = () => {} + } + } + + // eslint-disable-next-line @typescript-eslint/no-unused-vars + createQuery(params: P, callback: C): { unsubscribe: () => void } { + return { + unsubscribe: () => {} + } + } + + private needUpdate(params: FindMessagesParams, callback: MessagesQueryCallback): boolean { + if (!deepEqual(params, this.oldQuery)) return true + if (!deepEqual(callback.toString(), this.oldCallback?.toString())) return true + return false + } +} + +export class MessagesQuery extends BaseQuery { + override createQuery(params: FindMessagesParams, callback: MessagesQueryCallback): { unsubscribe: () => void } { + return this.lq.queryMessages(params, callback) + } +} + +export class NotificationsQuery extends BaseQuery { + override createQuery( + params: FindNotificationsParams, + callback: NotificationsQueryCallback + ): { + unsubscribe: () => void + } { + return this.lq.queryNotifications(params, callback) + } +} diff --git a/packages/client-query/tsconfig.json b/packages/client-query/tsconfig.json new file mode 100644 index 0000000000..3ae07cd3fa --- /dev/null +++ b/packages/client-query/tsconfig.json @@ -0,0 +1,9 @@ +{ + "extends": "../../tsconfig.json", + "compilerOptions": { + "jsx": "react-jsx", + "outDir": "./dist", + "rootDir": "./src" + }, + "include": ["src"] +} diff --git a/packages/client-sqlite/src/client.ts b/packages/client-sqlite/src/client.ts index fb0edd5c66..4b50d2e06c 100644 --- a/packages/client-sqlite/src/client.ts +++ b/packages/client-sqlite/src/client.ts @@ -17,7 +17,7 @@ import { type MessageCreatedEvent, type DbAdapter, EventType, - type BroadcastEvent, + type BroadcastEvent } from '@communication/sdk-types' import { createDbAdapter as createSqliteDbAdapter } from '@communication/sqlite-wasm' diff --git a/packages/query/package.json b/packages/query/package.json new file mode 100644 index 0000000000..f613302134 --- /dev/null +++ b/packages/query/package.json @@ -0,0 +1,19 @@ +{ + "name": "@communication/query", + "version": "0.1.0", + "main": "src/index.ts", + "module": "src/index.ts", + "type": "module", + "devDependencies": { + "@types/bun": "^1.1.14", + "@types/crypto-js": "^4.2.2" + }, + "dependencies": { + "@communication/types": "workspace:*", + "@communication/sdk-types": "workspace:*", + "fast-equals": "^5.0.1" + }, + "peerDependencies": { + "typescript": "^5.6.3" + } +} diff --git a/packages/query/src/index.ts b/packages/query/src/index.ts new file mode 100644 index 0000000000..57ad51bd4b --- /dev/null +++ b/packages/query/src/index.ts @@ -0,0 +1 @@ +export * from './lq.ts' diff --git a/packages/query/src/lq.ts b/packages/query/src/lq.ts new file mode 100644 index 0000000000..d19baf7c8a --- /dev/null +++ b/packages/query/src/lq.ts @@ -0,0 +1,150 @@ +import { type FindMessagesParams, type FindNotificationsParams } from '@communication/types' +import { deepEqual } from 'fast-equals' +import type { + Client, + MessagesQueryCallback, + NotificationsQueryCallback, + BroadcastEvent +} from '@communication/sdk-types' + +import type { Query, QueryId } from './types' +import { MessagesQuery } from './messages/query' +import { NotificationQuery } from './notifications/query' + +interface CreateQueryResult { + unsubscribe: () => void +} + +const maxQueriesCache = 10 + +export class LiveQueries { + private readonly client: Client + private readonly queries = new Map() + private readonly unsubscribed = new Set() + private counter: number = 0 + + constructor(client: Client) { + this.client = client + } + + async onEvent(event: BroadcastEvent): Promise { + for (const q of this.queries.values()) { + await q.onEvent(event) + } + } + + queryMessages(params: FindMessagesParams, callback: MessagesQueryCallback): CreateQueryResult { + const query = this.createMessagesQuery(params, callback) + this.queries.set(query.id, query) + + return { + unsubscribe: () => { + this.unsubscribeQuery(query) + } + } + } + + queryNotifications(params: FindNotificationsParams, callback: NotificationsQueryCallback): CreateQueryResult { + const query = this.createNotificationQuery(params, callback) + this.queries.set(query.id, query) + + return { + unsubscribe: () => { + this.unsubscribeQuery(query) + } + } + } + + private createMessagesQuery(params: FindMessagesParams, callback: MessagesQueryCallback): MessagesQuery { + const id = ++this.counter + const exists = this.findMessagesQuery(params) + + if (exists !== undefined) { + if (this.unsubscribed.has(id)) { + this.unsubscribed.delete(id) + exists.setCallback(callback) + return exists + } else { + const result = exists.copyResult() + return new MessagesQuery(this.client, id, params, callback, result) + } + } + + return new MessagesQuery(this.client, id, params, callback) + } + + private createNotificationQuery( + params: FindNotificationsParams, + callback: NotificationsQueryCallback + ): NotificationQuery { + const id = ++this.counter + const exists = this.findNotificationQuery(params) + + if (exists !== undefined) { + if (this.unsubscribed.has(id)) { + this.unsubscribed.delete(id) + exists.setCallback(callback) + return exists + } else { + const result = exists.copyResult() + return new NotificationQuery(this.client, id, params, callback, result) + } + } + + return new NotificationQuery(this.client, id, params, callback) + } + + private findMessagesQuery(params: FindMessagesParams): MessagesQuery | undefined { + for (const query of this.queries.values()) { + if (query instanceof MessagesQuery) { + if (!this.queryCompare(params, query.params)) continue + return query + } + } + } + + private findNotificationQuery(params: FindMessagesParams): NotificationQuery | undefined { + for (const query of this.queries.values()) { + if (query instanceof NotificationQuery) { + if (!this.queryCompare(params, query.params)) continue + return query + } + } + } + + private queryCompare(q1: FindMessagesParams, q2: FindMessagesParams): boolean { + if (Object.keys(q1).length !== Object.keys(q2).length) { + return false + } + return deepEqual(q1, q2) + } + + private removeOldQueries(): void { + const unsubscribed = Array.from(this.unsubscribed) + for (let i = 0; i < this.unsubscribed.size / 2; i++) { + const id = unsubscribed.shift() + if (id === undefined) return + this.unsubscribe(id) + } + } + + private unsubscribe(id: QueryId): void { + const query = this.queries.get(id) + if (query == null) return + void query.unsubscribe() + this.queries.delete(id) + this.unsubscribed.delete(id) + } + + private unsubscribeQuery(query: Query): void { + this.unsubscribed.add(query.id) + query.removeCallback() + if (this.unsubscribed.size > maxQueriesCache) { + this.removeOldQueries() + } + } + + close(): void { + this.queries.clear() + } +} diff --git a/packages/query/src/messages/query.ts b/packages/query/src/messages/query.ts new file mode 100644 index 0000000000..794b346c26 --- /dev/null +++ b/packages/query/src/messages/query.ts @@ -0,0 +1,202 @@ +import { + type CardID, + type FindMessagesParams, + type ID, + type Message, + type Patch, + SortOrder +} from '@communication/types' +import { + type AttachmentCreatedEvent, + type MessageCreatedEvent, + type PatchCreatedEvent, + type ReactionCreatedEvent, + EventType, + type BroadcastEvent, + type AttachmentRemovedEvent, + type MessageRemovedEvent, + type ReactionRemovedEvent +} from '@communication/sdk-types' + +import { BaseQuery } from '../query' + +export class MessagesQuery extends BaseQuery { + override async find(params: FindMessagesParams): Promise { + return this.client.findMessages(params, this.id) + } + + override getObjectId(object: Message): ID { + return object.id + } + + override getObjectDate(object: Message): Date { + return object.created + } + + override async onEvent(event: BroadcastEvent): Promise { + switch (event.type) { + case EventType.MessageCreated: + return await this.onCreateMessageEvent(event) + case EventType.MessageRemoved: + return await this.onRemoveMessageEvent(event) + case EventType.PatchCreated: + return await this.onCreatePatchEvent(event) + case EventType.ReactionCreated: + return await this.onCreateReactionEvent(event) + case EventType.ReactionRemoved: + return await this.onRemoveReactionEvent(event) + case EventType.AttachmentCreated: + return await this.onCreateAttachmentEvent(event) + case EventType.AttachmentRemoved: + return await this.onRemoveAttachmentEvent(event) + } + } + + async onCreateMessageEvent(event: MessageCreatedEvent): Promise { + if (this.result instanceof Promise) { + this.result = await this.result + } + + const message = { + ...event.message, + edited: new Date(event.message.edited), + created: new Date(event.message.created) + } + const exists = this.result.get(message.id) + + if (exists !== undefined) return + if (!this.match(message, event.card)) return + + if (this.result.isTail()) { + if (this.params.sort === SortOrder.Asc) { + this.result.push(message) + } else { + this.result.unshift(message) + } + await this.notify() + } + } + + private match(message: Message, card: CardID): boolean { + if (this.params.id != null && this.params.id !== message.id) { + return false + } + if (this.params.card != null && this.params.card !== card) { + return false + } + return true + } + + private async onCreatePatchEvent(event: PatchCreatedEvent): Promise { + if (this.result instanceof Promise) { + this.result = await this.result + } + + const patch = { + ...event.patch, + created: new Date(event.patch.created) + } + + const message = this.result.get(patch.message) + + if (message === undefined) return + + if (message.created < patch.created) { + this.result.update(this.applyPatch(message, patch)) + await this.notify() + } + } + + private async onRemoveMessageEvent(event: MessageRemovedEvent): Promise { + if (this.result instanceof Promise) { + this.result = await this.result + } + + const deleted = this.result.delete(event.message) + + if (deleted !== undefined) { + await this.notify() + } + } + + private async onCreateReactionEvent(event: ReactionCreatedEvent): Promise { + if (this.result instanceof Promise) { + this.result = await this.result + } + + const reaction = { + ...event.reaction, + created: new Date(event.reaction.created) + } + const message = this.result.get(reaction.message) + if (message === undefined) return + + message.reactions.push(reaction) + this.result.update(message) + await this.notify() + } + + private async onRemoveReactionEvent(event: ReactionRemovedEvent): Promise { + if (this.result instanceof Promise) { + this.result = await this.result + } + + const message = this.result.get(event.message) + if (message === undefined) return + + const reactions = message.reactions.filter((it) => it.reaction !== event.reaction && it.creator !== event.creator) + if (reactions.length === message.reactions.length) return + + const updated = { + ...message, + reactions + } + this.result.update(updated) + await this.notify() + } + + private async onCreateAttachmentEvent(event: AttachmentCreatedEvent): Promise { + if (this.result instanceof Promise) { + this.result = await this.result + } + + const attachment = { + ...event.attachment, + created: new Date(event.attachment.created) + } + const message = this.result.get(attachment.message) + if (message === undefined) return + + message.attachments.push(attachment) + this.result.update(message) + await this.notify() + } + + private async onRemoveAttachmentEvent(event: AttachmentRemovedEvent): Promise { + if (this.result instanceof Promise) { + this.result = await this.result + } + + const message = this.result.get(event.message) + if (message === undefined) return + + const attachments = message.attachments.filter((it) => it.card !== event.card) + if (attachments.length === message.attachments.length) return + + const updated = { + ...message, + attachments + } + this.result.update(updated) + await this.notify() + } + + private applyPatch(message: Message, patch: Patch): Message { + return { + ...message, + content: patch.content, + creator: patch.creator, + created: patch.created + } + } +} diff --git a/packages/query/src/notifications/query.ts b/packages/query/src/notifications/query.ts new file mode 100644 index 0000000000..8e9d39cfe3 --- /dev/null +++ b/packages/query/src/notifications/query.ts @@ -0,0 +1,129 @@ +import { + type FindNotificationsParams, + SortOrder, + type Notification, + type ID, +} from '@communication/types' +import { + type NotificationCreatedEvent, + EventType, + type BroadcastEvent, + type NotificationContextRemovedEvent, + type NotificationRemovedEvent, + type NotificationContextUpdatedEvent, +} from '@communication/sdk-types' + +import {BaseQuery} from '../query.ts'; + +export class NotificationQuery extends BaseQuery { + override async find(params: FindNotificationsParams): Promise { + return this.client.findNotifications(params, this.id) + } + + override getObjectId(object: Notification): ID { + return object.message.id + } + + override getObjectDate(object: Notification): Date { + return object.message.created + } + + override async onEvent(event: BroadcastEvent): Promise { + switch (event.type) { + case EventType.NotificationCreated: + return await this.onCreateNotificationEvent(event) + case EventType.NotificationRemoved: + return await this.onRemoveNotificationEvent(event) + case EventType.NotificationContextUpdated: + return await this.onUpdateNotificationContextEvent(event) + case EventType.NotificationContextRemoved: + return await this.onRemoveNotificationContextEvent(event) + } + } + + async onCreateNotificationEvent(event: NotificationCreatedEvent): Promise { + if (this.result instanceof Promise) { + this.result = await this.result + } + + const exists = this.result.get(event.notification.message.id) + if (exists !== undefined) return + + if (this.params.message != null && this.params.message !== event.notification.message.id) return + if (this.params.context != null && this.params.context !== event.notification.context) return + + if (this.result.isTail()) { + if (this.params.sort === SortOrder.Asc) { + this.result.push(event.notification) + } else { + this.result.unshift(event.notification) + } + await this.notify() + } + } + + + private async onUpdateNotificationContextEvent(event: NotificationContextUpdatedEvent): Promise { + if (this.result instanceof Promise) { + this.result = await this.result + } + + if (this.params.context != null && this.params.context !== event.context) return + if (event.update.lastView === undefined && event.update.archivedFrom === undefined) return + + const toUpdate = this.params.context === event.context ? + this.result.getResult() + : this.result.getResult().filter(it => it.context === event.context) + if (toUpdate.length === 0) return + + for (const notification of toUpdate) { + this.result.update({ + ...notification, + ...event.update.lastView !== undefined ? { + read: event.update.lastView < notification.message.created + } : {}, + ...event.update.archivedFrom !== undefined ? { + archived: event.update.archivedFrom < notification.message.created + } : {} + }) + } + } + + private async onRemoveNotificationEvent(event: NotificationRemovedEvent): Promise { + if (this.result instanceof Promise) { + this.result = await this.result + } + + const deleted = this.result.delete(event.message) + + if (deleted !== undefined) { + await this.notify() + } + } + + private async onRemoveNotificationContextEvent(event: NotificationContextRemovedEvent): Promise { + if (this.result instanceof Promise) { + this.result = await this.result + } + + if (this.params.context != null && this.params.context !== event.context) return + + if (event.context === this.params.context) { + if (this.result.length === 0) return + this.result.deleteAll() + this.result.setHead(true) + this.result.setTail(true) + await this.notify() + } else { + const toRemove = this.result.getResult().filter(it => it.context === event.context) + if (toRemove.length === 0) return + + for (const notification of toRemove) { + this.result.delete(notification.message.id) + } + await this.notify() + } + + } + +} diff --git a/packages/query/src/query.ts b/packages/query/src/query.ts new file mode 100644 index 0000000000..bec8fc200f --- /dev/null +++ b/packages/query/src/query.ts @@ -0,0 +1,198 @@ +import { Direction, type ID, SortOrder } from '@communication/types' +import { type BroadcastEvent, type QueryCallback, type Client } from '@communication/sdk-types' + +import { QueryResult } from './result' +import { defaultQueryParams, type FindParams, type Query, type QueryId } from './types' +import { WindowImpl } from './window' + +export class BaseQuery implements Query { + protected result: QueryResult | Promise> + private forward: Promise | T[] = [] + private backward: Promise | T[] = [] + + constructor( + protected readonly client: Client, + public readonly id: QueryId, + public readonly params: P, + private callback?: QueryCallback, + initialResult?: QueryResult + ) { + if (initialResult !== undefined) { + this.result = initialResult + void this.notify() + } else { + const limit = this.params.limit ?? defaultQueryParams.limit + const findParams = { + ...this.params, + excluded: this.params.excluded ?? defaultQueryParams.excluded, + direction: this.params.direction ?? defaultQueryParams.direction, + sort: this.params.sort ?? defaultQueryParams.sort, + limit: limit + 1 + } + + const findPromise = this.find(findParams) + this.result = findPromise.then((res) => { + const isTail = params.from ? res.length <= limit : params.sort === SortOrder.Desc + const isHead = params.from === undefined && params.sort === SortOrder.Asc + if (!isTail) { + res.pop() + } + const qResult = new QueryResult(res, this.getObjectId) + qResult.setTail(isTail) + qResult.setHead(isHead) + + return qResult + }) + this.result + .then(async () => { + await this.notify() + }) + .catch((err: any) => { + console.error('Failed to update Live query: ', err) + }) + } + } + + // eslint-disable-next-line @typescript-eslint/no-unused-vars + protected async find(params: FindParams): Promise { + /*Implement in subclass*/ + return [] as T[] + } + + // eslint-disable-next-line @typescript-eslint/no-unused-vars + protected getObjectId(object: T): ID { + /*Implement in subclass*/ + return '' as ID + } + + // eslint-disable-next-line @typescript-eslint/no-unused-vars + protected getObjectDate(object: T): Date { + /*Implement in subclass*/ + return new Date(0) as Date + } + + // eslint-disable-next-line @typescript-eslint/no-unused-vars + async onEvent(event: BroadcastEvent): Promise { + /*Implement in subclass*/ + } + + setCallback(callback: QueryCallback): void { + this.callback = callback + void this.notify() + } + + removeCallback(): void { + this.callback = () => {} + } + + protected async notify(): Promise { + if (this.callback === undefined) return + if (this.result instanceof Promise) { + this.result = await this.result + } + + const result = this.result.getResult() + const isTail = this.result.isTail() + const isHead = this.result.isHead() + + const window = new WindowImpl(result, isTail, isHead, this) + this.callback(window) + } + + async loadForward() { + if (this.result instanceof Promise) { + this.result = await this.result + } + if (this.forward instanceof Promise) { + this.forward = await this.forward + } + + if (this.result.isTail()) return + + const last = this.result.getLast() + if (last === undefined) return + + const limit = this.params.limit ?? defaultQueryParams.limit + const findParams: FindParams = { + ...this.params, + from: this.getObjectDate(last), + excluded: true, + direction: Direction.Forward, + limit: limit + 1, + sort: SortOrder.Asc + } + + const forward = this.find(findParams) + + this.forward = forward.then(async (res) => { + if (this.result instanceof Promise) { + this.result = await this.result + } + const isTail = res.length <= limit + if (!isTail) { + res.pop() + } + this.result.append(res) + this.result.setTail(isTail) + await this.notify() + return res + }) + } + + async loadBackward() { + if (this.result instanceof Promise) { + this.result = await this.result + } + if (this.backward instanceof Promise) { + this.backward = await this.backward + } + + if (this.result.isHead()) return + + const first = this.params.sort === SortOrder.Asc ? this.result.getFirst() : this.result.getLast() + if (first === undefined) return + + const limit = this.params.limit ?? defaultQueryParams.limit + const findParams: FindParams = { + ...this.params, + from: this.getObjectDate(first), + excluded: true, + direction: Direction.Backward, + limit: limit + 1, + sort: SortOrder.Desc + } + + const backward = this.find(findParams) + this.backward = backward.then(async (res) => { + if (this.result instanceof Promise) { + this.result = await this.result + } + const isHead = res.length <= limit + if (!isHead) { + res.pop() + } + + if (this.params.sort === SortOrder.Asc) { + const reversed = res.reverse() + this.result.prepend(reversed) + } else { + this.result.append(res) + } + this.result.setHead(isHead) + await this.notify() + return res + }) + } + + copyResult(): QueryResult | undefined { + if (this.result instanceof Promise) { + return undefined + } + + return this.result.copy() + } + + async unsubscribe(): Promise { + await this.client.unsubscribeQuery(this.id) + } +} diff --git a/packages/query/src/result.ts b/packages/query/src/result.ts new file mode 100644 index 0000000000..36395f74d5 --- /dev/null +++ b/packages/query/src/result.ts @@ -0,0 +1,92 @@ +import type { ID } from '@communication/types' + +export class QueryResult { + private objectById: Map + + private tail: boolean = false + private head: boolean = false + + get length(): number { + return this.objectById.size + } + + constructor( + messages: T[], + private readonly getId: (it: T) => ID + ) { + this.objectById = new Map(messages.map((it) => [getId(it), it])) + } + + isTail(): boolean { + return this.tail + } + + isHead(): boolean { + return this.head + } + + setHead(head: boolean) { + this.head = head + } + + setTail(tail: boolean) { + this.tail = tail + } + + getResult(): T[] { + return Array.from(this.objectById.values()) + } + + get(id: ID): Readonly | undefined { + return this.objectById.get(id) + } + + delete(id: ID): T | undefined { + const object = this.objectById.get(id) + this.objectById.delete(id) + return object + } + + deleteAll() { + this.objectById.clear() + } + + push(object: T): void { + this.objectById.set(this.getId(object), object) + } + + unshift(object: T): void { + this.objectById = new Map([[this.getId(object), object], ...this.objectById]) + } + + update(object: T): void { + this.objectById.set(this.getId(object), object) + } + + getFirst(): T | undefined { + return Array.from(this.objectById.values())[0] + } + + getLast(): T | undefined { + return Array.from(this.objectById.values())[this.objectById.size - 1] + } + + prepend(objects: T[]) { + const current = Array.from(this.objectById.entries()) + this.objectById = new Map([...objects.map<[ID, T]>((object) => [this.getId(object), object]), ...current]) + } + + append(objects: T[]) { + for (const object of objects) { + this.objectById.set(this.getId(object), object) + } + } + + copy(): QueryResult { + const copy = new QueryResult(Array.from(this.objectById.values()), this.getId) + + copy.setHead(this.head) + copy.setTail(this.tail) + return copy + } +} diff --git a/packages/query/src/types.ts b/packages/query/src/types.ts new file mode 100644 index 0000000000..895e21d3d9 --- /dev/null +++ b/packages/query/src/types.ts @@ -0,0 +1,33 @@ +import { type BroadcastEvent } from '@communication/sdk-types' +import { Direction, SortOrder, type Window } from '@communication/types' + +import { QueryResult } from './result.ts' + +export type QueryId = number + +export const defaultQueryParams = { + limit: 50, + excluded: false, + direction: Direction.Forward, + sort: SortOrder.Desc +} + +export type FindParams = Partial & { + from?: Date +} + +export interface Query { + readonly id: QueryId + readonly params: P + + onEvent(event: BroadcastEvent): Promise + + loadForward(): Promise + loadBackward(): Promise + + unsubscribe(): Promise + + setCallback(callback: (window: Window) => void): void + removeCallback(): void + copyResult(): QueryResult | undefined +} diff --git a/packages/query/src/window.ts b/packages/query/src/window.ts new file mode 100644 index 0000000000..508c2b5d89 --- /dev/null +++ b/packages/query/src/window.ts @@ -0,0 +1,34 @@ +import type { Window } from '@communication/types' + +import type { Query } from './types' + +export class WindowImpl implements Window { + constructor( + private readonly result: T[], + private readonly isTail: boolean, + private readonly isHead: boolean, + private readonly query: Query + ) {} + + getResult(): T[] { + return this.result + } + + async loadNextPage(): Promise { + if (!this.hasNextPage()) return + await this.query.loadForward() + } + + async loadPrevPage(): Promise { + if (!this.hasPrevPage()) return + await this.query.loadBackward() + } + + hasNextPage(): boolean { + return !this.isTail + } + + hasPrevPage(): boolean { + return !this.isHead + } +} diff --git a/packages/query/tsconfig.json b/packages/query/tsconfig.json new file mode 100644 index 0000000000..3ae07cd3fa --- /dev/null +++ b/packages/query/tsconfig.json @@ -0,0 +1,9 @@ +{ + "extends": "../../tsconfig.json", + "compilerOptions": { + "jsx": "react-jsx", + "outDir": "./dist", + "rootDir": "./src" + }, + "include": ["src"] +} diff --git a/packages/sdk-types/src/client.ts b/packages/sdk-types/src/client.ts index dce6d468e9..af3edce94b 100644 --- a/packages/sdk-types/src/client.ts +++ b/packages/sdk-types/src/client.ts @@ -42,4 +42,3 @@ export interface Client { unsubscribeQuery(id: number): Promise close(): void } - diff --git a/packages/sdk-types/src/index.ts b/packages/sdk-types/src/index.ts index 38c4587b9e..01596627a6 100644 --- a/packages/sdk-types/src/index.ts +++ b/packages/sdk-types/src/index.ts @@ -1,4 +1,5 @@ export * from './db' export * from './event' export * from './ws' -export * from './client' \ No newline at end of file +export * from './client' +export * from './query' diff --git a/packages/sdk-types/src/query.ts b/packages/sdk-types/src/query.ts new file mode 100644 index 0000000000..c09a164af5 --- /dev/null +++ b/packages/sdk-types/src/query.ts @@ -0,0 +1,6 @@ +import type { Message, Window, Notification } from '@communication/types' + +export type QueryCallback = (window: Window) => void + +export type MessagesQueryCallback = QueryCallback +export type NotificationsQueryCallback = QueryCallback