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
6 changes: 3 additions & 3 deletions .github/workflows/gh-pages.yml
Original file line number Diff line number Diff line change
Expand Up @@ -9,12 +9,12 @@ jobs:
deploy:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@main
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1

# needed for depext to work
- run: sudo apt-get update && sudo apt-get install mccs

- uses: ocaml/setup-ocaml@v3
- uses: ocaml/setup-ocaml@15d660006c1d3110d77c34b7faa3bddefe8b82f0 # v3.7.0
with:
ocaml-compiler: '5.1.x'
dune-cache: true
Expand All @@ -27,7 +27,7 @@ jobs:
run: opam exec -- odig odoc --cache-dir=_doc/ opentelemetry opentelemetry-lwt opentelemetry-client-ocurl opentelemetry-cohttp-lwt

- name: Deploy
uses: peaceiris/actions-gh-pages@v3
uses: peaceiris/actions-gh-pages@373f7f263a76c20808c831209c920827a82a2847 # v3.9.3
with:
github_token: ${{ secrets.GITHUB_TOKEN }}
publish_dir: ./_doc/html
Expand Down
4 changes: 2 additions & 2 deletions .github/workflows/main.yml
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@ jobs:

steps:
- name: Checkout code
uses: actions/checkout@v4
uses: actions/checkout@11d5960a326750d5838078e36cf38b85af677262 # v4.4.0
with:
submodules: recursive

Expand All @@ -34,7 +34,7 @@ jobs:
if: ${{ matrix.os == 'ubuntu-latest' }}

- name: Use OCaml ${{ matrix.ocaml-compiler }}
uses: ocaml/setup-ocaml@v3
uses: ocaml/setup-ocaml@15d660006c1d3110d77c34b7faa3bddefe8b82f0 # v3.7.0
with:
ocaml-compiler: ${{ matrix.ocaml-compiler }}
opam-depext-flags: --with-test
Expand Down
4 changes: 2 additions & 2 deletions .github/workflows/nix.yml
Original file line number Diff line number Diff line change
Expand Up @@ -7,8 +7,8 @@ jobs:
runs-on: ubuntu-latest
steps:
- name: Checkout tree
uses: actions/checkout@v4
uses: actions/checkout@11d5960a326750d5838078e36cf38b85af677262 # v4.4.0
with:
submodules: true
- uses: cachix/install-nix-action@v30
- uses: cachix/install-nix-action@08dcb3a5e62fa31e2da3d490afc4176ef55ecd72 # v30
- run: nix develop -L .# -c dune build @runtest @check
95 changes: 77 additions & 18 deletions src/client-cohttp-eio/opentelemetry_client_cohttp_eio.ml
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,15 @@ let n_errors = Atomic.make 0

let n_dropped = Atomic.make 0

let tick_interval_s = 0.5

(* Longer wait after a failed export, to avoid hammering a down endpoint. *)
let export_backoff_s = 3.0

(* Per-POST cap: [Httpc.send] has no timeout, so a black-hole endpoint (accepts
the connection, never replies) would otherwise wedge a tick or process exit. *)
let send_timeout_s = 5.0

let report_err_ = function
| `Sysbreak -> Printf.eprintf "opentelemetry: ctrl-c captured, stopping\n%!"
| `Failure msg ->
Expand Down Expand Up @@ -214,6 +223,9 @@ module type EMITTER = sig

val tick : unit -> unit

(** Background flush loop; forked into the caller's switch by {!create_backend}. *)
val run_ticker : unit -> unit

val cleanup : on_done:(unit -> unit) -> unit -> unit
end

Expand All @@ -222,26 +234,55 @@ end
exceptions inside should be caught, see
https://opentelemetry.io/docs/reference/specification/error-handling/ *)
let mk_emitter ~stop ~net (config : Config.t) : (module EMITTER) =
(* One-shot shutdown signal. A promise, not a condition: level-triggered
([await] after resolution returns at once, so no lost-wakeup window) and it
wakes awaiters on any domain. *)
let shutdown, resolve_shutdown = Eio.Promise.create () in
let signal_stop () =
(* guarded: Eio.Promise.resolve raises if the promise is already resolved *)
if not (Atomic.exchange stop true) then
Eio.Promise.resolve resolve_shutdown ()
in
let wait_or_shutdown d =
Fiber.first
(fun () -> Eio.Promise.await shutdown)
(fun () -> Eio_unix.sleep d)
in
(* [send_http] sets this on failure; the ticker reads+resets it to pick the next
wait (backoff vs cadence). *)
let send_failed = Atomic.make false in
(* local helpers *)
let open struct
let client =
(* Prime RNG state for TLS *)
Mirage_crypto_rng_unix.use_default ();
Httpc.create net

(* No backoff here on purpose — the ticker owns it; a failed export returns at
once, bounded by [send_timeout_s]. *)
let send_http ~url data : unit =
let r = Httpc.send client ~url ~decode:(`Ret ()) data in
match r with
| Ok () -> ()
| Error `Sysbreak ->
let outcome =
Fiber.first
(fun () -> `Sent (Httpc.send client ~url ~decode:(`Ret ()) data))
(fun () ->
Eio_unix.sleep send_timeout_s;
`Timed_out)
in
match outcome with
| `Sent (Ok ()) -> ()
| `Sent (Error `Sysbreak) ->
Printf.eprintf "ctrl-c captured, stopping\n%!";
Atomic.set stop true
| Error err ->
(* TODO: log error _via_ otel? *)
signal_stop ()
| `Sent (Error err) ->
Atomic.incr n_errors;
Atomic.set send_failed true;
(* once shutting down, a failed export is expected teardown noise *)
if not (Atomic.get stop) then report_err_ err
| `Timed_out ->
Atomic.incr n_errors;
report_err_ err;
(* avoid crazy error loop *)
Eio_unix.sleep 3.
Atomic.set send_failed true;
if not (Atomic.get stop) then
report_err_ (`Failure (spf "export POST timed out after %.0fs" send_timeout_s))

let timeout =
if config.batch_timeout_ms > 0 then
Expand Down Expand Up @@ -333,13 +374,32 @@ let mk_emitter ~stop ~net (config : Config.t) : (module EMITTER) =
sample_gc_metrics_if_needed ();
emit_all ~force:false

(* Because exports never block (see [send_http]), no [tick] blocks — so a
synchronous caller that ticks before cleanup ([Collector.remove_backend])
cannot stall here. *)
let run_ticker () =
let rec loop prev_failed =
wait_or_shutdown
(if prev_failed then export_backoff_s else tick_interval_s);
if not (Atomic.get stop) then (
Atomic.set send_failed false;
tick ();
loop (Atomic.get send_failed))
in
loop false

let cleanup ~on_done () =
if Config.Env.get_debug () then
Printf.eprintf "opentelemetry: exiting…\n%!";
Atomic.set stop true;
(* Signal before flushing. The ticker is non-daemon on purpose: the switch
awaits its in-flight export instead of cancelling it mid-send, which would
drop a batch it had already popped. *)
signal_stop ();
run_tick_callbacks ();
sample_gc_metrics_if_needed ();
emit_all ~force:true;
(* [Cancel.protect]: a switch teardown racing in must not abort this final
flush mid-send. *)
Eio.Cancel.protect (fun () -> emit_all ~force:true);
on_done ()
end in
(module M : EMITTER)
Expand Down Expand Up @@ -449,12 +509,11 @@ let create_backend ~sw ?(stop = Atomic.make false) ?(config = Config.make ())

NOTE: This cannot be located inside the [Backend], because switches
are not thread safe, and cannot be used accross domains, but the
backend is accessed across domains. *)
Eio.Fiber.fork ~sw (fun () ->
while not @@ Atomic.get stop do
Eio.Time.sleep env#clock 0.5;
B.tick ()
done);
backend is accessed across domains.

Non-daemon so the switch awaits an in-flight export at teardown (see
[cleanup]); its waits are interruptible, so it still exits promptly. *)
Eio.Fiber.fork ~sw E.run_ticker;

(module B)

Expand Down
Loading