116 lines
2.9 KiB
Dart
116 lines
2.9 KiB
Dart
|
|
// The byte pipe to a smolmaild server: raw TCP (dart:io), framed elsewhere.
|
||
|
|
// Mobile is this app's only target platform, so dart:io is fine here.
|
||
|
|
|
||
|
|
import "dart:async";
|
||
|
|
import "dart:io";
|
||
|
|
import "dart:typed_data";
|
||
|
|
|
||
|
|
import "package:smol_mail/smol/errors.dart";
|
||
|
|
import "package:smol_mail/smol/proto.dart";
|
||
|
|
|
||
|
|
/// Buffers the socket's stream and serves exact-length reads, so protocol
|
||
|
|
/// code never sees a partial frame.
|
||
|
|
class TcpWire implements Wire {
|
||
|
|
final Socket _socket;
|
||
|
|
final _chunks = <Uint8List>[];
|
||
|
|
final _waiters = <_ReadRequest>[];
|
||
|
|
int _buffered = 0;
|
||
|
|
Object? _closed;
|
||
|
|
late final StreamSubscription<Uint8List> _subscription;
|
||
|
|
|
||
|
|
TcpWire(this._socket) {
|
||
|
|
_subscription = _socket.listen(_onData,
|
||
|
|
onError: (Object error) => _fail(error),
|
||
|
|
onDone: () => _fail(const SmolError("server closed the connection")));
|
||
|
|
}
|
||
|
|
|
||
|
|
static Future<TcpWire> connect(String host, int port) async {
|
||
|
|
try {
|
||
|
|
// Mobile networks routinely need longer than a LAN handshake; 30s keeps
|
||
|
|
// flaky handovers from surfacing as user-facing timeouts.
|
||
|
|
return TcpWire(await Socket.connect(host, port,
|
||
|
|
timeout: const Duration(seconds: 30)));
|
||
|
|
} on SocketException catch (error) {
|
||
|
|
throw SmolError("cannot reach $host:$port (${error.message})");
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
void _onData(Uint8List data) {
|
||
|
|
_chunks.add(data);
|
||
|
|
_buffered += data.length;
|
||
|
|
_wake();
|
||
|
|
}
|
||
|
|
|
||
|
|
void _wake() {
|
||
|
|
_waiters.removeWhere((w) {
|
||
|
|
if (_closed != null) {
|
||
|
|
w.completer.completeError(_closed!);
|
||
|
|
return true;
|
||
|
|
}
|
||
|
|
if (_buffered >= w.need) {
|
||
|
|
w.completer.complete();
|
||
|
|
return true;
|
||
|
|
}
|
||
|
|
return false;
|
||
|
|
});
|
||
|
|
}
|
||
|
|
|
||
|
|
void _fail(Object error) {
|
||
|
|
_closed = error;
|
||
|
|
for (final w in _waiters) {
|
||
|
|
w.completer.completeError(error);
|
||
|
|
}
|
||
|
|
_waiters.clear();
|
||
|
|
}
|
||
|
|
|
||
|
|
@override
|
||
|
|
void send(Uint8List bytes) {
|
||
|
|
if (_closed != null) throw _closed!;
|
||
|
|
_socket.add(bytes);
|
||
|
|
}
|
||
|
|
|
||
|
|
@override
|
||
|
|
void close() {
|
||
|
|
_subscription.cancel();
|
||
|
|
_socket.destroy();
|
||
|
|
_fail(const SmolError("connection closed"));
|
||
|
|
}
|
||
|
|
|
||
|
|
@override
|
||
|
|
Future<Uint8List> readExact(int n) async {
|
||
|
|
if (_closed != null) throw _closed!;
|
||
|
|
if (_buffered < n) {
|
||
|
|
final request = _ReadRequest(n);
|
||
|
|
_waiters.add(request);
|
||
|
|
try {
|
||
|
|
await request.completer.future;
|
||
|
|
} finally {
|
||
|
|
_waiters.remove(request);
|
||
|
|
}
|
||
|
|
if (_closed != null) throw _closed!;
|
||
|
|
}
|
||
|
|
final out = Uint8List(n);
|
||
|
|
var off = 0;
|
||
|
|
while (off < n) {
|
||
|
|
final chunk = _chunks.first;
|
||
|
|
final take = chunk.length < n - off ? chunk.length : n - off;
|
||
|
|
out.setRange(off, off + take, chunk);
|
||
|
|
if (take == chunk.length) {
|
||
|
|
_chunks.removeAt(0);
|
||
|
|
} else {
|
||
|
|
_chunks[0] = Uint8List.sublistView(chunk, take);
|
||
|
|
}
|
||
|
|
off += take;
|
||
|
|
_buffered -= take;
|
||
|
|
}
|
||
|
|
return out;
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
class _ReadRequest {
|
||
|
|
final int need;
|
||
|
|
final completer = Completer<void>();
|
||
|
|
|
||
|
|
_ReadRequest(this.need);
|
||
|
|
}
|