123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589openCoreopenAsynctypecommon_error=[`Connection_closed|`Unexpected][@@derivingshow,eq]typeresponse=(Resp.t,common_error)resulttypecommand=stringlisttyperequest={command:command;waiter:responseIvar.t}letconstruct_requestcommands=commands|>List.map~f:(funcmd->Resp.Bulkcmd)|>(funxs->Resp.Arrayxs)|>Resp.encodetypet={(* Need a queue of waiter Ivars. Need some way of closing the connection *)waiters:responseIvar.tQueue.t;reader:requestPipe.Reader.t;writer:requestPipe.Writer.t}letinitreaderwriter=letwaiters=Queue.create()inletrecrecv_loopreader=match%bindMonitor.try_with_or_error@@fun()->Parser.read_respreaderwith|Error_|Ok(Error_)->return()|Ok(Okrasresult)->(matchQueue.dequeuewaiterswith|NonewhenReader.is_closedreader->return()|None->failwithf"No waiters are waiting for this message: %s"(Resp.showr)()|Somewaiter->Ivar.fillwaiterresult;recv_loopreader)in(* Requests are posted to a pipe, and requests are processed in sequence *)letrequest_reader,request_writer=Pipe.create()inlethandle_request{command;waiter}=Queue.enqueuewaiterswaiter;letrequest=construct_requestcommandinreturn@@Writer.writewriterrequestin(* Start redis receiver. Processing ends if the connection is closed. *)don't_wait_for(let%bind()=recv_loopreaderinreturn@@Pipe.closerequest_writer);(* Start processing requests. Once the pipe is closed, we signal
closed to all outstanding waiters after closing the underlying
socket *)don't_wait_for(let%bind()=Pipe.iterrequest_reader~f:handle_requestinlet%bind()=Writer.closewriterinlet%bind()=Reader.closereaderin(* Signal this to all waiters. As the pipe has been closed, we
know that no new waiters will arrive *)Queue.iterwaiters~f:(funwaiter->Ivar.fillwaiter@@Error`Connection_closed);return@@Queue.clearwaiters);{waiters;reader=request_reader;writer=request_writer}letconnect?(port=6379)~host=letwhere=Tcp.Where_to_connect.of_host_and_port@@Host_and_port.create~host~portinlet%bind_socket,reader,writer=Tcp.connectwhereinreturn@@initreaderwriterletclose{writer;_}=return@@Pipe.closewriterletrequesttcommand=matchPipe.is_closedt.writerwith|true->return@@Error`Connection_closed|false->(letwaiter=Ivar.create()inlet%bind()=Pipe.writet.writer{command;waiter}in(* Type coercion: [common_error] -> [> common_error] *)match%mapIvar.readwaiterwith|Ok_asres->res|Error`Connection_closed->Error`Connection_closed|Error`Unexpected->Error`Unexpected)letechotmessage=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["ECHO";message]with|Resp.Bulkv->returnv|_->Deferred.return@@Error`Unexpectedtypeexist=|Always|Not_if_exists|Only_if_existsletsett~key?expire?(exist=Always)value=letopenDeferred.Result.Let_syntaxinletexpiry=matchexpirewith|None->[]|Somespan->["PX";span|>Time.Span.to_ms|>int_of_float|>string_of_int]inletexistence=matchexistwith|Always->[]|Not_if_exists->["NX"]|Only_if_exists->["XX"]inletcommand=["SET";key;value]@expiry@existenceinmatch%bindrequesttcommandwith|Resp.Null->returnfalse|Resp.String"OK"->returntrue|_->Deferred.return@@Error`Unexpectedletgettkey=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["GET";key]with|Resp.Bulkv->return@@Somev|Resp.Null->return@@None|_->Deferred.return@@Error`Unexpectedletgetranget~start~end'key=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["GETRANGE";key;string_of_intstart;string_of_intend']with|Resp.Bulkv->returnv|_->Deferred.return@@Error`Unexpectedletgetsett~keyvalue=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["GETSET";key;value]with|Resp.Bulkv->return(Somev)|Resp.Null->returnNone|_->Deferred.return@@Error`Unexpectedletstrlentkey=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["STRLEN";key]with|Resp.Integerv->returnv|_->Deferred.return@@Error`Unexpectedletmgettkeys=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt("MGET"::keys)with|Resp.Arrayxs->xs|>List.fold_right~init:(Ok[])~f:(funitemacc->matchaccwith|Error_->acc|Okacc->(matchitemwith|Resp.Null->Ok(None::acc)|Resp.Bulks->Ok(Somes::acc)|_->Error`Unexpected))|>Deferred.return|_->Deferred.return@@Error`Unexpectedletmsettalist=letopenDeferred.Result.Let_syntaxinletpayload=alist|>List.map~f:(fun(k,v)->[k;v])|>List.concatinmatch%bindrequestt("MSET"::payload)with|Resp.String"OK"->return()|_->Deferred.return@@Error`Unexpectedletmsetnxtalist=letopenDeferred.Result.Let_syntaxinletpayload=alist|>List.map~f:(fun(k,v)->[k;v])|>List.concatinmatch%bindrequestt("MSETNX"::payload)with|Resp.Integer1->returntrue|Resp.Integer0->returnfalse|_->Deferred.return@@Error`Unexpectedletlpusht~keyvalue=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["LPUSH";key;value]with|Resp.Integern->returnn|_->Deferred.return@@Error`Unexpectedletlranget~key~start~stop=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["LRANGE";key;string_of_intstart;string_of_intstop]with|Resp.Arrayxs->List.mapxs~f:(function|Resp.Bulkv->Okv|_->Error`Unexpected)|>Result.all|>Deferred.return|_->Deferred.return@@Error`Unexpectedletappendt~keyvalue=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["APPEND";key;value]with|Resp.Integern->returnn|_->Deferred.return@@Error`Unexpectedletauthtpassword=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["AUTH";password]with|Resp.String"OK"->return()|Resp.Errore->Deferred.return@@Error(`Redis_errore)|_->Deferred.return@@Error`Unexpectedletbgrewriteaoft=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["BGREWRITEAOF"]with(* the documentation says it returns OK, but that's not true *)|Resp.Stringv->returnv|_->Deferred.return@@Error`Unexpectedletbgsavet=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["BGSAVE"]with|Resp.Stringv->returnv|_->Deferred.return@@Error`Unexpectedletbitcountt?rangekey=letopenDeferred.Result.Let_syntaxinletrange=matchrangewith|None->[]|Some(start,end_)->[string_of_intstart;string_of_intend_]inmatch%bindrequestt(["BITCOUNT";key]@range)with|Resp.Integern->returnn|_->Deferred.return@@Error`Unexpectedtypeoverflow=|Wrap|Sat|Failletstring_of_overflow=function|Wrap->"WRAP"|Sat->"SAT"|Fail->"FAIL"(* Declaration of type of the integer *)typeintsize=|Signedofint|Unsignedofintletstring_of_intsize=function|Signedv->Printf.sprintf"i%d"v|Unsignedv->Printf.sprintf"u%d"vtypeoffset=|Absoluteofint|Relativeofintletstring_of_offset=function|Absolutev->string_of_intv|Relativev->Printf.sprintf"#%d"vtypefieldop=|Getofintsize*offset|Setofintsize*offset*int|Incrbyofintsize*offset*intletbitfieldt?overflowkeyops=letopenDeferred.Result.Let_syntaxinletops=ops|>List.map~f:(function|Get(size,offset)->["GET";string_of_intsizesize;string_of_offsetoffset]|Set(size,offset,value)->["SET";string_of_intsizesize;string_of_offsetoffset;string_of_intvalue]|Incrby(size,offset,increment)->["INCRBY";string_of_intsizesize;string_of_offsetoffset;string_of_intincrement])|>List.concatinletoverflow=matchoverflowwith|None->[]|Somebehaviour->["OVERFLOW";string_of_overflowbehaviour]inmatch%bindrequestt(["BITFIELD";key]@overflow@ops)with|Resp.Arrayxs->letopenResult.Let_syntaxinxs|>List.fold~init:(Ok[])~f:(funaccv->matchacc,vwith|Error_,_->acc|Okacc,Resp.Integeri->Ok(Somei::acc)|Okacc,Resp.Null->Ok(None::acc)|Ok_,_->Error`Unexpected)>>|List.rev|>Deferred.return|_->Deferred.return@@Error`Unexpectedtypebitop=|AND|OR|XOR|NOTletstring_of_bitop=function|AND->"AND"|OR->"OR"|XOR->"XOR"|NOT->"NOT"letbitopt~destkey?(keys=[])~keyop=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt(["BITOP";string_of_bitopop;destkey;key]@keys)with|Resp.Integern->returnn|_->Deferred.return@@Error`Unexpectedtypebit=|Zero|One[@@derivingshow,eq]letstring_of_bit=function|Zero->"0"|One->"1"letbitpost?start?end'keybit=letopenDeferred.Result.Let_syntaxinlet%bindrange=matchstart,end'with|Somes,Somee->return[string_of_ints;string_of_inte]|Somes,None->return[string_of_ints]|None,None->return[]|None,Some_->raise(Invalid_argument"Can't specify end without start")inmatch%bindrequestt(["BITPOS";key;string_of_bitbit]@range)with|Resp.Integer-1->returnNone|Resp.Integern->return@@Somen|_->Deferred.return@@Error`Unexpectedletgetbittkeyoffset=letopenDeferred.Result.Let_syntaxinletoffset=string_of_intoffsetinmatch%bindrequestt["GETBIT";key;offset]with|Resp.Integer0->returnZero|Resp.Integer1->returnOne|_->Deferred.return@@Error`Unexpectedletsetbittkeyoffsetvalue=letopenDeferred.Result.Let_syntaxinletoffset=string_of_intoffsetinletvalue=string_of_bitvalueinmatch%bindrequestt["SETBIT";key;offset;value]with|Resp.Integer0->returnZero|Resp.Integer1->returnOne|_->Deferred.return@@Error`Unexpectedletdecrtkey=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["DECR";key]with|Resp.Integern->returnn|_->Deferred.return@@Error`Unexpectedletdecrbytkeydecrement=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["DECRBY";key;string_of_intdecrement]with|Resp.Integern->returnn|_->Deferred.return@@Error`Unexpectedletincrtkey=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["INCR";key]with|Resp.Integern->returnn|_->Deferred.return@@Error`Unexpectedletincrbytkeyincrement=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["INCRBY";key;string_of_intincrement]with|Resp.Integern->returnn|_->Deferred.return@@Error`Unexpectedletincrbyfloattkeyincrement=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["INCRBYFLOAT";key;string_of_floatincrement]with|Resp.Bulkv->return@@float_of_stringv|_->Deferred.return@@Error`Unexpectedletselecttindex=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["SELECT";string_of_intindex]with|Resp.String"OK"->return()|_->Deferred.return@@Error`Unexpectedletdelt?(keys=[])key=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt(["DEL";key]@keys)with|Resp.Integern->returnn|_->Deferred.return@@Error`Unexpectedletexistst?(keys=[])key=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt(["EXISTS";key]@keys)with|Resp.Integern->returnn|_->Deferred.return@@Error`Unexpectedletexpiretkeyspan=letopenDeferred.Result.Let_syntaxinletmilliseconds=Time.Span.to_msspanin(* rounded to nearest millisecond *)letexpire=Printf.sprintf"%.0f"millisecondsinmatch%bindrequestt["PEXPIRE";key;expire]with|Resp.Integern->returnn|_->Deferred.return@@Error`Unexpectedletexpireattkeydt=letopenDeferred.Result.Let_syntaxinletsince_epoch=dt|>Time.to_span_since_epoch|>Time.Span.to_msinletexpire=Printf.sprintf"%.0f"since_epochinmatch%bindrequestt["PEXPIREAT";key;expire]with|Resp.Integern->returnn|_->Deferred.return@@Error`Unexpectedletkeystpattern=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["KEYS";pattern]with|Resp.Arrayxs->List.mapxs~f:(function|Resp.Bulkkey->Okkey|_->Error`Unexpected)|>Result.all|>Deferred.return|_->Deferred.return@@Error`Unexpectedletscan?pattern?countt=letpattern=matchpatternwith|Somepattern->["MATCH";pattern]|None->[]inletcount=matchcountwith|Somecount->["COUNT";string_of_intcount]|None->[]inPipe.create_reader~close_on_exception:false@@funwriter->Deferred.repeat_until_finished"0"@@funcursor->match%bindrequestt(["SCAN";cursor]@pattern@count)with|Ok(Resp.Array[Resp.Bulkcursor;Resp.Arrayfrom])->(letfrom=from|>List.map~f:(function|Resp.Bulks->s|_->failwith"unexpected")|>Queue.of_listinlet%bind()=Pipe.transfer_inwriter~frominmatchcursorwith|"0"->return@@`Finished()|cursor->return@@`Repeatcursor)|_->failwith"unexpected"letmovetkeydb=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["MOVE";key;string_of_intdb]with|Resp.Integer0->returnfalse|Resp.Integer1->returntrue|_->Deferred.return@@Error`Unexpectedletpersisttkey=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["PERSIST";key]with|Resp.Integer0->returnfalse|Resp.Integer1->returntrue|_->Deferred.return@@Error`Unexpectedletrandomkeyt=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["RANDOMKEY"]with|Resp.Bulks->returns|_->Deferred.return@@Error`Unexpectedletrenametkeynewkey=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["RENAME";key;newkey]with|Resp.String"OK"->return()|_->Deferred.return@@Error`Unexpectedletrenamenxt~keynewkey=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["RENAMENX";key;newkey]with|Resp.Integer0->returnfalse|Resp.Integer1->returntrue|_->Deferred.return@@Error`Unexpectedtypeorder=|Asc|Descletsortt?by?limit?get?order?alpha?storekey=letopenDeferred.Result.Let_syntaxinletby=matchbywith|None->[]|Someby->["BY";by]inletlimit=matchlimitwith|None->[]|Some(offset,count)->["LIMIT";string_of_intoffset;string_of_intcount]inletget=matchgetwith|None->[]|Somepatterns->patterns|>List.map~f:(funpattern->["GET";pattern])|>List.concatinletorder=matchorderwith|None->[]|SomeAsc->["ASC"]|SomeDesc->["DESC"]inletalpha=matchalphawith|None->[]|Somefalse->[]|Sometrue->["ALPHA"]inletstore=matchstorewith|None->[]|Somedestination->["STORE";destination]inletq=[["SORT";key];by;limit;get;order;alpha;store]|>List.concatinmatch%bindrequesttqwith|Resp.Integercount->return@@`Countcount|Resp.Arraysorted->sorted|>List.map~f:(function|Resp.Bulkv->Okv|_->Error`Unexpected)|>Result.all|>Result.map~f:(funx->`Sortedx)|>Deferred.return|_->Deferred.return@@Error`Unexpectedletttltkey=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["PTTL";key]with|Resp.Integer-2->Deferred.return@@Error(`No_such_keykey)|Resp.Integer-1->Deferred.return@@Error(`Not_expiringkey)|Resp.Integerms->ms|>float_of_int|>Time.Span.of_ms|>return|_->Deferred.return@@Error`Unexpectedlettype'tkey=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["TYPE";key]with|Resp.String"none"->returnNone|Resp.Strings->return@@Somes|_->Deferred.return@@Error`Unexpectedletdumptkey=letopenDeferred.Result.Let_syntaxinmatch%bindrequestt["DUMP";key]with|Resp.Bulkbulk->return@@Somebulk|Resp.Null->returnNone|_->Deferred.return@@Error`Unexpectedletrestoret~key?ttl?replacevalue=letopenDeferred.Result.Let_syntaxinletttl=matchttlwith|None->"0"|Somespan->span|>Time.Span.to_ms|>Printf.sprintf".0%f"inletreplace=matchreplacewith|Sometrue->["REPLACE"]|Somefalse|None->[]inmatch%bindrequestt(["RESTORE";key;ttl;value]@replace)with|Resp.String"OK"->return()|_->Deferred.return@@Error`Unexpectedletwith_connection?(port=6379)~hostf=letwhere=Tcp.Where_to_connect.of_host_and_port@@Host_and_port.create~host~portinTcp.with_connectionwhere@@fun_socketreaderwriter->lett=initreaderwriterinft