package ocaml-ai-sdk
sectionYPositions = computeSectionYPositions($el), 10)"
x-init="setTimeout(() => sectionYPositions = computeSectionYPositions($el), 10)"
>
OCaml AI SDK - Provider abstraction for AI models
Install
dune-project
Dependency
Authors
Maintainers
Sources
ocaml-ai-sdk-0.1.tbz
sha256=467ed85a42617ce399d0032d81f27f8dd2e8ca9b0753daeb2cbd84bf46761370
sha512=f5e96feea16c6ffdb45f49ee3e84e9cdb5426245c1ac32acc6434c15b949cc576c1007ed463bf1362dce507d10700eb96acf5a42fa7eed4e8e491ae0214a2e4a
doc/src/ocaml-ai-sdk.ai_core/ui_message_stream_writer.ml.html
Source file ui_message_stream_writer.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 64type t = { push : Ui_message_chunk.t option -> unit; in_flight : int ref; all_done : unit Lwt.t; wake_all_done : unit Lwt.u; } let write t chunk = t.push (Some chunk) let decr_in_flight t = t.in_flight := !(t.in_flight) - 1; if !(t.in_flight) = 0 then Lwt.wakeup_later t.wake_all_done () let merge t stream = t.in_flight := !(t.in_flight) + 1; (* Lwt.async safety: unbounded push target, Lwt.catch wraps all errors, in_flight counter prevents premature stream close. See .mli for details. *) Lwt.async (fun () -> Lwt.catch (fun () -> let%lwt () = Lwt_stream.iter (fun chunk -> t.push (Some chunk)) stream in decr_in_flight t; Lwt.return_unit) (fun exn -> t.push (Some (Ui_message_chunk.Error { error_text = Printexc.to_string exn })); decr_in_flight t; Lwt.return_unit)) let create_ui_message_stream ?message_id ?(on_error = Printexc.to_string) ?on_finish ~execute () = let stream, push = Lwt_stream.create () in let all_done, wake_all_done = Lwt.wait () in let writer = { push; in_flight = ref 0; all_done; wake_all_done } in push (Some (Ui_message_chunk.Start { message_id; message_metadata = None })); (* Lwt.async safety: same rationale as merge — see .mli *) Lwt.async (fun () -> let%lwt is_aborted = Lwt.catch (fun () -> let%lwt () = execute writer in let%lwt () = if !(writer.in_flight) > 0 then writer.all_done else Lwt.return_unit in Lwt.return_false) (fun exn -> let error_text = on_error exn in push (Some (Ui_message_chunk.Error { error_text })); Lwt.return_true) in push (Some (Ui_message_chunk.Finish { finish_reason = None; message_metadata = None })); let%lwt () = match on_finish with | Some f -> Lwt.catch (fun () -> f ~finish_reason:None ~is_aborted) (fun _exn -> Lwt.return_unit) | None -> Lwt.return_unit in push None; Lwt.return_unit); stream let create_ui_message_stream_response ?(status = `OK) ?(headers = []) ?(cors = true) chunks = let sse_stream = Ui_message_stream.stream_to_sse chunks in let extra_headers = match cors with | true -> Server_handler.cors_headers @ headers | false -> headers in Server_handler.make_sse_response ~status ~extra_headers sse_stream
sectionYPositions = computeSectionYPositions($el), 10)"
x-init="setTimeout(() => sectionYPositions = computeSectionYPositions($el), 10)"
>