From 8a3428cc09beb5d42ca8344b62d0f1cebed12099 Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Sun, 4 Oct 2026 23:39:07 +0300 Subject: [PATCH] [ffi] add a C api and go wrapper for embedding hydrant Signed-off-by: dawn <90008@klbr.net> --- AGENTS.md | 3 + Cargo.lock | 11 ++ Cargo.toml | 3 + default.nix | 2 +- docs/README.md | 1 + docs/embedding.md | 62 ++++++++++ ffi/Cargo.toml | 16 +++ ffi/go/go.mod | 3 + ffi/go/hydrant.go | 85 ++++++++++++++ ffi/go/hydrant_test.go | 255 +++++++++++++++++++++++++++++++++++++++++ ffi/include/hydrant.h | 31 +++++ ffi/src/lib.rs | 216 ++++++++++++++++++++++++++++++++++ tests/ffi_embed.nu | 37 ++++++ 13 files changed, 724 insertions(+), 1 deletion(-) create mode 100644 docs/embedding.md create mode 100644 ffi/Cargo.toml create mode 100644 ffi/go/go.mod create mode 100644 ffi/go/hydrant.go create mode 100644 ffi/go/hydrant_test.go create mode 100644 ffi/include/hydrant.h create mode 100644 ffi/src/lib.rs create mode 100644 tests/ffi_embed.nu diff --git a/AGENTS.md b/AGENTS.md index 0f2d310..d6fdfb6 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -68,6 +68,8 @@ Hydrant can be used as an embedded library. The public surface is exposed via `s See `examples/statusphere.rs` for a usage example. +Non-rust programs embed hydrant through the C api in `ffi/` (`hydrant-ffi`, a workspace member built as a staticlib), which starts hydrant from the binary's `HYDRANT_*` settings via `Config::from_lookup` and leaves all data access to the HTTP API, ideally on a `unix:` bind. `ffi/go` wraps it for go. See `docs/embedding.md`. + ## General conventions ### Correctness over convenience @@ -152,6 +154,7 @@ Hydrant uses multiple `fjall` keyspaces: - `nu tests/stream_cursor_replay.nu` - Tests `/stream?cursor=0` historical event replay. - `nu tests/stream_ping.nu` - Tests ping/pong handling on `/stream`. - `nu tests/api_unix_socket.nu` - Tests serving the API on a `unix:` bind only, including its permissions and a `/stream` upgrade over the socket. +- `nu tests/ffi_embed.nu` - Builds the `hydrant-ffi` C api and runs the go tests in `ffi/go` that embed hydrant through it. Needs `go`. - `nu tests/authenticated_stream_single_relay.nu` - Tests authenticated event streaming through one relay. Requires `TEST_REPO` and `TEST_PASSWORD` in `.env`. - `nu tests/authenticated_stream_multi_relay.nu` - Tests authenticated event streaming through multiple relays. Requires `TEST_REPO` and `TEST_PASSWORD` in `.env`. - `nu tests/relay_subscribe_repos_ping.nu` - Tests ping/pong handling on relay-mode `subscribeRepos`. diff --git a/Cargo.lock b/Cargo.lock index 57fc057..989c8a0 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1664,6 +1664,17 @@ dependencies = [ "zstd", ] +[[package]] +name = "hydrant-ffi" +version = "0.1.0" +dependencies = [ + "hydrant", + "miette", + "serde_json", + "tokio", + "tracing-subscriber", +] + [[package]] name = "hyper" version = "1.9.0" diff --git a/Cargo.toml b/Cargo.toml index 2a577c8..a8ff61e 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,3 +1,6 @@ +[workspace] +members = ["ffi"] + [package] name = "hydrant" version = "0.1.0" diff --git a/default.nix b/default.nix index bd91da2..796eb11 100644 --- a/default.nix +++ b/default.nix @@ -11,7 +11,7 @@ rustPlatform.buildRustPackage { src = lib.fileset.toSource { root = ./.; fileset = lib.fileset.unions [ - ./src ./examples ./benches ./Cargo.toml ./Cargo.lock + ./src ./examples ./benches ./ffi ./Cargo.toml ./Cargo.lock ]; }; nativeBuildInputs = [cmake]; diff --git a/docs/README.md b/docs/README.md index d088de4..462ef29 100644 --- a/docs/README.md +++ b/docs/README.md @@ -12,6 +12,7 @@ you can see [random.wisp.place](https://tangled.org/did:plc:dfl62fgb7wtjj3fcbb72 - [getting started](getting-started.md): building, running, reverse proxying - [configuration](configuration.md): all environment variables - [build features](build-features.md): optional cargo features (`relay`, `backlinks`, etc.) +- [embedding](embedding.md): running hydrant inside a go (or other non-rust) program - [concepts](concepts/README.md): how the stream works, relay comparison, multi-relay support - [rest api](api/README.md): management API reference - [xrpc](xrpc/README.md): data access via XRPC diff --git a/docs/embedding.md b/docs/embedding.md new file mode 100644 index 0000000..6645450 --- /dev/null +++ b/docs/embedding.md @@ -0,0 +1,62 @@ +--- +title: embedding +--- + +hydrant can run inside a program that isn't written in rust, through a small C api in `ffi/`. the embedded hydrant still serves its usual [rest api](api/README.md) and [xrpc](xrpc/README.md), and the program talks to it over that, ideally on a unix socket. the C side only starts hydrant and learns when it stops, so no events or callbacks cross the language boundary. + +rust programs don't need any of this, use the library directly (see the [statusphere example](https://tangled.org/did:plc:dfl62fgb7wtjj3fcbb72naae/hydrant/blob/main/examples/statusphere.rs)). + +## building + +```bash +cargo build --release -p hydrant-ffi +``` + +this produces the static library `libhydrant_ffi.a` in `target/release` (or `target//release` when a target is set). the header is `ffi/include/hydrant.h`. + +## settings + +`hydrant_start` takes a json object with the same settings the binary reads from its environment (see [configuration](configuration.md)), plus `RUST_LOG`. the process environment and any `.env` file are ignored, so the embedding program's own environment can't leak in. + +`HYDRANT_API_BIND` is required, because the binary's default listens on every interface. use a unix socket, or `none` if the program doesn't need the api: + +```json +{ + "HYDRANT_DATABASE_PATH": "/var/lib/myapp/hydrant", + "HYDRANT_API_BIND": "unix:/run/myapp/hydrant.sock", + "HYDRANT_FILTER_SIGNALS": "sh.tangled.repo" +} +``` + +the socket is created owner-only (`0600`), so only the embedding program's user can reach the management endpoints, which are served on the socket and never on tcp unless `HYDRANT_API_TCP_MANAGEMENT` is set. `HYDRANT_API_SOCKET_MODE` loosens that, eg. `0660` for the socket's group too. `HYDRANT_ENABLE_DEBUG` and `HYDRANT_DEBUG_PORT` have no effect when embedded. + +## lifecycle + +- `hydrant_start` returns once the database is open and hydrant is starting. the api binds in the background, so poll it (eg. `GET /stats`) until it answers. +- `hydrant_wait` blocks until hydrant stops, which only happens when it fails (for example when its socket is already in use), and gives the reason. +- there is no stop: hydrant's background workers have no shutdown path, so it runs until the process exits. when `hydrant_wait` returns, treat it as fatal and exit, since some of that background work outlives the failure. +- only one hydrant can use a database at a time, a second `hydrant_start` on it fails. + +## go + +`ffi/go` wraps the C api for go programs. point `CGO_LDFLAGS` at the directory with the static library: + +```bash +CGO_LDFLAGS=-L$PWD/target/release go build ./... +``` + +```go +h, err := hydrant.Start(map[string]string{ + "HYDRANT_DATABASE_PATH": "/var/lib/myapp/hydrant", + "HYDRANT_API_BIND": "unix:/run/myapp/hydrant.sock", +}) +if err != nil { + return err +} +go func() { + <-h.Done() + log.Fatalf("hydrant stopped: %v", h.Err()) +}() +``` + +then reach the api with an `http.Client` whose transport dials the socket. `nu tests/ffi_embed.nu` builds the library and runs the go package's tests, which double as an example. diff --git a/ffi/Cargo.toml b/ffi/Cargo.toml new file mode 100644 index 0000000..ea12eaa --- /dev/null +++ b/ffi/Cargo.toml @@ -0,0 +1,16 @@ +[package] +name = "hydrant-ffi" +version = "0.1.0" +edition = "2024" +publish = false + +[lib] +crate-type = ["staticlib"] + +[dependencies] +# no allocator features: the embedding program owns the process, its allocator included +hydrant = { path = "..", default-features = false, features = ["indexer", "indexer_stream"] } +tokio = { version = "1.0", features = ["rt-multi-thread"] } +tracing-subscriber = { version = "0.3", features = ["env-filter"] } +serde_json = "1.0" +miette = "7" diff --git a/ffi/go/go.mod b/ffi/go/go.mod new file mode 100644 index 0000000..2a6e507 --- /dev/null +++ b/ffi/go/go.mod @@ -0,0 +1,3 @@ +module tangled.org/ptr.pet/hydrant/ffi/go + +go 1.24 diff --git a/ffi/go/hydrant.go b/ffi/go/hydrant.go new file mode 100644 index 0000000..11df5f3 --- /dev/null +++ b/ffi/go/hydrant.go @@ -0,0 +1,85 @@ +// Package hydrant runs hydrant inside a go program through its C api. +// +// hydrant keeps serving its usual http api, so after Start you talk to it over the bind you +// gave it, ideally a unix socket. link against libhydrant_ffi.a by putting its directory in +// CGO_LDFLAGS, eg. CGO_LDFLAGS=-L/path/to/target/release. +package hydrant + +/* +#cgo CFLAGS: -I${SRCDIR}/../include +#cgo LDFLAGS: -lhydrant_ffi +#cgo darwin LDFLAGS: -framework SystemConfiguration -framework Security -framework CoreFoundation -liconv -lm +#cgo linux LDFLAGS: -lgcc_s -lutil -lrt -lpthread -lm -ldl +#include +#include "hydrant.h" +*/ +import "C" + +import ( + "encoding/json" + "errors" + "fmt" + "unsafe" +) + +// Hydrant is a hydrant running in this process. it can't be stopped, because hydrant's +// background threads have no shutdown path, so it runs until the process exits. +type Hydrant struct { + done chan struct{} + err error +} + +// Start runs hydrant with the settings its binary reads from the environment, like +// HYDRANT_DATABASE_PATH and RUST_LOG, given as a map instead. the process environment is +// ignored. HYDRANT_API_BIND is required, "none" runs without an api. hydrant binds it in +// the background, so the api may need a moment before it answers. +func Start(settings map[string]string) (*Hydrant, error) { + raw, err := json.Marshal(settings) + if err != nil { + return nil, fmt.Errorf("encode hydrant settings: %w", err) + } + cSettings := C.CString(string(raw)) + defer C.free(unsafe.Pointer(cSettings)) + + var cErr *C.char + handle := C.hydrant_start(cSettings, &cErr) + if handle == nil { + return nil, takeErr(cErr) + } + + h := &Hydrant{done: make(chan struct{})} + go func() { + var cErr *C.char + if C.hydrant_wait(handle, &cErr) != 0 { + h.err = takeErr(cErr) + } else { + h.err = errors.New("hydrant exited") + } + close(h.done) + }() + return h, nil +} + +// Done is closed once hydrant stops, which only happens when it fails. treat it as fatal and +// exit, since some of hydrant's background work outlives the failure and can't be stopped. +func (h *Hydrant) Done() <-chan struct{} { + return h.done +} + +// Err says why hydrant stopped. it's nil while hydrant is running. +func (h *Hydrant) Err() error { + select { + case <-h.done: + return h.err + default: + return nil + } +} + +func takeErr(cErr *C.char) error { + if cErr == nil { + return errors.New("hydrant failed without saying why") + } + defer C.hydrant_free_string(cErr) + return errors.New(C.GoString(cErr)) +} diff --git a/ffi/go/hydrant_test.go b/ffi/go/hydrant_test.go new file mode 100644 index 0000000..af9a182 --- /dev/null +++ b/ffi/go/hydrant_test.go @@ -0,0 +1,255 @@ +package hydrant + +import ( + "bufio" + "context" + "encoding/json" + "fmt" + "io" + "net" + "net/http" + "os" + "path/filepath" + "strings" + "testing" + "time" +) + +// settings for a hydrant that never touches the network, so the test runs offline +func offlineSettings(dir string) map[string]string { + return map[string]string{ + "HYDRANT_DATABASE_PATH": filepath.Join(dir, "db"), + "HYDRANT_API_BIND": "unix:" + filepath.Join(dir, "api.sock"), + "HYDRANT_RELAY_HOSTS": "", + "HYDRANT_SEED_HOSTS": "", + "HYDRANT_CRAWLER_URLS": "", + "HYDRANT_ENABLE_CRAWLER": "false", + "RUST_LOG": "warn", + } +} + +// a short temp dir, because unix socket paths are capped at ~104 bytes on darwin +func shortTempDir(t *testing.T) string { + t.Helper() + dir, err := os.MkdirTemp("/tmp", "hydrant-ffi-") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { os.RemoveAll(dir) }) + return dir +} + +func unixClient(socket string) *http.Client { + return &http.Client{ + Timeout: 5 * time.Second, + Transport: &http.Transport{ + DialContext: func(ctx context.Context, _, _ string) (net.Conn, error) { + var d net.Dialer + return d.DialContext(ctx, "unix", socket) + }, + }, + } +} + +func waitReady(t *testing.T, h *Hydrant, client *http.Client) { + t.Helper() + deadline := time.Now().Add(30 * time.Second) + for time.Now().Before(deadline) { + select { + case <-h.Done(): + t.Fatalf("hydrant stopped before its api came up: %v", h.Err()) + default: + } + if resp, err := client.Get("http://hydrant/stats"); err == nil { + resp.Body.Close() + if resp.StatusCode == http.StatusOK { + return + } + } + time.Sleep(100 * time.Millisecond) + } + t.Fatal("hydrant api never answered on the socket") +} + +func request(t *testing.T, client *http.Client, method, path, body string) string { + t.Helper() + req, err := http.NewRequest(method, "http://hydrant"+path, strings.NewReader(body)) + if err != nil { + t.Fatal(err) + } + if body != "" { + req.Header.Set("Content-Type", "application/json") + } + resp, err := client.Do(req) + if err != nil { + t.Fatalf("%s %s: %v", method, path, err) + } + defer resp.Body.Close() + out, _ := io.ReadAll(resp.Body) + if resp.StatusCode != http.StatusOK { + t.Fatalf("%s %s: status %d: %s", method, path, resp.StatusCode, out) + } + return string(out) +} + +func TestEmbeddedHydrantServesOverUnixSocket(t *testing.T) { + dir := shortTempDir(t) + settings := offlineSettings(dir) + socket := strings.TrimPrefix(settings["HYDRANT_API_BIND"], "unix:") + + // the embedded config must come only from the settings map + envDB := filepath.Join(dir, "from-env") + t.Setenv("HYDRANT_DATABASE_PATH", envDB) + + h, err := Start(settings) + if err != nil { + t.Fatalf("start: %v", err) + } + client := unixClient(socket) + waitReady(t, h, client) + + if _, err := os.Stat(settings["HYDRANT_DATABASE_PATH"]); err != nil { + t.Errorf("database not created at the configured path: %v", err) + } + if _, err := os.Stat(envDB); err == nil { + t.Errorf("hydrant read HYDRANT_DATABASE_PATH from the process environment") + } + + info, err := os.Stat(socket) + if err != nil { + t.Fatal(err) + } + if mode := info.Mode().Perm(); mode != 0o600 { + t.Errorf("socket mode is %o, want 600", mode) + } + + request(t, client, http.MethodPatch, "/filter", `{"mode":"filter","signals":["sh.tangled.repo"]}`) + var filter struct { + Mode string `json:"mode"` + Signals []string `json:"signals"` + } + if err := json.Unmarshal([]byte(request(t, client, http.MethodGet, "/filter", "")), &filter); err != nil { + t.Fatal(err) + } + if filter.Mode != "filter" || len(filter.Signals) != 1 || filter.Signals[0] != "sh.tangled.repo" { + t.Errorf("filter didn't stick: %+v", filter) + } + + request(t, client, http.MethodGet, "/stream/head", "") + + conn, err := net.Dial("unix", socket) + if err != nil { + t.Fatal(err) + } + defer conn.Close() + fmt.Fprint(conn, "GET /stream HTTP/1.1\r\nHost: hydrant\r\nConnection: Upgrade\r\nUpgrade: websocket\r\n"+ + "Sec-WebSocket-Version: 13\r\nSec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==\r\n\r\n") + conn.SetReadDeadline(time.Now().Add(5 * time.Second)) + status, err := bufio.NewReader(conn).ReadString('\n') + if err != nil { + t.Fatalf("read upgrade response: %v", err) + } + if !strings.Contains(status, "101") { + t.Errorf("/stream didn't upgrade over the socket: %q", status) + } + + t.Run("a second hydrant on the same database is refused", func(t *testing.T) { + again := offlineSettings(dir) + again["HYDRANT_API_BIND"] = "unix:" + filepath.Join(dir, "other.sock") + _, err := Start(again) + if err == nil || !strings.Contains(err.Error(), "Locked") { + t.Fatalf("got %v, want the database lock to refuse it", err) + } + }) + + t.Run("a socket that is still served is not taken over", func(t *testing.T) { + // not shortTempDir: the refused hydrant's workers keep using its database after it + // fails, and deleting it underneath them only makes noise + otherDir, err := os.MkdirTemp("/tmp", "hydrant-ffi-") + if err != nil { + t.Fatal(err) + } + other := offlineSettings(otherDir) + other["HYDRANT_API_BIND"] = settings["HYDRANT_API_BIND"] + h2, err := Start(other) + if err != nil { + t.Fatalf("start: %v", err) + } + select { + case <-h2.Done(): + if !strings.Contains(h2.Err().Error(), "listening on this socket") { + t.Errorf("unexpected error: %v", h2.Err()) + } + case <-time.After(30 * time.Second): + t.Fatal("second hydrant kept running on a socket that was in use") + } + // the first one must still be serving + request(t, client, http.MethodGet, "/stats", "") + }) +} + +func TestSocketModeIsConfigurable(t *testing.T) { + dir := shortTempDir(t) + settings := offlineSettings(dir) + settings["HYDRANT_API_SOCKET_MODE"] = "0660" + socket := strings.TrimPrefix(settings["HYDRANT_API_BIND"], "unix:") + h, err := Start(settings) + if err != nil { + t.Fatal(err) + } + waitReady(t, h, unixClient(socket)) + + info, err := os.Stat(socket) + if err != nil { + t.Fatal(err) + } + if mode := info.Mode().Perm(); mode != 0o660 { + t.Fatalf("socket mode is %o, want 660", mode) + } +} + +func TestNoApiKeepsRunning(t *testing.T) { + dir := shortTempDir(t) + settings := offlineSettings(dir) + settings["HYDRANT_API_BIND"] = "none" + h, err := Start(settings) + if err != nil { + t.Fatal(err) + } + // with nothing to serve, a serve future that returned right away would end hydrant here + select { + case <-h.Done(): + t.Fatalf("hydrant stopped without an api: %v", h.Err()) + case <-time.After(2 * time.Second): + } + entries, err := os.ReadDir(dir) + if err != nil { + t.Fatal(err) + } + for _, e := range entries { + if e.Type()&os.ModeSocket != 0 { + t.Fatalf("no api should bind nothing, found socket %s", e.Name()) + } + } +} + +func TestStartRejectsBadSettings(t *testing.T) { + for name, tc := range map[string]struct { + mutate func(map[string]string) + want string + }{ + "missing bind": {func(s map[string]string) { delete(s, "HYDRANT_API_BIND") }, "HYDRANT_API_BIND is required"}, + "bad bind": {func(s map[string]string) { s["HYDRANT_API_BIND"] = "nope" }, "invalid HYDRANT_API_BIND"}, + "bad setting": {func(s map[string]string) { s["HYDRANT_HISTORY_TTL"] = "soon" }, "HYDRANT_HISTORY_TTL"}, + "bad mode": {func(s map[string]string) { s["HYDRANT_API_SOCKET_MODE"] = "0800" }, "HYDRANT_API_SOCKET_MODE"}, + } { + t.Run(name, func(t *testing.T) { + settings := offlineSettings(shortTempDir(t)) + tc.mutate(settings) + _, err := Start(settings) + if err == nil || !strings.Contains(err.Error(), tc.want) { + t.Fatalf("got %v, want an error mentioning %q", err, tc.want) + } + }) + } +} diff --git a/ffi/include/hydrant.h b/ffi/include/hydrant.h new file mode 100644 index 0000000..6c1aba1 --- /dev/null +++ b/ffi/include/hydrant.h @@ -0,0 +1,31 @@ +#ifndef HYDRANT_H +#define HYDRANT_H + +/* run hydrant inside another program. hydrant serves its usual http api, ideally on a + * `unix:` bind, and the embedder talks to it over that. see ffi/src/lib.rs for details. */ + +#ifdef __cplusplus +extern "C" { +#endif + +typedef struct HydrantHandle hydrant_t; + +/* starts hydrant from a json object of the HYDRANT_* settings the binary reads from its + * environment, plus RUST_LOG. HYDRANT_API_BIND is required, `none` runs without an api. + * returns NULL and sets *err on failure. the handle lives until the process exits, there is + * no stop. */ +hydrant_t *hydrant_start(const char *settings_json, char **err); + +/* blocks until hydrant stops, which only happens when it fails: returns -1 with *err set. + * returns 0 if it exited cleanly. treat either as fatal and exit, since some of hydrant's + * background work outlives the failure and can't be stopped. */ +int hydrant_wait(const hydrant_t *handle, char **err); + +/* frees an error string from this library. */ +void hydrant_free_string(char *s); + +#ifdef __cplusplus +} +#endif + +#endif diff --git a/ffi/src/lib.rs b/ffi/src/lib.rs new file mode 100644 index 0000000..a48055d --- /dev/null +++ b/ffi/src/lib.rs @@ -0,0 +1,216 @@ +//! a C api for running hydrant inside another program. +//! +//! hydrant keeps serving its usual http api, and the embedder talks to it over that, ideally +//! on a `unix:` bind. the C side only starts it and learns when it stops, so no events or +//! callbacks ever cross the boundary. + +use std::any::Any; +use std::collections::HashMap; +use std::ffi::{CStr, CString, c_char, c_int}; +use std::io::IsTerminal; +use std::panic::{AssertUnwindSafe, catch_unwind}; +use std::sync::{Arc, Condvar, Mutex, PoisonError}; + +use hydrant::config::Config; +use hydrant::control::{ApiBinds, Hydrant}; +use hydrant::deps::futures::FutureExt; + +/// a running hydrant. it has no way to stop, since hydrant's background threads have no +/// shutdown path, so it lives until the process exits. +pub struct HydrantHandle { + exit: Arc, + // never read, but dropping it would take hydrant's tasks down with it + _runtime: tokio::runtime::Runtime, +} + +#[derive(Default)] +struct Exit { + result: Mutex>>, + changed: Condvar, +} + +/// starts hydrant from a json object of the `HYDRANT_*` settings the binary reads from its +/// environment, plus `RUST_LOG`. the process environment and `.env` are ignored. +/// `HYDRANT_API_BIND` is required, since the binary's default would expose the api on every +/// interface. set it to `none` to run without an api. +/// +/// returns null and sets `*err` (free it with `hydrant_free_string`) on failure. +/// +/// # Safety +/// `settings_json` must be a valid nul-terminated string, and `err` null or writable. +#[unsafe(no_mangle)] +pub unsafe extern "C" fn hydrant_start( + settings_json: *const c_char, + err: *mut *mut c_char, +) -> *mut HydrantHandle { + let started = guard(|| { + if settings_json.is_null() { + return Err("settings_json is null".to_string()); + } + // SAFETY: the caller promises a valid nul-terminated string + let settings = unsafe { CStr::from_ptr(settings_json) }; + start( + settings + .to_str() + .map_err(|e| format!("settings_json: {e}"))?, + ) + }); + match started { + Ok(handle) => Box::into_raw(Box::new(handle)), + Err(e) => { + // SAFETY: the caller promises `err` is null or writable + unsafe { set_err(err, e) }; + std::ptr::null_mut() + } + } +} + +/// blocks until hydrant stops, which only happens when it fails, and returns -1 with `*err` +/// set to why. returns 0 if hydrant somehow exited cleanly. safe to call from several threads, +/// and again after it returned. +/// +/// treat a return as fatal and exit: some of hydrant's background work outlives the failure +/// and can't be stopped, so the process can't just start another one on the same database. +/// +/// # Safety +/// `handle` must come from `hydrant_start`, and `err` must be null or writable. +#[unsafe(no_mangle)] +pub unsafe extern "C" fn hydrant_wait( + handle: *const HydrantHandle, + err: *mut *mut c_char, +) -> c_int { + let result = guard(|| { + // SAFETY: the caller promises a handle from `hydrant_start`, and handles are never freed + let handle = unsafe { handle.as_ref() }.ok_or("handle is null")?; + let exit = &handle.exit; + let result = exit.result.lock().unwrap_or_else(PoisonError::into_inner); + let result = exit + .changed + .wait_while(result, |result| result.is_none()) + .unwrap_or_else(PoisonError::into_inner); + result + .clone() + .expect("wait_while only returns once a result is set") + }); + match result { + Ok(()) => 0, + Err(e) => { + // SAFETY: the caller promises `err` is null or writable + unsafe { set_err(err, e) }; + -1 + } + } +} + +/// frees an error string returned by this library. +/// +/// # Safety +/// `s` must be null or a string from this library that wasn't freed yet. +#[unsafe(no_mangle)] +pub unsafe extern "C" fn hydrant_free_string(s: *mut c_char) { + if !s.is_null() { + // SAFETY: the caller promises the string came from `CString::into_raw` in `set_err` + drop(unsafe { CString::from_raw(s) }); + } +} + +fn start(settings_json: &str) -> Result { + let settings: HashMap = serde_json::from_str(settings_json) + .map_err(|e| format!("settings must be a json object of strings: {e}"))?; + let lookup = |key: &str| settings.get(key).cloned(); + if !settings.contains_key("HYDRANT_API_BIND") { + return Err("HYDRANT_API_BIND is required, set it to `none` for no api".to_string()); + } + let binds = ApiBinds::from_lookup(lookup, None).map_err(report)?; + let config = Config::from_lookup(lookup).map_err(report)?; + + init_logging(settings.get("RUST_LOG").map_or("info", String::as_str)); + // fails when the embedder or an earlier start already installed one, which is fine + let _ = hydrant::deps::rustls::crypto::aws_lc_rs::default_provider().install_default(); + + let runtime = tokio::runtime::Builder::new_multi_thread() + .enable_all() + .thread_name("hydrant") + .build() + .map_err(|e| format!("failed to start the tokio runtime: {e}"))?; + let hydrant = runtime.block_on(Hydrant::new(config)).map_err(report)?; + + let exit = Arc::new(Exit::default()); + runtime.spawn({ + let exit = exit.clone(); + async move { + let serve = binds + .map(|binds| hydrant.serve(binds).boxed()) + .unwrap_or_else(|| std::future::pending().boxed()); + let result = AssertUnwindSafe(async { + tokio::select! { + r = hydrant.run()? => r, + r = serve => r, + } + }) + // tokio would swallow a panic here and hydrant_wait would block forever + .catch_unwind() + .await + .unwrap_or_else(|panic| { + Err(miette::miette!( + "hydrant panicked: {}", + panic_message(&*panic) + )) + }); + *exit.result.lock().unwrap_or_else(PoisonError::into_inner) = + Some(result.map_err(report)); + exit.changed.notify_all(); + } + }); + + Ok(HydrantHandle { + exit, + _runtime: runtime, + }) +} + +fn init_logging(filter: &str) { + let filter = tracing_subscriber::EnvFilter::builder() + .with_default_directive(tracing_subscriber::filter::LevelFilter::INFO.into()) + .parse_lossy(filter); + // a global subscriber can only be set once per process, so a second start (or an embedder + // with its own) keeps the first one + let _ = tracing_subscriber::fmt() + .with_env_filter(filter) + .with_writer(std::io::stderr) + .with_ansi(std::io::stderr().is_terminal()) + .try_init(); +} + +/// plain text, because miette's debug rendering is meant for terminals and may carry colors +fn report(e: miette::Report) -> String { + e.chain() + .map(ToString::to_string) + .collect::>() + .join(": ") +} + +/// a panic must not unwind into C, that's undefined behaviour, so it becomes an error instead. +fn guard(f: impl FnOnce() -> Result) -> Result { + catch_unwind(AssertUnwindSafe(f)) + .unwrap_or_else(|panic| Err(format!("hydrant panicked: {}", panic_message(&*panic)))) +} + +fn panic_message(panic: &(dyn Any + Send)) -> String { + panic + .downcast_ref::<&str>() + .map(|s| s.to_string()) + .or_else(|| panic.downcast_ref::().cloned()) + .unwrap_or_else(|| "unknown panic".to_string()) +} + +/// # Safety +/// `err` must be null or writable. +unsafe fn set_err(err: *mut *mut c_char, msg: String) { + if err.is_null() { + return; + } + let msg = CString::new(msg.replace('\0', "\\0")).expect("nul bytes were escaped"); + // SAFETY: checked non-null above, the caller promises it's writable + unsafe { *err = msg.into_raw() }; +} diff --git a/tests/ffi_embed.nu b/tests/ffi_embed.nu new file mode 100644 index 0000000..754ef36 --- /dev/null +++ b/tests/ffi_embed.nu @@ -0,0 +1,37 @@ +#!/usr/bin/env nu +# builds the C api and runs the go tests that embed hydrant through it. + +def main [] { + print "building hydrant-ffi..." + let out = (^cargo build -p hydrant-ffi --message-format json | complete) + if $out.exit_code != 0 { + print $out.stderr + print "=== TEST FAILED ===" + exit 1 + } + let lib = ( + $out.stdout + | lines + | each {|line| try { $line | from json } catch { null } } + | where {|m| $m != null and $m.reason? == "compiler-artifact" and $m.target?.name? == "hydrant_ffi" } + | get filenames + | flatten + | where {|f| $f | str ends-with ".a" } + | first + ) + print $"linking against ($lib)" + + let result = (with-env { CGO_ENABLED: "1", CGO_LDFLAGS: $"-L($lib | path dirname)" } { + cd ffi/go + ^go test -count=1 -v ./... | complete + }) + print $result.stdout + print $result.stderr + + if $result.exit_code == 0 { + print "=== TEST PASSED ===" + } else { + print "=== TEST FAILED ===" + exit 1 + } +} -- 2.51.2