diff --git a/__tests__/services/mopidy/library-tracks.test.js b/__tests__/services/mopidy/library-tracks.test.js new file mode 100644 index 000000000..094f18cd8 --- /dev/null +++ b/__tests__/services/mopidy/library-tracks.test.js @@ -0,0 +1,216 @@ +import { applyMiddleware, combineReducers, createStore } from 'redux'; +import Mopidy from 'mopidy'; +import coreReducer from '../../../src/js/services/core/reducer'; +import mopidyMiddleware from '../../../src/js/services/mopidy/middleware'; +import mopidyReducer from '../../../src/js/services/mopidy/reducer'; +import { getLibraryTracks } from '../../../src/js/services/mopidy/actions'; +import * as uiActions from '../../../src/js/services/ui/actions'; +import uiReducer from '../../../src/js/services/ui/reducer'; + +jest.mock('mopidy', () => ({ + __esModule: true, + default: jest.fn(), +})); + +const libraryUri = 'local:directory?type=track'; +const processKey = 'MOPIDY_GET_LIBRARY_TRACKS'; + +const makeTrack = (uri) => ({ + uri, + name: `Track ${uri}`, + artists: [], +}); + +const makeSocket = (browseResponse, lookup) => ({ + on: jest.fn(), + close: jest.fn(), + library: { + browse: jest.fn(() => Promise.resolve(browseResponse)), + lookup, + }, +}); + +const recordActions = (actions) => () => (next) => (action) => { + actions.push(action); + return next(action); +}; + +const makeStore = (socket, actions = []) => { + Mopidy.mockImplementation(() => socket); + + const store = createStore( + combineReducers({ + core: coreReducer, + mopidy: mopidyReducer, + ui: uiReducer, + }), + { + core: { + items: {}, + libraries: {}, + }, + mopidy: { + connected: false, + host: 'localhost', + port: 6680, + ssl: false, + }, + ui: { + allow_reporting: false, + load_queue: {}, + processes: {}, + }, + }, + applyMiddleware(mopidyMiddleware, recordActions(actions)), + ); + + store.dispatch({ type: 'MOPIDY_CONNECT' }); + store.dispatch({ type: 'MOPIDY_CONNECTED' }); + + return store; +}; + +const waitForRequests = () => new Promise((resolve) => setTimeout(resolve, 0)); + +describe('Mopidy library track loading', () => { + beforeEach(() => { + Mopidy.mockReset(); + }); + + it('loads ordered tracks with bounded item commits and accurate progress', async () => { + const trackUris = Array.from( + { length: 5005 }, + (_, index) => `test:track:${index}`, + ); + const browseResponse = trackUris.map((uri) => ({ uri })); + const lookup = jest.fn(({ uris }) => Promise.resolve( + uris.reduce((response, uri) => ({ + ...response, + [uri]: [makeTrack(uri)], + }), {}), + )); + const actions = []; + const store = makeStore(makeSocket(browseResponse, lookup), actions); + + store.dispatch(uiActions.startLoading(libraryUri, libraryUri)); + store.dispatch(getLibraryTracks(libraryUri)); + await waitForRequests(); + + const lookupBatchSizes = lookup.mock.calls.map(([{ uris }]) => uris.length); + expect(lookupBatchSizes).toHaveLength(51); + expect(lookupBatchSizes.slice(0, -1).every((size) => size === 100)).toBe(true); + expect(lookupBatchSizes[lookupBatchSizes.length - 1]).toBe(5); + + const itemActions = actions.filter(({ type }) => type === 'ITEMS_LOADED'); + expect(itemActions.map(({ items }) => items.length)).toEqual([5000, 5]); + + const progressUpdates = actions + .filter(({ type }) => type === 'UPDATE_PROCESS') + .map(({ process }) => process); + expect(progressUpdates[0]).toEqual({ remaining: 5005, total: 5005 }); + expect(progressUpdates[progressUpdates.length - 1]).toEqual({ remaining: 0 }); + + const state = store.getState(); + expect(state.ui.processes[processKey].status).toBe('finished'); + expect(state.core.libraries[libraryUri].items_uris).toEqual(trackUris); + expect(Object.keys(state.core.items)).toHaveLength(trackUris.length); + expect(new Set(Object.keys(state.core.items)).size).toBe(trackUris.length); + }); + + it('caps commits when missing results cross the item boundary', async () => { + const trackUris = Array.from( + { length: 5100 }, + (_, index) => `test:track:${index}`, + ); + const missingUri = trackUris[4900]; + const browseResponse = trackUris.map((uri) => ({ uri })); + const lookup = jest.fn(({ uris }) => Promise.resolve( + uris.reduce((response, uri) => { + if (uri === missingUri) return response; + return { + ...response, + [uri]: [makeTrack(uri)], + }; + }, {}), + )); + const actions = []; + const store = makeStore(makeSocket(browseResponse, lookup), actions); + + store.dispatch(getLibraryTracks(libraryUri)); + await waitForRequests(); + + const itemActions = actions.filter(({ type }) => type === 'ITEMS_LOADED'); + expect(itemActions.every(({ items }) => items.length <= 5000)).toBe(true); + expect(itemActions.map(({ items }) => items.length)).toEqual([5000, 99]); + + const progressUpdates = actions + .filter(({ type }) => type === 'UPDATE_PROCESS') + .map(({ process }) => process); + expect(progressUpdates[0]).toEqual({ remaining: 5100, total: 5100 }); + expect(progressUpdates[progressUpdates.length - 1]).toEqual({ remaining: 0 }); + + const state = store.getState(); + expect(state.ui.processes[processKey].status).toBe('finished'); + expect(state.core.libraries[libraryUri].items_uris).toEqual(trackUris); + expect(Object.keys(state.core.items)).toHaveLength(5099); + expect(trackUris + .map((uri) => state.core.items[uri]) + .filter(Boolean) + .map((item) => item.uri)).toEqual(trackUris.filter((uri) => uri !== missingUri)); + }); + + it('keeps the browse order while ignoring missing lookup results', async () => { + const trackUris = [ + 'test:track:first', + 'test:track:missing', + 'test:track:last', + ]; + const browseResponse = trackUris.map((uri) => ({ uri })); + const lookup = jest.fn(({ uris }) => { + const response = {}; + response[uris[2]] = [makeTrack(uris[2])]; + response[uris[1]] = []; + response[uris[0]] = [makeTrack(uris[0])]; + return Promise.resolve(response); + }); + const store = makeStore(makeSocket(browseResponse, lookup)); + + store.dispatch(getLibraryTracks(libraryUri)); + await waitForRequests(); + + const { items, libraries } = store.getState().core; + expect(libraries[libraryUri].items_uris).toEqual(trackUris); + expect(trackUris.map((uri) => items[uri]).filter(Boolean).map((item) => item.uri)) + .toEqual([trackUris[0], trackUris[2]]); + expect(items[trackUris[1]]).toBeUndefined(); + }); + + it('cancels an in-flight load without requesting or publishing the next batch', async () => { + const trackUris = Array.from( + { length: 250 }, + (_, index) => `test:track:${index}`, + ); + const browseResponse = trackUris.map((uri) => ({ uri })); + let resolveLookup; + const lookup = jest.fn(() => new Promise((resolve) => { + resolveLookup = resolve; + })); + const store = makeStore(makeSocket(browseResponse, lookup)); + + store.dispatch(uiActions.startLoading(libraryUri, libraryUri)); + store.dispatch(getLibraryTracks(libraryUri)); + await waitForRequests(); + expect(lookup).toHaveBeenCalledTimes(1); + + store.dispatch(uiActions.cancelProcess(processKey)); + resolveLookup({ [trackUris[0]]: [makeTrack(trackUris[0])] }); + await waitForRequests(); + + const state = store.getState(); + expect(lookup).toHaveBeenCalledTimes(1); + expect(state.ui.processes[processKey].status).toBe('cancelled'); + expect(state.ui.load_queue).toEqual({}); + expect(state.core.libraries[libraryUri]).toBeUndefined(); + expect(state.core.items).toEqual({}); + }); +}); diff --git a/src/js/services/mopidy/middleware.js b/src/js/services/mopidy/middleware.js index e297c6d1e..fcb2b3ae3 100755 --- a/src/js/services/mopidy/middleware.js +++ b/src/js/services/mopidy/middleware.js @@ -44,6 +44,9 @@ const lastfmActions = require('../lastfm/actions.js'); const geniusActions = require('../genius/actions.js'); const discogsActions = require('../discogs/actions.js'); +const LIBRARY_TRACK_LOOKUP_BATCH_SIZE = 100; +const LIBRARY_TRACK_ITEMS_BATCH_SIZE = 5000; + /** * Fetch the method of our Mopidy API that is being called, by string * @@ -2291,6 +2294,33 @@ const MopidyMiddleware = (function () { request(store, 'library.browse', { uri: action.uri }) .then((browseResponse) => { const allUris = arrayOf('uri', browseResponse); + let nextUriIndex = 0; + let libraryItemsBuffer = []; + + const cancel = () => { + store.dispatch(uiActions.processCancelled(action.type)); + store.dispatch(uiActions.stopLoading(action.uri || 'mopidy:library:tracks')); + }; + + const isCancelling = () => { + const processor = store.getState().ui.processes[action.type]; + if (processor && processor.status === 'cancelling') { + cancel(); + return true; + } + return false; + }; + + const flushItems = (flushRemainder = false) => { + while ( + libraryItemsBuffer.length >= LIBRARY_TRACK_ITEMS_BATCH_SIZE + || (flushRemainder && libraryItemsBuffer.length) + ) { + const items = libraryItemsBuffer.slice(0, LIBRARY_TRACK_ITEMS_BATCH_SIZE); + libraryItemsBuffer = libraryItemsBuffer.slice(items.length); + store.dispatch(coreActions.itemsLoaded(items)); + } + }; store.dispatch(uiActions.updateProcess( action.type, @@ -2301,38 +2331,45 @@ const MopidyMiddleware = (function () { )); const run = () => { - if (allUris.length) { - const uris = allUris.splice(0, 100); - const processor = store.getState().ui.processes[action.type]; + if (isCancelling()) return; - if (processor && processor.status === 'cancelling') { - store.dispatch(uiActions.processCancelled(action.type)); - store.dispatch(uiActions.stopLoading('mopidy:library:tracks')); - return; - } - store.dispatch(uiActions.updateProcess(action.type, { remaining: allUris.length })); + if (nextUriIndex < allUris.length) { + const uris = allUris.slice( + nextUriIndex, + nextUriIndex + LIBRARY_TRACK_LOOKUP_BATCH_SIZE, + ); + nextUriIndex += uris.length; + store.dispatch(uiActions.updateProcess( + action.type, + { remaining: allUris.length - nextUriIndex }, + )); request(store, 'library.lookup', { uris }) .then( (lookupResponse) => { + if (isCancelling()) return; + const libraryItems = compact( indexToArray(lookupResponse).map( (results) => (results.length ? formatTrack(results[0]) : null), ), ); - if (libraryItems.length) { - store.dispatch(coreActions.itemsLoaded(libraryItems)); + libraryItemsBuffer.push(...libraryItems); + if (libraryItemsBuffer.length >= LIBRARY_TRACK_ITEMS_BATCH_SIZE) { + flushItems(); } run(); }, ); } else { + flushItems(true); + if (isCancelling()) return; store.dispatch(uiActions.processFinished(action.type)); store.dispatch(coreActions.libraryLoaded({ uri: action.uri, type: 'tracks', - items_uris: arrayOf('uri', browseResponse), + items_uris: allUris, })); } };