package tiny_libs

  1. Overview
  2. Docs
From-scratch libraries for teaching: graphics, audio, compression, crypto, networking and more

Install

dune-project
 Dependency

Authors

Maintainers

Sources

0.3.6.tar.gz
md5=7c636383d146d30ac6f2fa234a6253c8
sha512=c79f3823c5f8f57e5038eb640d487c61168b84aa07c61999d6622ef9fd0c890e2b03b4c6a7cdbbe9352a49e25dda00ac7bb14693cee8e3d7beeed251351a2af0

doc/src/tiny_libs.networking_unix/Server.ml.html

Source file Server.ml

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
(* 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 *)

type client = {
  id : int;
  fd : Unix.file_descr;
  mutable inbox : string; (* bytes read, not yet understood *)
  mutable outbox : string; (* bytes to write, when the socket takes them *)
  mutable upgraded : bool; (* past the WebSocket handshake *)
  mutable closing : bool; (* closed once its outbox is written *)
}

type event = Joined of int | Message of int * string | Left of int
type t = {
  listener : Unix.file_descr;
  lines : bool; (* plain TCP, a message a line, rather than WebSocket *)
  mutable clients : client list;
  mutable next_id : int;
}

let listen (caps : < Cap.network ; .. >) ?(lines = false) ~(bind : string) ~(port : int) () : t * int =
  let (_ : Cap.Network.t) = caps#network bind in
  let fd = Unix.socket Unix.PF_INET Unix.SOCK_STREAM 0 in
  Unix.setsockopt fd Unix.SO_REUSEADDR true;
  Unix.bind fd (Unix.ADDR_INET (Unix.inet_addr_of_string bind, port));
  Unix.listen fd 16;
  Unix.set_nonblock fd;
  let port = match Unix.getsockname fd with Unix.ADDR_INET (_, p) -> p | _ -> port in
  ({ listener = fd; lines; clients = []; next_id = 0 }, port)

let clients (t : t) : int list = List.filter_map (fun c -> if c.upgraded && not c.closing then Some c.id else None) t.clients
let frame (c : client) (f : Websocket.frame) : unit = c.outbox <- c.outbox ^ Websocket.encode f

let send (t : t) (id : int) (payload : string) : unit =
  List.iter
    (fun c ->
      if c.id = id && not c.closing then
        if t.lines then c.outbox <- c.outbox ^ payload ^ "\r\n" else frame c { fin = true; opcode = Binary; payload })
    t.clients

let close (t : t) (id : int) : unit =
  List.iter
    (fun c ->
      if c.id = id && not c.closing then begin
        if not t.lines then frame c { fin = true; opcode = Close; payload = "" };
        c.closing <- true
      end)
    t.clients

(*****************************************************************************)
(* One step of each connection *)
(*****************************************************************************)

let accept_all (t : t) : unit =
  let rec go () =
    match Unix.accept t.listener with
    | fd, _ ->
        Unix.set_nonblock fd;
        (* 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.setsockopt fd Unix.TCP_NODELAY true;
        t.clients <- t.clients @ [ { id = t.next_id; fd; inbox = ""; outbox = ""; upgraded = false; closing = false } ];
        t.next_id <- t.next_id + 1;
        go ()
    | exception Unix.Unix_error ((Unix.EAGAIN | Unix.EWOULDBLOCK), _, _) -> ()
  in
  go ()

(* what arrived; a connection ended or broken is closing *)
let read (c : client) : unit =
  let buf = Bytes.create 65536 in
  let rec go () =
    match Unix.read c.fd buf 0 (Bytes.length buf) with
    | 0 -> c.closing <- true
    | n ->
        c.inbox <- c.inbox ^ Bytes.sub_string buf 0 n;
        go ()
    | exception Unix.Unix_error ((Unix.EAGAIN | Unix.EWOULDBLOCK), _, _) -> ()
    | exception Unix.Unix_error _ -> c.closing <- true
  in
  if not c.closing then go ()

(* the handshake: the connection becomes a client *)
let upgrade (c : client) : event list =
  match Websocket.handshake c.inbox with
  | None -> []
  | Some (headers, stop) -> (
      c.inbox <- String.sub c.inbox stop (String.length c.inbox - stop);
      match List.assoc_opt "sec-websocket-key" headers with
      | None ->
          c.closing <- true;
          []
      | Some key ->
          c.outbox <- c.outbox ^ Websocket.response ~key;
          c.upgraded <- true;
          [ Joined c.id ])

(* in lines mode: every whole line, its CR LF (or LF: a telnet on Unix)
   taken off *)
let rec lines_of (c : client) (acc : event list) : event list =
  match String.index_opt c.inbox '\n' with
  | None -> List.rev acc
  | Some i ->
      let line = String.sub c.inbox 0 i in
      let line = if line <> "" && line.[String.length line - 1] = '\r' then String.sub line 0 (String.length line - 1) else line in
      c.inbox <- String.sub c.inbox (i + 1) (String.length c.inbox - i - 1);
      lines_of c (Message (c.id, line) :: acc)

(* every whole frame: a message, a ping answered, a close *)
let rec frames (c : client) (acc : event list) : event list =
  match Websocket.decode c.inbox with
  | Incomplete -> List.rev acc
  | Bad _ ->
      c.closing <- true;
      List.rev acc
  | Frame (f, n) -> (
      c.inbox <- String.sub c.inbox n (String.length c.inbox - n);
      match f.opcode with
      | Binary -> frames c (Message (c.id, f.payload) :: acc)
      | Ping ->
          frame c { fin = true; opcode = Pong; payload = f.payload };
          frames c acc
      | Close ->
          c.closing <- true;
          List.rev acc
      | _ -> frames c acc)

(* as much of the outbox as the socket takes now *)
let write (c : client) : unit =
  if c.outbox <> "" then
    match Unix.write_substring c.fd c.outbox 0 (String.length c.outbox) with
    | n -> c.outbox <- String.sub c.outbox n (String.length c.outbox - n)
    | exception Unix.Unix_error ((Unix.EAGAIN | Unix.EWOULDBLOCK), _, _) -> ()
    | exception Unix.Unix_error _ ->
        c.outbox <- "";
        c.closing <- true

let flush (t : t) : unit = List.iter write t.clients

let step (t : t) : event list =
  accept_all t;
  let events =
    List.concat_map
      (fun c ->
        read c;
        if t.lines then (
          (* no handshake: a client as soon as it is accepted *)
          let joined = if c.upgraded then [] else (c.upgraded <- true; [ Joined c.id ]) in
          joined @ lines_of c [])
        else
          let joined = if c.upgraded then [] else upgrade c in
          joined @ if c.upgraded then frames c [] else [])
      t.clients
  in
  List.iter write t.clients;
  (* the connections ended, once they have said what they had to *)
  let gone, kept = List.partition (fun c -> c.closing && c.outbox = "") t.clients in
  List.iter (fun c -> try Unix.close c.fd with Unix.Unix_error _ -> ()) gone;
  t.clients <- kept;
  events @ List.filter_map (fun c -> if c.upgraded then Some (Left c.id) else None) gone

let wait (t : t) (timeout : float) : unit =
  let reads = t.listener :: List.map (fun c -> c.fd) t.clients in
  let writes = List.filter_map (fun c -> if c.outbox <> "" then Some c.fd else None) t.clients in
  try ignore (Unix.select reads writes [] timeout) with Unix.Unix_error (Unix.EINTR, _, _) -> ()