From b75edeedc3be5861bed9b340685e622c2f17924c Mon Sep 17 00:00:00 2001 From: nils Date: Fri, 28 Aug 2026 23:48:23 +0200 Subject: [PATCH] Add ohm, an OpenHome Songcast sender audio backend Streams AirPlay audio to OpenHome/Songcast receivers as OHU unicast PCM (16/24-bit). Sends are paced on the PTP master clock so receivers slave to the AirPlay group clock through their buffer-fill servo; MediaLatency carries the phase, matching the wire behaviour of Linn's own sender (zero flags and timestamps). delay() reports the newest frame's playtime lead so shairport feeds a receiver latency ahead. Handles Join, Listen, Leave and Resend with a frame history, and forwards track metadata and metatext from the metadata hub. With -r , ohm_upnp controls the receiver over UPnP: discovered by friendly name via SSDP at time of use, woken, switched to its Songcast source and pointed at this sender when a session starts; volume follows the AirPlay slider capped at the device's VolumeLimit; standby follows some minutes after the session ends. --- .dockerignore | 5 +- Makefile.am | 4 + audio.c | 6 + audio_ohm.c | 1560 +++++++++++++++++++++++++++++++++++++++++ configure.ac | 6 + docker/Dockerfile.ohm | 66 ++ metadata/hub.h | 2 +- ohm_upnp.c | 579 +++++++++++++++ ohm_upnp.h | 58 ++ 9 files changed, 2284 insertions(+), 2 deletions(-) create mode 100644 audio_ohm.c create mode 100644 docker/Dockerfile.ohm create mode 100644 ohm_upnp.c create mode 100644 ohm_upnp.h diff --git a/.dockerignore b/.dockerignore index f5e8e83a..6b15e1ae 100644 --- a/.dockerignore +++ b/.dockerignore @@ -1,4 +1,7 @@ .github documents docker/Dockerfile -docker/classic/Dockerfile \ No newline at end of file +docker/classic/Dockerfile +reference +ohm.txt +out-ohm \ No newline at end of file diff --git a/Makefile.am b/Makefile.am index 96bf0cb4..5cd457b4 100644 --- a/Makefile.am +++ b/Makefile.am @@ -112,6 +112,10 @@ if USE_DUMMY shairport_sync_SOURCES += audio_dummy.c endif +if USE_OHM +shairport_sync_SOURCES += audio_ohm.c ohm_upnp.c +endif + if USE_AO shairport_sync_SOURCES += audio_ao.c endif diff --git a/audio.c b/audio.c index c1d27348..cc9fd49f 100644 --- a/audio.c +++ b/audio.c @@ -62,6 +62,9 @@ extern audio_output audio_pipe; #ifdef CONFIG_STDOUT extern audio_output audio_stdout; #endif +#ifdef CONFIG_OHM +extern audio_output audio_ohm; +#endif static audio_output *outputs[] = { #ifdef CONFIG_ALSA @@ -93,6 +96,9 @@ static audio_output *outputs[] = { #endif #ifdef CONFIG_DUMMY &audio_dummy, +#endif +#ifdef CONFIG_OHM + &audio_ohm, #endif NULL}; diff --git a/audio_ohm.c b/audio_ohm.c new file mode 100644 index 00000000..a7134bce --- /dev/null +++ b/audio_ohm.c @@ -0,0 +1,1560 @@ +/* + * 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}; diff --git a/configure.ac b/configure.ac index 010c1946..5e7e3649 100644 --- a/configure.ac +++ b/configure.ac @@ -98,6 +98,12 @@ if test "x$with_pipe" = "xyes" ; then fi AM_CONDITIONAL([USE_PIPE], [test "x$with_pipe" = "xyes" ]) +AC_ARG_WITH([ohm],[AS_HELP_STRING([--with-ohm],[build the OpenHome Songcast OHM/OHU unicast sender audio backend into Shairport Sync])]) +if test "x$with_ohm" = "xyes" ; then + AC_DEFINE([CONFIG_OHM], 1, [Include the OpenHome Songcast OHM/OHU unicast sender audio backend.]) +fi +AM_CONDITIONAL([USE_OHM], [test "x$with_ohm" = "xyes"]) + # Check to see if we should include the System V initscript AC_ARG_WITH([systemv-startup],[AS_HELP_STRING([--with-systemv-startup],[install a System V startup script during a make install])]) diff --git a/docker/Dockerfile.ohm b/docker/Dockerfile.ohm new file mode 100644 index 00000000..57012c1b --- /dev/null +++ b/docker/Dockerfile.ohm @@ -0,0 +1,66 @@ +# Minimal cross-build of shairport-sync with the OHM Songcast sender backend, +# for linux/arm64 (Raspberry Pi OS aarch64). Produces a glibc binary compatible +# with Debian bullseye (glibc 2.31) and later (e.g. bookworm). +# +# Usage from the repo root: +# docker buildx build --platform linux/arm64 --target binary \ +# -f docker/Dockerfile.ohm --output type=local,dest=./out-ohm . +# +# Then copy ./out-ohm/shairport-sync to the Pi. +# +# FFmpeg is taken from Debian packages (libavcodec-dev etc.) rather than built +# from source, to keep the cross-build fast under emulation. + +ARG DEBIAN_VERSION=bullseye + +FROM --platform=linux/arm64 debian:${DEBIAN_VERSION} AS builder + +ENV DEBIAN_FRONTEND=noninteractive + +RUN apt-get update && apt-get install -y --no-install-recommends \ + build-essential \ + autoconf \ + automake \ + libtool \ + pkg-config \ + git \ + ca-certificates \ + xxd \ + libasound2-dev \ + libplist-dev \ + libplist-utils \ + libssl-dev \ + libsodium-dev \ + libgcrypt-dev \ + uuid-dev \ + libavcodec-dev \ + libavformat-dev \ + libswresample-dev \ + libavutil-dev \ + libavahi-client-dev \ + avahi-daemon \ + libdbus-1-dev \ + libglib2.0-dev \ + libconfig-dev \ + libpopt-dev \ + libsndfile-dev \ + libsoxr-dev \ + xmltoman \ + && rm -rf /var/lib/apt/lists/* + +WORKDIR /shairport-sync +COPY . . + +WORKDIR /shairport-sync/build +RUN autoreconf -i ../ +RUN ../configure --sysconfdir=/etc --with-alsa \ + --with-avahi --with-ssl=openssl \ + --with-airplay-2 --with-ffmpeg --with-metadata \ + --with-ohm --with-dummy --with-pipe --with-stdout \ + --with-dbus-interface --with-mpris-interface \ + --with-convolution --with-soxr +RUN make -j"$(nproc)" + +# Target that yields just the binary on the host filesystem via --output. +FROM --platform=linux/arm64 scratch AS binary +COPY --from=builder /shairport-sync/build/shairport-sync /shairport-sync diff --git a/metadata/hub.h b/metadata/hub.h index 2d0670b2..babc662c 100644 --- a/metadata/hub.h +++ b/metadata/hub.h @@ -4,7 +4,7 @@ #include "rtsp.h" #include -#define number_of_watchers 2 +#define number_of_watchers 3 typedef enum { PS_NOT_AVAILABLE = 0, diff --git a/ohm_upnp.c b/ohm_upnp.c new file mode 100644 index 00000000..8944435a --- /dev/null +++ b/ohm_upnp.c @@ -0,0 +1,579 @@ +/* + * OpenHome receiver control for the ohm audio backend. + * + * 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. + */ + +#include "ohm_upnp.h" + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#ifdef TEST_OHM +#include +static void debug(int level, const char *fmt, ...) { + if (level > 1) + return; + va_list ap; + va_start(ap, fmt); + fprintf(stderr, "upnp: "); + vfprintf(stderr, fmt, ap); + fprintf(stderr, "\n"); + va_end(ap); +} +#define inform(...) debug(1, __VA_ARGS__) +#else +#include "config.h" +#include "common.h" +#endif + +#define NS_PER_SEC 1000000000LL +#define NS_PER_MIN (60LL * NS_PER_SEC) + +#define OHM_CTL_HTTP_MAX (128 * 1024) +#define OHM_CTL_TIMEOUT_S 3 +#define OHM_CTL_RETRY_NS (10LL * NS_PER_SEC) + +#define SSDP_ADDR "239.255.255.250" +#define SSDP_PORT 1900 +#define SSDP_LOCATIONS_MAX 8 +#define SSDP_LOCATION_LEN 192 + +#define CDB_PER_DB 100 /* shairport volume unit: dB * 100 */ +#define BMDB_PER_DB 1024 /* OpenHome volume unit: binary milli-dB */ + +typedef struct { + char type[80]; /* full urn, e.g. "urn:av-openhome-org:service:Volume:4" */ + char url[160]; /* controlURL path */ +} ohm_svc_t; + +/* Settings, copied in ohm_upnp_configure(). */ +static struct { + char *receiver; + char *iface; + char *name; + int standby_min; + int ohu_port; +} set; + +/* The discovered device, valid while .found. Tick thread only. */ +static struct { + int found; + char host[64]; + int port; + ohm_svc_t receiver, volume, product; + int vol_limit; + int vol_step; /* binary milli-dB per volume step */ +} dev; + +/* Requests from the backend callbacks to the tick thread. */ +static struct { + volatile int session; + volatile int activate; + volatile int vol_dirty; + volatile long vol_cdb; /* dB * 100, <= 0 */ + _Atomic uint64_t standby_at; /* 64-bit accesses must not tear on 32-bit targets */ +} req; + +static char *scratch; /* HTTP response buffer, tick thread only */ +static uint64_t retry_at; + +size_t ohm_xml_escape(char *dst, size_t dst_cap, const char *src) { + size_t w = 0; + dst[0] = '\0'; + if (!src) + return 0; + for (const char *s = src; *s && w + 8 < dst_cap; s++) { + switch (*s) { + case '&': + memcpy(dst + w, "&", 5); + w += 5; + break; + case '<': + memcpy(dst + w, "<", 4); + w += 4; + break; + case '>': + memcpy(dst + w, ">", 4); + w += 4; + break; + case '"': + memcpy(dst + w, """, 6); + w += 6; + break; + default: + dst[w++] = *s; + } + } + dst[w] = '\0'; + return w; +} + +static uint64_t ohm_ctl_mono_ns(void) { + struct timespec t; + clock_gettime(CLOCK_MONOTONIC, &t); + return (uint64_t)t.tv_sec * NS_PER_SEC + (uint64_t)t.tv_nsec; +} + +static long ohm_ctl_lround(double d) { return (long)(d < 0 ? d - 0.5 : d + 0.5); } + +static const char *ohm_ctl_find_ci(const char *hay, const char *needle) { + size_t n = strlen(needle); + for (; *hay; hay++) + if (strncasecmp(hay, needle, n) == 0) + return hay; + return NULL; +} + +/* ---- XML scraping ---- */ + +/* Copy text from p up to the next '<' into out. */ +static int ohm_ctl_span(const char *p, char *out, size_t cap) { + const char *e = strchr(p, '<'); + if (!e || (size_t)(e - p) >= cap) + return -1; + memcpy(out, p, e - p); + out[e - p] = '\0'; + return 0; +} + +/* Copy the text content of the first into out. */ +static int ohm_ctl_tag(const char *xml, const char *tag, char *out, size_t cap) { + char key[96]; + snprintf(key, sizeof(key), "<%s>", tag); + const char *p = strstr(xml, key); + return p ? ohm_ctl_span(p + strlen(key), out, cap) : -1; +} + +/* Fill svc from the element whose serviceType starts with prefix. */ +static int ohm_ctl_service(const char *xml, const char *prefix, ohm_svc_t *svc) { + const char *p = strstr(xml, prefix); + if (!p || ohm_ctl_span(p, svc->type, sizeof(svc->type)) != 0) + return -1; + const char *end = strstr(p, ""); + const char *cu = strstr(p, ""); + if (!cu || (end && cu > end)) + return -1; + return ohm_ctl_span(cu + strlen(""), svc->url, sizeof(svc->url)); +} + +/* ---- HTTP ---- */ + +/* Collapse chunked transfer-encoding in place; a chunk boundary must never + * split a tag. Returns the new body length. */ +static size_t ohm_ctl_dechunk(char *body, size_t body_len) { + char *end = body + body_len; + char *src = body; + char *dst = body; + while (src < end) { + char *num_end; + long sz = strtol(src, &num_end, 16); + if (num_end == src || sz <= 0) + break; + src = strstr(num_end, "\r\n"); + if (src == NULL) + break; + src += 2; + if (src + sz > end) + sz = (long)(end - src); + memmove(dst, src, (size_t)sz); + dst += sz; + src += sz; + if (src + 2 <= end && src[0] == '\r') + src += 2; + } + *dst = '\0'; + return (size_t)(dst - body); +} + +/* One request/response exchange into scratch. Returns 0 on HTTP 200. */ +static int ohm_ctl_http(const char *host, int port, const char *req_buf, size_t req_len) { + int fd = socket(AF_INET, SOCK_STREAM, 0); + if (fd < 0) + return -1; + struct timeval tv = {.tv_sec = OHM_CTL_TIMEOUT_S}; + setsockopt(fd, SOL_SOCKET, SO_RCVTIMEO, &tv, sizeof(tv)); + setsockopt(fd, SOL_SOCKET, SO_SNDTIMEO, &tv, sizeof(tv)); + struct sockaddr_in sin = {.sin_family = AF_INET, .sin_port = htons(port)}; + if (inet_pton(AF_INET, host, &sin.sin_addr) != 1 || + connect(fd, (struct sockaddr *)&sin, sizeof(sin)) != 0 || + send(fd, req_buf, req_len, 0) != (ssize_t)req_len) { + close(fd); + return -1; + } + size_t len = 0; + for (;;) { + ssize_t n = recv(fd, scratch + len, OHM_CTL_HTTP_MAX - 1 - len, 0); + if (n <= 0) + break; + len += (size_t)n; + if (len >= OHM_CTL_HTTP_MAX - 1) + break; + } + close(fd); + scratch[len] = '\0'; + + const char *sp = strchr(scratch, ' '); + if (sp == NULL || atoi(sp + 1) != 200) + return -1; + char *body = strstr(scratch, "\r\n\r\n"); + if (body == NULL) + return -1; + body += 4; + + body[-4] = '\0'; /* terminate the headers for the search */ + int chunked = ohm_ctl_find_ci(scratch, "transfer-encoding: chunked") != NULL; + body[-4] = '\r'; + if (chunked) + ohm_ctl_dechunk(body, len - (size_t)(body - scratch)); + return 0; +} + +static int ohm_ctl_get(const char *host, int port, const char *path) { + char req_buf[256]; + int n = snprintf(req_buf, sizeof(req_buf), + "GET %s HTTP/1.1\r\nHost: %s:%d\r\nConnection: close\r\n\r\n", path, host, port); + return ohm_ctl_http(host, port, req_buf, (size_t)n); +} + +static int ohm_ctl_soap(const ohm_svc_t *svc, const char *action, const char *args) { + char body[6144]; + int blen = snprintf(body, sizeof(body), + "" + "" + "%s", + action, svc->type, args, action); + char req_buf[8192]; + int rlen = snprintf(req_buf, sizeof(req_buf), + "POST %s HTTP/1.1\r\nHost: %s:%d\r\n" + "Content-Type: text/xml; charset=\"utf-8\"\r\n" + "SOAPACTION: \"%s#%s\"\r\nContent-Length: %d\r\nConnection: close\r\n\r\n%s", + svc->url, dev.host, dev.port, svc->type, action, blen, body); + return ohm_ctl_http(dev.host, dev.port, req_buf, (size_t)rlen); +} + +/* ---- Discovery ---- */ + +/* M-SEARCH for OpenHome receivers; collect response LOCATION urls. */ +static int ohm_ctl_ssdp(char locations[][SSDP_LOCATION_LEN], int max) { + int fd = socket(AF_INET, SOCK_DGRAM, 0); + if (fd < 0) + return 0; + struct timeval tv = {.tv_usec = 500000}; + setsockopt(fd, SOL_SOCKET, SO_RCVTIMEO, &tv, sizeof(tv)); + if (set.iface) { + struct in_addr mcast_if; + if (inet_pton(AF_INET, set.iface, &mcast_if) == 1) + setsockopt(fd, IPPROTO_IP, IP_MULTICAST_IF, &mcast_if, sizeof(mcast_if)); + } + const char *msearch = "M-SEARCH * HTTP/1.1\r\nHOST: " SSDP_ADDR ":1900\r\n" + "MAN: \"ssdp:discover\"\r\nMX: 1\r\n" + "ST: urn:av-openhome-org:service:Receiver:1\r\n\r\n"; + struct sockaddr_in dst = {.sin_family = AF_INET, .sin_port = htons(SSDP_PORT)}; + inet_pton(AF_INET, SSDP_ADDR, &dst.sin_addr); + sendto(fd, msearch, strlen(msearch), 0, (struct sockaddr *)&dst, sizeof(dst)); + + int n = 0; + char pkt[2048]; + for (int i = 0; i < max && n < max; i++) { + ssize_t r = recv(fd, pkt, sizeof(pkt) - 1, 0); + if (r <= 0) + break; + pkt[r] = '\0'; + const char *l = ohm_ctl_find_ci(pkt, "location:"); + if (!l) + continue; + l += strlen("location:"); + while (*l == ' ') + l++; + size_t len = strcspn(l, "\r\n"); + if (len == 0 || len >= SSDP_LOCATION_LEN) + continue; + memcpy(locations[n], l, len); + locations[n][len] = '\0'; + n++; + } + close(fd); + return n; +} + +/* Fetch one device description; on a friendly-name match fill dev. */ +static int ohm_ctl_probe(const char *location) { + char host[64], path[128], fname[128]; + int port = 80; + if (sscanf(location, "http://%63[^:/]:%d%127s", host, &port, path) != 3 && + sscanf(location, "http://%63[^:/]%127s", host, path) != 2) + return -1; + if (ohm_ctl_get(host, port, path) != 0 || + ohm_ctl_tag(scratch, "friendlyName", fname, sizeof(fname)) != 0 || + !ohm_ctl_find_ci(fname, set.receiver)) + return -1; + if (ohm_ctl_service(scratch, "urn:av-openhome-org:service:Receiver:", &dev.receiver) != 0 || + ohm_ctl_service(scratch, "urn:av-openhome-org:service:Volume:", &dev.volume) != 0 || + ohm_ctl_service(scratch, "urn:av-openhome-org:service:Product:", &dev.product) != 0) + return -1; + snprintf(dev.host, sizeof(dev.host), "%s", host); + dev.port = port; + inform("ohm: controlling receiver \"%s\" at %s:%d", fname, dev.host, dev.port); + return 0; +} + +static void ohm_ctl_read_volume_caps(void) { + dev.vol_limit = 100; + dev.vol_step = BMDB_PER_DB; + char v[32]; + if (ohm_ctl_soap(&dev.volume, "Characteristics", "") == 0) { + if (ohm_ctl_tag(scratch, "VolumeMax", v, sizeof(v)) == 0) + dev.vol_limit = atoi(v); + if (ohm_ctl_tag(scratch, "VolumeMilliDbPerStep", v, sizeof(v)) == 0 && atoi(v) > 0) + dev.vol_step = atoi(v); + } + if (ohm_ctl_soap(&dev.volume, "VolumeLimit", "") == 0 && + ohm_ctl_tag(scratch, "Value", v, sizeof(v)) == 0 && atoi(v) > 0) + dev.vol_limit = atoi(v); + debug(1, "ohm: receiver volume limit %d, %d mdB/step", dev.vol_limit, dev.vol_step); +} + +static int ohm_ctl_discover(void) { + char locations[SSDP_LOCATIONS_MAX][SSDP_LOCATION_LEN]; + int n = ohm_ctl_ssdp(locations, SSDP_LOCATIONS_MAX); + for (int i = 0; i < n; i++) { + if (ohm_ctl_probe(locations[i]) == 0) { + ohm_ctl_read_volume_caps(); + dev.found = 1; + return 0; + } + } + return -1; +} + +/* shairport delivers attenuation in cdB over the declared range; the device + * runs 0..vol_limit steps. */ +static long ohm_ctl_steps_from_cdb(long cdb) { + double step_db = (double)dev.vol_step / BMDB_PER_DB; + long steps = dev.vol_limit + ohm_ctl_lround((double)cdb / CDB_PER_DB / step_db); + if (steps < 0) + steps = 0; + if (steps > dev.vol_limit) + steps = dev.vol_limit; + return steps; +} + +/* ---- Actions ---- */ + +/* ohu:// URI of this sender as reachable from the receiver. */ +static int ohm_ctl_our_uri(char *out, size_t cap) { + struct sockaddr_in sin = {.sin_family = AF_INET, .sin_port = htons(SSDP_PORT)}; + inet_pton(AF_INET, dev.host, &sin.sin_addr); + int fd = socket(AF_INET, SOCK_DGRAM, 0); + if (fd < 0) + return -1; + struct sockaddr_in local; + socklen_t slen = sizeof(local); + int ok = connect(fd, (struct sockaddr *)&sin, sizeof(sin)) == 0 && + getsockname(fd, (struct sockaddr *)&local, &slen) == 0; + close(fd); + if (!ok) + return -1; + char ip[INET_ADDRSTRLEN]; + inet_ntop(AF_INET, &local.sin_addr, ip, sizeof(ip)); + snprintf(out, cap, "ohu://%s:%d", ip, set.ohu_port); + return 0; +} + +/* Switch the receiver's product to its Songcast source. The SourceXml body is + * XML-escaped inside the SOAP response, so count escaped elements up + * to the one of type Receiver. */ +static void ohm_ctl_select_source(void) { + if (ohm_ctl_soap(&dev.product, "SourceXml", "") != 0) + return; + const char *rcv = strstr(scratch, "<Type>Receiver</Type>"); + if (!rcv) + return; + int idx = -1; + for (const char *p = scratch; (p = strstr(p, "<Source>")) != NULL && p < rcv; p++) + idx++; + if (idx < 0) + return; + char arg[48]; + snprintf(arg, sizeof(arg), "%d", idx); + ohm_ctl_soap(&dev.product, "SetSourceIndex", arg); +} + +static int ohm_ctl_set_sender(const char *uri) { + char name[256], didl[768], meta[2304], args[2560]; + ohm_xml_escape(name, sizeof(name), set.name ? set.name : "Shairport Sync"); + snprintf(didl, sizeof(didl), + "" + "%s" + "%s" + "object.item.audioItem", + name, uri); + ohm_xml_escape(meta, sizeof(meta), didl); + snprintf(args, sizeof(args), "%s%s", uri, meta); + return ohm_ctl_soap(&dev.receiver, "SetSender", args); +} + +static int ohm_ctl_activate(void) { + char uri[128]; + if (ohm_ctl_soap(&dev.product, "SetStandby", "0") != 0) + return -1; + ohm_ctl_select_source(); + if (ohm_ctl_our_uri(uri, sizeof(uri)) != 0 || ohm_ctl_set_sender(uri) != 0 || + ohm_ctl_soap(&dev.receiver, "Play", "") != 0) + return -1; + debug(1, "ohm: receiver pointed at %s and playing", uri); + return 0; +} + +static int ohm_ctl_push_volume(void) { + long steps = ohm_ctl_steps_from_cdb(req.vol_cdb); + char arg[48]; + snprintf(arg, sizeof(arg), "%ld", steps); + if (ohm_ctl_soap(&dev.volume, "SetVolume", arg) != 0) + return -1; + debug(2, "ohm: receiver volume %ld/%d", steps, dev.vol_limit); + return 0; +} + +/* Standby only if the receiver is still pointed at us. */ +static int ohm_ctl_standby(void) { + char uri[128]; + if (ohm_ctl_soap(&dev.receiver, "Sender", "") != 0) + return -1; + if (ohm_ctl_our_uri(uri, sizeof(uri)) == 0 && !strstr(scratch, uri)) { + debug(1, "ohm: receiver plays another sender, skipping standby"); + return 0; + } + if (ohm_ctl_soap(&dev.product, "SetStandby", "1") != 0) + return -1; + debug(1, "ohm: receiver put into standby"); + return 0; +} + +/* Run every due request in order; -1 invalidates the device. */ +static int ohm_ctl_run(uint64_t now) { + if (req.activate) { + if (ohm_ctl_activate() != 0) + return -1; + req.activate = 0; + } + if (req.vol_dirty) { + req.vol_dirty = 0; /* a newer volume arriving during the push re-sets it */ + if (ohm_ctl_push_volume() != 0) { + req.vol_dirty = 1; + return -1; + } + } + if (req.standby_at && now >= req.standby_at) { + if (ohm_ctl_standby() != 0) + return -1; + req.standby_at = 0; + } + return 0; +} + +/* ---- Backend interface ---- */ + +void ohm_upnp_configure(const ohm_upnp_config_t *config) { + free(set.receiver); + free(set.iface); + free(set.name); + set.receiver = config->receiver_name ? strdup(config->receiver_name) : NULL; + set.iface = config->interface_ip ? strdup(config->interface_ip) : NULL; + set.name = config->display_name ? strdup(config->display_name) : NULL; + set.standby_min = config->standby_minutes; + set.ohu_port = config->ohu_port; +} + +int ohm_upnp_enabled(void) { return set.receiver != NULL; } + +void ohm_upnp_session_begin(void) { + if (set.receiver == NULL || req.session) + return; + req.session = 1; + req.standby_at = 0; + req.activate = 1; +} + +void ohm_upnp_session_end(void) { + if (set.receiver == NULL) + return; + req.session = 0; + req.standby_at = + set.standby_min > 0 ? ohm_ctl_mono_ns() + (uint64_t)set.standby_min * NS_PER_MIN : 0; +} + +void ohm_upnp_set_volume(long vol_cdb) { + req.vol_cdb = vol_cdb; + req.vol_dirty = 1; +} + +void ohm_upnp_tick(void) { + if (set.receiver == NULL) + return; + uint64_t now = ohm_ctl_mono_ns(); + int standby_due = req.standby_at && now >= req.standby_at; + if (!req.activate && !req.vol_dirty && !standby_due) + return; + if (now < retry_at) + return; + if (scratch == NULL && (scratch = malloc(OHM_CTL_HTTP_MAX)) == NULL) + return; + if (!dev.found && ohm_ctl_discover() != 0) { + retry_at = now + OHM_CTL_RETRY_NS; + debug(1, "ohm: receiver \"%s\" not found, retrying", set.receiver); + return; + } + if (ohm_ctl_run(now) != 0) { + dev.found = 0; + retry_at = now + OHM_CTL_RETRY_NS; + } +} + +void ohm_upnp_deinit(void) { + free(set.receiver); + set.receiver = NULL; + free(set.iface); + set.iface = NULL; + free(set.name); + set.name = NULL; + free(scratch); + scratch = NULL; + dev.found = 0; + req.session = 0; + req.activate = 0; + req.vol_dirty = 0; + req.standby_at = 0; +} diff --git a/ohm_upnp.h b/ohm_upnp.h new file mode 100644 index 00000000..c0840e01 --- /dev/null +++ b/ohm_upnp.h @@ -0,0 +1,58 @@ +/* + * OpenHome receiver control for the ohm audio backend. + * + * 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. + */ + +#pragma once + +#include + +/* XML-escape src into dst (NUL-terminated). Returns chars written (excluding NUL). */ +size_t ohm_xml_escape(char *dst, size_t dst_cap, const char *src); + +/* The backend declares volume as attenuation -OHM_UPNP_VOLUME_RANGE_CDB..0 + * (dB * 100); shairport stretches the AirPlay slider over exactly this range. */ +#define OHM_UPNP_VOLUME_RANGE_CDB 6000 + +typedef struct { + const char *receiver_name; /* friendly name to match; NULL disables the module */ + const char *interface_ip; + const char *display_name; + int standby_minutes; /* 0 = never */ + int ohu_port; +} ohm_upnp_config_t; + +/* Strings are copied. */ +void ohm_upnp_configure(const ohm_upnp_config_t *config); +int ohm_upnp_enabled(void); + +/* Callback side: cheap, thread-safe, never block. */ +void ohm_upnp_session_begin(void); +void ohm_upnp_session_end(void); +void ohm_upnp_set_volume(long vol_cdb); /* dB * 100, <= 0 */ + +/* Does the network work (SSDP discovery, SOAP); call periodically from one + * thread. */ +void ohm_upnp_tick(void); + +void ohm_upnp_deinit(void);