diff --git a/packages/dart_client/lib/src/db/stream_chat_database.dart b/packages/dart_client/lib/src/db/stream_chat_database.dart new file mode 100644 index 00000000..87ab828e --- /dev/null +++ b/packages/dart_client/lib/src/db/stream_chat_database.dart @@ -0,0 +1,97 @@ +import 'package:stream_chat/src/api/requests.dart'; +import 'package:stream_chat/src/models/channel_model.dart'; +import 'package:stream_chat/src/models/channel_state.dart'; +import 'package:stream_chat/src/models/event.dart'; +import 'package:stream_chat/src/models/member.dart'; +import 'package:stream_chat/src/models/message.dart'; +import 'package:stream_chat/src/models/reaction.dart'; +import 'package:stream_chat/src/models/read.dart'; +import 'package:stream_chat/src/models/user.dart'; + +/// +abstract class StreamChatDatabase { + /// Creates a new connection to the database + Future connect({ + bool connectBackground = false, + bool logStatements = false, + }); + + /// Closes the database instance + /// If [flush] is true, the database data will be deleted + Future disconnect({bool flush = false}); + + /// Get stored replies by messageId + Future> getReplies( + String parentId, { + String lessThan, + }); + + /// Get stored connection event + Future getConnectionInfo(); + + /// Get stored lastSyncAt + Future getLastSyncAt(); + + /// Update stored connection event + Future updateConnectionInfo(Event event); + + /// Update stored lastSyncAt + Future updateLastSyncAt(DateTime lastSyncAt); + + /// Get the channel cids saved in the offline storage + Future> getChannelCids(); + + Future getChannelByCid(String cid); + + Future> getMembersByCid(String cid); + + Future> getReadsByCid(String cid); + + Future> getMessagesByCid( + String cid, { + int limit = 20, + String messageLessThan, + String messageGreaterThan, + }); + + /// Get list of channels by filter, sort and paginationParams + Future> getChannelStates({ + Map filter, + List sort = const [], + PaginationParams paginationParams, + }); + + /// Update list of channel queries + /// If [clearQueryCache] is true before the insert + /// the list of matching rows will be deleted + Future updateChannelQueries( + Map filter, + List cids, + bool clearQueryCache, + ); + + /// Remove a message by message id + Future deleteMessageByIds(List messageIds); + + /// Remove a message by message id + Future deleteMessageByCids(List cids); + + /// Remove a channel by cid + Future deleteChannelByCids(List cids); + + /// Update messages data from a list + Future updateMessages(String cid, List messages); + + /// Get the info about channel threads + Future>> getChannelThreads(String cid); + + Future updateChannels(List channels); + + Future updateMembers(String cid, List members); + + Future updateReads(String cid, List reads); + + Future updateUsers(List users); + + Future updateReactions(List reactions); +} diff --git a/packages/dart_client/lib/src/utils/result.dart b/packages/dart_client/lib/src/utils/result.dart new file mode 100644 index 00000000..a2961ec4 --- /dev/null +++ b/packages/dart_client/lib/src/utils/result.dart @@ -0,0 +1,45 @@ +import 'package:flutter/foundation.dart'; +import 'package:freezed_annotation/freezed_annotation.dart'; + +import '../exceptions.dart'; + +part 'result.freezed.dart'; + +/// +@freezed +abstract class Result with _$Result { + /// + const factory Result.success({@required T data}) = Success; + + /// + const factory Result.error({@required ApiError error}) = Error; +} + +/// +extension ResultX on Result { + /// + bool get isSuccess => this is Success; + + /// + bool get isError => this is Error; + + /// + T get data { + if (isError) { + throw Exception( + 'Result is not successful. Check result.isSuccess before reading the data.', + ); + } + return (this as Success).data; + } + + /// + ApiError get error { + if (isSuccess) { + throw Exception( + 'Result is successful, not an error. Check result.isError before reading the error.', + ); + } + return (this as Error).error; + } +} diff --git a/packages/dart_client/lib/stream_chat.dart b/packages/dart_client/lib/stream_chat.dart index 101e9440..8bac944f 100644 --- a/packages/dart_client/lib/stream_chat.dart +++ b/packages/dart_client/lib/stream_chat.dart @@ -2,7 +2,7 @@ library stream_chat; export 'package:dio/src/dio_error.dart'; export 'package:dio/src/multipart_file.dart'; -export 'package:logging/src/level.dart'; +export 'package:logging/logging.dart' show Logger, Level; export './src/api/channel.dart'; export './src/api/connection_status.dart'; @@ -27,3 +27,5 @@ export './src/models/reaction.dart'; export './src/models/read.dart'; export './src/models/user.dart'; export './src/notifications.dart'; +export './src/db/stream_chat_database.dart'; +export './src/utils/result.dart'; diff --git a/packages/dart_client/pubspec.yaml b/packages/dart_client/pubspec.yaml index 9fd0d9e0..d4734e44 100644 --- a/packages/dart_client/pubspec.yaml +++ b/packages/dart_client/pubspec.yaml @@ -26,6 +26,9 @@ dependencies: collection: ^1.14.12 sqlite3_flutter_libs: ^0.3.0 pedantic: ^1.9.2 + freezed: ^0.12.7 + stream_chat_persistence: + path: ../stream_chat_persistence dev_dependencies: build_runner: ^1.10.0 diff --git a/packages/stream_chat_persistence/.gitignore b/packages/stream_chat_persistence/.gitignore new file mode 100644 index 00000000..1985397a --- /dev/null +++ b/packages/stream_chat_persistence/.gitignore @@ -0,0 +1,74 @@ +# Miscellaneous +*.class +*.log +*.pyc +*.swp +.DS_Store +.atom/ +.buildlog/ +.history +.svn/ + +# IntelliJ related +*.iml +*.ipr +*.iws +.idea/ + +# The .vscode folder contains launch configuration and tasks you configure in +# VS Code which you may wish to be included in version control, so this line +# is commented out by default. +#.vscode/ + +# Flutter/Dart/Pub related +**/doc/api/ +.dart_tool/ +.flutter-plugins +.flutter-plugins-dependencies +.packages +.pub-cache/ +.pub/ +build/ + +# Android related +**/android/**/gradle-wrapper.jar +**/android/.gradle +**/android/captures/ +**/android/gradlew +**/android/gradlew.bat +**/android/local.properties +**/android/**/GeneratedPluginRegistrant.java + +# iOS/XCode related +**/ios/**/*.mode1v3 +**/ios/**/*.mode2v3 +**/ios/**/*.moved-aside +**/ios/**/*.pbxuser +**/ios/**/*.perspectivev3 +**/ios/**/*sync/ +**/ios/**/.sconsign.dblite +**/ios/**/.tags* +**/ios/**/.vagrant/ +**/ios/**/DerivedData/ +**/ios/**/Icon? +**/ios/**/Pods/ +**/ios/**/.symlinks/ +**/ios/**/profile +**/ios/**/xcuserdata +**/ios/.generated/ +**/ios/Flutter/App.framework +**/ios/Flutter/Flutter.framework +**/ios/Flutter/Flutter.podspec +**/ios/Flutter/Generated.xcconfig +**/ios/Flutter/app.flx +**/ios/Flutter/app.zip +**/ios/Flutter/flutter_assets/ +**/ios/Flutter/flutter_export_environment.sh +**/ios/ServiceDefinitions.json +**/ios/Runner/GeneratedPluginRegistrant.* + +# Exceptions to above rules. +!**/ios/**/default.mode1v3 +!**/ios/**/default.mode2v3 +!**/ios/**/default.pbxuser +!**/ios/**/default.perspectivev3 diff --git a/packages/stream_chat_persistence/.metadata b/packages/stream_chat_persistence/.metadata new file mode 100644 index 00000000..5eb50347 --- /dev/null +++ b/packages/stream_chat_persistence/.metadata @@ -0,0 +1,10 @@ +# This file tracks properties of this Flutter project. +# Used by Flutter tool to assess capabilities and perform upgrades etc. +# +# This file should be version controlled and should not be manually edited. + +version: + revision: 78910062997c3a836feee883712c241a5fd22983 + channel: stable + +project_type: package diff --git a/packages/stream_chat_persistence/CHANGELOG.md b/packages/stream_chat_persistence/CHANGELOG.md new file mode 100644 index 00000000..ac071598 --- /dev/null +++ b/packages/stream_chat_persistence/CHANGELOG.md @@ -0,0 +1,3 @@ +## [0.0.1] - TODO: Add release date. + +* TODO: Describe initial release. diff --git a/packages/stream_chat_persistence/LICENSE b/packages/stream_chat_persistence/LICENSE new file mode 100644 index 00000000..ba75c69f --- /dev/null +++ b/packages/stream_chat_persistence/LICENSE @@ -0,0 +1 @@ +TODO: Add your license here. diff --git a/packages/stream_chat_persistence/README.md b/packages/stream_chat_persistence/README.md new file mode 100644 index 00000000..3e3f0c7c --- /dev/null +++ b/packages/stream_chat_persistence/README.md @@ -0,0 +1,14 @@ +# stream_chat_persistence + +A new Flutter package. + +## Getting Started + +This project is a starting point for a Dart +[package](https://flutter.dev/developing-packages/), +a library module containing code that can be shared easily across +multiple Flutter or Dart projects. + +For help getting started with Flutter, view our +[online documentation](https://flutter.dev/docs), which offers tutorials, +samples, guidance on mobile development, and a full API reference. diff --git a/packages/stream_chat_persistence/build.yaml b/packages/stream_chat_persistence/build.yaml new file mode 100644 index 00000000..b07425b6 --- /dev/null +++ b/packages/stream_chat_persistence/build.yaml @@ -0,0 +1,6 @@ +targets: + $default: + builders: + moor_generator: + options: + generate_connect_constructor: true diff --git a/packages/stream_chat_persistence/lib/src/converter/list_converter.dart b/packages/stream_chat_persistence/lib/src/converter/list_converter.dart new file mode 100644 index 00000000..5d0a628f --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/converter/list_converter.dart @@ -0,0 +1,22 @@ +import 'dart:convert'; + +import 'package:moor/moor.dart'; + +/// +class ListConverter extends TypeConverter, String> { + @override + List mapToDart(fromDb) { + if (fromDb == null) { + return null; + } + return List.from(jsonDecode(fromDb) ?? []); + } + + @override + String mapToSql(value) { + if (value == null) { + return null; + } + return jsonEncode(value); + } +} diff --git a/packages/stream_chat_persistence/lib/src/converter/map_converter.dart b/packages/stream_chat_persistence/lib/src/converter/map_converter.dart new file mode 100644 index 00000000..9fde0003 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/converter/map_converter.dart @@ -0,0 +1,22 @@ +import 'dart:convert'; + +import 'package:moor/moor.dart'; + +/// +class MapConverter extends TypeConverter, String> { + @override + Map mapToDart(fromDb) { + if (fromDb == null) { + return null; + } + return Map.from(jsonDecode(fromDb) ?? {}); + } + + @override + String mapToSql(value) { + if (value == null) { + return null; + } + return jsonEncode(value); + } +} diff --git a/packages/stream_chat_persistence/lib/src/converter/message_sending_status_converter.dart b/packages/stream_chat_persistence/lib/src/converter/message_sending_status_converter.dart new file mode 100644 index 00000000..b8d92b82 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/converter/message_sending_status_converter.dart @@ -0,0 +1,50 @@ +import 'package:moor/moor.dart'; +import 'package:stream_chat/stream_chat.dart'; + +/// +class MessageSendingStatusConverter + extends TypeConverter { + @override + MessageSendingStatus mapToDart(int fromDb) { + switch (fromDb) { + case 0: + return MessageSendingStatus.sending; + case 1: + return MessageSendingStatus.sent; + case 2: + return MessageSendingStatus.failed; + case 3: + return MessageSendingStatus.updating; + case 4: + return MessageSendingStatus.failed_update; + case 5: + return MessageSendingStatus.deleting; + case 6: + return MessageSendingStatus.failed_delete; + default: + return null; + } + } + + @override + int mapToSql(MessageSendingStatus value) { + switch (value) { + case MessageSendingStatus.sending: + return 0; + case MessageSendingStatus.sent: + return 1; + case MessageSendingStatus.failed: + return 2; + case MessageSendingStatus.updating: + return 3; + case MessageSendingStatus.failed_update: + return 4; + case MessageSendingStatus.deleting: + return 5; + case MessageSendingStatus.failed_delete: + return 6; + default: + return null; + } + } +} diff --git a/packages/stream_chat_persistence/lib/src/dao/channel_dao.dart b/packages/stream_chat_persistence/lib/src/dao/channel_dao.dart new file mode 100644 index 00000000..fa902f52 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/dao/channel_dao.dart @@ -0,0 +1,57 @@ +import 'package:moor/moor.dart'; +import 'package:stream_chat/stream_chat.dart'; +import 'package:stream_chat_persistence/src/db/moor_chat_database.dart'; +import 'package:stream_chat_persistence/src/entity/channels.dart'; +import 'package:stream_chat_persistence/src/entity/users.dart'; +import '../mapper/mapper.dart'; + +part 'channel_dao.g.dart'; + +/// +@UseDao(tables: [Channels, Users]) +class ChannelDao extends DatabaseAccessor + with _$ChannelDaoMixin { + /// + ChannelDao(MoorChatDatabase db) : super(db); + + /// Get channel by cid + Future getChannelByCid(String cid) async { + return (select(channels)..where((c) => c.cid.equals(cid))).join([ + leftOuterJoin(users, channels.createdById.equalsExp(users.id)), + ]).map((rows) { + final channel = rows.readTable(channels); + final createdBy = rows.readTable(users); + return channel.toChannelModel(createdBy: createdBy?.toUser()); + }).getSingle(); + } + + /// Delete all channels by matching cid in [cids] + /// + /// This will automatically delete the following linked records + /// 1. Channel Reads + /// 2. Channel Members + /// 3. Channel Messages -> Messages Reactions + Future deleteChannelByCids(List cids) async { + return (delete(channels)..where((tbl) => tbl.cid.isIn(cids))).go(); + } + + /// Get the channel cids saved in the storage + Future> get cids { + return (select(channels) + ..orderBy([(c) => OrderingTerm.desc(c.lastMessageAt)]) + ..limit(250)) + .map((c) => c.cid) + .get(); + } + + /// + Future updateChannels(List channelList) { + return batch( + (it) => it.insertAll( + channels, + channelList.map((c) => c.toEntity()).toList(), + mode: InsertMode.insertOrReplace, + ), + ); + } +} diff --git a/packages/stream_chat_persistence/lib/src/dao/channel_query_dao.dart b/packages/stream_chat_persistence/lib/src/dao/channel_query_dao.dart new file mode 100644 index 00000000..a42dd01e --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/dao/channel_query_dao.dart @@ -0,0 +1,126 @@ +import 'dart:convert'; + +import 'package:moor/moor.dart'; +import 'package:stream_chat/stream_chat.dart'; +import 'package:stream_chat_persistence/src/db/moor_chat_database.dart'; +import 'package:stream_chat_persistence/src/entity/channel_queries.dart'; +import 'package:stream_chat_persistence/src/entity/channels.dart'; +import 'package:stream_chat_persistence/src/entity/users.dart'; +import '../mapper/mapper.dart'; + +part 'channel_query_dao.g.dart'; + +/// +@UseDao(tables: [ChannelQueries, Channels, Users]) +class ChannelQueryDao extends DatabaseAccessor + with _$ChannelQueryDaoMixin { + /// + ChannelQueryDao(this._db) : super(_db); + + final MoorChatDatabase _db; + + String _computeHash(Map filter) { + if (filter == null) { + return 'allchannels'; + } + final hash = base64Encode(utf8.encode('filter: ${jsonEncode(filter)}')); + return hash; + } + + /// Update list of channel queries + /// If [clearQueryCache] is true before the insert + /// the list of matching rows will be deleted + Future updateChannelQueries( + Map filter, + List cids, + bool clearQueryCache, + ) async { + final hash = _computeHash(filter); + if (clearQueryCache) { + await (delete(channelQueries) + ..where((query) => query.queryHash.equals(hash))) + .go(); + } + + return batch((batch) { + batch.insertAll( + channelQueries, + cids.map((cid) { + return ChannelQueryEntity( + queryHash: hash, + channelCid: cid, + ); + }).toList(), + mode: InsertMode.insertOrReplace, + ); + }); + } + + /// Get list of channels by filter, sort and paginationParams + Future> getChannelStates({ + Map filter, + List sort = const [], + PaginationParams paginationParams, + }) async { + final hash = _computeHash(filter); + final cachedChannels = await Future.wait(await (select(channelQueries) + ..where((c) => c.queryHash.equals(hash))) + .get() + .then((channelQueries) { + final cids = channelQueries.map((c) => c.channelCid).toList(); + final query = select(channels)..where((c) => c.cid.isIn(cids)); + + sort = sort + ?.where((s) => ChannelModel.topLevelFields.contains(s.field)) + ?.toList(); + + if (sort != null && sort.isNotEmpty) { + query.orderBy(sort.map((s) { + final orderExpression = CustomExpression('channels.${s.field}'); + return (c) => OrderingTerm( + expression: orderExpression, + mode: s.direction == 1 ? OrderingMode.asc : OrderingMode.desc, + ); + }).toList()); + } + + if (paginationParams != null) { + query.limit( + paginationParams.limit ?? 10, + offset: paginationParams.offset, + ); + } + + return query.join([ + leftOuterJoin(users, channels.createdById.equalsExp(users.id)), + ]).map((row) async { + final userEntity = row.readTable(users); + final channelEntity = row.readTable(channels); + + final cid = channelEntity.cid; + final members = await _db.memberDao.getMembersByCid(cid); + final reads = await _db.readDao.getReadsByCid(cid); + final messages = await _db.messageDao.getMessagesByCid(cid); + + return channelEntity.toChannelState( + createdBy: userEntity?.toUser(), + members: members, + reads: reads, + messages: messages, + ); + }).get(); + })); + + if (sort?.isEmpty != false && cachedChannels?.isNotEmpty == true) { + cachedChannels + .sort((a, b) => b.channel.updatedAt.compareTo(a.channel.updatedAt)); + cachedChannels.sort((a, b) { + final dateA = a.channel.lastMessageAt ?? a.channel.createdAt; + final dateB = b.channel.lastMessageAt ?? b.channel.createdAt; + return dateB.compareTo(dateA); + }); + } + + return cachedChannels; + } +} diff --git a/packages/stream_chat_persistence/lib/src/dao/connection_event_dao.dart b/packages/stream_chat_persistence/lib/src/dao/connection_event_dao.dart new file mode 100644 index 00000000..59d22010 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/dao/connection_event_dao.dart @@ -0,0 +1,58 @@ +import 'package:moor/moor.dart'; +import 'package:stream_chat/stream_chat.dart'; +import 'package:stream_chat_persistence/src/db/moor_chat_database.dart'; +import 'package:stream_chat_persistence/src/entity/connection_events.dart'; +import 'package:stream_chat_persistence/src/entity/users.dart'; +import '../mapper/mapper.dart'; + +part 'connection_event_dao.g.dart'; + +/// +@UseDao(tables: [ConnectionEvents, Users]) +class ConnectionEventDao extends DatabaseAccessor + with _$ConnectionEventDaoMixin { + /// + ConnectionEventDao(MoorChatDatabase db) : super(db); + + /// Get the latest stored connection event + Future get connectionEvent { + return select(connectionEvents).join([ + leftOuterJoin(users, connectionEvents.ownUserId.equalsExp(users.id)), + ]).map((rows) { + final event = rows.readTable(connectionEvents); + final user = rows.readTable(users); + return event.toEvent(user: user?.toUser()); + }).getSingle(); + } + + /// Get the latest stored lastSyncAt + Future get lastSyncAt { + return select(connectionEvents).getSingle().then((r) => r?.lastSyncAt); + } + + /// Update stored connection event with latest data + Future updateConnectionEvent(Event event) async { + final connectionInfo = await select(connectionEvents).getSingle(); + return into(connectionEvents).insert( + ConnectionEventEntity( + id: 1, + lastSyncAt: connectionInfo?.lastSyncAt, + lastEventAt: event.createdAt ?? connectionInfo?.lastEventAt, + totalUnreadCount: + event.totalUnreadCount ?? connectionInfo?.totalUnreadCount, + ownUserId: event.me?.id ?? connectionInfo?.ownUserId, + unreadChannels: event.unreadChannels ?? connectionInfo?.unreadChannels, + ), + mode: InsertMode.insertOrReplace, + ); + } + + /// Update stored lastSyncAt with latest data + Future updateLastSyncAt(DateTime lastSyncAt) async { + return (update(connectionEvents)..where((tbl) => tbl.id.equals(1))).write( + ConnectionEventsCompanion( + lastSyncAt: Value(lastSyncAt), + ), + ); + } +} diff --git a/packages/stream_chat_persistence/lib/src/dao/dao.dart b/packages/stream_chat_persistence/lib/src/dao/dao.dart new file mode 100644 index 00000000..6f1e8221 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/dao/dao.dart @@ -0,0 +1,8 @@ +export 'user_dao.dart'; +export 'channel_dao.dart'; +export 'message_dao.dart'; +export 'member_dao.dart'; +export 'connection_event_dao.dart'; +export 'reaction_dao.dart'; +export 'read_dao.dart'; +export 'channel_query_dao.dart'; diff --git a/packages/stream_chat_persistence/lib/src/dao/member_dao.dart b/packages/stream_chat_persistence/lib/src/dao/member_dao.dart new file mode 100644 index 00000000..f21ab334 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/dao/member_dao.dart @@ -0,0 +1,43 @@ +import 'package:moor/moor.dart'; +import 'package:stream_chat/stream_chat.dart'; +import 'package:stream_chat_persistence/src/db/moor_chat_database.dart'; + +import 'package:stream_chat_persistence/src/entity/members.dart'; +import 'package:stream_chat_persistence/src/entity/users.dart'; + +import '../mapper/mapper.dart'; + +part 'member_dao.g.dart'; + +/// +@UseDao(tables: [Members, Users]) +class MemberDao extends DatabaseAccessor + with _$MemberDaoMixin { + /// + MemberDao(MoorChatDatabase db) : super(db); + + /// Get all members where [members.channelCid] matches [cid] + Future> getMembersByCid(String cid) async { + return (select(members).join([ + leftOuterJoin(users, members.userId.equalsExp(users.id)), + ]) + ..where(members.channelCid.equals(cid)) + ..orderBy([OrderingTerm.asc(members.createdAt)])) + .map((row) { + final userEntity = row.readTable(users); + final memberEntity = row.readTable(members); + return memberEntity.toMember(user: userEntity?.toUser()); + }).get(); + } + + /// + Future updateMembers(String cid, List memberList) async { + return batch( + (it) => it.insertAll( + members, + memberList.map((m) => m.toEntity(cid: cid)).toList(), + mode: InsertMode.insertOrReplace, + ), + ); + } +} diff --git a/packages/stream_chat_persistence/lib/src/dao/message_dao.dart b/packages/stream_chat_persistence/lib/src/dao/message_dao.dart new file mode 100644 index 00000000..700f8bed --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/dao/message_dao.dart @@ -0,0 +1,143 @@ +import 'package:moor/moor.dart'; +import 'package:stream_chat/stream_chat.dart'; +import 'package:stream_chat_persistence/src/db/moor_chat_database.dart'; +import 'package:stream_chat_persistence/src/entity/messages.dart'; +import 'package:stream_chat_persistence/src/entity/users.dart'; +import '../mapper/mapper.dart'; + +part 'message_dao.g.dart'; + +/// +@UseDao(tables: [Messages, Users]) +class MessageDao extends DatabaseAccessor + with _$MessageDaoMixin { + /// + MessageDao(this._db) : super(_db); + + final MoorChatDatabase _db; + + /// Removes all the messages by matching [messages.id] in [messageIds] + /// + /// This will automatically delete the following linked records + /// 1. Message Reactions + Future deleteMessageByIds(List messageIds) { + return (delete(messages)..where((tbl) => tbl.id.isIn(messageIds))).go(); + } + + /// Removes all the messages by matching [messages.channelCid] in [cids] + /// + /// This will automatically delete the following linked records + /// 1. Message Reactions + Future deleteMessageByCids(List cids) async { + return (delete(messages)..where((tbl) => tbl.channelCid.isIn(cids))).go(); + } + + Future _messageFromJoinRow(TypedResult rows) async { + final userEntity = rows.readTable(users); + final msgEntity = rows.readTable(messages); + final latestReactions = await _db.reactionDao.getReactions(msgEntity.id); + final ownReactions = await _db.reactionDao.getOwnReactions( + msgEntity.id, + userEntity.id, + ); + Message quotedMessage; + if (msgEntity.quotedMessageId != null) { + quotedMessage = await getMessageById(msgEntity.quotedMessageId); + } + return msgEntity.toMessage( + user: userEntity?.toUser(), + latestReactions: latestReactions, + ownReactions: ownReactions, + quotedMessage: quotedMessage, + ); + } + + /// + Future getMessageById(String id) async { + return await (select(messages).join([ + leftOuterJoin(users, messages.userId.equalsExp(users.id)), + ]) + ..where(messages.id.equals(id))) + .map(_messageFromJoinRow) + .getSingle(); + } + + /// + Future> getThreadMessages(String cid) async { + return Future.wait(await (select(messages).join([ + leftOuterJoin(users, messages.userId.equalsExp(users.id)), + ]) + ..where(messages.channelCid.equals(cid)) + ..where(isNotNull(messages.parentId)) + ..orderBy([OrderingTerm.asc(messages.createdAt)])) + .map(_messageFromJoinRow) + .get()); + } + + /// + Future> getThreadMessagesByParentId( + String parentId, { + String lessThan, + }) async { + final msgList = await Future.wait(await (select(messages).join([ + innerJoin(users, messages.userId.equalsExp(users.id)), + ]) + ..where(messages.parentId.equals(parentId)) + ..orderBy([OrderingTerm.asc(messages.createdAt)])) + .map(_messageFromJoinRow) + .get()); + + if (lessThan != null) { + final lessThanIndex = msgList.indexWhere((m) => m.id == lessThan); + msgList.removeRange(lessThanIndex, msgList.length); + } + return msgList; + } + + /// + Future> getMessagesByCid( + String cid, { + int limit = 20, + String messageLessThan, + String messageGreaterThan, + }) async { + final msgList = await Future.wait(await (select(messages).join([ + leftOuterJoin(users, messages.userId.equalsExp(users.id)), + ]) + ..where(messages.channelCid.equals(cid)) + ..where( + isNull(messages.parentId) | messages.showInChannel.equals(true)) + ..orderBy([OrderingTerm.asc(messages.createdAt)])) + .map(_messageFromJoinRow) + .get()); + + if (messageLessThan != null) { + final lessThanIndex = msgList.indexWhere((m) => m.id == messageLessThan); + if (lessThanIndex != -1) { + msgList.removeRange(lessThanIndex, msgList.length); + } + } + if (messageGreaterThan != null) { + final greaterThanIndex = + msgList.indexWhere((m) => m.id == messageGreaterThan); + if (greaterThanIndex != -1) { + msgList.removeRange(0, greaterThanIndex); + } + } + if (limit != null) { + return msgList.take(limit).toList(); + } + return msgList; + } + + /// + Future updateMessages(String cid, List messageList) { + return batch((batch) { + batch.insertAll( + messages, + messageList.map((it) => it.toEntity(cid: cid)).toList(), + mode: InsertMode.insertOrReplace, + ); + }); + } +} diff --git a/packages/stream_chat_persistence/lib/src/dao/reaction_dao.dart b/packages/stream_chat_persistence/lib/src/dao/reaction_dao.dart new file mode 100644 index 00000000..8df23d44 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/dao/reaction_dao.dart @@ -0,0 +1,50 @@ +import 'package:moor/moor.dart'; +import 'package:stream_chat/stream_chat.dart'; +import 'package:stream_chat_persistence/src/db/moor_chat_database.dart'; +import 'package:stream_chat_persistence/src/entity/reactions.dart'; +import 'package:stream_chat_persistence/src/entity/users.dart'; +import '../mapper/mapper.dart'; + +part 'reaction_dao.g.dart'; + +/// +@UseDao(tables: [Reactions, Users]) +class ReactionDao extends DatabaseAccessor + with _$ReactionDaoMixin { + /// + ReactionDao(MoorChatDatabase db) : super(db); + + /// + Future> getReactions(String messageId) { + return (select(reactions).join([ + leftOuterJoin(users, reactions.userId.equalsExp(users.id)), + ]) + ..where(reactions.messageId.equals(messageId)) + ..orderBy([OrderingTerm.asc(reactions.createdAt)])) + .map((rows) { + final userEntity = rows.readTable(users); + final reactionEntity = rows.readTable(reactions); + return reactionEntity.toReaction(user: userEntity?.toUser()); + }).get(); + } + + /// + Future> getOwnReactions( + String messageId, + String userId, + ) async { + final reactions = await getReactions(messageId); + return reactions.where((it) => it.userId == userId).toList(); + } + + /// + Future updateReactions(List reactionList) { + return batch((it) { + it.insertAll( + reactions, + reactionList.map((r) => r.toEntity()).toList(), + mode: InsertMode.insertOrReplace, + ); + }); + } +} diff --git a/packages/stream_chat_persistence/lib/src/dao/read_dao.dart b/packages/stream_chat_persistence/lib/src/dao/read_dao.dart new file mode 100644 index 00000000..cd4da0f6 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/dao/read_dao.dart @@ -0,0 +1,42 @@ +import 'package:moor/moor.dart'; +import 'package:stream_chat/stream_chat.dart'; +import 'package:stream_chat_persistence/src/db/moor_chat_database.dart'; +import 'package:stream_chat_persistence/src/entity/reads.dart'; +import 'package:stream_chat_persistence/src/entity/users.dart'; +import '../mapper/mapper.dart'; + +part 'read_dao.g.dart'; + +/// +@UseDao(tables: [Reads, Users]) +class ReadDao extends DatabaseAccessor with _$ReadDaoMixin { + /// + ReadDao(MoorChatDatabase db) : super(db); + + /// Get all reads where [reads.channelCid] matches [cid] + Future> getReadsByCid(String cid) async { + return (select(reads).join([ + leftOuterJoin(users, reads.userId.equalsExp(users.id)), + ]) + ..where(reads.channelCid.equals(cid)) + ..orderBy([ + OrderingTerm.asc(reads.lastRead), + ])) + .map((row) { + final userEntity = row.readTable(users); + final readEntity = row.readTable(reads); + return readEntity.toRead(user: userEntity?.toUser()); + }).get(); + } + + /// + Future updateReads(String cid, List readList) { + return batch( + (it) => it.insertAll( + reads, + readList.map((r) => r.toEntity(cid: cid)).toList(), + mode: InsertMode.insertOrReplace, + ), + ); + } +} diff --git a/packages/stream_chat_persistence/lib/src/dao/user_dao.dart b/packages/stream_chat_persistence/lib/src/dao/user_dao.dart new file mode 100644 index 00000000..630a8252 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/dao/user_dao.dart @@ -0,0 +1,25 @@ +import 'package:moor/moor.dart'; +import 'package:stream_chat/stream_chat.dart'; +import 'package:stream_chat_persistence/src/db/moor_chat_database.dart'; +import 'package:stream_chat_persistence/src/entity/users.dart'; +import '../mapper/user_mapper.dart'; + +part 'user_dao.g.dart'; + +/// +@UseDao(tables: [Users]) +class UserDao extends DatabaseAccessor with _$UserDaoMixin { + /// + UserDao(MoorChatDatabase db) : super(db); + + /// + Future updateUsers(List userList) { + return batch( + (it) => it.insertAll( + users, + userList.map((u) => u.toEntity()).toList(), + mode: InsertMode.insertOrReplace, + ), + ); + } +} diff --git a/packages/stream_chat_persistence/lib/src/db/moor_chat_database.dart b/packages/stream_chat_persistence/lib/src/db/moor_chat_database.dart new file mode 100644 index 00000000..21447981 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/db/moor_chat_database.dart @@ -0,0 +1,83 @@ +import 'package:moor/isolate.dart'; +import 'package:moor/moor.dart'; +import 'package:stream_chat/stream_chat.dart'; + +import '../entity/entity.dart'; +import '../dao/dao.dart'; +import 'shared/shared_db.dart'; + +part 'moor_chat_database.g.dart'; + +LazyDatabase _openConnection( + String dbName, { + logStatements = false, +}) { + return LazyDatabase(() async { + return await SharedDB.constructDatabase( + dbName, + logStatements: logStatements, + ); + }); +} + +/// +@UseMoor(tables: [ + Channels, + Messages, + Reactions, + Users, + Members, + Reads, + ChannelQueries, + ConnectionEvents, +], daos: [ + UserDao, + ChannelDao, + MessageDao, + MemberDao, + ReactionDao, + ReadDao, + ChannelQueryDao, + ConnectionEventDao, +]) +class MoorChatDatabase extends _$MoorChatDatabase { + /// Instantiate a new database instance + MoorChatDatabase( + String dbName, { + logStatements = false, + }) : super(_openConnection( + dbName, + logStatements: logStatements, + )); + + /// Instantiate a new database instance + MoorChatDatabase.connect( + this._isolate, + DatabaseConnection connection, + ) : super.connect(connection); + + MoorIsolate _isolate; + + // you should bump this number whenever you change or add a table definition. + @override + int get schemaVersion => 1; + + @override + MigrationStrategy get migration => MigrationStrategy( + onUpgrade: (openingDetails, before, after) async { + if (before != after) { + final m = createMigrator(); + for (final table in allTables) { + await m.deleteTable(table.actualTableName); + await m.createTable(table); + } + } + }, + ); + + /// Closes the database instance + Future disconnect() async { + await _isolate?.shutdownAll(); + await close(); + } +} diff --git a/packages/stream_chat_persistence/lib/src/db/shared/native_db.dart b/packages/stream_chat_persistence/lib/src/db/shared/native_db.dart new file mode 100644 index 00000000..198b9010 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/db/shared/native_db.dart @@ -0,0 +1,87 @@ +//ignore_for_file: public_member_api_docs + +import 'dart:io'; +import 'dart:isolate'; +import 'package:moor/ffi.dart'; +import 'package:moor/isolate.dart'; +import 'package:moor/moor.dart'; +import 'package:path/path.dart'; +import 'package:path_provider/path_provider.dart'; + +import '../moor_chat_database.dart'; + +class SharedDB { + static Future constructDatabase( + String dbName, { + bool logStatements = false, + }) async { + if (Platform.isIOS || Platform.isAndroid) { + final dir = await getApplicationDocumentsDirectory(); + final path = join(dir.path, '$dbName.sqlite'); + final file = File(path); + return VmDatabase(file, logStatements: logStatements); + } + if (Platform.isMacOS || Platform.isLinux) { + final file = File('$dbName.sqlite'); + return VmDatabase(file, logStatements: logStatements); + } + return VmDatabase.memory(logStatements: logStatements); + } + + static void _startBackground(_IsolateStartRequest request) { + final executor = LazyDatabase(() async { + return VmDatabase( + File(request.targetPath), + logStatements: request.logStatements, + ); + }); + final moorIsolate = MoorIsolate.inCurrent( + () => DatabaseConnection.fromExecutor(executor), + ); + request.sendMoorIsolate.send(moorIsolate); + } + + static Future _createMoorIsolate( + String dbName, { + bool logStatements = false, + }) async { + final dir = await getApplicationDocumentsDirectory(); + final path = join(dir.path, '$dbName.sqlite'); + + final receivePort = ReceivePort(); + await Isolate.spawn( + _startBackground, + _IsolateStartRequest( + receivePort.sendPort, + path, + logStatements: logStatements, + ), + ); + + return (await receivePort.first as MoorIsolate); + } + + static Future constructOfflineStorage( + String dbName, { + bool logStatements = false, + }) async { + final isolate = await _createMoorIsolate( + dbName, + logStatements: logStatements, + ); + final connection = await isolate.connect(); + return MoorChatDatabase.connect(isolate, connection); + } +} + +class _IsolateStartRequest { + final SendPort sendMoorIsolate; + final String targetPath; + final bool logStatements; + + const _IsolateStartRequest( + this.sendMoorIsolate, + this.targetPath, { + this.logStatements = false, + }); +} diff --git a/packages/stream_chat_persistence/lib/src/db/shared/shared_db.dart b/packages/stream_chat_persistence/lib/src/db/shared/shared_db.dart new file mode 100644 index 00000000..25120fcf --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/db/shared/shared_db.dart @@ -0,0 +1,3 @@ +export 'unsupported_db.dart' + if (dart.library.io) 'native_db.dart' // implementation using dart:io + if (dart.library.html) 'web_db.dart'; diff --git a/packages/stream_chat_persistence/lib/src/db/shared/unsupported_db.dart b/packages/stream_chat_persistence/lib/src/db/shared/unsupported_db.dart new file mode 100644 index 00000000..7415fdba --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/db/shared/unsupported_db.dart @@ -0,0 +1,18 @@ +//ignore_for_file: public_member_api_docs +//ignore_for_file: always_declare_return_types + +class SharedDB { + static constructDatabase( + String dbName, { + bool logStatements = false, + }) { + throw 'Unsupported Platform'; + } + + static constructOfflineStorage( + String dbName, { + logStatements = false, + }) { + throw 'Unsupported Platform'; + } +} diff --git a/packages/stream_chat_persistence/lib/src/db/shared/web_db.dart b/packages/stream_chat_persistence/lib/src/db/shared/web_db.dart new file mode 100644 index 00000000..31b4f4b6 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/db/shared/web_db.dart @@ -0,0 +1,21 @@ +//ignore_for_file: public_member_api_docs +//ignore_for_file: always_declare_return_types +import 'package:moor/moor_web.dart'; + +import '../moor_chat_database.dart'; + +class SharedDB { + static constructDatabase( + String dbName, { + bool logStatements = false, + }) async { + return WebDatabase(dbName, logStatements: logStatements); + } + + static Future constructOfflineStorage( + String dbName, { + bool logStatements = false, + }) async { + return MoorChatDatabase(dbName, logStatements: logStatements); + } +} diff --git a/packages/stream_chat_persistence/lib/src/db/stream_chat_database_impl.dart b/packages/stream_chat_persistence/lib/src/db/stream_chat_database_impl.dart new file mode 100644 index 00000000..de9bbe45 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/db/stream_chat_database_impl.dart @@ -0,0 +1,222 @@ +import 'package:stream_chat/stream_chat.dart'; + +import 'moor_chat_database.dart'; +import 'shared/shared_db.dart'; + +/// +class StreamChatDatabaseImpl implements StreamChatDatabase { + /// + StreamChatDatabaseImpl( + this._userId, { + Logger logger, + }) : _logger = logger, + assert(_userId != null); + + final String _userId; + final Logger _logger; + MoorChatDatabase _db; + + bool get _debugAssertConnected { + assert(() { + if (_db == null) { + throw Exception( + 'A $runtimeType was used after being disconnected.\n' + 'Once you have called disconnect() on a $runtimeType, it can no longer be used.', + ); + } + return true; + }()); + return true; + } + + @override + Future connect({ + bool connectBackground = false, + bool logStatements = false, + }) async { + if (_db != null) { + throw Exception( + 'An instance of StreamChatDatabase is already connected.\n' + 'disconnect the previous instance before connecting again.', + ); + } + + final dbName = 'db_$_userId'; + if (connectBackground) { + _logger?.info('Connecting on background isolate'); + _db = await SharedDB.constructOfflineStorage( + dbName, + logStatements: logStatements, + ); + } else { + _logger?.info('Connecting on a regular isolate'); + _db = MoorChatDatabase(dbName, logStatements: logStatements); + } + } + + @override + Future getConnectionInfo() { + return _db.connectionEventDao.connectionEvent; + } + + @override + Future updateConnectionInfo(Event event) { + return _db.connectionEventDao.updateConnectionEvent(event); + } + + @override + Future updateLastSyncAt(DateTime lastSyncAt) { + return _db.connectionEventDao.updateLastSyncAt(lastSyncAt); + } + + @override + Future getLastSyncAt() { + return _db.connectionEventDao.lastSyncAt; + } + + @override + Future deleteChannelByCids(List cids) { + return _db.channelDao.deleteChannelByCids(cids); + } + + @override + Future> getChannelCids() => _db.channelDao.cids; + + @override + Future deleteMessageByIds(List messageIds) { + return _db.messageDao.deleteMessageByIds(messageIds); + } + + @override + Future deleteMessageByCids(List cids) { + return _db.messageDao.deleteMessageByCids(cids); + } + + @override + Future> getMembersByCid(String cid) { + return _db.memberDao.getMembersByCid(cid); + } + + @override + Future getChannelByCid(String cid) { + return _db.channelDao.getChannelByCid(cid); + } + + @override + Future> getMessagesByCid( + String cid, { + int limit = 20, + String messageLessThan, + String messageGreaterThan, + }) { + return _db.messageDao.getMessagesByCid( + cid, + limit: limit, + messageLessThan: messageLessThan, + messageGreaterThan: messageGreaterThan, + ); + } + + @override + Future> getReadsByCid(String cid) { + return _db.readDao.getReadsByCid(cid); + } + + @override + Future>> getChannelThreads(String cid) async { + final messages = await _db.messageDao.getThreadMessages(cid); + final messageByParentIdDictionary = >{}; + for (final message in messages) { + final parentId = message.parentId; + messageByParentIdDictionary[parentId] = [ + ...messageByParentIdDictionary[parentId] ?? [], + message + ]; + } + return messageByParentIdDictionary; + } + + @override + Future> getReplies( + String parentId, { + String lessThan, + }) { + return _db.messageDao.getThreadMessagesByParentId( + parentId, + lessThan: lessThan, + ); + } + + @override + Future> getChannelStates({ + Map filter, + List sort = const [], + PaginationParams paginationParams, + }) { + return _db.channelQueryDao.getChannelStates( + filter: filter, + sort: sort, + paginationParams: paginationParams, + ); + } + + @override + Future updateChannelQueries( + Map filter, + List cids, + bool clearQueryCache, + ) { + return _db.channelQueryDao.updateChannelQueries( + filter, + cids, + clearQueryCache, + ); + } + + @override + Future updateChannels(List channels) { + return _db.channelDao.updateChannels(channels); + } + + @override + Future updateMembers(String cid, List members) { + return _db.memberDao.updateMembers(cid, members); + } + + @override + Future updateMessages(String cid, List messages) { + return _db.messageDao.updateMessages(cid, messages); + } + + @override + Future updateReactions(List reactions) { + return _db.reactionDao.updateReactions(reactions); + } + + @override + Future updateReads(String cid, List reads) { + return _db.readDao.updateReads(cid, reads); + } + + @override + Future updateUsers(List users) { + return _db.userDao.updateUsers(users); + } + + @override + Future disconnect({bool flush = false}) async { + if (_db != null) { + _logger?.info('Disconnecting'); + if (flush) { + _logger?.info('Flushing'); + await _db.batch((batch) { + _db.allTables.forEach((table) { + _db.delete(table).go(); + }); + }); + } + await _db.disconnect(); + _db = null; + } + } +} diff --git a/packages/stream_chat_persistence/lib/src/entity/channel_queries.dart b/packages/stream_chat_persistence/lib/src/entity/channel_queries.dart new file mode 100644 index 00000000..fefdc925 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/entity/channel_queries.dart @@ -0,0 +1,14 @@ +import 'package:moor/moor.dart'; + +@DataClassName('ChannelQueryEntity') +class ChannelQueries extends Table { + TextColumn get queryHash => text()(); + + TextColumn get channelCid => text()(); + + @override + Set get primaryKey => { + queryHash, + channelCid, + }; +} diff --git a/packages/stream_chat_persistence/lib/src/entity/channels.dart b/packages/stream_chat_persistence/lib/src/entity/channels.dart new file mode 100644 index 00000000..9f652ed1 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/entity/channels.dart @@ -0,0 +1,32 @@ +import 'package:moor/moor.dart'; +import 'package:stream_chat_persistence/src/converter/map_converter.dart'; + +@DataClassName('ChannelEntity') +class Channels extends Table { + TextColumn get id => text()(); + + TextColumn get type => text()(); + + TextColumn get cid => text()(); + + TextColumn get config => text().map(MapConverter())(); + + BoolColumn get frozen => boolean().withDefault(Constant(false))(); + + DateTimeColumn get lastMessageAt => dateTime().nullable()(); + + DateTimeColumn get createdAt => dateTime().nullable()(); + + DateTimeColumn get updatedAt => dateTime().nullable()(); + + DateTimeColumn get deletedAt => dateTime().nullable()(); + + IntColumn get memberCount => integer().nullable()(); + + TextColumn get createdById => text().nullable()(); + + TextColumn get extraData => text().nullable().map(MapConverter())(); + + @override + Set get primaryKey => {cid}; +} diff --git a/packages/stream_chat_persistence/lib/src/entity/connection_events.dart b/packages/stream_chat_persistence/lib/src/entity/connection_events.dart new file mode 100644 index 00000000..311f2da0 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/entity/connection_events.dart @@ -0,0 +1,19 @@ +import 'package:moor/moor.dart'; + +@DataClassName('ConnectionEventEntity') +class ConnectionEvents extends Table { + IntColumn get id => integer()(); + + TextColumn get ownUserId => text().nullable()(); + + IntColumn get totalUnreadCount => integer().nullable()(); + + IntColumn get unreadChannels => integer().nullable()(); + + DateTimeColumn get lastEventAt => dateTime().nullable()(); + + DateTimeColumn get lastSyncAt => dateTime().nullable()(); + + @override + Set get primaryKey => {id}; +} diff --git a/packages/stream_chat_persistence/lib/src/entity/entity.dart b/packages/stream_chat_persistence/lib/src/entity/entity.dart new file mode 100644 index 00000000..8751f538 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/entity/entity.dart @@ -0,0 +1,8 @@ +export 'channels.dart'; +export 'messages.dart'; +export 'reactions.dart'; +export 'users.dart'; +export 'members.dart'; +export 'reads.dart'; +export 'channel_queries.dart'; +export 'connection_events.dart'; diff --git a/packages/stream_chat_persistence/lib/src/entity/members.dart b/packages/stream_chat_persistence/lib/src/entity/members.dart new file mode 100644 index 00000000..39e224e0 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/entity/members.dart @@ -0,0 +1,33 @@ +import 'package:moor/moor.dart'; + +@DataClassName('MemberEntity') +class Members extends Table { + TextColumn get userId => text()(); + + TextColumn get channelCid => + text().customConstraint('REFERENCES channels(cid) ON DELETE CASCADE')(); + + TextColumn get role => text().nullable()(); + + DateTimeColumn get inviteAcceptedAt => dateTime().nullable()(); + + DateTimeColumn get inviteRejectedAt => dateTime().nullable()(); + + BoolColumn get invited => boolean().nullable()(); + + BoolColumn get banned => boolean().nullable()(); + + BoolColumn get shadowBanned => boolean().nullable()(); + + BoolColumn get isModerator => boolean().nullable()(); + + DateTimeColumn get createdAt => dateTime()(); + + DateTimeColumn get updatedAt => dateTime().nullable()(); + + @override + Set get primaryKey => { + userId, + channelCid, + }; +} diff --git a/packages/stream_chat_persistence/lib/src/entity/messages.dart b/packages/stream_chat_persistence/lib/src/entity/messages.dart new file mode 100644 index 00000000..10c98315 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/entity/messages.dart @@ -0,0 +1,54 @@ +import 'package:moor/moor.dart'; +import 'package:stream_chat_persistence/src/converter/list_converter.dart'; +import 'package:stream_chat_persistence/src/converter/map_converter.dart'; +import 'package:stream_chat_persistence/src/converter/message_sending_status_converter.dart'; + +@DataClassName('MessageEntity') +class Messages extends Table { + TextColumn get id => text()(); + + TextColumn get messageText => text().nullable()(); + + TextColumn get attachments => + text().nullable().map(ListConverter())(); + + IntColumn get status => + integer().nullable().map(MessageSendingStatusConverter())(); + + TextColumn get type => text().nullable()(); + + TextColumn get mentionedUsers => + text().nullable().map(ListConverter())(); + + TextColumn get reactionCounts => text().nullable().map(MapConverter())(); + + TextColumn get reactionScores => text().nullable().map(MapConverter())(); + + TextColumn get parentId => text().nullable()(); + + TextColumn get quotedMessageId => text().nullable()(); + + IntColumn get replyCount => integer().nullable()(); + + BoolColumn get showInChannel => boolean().nullable()(); + + BoolColumn get shadowed => boolean().nullable()(); + + TextColumn get command => text().nullable()(); + + DateTimeColumn get createdAt => dateTime()(); + + DateTimeColumn get updatedAt => dateTime().nullable()(); + + DateTimeColumn get deletedAt => dateTime().nullable()(); + + TextColumn get userId => text().nullable()(); + + TextColumn get channelCid => text().nullable().customConstraint( + 'NULLABLE REFERENCES channels(cid) ON DELETE CASCADE')(); + + TextColumn get extraData => text().nullable().map(MapConverter())(); + + @override + Set get primaryKey => {id}; +} diff --git a/packages/stream_chat_persistence/lib/src/entity/reactions.dart b/packages/stream_chat_persistence/lib/src/entity/reactions.dart new file mode 100644 index 00000000..a20ffbeb --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/entity/reactions.dart @@ -0,0 +1,25 @@ +import 'package:moor/moor.dart'; +import 'package:stream_chat_persistence/src/converter/map_converter.dart'; + +@DataClassName('ReactionEntity') +class Reactions extends Table { + TextColumn get userId => text()(); + + TextColumn get messageId => + text().customConstraint('REFERENCES messages(id) ON DELETE CASCADE')(); + + TextColumn get type => text()(); + + DateTimeColumn get createdAt => dateTime()(); + + IntColumn get score => integer().nullable()(); + + TextColumn get extraData => text().nullable().map(MapConverter())(); + + @override + Set get primaryKey => { + messageId, + type, + userId, + }; +} diff --git a/packages/stream_chat_persistence/lib/src/entity/reads.dart b/packages/stream_chat_persistence/lib/src/entity/reads.dart new file mode 100644 index 00000000..c36dceef --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/entity/reads.dart @@ -0,0 +1,19 @@ +import 'package:moor/moor.dart'; + +@DataClassName('ReadEntity') +class Reads extends Table { + DateTimeColumn get lastRead => dateTime()(); + + TextColumn get userId => text()(); + + TextColumn get channelCid => + text().customConstraint('REFERENCES channels(cid) ON DELETE CASCADE')(); + + IntColumn get unreadMessages => integer().nullable()(); + + @override + Set get primaryKey => { + userId, + channelCid, + }; +} diff --git a/packages/stream_chat_persistence/lib/src/entity/users.dart b/packages/stream_chat_persistence/lib/src/entity/users.dart new file mode 100644 index 00000000..4beee396 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/entity/users.dart @@ -0,0 +1,24 @@ +import 'package:moor/moor.dart'; +import 'package:stream_chat_persistence/src/converter/map_converter.dart'; + +@DataClassName('UserEntity') +class Users extends Table { + TextColumn get id => text()(); + + TextColumn get role => text().nullable()(); + + DateTimeColumn get createdAt => dateTime().nullable()(); + + DateTimeColumn get updatedAt => dateTime().nullable()(); + + DateTimeColumn get lastActive => dateTime().nullable()(); + + BoolColumn get online => boolean().nullable()(); + + BoolColumn get banned => boolean().nullable()(); + + TextColumn get extraData => text().nullable().map(MapConverter())(); + + @override + Set get primaryKey => {id}; +} diff --git a/packages/stream_chat_persistence/lib/src/mapper/channel_mapper.dart b/packages/stream_chat_persistence/lib/src/mapper/channel_mapper.dart new file mode 100644 index 00000000..cdabf17f --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/mapper/channel_mapper.dart @@ -0,0 +1,61 @@ +import 'package:stream_chat/stream_chat.dart'; +import 'user_mapper.dart'; +import 'package:stream_chat_persistence/src/db/moor_chat_database.dart'; + +/// +extension ChannelEntityX on ChannelEntity { + /// + ChannelModel toChannelModel({User createdBy}) { + final config = ChannelConfig.fromJson(this.config ?? {}); + return ChannelModel( + id: id, + config: config, + type: type, + frozen: frozen, + createdAt: createdAt, + updatedAt: updatedAt, + memberCount: memberCount, + cid: cid, + lastMessageAt: lastMessageAt, + deletedAt: deletedAt, + extraData: extraData, + createdBy: createdBy, + ); + } + + /// + ChannelState toChannelState({ + User createdBy, + List members, + List reads, + List messages, + }) { + return ChannelState( + members: members, + read: reads, + messages: messages, + channel: toChannelModel(createdBy: createdBy), + ); + } +} + +/// +extension ChannelModelX on ChannelModel { + /// + ChannelEntity toEntity() { + return ChannelEntity( + id: id, + type: type, + cid: cid, + config: config.toJson(), + frozen: frozen, + lastMessageAt: lastMessageAt, + createdAt: createdAt, + updatedAt: updatedAt, + deletedAt: deletedAt, + memberCount: memberCount, + createdById: createdBy.id, + extraData: extraData, + ); + } +} diff --git a/packages/stream_chat_persistence/lib/src/mapper/event_mapper.dart b/packages/stream_chat_persistence/lib/src/mapper/event_mapper.dart new file mode 100644 index 00000000..41e90205 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/mapper/event_mapper.dart @@ -0,0 +1,15 @@ +import 'package:stream_chat/stream_chat.dart'; +import 'user_mapper.dart'; +import 'package:stream_chat_persistence/src/db/moor_chat_database.dart'; + +/// +extension ConnectionEventX on ConnectionEventEntity { + /// + Event toEvent({User user}) { + return Event( + me: user, + totalUnreadCount: totalUnreadCount, + unreadChannels: unreadChannels, + ); + } +} diff --git a/packages/stream_chat_persistence/lib/src/mapper/mapper.dart b/packages/stream_chat_persistence/lib/src/mapper/mapper.dart new file mode 100644 index 00000000..efedde22 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/mapper/mapper.dart @@ -0,0 +1,7 @@ +export 'user_mapper.dart'; +export 'reaction_mapper.dart'; +export 'channel_mapper.dart'; +export 'event_mapper.dart'; +export 'member_mapper.dart'; +export 'read_mapper.dart'; +export 'message_mapper.dart'; \ No newline at end of file diff --git a/packages/stream_chat_persistence/lib/src/mapper/member_mapper.dart b/packages/stream_chat_persistence/lib/src/mapper/member_mapper.dart new file mode 100644 index 00000000..42f463ba --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/mapper/member_mapper.dart @@ -0,0 +1,43 @@ +import 'package:stream_chat/stream_chat.dart'; +import 'package:stream_chat_persistence/src/db/moor_chat_database.dart'; +import 'user_mapper.dart'; + +/// +extension MemberEntityX on MemberEntity { + /// + Member toMember({User user}) { + return Member( + user: user, + userId: userId, + banned: banned, + shadowBanned: shadowBanned, + updatedAt: updatedAt, + createdAt: createdAt, + role: role, + inviteAcceptedAt: inviteAcceptedAt, + invited: invited, + inviteRejectedAt: inviteRejectedAt, + isModerator: isModerator, + ); + } +} + +/// +extension MemberX on Member { + /// + MemberEntity toEntity({String cid}) { + return MemberEntity( + userId: user?.id, + banned: banned, + shadowBanned: shadowBanned, + channelCid: cid, + createdAt: createdAt, + isModerator: isModerator, + inviteRejectedAt: inviteRejectedAt, + invited: invited, + inviteAcceptedAt: inviteAcceptedAt, + role: role, + updatedAt: updatedAt, + ); + } +} diff --git a/packages/stream_chat_persistence/lib/src/mapper/message_mapper.dart b/packages/stream_chat_persistence/lib/src/mapper/message_mapper.dart new file mode 100644 index 00000000..4ea001a1 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/mapper/message_mapper.dart @@ -0,0 +1,70 @@ +import 'dart:convert'; + +import 'package:stream_chat/stream_chat.dart'; +import 'package:stream_chat_persistence/src/db/moor_chat_database.dart'; + +/// +extension MessageEntityX on MessageEntity { + /// + Message toMessage({ + User user, + List latestReactions, + List ownReactions, + Message quotedMessage, + }) { + return Message( + shadowed: shadowed, + latestReactions: latestReactions, + ownReactions: ownReactions, + attachments: attachments?.map((it) { + final json = jsonDecode(it); + return Attachment.fromJson(json); + })?.toList(), + createdAt: createdAt, + extraData: extraData, + updatedAt: updatedAt, + id: id, + type: type, + status: status, + command: command, + parentId: parentId, + quotedMessageId: quotedMessageId, + quotedMessage: quotedMessage, + reactionCounts: reactionCounts, + reactionScores: reactionScores, + replyCount: replyCount, + showInChannel: showInChannel, + text: messageText, + user: user, + deletedAt: deletedAt, + ); + } +} + +/// +extension MessageX on Message { + /// + MessageEntity toEntity({String cid}) { + return MessageEntity( + id: id, + attachments: attachments.map((it) => jsonEncode(it)).toList(), + channelCid: cid, + type: type, + parentId: parentId, + quotedMessageId: quotedMessageId, + command: command, + createdAt: createdAt, + shadowed: shadowed, + showInChannel: showInChannel, + replyCount: replyCount, + reactionScores: reactionScores, + reactionCounts: reactionCounts, + status: status, + updatedAt: updatedAt, + extraData: extraData, + userId: user?.id, + deletedAt: deletedAt, + messageText: text, + ); + } +} diff --git a/packages/stream_chat_persistence/lib/src/mapper/reaction_mapper.dart b/packages/stream_chat_persistence/lib/src/mapper/reaction_mapper.dart new file mode 100644 index 00000000..d468d11c --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/mapper/reaction_mapper.dart @@ -0,0 +1,33 @@ +import 'package:stream_chat/stream_chat.dart'; +import 'package:stream_chat_persistence/src/db/moor_chat_database.dart'; + +/// +extension ReactionEntityX on ReactionEntity { + /// + Reaction toReaction({User user}) { + return Reaction( + extraData: extraData, + type: type, + createdAt: createdAt, + userId: userId, + user: user, + messageId: messageId, + score: score, + ); + } +} + +/// +extension ReactionX on Reaction { + /// + ReactionEntity toEntity() { + return ReactionEntity( + extraData: extraData, + type: type, + createdAt: createdAt, + userId: userId, + messageId: messageId, + score: score, + ); + } +} diff --git a/packages/stream_chat_persistence/lib/src/mapper/read_mapper.dart b/packages/stream_chat_persistence/lib/src/mapper/read_mapper.dart new file mode 100644 index 00000000..0f9ea560 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/mapper/read_mapper.dart @@ -0,0 +1,28 @@ +import 'package:stream_chat/stream_chat.dart'; +import 'package:stream_chat_persistence/src/db/moor_chat_database.dart'; +import 'user_mapper.dart'; + +/// +extension ReadEntityX on ReadEntity { + /// + Read toRead({User user}) { + return Read( + user: user, + lastRead: lastRead, + unreadMessages: unreadMessages, + ); + } +} + +/// +extension ReadX on Read { + /// + ReadEntity toEntity({String cid}) { + return ReadEntity( + lastRead: lastRead, + userId: user?.id, + channelCid: cid, + unreadMessages: unreadMessages, + ); + } +} diff --git a/packages/stream_chat_persistence/lib/src/mapper/user_mapper.dart b/packages/stream_chat_persistence/lib/src/mapper/user_mapper.dart new file mode 100644 index 00000000..cb775ab7 --- /dev/null +++ b/packages/stream_chat_persistence/lib/src/mapper/user_mapper.dart @@ -0,0 +1,36 @@ +import 'package:stream_chat/stream_chat.dart'; +import 'package:stream_chat_persistence/src/db/moor_chat_database.dart'; + +/// +extension UserEntityX on UserEntity { + /// + User toUser() { + return User( + id: id, + updatedAt: updatedAt, + role: role, + online: online, + lastActive: lastActive, + extraData: extraData, + banned: banned, + createdAt: createdAt, + ); + } +} + +/// +extension UserX on User { + /// + UserEntity toEntity() { + return UserEntity( + id: id, + role: role, + createdAt: createdAt, + updatedAt: updatedAt, + lastActive: lastActive, + online: online, + banned: banned, + extraData: extraData, + ); + } +} diff --git a/packages/stream_chat_persistence/lib/stream_chat_persistence.dart b/packages/stream_chat_persistence/lib/stream_chat_persistence.dart new file mode 100644 index 00000000..3388d36e --- /dev/null +++ b/packages/stream_chat_persistence/lib/stream_chat_persistence.dart @@ -0,0 +1,3 @@ +library stream_chat_persistence; + +export 'src/db/stream_chat_database_impl.dart'; diff --git a/packages/stream_chat_persistence/pubspec.yaml b/packages/stream_chat_persistence/pubspec.yaml new file mode 100644 index 00000000..9ba42d1a --- /dev/null +++ b/packages/stream_chat_persistence/pubspec.yaml @@ -0,0 +1,18 @@ +name: stream_chat_persistence +description: A new Flutter package. +version: 0.0.1 +author: +homepage: + +environment: + sdk: ">=2.7.0 <3.0.0" + +dependencies: + moor: ^3.4.0 + stream_chat: + path: ../dart_client + +dev_dependencies: + test: ^1.15.7 + build_runner: ^1.10.13 + moor_generator: ^3.4.1