From e884fbec2f549a76dfa363683cf6ffbc7300291d Mon Sep 17 00:00:00 2001 From: Kristina Date: Mon, 23 Dec 2024 17:26:39 +0400 Subject: [PATCH] Init cockroach adapter (#4) * Init cockroach adapter --- bun.lockb | Bin 56600 -> 57472 bytes packages/cockroach/package.json | 5 +- packages/cockroach/src/adapter.ts | 122 +++++++++ packages/cockroach/src/connection.ts | 104 ++++++++ packages/cockroach/src/db/base.ts | 41 +++ packages/cockroach/src/db/message.ts | 228 +++++++++++++++++ packages/cockroach/src/db/notification.ts | 239 ++++++++++++++++++ packages/cockroach/src/db/types.ts | 59 +++++ packages/cockroach/src/index.ts | 1 + packages/postgres/migrations/01_message.sql | 19 -- packages/postgres/migrations/03_reaction.sql | 11 - .../postgres/migrations/04_notification.sql | 7 - .../migrations/05_notificationContext.sql | 12 - packages/postgres/src/index.ts | 0 packages/sdk-types/package.json | 16 ++ packages/sdk-types/src/db.ts | 54 ++++ packages/sdk-types/src/index.ts | 1 + packages/sdk-types/tsconfig.json | 8 + 18 files changed, 877 insertions(+), 50 deletions(-) create mode 100644 packages/cockroach/src/adapter.ts create mode 100644 packages/cockroach/src/connection.ts create mode 100644 packages/cockroach/src/db/base.ts create mode 100644 packages/cockroach/src/db/message.ts create mode 100644 packages/cockroach/src/db/notification.ts create mode 100644 packages/cockroach/src/db/types.ts create mode 100644 packages/cockroach/src/index.ts delete mode 100644 packages/postgres/migrations/01_message.sql delete mode 100644 packages/postgres/migrations/03_reaction.sql delete mode 100644 packages/postgres/migrations/04_notification.sql delete mode 100644 packages/postgres/migrations/05_notificationContext.sql delete mode 100644 packages/postgres/src/index.ts create mode 100644 packages/sdk-types/package.json create mode 100644 packages/sdk-types/src/db.ts create mode 100644 packages/sdk-types/src/index.ts create mode 100644 packages/sdk-types/tsconfig.json diff --git a/bun.lockb b/bun.lockb index 22de1295d3aa3489aa57d4f6cf58200747fc3dcd..5db3fe48a99d18dac938dc6ad72f92e427c2bd1e 100755 GIT binary patch delta 8169 zcmeHMd3;pW^?!Gg33)>%E08=UkOUL9m~}`d5Xb{WAV5$o5ODz~WD*9G5Hb@aMa`r@ zB?7XXC`(94SVRQD5J4&^oBi0JptxZb2yH*sS{Bi&)K;+H^WHq9{%YIb=hHvE^ZCB> z&b{ZJd+xdCF7xi%mjZl80#>FDct5>y_S8#3<@W7g%h`)6e>qJ4b;1N^Qr$Iob&nMT zleS;!l(h8Js=~P1R~mFie_ABT*MZjvP?x97UE!6a`vN41^ACZBgO<7GRj~A4fBq&v z-3cnAyx8TcOmTap0`OtrNBQ%2n{_@aNY`fu4V?gi90-JfMu8@Qx;-hSWpj$8#a(p4 zVB~k@{L+d!RW3=oX43f*cSUhZb)|GLSm)blgov(fgAFWw2~uqCNZ8Kimy|iZURSa7 zpsQvcn|~RZEO!oaY^KXo?RA%VB$sDywX3Q|`V@s+&*iDCa-r4NAr}R{7^ZSJoipcD zu^$d0pRHeN)$)^lJCMf_!f9rBJ=WurM;2`VO81Ba!2q<@BV$ z#?sN?SxmtgK&5XMDBIzjt(h%-4!$S&v*a{I*FFKIk>Kh;S-RZc_#L3V!H@LkdtuSd z<>|=hdOjG&+T+2qw~lnzdpX;aQsXQylblthmCh=UYcKM-mlt~I^c*Pk)%t3pb$9TH zIpMx;P{e`~D0{&hBT1<1`$R~R0@?w}h8gah3O*kEFn@Uuns=758>~lFt>iPpYAq`AxNPgC3tR$azNpGM(}f`@MyuS@XF=KD--2>^xvQ$wmHB;C ztwBBZM2XYmoeS?&xvD+xk{am>tl{x*gSBkPHPEo3czMiU!Q(2Q=ck4WP_7VrFnyoCoFR zSA%l?<0$8D8FFRd+3%x4xmCDJqt(^!V!TR~b3ESCDg-r$rX-D&*7nmQV-YARW97r<4^eMCE$=%t7{16tC7B{D4pBn&7FE1S9TwFT9Y~h0d14IZcU8qg zs>kyXb>R68$swwoW|ky$$tP%1SBI#k`VdvTLLDKhsaFtLta)M@OcqTBUJetKuC2JCu%h%aiZK zNZyJ}v^Kb@g0a5~vyQ_Dp8Ju2^N9U@a(_P`fiol)M0@?fva6_mdJXb!5HwW@Vc@Pj7 za&J{6Q+{t%u8)zVak_3e=pE`nUM$w#yO4(vgd^%HKTb8DM&4AiT5`?%uq=($-iBZW zovyvL;%&0__7&bn^P8NDm4xloXW%N56I5}6@)K0K3syNjKzN7>z;U$cA>KxEAJzOd z=TU)@E2sC?2L-_nzf^$Z;Y1jK+YgR|!mJIo8S5?UK@cK*RB$|)@I+UK>;cCiV&-5H zXd_WI1;o?1#5|Er^@*w|qmD#XemVYz$rzH5gd0nU%^}|fPK64@YJx*<2FL2VP(*^m z{5Ck8wNMC=rDu4^YL^l#k~{s>d^hI`DLoJWv&DC?C)1 zRCmTfzJVbkk|HYXba>!c`44?1DFa%%vWHFG`jI6$&orSQjZ4mxcVTTGskIJU#1-mD zMilqgeJSWvvO|6poF1DT#aF;_G@+~Dh72&KFs3(LET^dQF68l;p*H&Sg`d-xitK@U z&$Z~4%lsVT0_|)9r;jJxFW#s6)LbXlJl=VvKmcop1efxXAz?LVo3zB#1i)*O05FTS z^zSIkA<85Z6{N+{eQCi|mKHpXt04fibTeh?;Q+S+UulUc=ilb1=&{D%4$6rsmtz*7 zl@n=kBk&^S^8<+~m%@`;Vj2Jx=&VM$9P5Ran3||pdUOQmV&-azX#f?ZM{8;kLz>P) zy-;Gx7IU|_);$0xrp({#r}u$!Vj4)#rAOCtsgtvi{u|1z%+#d+nX;Z@e?C(gDEISB zxxB*9Gi5!MI6-Pzc+mv-z{TlbQFfVEEB$B6ist$AnQ{o#0US1u0l~lufaO*JoS3pf ztEFUS`9;qlKWzZTkwaPsaK%PIqO%#gznYmtx|y=NjQ~rx0?cm$I5Fk?Edb|l1vvdH z%5vKPZhfbpwt;d2)&0d8+`ktAuJ9{>)6J9{dP&RGDf0&amVX&w{uO`|Q|1r(>0wZo zJL;#eg0lZ#lWdZNGn8}^;B+(XLcQ#9wQTo?0IT~5;Pl^6mOle9JqsA=|G6q;Y6^1w zwf`OEhJIQWbcg@6ENB7q)3Tt4!~cF+xcz^!EVO^)cqfHoT@xs##5&=j30F^#IeO3a z`}@hK-&^rU)8XOI%(O%e9x><6(AtC7k`H%nI-B}J_wUyKMt#_t*LHCE`~C0DJo>wx zO6i)ZTP~HnunP--kgkVa2Yrl=xC-evC00r*8IL3K+>%0?UuvcIz%8YO(n9J#%Sww& z$BX533fu?avSy9PzskBP5Y5@(0;F|(Mw6ikU+@}ZrZ6oJF z(VkefM4+kXJTq~^jPNeBHX(LV?E5AWwQ&wSU^V{izug!A-BD|XXm6QUt>6(EAU~+L7j185?{c(Su zvAOcjwG@DSl*BtNE8Oq5D`YJ-)3_A40lN)8y8&+iCx8RML4a++pG5vj$se(#D2iy9 zqom`F_r*N`ZyNYh%pWSGlQikc=vv3}9Vsa6P~=$vRg79C0gv<-ijF*Js<$vtdYg z?V;v%qq5isF_;C$hp6}<9)d{Z@Gy4+JwJK&%0EwR@Wco4K!+m7N7vRVHg41c?s(9E zm&d2qeHkpKXQbjcNNS*r^>!O~8lmpnSak8CbMz?*<%3ad!$rwMmXIt}G~Gqupj)WX%87R}S-# z{2eHV&r{JUO^406*vNS*r6J*1B)E${Fj7Z*K(h9Syb7)fm{1>!P&2iZ?a zX|;=DDrjZTm9{3!UxnzO{FhQA@;;?E{m;@^<1S%%_s+FL4h6GD?g}*GcHyo~nKt7# zVPe|RP0I0O-5@YXZ*mfC++?>Iw+uyd6Ne4&8kXZ%0Q+arrA>AWTEb>UoTXemKc^|1 z?Kb0v;mZ$i>E62i531G)J2&a2&71d%y)8~lZPLp=Z{X~ zZenS1L)*a*J9z%G0E}5fr_hUQbm1Ah*hX>N?Kb1WWXb~xVHZ2-ECib9>JFY< zx9_xzeN=#IxNJ$cWUZz=3mD8l+++!*D*6)WGj;QdEY!cCy$ZR)jK5uN1NU219B z1BWA>-$O&r*yz~qBpefQO8C?6^nQzD^~td+z`Md;A@q_v$*6^8g%FC-O+%uhC_neBf8~N?ap7-{RD0X|iX)aHBT9LaVjXgna zFW4gaB`rncQ~bVgz82ped-~4-mb?88R-U8qj}rSBmtMx%!#HW|_lu3~)cR8X(L$=k pfxczlt(N(ny2&d$p`kCup9${y!`dJCMXqq|app(_A%E~)^ADFPZczXL delta 7790 zcmeHMdtB93w*Tz|2l*if5MTT_DuVbx1vwn>ARs?q6@gMTB081|a=?QikAttYz>Ku< zk+7hgd>}qje5F!mQ)WJK$;$qIsn04#Gif`8WfAZp+d&~dP{?&=#-;db$vrA_m zZnF5^yb`+hV@WHG^?4CB=NfcI{|%6&+GafV0(E;!J&UR(DZxjQSl$F00$S=`vWT_U z8S-TYeGF7aeWBZ35$o|vL%;`v4>#m>R-NwvWqa2^gQnsk2nvCq9TrI%3hMF37MGS6 zNV9wBiq|3U$@1bwS}kPw9vhL3A=v|Le_f$dhDj#ySUm@>XqEy%Ej)g<-ui`i~LC)Ui+cf#8+EozrfuPV`u(%isWnS)1l9oWu9c~KJJNzt&ES6#8hQmon zY8T{tW8q`*0Pw8#Cs3HIO~+iYANeJk+tQoh!@;*vzNLRnA&iEC%L8TY$wucxK?i{k zGvt}@sC&I1hyZsz$0c(`u1jcjS?Qy24pdo~MGU{P}K7K?-{1k=BB8yx#}q`Z9M_ zu{+`WtSUo0&YPlqZ*?W2Q{`Ul^%O0a-heke|7+kaJJJpsoQ8*oj0RqJ*%Eh^Bu$5> zJfbjARMq|+MtB51;JHJjQ+2h+UG*e6eEZj&8=|{a0*X-9uE*e6i<1U1t}Ov&z7i{g z)8aTNcfSCX&zUn5l;f7!CR^|E_5Z_odao_pbA}ekuZtrT6QIf^n8GAz z$m9rci8N{rP{khV2v99w`q0Rp>0%&xda7a;wc@vlI`DguVggm!YL%qB(NIujPnXD| z)<9J>Qb(X_`HPiC+S0`s^4L`2rB<6NxA|dCG-Zd)C5BQ=kSg5d2~w>sJtQfeVgpkx zH+oQKP`cb38$&u&{ixFBl5@a~1!twMV3&0p*HKyjG~rL3!RfMsU?!WjB8NO7s@wn> zx0k3g#3f$>mjTY7+sjc1I?ui@!3}v1IJWIWT|qAR-r7r*FG9wxeON=-C?-^ue}YJ} zo`oFHdy-m1RdJp=P-DfUa_|DEE7T<>lc%?8t%odyz6?sW{4tO^d#B5hn1~5b@~6t) zF1Z9;3OHYGBJT!wzs?1_EMM8EvroD_3^O(wLi7if2f(pMJ*X?dB{zfPs9UHjHrXEx zQilhC!6jlth)DG|(j{Vc|@HW)IE3<>1%~{DbCE zaNN1>?SF#f{@@sFW?)a_Ux1T}ho@DvTfR>sMuXZRY5iBhfavW&`kb_8yzSvZGG7*WMF8N+?>6#{zNuZBt)$(Bkl|-kDKx&Ov#U$#8R%IG|+ht72 zZ^0RxhutNQh?FE1CXlOUo9{u6NSAdVI6XMu;tAP=Fl9&Sfk^Ym1N)`L!E!x{ zMh;Jx`(aa_t~-PfZ=}}Ys(6n&hO084Q9SiX91L+RxEWew9#j)J-4sv#hv3+T%t<0* zDJDi0E6Ic3Gt`P-AL_ubi(*EqqKrHvRVzG}f@oE_lLn8tSJQ?Z)fs(j6#wEol{kE{qx8i@TYq$#g1cV0?e}hx1YE{4P+| zx(DFGlsm&duN95bwD<^z5BRBulqqXXGkB(K9T(_o|3ul)d_&F@ zXP{JS@JzYB%;1@_opOUO=l)rt!Vuhv`k;>Ab+~k=Y-q72)hP~6sTSZ=dIInVHUg}- z3E;w%`>SULrJl#nf(C;UDBd-s%>Xxi7U05^^|lyvt3kJca_LT4Z#%&4cL2<{09=@| zd>6p-n%(%}awp1)&jIYv0fQa^<ruTe*)mropMJnYf_yu{|f;6HPWkiV8Ka% z3sdG_Gw3N$HhkKkXF=Jaw*W3odjJ=iLAeuU{VM?5{Q%&iQ+mzmN3S~lP1LZ;ReO0*<+*Jk}UBQwUp#hMv0BA z3$sK$%HEBSVzxC_B`rF(~w4?Newnpz-A&P0>;`w-yT(x+$jnBM(INB~o!-~;5L(s~4 zN8A0&ejC=dvc8v3fIlBxD2+sKHKyTI`}RA%uy6yw;o*~%i^=o0&M#qH_@p#-;Ez_G zGGxA>hjr2I6O^~~Ck+{&rreg#p`QW1Y`{EGUpG7pe43ta%8cRT;|<_VfX|>;0QLc) z;Y$_1SLACasSnwnF0bKJ?kK>!KcDY`03SH#fwzHofOmll03T!(KoPJ2@Bj;e0w5pY zx1ooCIRL*&fZv$r0{luhli!YV@xbqgcL8?;6M<9! z|2;`-0Gu$WPGAS{EU*jM0`Q9sAN^in zo`N5q&?cZBNM>s&jX%d+Mzy7F4)e)=0(>3N0I&}0vr}sUc97F>6DTLsMqmT54p}j5pxpb z61bm}v}&7TxhiPuwy*)_g{R_|7i}N0q!DWIapPF+mY_@9!jfhfhl%> zciA1uyoB{mU6?y>@?&*c(}XxDZtiH{cExU9;-)yVriZV(aS-(wXPg7`O*z{Y++i%l zZz0ufpC*=2`*x?@ynEi*axC{Bk6g^ve2(Y7nkcea5iOM3Tx>V5yp^}=14B0v{P5J({kNH!*(PjnYZ#Y!ZR*MbhW{K(+Tt5e$hAMzMfX{fL-fsY#bI( zIGuV<5kqLyK1C$Z=^ak7jOtqy3Tg?ro1Y0@5Bj&(d)f*)tPX@V33)@4Tby?DGsCi% z|LU@)+{?}(=y78uNwb0N{6LGM+Y@zkqeZb?v(lygVRrMo!ubs=gZf)OSY{CzI0s9b zO|d(jI3x0RDx!f_?3`vdKN5`DoBsMA>Ixn~Cs5@$h0*VJP7~+Iu}cxBNbK%Ib9Y6E z6siHw(Nm~!S4htg@eVetaAZk7T`)9t{o(=U{r-f$H#Vo7@@LyT0NBQj|Kqz8?B)%B z_Lws}mAB9KfkL9*`Ea_t+iCw(h<@W*P#HabLeJo2!vIE-P7!;Y|I_rLzKzQEDB?|8 zv&V@m#GLJY=;9u$$fOVUyok81d(J8JnR$jj#PbVe+pE}5V=ZAUwKu;m_+)Nx|3u9j z&E5NG;$8)R?U2a%8=IX5f*ZvpF%KLs0(xI`Xu zAoW5HXqjNqmY?9o(>H=UFfYjKUpjl%6XjU@{X$HlF|CT2PKWl5vYX5Chs!oZRMTJf zMT#r5YQG}BZ=qOjhSR=7BgIsDsx{Yceim@3yN9k^-76ffIpdsql75FybO*Z|(Pr5& zn?qanD+A0=46{!C(Yp2emm|KP5Iyyh%ufL83LB5U`qB9th5}ZC`8DCx=q(?m=7wcz zkal@v z?;=WXi{9l*tkTgl44O`g3)lMt+-iyVQuIaQ5L~xgVTkVf^JsK zk3`KwYe$c&a~H&OyvFhdH;)`g6#M^D-}{G0RFq#LEj^M-?;gpVe#<4nz=rVM@^u^KQ zfX|=D&`?Gkp$CsmpdH8J!)FiBm*(TiN$-kfLEp|iqO08ch&qqWZ(G{7a2VB{u(zon zec=<3d`+*+xJHTBqX#i#exor@+D(Sazh6Tac51r5I$-d}daL-4xfS~RVVh-!!rG6u NHZR^7BlvET{tK0; { + return await this.message.createMessage(content, creator, created) + } + + async placeMessage(message: MessageID, card: CardID, workspace: string): Promise { + return await this.message.placeMessage(message, card, workspace) + } + + async createPatch(message: MessageID, content: RichText, creator: SocialID, created: Date): Promise { + return await this.message.createPatch(message, content, creator, created) + } + + async removeMessage(message: MessageID): Promise { + return await this.message.removeMessage(message) + } + + async createReaction(message: MessageID, reaction: string, creator: SocialID, created: Date): Promise { + return await this.message.createReaction(message, reaction, creator, created) + } + + async removeReaction(message: MessageID, reaction: string, creator: SocialID): Promise { + return await this.message.removeReaction(message, reaction, creator) + } + + async createAttachment(message: MessageID, card: CardID, creator: SocialID, created: Date): Promise { + return await this.message.createAttachment(message, card, creator, created) + } + + async removeAttachment(message: MessageID, card: CardID): Promise { + return await this.message.removeAttachment(message, card) + } + + async findMessages(workspace: string, params: FindMessagesParams): Promise { + return await this.message.find(workspace, params) + } + + async createNotification(message: MessageID, context: ContextID): Promise { + return await this.notification.createNotification(message, context) + } + + async removeNotification(message: MessageID, context: ContextID): Promise { + return await this.notification.removeNotification(message, context) + } + + async createContext( + workspace: string, + card: CardID, + personWorkspace: string, + lastView?: Date, + lastUpdate?: Date + ): Promise { + return await this.notification.createContext(workspace, card, personWorkspace, lastView, lastUpdate) + } + + async updateContext(context: ContextID, update: NotificationContextUpdate): Promise { + return await this.notification.updateContext(context, update) + } + + async removeContext(context: ContextID): Promise { + return await this.notification.removeContext(context) + } + + async findContexts( + params: FindNotificationContextParams, + personWorkspaces: string[], + workspace?: string + ): Promise { + return await this.notification.findContexts(params, personWorkspaces, workspace) + } + + async findNotifications( + params: FindNotificationsParams, + personWorkspace: string, + workspace?: string + ): Promise { + return await this.notification.findNotifications(params, personWorkspace, workspace) + } + + close(): void { + this.db.close() + } +} + +export async function createDbAdapter(connectionString: string): Promise { + const db = connect(connectionString) + const sqlClient = await db.getClient() + + return new CockroachAdapter(db, sqlClient) +} diff --git a/packages/cockroach/src/connection.ts b/packages/cockroach/src/connection.ts new file mode 100644 index 0000000000..a9aa16e74b --- /dev/null +++ b/packages/cockroach/src/connection.ts @@ -0,0 +1,104 @@ +//Full copy from @hcengineering/postgres +import postgres from 'postgres' +import { v4 as uuid } from 'uuid' + +const connections = new Map() +const clientRefs = new Map() + +export interface PostgresClientReference { + getClient: () => Promise + close: () => void +} + +class PostgresClientReferenceImpl { + count: number + client: postgres.Sql | Promise + + constructor( + client: postgres.Sql | Promise, + readonly onclose: () => void + ) { + this.count = 0 + this.client = client + } + + async getClient(): Promise { + if (this.client instanceof Promise) { + this.client = await this.client + } + return this.client + } + + close(force: boolean = false): void { + this.count-- + if (this.count === 0 || force) { + if (force) { + this.count = 0 + } + void (async () => { + this.onclose() + const cl = await this.client + await cl.end() + console.log('Closed postgres connection') + })() + } + } + + addRef(): void { + this.count++ + console.log('Add postgres connection', this.count) + } +} + +export class ClientRef implements PostgresClientReference { + id = uuid() + constructor(readonly client: PostgresClientReferenceImpl) { + clientRefs.set(this.id, this) + } + + closed = false + async getClient(): Promise { + if (!this.closed) { + return await this.client.getClient() + } else { + throw new Error('DB client-query is already closed') + } + } + + close(): void { + // Do not allow double close of connection client-query + if (!this.closed) { + clientRefs.delete(this.id) + this.closed = true + this.client.close() + } + } +} + +export function connect(connectionString: string, database?: string): PostgresClientReference { + const extraOptions = JSON.parse(process.env.POSTGRES_OPTIONS ?? '{}') + const key = `${connectionString}${extraOptions}` + let existing = connections.get(key) + + if (existing === undefined) { + const sql = postgres(connectionString, { + connection: { + application_name: 'communication' + }, + database, + max: 10, + transform: { + undefined: null + }, + ...extraOptions + }) + + existing = new PostgresClientReferenceImpl(sql, () => { + connections.delete(key) + }) + connections.set(key, existing) + } + // Add reference and return once closable + existing.addRef() + return new ClientRef(existing) +} diff --git a/packages/cockroach/src/db/base.ts b/packages/cockroach/src/db/base.ts new file mode 100644 index 0000000000..671ffc80dc --- /dev/null +++ b/packages/cockroach/src/db/base.ts @@ -0,0 +1,41 @@ +import type postgres from 'postgres' + +export class BaseDb { + constructor( + readonly client: postgres.Sql + ) {} + + async insert(table: string, data: Record): Promise { + const keys = Object.keys(data) + const values = Object.values(data) + const sql = ` + INSERT INTO ${table} (${keys.map((k) => `"${k}"`).join(', ')}) + VALUES (${keys.map((_, idx) => `$${idx + 1}`).join(', ')}); + ` + await this.client.unsafe(sql, values) + } + + async insertWithReturn(table: string, data: Record, returnField : string): Promise { + const keys = Object.keys(data) + const values = Object.values(data) + const sql = ` + INSERT INTO ${table} (${keys.map((k) => `"${k}"`).join(', ')}) + VALUES (${keys.map((_, idx) => `$${idx + 1}`).join(', ')}) + RETURNING ${returnField};` + const result =await this.client.unsafe(sql, values) + + return result[0][returnField] + } + + async remove(table: string, where: Record): Promise { + const keys = Object.keys(where) + const values = Object.values(where) + + const sql = ` + DELETE + FROM ${table} + WHERE ${keys.map((k, idx) => `"${k}" = $${idx + 1}`).join(' AND ')};` + + await this.client.unsafe(sql, values) + } +} diff --git a/packages/cockroach/src/db/message.ts b/packages/cockroach/src/db/message.ts new file mode 100644 index 0000000000..2722f68bef --- /dev/null +++ b/packages/cockroach/src/db/message.ts @@ -0,0 +1,228 @@ +import { + type Message, + type MessageID, + type CardID, + type FindMessagesParams, + SortOrder, + type SocialID, + type RichText, + Direction, type Reaction, type Attachment +} from '@communication/types' + +import {BaseDb} from './base.ts' +import { + TableName, + type MessageDb, + type MessagePlaceDb, + type AttachmentDb, + type ReactionDb, + type PatchDb +} from './types.ts' + +export class MessagesDb extends BaseDb { + //Message + async createMessage(content: RichText, creator: SocialID, created: Date): Promise { + const dbData: MessageDb = { + content: content, + creator: creator, + created: created, + } + + const id = await this.insertWithReturn(TableName.Message, dbData, 'id') + + return id as MessageID + } + + async removeMessage(message: MessageID): Promise { + await this.remove(TableName.Message, {id: message}) + } + + async placeMessage(message: MessageID, card: CardID, workspace: string): Promise { + const dbData: MessagePlaceDb = { + workspace_id: workspace, + card_id: card, + message_id: message + } + await this.insert(TableName.MessagePlace, dbData) + } + + async createPatch(message: MessageID, content: RichText, creator: SocialID, created: Date): Promise { + const dbData: PatchDb = { + message_id: message, + content: content, + creator: creator, + created: created + } + + await this.insert(TableName.Patch, dbData) + } + + //Attachment + async createAttachment(message: MessageID, card: CardID, creator: SocialID, created: Date): Promise { + const dbData: AttachmentDb = { + message_id: message, + card_id: card, + creator: creator, + created: created + } + await this.insert(TableName.Attachment, dbData) + } + + async removeAttachment(message: MessageID, card: CardID): Promise { + await this.remove(TableName.Attachment, { + message_id: message, + card_id: card + }) + } + + //Reaction + async createReaction(message: MessageID, reaction: string, creator: SocialID, created: Date): Promise { + const dbData: ReactionDb = { + message_id: message, + reaction: reaction, + creator: creator, + created: created + } + await this.insert(TableName.Reaction, dbData) + } + + async removeReaction(message: MessageID, reaction: string, creator: SocialID): Promise { + await this.remove(TableName.Reaction, { + message_id: message, + reaction: reaction, + creator: creator + }) + } + + //Find messages + async find(workspace: string, params: FindMessagesParams): Promise { + //TODO: experiment with select to improve performance + const select = `SELECT m.id, + m.content, + m.creator, + m.created, + ${this.subSelectPatches()}, + ${this.subSelectAttachments()}, + ${this.subSelectReactions()} + FROM ${TableName.Message} m + INNER JOIN ${TableName.MessagePlace} mp ON m.id = mp.message_id` + + const {where, values} = this.buildMessageWhere(workspace, params) + const orderBy = params.sort ? `ORDER BY m.created ${params.sort === SortOrder.Asc ? 'ASC' : 'DESC'}` : '' + const limit = params.limit ? ` LIMIT ${params.limit}` : '' + const sql = [select, where, orderBy, limit].join(' ') + + const result = await this.client.unsafe(sql, values) + + return result.map(it => this.toMessage(it)) as Message[] + } + + buildMessageWhere(workspace: string, params: FindMessagesParams): { where: string, values: any[] } { + const where: string[] = ['mp.workspace_id = $1'] + const values: any[] = [workspace] + let index = 2 + for (const key of Object.keys(params)) { + const value = (params as any)[key] + switch (key) { + case 'id': { + where.push(`m.id = $${index++}`) + values.push(value) + break + } + case 'card': { + where.push(`mp.card_id = $${index++}`) + values.push(value) + break + } + case 'from': { + const exclude = params.excluded ?? false + const direction = params.direction ?? Direction.Forward + const getOperator = () => { + if (exclude) { + return direction === Direction.Forward ? '>' : '<' + } else { + return direction === Direction.Forward ? '>=' : '<=' + } + } + + where.push(`m.created ${getOperator()} $${index++}`) + values.push(value) + break + } + } + } + + return {where: `WHERE ${where.join(' AND ')}`, values} + } + + subSelectPatches(): string { + return `array( + SELECT jsonb_build_object( + 'content', p.content, + 'creator', p.creator, + 'created', p.created + ) + FROM ${TableName.Patch} p + WHERE p.message_id = m.id + ) AS patches` + } + + subSelectAttachments(): string { + return `array( + SELECT jsonb_build_object( + 'card_id', a.card_id, + 'message_id', a.message_id, + 'creator', a.creator, + 'created', a.created + ) + FROM ${TableName.Attachment} a + WHERE a.message_id = m.id + ) AS attachments` + } + + subSelectReactions(): string { + return `array( + SELECT jsonb_build_object( + 'message_id', r.message_id, + 'reaction', r.reaction, + 'creator', r.creator, + 'created', r.created + ) + FROM ${TableName.Reaction} r + WHERE r.message_id = m.id + ) AS reactions` + } + + toMessage(row: any): Message { + const lastPatch = row.patches?.[0] + + return { + id: row.id, + content: lastPatch?.content ?? row.content, + creator: row.creator, + created: new Date(row.created), + edited: new Date(lastPatch?.created ?? row.created), + reactions: (row.reactions ?? []).map(this.toReaction), + attachments: (row.attachments ?? []).map(this.toAttachment) + } + } + + toReaction(row: any): Reaction { + return { + message: row.message_id, + reaction: row.reaction, + creator: row.creator, + created: new Date(row.created) + } + } + + toAttachment(row: any): Attachment { + return { + message: row.message_id, + card: row.card_id, + creator: row.creator, + created: new Date(row.created) + } + } +} + diff --git a/packages/cockroach/src/db/notification.ts b/packages/cockroach/src/db/notification.ts new file mode 100644 index 0000000000..dfa9528d12 --- /dev/null +++ b/packages/cockroach/src/db/notification.ts @@ -0,0 +1,239 @@ +import { + type MessageID, + type ContextID, + type CardID, + type NotificationContext, + type FindNotificationContextParams, SortOrder, + type FindNotificationsParams, type Notification, + type NotificationContextUpdate +} from '@communication/types' + +import {BaseDb} from './base.ts' +import {TableName, type ContextDb, type NotificationDb} from './types.ts' + +export class NotificationsDb extends BaseDb { + async createNotification(message: MessageID, context: ContextID): Promise { + const dbData: NotificationDb = { + message_id: message, + context + } + await this.insert(TableName.Notification, dbData) + } + + async removeNotification(message: MessageID, context: ContextID): Promise { + await this.remove(TableName.Notification, { + message_id: message, + context + }) + } + + async createContext(workspace: string, card: CardID, personWorkspace: string, lastView?: Date, lastUpdate?: Date): Promise { + const dbData: ContextDb = { + workspace_id: workspace, + card_id: card, + person_workspace: personWorkspace, + last_view: lastView, + last_update: lastUpdate + } + return await this.insertWithReturn(TableName.NotificationContext, dbData, 'id') as ContextID + } + + async removeContext(context: ContextID): Promise { + await this.remove(TableName.NotificationContext, { + id: context + }) + } + + async updateContext(context: ContextID, update: NotificationContextUpdate): Promise { + const dbData: Partial = {} + + if (update.archivedFrom != null) { + dbData.archived_from = update.archivedFrom + } + if (update.lastView != null) { + dbData.last_view = update.lastView + } + if (update.lastUpdate != null) { + dbData.last_update = update.lastUpdate + } + + if (Object.keys(dbData).length === 0) { + return + } + + const keys = Object.keys(dbData) + const values = Object.values(dbData) + + const sql = `UPDATE ${TableName.NotificationContext} + SET ${keys.map((k, idx) => `"${k}" = $${idx + 1}`).join(', ')} + WHERE id =$${keys.length + 1}` + + await this.client.unsafe(sql, [values, context]) + } + + async findContexts( params: FindNotificationContextParams, personWorkspaces: string[], workspace?: string,): Promise { + const select = ` + SELECT nc.id, nc.card_id, nc.archived_from, nc.last_view, nc.last_update + FROM ${TableName.NotificationContext} nc`; + const {where, values} = this.buildContextWhere(params, personWorkspaces, workspace) + // const orderSql = `ORDER BY nc.created ${params.sort === SortOrder.Asc ? 'ASC' : 'DESC'}` + const limit = params.limit ? ` LIMIT ${params.limit}` : '' + const sql = [select, where, limit].join(' ') + + const result = await this.client.unsafe(sql, values); + + return result.map(this.toNotificationContext); + } + + + async findNotifications(params: FindNotificationsParams, personWorkspace: string, workspace?: string): Promise { + //TODO: experiment with select to improve performance, should join with attachments and reactions? + const select = ` + SELECT n.message_id, + n.context, + m.content AS message_content, + m.creator AS message_creator, + m.created AS message_created, + nc.card_id, + nc.archived_from, + nc.last_view, + nc.last_update, + (SELECT json_agg( + jsonb_build_object( + 'id', p.id, + 'content', p.content, + 'creator', p.creator, + 'created', p.created + ) + ) + FROM ${TableName.Patch} p + WHERE p.message_id = m.id) AS patches + FROM ${TableName.Notification} n + JOIN ${TableName.NotificationContext} nc ON n.context = nc.id + JOIN ${TableName.Message} m ON n.message_id = m.id + `; + const {where, values} = this.buildNotificationWhere(params, personWorkspace, workspace) + const orderBy = params.sort ? `ORDER BY m.created ${params.sort === SortOrder.Asc ? 'ASC' : 'DESC'}` : '' + const limit = params.limit ? ` LIMIT ${params.limit}` : '' + const sql = [select, where, orderBy, limit].join(' ') + + const result = await this.client.unsafe(sql, values); + + return result.map(this.toNotification); + } + + buildContextWhere(params: FindNotificationContextParams, personWorkspaces: string[], workspace?: string,): { + where: string, + values: any[] + } { + const where: string[] = [] + const values: any[] = [] + let index = 1 + + if(workspace != null) { + where.push(`nc.workspace_id = $${index++}`) + values.push(workspace) + } + + if(personWorkspaces.length > 0) { + where.push(`nc.person_workspace IN (${personWorkspaces.map((it) => `$${index++}`).join(', ')})`) + values.push(...personWorkspaces) + } + + for (const key of Object.keys(params)) { + const value = (params as any)[key] + switch (key) { + case 'card': { + where.push(`nc.card_id = $${index++}`) + values.push(value) + break + } + } + } + + return {where: `WHERE ${where.join(' AND ')}`, values} + } + + buildNotificationWhere(params: FindNotificationsParams, personWorkspace: string, workspace?: string): { + where: string, + values: any[] + } { + const where: string[] = ['nc.person_workspace = $1'] + const values: any[] = [personWorkspace] + let index = 2 + + if(workspace != null) { + where.push(`nc.workspace_id = $${index++}`) + values.push(workspace) + } + + for (const key of Object.keys(params)) { + const value = (params as any)[key] + switch (key) { + case 'context': { + where.push(`n.context = $${index++}`) + values.push(value) + break + } + case 'card': { + where.push(`nc.card_id = $${index++}`) + values.push(value) + break + } + case 'read': { + if (value === true) { + where.push(`nc.last_view IS NOT NULL AND nc.last_view >= m.created`) + } else if (value === false) { + where.push(`(nc.last_view IS NULL OR nc.last_view > m.created)`) + } + break + } + case 'archived': { + if (value === true) { + where.push(`nc.archived_from IS NOT NULL AND nc.archived_from >= m.created`) + } else if (value === false) { + where.push(`(nc.archived_from IS NULL OR nc.archived_from > m.created)`) + } + break + } + } + } + + return {where: `WHERE ${where.join(' AND ')}`, values} + } + + toNotificationContext(row: any): NotificationContext { + return { + id: row.id, + card: row.card_id, + workspace: row.workspace_id, + personWorkspace: row.person_workspace, + archivedFrom: row.archived_from ? new Date(row.archived_from) : undefined, + lastView: row.last_view ? new Date(row.last_view) : undefined, + lastUpdate: row.last_update ? new Date(row.last_update) : undefined + } + } + + toNotification(row: any): Notification { + const lastPatch = row.patches?.[0] + const lastView = row.last_view ? new Date(row.last_view) : undefined + const archivedFrom = row.archived_from ? new Date(row.archived_from) : undefined + const created = new Date(row.message_created) + + return { + message: { + id: row.id, + content: lastPatch?.content ?? row.message_content, + creator: row.message_creator, + created, + edited: new Date(lastPatch?.created ?? row.message_created), + reactions: row.reactions ?? [], + attachments: row.attachments ?? [] + }, + context: row.context, + read: lastView != null && lastView >= created, + archived: archivedFrom != null && archivedFrom >= created + } + } +} + diff --git a/packages/cockroach/src/db/types.ts b/packages/cockroach/src/db/types.ts new file mode 100644 index 0000000000..9dab08561a --- /dev/null +++ b/packages/cockroach/src/db/types.ts @@ -0,0 +1,59 @@ +import type {CardID, ContextID, MessageID, RichText, SocialID } from "@communication/types" + +export enum TableName { + Message = 'message', + Patch = 'patch', + MessagePlace = 'message_place', + Attachment = 'attachment', + Reaction = 'reaction', + Notification = 'notification', + NotificationContext = 'notification_context' +} + +export interface MessageDb { + content: RichText, + creator: SocialID, + created: Date, +} + +export interface PatchDb { + message_id: MessageID, + content: RichText, + creator: SocialID, + created: Date, +} + +export interface MessagePlaceDb { + workspace_id: string, + card_id: CardID, + message_id: MessageID +} + +export interface ReactionDb { + message_id: MessageID, + reaction: string, + creator: SocialID + created: Date +} + +export interface AttachmentDb { + message_id: MessageID, + card_id: CardID, + creator: SocialID + created: Date +} + +export interface NotificationDb { + message_id: MessageID, + context: ContextID +} + +export interface ContextDb { + workspace_id: string + card_id: CardID + person_workspace: string + + archived_from?: Date + last_view?: Date + last_update?: Date +} \ No newline at end of file diff --git a/packages/cockroach/src/index.ts b/packages/cockroach/src/index.ts new file mode 100644 index 0000000000..03eeab5ffa --- /dev/null +++ b/packages/cockroach/src/index.ts @@ -0,0 +1 @@ +export * from './adapter.ts' diff --git a/packages/postgres/migrations/01_message.sql b/packages/postgres/migrations/01_message.sql deleted file mode 100644 index 8aa19866f9..0000000000 --- a/packages/postgres/migrations/01_message.sql +++ /dev/null @@ -1,19 +0,0 @@ -CREATE TABLE IF NOT EXISTS message -( - id INT8 NOT NULL DEFAULT unique_rowid(), - content TEXT, - version INTEGER NOT NULL, - creator VARCHAR(255) NOT NULL, - created TIMESTAMPTZ NOT NULL DEFAULT now(), - - PRIMARY KEY (id, version) -); - -CREATE TABLE IF NOT EXISTS message_place -( - workspace_id UUID NOT NULL, - card_id UUID NOT NULL, - message_id INT8 NOT NULL, - - PRIMARY KEY (workspace_id, card_id, message_id) -); diff --git a/packages/postgres/migrations/03_reaction.sql b/packages/postgres/migrations/03_reaction.sql deleted file mode 100644 index 5dc21091f7..0000000000 --- a/packages/postgres/migrations/03_reaction.sql +++ /dev/null @@ -1,11 +0,0 @@ -CREATE TABLE IF NOT EXISTS reaction -( - message_id INT8 NOT NULL, - reaction INTEGER NOT NULL, - creator VARCHAR(255) NOT NULL, - created TIMESTAMPTZ NOT NULL DEFAULT now(), - - PRIMARY KEY (message_id, creator, reaction) -); - -CREATE INDEX IF NOT EXISTS reaction_message_idx ON reaction (message_id); diff --git a/packages/postgres/migrations/04_notification.sql b/packages/postgres/migrations/04_notification.sql deleted file mode 100644 index 96b53f1a6d..0000000000 --- a/packages/postgres/migrations/04_notification.sql +++ /dev/null @@ -1,7 +0,0 @@ -CREATE TABLE IF NOT EXISTS notification -( - social_id VARCHAR(255) NOT NULL, - message_id INT8 NOT NULL, - - PRIMARY KEY (social_id, message_id) -); diff --git a/packages/postgres/migrations/05_notificationContext.sql b/packages/postgres/migrations/05_notificationContext.sql deleted file mode 100644 index f8936f24cc..0000000000 --- a/packages/postgres/migrations/05_notificationContext.sql +++ /dev/null @@ -1,12 +0,0 @@ -CREATE TABLE IF NOT EXISTS notification_context -( - workspace_id UUID NOT NULL, - card_id UUID NOT NULL, - huly_id VARCHAR(255) NOT NULL, /* Or maybe account id or something else */ - - archived_from TIMESTAMPTZ, - last_view TIMESTAMPTZ, - last_update TIMESTAMPTZ, - - PRIMARY KEY (workspace_id, card_id, huly_id) -); diff --git a/packages/postgres/src/index.ts b/packages/postgres/src/index.ts deleted file mode 100644 index e69de29bb2..0000000000 diff --git a/packages/sdk-types/package.json b/packages/sdk-types/package.json new file mode 100644 index 0000000000..eccea9085e --- /dev/null +++ b/packages/sdk-types/package.json @@ -0,0 +1,16 @@ +{ + "name": "@communication/sdk-types", + "version": "0.1.0", + "main": "src/index.ts", + "module": "src/index.ts", + "type": "module", + "devDependencies": { + "@types/bun": "^1.1.14" + }, + "dependencies": { + "@communication/types": "workspace:*" + }, + "peerDependencies": { + "typescript": "^5.6.3" + } +} diff --git a/packages/sdk-types/src/db.ts b/packages/sdk-types/src/db.ts new file mode 100644 index 0000000000..0e5c9c38d3 --- /dev/null +++ b/packages/sdk-types/src/db.ts @@ -0,0 +1,54 @@ +import type { + CardID, + ContextID, + FindMessagesParams, + FindNotificationContextParams, + FindNotificationsParams, + Message, + MessageID, + NotificationContext, + NotificationContextUpdate, + RichText, + SocialID, + Notification +} from '@communication/types' + +export interface DbAdapter { + createMessage(content: RichText, creator: SocialID, created: Date): Promise + removeMessage(id: MessageID): Promise + + placeMessage(message: MessageID, card: CardID, workspace: string): Promise + createPatch(message: MessageID, content: RichText, creator: SocialID, created: Date): Promise + + createReaction(message: MessageID, reaction: string, creator: SocialID, created: Date): Promise + removeReaction(message: MessageID, reaction: string, creator: SocialID): Promise + + createAttachment(message: MessageID, card: CardID, creator: SocialID, created: Date): Promise + removeAttachment(message: MessageID, card: CardID): Promise + + findMessages(workspace: string, query: FindMessagesParams): Promise + + createNotification(message: MessageID, context: ContextID): Promise + removeNotification(message: MessageID, context: ContextID): Promise + createContext( + personWorkspace: string, + workspace: string, + card: CardID, + lastView?: Date, + lastUpdate?: Date + ): Promise + updateContext(context: ContextID, update: NotificationContextUpdate): Promise + removeContext(context: ContextID): Promise + findContexts( + params: FindNotificationContextParams, + personWorkspaces: string[], + workspace?: string + ): Promise + findNotifications( + params: FindNotificationsParams, + personWorkspace: string, + workspace?: string + ): Promise + + close(): void +} diff --git a/packages/sdk-types/src/index.ts b/packages/sdk-types/src/index.ts new file mode 100644 index 0000000000..1beb455f5e --- /dev/null +++ b/packages/sdk-types/src/index.ts @@ -0,0 +1 @@ +export * from './db' diff --git a/packages/sdk-types/tsconfig.json b/packages/sdk-types/tsconfig.json new file mode 100644 index 0000000000..49e05cea1e --- /dev/null +++ b/packages/sdk-types/tsconfig.json @@ -0,0 +1,8 @@ +{ + "extends": "../../tsconfig.json", + "compilerOptions": { + "outDir": "./dist", + "rootDir": "./src" + }, + "include": ["src"] +}