/* * ohm output driver. This file is part of Shairport Sync. * Copyright (c) Nils Schneider 2026 * * Permission is hereby granted, free of charge, to any person * obtaining a copy of this software and associated documentation * files (the "Software"), to deal in the Software without * restriction, including without limitation the rights to use, * copy, modify, merge, publish, distribute, sublicense, and/or * sell copies of the Software, and to permit persons to whom the * Software is furnished to do so, subject to the following conditions: * * The above copyright notice and this permission notice shall be * included in all copies or substantial portions of the Software. * * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, * EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF * MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND * NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR COPYRIGHT * HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, * WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING * FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR * OTHER DEALINGS IN THE SOFTWARE. */ #ifndef _GNU_SOURCE #define _GNU_SOURCE #endif #ifdef TEST_OHM #include "test/libconfig.h" #include "test/config.h" #else #include "audio.h" #include "common.h" #include "config.h" #endif #include "ohm_upnp.h" #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #ifndef TEST_OHM #include "metadata/hub.h" #ifdef CONFIG_AIRPLAY_2 #include "ptp-utilities.h" #define OHM_HAVE_PTP 1 #endif #endif #ifndef __linux__ #ifndef TIMER_ABSTIME #define TIMER_ABSTIME 1 #endif static inline int clock_nanosleep(clockid_t clk, int flags, const struct timespec *req, struct timespec *rem) { (void)clk; (void)rem; if (flags & TIMER_ABSTIME) { struct timespec now; clock_gettime(CLOCK_MONOTONIC, &now); struct timespec rel; rel.tv_sec = req->tv_sec - now.tv_sec; rel.tv_nsec = req->tv_nsec - now.tv_nsec; if (rel.tv_nsec < 0) { rel.tv_sec--; rel.tv_nsec += 1000000000L; } if (rel.tv_sec < 0) return 0; return nanosleep(&rel, NULL); } return nanosleep(req, rem); } #endif /* ---- Protocol constants (from ohm.txt / Ohm.h) ---- */ /* The OHM signature is "Ohm " (capital O). Note: the OpenHome wiki spec * (ohm.txt) lists 0x6f ('o'), but that is a typo -- the reference * implementation (ohSongcast/Ohm.cpp: kOhm("Ohm ")) and real hardware (Linn) * both send 0x4f ('O'). Verified from on-wire packet captures. */ static const uint8_t ohm_signature[4] = {0x4f, 0x68, 0x6d, 0x20}; /* "Ohm " */ #define OHM_VERSION 1 #define OHM_HEADER_BYTES 8 #define OHM_AUDIO_SUBHEADER_BYTES 50 /* the value of the AudioHeaderLength byte */ #define OHM_MAX_MSG 16392 /* spec max message size */ #define OHM_MSG_JOIN 0 #define OHM_MSG_LISTEN 1 #define OHM_MSG_LEAVE 2 #define OHM_MSG_AUDIO 3 #define OHM_MSG_TRACK 4 #define OHM_MSG_METATEXT 5 #define OHM_MSG_SLAVE 6 #define OHM_MSG_RESEND 7 #define OHM_FLAG_HALT 0x01 #define OHM_FLAG_LOSSLESS 0x02 #define OHM_FLAG_TIMESTAMPED 0x04 #define OHM_FLAG_RESENT 0x08 #define OHM_CODEC_NAME "PCM" #define OHM_CODEC_NAME_LEN 3 /* 8 (msg header) + 50 (audio subheader, includes AudioHeaderLength + CodecNameLength * bytes) + 3 ("PCM") = 61 bytes before the PCM payload. */ #define OHM_AUDIO_PREFIX (OHM_HEADER_BYTES + 50 + OHM_CODEC_NAME_LEN) /* = 61 */ /* Compile-time guard: the 13 audio fields (Flags..Channels) must sum to 47 so * that AudioHeaderLength(1) + fields(47) + Reserved(1) + CodecNameLength(1) = 50. */ #define OHM_AUDIO_FIELDS_BYTES (1 + 2 + 4 + 4 + 4 + 4 + 8 + 8 + 4 + 4 + 2 + 1 + 1) _Static_assert(OHM_AUDIO_FIELDS_BYTES == 47, "OHM audio field layout drifted"); _Static_assert(1 + OHM_AUDIO_FIELDS_BYTES + 1 + 1 == 50, "OHM audio subheader must be 50 bytes"); _Static_assert(OHM_AUDIO_PREFIX == 61, "OHM audio prefix must be 61 bytes"); /* ---- Timing / receiver bookkeeping (mirrors ohSongcast) ---- */ #define OHM_RECEIVER_EXPIRY_NS (10LL * 1000 * 1000 * 1000) /* 10 s: ~4 missed listens */ #define OHM_RECV_TIMEOUT_MS 1000 /* for the control loop sweep */ #define OHM_DEFAULT_PORT 51970 /* Ohm::kPort */ #define OHM_DEFAULT_LATENCY_MS 300 #define OHM_DEFAULT_MAX_RECEIVERS 32 #define OHM_DEFAULT_HISTORY 100 #define OHM_TICK_MS 5 #define OHM_MIN_SEND_GAP_NS 1000000 /* burst rate limit: 1 packet per ms */ /* ---- Sender time base ---- * All sender-side time flows through these two functions so the test harness * can substitute a virtual clock. ohm_test_capture/ohm_test_idle are hooks for * the simulator; they compile to nothing in the real build. */ #ifdef TEST_OHM static uint64_t ohm_mono_ns(void); static void ohm_sleep_until_ns(uint64_t deadline_ns); static void ohm_test_capture(const uint8_t *buf, size_t len, uint64_t send_ns); static void ohm_test_idle(void); #else static uint64_t ohm_mono_ns(void) { struct timespec t; clock_gettime(CLOCK_MONOTONIC, &t); return (uint64_t)t.tv_sec * 1000000000ULL + (uint64_t)t.tv_nsec; } static void ohm_sleep_until_ns(uint64_t deadline_ns) { struct timespec next = {.tv_sec = (time_t)(deadline_ns / 1000000000ULL), .tv_nsec = (long)(deadline_ns % 1000000000ULL)}; clock_nanosleep(CLOCK_MONOTONIC, TIMER_ABSTIME, &next, NULL); } #define ohm_test_capture(buf, len, send_ns) ((void)0) #define ohm_test_idle() ((void)0) #endif /* ---- State ---- */ typedef struct { struct sockaddr_storage addr; socklen_t addr_len; uint64_t last_heard_ns; int in_use; } ohm_receiver_t; typedef struct { uint32_t frame; size_t len; uint8_t data[OHM_MAX_MSG]; } ohm_history_entry_t; typedef struct { unsigned int sample_rate; unsigned int channels; unsigned int bytes_per_sample; unsigned int bit_depth; } ohm_format_t; enum ohm_state { OHM_HALTED, OHM_PLAYING }; /* Playtime of the sample that was samples_written'th into the queue. */ typedef struct { uint64_t sample; uint64_t playtime; } ohm_time_mark_t; #define OHM_TIME_MARKS 64 /* Audio path. Everything is guarded by .lock except buf, which only the * sender thread touches. */ static struct { pthread_mutex_t lock; pthread_cond_t cond; /* signalled when data arrives or a halt is requested */ pthread_cond_t not_full; /* signalled when queue space frees up */ uint8_t *queue; /* ring buffer of PCM frames */ size_t cap; /* in frames, as are len and rd */ size_t len; size_t rd; enum ohm_state state; int halt_requested; int configured; ohm_format_t fmt; uint32_t frame; uint64_t sample_start; uint64_t samples_written; /* total frames ever queued by play() */ ohm_time_mark_t marks[OHM_TIME_MARKS]; /* newest overwrite the oldest */ int mark_head; int mark_count; uint8_t buf[OHM_MAX_MSG]; /* packet scratch, sender thread only */ } audio = { .lock = PTHREAD_MUTEX_INITIALIZER, .cond = PTHREAD_COND_INITIALIZER, .not_full = PTHREAD_COND_INITIALIZER, .state = OHM_HALTED, }; /* Receivers, resend history, and the track/metatext cache, with their packet * scratch buffers. Everything is guarded by .lock. */ static struct { pthread_mutex_t lock; ohm_receiver_t *receivers; int receiver_count; struct sockaddr_storage *bcast_addrs; /* snapshot scratch, sized with receivers */ socklen_t *bcast_lens; ohm_history_entry_t *history; /* ring buffer */ int history_head; int history_count; char *track_uri; /* always "" for a live stream */ char *track_metadata; /* DIDL-Lite */ char *metatext; /* free text */ uint32_t track_seq; uint32_t metatext_seq; uint8_t track_buf[OHM_MAX_MSG]; uint8_t metatext_buf[OHM_MAX_MSG]; uint8_t resend_buf[OHM_MAX_MSG]; } control = { .lock = PTHREAD_MUTEX_INITIALIZER, }; /* Configuration, written once in init(). */ static struct { int port; char *interface_ip; /* bind address, NULL = INADDR_ANY */ int latency_ms; char *name; /* sender display name, default "Shairport Sync" */ int max_receivers; int history_size; char *receiver_name; /* UPnP-control the receiver with this friendly name */ int standby_min; } cfg = { .port = OHM_DEFAULT_PORT, .latency_ms = OHM_DEFAULT_LATENCY_MS, .max_receivers = OHM_DEFAULT_MAX_RECEIVERS, .history_size = OHM_DEFAULT_HISTORY, .standby_min = 15, }; /* Socket and thread lifecycle. */ static struct { int sock; pthread_t control_thread; int control_running; int control_created; pthread_t sender_thread; int sender_running; int sender_created; pthread_t upnp_thread; int upnp_running; int upnp_created; int watcher_registered; } rt = { .sock = -1, }; /* ---- Helpers ---- */ /* Same clock as the playtime_ns passed to play(): CLOCK_MONOTONIC_RAW on Linux. */ #define ohm_now_ns() get_absolute_time_in_ns() static void put_be16(uint8_t *p, uint16_t v) { p[0] = (uint8_t)(v >> 8); p[1] = (uint8_t)v; } static void put_be32(uint8_t *p, uint32_t v) { p[0] = (uint8_t)(v >> 24); p[1] = (uint8_t)(v >> 16); p[2] = (uint8_t)(v >> 8); p[3] = (uint8_t)v; } static void put_be64(uint8_t *p, uint64_t v) { put_be32(p, (uint32_t)(v >> 32)); put_be32(p + 4, (uint32_t)v); } static void ohm_write_header(uint8_t *p, uint8_t type, uint16_t total_bytes) { memcpy(p, ohm_signature, 4); p[4] = OHM_VERSION; p[5] = type; put_be16(p + 6, total_bytes); } /* Endpoint comparison for sockaddr_storage (family + address + port). */ static int ohm_addr_equal(const struct sockaddr_storage *a, socklen_t alen, const struct sockaddr_storage *b, socklen_t blen) { if (alen != blen) return 0; if (a->ss_family != b->ss_family) return 0; if (a->ss_family == AF_INET) { const struct sockaddr_in *x = (const struct sockaddr_in *)a; const struct sockaddr_in *y = (const struct sockaddr_in *)b; return x->sin_port == y->sin_port && x->sin_addr.s_addr == y->sin_addr.s_addr; } if (a->ss_family == AF_INET6) { const struct sockaddr_in6 *x = (const struct sockaddr_in6 *)a; const struct sockaddr_in6 *y = (const struct sockaddr_in6 *)b; return x->sin6_port == y->sin6_port && memcmp(&x->sin6_addr, &y->sin6_addr, sizeof(struct in6_addr)) == 0; } return memcmp(a, b, alen) == 0; } /* All of the following must be called with control.lock held. */ /* Returns index of matching receiver, or -1 if absent. */ static int ohm_find_receiver_locked(const struct sockaddr_storage *addr, socklen_t addr_len) { for (int i = 0; i < cfg.max_receivers; i++) { if (control.receivers[i].in_use && ohm_addr_equal(&control.receivers[i].addr, control.receivers[i].addr_len, addr, addr_len)) return i; } return -1; } /* Add or refresh a receiver. Returns index, or -1 if full. */ static int ohm_add_receiver_locked(const struct sockaddr_storage *addr, socklen_t addr_len) { int idx = ohm_find_receiver_locked(addr, addr_len); if (idx >= 0) { control.receivers[idx].last_heard_ns = ohm_now_ns(); return idx; } for (int i = 0; i < cfg.max_receivers; i++) { if (!control.receivers[i].in_use) { memcpy(&control.receivers[i].addr, addr, addr_len); control.receivers[i].addr_len = addr_len; control.receivers[i].last_heard_ns = ohm_now_ns(); control.receivers[i].in_use = 1; control.receiver_count++; debug(2, "ohm: receiver added at frame %u, %d now listening.", audio.frame, control.receiver_count); return i; } } return -1; } static void ohm_remove_receiver_index_locked(int idx) { if (idx < 0 || idx >= cfg.max_receivers || !control.receivers[idx].in_use) return; control.receivers[idx].in_use = 0; control.receiver_count--; debug(2, "ohm: receiver removed (%d now listening).", control.receiver_count); } static void ohm_send_all_locked(const uint8_t *buf, size_t len) { for (int i = 0; i < cfg.max_receivers; i++) { if (!control.receivers[i].in_use) continue; ssize_t rc = sendto(rt.sock, buf, len, 0, (struct sockaddr *)&control.receivers[i].addr, control.receivers[i].addr_len); if (rc < 0) { char h[NI_MAXHOST]; char s[NI_MAXSERV]; getnameinfo((struct sockaddr *)&control.receivers[i].addr, control.receivers[i].addr_len, h, sizeof(h), s, sizeof(s), NI_DGRAM | NI_NUMERICHOST | NI_NUMERICSERV); debug(3, "ohm: sendto %s:%s failed: %s.", h, s, strerror(errno)); } } } static void ohm_send_one_locked(const uint8_t *buf, size_t len, const struct sockaddr_storage *addr, socklen_t addr_len) { ssize_t rc = sendto(rt.sock, buf, len, 0, (const struct sockaddr *)addr, addr_len); if (rc < 0) { char h[NI_MAXHOST]; char s[NI_MAXSERV]; getnameinfo((const struct sockaddr *)addr, addr_len, h, sizeof(h), s, sizeof(s), NI_DGRAM | NI_NUMERICHOST | NI_NUMERICSERV); debug(1, "ohm: sendto %s:%s (len=%zu) failed: %s.", h, s, len, strerror(errno)); } } /* Drop receivers we haven't heard from in OHM_RECEIVER_EXPIRY_NS. */ static void ohm_sweep_expired_locked(void) { uint64_t now = ohm_now_ns(); for (int i = 0; i < cfg.max_receivers; i++) { if (control.receivers[i].in_use && (now - control.receivers[i].last_heard_ns) > OHM_RECEIVER_EXPIRY_NS) ohm_remove_receiver_index_locked(i); } } /* ---- Message builders (return total length) ---- */ /* Copy in-use receiver addresses into the caller's arrays; returns the count. * Caller holds control.lock. */ static int ohm_snapshot_receivers_locked(struct sockaddr_storage *addrs, socklen_t *lens) { int n_rx = control.receiver_count; int rx_count = 0; for (int i = 0; i < cfg.max_receivers && rx_count < n_rx; i++) { if (control.receivers[i].in_use) { addrs[rx_count] = control.receivers[i].addr; lens[rx_count] = control.receivers[i].addr_len; rx_count++; } } return rx_count; } /* MediaLatency is in MCLK ticks: 256 * 44100 or 256 * 48000 depending on the * sample rate family. */ static uint64_t ohm_mclk_hz(void) { return ((audio.fmt.sample_rate % 441) == 0 ? 44100ULL : 48000ULL) * 256ULL; } static size_t ohm_build_audio(uint8_t *buf, uint32_t samples, uint32_t frame, uint64_t sample_start, uint8_t flags, uint32_t network_ts, uint32_t media_ts) { uint32_t media_latency = (uint32_t)((uint64_t)cfg.latency_ms * ohm_mclk_hz() / 1000); uint32_t bit_rate = audio.fmt.sample_rate * audio.fmt.channels * audio.fmt.bit_depth; size_t pcm_bytes = (size_t)samples * audio.fmt.channels * audio.fmt.bytes_per_sample; size_t total = OHM_AUDIO_PREFIX + pcm_bytes; ohm_write_header(buf, OHM_MSG_AUDIO, (uint16_t)total); uint8_t *p = buf + OHM_HEADER_BYTES; p[0] = OHM_AUDIO_SUBHEADER_BYTES; p[1] = flags; put_be16(p + 2, (uint16_t)samples); put_be32(p + 4, frame); /* Frame */ put_be32(p + 8, network_ts); /* NetworkTimestamp */ put_be32(p + 12, media_latency); /* MediaLatency */ put_be32(p + 16, media_ts); /* MediaTimestamp */ put_be64(p + 20, sample_start); /* StartSample */ put_be64(p + 28, 0); /* TotalSamples (0 = unknown, live stream) */ put_be32(p + 36, audio.fmt.sample_rate); /* SampleRate */ put_be32(p + 40, bit_rate); /* BitRate */ put_be16(p + 44, 0); /* VolumeOffset (binary milli-dB) */ p[46] = (uint8_t)audio.fmt.bit_depth; /* BitDepth */ p[47] = (uint8_t)audio.fmt.channels; /* Channels */ p[48] = 0; /* Reserved */ p[49] = OHM_CODEC_NAME_LEN; /* CodecNameLength */ memcpy(p + 50, OHM_CODEC_NAME, OHM_CODEC_NAME_LEN); return total; } static size_t ohm_build_track(uint8_t *buf, uint32_t seq, const char *uri, const char *metadata) { size_t uri_len = uri ? strlen(uri) : 0; size_t md_len = metadata ? strlen(metadata) : 0; size_t cap = OHM_MAX_MSG - OHM_HEADER_BYTES - 12; if (uri_len > cap) uri_len = cap; if (md_len > cap - uri_len) md_len = cap - uri_len; size_t total = OHM_HEADER_BYTES + 12 + uri_len + md_len; ohm_write_header(buf, OHM_MSG_TRACK, (uint16_t)total); uint8_t *p = buf + OHM_HEADER_BYTES; put_be32(p, seq); put_be32(p + 4, (uint32_t)uri_len); put_be32(p + 8, (uint32_t)md_len); if (uri_len) memcpy(p + 12, uri, uri_len); if (md_len) memcpy(p + 12 + uri_len, metadata, md_len); return total; } static size_t ohm_build_metatext(uint8_t *buf, uint32_t seq, const char *text) { size_t t_len = text ? strlen(text) : 0; size_t cap = OHM_MAX_MSG - OHM_HEADER_BYTES - 8; if (t_len > cap) t_len = cap; size_t total = OHM_HEADER_BYTES + 8 + t_len; ohm_write_header(buf, OHM_MSG_METATEXT, (uint16_t)total); uint8_t *p = buf + OHM_HEADER_BYTES; put_be32(p, seq); put_be32(p + 4, (uint32_t)t_len); if (t_len) memcpy(p + 8, text, t_len); return total; } static size_t ohm_build_control(uint8_t *buf, uint8_t type) { ohm_write_header(buf, type, OHM_HEADER_BYTES); return OHM_HEADER_BYTES; } /* ---- DIDL-Lite / metatext from the metadata hub ---- */ /* Build a DIDL-Lite metadata string from the current track. Caller frees. */ static char *ohm_build_didl_lite_locked(void) { /* Called under the metadata hub write lock (from the watcher) or not at all. * Read fields directly from metadata_store. */ char title[1024], artist[1024], album[1024]; const char *t = metadata_store.track_name; const char *a = metadata_store.artist_name; const char *al = metadata_store.album_name; ohm_xml_escape(title, sizeof(title), t); ohm_xml_escape(artist, sizeof(artist), a); ohm_xml_escape(album, sizeof(album), al); size_t cap = 2048; char *out = malloc(cap); if (!out) return NULL; int n = snprintf(out, cap, "" "" "%s" "object.item.audioItem" "%s" "%s" "", title, artist, album); if (n < 0 || (size_t)n >= cap) { free(out); return NULL; } return out; } static char *ohm_build_metatext_string_locked(void) { const char *t = metadata_store.track_name; const char *a = metadata_store.artist_name; char title[512], artist[512]; ohm_xml_escape(title, sizeof(title), t); ohm_xml_escape(artist, sizeof(artist), a); size_t cap = 1100; char *out = malloc(cap); if (!out) return NULL; if (title[0] && artist[0]) snprintf(out, cap, "%s - %s", artist, title); else if (title[0]) snprintf(out, cap, "%s", title); else if (artist[0]) snprintf(out, cap, "%s", artist); else snprintf(out, cap, "%s", cfg.name ? cfg.name : "Shairport Sync"); return out; } /* Send the cached track + metatext to a single receiver (used on Join). */ static void ohm_send_track_metatext_to_locked(const struct sockaddr_storage *addr, socklen_t addr_len) { debug(3, "ohm: send_track_metatext_to: track_metadata=%p metatext=%p track_uri=%p", (void *)control.track_metadata, (void *)control.metatext, (void *)control.track_uri); if (control.track_metadata) { size_t len = ohm_build_track(control.track_buf, control.track_seq, control.track_uri ? control.track_uri : "", control.track_metadata); ohm_send_one_locked(control.track_buf, len, addr, addr_len); } else { debug(2, "ohm: track_metadata is NULL -- sending nothing on join."); } if (control.metatext) { size_t len = ohm_build_metatext(control.metatext_buf, control.metatext_seq, control.metatext); ohm_send_one_locked(control.metatext_buf, len, addr, addr_len); } else { debug(2, "ohm: metatext is NULL -- sending nothing on join."); } } static void ohm_send_track_to_all_locked(void) { if (!control.track_metadata) return; size_t len = ohm_build_track(control.track_buf, control.track_seq, control.track_uri ? control.track_uri : "", control.track_metadata); ohm_send_all_locked(control.track_buf, len); } static void ohm_send_metatext_to_all_locked(void) { if (!control.metatext) return; size_t len = ohm_build_metatext(control.metatext_buf, control.metatext_seq, control.metatext); ohm_send_all_locked(control.metatext_buf, len); } /* ---- Metadata watcher (runs under the hub write lock) ---- */ static void ohm_metadata_watcher(struct metadata_bundle *md, __attribute__((unused)) void *userdata) { if (!rt.control_running) return; /* Decide whether this is a track-identity change worth re-sending. */ int track_changed = md->track_name_changed || md->artist_name_changed || md->album_name_changed || md->item_id_changed || md->songtime_in_milliseconds_changed || md->cover_art_pathname_changed; if (!track_changed) return; pthread_mutex_lock(&control.lock); char *new_metadata = ohm_build_didl_lite_locked(); char *new_metatext = ohm_build_metatext_string_locked(); if (new_metadata) { free(control.track_metadata); control.track_metadata = new_metadata; control.track_seq++; /* increment => "new track" rather than "new receiver joined" */ } if (new_metatext) { free(control.metatext); control.metatext = new_metatext; control.metatext_seq++; } if (control.receiver_count > 0) { ohm_send_track_to_all_locked(); ohm_send_metatext_to_all_locked(); } pthread_mutex_unlock(&control.lock); } /* ---- Control thread: Join / Listen / Leave / Resend + expiry ---- */ static void ohm_handle_resend_locked(const uint8_t *payload, size_t payload_len, const struct sockaddr_storage *from, socklen_t from_len) { if (payload_len < 4) return; uint32_t count = ((uint32_t)payload[0] << 24) | ((uint32_t)payload[1] << 16) | ((uint32_t)payload[2] << 8) | (uint32_t)payload[3]; size_t avail = payload_len - 4; uint32_t provided = (uint32_t)(avail / 4); if (count > provided) count = provided; if (count == 0) return; for (uint32_t i = 0; i < count; i++) { const uint8_t *f = payload + 4 + i * 4; uint32_t frame = ((uint32_t)f[0] << 24) | ((uint32_t)f[1] << 16) | ((uint32_t)f[2] << 8) | (uint32_t)f[3]; for (int h = 0; h < control.history_count; h++) { int idx = (control.history_head - 1 - h + cfg.history_size) % cfg.history_size; if (control.history[idx].frame == frame) { size_t len = control.history[idx].len; if (len > 0 && len <= OHM_MAX_MSG) { memcpy(control.resend_buf, control.history[idx].data, len); control.resend_buf[OHM_HEADER_BYTES + 1] |= OHM_FLAG_RESENT; ohm_send_one_locked(control.resend_buf, len, from, from_len); } break; } } } } static void *ohm_control_thread(__attribute__((unused)) void *arg) { debug(2, "ohm: control thread entered, sock=%d, running=%d.", rt.sock, rt.control_running); uint8_t buf[OHM_MAX_MSG]; while (rt.control_running) { struct sockaddr_storage from; socklen_t from_len = sizeof(from); ssize_t n = recvfrom(rt.sock, buf, sizeof(buf), 0, (struct sockaddr *)&from, &from_len); if (n < 0) { if (errno == EINTR || errno == EAGAIN || errno == EWOULDBLOCK) { /* recv timeout: sweep expired receivers */ pthread_mutex_lock(&control.lock); ohm_sweep_expired_locked(); pthread_mutex_unlock(&control.lock); continue; } if (!rt.control_running) break; debug(1, "ohm: recvfrom error: %s.", strerror(errno)); continue; } if (n < OHM_HEADER_BYTES) { debug(1, "ohm: short packet (%zd bytes), ignoring.", n); continue; } if (memcmp(buf, ohm_signature, 4) != 0) { debug(1, "ohm: signature mismatch (%02x %02x %02x %02x), ignoring.", buf[0], buf[1], buf[2], buf[3]); continue; } if (buf[4] != OHM_VERSION) { debug(1, "ohm: version %u != %u, ignoring.", buf[4], OHM_VERSION); continue; } uint8_t type = buf[5]; debug(3, "ohm: received type=%u, len=%zd.", type, n); pthread_mutex_lock(&control.lock); switch (type) { case OHM_MSG_JOIN: { debug(2, "ohm: JOIN received -> adding receiver + sending track/metatext."); int idx = ohm_add_receiver_locked(&from, from_len); if (idx >= 0) { /* Per spec: send Track + Metatext to a joining receiver. */ ohm_send_track_metatext_to_locked(&from, from_len); } break; } case OHM_MSG_LISTEN: { int idx = ohm_find_receiver_locked(&from, from_len); if (idx >= 0) { control.receivers[idx].last_heard_ns = ohm_now_ns(); debug(2, "ohm: LISTEN from known receiver %d.", idx); } else { /* Unknown listener (e.g. briefly disconnected then returned): accept it. */ debug(2, "ohm: LISTEN from unknown receiver -> adding + sending track/metatext."); ohm_add_receiver_locked(&from, from_len); ohm_send_track_metatext_to_locked(&from, from_len); } break; } case OHM_MSG_LEAVE: { int idx = ohm_find_receiver_locked(&from, from_len); if (idx >= 0) { /* Acknowledge the Leave, then drop the receiver. */ size_t len = ohm_build_control(control.resend_buf, OHM_MSG_LEAVE); ohm_send_one_locked(control.resend_buf, len, &from, from_len); ohm_remove_receiver_index_locked(idx); } break; } case OHM_MSG_RESEND: { ohm_handle_resend_locked(buf + OHM_HEADER_BYTES, (size_t)n - OHM_HEADER_BYTES, &from, from_len); break; } default: /* Join/Listen/Leave/Audio/Track/Metatext/Slave addressed to us as a sender * are not expected; ignore. */ break; } pthread_mutex_unlock(&control.lock); } return NULL; } /* UPnP receiver control runs on its own thread: discovery and SOAP block for * seconds and must never stall Join/Listen/Resend handling. */ #define OHM_UPNP_TICK_MS 250 static void *ohm_upnp_thread(__attribute__((unused)) void *arg) { const struct timespec interval = {.tv_nsec = OHM_UPNP_TICK_MS * 1000000L}; while (rt.upnp_running) { ohm_upnp_tick(); nanosleep(&interval, NULL); } return NULL; } /* ---- Paced sender thread ---- */ static void ohm_audio_unlock_cleanup(__attribute__((unused)) void *arg) { pthread_mutex_unlock(&audio.lock); } static size_t ohm_frame_size(void) { return (size_t)audio.fmt.channels * audio.fmt.bytes_per_sample; } /* Pull n frames from the front of the ring queue into dst. Caller holds * audio.lock and has checked audio.len >= n. */ static void ohm_queue_pull_locked(uint8_t *dst, size_t n) { size_t fs = ohm_frame_size(); size_t first = audio.cap - audio.rd; if (first > n) first = n; memcpy(dst, audio.queue + audio.rd * fs, first * fs); if (n > first) memcpy(dst + first * fs, audio.queue, (n - first) * fs); audio.rd = (audio.rd + n) % audio.cap; audio.len -= n; } static void ohm_broadcast_packet(const uint8_t *buf, size_t len) { struct sockaddr_storage *rx_addrs = control.bcast_addrs; socklen_t *rx_lens = control.bcast_lens; int rx_count = 0; if (rx_addrs == NULL) return; pthread_mutex_lock(&control.lock); if (control.receiver_count > 0) rx_count = ohm_snapshot_receivers_locked(rx_addrs, rx_lens); pthread_mutex_unlock(&control.lock); for (int i = 0; i < rx_count; i++) sendto(rt.sock, buf, len, 0, (struct sockaddr *)&rx_addrs[i], rx_lens[i]); } /* Discard the queue and its timing map. Caller holds audio.lock. */ static void ohm_queue_reset_locked(void) { audio.len = 0; audio.rd = 0; audio.samples_written = 0; audio.mark_head = 0; audio.mark_count = 0; } /* Record a packet in the resend history and send it to all receivers. */ static void ohm_publish(const uint8_t *buf, size_t total, uint32_t frame) { ohm_test_capture(buf, total, ohm_mono_ns()); pthread_mutex_lock(&control.lock); if (cfg.history_size > 0) { ohm_history_entry_t *he = &control.history[control.history_head]; he->frame = frame; he->len = total; memcpy(he->data, buf, total); control.history_head = (control.history_head + 1) % cfg.history_size; if (control.history_count < cfg.history_size) control.history_count++; } pthread_mutex_unlock(&control.lock); ohm_broadcast_packet(buf, total); } /* Interpolated playtime of the sample queued as number `sample`, from the * nearest mark at or before it. 0 if unknown. Caller holds audio.lock. */ static uint64_t ohm_playtime_of_locked(uint64_t sample) { int idx = -1; for (int i = 0; i < audio.mark_count; i++) { idx = (audio.mark_head - 1 - i + OHM_TIME_MARKS) % OHM_TIME_MARKS; if (audio.marks[idx].sample <= sample) break; } if (idx < 0) return 0; /* Samples before the oldest mark (priming silence) extrapolate backwards. */ if (audio.marks[idx].sample <= sample) return audio.marks[idx].playtime + (sample - audio.marks[idx].sample) * 1000000000ULL / audio.fmt.sample_rate; return audio.marks[idx].playtime - (audio.marks[idx].sample - sample) * 1000000000ULL / audio.fmt.sample_rate; } /* Robust estimate of the playtime axis: median over all marks of * (playtime - sample duration), immune to individual re-anchor steps. * 0 if no marks. Caller holds audio.lock. */ static uint64_t ohm_pt_base_locked(void) { if (audio.mark_count == 0 || audio.fmt.sample_rate == 0) return 0; uint64_t base[OHM_TIME_MARKS]; int n = audio.mark_count; for (int i = 0; i < n; i++) { const ohm_time_mark_t *m = &audio.marks[i]; base[i] = m->playtime - m->sample * 1000000000ULL / audio.fmt.sample_rate; } for (int i = 1; i < n; i++) { /* insertion sort, n <= 64 */ uint64_t v = base[i]; int j = i - 1; while (j >= 0 && base[j] > v) { base[j + 1] = base[j]; j--; } base[j + 1] = v; } return base[n / 2]; } static void *ohm_sender_run(__attribute__((unused)) void *arg) { uint64_t last_send_mono = 0; /* actual send time of the previous packet */ uint64_t master_line = 0; /* send time of the head frame, PTP master domain */ int line_valid = 0; while (rt.sender_running) { pthread_mutex_lock(&audio.lock); if (audio.halt_requested) { audio.halt_requested = 0; int was_playing = (audio.state == OHM_PLAYING); audio.state = OHM_HALTED; if (was_playing) { uint32_t frame = audio.frame; audio.frame++; size_t total = ohm_build_audio(audio.buf, 0, frame, audio.sample_start, OHM_FLAG_HALT, 0, 0); pthread_mutex_unlock(&audio.lock); ohm_publish(audio.buf, total, frame); debug(2, "ohm: halt at f=%u", frame); } else { pthread_mutex_unlock(&audio.lock); } continue; } if (!audio.configured) { ohm_test_idle(); pthread_cond_wait(&audio.cond, &audio.lock); pthread_mutex_unlock(&audio.lock); continue; } unsigned int rate = audio.fmt.sample_rate; size_t tick_frames = (size_t)(rate * OHM_TICK_MS / 1000); size_t max_frames = (OHM_MAX_MSG - OHM_AUDIO_PREFIX) / ohm_frame_size(); if (tick_frames > max_frames) tick_frames = max_frames; if (audio.state == OHM_HALTED) { audio.state = OHM_PLAYING; audio.sample_start = 0; last_send_mono = 0; line_valid = 0; } /* A short ring means we outran play(), not that the stream ended: wait * for supply. A genuine end arrives as a halt via flush()/stop(). */ if (audio.len < tick_frames) { ohm_test_idle(); pthread_cond_wait(&audio.cond, &audio.lock); pthread_mutex_unlock(&audio.lock); continue; } /* Phase once, rate from the PTP master. The playtime axis is not a clean * clock (shairport re-anchors it in steps), so it only sets the anchor: * pt - L, carried into the master domain (master = local + offset, * rtp.c). From there the line advances strictly at the sample rate on the * master clock and is converted back per packet, so nqptp's discipline * flows into our cadence. Without timing marks, free-run at one tick per * tick. */ uint64_t pt_base = ohm_pt_base_locked(); uint64_t pt = pt_base ? pt_base + (audio.samples_written - audio.len) * 1000000000ULL / rate : 0; uint64_t raw_now = ohm_now_ns(); uint64_t mono_now = ohm_mono_ns(); uint64_t tick_ns = (uint64_t)tick_frames * 1000000000ULL / rate; uint64_t ptp_off = 0; #ifdef OHM_HAVE_PTP uint64_t ptp_id; ptp_get_clock_info(&ptp_id, NULL, &ptp_off, NULL); #endif uint64_t send_raw; uint64_t deadline; int64_t line_err = 0; if (pt != 0) { uint64_t target = pt - (uint64_t)cfg.latency_ms * 1000000ULL + ptp_off; if (!line_valid) { master_line = target; line_valid = 1; } line_err = (int64_t)target - (int64_t)master_line; if (line_err > 500000000LL || line_err < -500000000LL) { debug(1, "ohm: re-anchoring send line, drifted %lldms", (long long)(line_err / 1000000)); master_line = target; line_err = 0; } else { /* Trickle toward the median target: tracks slow meander of the * playtime axis at an inaudible rate, never in jumps. */ int64_t adj = line_err; int64_t max_adj = (int64_t)(tick_ns / 10000); /* 100 ppm */ if (adj > max_adj) adj = max_adj; if (adj < -max_adj) adj = -max_adj; master_line += adj; } send_raw = master_line - ptp_off; deadline = (uint64_t)((int64_t)mono_now + ((int64_t)send_raw - (int64_t)raw_now)); } else { send_raw = raw_now; deadline = last_send_mono ? last_send_mono + tick_ns : mono_now; } if (last_send_mono && deadline < last_send_mono + OHM_MIN_SEND_GAP_NS) deadline = last_send_mono + OHM_MIN_SEND_GAP_NS; if (deadline > mono_now) { pthread_mutex_unlock(&audio.lock); ohm_sleep_until_ns(deadline); continue; } size_t pull = (audio.len < tick_frames) ? audio.len : tick_frames; uint32_t frame = audio.frame; uint64_t sample = audio.sample_start; audio.frame++; audio.sample_start += pull; ohm_queue_pull_locked(audio.buf + OHM_AUDIO_PREFIX, pull); pthread_cond_signal(&audio.not_full); size_t queued = audio.len; last_send_mono = mono_now; if (line_valid) master_line += (uint64_t)pull * 1000000000ULL / rate; size_t total = ohm_build_audio(audio.buf, (uint32_t)pull, frame, sample, 0, 0, 0); pthread_mutex_unlock(&audio.lock); ohm_publish(audio.buf, total, frame); debug(3, "ohm: send f=%u s=%zu start=%llu q=%zu wt=%llums pt=%llums d=%lldms le=%lldms ptpoff=%llu", frame, pull, (unsigned long long)sample, queued, (unsigned long long)(raw_now / 1000000), (unsigned long long)(pt / 1000000), pt ? (long long)(((int64_t)pt - (int64_t)raw_now) / 1000000) : 0, (long long)(line_err / 1000000), (unsigned long long)ptp_off); } return NULL; } /* ---- Startup helpers ---- */ static void ohm_print_uris(int port) { if (cfg.interface_ip) { inform("ohm: OHU sender URI: ohu://%s:%d", cfg.interface_ip, port); return; } struct ifaddrs *ifa = NULL; if (getifaddrs(&ifa) != 0) { inform("ohm: OHU sender listening on port %d (could not enumerate interfaces).", port); return; } for (struct ifaddrs *p = ifa; p; p = p->ifa_next) { if (!p->ifa_addr || p->ifa_addr->sa_family != AF_INET) continue; if (p->ifa_flags & IFF_LOOPBACK) continue; if (!(p->ifa_flags & IFF_UP)) continue; char host[NI_MAXHOST]; if (getnameinfo(p->ifa_addr, sizeof(struct sockaddr_in), host, sizeof(host), NULL, 0, NI_DGRAM | NI_NUMERICHOST) == 0) { inform("ohm: OHU sender URI: ohu://%s:%d", host, port); } } freeifaddrs(ifa); } static int ohm_start_socket(void) { rt.sock = socket(AF_INET, SOCK_DGRAM, 0); if (rt.sock < 0) { warn("ohm: cannot create socket: %s.", strerror(errno)); return -1; } int yes = 1; setsockopt(rt.sock, SOL_SOCKET, SO_REUSEADDR, &yes, sizeof(yes)); struct sockaddr_in sin; memset(&sin, 0, sizeof(sin)); sin.sin_family = AF_INET; sin.sin_port = htons((uint16_t)cfg.port); if (cfg.interface_ip && *cfg.interface_ip) inet_pton(AF_INET, cfg.interface_ip, &sin.sin_addr); else sin.sin_addr.s_addr = htonl(INADDR_ANY); if (bind(rt.sock, (struct sockaddr *)&sin, sizeof(sin)) < 0) { warn("ohm: cannot bind port %d: %s.", cfg.port, strerror(errno)); close(rt.sock); rt.sock = -1; return -1; } socklen_t slen = sizeof(sin); getsockname(rt.sock, (struct sockaddr *)&sin, &slen); int bound_port = ntohs(sin.sin_port); struct timeval tv; tv.tv_sec = OHM_RECV_TIMEOUT_MS / 1000; tv.tv_usec = (OHM_RECV_TIMEOUT_MS % 1000) * 1000; setsockopt(rt.sock, SOL_SOCKET, SO_RCVTIMEO, &tv, sizeof(tv)); ohm_print_uris(bound_port); ohm_upnp_config_t upnp = { .receiver_name = cfg.receiver_name, .interface_ip = cfg.interface_ip, .display_name = cfg.name, .standby_minutes = cfg.standby_min, .ohu_port = bound_port, }; ohm_upnp_configure(&upnp); return 0; } /* ---- Backend interface ---- */ static void help(void) { printf(" -o UDP port to listen on (default %d, 0 = ephemeral)\n", OHM_DEFAULT_PORT); printf(" -i IP address to bind (default: all interfaces)\n"); printf(" -l sender media latency in ms (default %d)\n", OHM_DEFAULT_LATENCY_MS); printf(" -n sender display name (default \"Shairport Sync\")\n"); printf(" -m max simultaneous receivers (default %d)\n", OHM_DEFAULT_MAX_RECEIVERS); printf(" -H resend history frames (default %d)\n", OHM_DEFAULT_HISTORY); printf(" -r control the receiver with this friendly name (UPnP)\n"); printf(" -S minutes after stop before receiver standby (default 15, 0 = never)\n"); } static int init(__attribute__((unused)) int argc, __attribute__((unused)) char **argv) { config.audio_backend_buffer_desired_length = 0.15; cfg.port = OHM_DEFAULT_PORT; cfg.latency_ms = OHM_DEFAULT_LATENCY_MS; cfg.max_receivers = OHM_DEFAULT_MAX_RECEIVERS; cfg.history_size = OHM_DEFAULT_HISTORY; cfg.name = strdup("Shairport Sync"); /* settings file (libconfig), stanza "ohm". config.cfg may be NULL if no * configuration file was supplied / readable, so guard every lookup. * Command-line options are parsed afterwards and take priority. */ if (config.cfg != NULL) { const char *str; int value; if (config_lookup_non_empty_string(config.cfg, "ohm.port", &str)) { cfg.port = atoi(str); } else if (config_lookup_int(config.cfg, "ohm.port", &value)) { cfg.port = value; } if (config_lookup_non_empty_string(config.cfg, "ohm.interface", &str)) { free(cfg.interface_ip); cfg.interface_ip = strdup(str); } if (config_lookup_int(config.cfg, "ohm.latency", &value)) { cfg.latency_ms = value; } if (config_lookup_non_empty_string(config.cfg, "ohm.name", &str)) { free(cfg.name); cfg.name = strdup(str); } if (config_lookup_int(config.cfg, "ohm.max_receivers", &value)) { cfg.max_receivers = value; } if (config_lookup_int(config.cfg, "ohm.history", &value)) { cfg.history_size = value; } if (config_lookup_non_empty_string(config.cfg, "ohm.receiver", &str)) { free(cfg.receiver_name); cfg.receiver_name = strdup(str); } if (config_lookup_int(config.cfg, "ohm.standby_minutes", &value)) { cfg.standby_min = value; } } /* command line (shairport passes this backend's arguments as a slice) */ optind = 1; /* optind=0 is equivalent to optind=1 plus special behaviour */ argv--; /* shift so getopt() sees a conventional argv[0..argc-1] */ argc++; opterr = 0; int opt; while ((opt = getopt(argc, argv, "o:i:l:n:m:H:r:S:")) > 0) { switch (opt) { case 'o': cfg.port = atoi(optarg); break; case 'i': free(cfg.interface_ip); cfg.interface_ip = strdup(optarg); break; case 'l': cfg.latency_ms = atoi(optarg); break; case 'n': free(cfg.name); cfg.name = strdup(optarg); break; case 'm': cfg.max_receivers = atoi(optarg); break; case 'H': cfg.history_size = atoi(optarg); break; case 'r': free(cfg.receiver_name); cfg.receiver_name = strdup(optarg); break; case 'S': cfg.standby_min = atoi(optarg); break; default: warn("ohm: invalid audio option \"-%c\" specified -- ignored.", optopt); help(); break; } } if (cfg.max_receivers < 1) cfg.max_receivers = OHM_DEFAULT_MAX_RECEIVERS; if (cfg.history_size < 0) cfg.history_size = OHM_DEFAULT_HISTORY; if (cfg.latency_ms < 1) { warn("ohm: latency %dms is invalid, using %dms", cfg.latency_ms, OHM_DEFAULT_LATENCY_MS); cfg.latency_ms = OHM_DEFAULT_LATENCY_MS; } else if (cfg.latency_ms < 100) { warn("ohm: latency %dms is below the working range of typical receivers; expect no sync", cfg.latency_ms); } #ifndef CONFIG_AIRPLAY_2 parse_audio_options("ohm", (1 << SPS_FORMAT_S16_BE) | (1 << SPS_FORMAT_S24_3BE), (1 << SPS_RATE_44100) | (1 << SPS_RATE_48000), (1 << 2)); #else parse_audio_options("ohm", (1 << SPS_FORMAT_S16_BE) | (1 << SPS_FORMAT_S24_3BE), (1 << SPS_RATE_48000), (1 << 2)); #endif /* delay() never reports less than the receiver latency, so the feed gate * must sit above it; the margin is the ring's working level. * parse_audio_options() may have replaced the value, so add the latency * afterwards. */ config.audio_backend_buffer_desired_length += cfg.latency_ms / 1000.0; control.receivers = calloc(cfg.max_receivers, sizeof(ohm_receiver_t)); control.history = calloc(cfg.history_size > 0 ? cfg.history_size : 1, sizeof(ohm_history_entry_t)); control.bcast_addrs = calloc(cfg.max_receivers, sizeof(struct sockaddr_storage)); control.bcast_lens = calloc(cfg.max_receivers, sizeof(socklen_t)); if (!control.receivers || !control.history || !control.bcast_addrs || !control.bcast_lens) { warn("ohm: out of memory allocating receiver/history tables."); return -1; } if (ohm_start_socket() < 0) return -1; /* Seed the track cache with a sender-level DIDL-Lite so joining receivers * get something sensible before the first AirPlay track metadata arrives. */ { size_t cap = 1024; char *seed = malloc(cap); if (seed) snprintf(seed, cap, "" "" "%s" "object.item.audioItem" "", cfg.name ? cfg.name : "Shairport Sync"); control.track_metadata = seed; control.track_uri = strdup(""); control.metatext = strdup(cfg.name ? cfg.name : "Shairport Sync"); } rt.control_running = 1; debug(2, "ohm: starting control thread..."); if (pthread_create(&rt.control_thread, NULL, ohm_control_thread, NULL) != 0) { warn("ohm: cannot create control thread: %s.", strerror(errno)); rt.control_running = 0; return -1; } rt.control_created = 1; debug(2, "ohm: control thread started (created=%d).", rt.control_created); rt.sender_running = 1; if (pthread_create(&rt.sender_thread, NULL, ohm_sender_run, NULL) != 0) { warn("ohm: cannot create sender thread: %s.", strerror(errno)); rt.sender_running = 0; return -1; } rt.sender_created = 1; if (ohm_upnp_enabled()) { rt.upnp_running = 1; if (pthread_create(&rt.upnp_thread, NULL, ohm_upnp_thread, NULL) != 0) { warn("ohm: cannot create upnp thread: %s.", strerror(errno)); rt.upnp_running = 0; return -1; } rt.upnp_created = 1; } return 0; } static void deinit(void) { /* Cooperative shutdown: never pthread_cancel, a thread cancelled inside a * lock-held region would die holding the lock. */ rt.sender_running = 0; rt.control_running = 0; rt.upnp_running = 0; pthread_mutex_lock(&audio.lock); pthread_cond_broadcast(&audio.cond); pthread_cond_broadcast(&audio.not_full); pthread_mutex_unlock(&audio.lock); if (rt.sock >= 0) { shutdown(rt.sock, SHUT_RDWR); } if (rt.sender_created) { pthread_join(rt.sender_thread, NULL); rt.sender_created = 0; } if (rt.control_created) { pthread_join(rt.control_thread, NULL); rt.control_created = 0; } if (rt.upnp_created) { pthread_join(rt.upnp_thread, NULL); rt.upnp_created = 0; } /* The metadata hub has no unregister; the watcher bails out on * !control_running. Cycling the hub write lock waits out any watcher that * passed that check before the flag was cleared, so the frees below cannot * race it. */ #ifndef TEST_OHM if (rt.watcher_registered) { metadata_hub_modify_prolog(); metadata_hub_modify_epilog(0); } #endif if (rt.sock >= 0) { close(rt.sock); rt.sock = -1; } free(control.receivers); control.receivers = NULL; free(control.history); control.history = NULL; free(audio.queue); audio.queue = NULL; audio.cap = 0; audio.len = 0; free(control.track_uri); control.track_uri = NULL; free(control.track_metadata); control.track_metadata = NULL; free(control.metatext); control.metatext = NULL; free(cfg.interface_ip); cfg.interface_ip = NULL; free(cfg.name); cfg.name = NULL; free(cfg.receiver_name); cfg.receiver_name = NULL; ohm_upnp_deinit(); } static int prepare(void) { /* Not in init(): that runs before metadata_hub_init(), whose memset of * metadata_store would wipe the registration. */ if (!rt.watcher_registered) { add_metadata_watcher(ohm_metadata_watcher, NULL); rt.watcher_registered = 1; } pthread_mutex_lock(&audio.lock); audio.halt_requested = 1; pthread_cond_broadcast(&audio.cond); pthread_cond_broadcast(&audio.not_full); pthread_mutex_unlock(&audio.lock); return 0; } static int ohm_check_configuration(__attribute__((unused)) unsigned int channels, __attribute__((unused)) unsigned int rate, __attribute__((unused)) unsigned int format) { return (format == SPS_FORMAT_S16_BE || format == SPS_FORMAT_S24_3BE) ? 0 : -1; } static int32_t get_configuration(unsigned int channels, unsigned int rate, unsigned int format) { return search_for_suitable_configuration(channels, rate, format, ohm_check_configuration); } static int configure(int32_t requested_encoded_format, __attribute__((unused)) char **channel_map) { unsigned int fmt = FORMAT_FROM_ENCODED_FORMAT(requested_encoded_format); ohm_format_t new_fmt = { .sample_rate = RATE_FROM_ENCODED_FORMAT(requested_encoded_format), .channels = CHANNELS_FROM_ENCODED_FORMAT(requested_encoded_format), }; switch (fmt) { case SPS_FORMAT_S16_BE: new_fmt.bytes_per_sample = 2; new_fmt.bit_depth = 16; break; case SPS_FORMAT_S24_3BE: new_fmt.bytes_per_sample = 3; new_fmt.bit_depth = 24; break; default: debug(1, "ohm: unsupported output format %u in configure(); rejecting.", fmt); return -1; } double qlen = config.audio_backend_buffer_desired_length * 2.0; size_t new_cap = (size_t)(qlen * new_fmt.sample_rate); /* frames */ size_t new_bytes = new_cap * new_fmt.channels * new_fmt.bytes_per_sample; pthread_mutex_lock(&audio.lock); /* Re-configuring with an unchanged format must not disturb the stream: * shairport calls configure() repeatedly, and discarding the queue here * destroys its priming and forces an endless restart cycle. */ if (audio.configured && memcmp(&audio.fmt, &new_fmt, sizeof(new_fmt)) == 0) { pthread_mutex_unlock(&audio.lock); return 0; } size_t old_bytes = audio.cap * ohm_frame_size(); if (new_bytes != old_bytes) { uint8_t *newq = realloc(audio.queue, new_bytes); if (newq == NULL) { pthread_mutex_unlock(&audio.lock); warn("ohm: out of memory resizing the audio queue to %zu bytes.", new_bytes); return -1; } audio.queue = newq; } audio.cap = new_cap; ohm_queue_reset_locked(); audio.halt_requested = 1; audio.fmt = new_fmt; audio.configured = 1; pthread_cond_broadcast(&audio.not_full); pthread_cond_broadcast(&audio.cond); pthread_mutex_unlock(&audio.lock); debug(2, "ohm: configured for %u-bit, %u-channel, %u Hz (queue %zu frames).", audio.fmt.bit_depth, audio.fmt.channels, audio.fmt.sample_rate, audio.cap); return 0; } static int play(void *buf, int samples, int sample_type, __attribute__((unused)) uint32_t timestamp, uint64_t playtime) { if (!audio.configured || samples <= 0) return samples; ohm_upnp_session_begin(); uint64_t wallclock = ohm_now_ns(); if (sample_type == play_samples_are_untimed) debug(3, "ohm: play samples=%d untimed wall=%llums", samples, (unsigned long long)(wallclock / 1000000)); else debug(3, "ohm: play samples=%d playtime=%llums wall=%llums lead=%lldms", samples, (unsigned long long)(playtime / 1000000), (unsigned long long)(wallclock / 1000000), (long long)(((int64_t)playtime - (int64_t)wallclock) / 1000000)); const uint8_t *src = (const uint8_t *)buf; pthread_mutex_lock(&audio.lock); pthread_cleanup_push(ohm_audio_unlock_cleanup, NULL); if (sample_type != play_samples_are_untimed) { audio.marks[audio.mark_head] = (ohm_time_mark_t){.sample = audio.samples_written, .playtime = playtime}; audio.mark_head = (audio.mark_head + 1) % OHM_TIME_MARKS; if (audio.mark_count < OHM_TIME_MARKS) audio.mark_count++; } size_t fs = ohm_frame_size(); size_t remain = (rt.sender_running && audio.cap > 0) ? (size_t)samples : 0; while (remain > 0) { while (audio.len >= audio.cap && rt.sender_running) { pthread_cond_signal(&audio.cond); pthread_cond_wait(&audio.not_full, &audio.lock); } if (!rt.sender_running) break; size_t wr = (audio.rd + audio.len) % audio.cap; size_t space = audio.cap - audio.len; size_t first = audio.cap - wr; if (first > space) first = space; if (first > remain) first = remain; memcpy(audio.queue + wr * fs, src, first * fs); src += first * fs; audio.len += first; remain -= first; audio.samples_written += first; } pthread_cond_signal(&audio.cond); pthread_cleanup_pop(1); return samples; } /* Discard queued audio; the sender emits a HALT packet on its next tick. */ static void ohm_request_halt(void) { pthread_mutex_lock(&audio.lock); ohm_queue_reset_locked(); audio.halt_requested = 1; pthread_cond_broadcast(&audio.cond); pthread_cond_broadcast(&audio.not_full); pthread_mutex_unlock(&audio.lock); } static void stop(void) { ohm_request_halt(); ohm_upnp_session_end(); } static void flush(void) { ohm_request_halt(); } static int is_running(void) { return 0; } /* Time until audio handed to play() is audible. The newest queued frame is * audible at its playtime (the sender sends it at pt - L, the receiver holds * L), so the report is simply its playtime lead. Without timing marks fall * back to occupancy plus the receiver latency. */ static int ohm_delay(long *the_delay) { pthread_mutex_lock(&audio.lock); if (!audio.configured || audio.fmt.sample_rate == 0) { pthread_mutex_unlock(&audio.lock); return -1; } unsigned int rate = audio.fmt.sample_rate; /* Queued audio plus the receiver hold: exact when the sender is on * schedule or behind. */ uint64_t occupancy = audio.len + (uint64_t)cfg.latency_ms * rate / 1000; /* While the queue fills ahead of the schedule, the newest frame is audible * at its playtime, not after draining: its lead is the truth then. In * steady state both are identical. */ uint64_t lead = 0; if (audio.len > 0) { uint64_t pt = ohm_playtime_of_locked(audio.samples_written - 1); uint64_t now = ohm_now_ns(); if (pt > now) lead = (pt - now) * rate / 1000000000ULL; } *the_delay = (long)(lead > occupancy ? lead : occupancy); pthread_mutex_unlock(&audio.lock); return 0; } static void volume(double vol) { if (ohm_upnp_enabled()) ohm_upnp_set_volume((long)vol); else debug(2, "ohm: volume %.2f ignored, receiver controls volume.", vol); } static volume_range_t ohm_volume_range = {.minimum_volume_dB = -OHM_UPNP_VOLUME_RANGE_CDB, .maximum_volume_dB = 0}; static output_parameters_t ohm_parameters = {.volume_range = &ohm_volume_range}; static output_parameters_t *parameters(void) { return &ohm_parameters; } audio_output audio_ohm = {.name = "ohm", .help = &help, .init = &init, .deinit = &deinit, .prepare = &prepare, .get_configuration = &get_configuration, .configure = &configure, .start = NULL, .stop = &stop, .is_running = &is_running, .flush = &flush, .delay = &ohm_delay, .stats = NULL, .play = &play, .volume = &volume, .parameters = ¶meters, .mute = NULL};