Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion example/hello-tls/receive.dart
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import "dart:io";
import "package:dart_amqp/dart_amqp.dart";
import "package:universal_io/io.dart";

void main(List<String> args) async {
var useClientCert = false;
Expand Down
3 changes: 1 addition & 2 deletions example/hello-tls/send.dart
Original file line number Diff line number Diff line change
@@ -1,6 +1,5 @@
import "dart:io";

import "package:dart_amqp/dart_amqp.dart";
import "package:universal_io/io.dart";

void main(List<String> args) async {
var useClientCert = false;
Expand Down
2 changes: 1 addition & 1 deletion example/hello/receive.dart
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import "dart:io";
import "package:dart_amqp/dart_amqp.dart";
import "package:universal_io/io.dart";

void main() async {
Client client = Client();
Expand Down
2 changes: 1 addition & 1 deletion example/pubsub/receive_logs.dart
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import "dart:io";
import "package:dart_amqp/dart_amqp.dart";
import "package:universal_io/io.dart";

void main() async {
Client client = Client();
Expand Down
2 changes: 1 addition & 1 deletion example/routing/emit_log_direct.dart
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import "dart:io";
import "package:dart_amqp/dart_amqp.dart";
import "package:universal_io/io.dart";

void main(List<String> args) async {
if (args.length < 2 || !["info", "warning", "error"].contains(args[0])) {
Expand Down
2 changes: 1 addition & 1 deletion example/routing/receive_logs_direct.dart
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import "dart:io";
import "package:dart_amqp/dart_amqp.dart";
import "package:universal_io/io.dart";

void main(List<String> args) async {
if (args.isEmpty || !args.every(["info", "warning", "error"].contains)) {
Expand Down
2 changes: 1 addition & 1 deletion example/rpc/rpc_client.dart
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import "dart:io";
import "dart:async";
import "package:dart_amqp/dart_amqp.dart";
import "package:universal_io/io.dart";

class FibonacciRpcClient {
int _nextCorrelationId = 1;
Expand Down
2 changes: 1 addition & 1 deletion example/rpc/rpc_server.dart
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import "dart:io";
import "package:dart_amqp/dart_amqp.dart";
import "package:universal_io/io.dart";

// Slow implementation of fib
int fib(int n) {
Expand Down
2 changes: 1 addition & 1 deletion example/topics/emit_log_topic.dart
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import "dart:io";
import "package:dart_amqp/dart_amqp.dart";
import "package:universal_io/io.dart";

void main(List<String> args) async {
if (args.length < 2) {
Expand Down
2 changes: 1 addition & 1 deletion example/topics/receive_logs_topic.dart
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import "dart:io";
import "package:dart_amqp/dart_amqp.dart";
import "package:universal_io/io.dart";

void main(List<String> args) async {
if (args.isEmpty) {
Expand Down
2 changes: 1 addition & 1 deletion example/workers/worker.dart
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import "dart:io";
import "package:dart_amqp/dart_amqp.dart";
import "package:universal_io/io.dart";

void main() async {
Client client = Client();
Expand Down
11 changes: 6 additions & 5 deletions lib/src/client.dart
Original file line number Diff line number Diff line change
@@ -1,20 +1,21 @@
library dart_amqp.client;

import "dart:async";
import "dart:io";
import "package:dart_amqp/dart_amqp.dart";
import "package:universal_io/io.dart";
import "dart:typed_data";
import "dart:collection";

import 'package:async/async.dart';
import "package:web_socket_channel/web_socket_channel.dart";

// Internal lib dependencies
import "logging.dart";
import "exceptions.dart";
import "enums.dart";
import "protocol.dart";
import "authentication.dart";

part "client/connection_settings.dart";
part "client/base_connection_settings.dart";
part "client/impl/connection_settings.dart";
part "client/impl/websocket_connection_settings.dart";

// client interfaces
part "client/client.dart";
Expand Down
43 changes: 43 additions & 0 deletions lib/src/client/base_connection_settings.dart
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
part of "../client.dart";

abstract class BaseConnectionSettings {
// The host to connect to
String? host;

// The port to connect to
int? port;

// The uri to connect to websocket. Use either host and url
Uri? uri;

// The connection vhost that will be sent to the server
abstract String virtualHost;

// The max number of reconnection attempts before declaring a connection as unusable
abstract int maxConnectionAttempts;

// The time to wait before trying to reconnect
abstract Duration reconnectWaitTime;

// Authentication provider
abstract Authenticator authProvider;

// Protocol version
int amqpProtocolVersion = 0;
int amqpMajorVersion = 0;
int amqpMinorVersion = 9;
int amqpRevision = 1;

// Tuning settings
abstract TuningSettings tuningSettings;

// TLS settings (if TLS connection is required)
SecurityContext? tlsContext;
bool Function(X509Certificate)? onBadCertificate;

// Connection identifier
String? connectionName;

// The time to wait for socket connection to be established
Duration? connectTimeout;
}
4 changes: 2 additions & 2 deletions lib/src/client/client.dart
Original file line number Diff line number Diff line change
@@ -1,11 +1,11 @@
part of "../client.dart";

abstract class Client {
factory Client({ConnectionSettings? settings}) =>
factory Client({BaseConnectionSettings? settings}) =>
_ClientImpl(settings: settings);

// Configuration options
ConnectionSettings get settings;
BaseConnectionSettings get settings;
TuningSettings get tuningSettings;

/// Check if a connection is currently in handshake state
Expand Down
39 changes: 24 additions & 15 deletions lib/src/client/impl/channel_impl.dart
Original file line number Diff line number Diff line change
Expand Up @@ -50,25 +50,24 @@ class _ChannelImpl implements Channel {
_pendingOperationPayloads.add(this);

// Transmit handshake
_frameWriter
..writeProtocolHeader(
_client.settings.amqpProtocolVersion,
_client.settings.amqpMajorVersion,
_client.settings.amqpMinorVersion,
_client.settings.amqpRevision)
..pipe(_client._socket!);
_frameWriter.writeProtocolHeader(
_client.settings.amqpProtocolVersion,
_client.settings.amqpMajorVersion,
_client.settings.amqpMinorVersion,
_client.settings.amqpRevision);
_setPipe();
}

void writeHeartbeat() {
if (_channelClosed != null || _client._socket == null) {
if (_channelClosed != null ||
(_client._socket == null && _client._webSocketChannel == null)) {
return; // no-op
}

// Transmit heartbeat
try {
_frameWriter
..writeHeartbeat()
..pipe(_client._socket!);
_frameWriter.writeHeartbeat();
_setPipe();
} catch (_) {
// An exception will be raised if we attempt to send a hearbeat
// immediately after the connection to the server is lost. We can safely
Expand Down Expand Up @@ -104,10 +103,10 @@ class _ChannelImpl implements Channel {
_nextPublishSeqNo++;
}

_frameWriter
..writeMessage(channelId, message,
properties: properties, payloadContent: payloadContent)
..pipe(_client._socket!);
_frameWriter.writeMessage(channelId, message,
properties: properties, payloadContent: payloadContent);

_setPipe();

// If the noWait flag was specified, complete the future now. The broken
// will raise any errors asynchronously via the channel or connection.
Expand All @@ -116,6 +115,16 @@ class _ChannelImpl implements Channel {
}
}

void _setPipe() {
if (_client._webSocketChannel != null) {
_frameWriter.pipe(_client._webSocketChannel!.sink);
} else if (_client._socket != null) {
_frameWriter.pipe(_client._socket!);
} else {
throw FatalException("Couldn't set pipe");
}
}

/// Implement the handshake flow specified by the AMQP spec by
/// examining [serverFrame] and generating the appropriate response
void _processHandshake(DecodedMessage serverMessage) {
Expand Down
Loading