Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
34 changes: 33 additions & 1 deletion server/models/User.js
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,14 @@ class UserCache {

const userCache = new UserCache()

/**
* Pending media progress updates keyed by user id and media item
* Used to run updates for the same media item one at a time so concurrent requests
* (e.g. /api/session/local and /api/me/progress/batch/update) cannot each create a mediaProgress
* @type {Map<string, Promise<any>>}
*/
const mediaProgressUpdateQueues = new Map()

const { DataTypes, Model } = sequelize

/**
Expand Down Expand Up @@ -721,13 +729,37 @@ class User extends Model {
return mediaProgress?.getOldMediaProgress() || null
}

/**
* Updates for the same media item are queued so that concurrent requests do not create duplicate mediaProgress rows
*
* @param {ProgressUpdatePayload} progressPayload
* @returns {Promise<{ mediaProgress: import('./MediaProgress'), error: [string], statusCode: [number] }>}
*/
createUpdateMediaProgressFromPayload(progressPayload) {
const queueKey = `${this.id}:${progressPayload.episodeId || progressPayload.libraryItemId}`
const previousUpdate = mediaProgressUpdateQueues.get(queueKey) || Promise.resolve()
const update = previousUpdate.then(() => this.applyMediaProgressPayload(progressPayload))

// Keep the queue going if this update fails and remove the key once nothing else is queued
const queueTail = update.catch(() => {})
mediaProgressUpdateQueues.set(queueKey, queueTail)
queueTail.then(() => {
if (mediaProgressUpdateQueues.get(queueKey) === queueTail) {
mediaProgressUpdateQueues.delete(queueKey)
}
})

return update
}

/**
* TODO: Uses old model and should account for the different between ebook/audiobook progress
* Should only be called from createUpdateMediaProgressFromPayload
*
* @param {ProgressUpdatePayload} progressPayload
* @returns {Promise<{ mediaProgress: import('./MediaProgress'), error: [string], statusCode: [number] }>}
*/
async createUpdateMediaProgressFromPayload(progressPayload) {
async applyMediaProgressPayload(progressPayload) {
/** @type {import('./MediaProgress')|null} */
let mediaProgress = null
let mediaItemId = null
Expand Down
63 changes: 63 additions & 0 deletions test/server/models/User.test.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
const { expect } = require('chai')
const { Sequelize } = require('sequelize')
const sinon = require('sinon')

const Database = require('../../../server/Database')
const Logger = require('../../../server/Logger')

describe('User', () => {
beforeEach(async () => {
global.ServerSettings = {}
Database.sequelize = new Sequelize({ dialect: 'sqlite', storage: ':memory:', logging: false })
Database.sequelize.uppercaseFirst = (str) => (str ? `${str[0].toUpperCase()}${str.substr(1)}` : '')
await Database.buildModels()

sinon.stub(Logger, 'info')
})

afterEach(async () => {
sinon.restore()

// Clear all tables
await Database.sequelize.sync({ force: true })
})

describe('createUpdateMediaProgressFromPayload', () => {
let user
let book
let libraryItem

beforeEach(async () => {
const library = await Database.libraryModel.create({ name: 'Test Library', mediaType: 'book' })
const libraryFolder = await Database.libraryFolderModel.create({ path: '/test', libraryId: library.id })
book = await Database.bookModel.create({ title: 'Test Book', audioFiles: [], tags: [], narrators: [], genres: [], chapters: [] })
libraryItem = await Database.libraryItemModel.create({ libraryFiles: [], mediaId: book.id, mediaType: 'book', libraryId: library.id, libraryFolderId: libraryFolder.id })

await Database.userModel.create({ username: 'testuser', type: 'user', isActive: true })
user = await Database.userModel.findOne({ where: { username: 'testuser' }, include: Database.mediaProgressModel })
})

it('should create a single media progress when concurrent updates create progress for the same item', async () => {
const results = await Promise.all([user.createUpdateMediaProgressFromPayload({ libraryItemId: libraryItem.id, currentTime: 100, duration: 3600 }), user.createUpdateMediaProgressFromPayload({ libraryItemId: libraryItem.id, currentTime: 200, duration: 3600 }), user.createUpdateMediaProgressFromPayload({ libraryItemId: libraryItem.id, currentTime: 300, duration: 3600 })])

const mediaProgresses = await Database.mediaProgressModel.findAll({ where: { userId: user.id, mediaItemId: book.id } })
expect(mediaProgresses).to.have.length(1)
expect(mediaProgresses[0].currentTime).to.equal(300)
expect(results.map((r) => r.mediaProgress.id)).to.deep.equal([mediaProgresses[0].id, mediaProgresses[0].id, mediaProgresses[0].id])

expect(user.mediaProgresses).to.have.length(1)
expect(user.getMediaProgress(book.id).currentTime).to.equal(300)
})

it('should keep processing queued updates for the same item after an update fails', async () => {
const applyStub = sinon.stub(user, 'applyMediaProgressPayload').callThrough()
applyStub.onFirstCall().rejects(new Error('db error'))

const [failedUpdate, update] = await Promise.allSettled([user.createUpdateMediaProgressFromPayload({ libraryItemId: libraryItem.id, currentTime: 100, duration: 3600 }), user.createUpdateMediaProgressFromPayload({ libraryItemId: libraryItem.id, currentTime: 200, duration: 3600 })])

expect(failedUpdate.status).to.equal('rejected')
expect(update.status).to.equal('fulfilled')
expect(update.value.mediaProgress.currentTime).to.equal(200)
})
})
})