package moonpool
sectionYPositions = computeSectionYPositions($el), 10)"
x-init="setTimeout(() => sectionYPositions = computeSectionYPositions($el), 10)"
>
Pools of threads supported by a pool of domains
Install
dune-project
Dependency
Authors
Maintainers
Sources
moonpool-0.12.tbz
sha256=7641327370ece987c09accdd15eb31566931ecdd58d41189080b0b951b8addfc
sha512=65d45c44cda768ba6934ad4d6b9a397651a87b34a143c85fe1108228db861c369d6efc4832adfcac85368c6a861f33816af4d7c3e3d1238eaddfc564c0ea42db
doc/src/moonpool.dpool/moonpool_dpool.ml.html
Source file moonpool_dpool.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(* shared with {!Moonpool.Bb_queue}/{!Moonpool.Lock}; we can't depend on [moonpool] itself (layering would be circular), so these live in [moonpool.private] instead. *) module Bb_queue = Bb_queue_ module Lock = Lock_ type domain = Domain_.t type event = | Run of (unit -> unit) (** Run this function *) | Decr (** Decrease count *) (* State for a domain worker. It should not do too much except for starting new threads for pools. *) type worker_state = { q: event Bb_queue.t; th_count: int Atomic.t; (** Number of threads on this *) } type pool = (worker_state option * Domain_.t option) Lock.t array let max_number_of_domains_ = max 1 (Domain_.recommended_number ()) (* lazily created pool (for fork safety; safe to fork until the first domain pool is created). *) let domains_ : pool option Lock.t = Lock.create None (** main work loop for a domain worker. A domain worker does two things: - run functions it's asked to (mainly, to start new threads inside it) - decrease the refcount when one of these threads stops. The thread will notify the domain that it's exiting, so the domain can know how many threads are still using it. If all threads exit, the domain polls a bit (in case new threads are created really shortly after, which happens with a [Pool.with_] or [Pool.create() … Pool.shutdown()] in a tight loop), and if nothing happens it tries to stop to free resources. *) let work_ (domains : pool) idx (st : worker_state) : unit = Signals_.ignore_signals_ (); let main_loop () = let continue = ref true in while !continue do match Bb_queue.pop st.q with | Run f -> (try f () with exn -> (* that invariant is just too important *) Printf.eprintf "moonpool.dpool: fatal error, callback raised with %s\n%!" (Printexc.to_string exn); exit 1) | Decr -> if Atomic.fetch_and_add st.th_count (-1) = 1 then ( continue := false; (* wait a bit, we might be needed again in a short amount of time *) try for _n_attempt = 1 to 50 do Thread.delay 0.001; if Atomic.get st.th_count > 0 then ( (* needed again! *) continue := true; raise Exit ) done with Exit -> () ) done in while main_loop (); (* exit: try to remove ourselves from [domains]. If that fails, keep living. *) let is_alive = Lock.update_map domains.(idx) (function | None, _ -> assert false | Some _st', dom -> assert (st == _st'); if Atomic.get st.th_count > 0 then (* still alive! *) (Some st, dom), true else (None, dom), false) in is_alive do () done; () let init_domains_ () : pool = Lock.update_map domains_ @@ function | Some domains -> Some domains, domains | None -> assert (Domain_.is_main_domain ()); let domains = Array.init max_number_of_domains_ (fun _ -> Lock.create (None, None)) in let w = { th_count = Atomic.make 1; q = Bb_queue.create () } in domains.(0) <- Lock.create (Some w, None); ignore (Thread.create (fun () -> work_ domains 0 w) () : Thread.t); Some domains, domains let[@inline] max_number_of_domains () : int = max_number_of_domains_ let run_on (i : int) (f : unit -> unit) : unit = assert (i < max_number_of_domains_); let domains = init_domains_ () in let w : worker_state = Lock.update_map domains.(i) (function | (Some w, _) as st -> Atomic.incr w.th_count; st, w | None, dying_dom -> (* join previous dying domain, to free its resources, if any *) Option.iter Domain_.join dying_dom; let w = { th_count = Atomic.make 1; q = Bb_queue.create () } in let worker : domain = Domain_.spawn (fun () -> work_ domains i w) in (Some w, Some worker), w) in Bb_queue.push w.q (Run f) let decr_on (i : int) : unit = assert (i < max_number_of_domains_); let domains = match Lock.get domains_ with | Some domains -> domains | None -> assert false in match Lock.get domains.(i) with | Some st, _ -> Bb_queue.push st.q Decr | None, _ -> () let run_on_and_wait (i : int) (f : unit -> 'a) : 'a = let q = Bb_queue.create () in run_on i (fun () -> let x = f () in Bb_queue.push q x); Bb_queue.pop q
sectionYPositions = computeSectionYPositions($el), 10)"
x-init="setTimeout(() => sectionYPositions = computeSectionYPositions($el), 10)"
>