From 5b54f3f6360de72c6ad5d4dee7c3b505e5ed12b9 Mon Sep 17 00:00:00 2001 From: Andrew Gundersen Date: Sat, 21 Nov 2020 14:43:17 -0600 Subject: [PATCH] working v1 --- __pycache__/agent.cpython-38.pyc | Bin 0 -> 2536 bytes __pycache__/mail.cpython-38.pyc | Bin 0 -> 1488 bytes __pycache__/mailroom.cpython-38.pyc | Bin 0 -> 968 bytes __pycache__/message.cpython-38.pyc | Bin 0 -> 732 bytes __pycache__/messenger.cpython-38.pyc | Bin 0 -> 1021 bytes __pycache__/schemes.cpython-38.pyc | Bin 0 -> 1146 bytes __pycache__/sdk.cpython-38.pyc | Bin 0 -> 1775 bytes __pycache__/server.cpython-38.pyc | Bin 0 -> 678 bytes __pycache__/session.cpython-38.pyc | Bin 0 -> 1374 bytes __pycache__/typs.cpython-38.pyc | Bin 0 -> 1095 bytes agent.py | 98 +++++++++++++++++++++++++++ mail.py | 42 ++++++++++++ mailroom.py | 42 ++++++++++++ main.py | 18 +++++ message.py | 37 ++++++++++ schemes.py | 30 ++++++++ server.py | 27 ++++++++ session.py | 81 ++++++++++++++++++++++ typs.py | 28 ++++++++ 19 files changed, 403 insertions(+) create mode 100644 __pycache__/agent.cpython-38.pyc create mode 100644 __pycache__/mail.cpython-38.pyc create mode 100644 __pycache__/mailroom.cpython-38.pyc create mode 100644 __pycache__/message.cpython-38.pyc create mode 100644 __pycache__/messenger.cpython-38.pyc create mode 100644 __pycache__/schemes.cpython-38.pyc create mode 100644 __pycache__/sdk.cpython-38.pyc create mode 100644 __pycache__/server.cpython-38.pyc create mode 100644 __pycache__/session.cpython-38.pyc create mode 100644 __pycache__/typs.cpython-38.pyc create mode 100644 agent.py create mode 100644 mail.py create mode 100644 mailroom.py create mode 100644 main.py create mode 100644 message.py create mode 100644 schemes.py create mode 100644 server.py create mode 100644 session.py create mode 100644 typs.py diff --git a/__pycache__/agent.cpython-38.pyc b/__pycache__/agent.cpython-38.pyc new file mode 100644 index 0000000000000000000000000000000000000000..2fc5689b5cc0d30cf855f8f967b0b4c297a01ed6 GIT binary patch literal 2536 zcmZ`)O>f*p7@iN?>ufgL^qcQmfIurjf{?gWQB~AZsDx4=>VcJz<$5NmlU;i|c3QG> zE={f+IJcEH$NmJ){0m<>@fRwAc;1<8vYQrL^XZv)-uHQ)d3}Frsm<^_{pOF@L+>vd z%svu>n<(X%s1#Ft!1}J`eQw91FKjFO@>`~al7~#m1JU;s-(WrOZ|r01Db=&jiQNa7 z)r@Sx4)0560lPirAF{r$uBes@4|zXO3#yHomReLDw4qv3U9=0jy}hW;J?Cm!oj>&Y z9Xq~&@kMnB<0W+&eq3$*ueWqwjSn|;S!P8ZWqGB|W|Hb?voO&ulMRz9iR_|wJHqpm zPIdNJM`>ZOW;ahOY)b}NmFd#%NsHldH_y_f%JQw~Nme~PF=Gu_TR105^61{ZAEGD8 z-ikAwiWw}ZPvsU$No>@-W-&nCU@!T4k2meK$a8Iv?}^4M^V7J5#~b8xA;DP|9UgTd@2~xX)|;yLjYnIPffjKvk@K(Qs`?u4-XN(_R~7 zWd)u`H{KyLgn9t7#@8H7n{@KZ+Mr02L3sn6>FQU$RR|2XK+WSU)lCT4t#%sUZbMUm zuJN*UKQg);6?v(zN4&rb!T z#3aMAkz2az1%{|IU8 zBJly{R#8fVkclhOW~UAgRXQnKg7OM5=g(_$!~&Bs@rhp(xBF2&OEC+r(WuuygUwoZr{sUNI!zfvyhU!enhcJgUz`5+T8`46zv;5zC{)&dsToJ zdVrsEkT%|SS>%ncc88;qQo71gq;&HMww&>ZzPd2(lTqS=5yvMf-+janCosC&HTj%F z1UYT74y~Z$9CGeVkoM7&wLp(t^eDmC5*Jrc3Wszhq?;#0ZAFP|_bE&?{-8)yX-J}{ z0-mNS0<|A*PX+NAE&9LPX(P8n`WefdiUB#jvk^M2;h%b^*QG7wP#ibyI35;icR>A4 z96#Dk2FG7o0JKPfXGo`0GBl)idtclgBbrZrHy_jRL#j@1nC;R3#Lz#9ER`3i!@VH% z!*|=~!wy<+*|XoL*t8P5by?B0pXdi=k?!cKbO(`@H<@q3H5bMW`VW?LP{{F$Q2qyX C%OJV{ literal 0 HcmV?d00001 diff --git a/__pycache__/mail.cpython-38.pyc b/__pycache__/mail.cpython-38.pyc new file mode 100644 index 0000000000000000000000000000000000000000..6a79f0ad1e61fc959517db1db5396963092b1808 GIT binary patch literal 1488 zcmZux&2H2%5VjL5w{%Jl@LdSkU)YSD2GxZ*%N&X z4%;J-z?rx3l@qVP1!nASe*mvMYkS6%@0)Lu?a?S9Fs@hM_;yIh58Rv|a7LcM)c3#$ zB4|x=|4(!JoQOb#M?{35gFFy)i)8F4#F5Ps?fJoA^8}`T4Mvij3W7KC&}Yb*V4`g=}xQCDfm>|PRB_yc>S@(eUomyN~hAiYlR-CZ>qXZjVc;l7e>0S1IS<;T|1;=y);JqLV*e%139@9n81zSX6bjEDFDpeluO2ctylM`IL5_ZJ# zp_~g@wnAp1>cKtBTwH`jZ_F846`^y8&bRF+u4K7?8BR1F6MB^<^zaI%8=sFT9XpJG zfvB){PhsjiU;s7os102Pz$om3pM|ACOWgd z_G2TJx&khB70rJL=g?rjT(==YV_gW1;FrK}{jAk|3=noXP|OQ$hgCe4>n>~rT~L8+ zNd=AOB%>%8a1T$IHLa<34jeL&#R0qlH*K`xx+!DsS(Y`$RH|#>@v2=E4DtJxyay2; zf53u(u-_~kp5gZ~o7o88D9p&YO*o&nVqW8V$oYqPQJ;AFkhCp1SE!;Ip}B+x@u^#A zu-5qD z_A`tUijGCT^pul-5At9v@=%MJoQpgvMINg_g%=`ElvELXs$!MEuj$~hr_%4F)>Q98 z=6yxGBJ2M_X0bCOcB>Ji0d0wBo&5_Z#ZGV;UF8^Z+BWRs^v`J=w^?!13A-o)qvFISkQDO!SL^)?6SQ7mD z8Fi3!K|7f=YdGgqnFK2rj&)=D9Ao^Yxy2iOY+SNDL}0~|QVJTwbk|szgUj=ak3NEF zKy&VIKBQSN7at-22XPj3!LBg3aUctQ78jLM=x6RML+fy61I9@&&JC*clIjVL*ce6t zXh7yZRvDRla3kip9L&XY3)4P;S#g1bxt0yXjd6`BuG)s7_yTZxC{FFkZYq_cVZ2gL z`3~++W-@$Pw(cVI6uK%Y?CEq=R)%4G5nS|%cTB_quEnw{tgQbVC%3sfIL(zCy&k?m c>FY_Q+GG8OGg_7pwF%NthAGK7P4LA30InIpFaQ7m literal 0 HcmV?d00001 diff --git a/__pycache__/message.cpython-38.pyc b/__pycache__/message.cpython-38.pyc new file mode 100644 index 0000000000000000000000000000000000000000..77079d0b612704e5531ab0b8b26223cbb1873b52 GIT binary patch literal 732 zcmZWnF>ll`6n-yG?sDlJDy1MKmK(Zcpf@5RggSI611ACMa!VD?R+Nx9sqIKyl}_b1 zbU{d$k)OcKzwpY`zre(Eu7?g`>3z1}``+hgKc7z~5zv_CKQ7)efM0I2J3QK)km(}= z4jcD+zwI(pX?`@z9lfEbb|;6b8pQ$P)EjJ z2i)uM4XK|2!e!XU1|3{GHftdteg*D-V=j+&q07^zUpW>UT#s7JqYkgV*30+S;}&=H zPI|p@3!OXnztcES?ZN!&;rm9)Os3PfLTe#Yq3)0hzB)KdBrjZ`OtRu*_4(!0;&D8& zL9H@rZ0Ovos!Gc)8({-c)@EggiC)QNR@p!oLe!R}!dPD>rLcplF(hO6$aQhAX$K}z zjuK?j#aP{#nAzbrJS?fGBo&JX8{RhWJE%ugy;h-^^RPNo)Ul;1`+O~$09r+i`WpssYbf*rQK%(R8I(FUxhTOz=1ol_wHo6`>H-Yd{tJwDa32{ N%k+RC^g=&;_6KBktFiz9 literal 0 HcmV?d00001 diff --git a/__pycache__/messenger.cpython-38.pyc b/__pycache__/messenger.cpython-38.pyc new file mode 100644 index 0000000000000000000000000000000000000000..da20a91e4d7c921e00dff41d77702c50f951c023 GIT binary patch literal 1021 zcmZuv&2G~`5Z>|HiDObIRq26~#RWxzMB)ZQs1lVZherIU7fWa{*=^+DpJdmrXyl&S zw?H6A;=*(A4jlH%sgD3BX6>X25+lv0+1c6oc7BexwmgFJ`~5felM?dF95w}gWFJ-C zMc{-(Oak(PaEsGR!s&$-0Ef4vXJ6qx8IBBYiiL-LROKTGg3%={YpB4Fs?2}M_p7~K)n|;soQAU6) zD>1gAe~zkl5pq(&f|P_=kwMA{Qo^uj>vpKl(ovLY=XD_pp4CPD zh)=!~>0-V=j|ay}NHqxgxJWr#A{`tE^)b(;gOQAqFb^Llt9l04R)2b?UB;p`${925 zS2h9x59T}9<@Z+s|J3AQoC*Gcs=gxRIFbthtY-q1wIDBXNQdMc&~KrSPVCP(rZ2RP zfr(SK7DSzv##%>PrBzwnUcdo)cki6yyLM%jbOGmfX;0j$Q(6mXG6R{tnmwr5Iwh>V zxmVrb&iXua>+ZIO^Sw7knnvl^UzdhGP3Jp@Vaj83VOoq+eBx{9>sB7hvB{PsqTP5xo8 zWoB*u%9!>TOEO->#_lpUE5dkr(`JlkBgUlJhni`F-7K4|zhL&+)Tj)&U{lwT#@Dmz VN^ARn54HQWA%12WqRbPQ{smz0_pSf{ literal 0 HcmV?d00001 diff --git a/__pycache__/schemes.cpython-38.pyc b/__pycache__/schemes.cpython-38.pyc new file mode 100644 index 0000000000000000000000000000000000000000..05076b90382fa6061f82b193b77aef461e6597d5 GIT binary patch literal 1146 zcmZuw&2G~`5Z?9Lj%$)Olpk@7xKxoKxga5-3LzDVB0(-yFS)EtcAJ_ePS|y-)=E$9 zLm;6_<=98y%v;@r3u<+|{f$hv`6L>03S8yJW`4wX2uiU07^4 zQNvlLu+mxgt+ogGWZE5=be!bLa|FZjP@Ar;LU*ReE|N0M(p<`rLQ-tTi{>J?&eAw2 z4P?ZiH3j_%u^*svRB6sX^KYorxj5yV(GH1wG%Io7nHO_*A?Ex$e-AL*th%9mShivD z$-*m`L9b|TCz%?h*|3#nxz6&I;~nQynyyA5lA&$~E|?maou4IRZ73%bky?Se^bC>d z!s>llqDlpH+}J1o8oOf)2Iy}V!XxRlUE=d$U7ccmbsVde#mT*d_wlzu*b*JsSaGIE8Y@HbD+^oT|3$sa{?j0LS^A4I4 z^+TLn;)L;4UeuSlQ^{_nOjn5LEr{KOx=#6o$fER-@`=bjCFTNz`g0FDI1fr%7bGsI z3VY7a^-=D^JTb_%h_B%k;#F(oif^=i(?lF|pA>bptB;d(^l5S=(ToaaDxqnRdO(Vx z+=WOKJ%LzanDIKk$D&wU!d&!3uU(^drhbsp#ZrzZY6hC6t4sN1mW(Qkp#qpyQhyuf pMnEc&p*At}Gt#b%8hMe3?nQoiia)WO|M$#An^oM`^!7J|{Q>y4^WOjf literal 0 HcmV?d00001 diff --git a/__pycache__/sdk.cpython-38.pyc b/__pycache__/sdk.cpython-38.pyc new file mode 100644 index 0000000000000000000000000000000000000000..a4257df8e55121af1d02fb483a2b15143a7fa4f0 GIT binary patch literal 1775 zcmaJ>PjA~c6emT=mQ|-s+GX2igRsK@4VVqA*P>XFp<6ELu%Z`T1XZMMIX0z{R5J(4 zsX6Zx?BD>yfPRjC3to5HSJ-LqQFanH+bHPq@!$LX-XHSS-d;qY{h@x(p8ACRiG$mX z!Qfl?^*10m;dDXLuBRz=V=whw&rT~*%B>rD;j+_)IZ?|i^Da;Cj zw!jWeS!V%G`#iWLX}~|?eI8!Yw8sZLf|)+Q$A{2|Vla>R-fxuO=MOIZ^qw2>i(KpYY$jyJ>zc+$wil$FdEr74AuCzXPQ<+71wksCaJR+?FSq{?M(^4J-8HWkvurHpeK zzj*Oue3q{^rVH;@EKtY)2!4&(XlY9_2=Ytvj-DjcMnxs1aPWy|eJvIf1-T;ci0VQ2 z_v7&^EtDSTeA-A3nvmn~gg!M@JuW&zUoVA*>8Vg-%}+=5%7$51%F<-n1B7*a==l(0 zHQFxRAxUB9d)$rSNKauF9f!nI0qBqvX(H zMb~cK*_BXdbEoRInK?oWh-lC!I`=-vh4tsUk~ZMYvesDIyey!!H<|cPt|{#7-2DOp zu&fCkUSET52`*qy3DYRKLk-h}2OPB#C>^zq*czTP&{*q%hGc!i)_#jA@GgQjfQN~l zD6CEOF$jn-rrrjND(Bj=sW4FAt6D3(FiFqF-vv7j=bNg!Io-LpTmI--SB+nS1^22D z2t9r`{&gGRg?XOrVbCzYS!SawTUNYTpgzp9pPPKKS?Pn(s>m`Gf>AX`ng`oqO5FmCSL2RI2A^B YFlgi7jN9-@$Nwo_a*Z${h0hDwKdrWX{{R30 literal 0 HcmV?d00001 diff --git a/__pycache__/server.cpython-38.pyc b/__pycache__/server.cpython-38.pyc new file mode 100644 index 0000000000000000000000000000000000000000..ff30f0d620cc8f4ccfa5e5949a0e36dc6e4633d5 GIT binary patch literal 678 zcmY*W&59H;5Kg7jKRe8<3SK-2gC_^nzK91|#Dj+wWLBJAkOpbmPPQj(PckGKWo13- zTP(u#(MRy?Tjc7=R~S54>G5Z_puVaKsV`MuPDY~vL3{n?Te$&3euQGzC>BrA?L!1k zIJBhLMkt`KX+b$HA|4ftGai2?Ma&bPVoW&9$RzuTon%_suxCy$^V-No<*NG&d|IeRm>jPn?}O#J-*2mxs^{gO`9@*AUT+vs zMiZL+GNG)|uh>5=?%}aR)g4?Cj#^xjFYFTlanYK&*ny*6bVZM7`r-BqX>moVn(?VN z^JVwRhr9mA$J3N{ld-n3^)v+JE7jLbu&xEr+-mVa4>jl_>Jv{TQm8b3& zW|kokxC!P4#z|_zPnv51!vC7EvDpvMbA?MgTtQpY6J`es%~!2>68h~BA%l!&lmTPe F@HbUro)7>4 literal 0 HcmV?d00001 diff --git a/__pycache__/session.cpython-38.pyc b/__pycache__/session.cpython-38.pyc new file mode 100644 index 0000000000000000000000000000000000000000..5c968ecba9b57b5726bd3bc12ed691978a30856f GIT binary patch literal 1374 zcmZ`(J8u**5VpPE2O$X;pg>ga;3UYENC-h9A|4W?i6W$JtYz&ad+d2Hwl5*4qd-Yb z4+trgRMgb`i)|@T{sI*3b+kg$`x1P5lKT8sTjscC4hFJjyJS(MP?YI@H=`S zd?f~1QNFBRFN-F7wRF`jv5y$T5t#EBZLB~Lr@=t}^c?NxIaWLqoO95JpiY4Cq)(68 zgiPpXMuGVWk%8xm9a2~ks61GIG`>15Su0`mKJ3yE(GBYmN-@3;SDPV`vM`00vYyLr zG01!-3sr|hTNvMEHFvt)0@ojMz7UIT@#{`+(`oB^nc5y11v71WYudf?)v(tB-Nu@= z+19q#?!G)cupwk$26fm!gf6m(jXUi%Ezb^B#pF?8yu#Pj1&vD5YLjG}S!)OeucuK4-kJ<7KS-6&ncE0+yNzzq6r z413nxWCNWp?7)?ky7km?<_=6=RCVY|t%qNjO&~Un+C*fkq03%Gh>ohx0lVw))nNwV zi?kDT*!bFviN8bt&&V@KpO8Lr*E}?ZI#c;BG8LRM1uK4z8I1oIGbc#Dod?qmE#8w| zL~{vDI^XuLM4r+m?A~)l*$#&ORr$QqWdW~EV)6|La0pI#L_3r*sGW_k{;B*-=z2G< zp+k{+>`*BOl^Oy4vYwZ6ca)W_CzeuGxs*+J*(E%|6*Q;EXl8f|7TXCY4e<%(QIaGJ z%~$wga9<2`Q*FxetQH$Dja?3pT)Q;hm`*G{$C^LK2*mr>Ms;#;dfNMl;qU+nO&BOm Fg5Q^SQaJzs literal 0 HcmV?d00001 diff --git a/__pycache__/typs.cpython-38.pyc b/__pycache__/typs.cpython-38.pyc new file mode 100644 index 0000000000000000000000000000000000000000..d33211da9374c4f099654b5cde974a659db3b185 GIT binary patch literal 1095 zcmbVKJ8#rL5Z;#`IdYfq5JD96F32^MC^|tTM8QcYO5AAz(j6@P(>89UBdJQ~KDZ|B>`Xy)^r4Tl2)_I36}exrnZN76SAWCCuF0TfYm zPSSRzDSbs0Q~Zi3{*k2&m{0+*02l`rst8yFOelIvlK2PYAk!SNZw$x;+&%!%B&CX^ zOi{(INXiw5m%@uTSuxKAT<8pAeEAaGJ_Oj%4UuHS0O^JUGDS|wmQEAygX*Hx#)mF5 zGwr-^dgU;y)${iTGnoJ(1lST2!}9al_{?f!$C;YdrGiM8;}dN^xN14hO|i&a_H>~w zY|peAyY z2lVprrp%p!=FaV_-`*d<0eYtm#(!voB~S3uV}NexG*|Q!3?2-%fT0yIdVvpLmrj>1 zp*~z@CR^AX-*$%atyI`ee++lxu(ng2@PzkT$Bn=Bk5$=b+;Mmg+7fonCxO|=mqUaR zfETLB-M{DMKHPZ(-f4mylH7yEpw`U*8C1&PQLm&%T>LLSe%?;$G4^b6Kc+Ds#iPAn Doj=6) literal 0 HcmV?d00001 diff --git a/agent.py b/agent.py new file mode 100644 index 0000000..3d9c283 --- /dev/null +++ b/agent.py @@ -0,0 +1,98 @@ +import json +import asyncio +import websockets + +import typs + + +class Agent: + """Session interface for Crimata Agent + + Receive core functionalities for communicating with Crimata + Agent in an OOP way. + + """ + def __init__(self, connection): + self.connection = connection + + # Submit a request for info, return response. + # Must disclose what service it's for. + async def fetch(self, entities, service): + + # Assert list. + if type(entities) == str: + entities = [entities] + + print(f"Fetching {[e for e in entities]} for {service}") + + # Build intent. + params = { + "service": service, + "entities": entities + }; intent = typs.Intent("fetch", params) + + # Send intent to agent, wait for response. + await self.send_agent_intent(intent) + response = await self.recv_agent_intent() + + # Parse response and return entities. + entities = response.params.get("found") + print(f"Fetch response: {entities}") + return entities + + # Text the client something. Also, let agent know + # when your service is complete via the complete + # attribute. + # Must disclose what service it's for. + async def notify(self, text, end=False): + + print(f"Notifying: {text} for {self.service}") + + params = { + "service": self.service, + "text": text, + "end": end + } + + # Build the intent + intent = typs.Intent("notify", params) + + # Send intent to agent. + await self.send_agent_intent(intent) + + # Get user info. + async def sync(self): + + intent = typs.Intent("sync", params={}) + + # Send intent to agent. + await self.send_agent_intent(intent) + + # Recv and return profile. + intent = await self.recv_agent_intent() + profile = intent.params.get("profile") + return profile + + async def recv_agent_intent(self): + package = await self.connection.recv() + intent = self.__decode(package) + return intent + + async def send_agent_intent(self, intent): + package = self.__encode(intent) + await self.connection.send(package) + + def __encode(self, intent: typs.Intent): + package = json.dumps(intent.__dict__) + return package + + def __decode(self, package) -> typs.Intent: + jpackage = json.loads(package) + name = jpackage.get("name") + params = jpackage.get("params") + intent = typs.Intent(name, params) + return intent + + + + diff --git a/mail.py b/mail.py new file mode 100644 index 0000000..489b2a0 --- /dev/null +++ b/mail.py @@ -0,0 +1,42 @@ +# mail.py + +import typs +import mailroom + + +class Mail: + """Session interface for mailroom. + + Two main IO methods. Will translate intents to mail and vice versa. + """ + def __init__(self): + pass + + # Receive intent. + async def mailbox_recv(self): + mail = await mailroom.get_mail(self.crimata_id) + intent = self.__decode(mail) + return intent + + # Send intent. + def mailbox_send(self, intent): + mail = self.__encode(intent) + mailroom.put_mail(self.crimata_id, mail) + + def __encode(self, intent: typs.Intent) -> typs.Mail: + p = intent.params + owner = self.crimata_id + target = p.get("target") + text = p.get("text") + mail = typs.Mail(owner, target, text) + return mail + + @staticmethod + def __decode(mail: typs.Mail) -> typs.Intent: + name = "notify" + params = { + "text": mail.text #dev only + } + intent = typs.Intent(name, params) + return intent + \ No newline at end of file diff --git a/mailroom.py b/mailroom.py new file mode 100644 index 0000000..7e0d0a2 --- /dev/null +++ b/mailroom.py @@ -0,0 +1,42 @@ +# mailroom.py + +import queue +import asyncio + +import typs + +# Subcribers (Crimata IDs) +subs = [] + +# Mailroom (all mailboxes) +que = queue.Queue() +ref = {} + +# Mailroom functions +# ------------------ + +# Create mailbox and add to mailroom. +def create_mailbox(crimata_id): + mailbox = typs.MailBox(name=crimata_id) + que.put(mailbox) + ref.update({crimata_id: mailbox}) + subs.append(crimata_id) + +# Get mailbox by Crimata ID. +def get_mailbox(crimata_id): + if not crimata_id in subs: + create_mailbox(crimata_id) + mailbox = ref.get(crimata_id) + return mailbox + +# Inbox -> receiving mail. +async def get_mail(crimata_id): + mailbox = get_mailbox(crimata_id) + mail = await mailbox.inbox.get() + return mail + +# Outbox -> sending mail. +def put_mail(crimata_id, mail: typs.Mail): + mailbox = get_mailbox(crimata_id) + mailbox.outbox.put(mail) + \ No newline at end of file diff --git a/main.py b/main.py new file mode 100644 index 0000000..c6ac2fa --- /dev/null +++ b/main.py @@ -0,0 +1,18 @@ +# server.py + +import typs +import queue +import asyncio +import websockets + +import server +import message + + +# Run the ws server and the messenger concurrently. +async def main(): + print(f"Crimata Messenger") + await asyncio.gather(server.lift(), message.run()) + +if __name__ == '__main__': + asyncio.run(main()) diff --git a/message.py b/message.py new file mode 100644 index 0000000..76a1c02 --- /dev/null +++ b/message.py @@ -0,0 +1,37 @@ +# messenge.py + +import asyncio + +import mailroom + + +# Independent loop that updates mailboxes. +async def messenger(): + + print("Running Messenger") + + while True: + + # Yield if no mailboxes. + if mailroom.que.empty(): + await asyncio.sleep(1) + continue + + # Grab a mailbox. + mailbox = mailroom.que.get() + print(f"Handling mailbox: {mailbox.name}.") + + # Handle every message in outbox. + while not mailbox.outbox.empty(): + message = mailbox.outbox.get() + + # Get mailbox of target and put mail there. + target_mailbox = mailroom.get_mailbox(message.target) + await target_mailbox.inbox.put(message) + + # Put the mailbox back + mailroom.que.put(mailbox) + await asyncio.sleep(1) + +async def run(): + await messenger() \ No newline at end of file diff --git a/schemes.py b/schemes.py new file mode 100644 index 0000000..1703a86 --- /dev/null +++ b/schemes.py @@ -0,0 +1,30 @@ +# schemes.py + + +class Schemes: + + def __init__(self): + self.service = None + + async def handle_intent(self, intent): + print(f"Handling intent {intent.name}.") + + # Set the service + self.service = intent.name + + if intent.name == "init": + await self.init(intent) + if intent.name == "message": + self.message(intent) + + async def init(self, intent): + self.crimata_id = intent.params.get("crimata_id") + await self.notify("Messaging is live.") + + def message(self, intent): + text = intent.params.get("text") + target = intent.params.get("target") #crimata_id + print(f"Messaging {target}: '{text}'") + + # Uses mail.Mail interface. + self.mailbox_send(intent) \ No newline at end of file diff --git a/server.py b/server.py new file mode 100644 index 0000000..58ce7b8 --- /dev/null +++ b/server.py @@ -0,0 +1,27 @@ +# server.py + +import asyncio +import websockets + +import session + +HOST = "localhost" +PORT = 8764 + + +# Makes instance of a Messenger session for every client that connects. +async def launch_session(connection, path): + s = session.Session(connection) + + await asyncio.gather( + + s.do_agent_intents(), + + s.deliver_mail() + + ) + +# Run run for every new connection. +async def lift(): + print(f"Listening for connections on {HOST}:{PORT}") + await websockets.serve(launch_session, HOST, PORT) \ No newline at end of file diff --git a/session.py b/session.py new file mode 100644 index 0000000..29a1d37 --- /dev/null +++ b/session.py @@ -0,0 +1,81 @@ +# session.py + +import time +import asyncio + +import mail +import agent +import schemes + + +class Session(schemes.Schemes, agent.Agent, mail.Mail): + """Launch instance for every client connection. + + Will recv messages from client and push them to mailbox. + Also, will pull messages from mailbox and send to client. + + """ + def __init__(self, connection): + agent.Agent.__init__(self, connection) + + self.crimata_id = False + + print("Launched new session") + + # Recv intents run desired endpoint. + async def do_agent_intents(self): + while True: + + # Recv intent from agent. + intent = await self.recv_agent_intent() + print(f"Intent: {intent.name}") + + # Call corresponding endpoint. + await self.handle_intent(intent) + + await asyncio.sleep(0.1) + + # Send mail back to agent. + async def deliver_mail(self): + while True: + + # Sleep until Crimata ID. + if not self.crimata_id: + await asyncio.sleep(1) + continue + + intent = await self.mailbox_recv() + await self.send_agent_intent(intent) + + await asyncio.sleep(0.1) + + + + + + + + + + + + + + + + + + # async def hello(self): + + # intent = await self.recv_agent_intent() + + # response = await self.fetch("crimata_id", service="message") + # self.profile = response.get("profile") + + # target = intent.params.get("target") + # await self.notify(f"Messenger APP: Received Message Intent", + # service="message", complete=True) + + # print("Shutting down in 10s") + + # await asyncio.sleep(10) diff --git a/typs.py b/typs.py new file mode 100644 index 0000000..2ef992f --- /dev/null +++ b/typs.py @@ -0,0 +1,28 @@ +# typs.py + +import queue +import asyncio + + +class Mail: + + def __init__(self, owner, target, text): + self.owner = owner + self.target = target + self.text = text + + +class MailBox: + + def __init__(self, name): + + self.name = name # just for logging purposes. + self.inbox = asyncio.Queue() + self.outbox = queue.Queue() + + +class Intent: + + def __init__(self, name, params: dict): + self.name = name + self.params = params \ No newline at end of file -- 2.43.0