123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177(* Claude Code
*
* Copyright (C) 2026 Yoann Padioleau
*
* This library is free software; you can redistribute it and/or
* modify it under the terms of the GNU Library General Public License
* (LGPL) as published by the Free Software Foundation; either version
* 2 of the License, or (at your option) any later version.
*)(* See Server.mli *)typeclient={id:int;fd:Unix.file_descr;mutableinbox:string;(* bytes read, not yet understood *)mutableoutbox:string;(* bytes to write, when the socket takes them *)mutableupgraded:bool;(* past the WebSocket handshake *)mutableclosing:bool;(* closed once its outbox is written *)}typeevent=Joinedofint|Messageofint*string|Leftofinttypet={listener:Unix.file_descr;lines:bool;(* plain TCP, a message a line, rather than WebSocket *)mutableclients:clientlist;mutablenext_id:int;}letlisten(caps:<Cap.network;..>)?(lines=false)~(bind:string)~(port:int)():t*int=let(_:Cap.Network.t)=caps#networkbindinletfd=Unix.socketUnix.PF_INETUnix.SOCK_STREAM0inUnix.setsockoptfdUnix.SO_REUSEADDRtrue;Unix.bindfd(Unix.ADDR_INET(Unix.inet_addr_of_stringbind,port));Unix.listenfd16;Unix.set_nonblockfd;letport=matchUnix.getsocknamefdwithUnix.ADDR_INET(_,p)->p|_->portin({listener=fd;lines;clients=[];next_id=0},port)letclients(t:t):intlist=List.filter_map(func->ifc.upgraded&¬c.closingthenSomec.idelseNone)t.clientsletframe(c:client)(f:Websocket.frame):unit=c.outbox<-c.outbox^Websocket.encodefletsend(t:t)(id:int)(payload:string):unit=List.iter(func->ifc.id=id&¬c.closingthenift.linesthenc.outbox<-c.outbox^payload^"\r\n"elseframec{fin=true;opcode=Binary;payload})t.clientsletclose(t:t)(id:int):unit=List.iter(func->ifc.id=id&¬c.closingthenbeginifnott.linesthenframec{fin=true;opcode=Close;payload=""};c.closing<-trueend)t.clients(*****************************************************************************)(* One step of each connection *)(*****************************************************************************)letaccept_all(t:t):unit=letrecgo()=matchUnix.acceptt.listenerwith|fd,_->Unix.set_nonblockfd;(* claude: no Nagle: a small write after another (the welcome
after the handshake's answer, a tick's packet after the one
before) otherwise waits for the first one's ACK, which the
other side may delay: tens of milliseconds a packet, a
game's lag *)Unix.setsockoptfdUnix.TCP_NODELAYtrue;t.clients<-t.clients@[{id=t.next_id;fd;inbox="";outbox="";upgraded=false;closing=false}];t.next_id<-t.next_id+1;go()|exceptionUnix.Unix_error((Unix.EAGAIN|Unix.EWOULDBLOCK),_,_)->()ingo()(* what arrived; a connection ended or broken is closing *)letread(c:client):unit=letbuf=Bytes.create65536inletrecgo()=matchUnix.readc.fdbuf0(Bytes.lengthbuf)with|0->c.closing<-true|n->c.inbox<-c.inbox^Bytes.sub_stringbuf0n;go()|exceptionUnix.Unix_error((Unix.EAGAIN|Unix.EWOULDBLOCK),_,_)->()|exceptionUnix.Unix_error_->c.closing<-trueinifnotc.closingthengo()(* the handshake: the connection becomes a client *)letupgrade(c:client):eventlist=matchWebsocket.handshakec.inboxwith|None->[]|Some(headers,stop)->(c.inbox<-String.subc.inboxstop(String.lengthc.inbox-stop);matchList.assoc_opt"sec-websocket-key"headerswith|None->c.closing<-true;[]|Somekey->c.outbox<-c.outbox^Websocket.response~key;c.upgraded<-true;[Joinedc.id])(* in lines mode: every whole line, its CR LF (or LF: a telnet on Unix)
taken off *)letreclines_of(c:client)(acc:eventlist):eventlist=matchString.index_optc.inbox'\n'with|None->List.revacc|Somei->letline=String.subc.inbox0iinletline=ifline<>""&&line.[String.lengthline-1]='\r'thenString.subline0(String.lengthline-1)elselineinc.inbox<-String.subc.inbox(i+1)(String.lengthc.inbox-i-1);lines_ofc(Message(c.id,line)::acc)(* every whole frame: a message, a ping answered, a close *)letrecframes(c:client)(acc:eventlist):eventlist=matchWebsocket.decodec.inboxwith|Incomplete->List.revacc|Bad_->c.closing<-true;List.revacc|Frame(f,n)->(c.inbox<-String.subc.inboxn(String.lengthc.inbox-n);matchf.opcodewith|Binary->framesc(Message(c.id,f.payload)::acc)|Ping->framec{fin=true;opcode=Pong;payload=f.payload};framescacc|Close->c.closing<-true;List.revacc|_->framescacc)(* as much of the outbox as the socket takes now *)letwrite(c:client):unit=ifc.outbox<>""thenmatchUnix.write_substringc.fdc.outbox0(String.lengthc.outbox)with|n->c.outbox<-String.subc.outboxn(String.lengthc.outbox-n)|exceptionUnix.Unix_error((Unix.EAGAIN|Unix.EWOULDBLOCK),_,_)->()|exceptionUnix.Unix_error_->c.outbox<-"";c.closing<-trueletflush(t:t):unit=List.iterwritet.clientsletstep(t:t):eventlist=accept_allt;letevents=List.concat_map(func->readc;ift.linesthen((* no handshake: a client as soon as it is accepted *)letjoined=ifc.upgradedthen[]else(c.upgraded<-true;[Joinedc.id])injoined@lines_ofc[])elseletjoined=ifc.upgradedthen[]elseupgradecinjoined@ifc.upgradedthenframesc[]else[])t.clientsinList.iterwritet.clients;(* the connections ended, once they have said what they had to *)letgone,kept=List.partition(func->c.closing&&c.outbox="")t.clientsinList.iter(func->tryUnix.closec.fdwithUnix.Unix_error_->())gone;t.clients<-kept;events@List.filter_map(func->ifc.upgradedthenSome(Leftc.id)elseNone)goneletwait(t:t)(timeout:float):unit=letreads=t.listener::List.map(func->c.fd)t.clientsinletwrites=List.filter_map(func->ifc.outbox<>""thenSomec.fdelseNone)t.clientsintryignore(Unix.selectreadswrites[]timeout)withUnix.Unix_error(Unix.EINTR,_,_)->()