Files
shairport-sync/audio_ohm.c
T
nils b75edeedc3
Set variables for a Docker build (used by another workflow). / Set variables required for docker build. (push) Successful in 1s
Build and conditionally push docker image / docker-vars (push) Successful in 2s
Build and conditionally push docker image / Build and conditionally push docker image (SHAIRPORT_SYNC_BRANCH=. , ./docker/classic/Dockerfile, classic) (push) Failing after 31m52s
Build and conditionally push docker image / Build and conditionally push docker image (SHAIRPORT_SYNC_BRANCH=. NQPTP_BRANCH=${{ needs.docker-vars.outputs.nqptp_branch }} , ./docker/Dockerfile, main) (push) Failing after 31m49s
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 <name>, 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.
2026-08-28 23:54:16 +02:00

1561 lines
54 KiB
C

/*
* 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 <arpa/inet.h>
#include <errno.h>
#include <ifaddrs.h>
#include <net/if.h>
#include <netdb.h>
#include <netinet/in.h>
#include <pthread.h>
#include <stdint.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <sys/socket.h>
#include <sys/types.h>
#include <time.h>
#include <unistd.h>
#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,
"<DIDL-Lite xmlns:dc=\"http://purl.org/dc/elements/1.1/\" "
"xmlns:upnp=\"urn:schemas-upnp-org:metadata-1-0/upnp/\" "
"xmlns=\"urn:schemas-upnp-org:metadata-1-0/DIDL-Lite/\">"
"<item id=\"0\" restricted=\"True\">"
"<dc:title>%s</dc:title>"
"<upnp:class>object.item.audioItem</upnp:class>"
"<upnp:artist>%s</upnp:artist>"
"<upnp:album>%s</upnp:album>"
"</item></DIDL-Lite>",
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 <port> UDP port to listen on (default %d, 0 = ephemeral)\n",
OHM_DEFAULT_PORT);
printf(" -i <ip> IP address to bind (default: all interfaces)\n");
printf(" -l <ms> sender media latency in ms (default %d)\n", OHM_DEFAULT_LATENCY_MS);
printf(" -n <name> sender display name (default \"Shairport Sync\")\n");
printf(" -m <n> max simultaneous receivers (default %d)\n",
OHM_DEFAULT_MAX_RECEIVERS);
printf(" -H <n> resend history frames (default %d)\n", OHM_DEFAULT_HISTORY);
printf(" -r <name> control the receiver with this friendly name (UPnP)\n");
printf(" -S <min> 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,
"<DIDL-Lite xmlns:dc=\"http://purl.org/dc/elements/1.1/\" "
"xmlns:upnp=\"urn:schemas-upnp-org:metadata-1-0/upnp/\" "
"xmlns=\"urn:schemas-upnp-org:metadata-1-0/DIDL-Lite/\">"
"<item id=\"0\" restricted=\"True\">"
"<dc:title>%s</dc:title>"
"<upnp:class>object.item.audioItem</upnp:class>"
"</item></DIDL-Lite>",
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 = &parameters,
.mute = NULL};