network addition
This commit is contained in:
parent
3bd713d667
commit
2099d3d64d
27 changed files with 2237 additions and 40 deletions
273
lib/remote/host_server.dart
Normal file
273
lib/remote/host_server.dart
Normal file
|
|
@ -0,0 +1,273 @@
|
|||
// `controller` can't be a private-named parameter, so an initializing formal
|
||||
// isn't possible for `_controller` (same reason as playback_engine.dart).
|
||||
// ignore_for_file: prefer_initializing_formals
|
||||
|
||||
import 'dart:async';
|
||||
import 'dart:io';
|
||||
|
||||
import '../playback/playback_engine.dart';
|
||||
import 'command_dispatch.dart';
|
||||
import 'messages.dart';
|
||||
|
||||
/// Turns the local device into a controllable player: runs an in-process
|
||||
/// WebSocket server, authenticates remotes against the shared server password,
|
||||
/// streams playback snapshots to them, and applies the commands they send to
|
||||
/// the local [PlaybackController].
|
||||
///
|
||||
/// Lifecycle is owned by `remote_providers.dart`, which starts the host while
|
||||
/// this device is playable and advertises its [port] over mDNS
|
||||
/// (`discovery.dart`). The server binds an ephemeral port on all IPv4
|
||||
/// interfaces; discovery carries the port to remotes.
|
||||
class RemoteHost {
|
||||
RemoteHost({
|
||||
required PlaybackController controller,
|
||||
required this.serverKey,
|
||||
required this.password,
|
||||
required this.deviceName,
|
||||
}) : _controller = controller;
|
||||
|
||||
final PlaybackController _controller;
|
||||
|
||||
/// Only remotes reporting this same [serverKey] (i.e. signed into the same
|
||||
/// Subsonic server) are allowed to connect — song ids must resolve on both.
|
||||
final String serverKey;
|
||||
|
||||
/// Shared secret for the HMAC handshake — never sent, only proven.
|
||||
final String password;
|
||||
|
||||
/// Friendly name shown to remotes ("Living Room iPad").
|
||||
final String deviceName;
|
||||
|
||||
HttpServer? _server;
|
||||
void Function()? _removeListener;
|
||||
|
||||
final Set<WebSocket> _clients = {};
|
||||
final StreamController<int> _clientCount =
|
||||
StreamController<int>.broadcast();
|
||||
|
||||
/// Monotonic snapshot sequence so a remote can drop stale/out-of-order state.
|
||||
int _seq = 0;
|
||||
|
||||
/// The last state actually put on the wire, for cheap change detection.
|
||||
PlaybackState? _sentState;
|
||||
Timer? _throttle;
|
||||
PlaybackState? _pending;
|
||||
|
||||
/// Bound port once [start] has run, else null.
|
||||
int? get port => _server?.port;
|
||||
|
||||
/// Emits the number of connected remotes (drives the "N remotes" UI).
|
||||
Stream<int> get clientCountStream => _clientCount.stream;
|
||||
|
||||
int get clientCount => _clients.length;
|
||||
|
||||
/// Bind and begin serving. Returns the chosen ephemeral port. Idempotent —
|
||||
/// a second call returns the already-bound port.
|
||||
Future<int> start() async {
|
||||
final existing = _server;
|
||||
if (existing != null) return existing.port;
|
||||
|
||||
final server = await HttpServer.bind(InternetAddress.anyIPv4, 0);
|
||||
_server = server;
|
||||
server.listen((req) async {
|
||||
if (!WebSocketTransformer.isUpgradeRequest(req)) {
|
||||
req.response
|
||||
..statusCode = HttpStatus.forbidden
|
||||
..close();
|
||||
return;
|
||||
}
|
||||
try {
|
||||
final ws = await WebSocketTransformer.upgrade(req);
|
||||
_handleSocket(ws);
|
||||
} catch (_) {
|
||||
// Failed upgrade — nothing to clean up.
|
||||
}
|
||||
});
|
||||
|
||||
// fireImmediately:false — the per-client welcome snapshot covers the
|
||||
// initial state; we only want deltas here.
|
||||
_removeListener =
|
||||
_controller.addListener(_onStateChanged, fireImmediately: false);
|
||||
return server.port;
|
||||
}
|
||||
|
||||
/// Stop serving and drop every remote. Safe to call when not started.
|
||||
Future<void> stop() async {
|
||||
_removeListener?.call();
|
||||
_removeListener = null;
|
||||
_throttle?.cancel();
|
||||
_throttle = null;
|
||||
_pending = null;
|
||||
_sentState = null;
|
||||
for (final c in _clients.toList()) {
|
||||
await c.close().catchError((_) {});
|
||||
}
|
||||
_clients.clear();
|
||||
if (!_clientCount.isClosed) _clientCount.add(0);
|
||||
await _server?.close(force: true);
|
||||
_server = null;
|
||||
}
|
||||
|
||||
Future<void> dispose() async {
|
||||
await stop();
|
||||
await _clientCount.close();
|
||||
}
|
||||
|
||||
// ---- Handshake + per-client loop --------------------------------------
|
||||
|
||||
void _handleSocket(WebSocket socket) {
|
||||
final nonce = RemoteAuth.newNonce();
|
||||
var authed = false;
|
||||
|
||||
// Drop a peer that connects but never authenticates.
|
||||
final authTimeout = Timer(const Duration(seconds: 10), () {
|
||||
if (!authed) socket.close().catchError((_) {});
|
||||
});
|
||||
|
||||
_send(
|
||||
socket,
|
||||
ChallengeMessage(
|
||||
version: kRemoteProtocolVersion,
|
||||
serverKey: serverKey,
|
||||
nonce: nonce,
|
||||
device: deviceName,
|
||||
),
|
||||
);
|
||||
|
||||
socket.listen(
|
||||
(data) async {
|
||||
RemoteMessage msg;
|
||||
try {
|
||||
msg = RemoteMessage.decode(data as String);
|
||||
} catch (_) {
|
||||
return; // ignore a malformed frame
|
||||
}
|
||||
|
||||
if (!authed) {
|
||||
final denyReason = _rejectReason(msg, nonce);
|
||||
if (denyReason != null) {
|
||||
_send(socket, DenyMessage(reason: denyReason));
|
||||
await socket.close().catchError((_) {});
|
||||
return;
|
||||
}
|
||||
authed = true;
|
||||
authTimeout.cancel();
|
||||
_send(socket, WelcomeMessage(device: deviceName));
|
||||
_clients.add(socket);
|
||||
if (!_clientCount.isClosed) _clientCount.add(_clients.length);
|
||||
// Immediate full snapshot so the remote paints without waiting for
|
||||
// the next state change.
|
||||
_send(socket, SnapshotMessage(_snapshotOf(_controller.currentState)));
|
||||
return;
|
||||
}
|
||||
|
||||
if (msg is CommandMessage) {
|
||||
try {
|
||||
await applyRemoteCommand(_controller, msg.command);
|
||||
} catch (_) {
|
||||
// A command that fails locally (e.g. offline host) surfaces via the
|
||||
// next snapshot's error field — no need to tear down the client.
|
||||
}
|
||||
}
|
||||
},
|
||||
onDone: () {
|
||||
authTimeout.cancel();
|
||||
_removeClient(socket);
|
||||
},
|
||||
onError: (_) {
|
||||
authTimeout.cancel();
|
||||
_removeClient(socket);
|
||||
},
|
||||
cancelOnError: true,
|
||||
);
|
||||
}
|
||||
|
||||
/// Null if the auth message is acceptable, else a human-readable deny reason.
|
||||
String? _rejectReason(RemoteMessage msg, String nonce) {
|
||||
if (msg is! AuthMessage) return 'expected auth';
|
||||
if (msg.version != kRemoteProtocolVersion) return 'protocol version';
|
||||
if (msg.serverKey != serverKey) return 'different server';
|
||||
if (!RemoteAuth.verify(
|
||||
password: password,
|
||||
nonce: nonce,
|
||||
proof: msg.proof,
|
||||
)) {
|
||||
return 'authentication failed';
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
void _removeClient(WebSocket socket) {
|
||||
if (_clients.remove(socket) && !_clientCount.isClosed) {
|
||||
_clientCount.add(_clients.length);
|
||||
}
|
||||
}
|
||||
|
||||
// ---- Broadcast --------------------------------------------------------
|
||||
|
||||
/// Send discrete changes (play/pause, track, queue, shuffle/loop, error)
|
||||
/// immediately; coalesce position-only ticks to at most one per second (the
|
||||
/// remote interpolates position between snapshots).
|
||||
void _onStateChanged(PlaybackState s) {
|
||||
if (_clients.isEmpty) {
|
||||
_sentState = s;
|
||||
return;
|
||||
}
|
||||
final prev = _sentState;
|
||||
if (prev == null || _discreteChanged(prev, s)) {
|
||||
_flush(s);
|
||||
} else {
|
||||
_pending = s;
|
||||
_throttle ??= Timer(const Duration(seconds: 1), () {
|
||||
_throttle = null;
|
||||
final p = _pending;
|
||||
_pending = null;
|
||||
if (p != null) _flush(p);
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
void _flush(PlaybackState s) {
|
||||
_throttle?.cancel();
|
||||
_throttle = null;
|
||||
_pending = null;
|
||||
_sentState = s;
|
||||
final frame = SnapshotMessage(_snapshotOf(s)).encode();
|
||||
for (final c in _clients.toList()) {
|
||||
try {
|
||||
c.add(frame);
|
||||
} catch (_) {
|
||||
_removeClient(c);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
static bool _discreteChanged(PlaybackState a, PlaybackState b) =>
|
||||
a.playing != b.playing ||
|
||||
a.currentIndex != b.currentIndex ||
|
||||
a.shuffle != b.shuffle ||
|
||||
a.loop != b.loop ||
|
||||
a.error != b.error ||
|
||||
// copyWith reuses the same queue list unless it actually changed, so an
|
||||
// identity check cheaply distinguishes a queue edit from a position tick.
|
||||
!identical(a.queue, b.queue);
|
||||
|
||||
StateSnapshot _snapshotOf(PlaybackState s) => StateSnapshot(
|
||||
queue: s.queue,
|
||||
currentIndex: s.currentIndex,
|
||||
playing: s.playing,
|
||||
positionMs: s.position.inMilliseconds,
|
||||
durationMs: s.effectiveDuration.inMilliseconds,
|
||||
shuffle: s.shuffle,
|
||||
loop: s.loop.name,
|
||||
seq: ++_seq,
|
||||
);
|
||||
|
||||
void _send(WebSocket socket, RemoteMessage msg) {
|
||||
try {
|
||||
socket.add(msg.encode());
|
||||
} catch (_) {
|
||||
// Socket already gone; the listen callbacks will clean it up.
|
||||
}
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue