package ohttp-piaf

  1. Overview
  2. Docs
Legend:
Page
Library
Module
Module type
Parameter
Class
Class type
Source

Source file ohttp_piaf.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
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
(* Oblivious HTTP over Piaf. *)

module Service = Ohttp.Service

type 'ctx handler = 'ctx Piaf.Server.ctx -> Piaf.Response.t
type forward = Bhttp.Request.t -> Bhttp.Response.t
type error = [ `Ohttp of Ohttp.Error.t | `Piaf of Piaf.Error.t ]

let pp_error ppf = function
  | `Ohttp e -> Ohttp.Error.pp ppf e
  | `Piaf e -> Piaf.Error.pp_hum ppf e

let fields headers =
  Bhttp.Field.without_connection_specific
    (Bhttp.Field.lowercase (Piaf.Headers.to_list headers))

let respond (response : Service.response) =
  Piaf.Response.of_string
    ~headers:(Piaf.Headers.of_list response.headers)
    ~body:response.body
    (Piaf.Status.of_code response.status)

let plain status = respond { status; headers = []; body = "" }

(* The content of a body, unless it is longer than [max_size]: then the rest is
   read without being kept, since neither a server nor a client is done with a
   message until its body has been read. *)
let read ~max_size headers body =
  let too_long () = Result.map (fun () -> None) (Piaf.Body.drain body) in
  if Service.exceeds ~max_size (Piaf.Headers.to_list headers) then too_long ()
  else
    let buffer = Buffer.create 1024 in
    let stream = Piaf.Body.to_string_stream body in
    let rec loop () =
      match Piaf.Stream.take stream with
      | None ->
          Result.map
            (fun () -> Some (Buffer.contents buffer))
            (Piaf.Body.closed body)
      | Some chunk when Buffer.length buffer + String.length chunk > max_size ->
          too_long ()
      | Some chunk ->
          Buffer.add_string buffer chunk;
          loop ()
    in
    loop ()

(* [f ()], if [admit] lets it in, and [busy] otherwise. *)
let admitted ~admit ~release f =
  if admit () then Fun.protect f ~finally:release else respond Service.busy

(* What failed is answered, unless the fiber is being cancelled. *)
let or_else f ~default =
  try f () with Eio.Cancel.Cancelled _ as e -> raise e | _ -> default ()

module Gateway = struct
  let key_configs service (_ : _ Piaf.Server.ctx) =
    respond (Service.Gateway.key_configs service)

  let requests ?(now = Unix.gettimeofday) service forward
      ({ request; _ } : _ Piaf.Server.ctx) =
    admitted
      ~admit:(fun () -> Service.Gateway.admit service)
      ~release:(fun () -> Service.Gateway.release service)
    @@ fun () ->
    match
      read
        ~max_size:(Service.Gateway.max_request_size service)
        request.headers request.body
    with
    | Error _ -> plain 400
    | Ok None -> respond Service.content_too_large
    | Ok (Some content) -> (
        match
          Service.Gateway.receive service ~now:(now ())
            ~meth:(Piaf.Method.to_string request.meth)
            ~headers:(Piaf.Headers.to_list request.headers)
            content
        with
        | Respond response -> respond response
        | Forward (inner, seal) ->
            respond
              (seal
                 (or_else
                    (fun () -> forward inner)
                    ~default:(fun () -> Bhttp.Response.make ~status:500 ()))))

  let handler ?(path = "/gateway") ?now service forward
      ({ request; _ } as ctx : _ Piaf.Server.ctx) =
    let requested = Uri.path (Uri.of_string request.target) in
    if
      request.meth = `GET
      && requested = Ohttp.Http_binding.well_known_gateway_path
    then key_configs service ctx
    else if requested = path then requests ?now service forward ctx
    else plain 404
end

(* A request that is over when its response has been read, or found too long:
   leaving the switch closes the connection. *)
let fetch ?config ?(headers = []) ?body env ~max_size ~meth uri =
  Eio.Switch.run @@ fun sw ->
  match
    Piaf.Client.Oneshot.request ?config ~headers ?body ~sw env ~meth uri
  with
  | Error e -> Error (`Piaf e)
  | Ok (response : Piaf.Response.t) -> (
      match read ~max_size response.headers response.body with
      | Error e -> Error (`Piaf e)
      | Ok None -> Error (`Ohttp (Ohttp.Error.Content_too_large max_size))
      | Ok (Some content) -> Ok (response, content))

let post ?config env ~max_size uri (request : Service.request) =
  fetch ?config ~headers:request.headers
    ~body:(Piaf.Body.of_string request.body)
    env ~max_size ~meth:`POST uri

let ohttp r = Result.map_error (fun e -> `Ohttp e) r

module Client = struct
  let key_configs ?config
      ?(max_response_size = Service.default_max_response_size) env uri =
    Result.bind
      (fetch ?config
         ~headers:Ohttp.Http_binding.Client.key_config_request_headers env
         ~max_size:max_response_size ~meth:`GET uri)
      (fun ((response : Piaf.Response.t), content) ->
        ohttp
          (Service.Client.key_configs
             ~status:(Piaf.Status.to_code response.status)
             ~headers:(Piaf.Headers.to_list response.headers)
             content))

  let call ?config ?(max_response_size = Service.default_max_response_size) env
      ~rng ?preference ?framing ?padding ?(date = true)
      ?(now = Unix.gettimeofday) ~relay key_config request =
    let rec send (request, exchange) =
      Result.bind (post ?config env ~max_size:max_response_size relay request)
        (fun ((response : Piaf.Response.t), content) ->
          match
            Service.Client.finish exchange
              ~status:(Piaf.Status.to_code response.status)
              ~headers:(Piaf.Headers.to_list response.headers)
              content
          with
          | Error e -> Error (`Ohttp e)
          | Ok (Response response) -> Ok response
          | Ok (Retry (request, exchange)) -> send (request, exchange))
    in
    Result.bind
      (ohttp
         (Service.Client.start ~rng ?preference ?framing ?padding
            ?now:(if date then Some (now ()) else None)
            key_config request))
      send
end

module Relay = struct
  let handler ?config env relay ~gateway ({ request; _ } : _ Piaf.Server.ctx) =
    admitted
      ~admit:(fun () -> Service.Relay.admit relay)
      ~release:(fun () -> Service.Relay.release relay)
    @@ fun () ->
    match
      read
        ~max_size:(Service.Relay.max_request_size relay)
        request.headers request.body
    with
    | Error _ -> plain 400
    | Ok None -> respond Service.content_too_large
    | Ok (Some content) -> (
        match
          Service.Relay.request
            ~meth:(Piaf.Method.to_string request.meth)
            ~headers:(Piaf.Headers.to_list request.headers)
            content
        with
        | Error response -> respond response
        | Ok forwarded ->
            respond
              (match
                 or_else
                   (fun () ->
                     post ?config env
                       ~max_size:(Service.Relay.max_response_size relay)
                       gateway forwarded)
                   ~default:(fun () -> Error (`Piaf (`Msg "no answer")))
               with
              | Ok ((response : Piaf.Response.t), content) ->
                  Service.Relay.response
                    ~status:(Piaf.Status.to_code response.status)
                    ~headers:(Piaf.Headers.to_list response.headers)
                    content
              | Error _ -> Service.Relay.unreachable))
end

module Target = struct
  let request_headers (request : Bhttp.Request.t) =
    let fields = Bhttp.Field.without_connection_specific request.headers in
    if request.authority = "" || Bhttp.Field.get "host" fields <> None then
      fields
    else ("host", request.authority) :: fields

  let forward ?config ?(max_response_size = Service.default_max_response_size)
      env ~targets (request : Bhttp.Request.t) =
    let targets = List.map (fun (a, uri) -> (a, Uri.to_string uri)) targets in
    match Service.Gateway.target ~targets request with
    | Error response -> response
    | Ok uri -> (
        match
          or_else
            (fun () ->
              fetch ?config ~headers:(request_headers request)
                ~body:(Piaf.Body.of_string request.content)
                env ~max_size:max_response_size
                ~meth:(Piaf.Method.of_string request.meth)
                (Uri.of_string uri))
            ~default:(fun () -> Error (`Piaf (`Msg "no answer")))
        with
        | Ok ((response : Piaf.Response.t), content) ->
            Bhttp.Response.make
              ~status:(Piaf.Status.to_code response.status)
              ~headers:(fields response.headers) ~content ()
        (* Too long to seal whole: the target failed to answer. *)
        | Error (`Ohttp (Ohttp.Error.Content_too_large _)) ->
            Bhttp.Response.make ~status:502 ()
        (* No answer from the target is an answer to the client, and is sealed
           like any other (RFC 9458 Section 5). *)
        | Error _ -> Bhttp.Response.make ~status:504 ())
end