2017-11-11 06:34:12 +04:00
|
|
|
(**************************************************************************)
|
|
|
|
(* *)
|
2017-11-14 03:36:14 +04:00
|
|
|
(* Copyright (c) 2014 - 2017. *)
|
2017-11-11 06:34:12 +04:00
|
|
|
(* Dynamic Ledger Solutions, Inc. <contact@tezos.com> *)
|
|
|
|
(* *)
|
|
|
|
(* All rights reserved. No warranty, explicit or implicit, provided. *)
|
|
|
|
(* *)
|
|
|
|
(**************************************************************************)
|
|
|
|
|
|
|
|
open Lwt.Infix
|
|
|
|
|
|
|
|
exception Closed
|
|
|
|
|
|
|
|
type 'a t =
|
|
|
|
{ mutable data : 'a option ;
|
|
|
|
mutable closed : bool ;
|
2017-12-29 18:29:29 +04:00
|
|
|
mutable put_waiter : (unit Lwt.t * unit Lwt.u) option ;
|
2017-11-11 06:34:12 +04:00
|
|
|
}
|
|
|
|
|
|
|
|
let create () =
|
|
|
|
{ data = None ;
|
|
|
|
closed = false ;
|
|
|
|
put_waiter = None ;
|
|
|
|
}
|
|
|
|
|
|
|
|
let notify_put dropbox =
|
|
|
|
match dropbox.put_waiter with
|
|
|
|
| None -> ()
|
2017-12-29 18:29:29 +04:00
|
|
|
| Some (_waiter, wakener) ->
|
2017-11-11 06:34:12 +04:00
|
|
|
dropbox.put_waiter <- None ;
|
2017-12-29 18:29:29 +04:00
|
|
|
Lwt.wakeup_later wakener ()
|
2017-11-11 06:34:12 +04:00
|
|
|
|
2017-11-13 17:29:28 +04:00
|
|
|
let put dropbox elt =
|
2017-11-11 06:34:12 +04:00
|
|
|
if dropbox.closed then
|
|
|
|
raise Closed
|
|
|
|
else begin
|
|
|
|
dropbox.data <- Some elt ;
|
|
|
|
notify_put dropbox
|
|
|
|
end
|
|
|
|
|
|
|
|
let peek dropbox = dropbox.data
|
|
|
|
|
|
|
|
let close dropbox =
|
|
|
|
if not dropbox.closed then begin
|
|
|
|
dropbox.closed <- true ;
|
|
|
|
notify_put dropbox ;
|
|
|
|
end
|
|
|
|
|
|
|
|
let wait_put ~timeout dropbox =
|
|
|
|
match dropbox.put_waiter with
|
2017-12-29 18:29:29 +04:00
|
|
|
| Some (waiter, _wakener) ->
|
2017-11-11 06:34:12 +04:00
|
|
|
Lwt.choose [
|
|
|
|
timeout ;
|
2017-12-29 18:29:29 +04:00
|
|
|
Lwt.protected waiter
|
2017-11-11 06:34:12 +04:00
|
|
|
]
|
|
|
|
| None ->
|
|
|
|
let waiter, wakener = Lwt.wait () in
|
2017-12-29 18:29:29 +04:00
|
|
|
dropbox.put_waiter <- Some (waiter, wakener) ;
|
2017-11-11 06:34:12 +04:00
|
|
|
Lwt.choose [
|
|
|
|
timeout ;
|
|
|
|
Lwt.protected waiter ;
|
|
|
|
]
|
|
|
|
|
|
|
|
let rec take dropbox =
|
|
|
|
match dropbox.data with
|
|
|
|
| Some elt ->
|
|
|
|
dropbox.data <- None ;
|
|
|
|
Lwt.return elt
|
|
|
|
| None ->
|
|
|
|
if dropbox.closed then
|
|
|
|
Lwt.fail Closed
|
|
|
|
else
|
|
|
|
wait_put ~timeout:Lwt_utils.never_ending dropbox >>= fun () ->
|
|
|
|
take dropbox
|
|
|
|
|
|
|
|
let rec take_with_timeout timeout dropbox =
|
|
|
|
match dropbox.data with
|
|
|
|
| Some elt ->
|
|
|
|
Lwt.cancel timeout ;
|
|
|
|
dropbox.data <- None ;
|
|
|
|
Lwt.return (Some elt)
|
|
|
|
| None ->
|
|
|
|
if Lwt.is_sleeping timeout then
|
|
|
|
if dropbox.closed then
|
|
|
|
Lwt.fail Closed
|
|
|
|
else
|
|
|
|
wait_put ~timeout dropbox >>= fun () ->
|
|
|
|
take_with_timeout timeout dropbox
|
|
|
|
else
|
|
|
|
Lwt.return_none
|
|
|
|
|
|
|
|
let take_with_timeout timeout dropbox =
|
|
|
|
take_with_timeout (Lwt_unix.sleep timeout) dropbox
|