2016-01-16 06:26:07 +08:00
|
|
|
_ = require 'underscore'
|
2016-01-09 06:31:33 +08:00
|
|
|
Rx = require 'rx-lite'
|
2016-01-16 06:26:07 +08:00
|
|
|
|
|
|
|
{ObservableListDataSource,
|
2016-01-09 06:31:33 +08:00
|
|
|
DatabaseStore,
|
2016-01-16 06:26:07 +08:00
|
|
|
Message,
|
2016-01-09 06:31:33 +08:00
|
|
|
QueryResultSet,
|
2016-01-16 06:26:07 +08:00
|
|
|
QuerySubscription} = require 'nylas-exports'
|
2016-01-09 06:31:33 +08:00
|
|
|
|
|
|
|
_flatMapJoiningMessages = ($threadsResultSet) =>
|
|
|
|
# DatabaseView leverages `QuerySubscription` for threads /and/ for the
|
|
|
|
# messages on each thread, which are passed to out as `thread.metadata`.
|
|
|
|
|
|
|
|
$messagesResultSets = {}
|
|
|
|
|
|
|
|
# 2. when we receive a set of threads, we check to see if we have message
|
|
|
|
# observables for each thread. If threads have been added to the result set,
|
|
|
|
# we make a single database query and load /all/ the message metadata for
|
|
|
|
# the new threads at once. (This is a performance optimization -it's about
|
|
|
|
# ~80msec faster than making 100 queries for 100 new thread ids separately.)
|
2016-04-29 06:53:29 +08:00
|
|
|
$threadsResultSet.flatMapLatest (threadsResultSet) =>
|
2016-01-09 06:31:33 +08:00
|
|
|
missingIds = threadsResultSet.ids().filter (id) -> not $messagesResultSets[id]
|
2016-04-29 06:53:29 +08:00
|
|
|
if missingIds.length is 0
|
|
|
|
promise = Promise.resolve([threadsResultSet, []])
|
|
|
|
else
|
|
|
|
promise = DatabaseStore.findAll(Message, threadId: missingIds).then (messages) =>
|
|
|
|
Promise.resolve([threadsResultSet, messages])
|
|
|
|
|
|
|
|
return Rx.Observable.fromPromise(promise)
|
2016-01-09 06:31:33 +08:00
|
|
|
|
|
|
|
# 3. when that finishes, we group the loaded messsages by threadId and create
|
|
|
|
# the missing observables. Creating a query subscription would normally load
|
|
|
|
# an initial result set. To avoid that, we just hand new subscriptions the
|
|
|
|
# results we loaded in #2.
|
|
|
|
.flatMapLatest ([threadsResultSet, messagesForNewThreads]) =>
|
|
|
|
messagesGrouped = {}
|
|
|
|
for message in messagesForNewThreads
|
|
|
|
messagesGrouped[message.threadId] ?= []
|
|
|
|
messagesGrouped[message.threadId].push(message)
|
|
|
|
|
|
|
|
oldSets = $messagesResultSets
|
|
|
|
$messagesResultSets = {}
|
|
|
|
|
|
|
|
sets = threadsResultSet.ids().map (id) =>
|
|
|
|
$messagesResultSets[id] = oldSets[id] || _observableForThreadMessages(id, messagesGrouped[id])
|
|
|
|
$messagesResultSets[id]
|
|
|
|
sets.unshift(Rx.Observable.from([threadsResultSet]))
|
|
|
|
|
|
|
|
# 4. We use `combineLatest` to merge the message observables into a single
|
|
|
|
# stream (like Promise.all). When /any/ of them emit a new result set, we
|
|
|
|
# trigger.
|
|
|
|
Rx.Observable.combineLatest(sets)
|
|
|
|
|
|
|
|
.flatMapLatest ([threadsResultSet, messagesResultSets...]) =>
|
|
|
|
threadsWithMetadata = {}
|
|
|
|
threadsResultSet.models().map (thread, idx) ->
|
|
|
|
thread = new thread.constructor(thread)
|
|
|
|
thread.metadata = messagesResultSets[idx]?.models()
|
|
|
|
threadsWithMetadata[thread.id] = thread
|
|
|
|
|
|
|
|
Rx.Observable.from([QueryResultSet.setByApplyingModels(threadsResultSet, threadsWithMetadata)])
|
|
|
|
|
|
|
|
_observableForThreadMessages = (id, initialModels) ->
|
|
|
|
subscription = new QuerySubscription(DatabaseStore.findAll(Message, threadId: id), {
|
2016-04-20 02:32:33 +08:00
|
|
|
emitResultSet: true,
|
2016-01-09 06:31:33 +08:00
|
|
|
initialModels: initialModels
|
|
|
|
})
|
2016-02-05 06:14:24 +08:00
|
|
|
Rx.Observable.fromNamedQuerySubscription('message-'+id, subscription)
|
2016-01-09 06:31:33 +08:00
|
|
|
|
|
|
|
|
2016-01-16 06:26:07 +08:00
|
|
|
class ThreadListDataSource extends ObservableListDataSource
|
|
|
|
constructor: (subscription) ->
|
2016-02-05 06:14:24 +08:00
|
|
|
$resultSetObservable = Rx.Observable.fromNamedQuerySubscription('thread-list', subscription)
|
2016-01-16 06:26:07 +08:00
|
|
|
$resultSetObservable = _flatMapJoiningMessages($resultSetObservable)
|
2016-04-09 05:10:05 +08:00
|
|
|
super($resultSetObservable, subscription.replaceRange.bind(subscription))
|
2016-01-09 06:31:33 +08:00
|
|
|
|
2016-01-16 06:26:07 +08:00
|
|
|
module.exports = ThreadListDataSource
|