ligo/src/lib_stdlib/lwt_watcher.ml

76 lines
3.0 KiB
OCaml
Raw Normal View History

2018-06-29 16:08:08 +04:00
(*****************************************************************************)
(* *)
(* Open Source License *)
(* Copyright (c) 2018 Dynamic Ledger Solutions, Inc. <contact@tezos.com> *)
(* *)
(* Permission is hereby granted, free of charge, to any person obtaining a *)
(* copy of this software and associated documentation files (the "Software"),*)
(* to deal in the Software without restriction, including without limitation *)
(* the rights to use, copy, modify, merge, publish, distribute, sublicense, *)
(* and/or sell copies of the Software, and to permit persons to whom the *)
(* Software is furnished to do so, subject to the following conditions: *)
(* *)
(* The above copyright notice and this permission notice shall be included *)
(* in all copies or substantial portions of the Software. *)
(* *)
(* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR*)
(* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, *)
(* FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL *)
(* THE AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER*)
(* LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING *)
(* FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER *)
(* DEALINGS IN THE SOFTWARE. *)
(* *)
(*****************************************************************************)
2017-02-17 21:15:46 +04:00
type 'a inner_stopper = {
id: int ;
push: ('a option -> unit) ;
mutable active : bool;
input : 'a input;
}
and 'a input =
{ mutable watchers : 'a inner_stopper list;
mutable cpt : int; }
type stopper = unit -> unit
let create_input () =
{ watchers = [];
cpt = 0 }
let shutdown_input input =
let { watchers ; _ } = input in
List.iter (fun w ->
w.active <- false ;
w.push None
) watchers ;
input.cpt <- 0 ;
input.watchers <- []
2017-02-17 21:15:46 +04:00
let create_fake_stream () =
let str, push = Lwt_stream.create () in
str, (fun () -> push None)
let notify input info =
List.iter (fun w -> w.push (Some info)) input.watchers
let shutdown_output output =
if output.active then begin
output.active <- false;
output.push None;
output.input.watchers <-
List.filter (fun w -> w.id <> output.id) output.input.watchers;
end
let create_stream input =
input.cpt <- input.cpt + 1;
let id = input.cpt in
let stream, push = Lwt_stream.create () in
let output = { id; push; input; active = true } in
input.watchers <- output :: input.watchers;
stream, (fun () -> shutdown_output output)
let shutdown f = f ()