diff --git a/cache.ml b/cache.ml index a043e1e..830527e 100644 --- a/cache.ml +++ b/cache.ml @@ -383,17 +383,24 @@ let reuse create reset = { cache = Stack.create (); create; reset; } let use t = if Stack.is_empty t.cache then t.create () else Stack.pop t.cache let recycle t x = t.reset x; Stack.push x t.cache -module ReuseLocked(L : Lock)(T : sig type t val create : unit -> t val reset : t -> unit end) : sig +module Reuse(T : sig type t val create : unit -> t val reset : t -> unit end) : sig type t = T.t val get : unit -> t val release : t -> unit end = struct type t = T.t -type cache = { cache : t Stack.t; lock : L.t } -let cache = { cache = Stack.create (); lock = L.create () } -let get' () = if Stack.is_empty cache.cache then T.create () else Stack.pop cache.cache -let get () = L.locked cache.lock get' -let release x = L.locked cache.lock (fun () -> T.reset x; Stack.push x cache.cache) +type cache = { cache : t list Atomic.t; } [@@unboxed] +let cache = { cache = Atomic.make []; } +let rec get () = match Atomic.get cache.cache with +| [] -> T.create() +| x :: tl as old -> + if Atomic.compare_and_set cache.cache old tl then x else get () +let release x = + T.reset x; + let continue = ref true in + while !continue do + let old = Atomic.get cache.cache in + continue := not (Atomic.compare_and_set cache.cache old (x::old)) + done end -module Reuse = ReuseLocked(NoLock) diff --git a/control.ml b/control.ml index d829d37..14fa5eb 100644 --- a/control.ml +++ b/control.ml @@ -26,15 +26,19 @@ let with_output_txt name k = with_open_out_txt name (fun ch -> bracket (IO.outpu let with_opendir dir = bracket (Unix.opendir dir) Unix.closedir -(* token bucket +(* token bucket, domain-safe. https://en.wikipedia.org/wiki/Token_bucket *) module Rate_limit = struct + type bucket = { + tokens: float; + last_update: float; + } + type t = | Unlimited | RL of { - mutable tokens: float; - mutable count_silenced: int; - mutable last_update: float; + bucket: bucket Atomic.t; (** current state. avoid mutex for reentrancy *) + count_silenced: int Atomic.t; capacity: float; rate: float; (** new tokens/sec *) } @@ -48,33 +52,36 @@ module Rate_limit = struct if burst_factor < 1 then invalid_arg "Rate_limit.create: burst factor must be >= 1"; let capacity = max 1. (float burst_factor *. allowed_per_sec) in RL { - tokens=capacity; last_update=Time.now(); count_silenced=0; capacity; + bucket = Atomic.make { tokens=capacity; last_update=Time.now() }; + count_silenced = Atomic.make 0; + capacity; rate=allowed_per_sec; } let take_rate_limited_count = function | Unlimited -> 0 - | RL rl -> - let n = rl.count_silenced in - rl.count_silenced <- 0; - n + | RL rl -> Atomic.exchange rl.count_silenced 0 - let attempt = function + let rec attempt_rec now = function | Unlimited -> true - | RL rl -> - let now = Time.now() in - - if now > rl.last_update then ( - rl.tokens <- min rl.capacity - (rl.tokens +. rl.rate *. (now -. rl.last_update)); - rl.last_update <- now; - ); - - if rl.tokens >= 1. then ( - rl.tokens <- rl.tokens -. 1.; - true - ) else ( - rl.count_silenced <- 1 + rl.count_silenced; + | RL rl as rate_limiter -> + let old = Atomic.get rl.bucket in + let b = + let time_since_last_refill = now -. old.last_update in + if time_since_last_refill > 1e-3 then + (* lazily refill, avoid small float precision errors *) + let tokens = min rl.capacity (old.tokens +. rl.rate *. time_since_last_refill) in + { tokens; last_update = now } + else + old + in + if b.tokens >= 1. then + let ok = Atomic.compare_and_set rl.bucket old { b with tokens = b.tokens -. 1. } in + if ok then true (* done *) else attempt_rec now rate_limiter + else begin + Atomic.incr rl.count_silenced; false - ) + end + + let attempt rl = attempt_rec (Time.now()) rl end diff --git a/log.ml b/log.ml index 465def5..0c13685 100644 --- a/log.ml +++ b/log.ml @@ -31,6 +31,21 @@ or Output only messages of warning level or higher for all facilities {[Log.set_filter `Warn]} +{2 Domains} + +Logging (the [logger] methods, [Logger.t.put]) is safe from any domain. +Output lines are not interleaved, and the shared {!main_rate_limiter} is +domain-safe. + +Everything else is main-domain only: creating facilities ({!facility}, +{!from}, [new logger]) and configuration ({!set_filter}, {!set_loglevels}, +{!read_env_config}, {!set_utc}, [State.set_plaintext], [State.set_logfmt]). +These raise [Failure] when called from another domain. Do them at startup, +before spawning domains; later changes become visible to other domains eventually. + +[State.hook] and [State.logger_target] must likewise only be set from the main domain, +but the hook is {e called} from whichever domain logs, so it must be domain-safe itself. + {2 API} *) @@ -40,12 +55,17 @@ open Prelude (** Global logger state *) module State = struct + let check_main_domain name = + if not (Domain.is_main_domain ()) then + Exn.fail "Log.%s: must be called from the main domain" name + let all = Hashtbl.create 10 let default_level = ref (`Info : Logger.level) let utc_timezone = ref false let facility name = + check_main_domain "facility"; try Hashtbl.find all name with @@ -55,6 +75,7 @@ module State = struct x let set_filter ?name level = + check_main_domain "set_filter"; match name with | None -> default_level := level; Hashtbl.iter (fun _ x -> Logger.set_filter x level) all | Some name when Stre.ends_with name "*" -> @@ -63,6 +84,7 @@ module State = struct | Some name -> Logger.set_filter (facility name) level let set_loglevels s = + check_main_domain "set_loglevels"; Stre.nsplitc s ',' |> List.iter begin fun spec -> match Stre.nsplitc spec '=' with | name :: l :: [] -> set_filter ~name (Logger.level l) @@ -114,8 +136,8 @@ module State = struct end let get_cur_format () = Atomic.get cur_format let is_structured_format () = match get_cur_format () with `Plain, _ -> false | `Logfmt, _ -> true - let set_plaintext () = set_cur_format (`Plain, format_simple_full) - let set_logfmt () = set_cur_format (`Logfmt, format_logfmt) + let set_plaintext () = check_main_domain "set_plaintext"; set_cur_format (`Plain, format_simple_full) + let set_logfmt () = check_main_domain "set_logfmt"; set_cur_format (`Logfmt, format_logfmt) let format level facil ts pairs msg = (snd (Atomic.get cur_format)) level facil ts pairs msg @@ -137,6 +159,7 @@ module State = struct let logger = Logger.put_simple logger_target let self = "lib" + let self_facil = facility self (* we open the new fd, then dup it to stderr and close afterwards @@ -153,14 +176,15 @@ module State = struct with e -> let now = (Unix.gettimeofday ()) in - logger.put `Warn (facility self) now [] (sprintf "reopen_log_ch(%s) failed : %s" file (Printexc.to_string e)) + (* this might run from any domain, so use [self_facil] since [facility] is not domain-safe *) + logger.put `Warn self_facil now [] (sprintf "reopen_log_ch(%s) failed : %s" file (Printexc.to_string e)) end let facility = State.facility let set_filter = State.set_filter let set_loglevels = State.set_loglevels -let set_utc () = State.utc_timezone := true +let set_utc () = State.check_main_domain "set_utc"; State.utc_timezone := true (** Update facilities configuration from the environment.