Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

Add basic WAMP client #4

Merged
merged 10 commits into from
May 2, 2024
Merged
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
3 changes: 3 additions & 0 deletions lib/exports.dart
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
export "src/client.dart" show Client;
export "src/session.dart" show Session;
export "src/types.dart" show Event, Invocation, Registration, Result, Subscription;
22 changes: 22 additions & 0 deletions lib/src/client.dart
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
import "package:wamp/src/session.dart";
import "package:wamp/src/types.dart";
import "package:wamp/src/wsjoiner.dart";
import "package:wampproto/auth.dart";
import "package:wampproto/serializers.dart";

class Client {
Client({IClientAuthenticator? authenticator, Serializer? serializer}) {
_authenticator = authenticator;
_serializer = serializer;
}

IClientAuthenticator? _authenticator;
Serializer? _serializer;

Future<Session> connect(String url, String realm) async {
WAMPSessionJoiner joiner = WAMPSessionJoiner(authenticator: _authenticator, serializer: _serializer);
BaseSession baseSession = await joiner.join(url, realm);

return Session(baseSession);
}
}
28 changes: 28 additions & 0 deletions lib/src/helpers.dart
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
import "package:wamp/src/wsjoiner.dart";
import "package:wampproto/messages.dart";
import "package:wampproto/serializers.dart";

String getSubProtocol(Serializer serializer) {
if (serializer is JSONSerializer) {
return WAMPSessionJoiner.jsonSubProtocol;
} else if (serializer is CBORSerializer) {
return WAMPSessionJoiner.cborSubProtocol;
} else if (serializer is MsgPackSerializer) {
return WAMPSessionJoiner.msgpackSubProtocol;
} else {
throw ArgumentError("invalid serializer");
}
}

String wampErrorString(Error err) {
String errStr = err.uri;
if (err.args.isNotEmpty) {
String args = err.args.map((arg) => arg.toString()).join(", ");
errStr += ": $args";
}
if (err.kwargs.isNotEmpty) {
String kwargs = err.kwargs.entries.map((entry) => "${entry.key}=${entry.value}").join(", ");
errStr += ": $kwargs";
}
return errStr;
}
192 changes: 192 additions & 0 deletions lib/src/session.dart
Original file line number Diff line number Diff line change
@@ -0,0 +1,192 @@
import "dart:async";
import "dart:typed_data";

import "package:wamp/src/helpers.dart";
import "package:wamp/src/types.dart";
import "package:wampproto/idgen.dart";
import "package:wampproto/messages.dart" as msg;
import "package:wampproto/session.dart";

class Session {
Session(this._baseSession) {
_wampSession = WAMPSession(serializer: _baseSession.serializer);
Future.microtask(() async {
while (true) {
var message = await _baseSession.receive();
var decodedMessage = Uint8List.fromList((message as String).codeUnits);
_processIncomingMessage(_wampSession.receive(decodedMessage));
}
});
}

final BaseSession _baseSession;
late WAMPSession _wampSession;

final SessionScopeIDGenerator _idGen = SessionScopeIDGenerator();

int get _nextID => _idGen.next();

Future<void> close() async {
await _baseSession.close();
}

final Map<int, Completer<Result>> _callRequests = {};
final Map<int, RegisterRequest> _registerRequests = {};
final Map<int, Result Function(Invocation)> _registrations = {};
final Map<int, UnregisterRequest> _unregisterRequests = {};
final Map<int, Completer<void>> _publishRequests = {};
final Map<int, SubscribeRequest> _subscribeRequests = {};
final Map<int, void Function(Event)> _subscriptions = {};
final Map<int, UnsubscribeRequest> _unsubscribeRequests = {};

void _processIncomingMessage(msg.Message message) {
if (message is msg.Result) {
var request = _callRequests.remove(message.requestID);
if (request != null) {
request.complete(Result(args: message.args, kwargs: message.kwargs, details: message.details));
}
} else if (message is msg.Registered) {
var request = _registerRequests.remove(message.requestID);
if (request != null) {
_registrations[message.registrationID] = request.endpoint;
request.future.complete(Registration(message.registrationID));
}
} else if (message is msg.Invocation) {
var endpoint = _registrations[message.registrationID];
if (endpoint != null) {
Result result = endpoint(Invocation(args: message.args, kwargs: message.kwargs, details: message.details));
Uint8List data = _wampSession.sendMessage(
msg.Yield(message.requestID, args: result.args, kwargs: result.kwargs, options: result.details),
);
_baseSession.send(data);
}
} else if (message is msg.UnRegistered) {
var request = _unregisterRequests.remove(message.requestID);
if (request != null) {
_registrations.remove(request.registrationID);
request.future.complete();
}
} else if (message is msg.Published) {
var request = _publishRequests.remove(message.requestID);
if (request != null) {
request.complete();
}
} else if (message is msg.Subscribed) {
var request = _subscribeRequests.remove(message.requestID);
if (request != null) {
_subscriptions[message.subscriptionID] = request.endpoint;
request.future.complete(Subscription(message.subscriptionID));
}
} else if (message is msg.Event) {
var endpoint = _subscriptions[message.subscriptionID];
if (endpoint != null) {
endpoint(Event(args: message.args, kwargs: message.kwargs, details: message.details));
}
} else if (message is msg.UnSubscribed) {
var request = _unsubscribeRequests.remove(message.requestID);
if (request != null) {
_subscriptions.remove(request.subscriptionId);
request.future.complete();
}
} else if (message is msg.Error) {
switch (message.msgType) {
case msg.Call.id:
_callRequests.remove(message.requestID);

case msg.Register.id:
_registerRequests.remove(message.requestID);

case msg.UnRegister.id:
_unregisterRequests.remove(message.requestID);

case msg.Subscribe.id:
_subscribeRequests.remove(message.requestID);

case msg.UnSubscribe.id:
_unsubscribeRequests.remove(message.requestID);

case msg.Publish.id:
_publishRequests.remove(message.requestID);
}

throw Exception(wampErrorString(message));
}
}

Future<Result> call(
String procedure, {
List<dynamic>? args,
Map<String, dynamic>? kwargs,
Map<String, dynamic>? options,
}) {
var call = msg.Call(_nextID, procedure, args: args, kwargs: kwargs, options: options);

var completer = Completer<Result>();
_callRequests[call.requestID] = completer;

_baseSession.send(_wampSession.sendMessage(call));

return completer.future;
}

Future<Registration> register(String procedure, Result Function(Invocation) endpoint) {
var register = msg.Register(_nextID, procedure);

var completer = Completer<Registration>();
_registerRequests[register.requestID] = RegisterRequest(completer, endpoint);

_baseSession.send(_wampSession.sendMessage(register));

return completer.future;
}

Future<void> unregister(Registration reg) {
var unregister = msg.UnRegister(_nextID, reg.registrationID);

var completer = Completer();
_unregisterRequests[unregister.requestID] = UnregisterRequest(completer, reg.registrationID);

_baseSession.send(_wampSession.sendMessage(unregister));

return completer.future;
}

Future<void>? publish(
String topic, {
List<dynamic>? args,
Map<String, dynamic>? kwargs,
Map<String, dynamic>? options,
}) {
var publish = msg.Publish(_nextID, topic, args: args, kwargs: kwargs, options: options);

var completer = Completer<void>();
_publishRequests[publish.requestID] = completer;
_baseSession.send(_wampSession.sendMessage(publish));

if (options != null && options["acknowledge"]) {
return completer.future;
}

return null;
}

Future<Subscription> subscribe(String topic, void Function(Event) endpoint) {
var subscribe = msg.Subscribe(_nextID, topic);

var completer = Completer<Subscription>();
_subscribeRequests[subscribe.requestID] = SubscribeRequest(completer, endpoint);
_baseSession.send(_wampSession.sendMessage(subscribe));

return completer.future;
}

Future<void> unsubscribe(Subscription sub) {
var unsubscribe = msg.UnSubscribe(_nextID, sub.subscriptionId);

var completer = Completer<void>();
_unsubscribeRequests[unsubscribe.requestID] = UnsubscribeRequest(completer, sub.subscriptionId);
_baseSession.send(_wampSession.sendMessage(unsubscribe));

return completer.future;
}
}
108 changes: 108 additions & 0 deletions lib/src/types.dart
Original file line number Diff line number Diff line change
@@ -0,0 +1,108 @@
import "dart:async";
import "dart:io";

import "package:wampproto/serializers.dart";
import "package:wampproto/session.dart";

class BaseSession {
BaseSession(this._ws, this._wsStreamController, this.sessionDetails, this.serializer);

final WebSocket _ws;
final StreamController _wsStreamController;
SessionDetails sessionDetails;
Serializer serializer;

void send(Object data) {
_ws.add(data);
}

Future<Object> receive() async {
return await _wsStreamController.stream.first;
}

Future<void> close() async {
await _ws.close();
}
}

class Result {
Result({
List<dynamic>? args,
Map<String, dynamic>? kwargs,
Map<String, dynamic>? details,
}) : args = args ?? [],
kwargs = kwargs ?? {},
details = details ?? {};

final List<dynamic> args;
final Map<String, dynamic> kwargs;
final Map<String, dynamic> details;
}

class Registration {
Registration(this.registrationID);

final int registrationID;
}

class RegisterRequest {
RegisterRequest(this.future, this.endpoint);

final Completer<Registration> future;
final Result Function(Invocation) endpoint;
}

class Invocation {
Invocation({
List<dynamic>? args,
Map<String, dynamic>? kwargs,
Map<String, dynamic>? details,
}) : args = args ?? [],
kwargs = kwargs ?? {},
details = details ?? {};

final List<dynamic> args;
final Map<String, dynamic> kwargs;
final Map<String, dynamic> details;
}

class UnregisterRequest {
UnregisterRequest(this.future, this.registrationID);

final Completer<void> future;
final int registrationID;
}

class Subscription {
Subscription(this.subscriptionId);

final int subscriptionId;
}

class SubscribeRequest {
SubscribeRequest(this.future, this.endpoint);

final Completer<Subscription> future;
final void Function(Event) endpoint;
}

class Event {
Event({
List<dynamic>? args,
Map<String, dynamic>? kwargs,
Map<String, dynamic>? details,
}) : args = args ?? [],
kwargs = kwargs ?? {},
details = details ?? {};

final List<dynamic> args;
final Map<String, dynamic> kwargs;
final Map<String, dynamic> details;
}

class UnsubscribeRequest {
UnsubscribeRequest(this.future, this.subscriptionId);

final Completer<void> future;
final int subscriptionId;
}
Loading