Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 14 additions & 7 deletions cache.ml
Original file line number Diff line number Diff line change
Expand Up @@ -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)
57 changes: 32 additions & 25 deletions control.ml
Original file line number Diff line number Diff line change
Expand Up @@ -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 *)
}
Expand All @@ -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
32 changes: 28 additions & 4 deletions log.ml
Original file line number Diff line number Diff line change
Expand Up @@ -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}
*)

Expand All @@ -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
Expand All @@ -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 "*" ->
Expand All @@ -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)
Expand Down Expand Up @@ -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
Expand All @@ -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
Expand All @@ -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.

Expand Down
Loading