-
Notifications
You must be signed in to change notification settings - Fork 0
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Merge pull request #18 from momentohq/subscribe
feat: topic subscribe
- Loading branch information
Showing
6 changed files
with
168 additions
and
16 deletions.
There are no files selected for viewing
This file was deleted.
Oops, something went wrong.
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,48 @@ | ||
import 'dart:io'; | ||
|
||
import 'package:client_sdk_dart/client_sdk_dart.dart'; | ||
import 'package:client_sdk_dart/src/auth/credential_provider.dart'; | ||
import 'package:client_sdk_dart/src/messages/responses/topics/topic_publish.dart'; | ||
import 'package:client_sdk_dart/src/messages/responses/topics/topic_subscribe.dart'; | ||
import 'package:client_sdk_dart/src/messages/responses/topics/topic_subscription_item.dart'; | ||
import 'package:client_sdk_dart/src/messages/values.dart'; | ||
import 'package:logging/logging.dart'; | ||
|
||
void main() async { | ||
Logger.root.level = Level.ALL; // defaults to Level.INFO | ||
Logger.root.onRecord.listen((record) { | ||
print('${record.level.name}: ${record.time}: ${record.message}'); | ||
}); | ||
|
||
var topicClient = TopicClient( | ||
CredentialProvider.fromEnvironmentVariable("MOMENTO_API_KEY")); | ||
|
||
var result = await topicClient.publish("cache", "topic", StringValue("hi")); | ||
switch (result) { | ||
case TopicPublishSuccess(): | ||
print("Successful publish!"); | ||
case TopicPublishError(): | ||
print("Publish error: ${result.errorCode} ${result.message}"); | ||
} | ||
|
||
var sub = await topicClient.subscribe("cache", "topic"); | ||
switch (sub) { | ||
case TopicSubscription(): | ||
print("Successful subscription!"); | ||
await for (final msg in sub.stream) { | ||
switch (msg) { | ||
case TopicSubscriptionItemBinary(): | ||
print("Binary value: ${msg.value}"); | ||
case TopicSubscriptionItemText(): | ||
print("String value: ${msg.value}"); | ||
case TopicSubscriptionItemError(): | ||
print("Error receiving message: ${msg.errorCode}"); | ||
} | ||
} | ||
case TopicSubscribeError(): | ||
print("Subscribe error: ${sub.errorCode} ${sub.message}"); | ||
} | ||
|
||
print("End of Momento topics example"); | ||
exit(0); | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,38 @@ | ||
import 'package:client_sdk_dart/generated/cachepubsub.pbgrpc.dart'; | ||
import 'package:client_sdk_dart/src/messages/responses/responses_base.dart'; | ||
import 'package:grpc/grpc.dart'; | ||
import 'package:logging/logging.dart'; | ||
|
||
import 'topic_subscription_item.dart'; | ||
|
||
sealed class TopicSubscribeResponse {} | ||
|
||
class TopicSubscription implements TopicSubscribeResponse { | ||
final ResponseStream<SubscriptionItem_> _stream; | ||
|
||
TopicSubscription(this._stream); | ||
|
||
Stream<TopicSubscriptionItemResponse> get stream => | ||
_stream.map(_processResult).where((item) => item != null).cast(); | ||
|
||
TopicSubscriptionItemResponse? _processResult(SubscriptionItem_ item) { | ||
final logger = Logger("TopicSubscribeResponse"); | ||
switch (item.whichKind()) { | ||
case SubscriptionItem__Kind.item: | ||
return createTopicItemResponse(item.item); | ||
case SubscriptionItem__Kind.heartbeat: | ||
logger.fine("topic client received a heartbeat"); | ||
case SubscriptionItem__Kind.discontinuity: | ||
logger.fine("topic client received a discontinuity"); | ||
default: | ||
logger.shout( | ||
"topic client received unknown subscription item: ${item.whichKind()}"); | ||
} | ||
return null; | ||
} | ||
} | ||
|
||
class TopicSubscribeError extends ErrorResponseBase | ||
implements TopicSubscribeResponse { | ||
TopicSubscribeError(super.exception); | ||
} |
34 changes: 34 additions & 0 deletions
34
lib/src/messages/responses/topics/topic_subscription_item.dart
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,34 @@ | ||
import 'package:client_sdk_dart/generated/cachepubsub.pb.dart'; | ||
import 'package:client_sdk_dart/src/errors/errors.dart'; | ||
import 'package:client_sdk_dart/src/messages/responses/responses_base.dart'; | ||
|
||
sealed class TopicSubscriptionItemResponse {} | ||
|
||
class TopicSubscriptionItemText implements TopicSubscriptionItemResponse { | ||
final String _value; | ||
TopicSubscriptionItemText(this._value); | ||
String get value => _value; | ||
} | ||
|
||
class TopicSubscriptionItemBinary implements TopicSubscriptionItemResponse { | ||
final List<int> _value; | ||
TopicSubscriptionItemBinary(this._value); | ||
List<int> get value => _value; | ||
} | ||
|
||
class TopicSubscriptionItemError extends ErrorResponseBase | ||
implements TopicSubscriptionItemResponse { | ||
TopicSubscriptionItemError(super.exception); | ||
} | ||
|
||
TopicSubscriptionItemResponse createTopicItemResponse(TopicItem_ item) { | ||
switch (item.value.whichKind()) { | ||
case TopicValue__Kind.text: | ||
return TopicSubscriptionItemText(item.value.text); | ||
case TopicValue__Kind.binary: | ||
return TopicSubscriptionItemBinary(item.value.binary); | ||
default: | ||
return TopicSubscriptionItemError(UnknownException( | ||
"unknown TopicItemResponse value: ${item.value}", null, null)); | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters