-
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.
- Loading branch information
Showing
4 changed files
with
110 additions
and
11 deletions.
There are no files selected for viewing
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,36 @@ | ||
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); | ||
|
||
TopicSubscriptionItemResponse? _processResult(SubscriptionItem_ item) { | ||
final logger = Logger("TopicSubscribeResponse"); | ||
switch (item.runtimeType) { | ||
case TopicItem_: | ||
return createTopicItemResponse(item as TopicItem_); | ||
case Heartbeat_: | ||
logger.info("topic client received a heartbeat"); | ||
case Discontinuity_: | ||
logger.info("topic client received a discontinuity"); | ||
default: | ||
logger.shout("topic client received unknown subscription item: ", item.runtimeType); | ||
} | ||
return null; | ||
} | ||
} | ||
|
||
class TopicSubscribeError extends ErrorResponseBase | ||
implements TopicSubscribeResponse { | ||
TopicSubscribeError(super.exception); | ||
} |
33 changes: 33 additions & 0 deletions
33
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,33 @@ | ||
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