(*****************************************************************************) (* *) (* Open Source License *) (* Copyright (c) 2018 Dynamic Ledger Solutions, Inc. *) (* *) (* 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. *) (* *) (*****************************************************************************) 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 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 ()