123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334(*
* Copyright (C) 2015 David Scott <dave.scott@unikernel.com>
*
* Permission to use, copy, modify, and distribute this software for any
* purpose with or without fee is hereby granted, provided that the above
* copyright notice and this permission notice appear in all copies.
*
* THE SOFTWARE IS PROVIDED "AS IS" AND THE AUTHOR DISCLAIMS ALL WARRANTIES
* WITH REGARD TO THIS SOFTWARE INCLUDING ALL IMPLIED WARRANTIES OF
* MERCHANTABILITY AND FITNESS. IN NO EVENT SHALL THE AUTHOR BE LIABLE FOR
* ANY SPECIAL, DIRECT, INDIRECT, OR CONSEQUENTIAL DAMAGES OR ANY DAMAGES
* WHATSOEVER RESULTING FROM LOSS OF USE, DATA OR PROFITS, WHETHER IN AN
* ACTION OF CONTRACT, NEGLIGENCE OR OTHER TORTIOUS ACTION, ARISING OUT OF
* OR IN CONNECTION WITH THE USE OR PERFORMANCE OF THIS SOFTWARE.
*
*)openResultopenProtocol_9p_infixopenProtocol_9p_infomoduleError=Protocol_9p_erroropenErrormoduleTypes=Protocol_9p_typesmoduleRequest=Protocol_9p_requestmoduleResponse=Protocol_9p_responsetypeexn_converter=Protocol_9p_info.t->exn->Protocol_9p_response.payloadmoduleMake(Log:Protocol_9p_s.LOG)(FLOW:Mirage_flow_lwt.S)(Filesystem:Protocol_9p_filesystem.S)=structmoduleReader=Protocol_9p_buffered9PReader.Make(Log)(FLOW)openLogtypet={write_lock:Lwt_mutex.t;reader:Reader.t;writer:FLOW.flow;info:Protocol_9p_info.t;root_qid:Types.Qid.t;(* Press the "cancel button" by setting the ref to true (if present in the
map) *)mutablecancel_buttons:unitLwt.uTypes.Tag.Map.t;mutableplease_shutdown:bool;shutdown_complete_t:unitLwt.t;}letget_infot=t.infoletdefault_exn_converterinfoexn=Response.Err{Response.Err.ename=Printexc.to_stringexn;errno=None;}(* For converting flow errors *)let(>>|=)mf=letopenLwtinm>>=function|Okx->fx|Error`Closed->return(error_msg"Writing to closed FLOW")|Errore->return(error_msg"Unexpected error on underlying FLOW: %a"FLOW.pp_write_errore)letdisconnectt=t.please_shutdown<-true;t.shutdown_complete_tletafter_disconnectt=t.shutdown_complete_tletwrite_one_packet?write_lockwriterresponse=debug(funf->f"S %a"Response.ppresponse);letsizeof=Response.sizeofresponseinletbuffer=Cstruct.createsizeofinLwt.return(Response.writeresponsebuffer)>>*=fun_->(matchwrite_lockwith|Somem->Lwt_mutex.with_lockm(fun()->FLOW.writewriterbuffer)|None->FLOW.writewriterbuffer)>>|=fun()->Lwt.return(Ok())letread_one_packetreader=letopenLwtinReader.readreader>>=function|Error(`Msg_)ase->Lwt.returne|Okbuffer->Lwt.returnbeginmatchRequest.readbufferwith|Error(`Msgename)->Error(`Parse(ename,buffer))|Ok(request,_)->debug(funf->f"C %a"Request.pprequest);Okrequestendleterror_responsetagename={Response.tag;payload=Response.(Err{Err.ename;errno=None;});}letrecdispatcher_tinfoexn_convertershutdown_complete_wakenerreceive_cbt=ift.please_shutdownthenbeginLwt.wakeup_latershutdown_complete_wakener();Lwt.return(Ok())endelsebeginletopenLwtinread_one_packett.reader>>=function|Error(`Msgmessage)->debug(funf->f"S error reading: %s"message);debug(funf->f"Disconnecting client");t.please_shutdown<-true;dispatcher_tinfoexn_convertershutdown_complete_wakenerreceive_cbt|Error(`Parse(ename,buffer))->beginmatchRequest.read_headerbufferwith|Error(`Msg_)->debug(funf->f"C sent bad header: %s"ename);dispatcher_tinfoexn_convertershutdown_complete_wakenerreceive_cbt|Ok(_,tag,_)->debug(funf->f"C error: %s"ename);letresponse=error_responsetagenameinwrite_one_packet~write_lock:t.write_lockt.writerresponse>>*=fun()->dispatcher_tinfoexn_convertershutdown_complete_wakenerreceive_cbtend|Ok{Request.tag;payload=Request.Flush{Request.Flush.oldtag}}->Lwt_mutex.with_lockt.write_lock(fun()->ifTypes.Tag.Map.memoldtagt.cancel_buttonsthenbeginletcancel_u=Types.Tag.Map.findoldtagt.cancel_buttonsinLwt.wakeup_latercancel_u();debug(funf->f"S will suppress response for tag %s"(Sexplib.Sexp.to_string(Types.Tag.sexp_of_toldtag)));t.cancel_buttons<-Types.Tag.Map.removeoldtagt.cancel_buttons;end;write_one_packett.writer{Response.tag;payload=Response.Flush()})>>=beginfunction|Ok()->dispatcher_tinfoexn_convertershutdown_complete_wakenerreceive_cbt|Error(`Msgm)->disconnectt>>=fun()->Lwt.return(Error(`Msgm))end|Okrequest->letcancel_t,cancel_u=Lwt.task()int.cancel_buttons<-Types.Tag.Map.addrequest.Request.tagcancel_ut.cancel_buttons;Lwt.async(fun()->Lwt.catch(fun()->receive_cb~cancel:cancel_trequest.Request.payload)(funexn->letbacktrace=Printexc.get_raw_backtrace()inLog.err(funf->f"Uncaught exception handling %a: %a"Request.pprequestFmt.exn_backtrace(exn,backtrace));Lwt.return(Result.Ok(exn_converterinfoexn)))>>=beginfunction|Error(`Msgmessage)->Lwt.return(error_responserequest.Request.tagmessage)|Okresponse_payload->Lwt.return{Response.tag=request.Request.tag;payload=response_payload;}end>>=funresponse->Lwt_mutex.with_lockt.write_lock(fun()->ifLwt.statecancel_t=Lwt.Sleepthenbegin(* It's safe to unbind the tag because the flush hasn't been
transmitted yet. *)t.cancel_buttons<-Types.Tag.Map.removerequest.Request.tagt.cancel_buttons;write_one_packett.writerresponseendelsebeginLwt.return(Ok())end)>>=beginfunction|Error(`Msgmessage)->debug(funf->f"S error writing: %s"message);debug(funf->f"Disconnecting client");disconnectt|Ok()->Lwt.return()end);dispatcher_tinfoexn_convertershutdown_complete_wakenerreceive_cbtendmoduleLowLevel=structletreturn_error~write_lockwriterrequestename=write_one_packet~write_lockwriter{Response.tag=request.Request.tag;payload=Response.ErrResponse.Err.({ename;errno=None})}>>*=fun()->Lwt.return(Error(`Msgename))letexpect_version~write_lockreaderwriter=Reader.readreader>>*=funbuffer->Lwt.return(Request.readbuffer)>>*=function|({Request.payload=Request.Versionv;tag}asreq,_)->debug(funf->f"C %a"Request.ppreq);Lwt.return(Ok(tag,v))|request,_->return_error~write_lockwriterrequest"Expected Version message"letexpect_attach~write_lockreaderwriter=Reader.readreader>>*=funbuffer->Lwt.return(Request.readbuffer)>>*=function|({Request.payload=(Request.Attacha)aspayload;tag}asreq,_)->debug(funf->f"C %a"Request.ppreq);Lwt.return(Ok(tag,a,payload))|request,_->return_error~write_lockwriterrequest"Expected Attach message"endletconnectfsflow?(msize=16384l)?(exn_converter=default_exn_converter)()=letwrite_lock=Lwt_mutex.create()inletreader=Reader.createflowinletwriter=flowinLowLevel.expect_version~write_lockreaderwriter>>*=fun(tag,v)->letmsize=minmsizev.Request.Version.msizeinletmsize=minmsizeProtocol_9p_buffered9PReader.max_message_sizeinifv.Request.Version.version=Types.Version.unknownthenbeginerr(funf->f"Client sent a 9P version string we couldn't understand");Lwt.return(Error(`Msg"Received unknown 9P version string"))endelsebeginletversion=v.Request.Version.versioninwrite_one_packet~write_lockflow{Response.tag;payload=Response.VersionResponse.Version.({msize;version});}>>*=fun()->info(funf->f"Using protocol %s msize %ld"(Sexplib.Sexp.to_string(Types.Version.sexp_of_tversion))msize);LowLevel.expect_attach~write_lockreaderwriter>>*=fun(tag,a,payload)->letcancel_buttons=Types.Tag.Map.emptyinletplease_shutdown=falseinletshutdown_complete_t,shutdown_complete_wakener=Lwt.task()inletroot=a.Request.Attach.fidinletaname=a.Request.Attach.anameinletinfo={root;version;aname;msize}inletconnection=Filesystem.connectfsinfoinletreceive_cb~cancel=letis_unix=(info.version=Types.Version.unix)inletadjust_errnoerr=ifnotis_unixthen{errwithResponse.Err.errno=None}elsematcherr.Response.Err.errnowith|Some_->err|None->{errwithResponse.Err.errno=Some0l}inletwrapfnxresult=letopenLwt.Infixinfnconnection~cancelx>|=function|Okresponse->Ok(resultresponse)|Errorerr->Ok(Response.Err(adjust_errnoerr))inRequest.(function|Attachx->wrapFilesystem.attachx(funx->Response.Attachx)|Walkx->wrapFilesystem.walkx(funx->Response.Walkx)|Openx->wrapFilesystem.open_x(funx->Response.Openx)|Readx->wrapFilesystem.readx(funx->Response.Readx)|Clunkx->wrapFilesystem.clunkx(funx->Response.Clunkx)|Statx->wrapFilesystem.statx(funx->Response.Statx)|Createx->wrapFilesystem.createx(funx->Response.Createx)|Writex->wrapFilesystem.writex(funx->Response.Writex)|Removex->wrapFilesystem.removex(funx->Response.Removex)|Wstatx->wrapFilesystem.wstatx(funx->Response.Wstatx)|Version_|Auth_|Flush_->leterr={Response.Err.ename="Function not implemented";errno=None}inLwt.return(Result.Ok(Response.Err(adjust_errnoerr))))inletopenLwtinletcancel_t,cancel_u=Lwt.task()inreceive_cb~cancel:cancel_tpayload>>=beginfunction|Error(`Msgmessage)->letresponse=error_responsetagmessageinLwt_mutex.with_lockwrite_lock(fun()->write_one_packetwriterresponse)>>*=fun()->Lwt.return(Error(`Msgmessage))|Ok(Response.Attachaaspayload)->letresponse={Response.tag;payload;}inLwt_mutex.with_lockwrite_lock(fun()->write_one_packetwriterresponse)>>*=fun()->return(Oka.Response.Attach.qid)|Ok_->letmessage="expected Attach reply"inletresponse=error_responsetagmessageinLwt_mutex.with_lockwrite_lock(fun()->write_one_packetwriterresponse)>>*=fun()->Lwt.return(Error(`Msgmessage))end>>*=funroot_qid->lett={reader;writer;info;root_qid;cancel_buttons;please_shutdown;shutdown_complete_t;write_lock;}inLwt.async(fun()->Lwt.catch(fun()->letopenLwt.Infixindispatcher_tinfoexn_convertershutdown_complete_wakenerreceive_cbt>>=function|Result.Error(`Msgm)->err(funf->f"dispatcher caught %s: no more requests will be handled"m);Lwt.return()|Result.Ok()->Lwt.return())(fune->err(funf->f"dispatcher caught %s: no more requests will be handled"(Printexc.to_stringe));Lwt.wakeup_latershutdown_complete_wakener();Lwt.return()));Lwt.return(Okt)endend