Set variables for a Docker build (used by another workflow). / Set variables required for docker build. (push) Successful in 28s
Build and conditionally push docker image / docker-vars (push) Successful in 28s
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 31m53s
Build and conditionally push docker image / Build and conditionally push docker image (SHAIRPORT_SYNC_BRANCH=.
, ./docker/classic/Dockerfile, classic) (push) Failing after 32m4s
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.
1494 lines
51 KiB
C
1494 lines
51 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 <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;
|
|
} cfg = {
|
|
.port = OHM_DEFAULT_PORT,
|
|
.latency_ms = OHM_DEFAULT_LATENCY_MS,
|
|
.max_receivers = OHM_DEFAULT_MAX_RECEIVERS,
|
|
.history_size = OHM_DEFAULT_HISTORY,
|
|
};
|
|
|
|
/* 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;
|
|
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 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 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 ---- */
|
|
|
|
/* XML-escape src into dst (NUL-terminated). Returns chars written (excluding NUL). */
|
|
static size_t ohm_xml_escape(char *dst, size_t dst_cap, const char *src) {
|
|
size_t w = 0;
|
|
if (dst == NULL || dst_cap == 0)
|
|
return 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;
|
|
}
|
|
|
|
/* 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) {
|
|
/* 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;
|
|
}
|
|
|
|
/* ---- 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);
|
|
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);
|
|
}
|
|
|
|
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");
|
|
|
|
/* 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:")) > 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;
|
|
default:
|
|
warn("ohm: invalid audio option \"-%c\" specified -- ignored.", optopt);
|
|
help();
|
|
break;
|
|
}
|
|
}
|
|
|
|
/* settings file (libconfig), stanza "ohm". config.cfg may be NULL if no
|
|
* configuration file was supplied / readable, so guard every lookup. */
|
|
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 (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 < 100)
|
|
warn("ohm: latency %dms is below the working range of typical receivers; expect no sync", cfg.latency_ms);
|
|
|
|
/* delay() never reports less than the receiver latency, so the feed gate
|
|
* must sit above it; the margin is the ring's working level. */
|
|
config.audio_backend_buffer_desired_length += cfg.latency_ms / 1000.0;
|
|
|
|
#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
|
|
|
|
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;
|
|
return 0;
|
|
}
|
|
|
|
static void deinit(void) {
|
|
rt.sender_running = 0;
|
|
rt.control_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_cancel(rt.sender_thread);
|
|
pthread_join(rt.sender_thread, NULL);
|
|
rt.sender_created = 0;
|
|
}
|
|
if (rt.control_created) {
|
|
pthread_cancel(rt.control_thread);
|
|
pthread_join(rt.control_thread, NULL);
|
|
rt.control_created = 0;
|
|
}
|
|
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;
|
|
}
|
|
|
|
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;
|
|
|
|
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(); }
|
|
|
|
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) {
|
|
if (the_delay == NULL)
|
|
return -1;
|
|
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) {
|
|
debug(2, "ohm: volume %.2f ignored, receiver controls volume.", vol);
|
|
}
|
|
|
|
static volume_range_t ohm_volume_range = {.minimum_volume_dB = -9630, .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};
|