Viewers federation protocol v2
More efficient than the current one where instance is not fast enough to send all viewers if a video becomes popular The new protocol can be enabled by setting env USE_VIEWERS_FEDERATION_V2='true' Introduce a result field in View activity that contains the number of viewers. This field is used by the origin instance to send the total viewers on the video to remote instances. The difference with the current protocol is that we don't have to send viewers individually to remote instances. There are 4 cases: * View activity from federation on Remote Video -> instance replaces all current viewers by a new viewer that contains the result counter * View activity from federation on Local Video -> instance adds the viewer without considering the result counter * Local view on Remote Video -> instance adds the viewer and send it to the origin instance * Local view on Local Video -> instance adds the viewer Periodically PeerTube cleanups expired viewers. On local videos, the instance sends to remote instances a View activity with the result counter so they can update their viewers counter for that particular video
This commit is contained in:
parent
a73f476c8a
commit
b4f4432459
13 changed files with 327 additions and 172 deletions
|
@ -116,6 +116,11 @@ export interface ActivityView extends BaseActivity {
|
||||||
|
|
||||||
// If sending a "viewer" event
|
// If sending a "viewer" event
|
||||||
expires?: string
|
expires?: string
|
||||||
|
result?: {
|
||||||
|
type: 'InteractionCounter'
|
||||||
|
interactionType: 'WatchAction'
|
||||||
|
userInteractionCount: number
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
export interface ActivityDislike extends BaseActivity {
|
export interface ActivityDislike extends BaseActivity {
|
||||||
|
|
|
@ -18,6 +18,7 @@ export interface VideoObject {
|
||||||
licence: ActivityIdentifierObject
|
licence: ActivityIdentifierObject
|
||||||
language: ActivityIdentifierObject
|
language: ActivityIdentifierObject
|
||||||
subtitleLanguage: ActivityIdentifierObject[]
|
subtitleLanguage: ActivityIdentifierObject[]
|
||||||
|
|
||||||
views: number
|
views: number
|
||||||
|
|
||||||
sensitive: boolean
|
sensitive: boolean
|
||||||
|
|
|
@ -56,3 +56,7 @@ export function isProdInstance () {
|
||||||
export function getAppNumber () {
|
export function getAppNumber () {
|
||||||
return process.env.NODE_APP_INSTANCE || ''
|
return process.env.NODE_APP_INSTANCE || ''
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export function isUsingViewersFederationV2 () {
|
||||||
|
return process.env.USE_VIEWERS_FEDERATION_V2 === 'true'
|
||||||
|
}
|
||||||
|
|
|
@ -21,12 +21,7 @@ describe('Test video views/viewers counters', function () {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
before(async function () {
|
function runTests () {
|
||||||
this.timeout(120000)
|
|
||||||
|
|
||||||
servers = await prepareViewsServers()
|
|
||||||
})
|
|
||||||
|
|
||||||
describe('Test views counter on VOD', function () {
|
describe('Test views counter on VOD', function () {
|
||||||
let videoUUID: string
|
let videoUUID: string
|
||||||
|
|
||||||
|
@ -92,12 +87,17 @@ describe('Test video views/viewers counters', function () {
|
||||||
it('Should view twice and display 1 view/viewer', async function () {
|
it('Should view twice and display 1 view/viewer', async function () {
|
||||||
this.timeout(30000)
|
this.timeout(30000)
|
||||||
|
|
||||||
|
for (let i = 0; i < 3; i++) {
|
||||||
await servers[0].views.simulateViewer({ id: liveVideoId, currentTimes: [ 0, 35 ] })
|
await servers[0].views.simulateViewer({ id: liveVideoId, currentTimes: [ 0, 35 ] })
|
||||||
await servers[0].views.simulateViewer({ id: liveVideoId, currentTimes: [ 0, 35 ] })
|
await servers[0].views.simulateViewer({ id: liveVideoId, currentTimes: [ 0, 35 ] })
|
||||||
await servers[0].views.simulateViewer({ id: vodVideoId, currentTimes: [ 0, 5 ] })
|
await servers[0].views.simulateViewer({ id: vodVideoId, currentTimes: [ 0, 5 ] })
|
||||||
await servers[0].views.simulateViewer({ id: vodVideoId, currentTimes: [ 0, 5 ] })
|
await servers[0].views.simulateViewer({ id: vodVideoId, currentTimes: [ 0, 5 ] })
|
||||||
|
|
||||||
|
await wait(1000)
|
||||||
|
}
|
||||||
|
|
||||||
await waitJobs(servers)
|
await waitJobs(servers)
|
||||||
|
|
||||||
await checkCounter('viewers', liveVideoId, 1)
|
await checkCounter('viewers', liveVideoId, 1)
|
||||||
await checkCounter('viewers', vodVideoId, 1)
|
await checkCounter('viewers', vodVideoId, 1)
|
||||||
|
|
||||||
|
@ -108,21 +108,31 @@ describe('Test video views/viewers counters', function () {
|
||||||
})
|
})
|
||||||
|
|
||||||
it('Should wait and display 0 viewers but still have 1 view', async function () {
|
it('Should wait and display 0 viewers but still have 1 view', async function () {
|
||||||
this.timeout(30000)
|
this.timeout(45000)
|
||||||
|
|
||||||
await wait(12000)
|
let error = false
|
||||||
await waitJobs(servers)
|
|
||||||
|
|
||||||
|
do {
|
||||||
|
try {
|
||||||
await checkCounter('views', liveVideoId, 1)
|
await checkCounter('views', liveVideoId, 1)
|
||||||
await checkCounter('viewers', liveVideoId, 0)
|
await checkCounter('viewers', liveVideoId, 0)
|
||||||
|
|
||||||
await checkCounter('views', vodVideoId, 1)
|
await checkCounter('views', vodVideoId, 1)
|
||||||
await checkCounter('viewers', vodVideoId, 0)
|
await checkCounter('viewers', vodVideoId, 0)
|
||||||
|
|
||||||
|
error = false
|
||||||
|
await wait(2500)
|
||||||
|
} catch {
|
||||||
|
error = true
|
||||||
|
}
|
||||||
|
} while (error)
|
||||||
})
|
})
|
||||||
|
|
||||||
it('Should view on a remote and on local and display 2 viewers and 3 views', async function () {
|
it('Should view on a remote and on local and display appropriate views/viewers', async function () {
|
||||||
this.timeout(30000)
|
this.timeout(30000)
|
||||||
|
|
||||||
|
await servers[0].views.simulateViewer({ id: vodVideoId, xForwardedFor: '0.0.0.1,127.0.0.1', currentTimes: [ 0, 5 ] })
|
||||||
|
await servers[0].views.simulateViewer({ id: vodVideoId, xForwardedFor: '0.0.0.1,127.0.0.1', currentTimes: [ 0, 5 ] })
|
||||||
await servers[0].views.simulateViewer({ id: vodVideoId, currentTimes: [ 0, 5 ] })
|
await servers[0].views.simulateViewer({ id: vodVideoId, currentTimes: [ 0, 5 ] })
|
||||||
await servers[1].views.simulateViewer({ id: vodVideoId, currentTimes: [ 0, 5 ] })
|
await servers[1].views.simulateViewer({ id: vodVideoId, currentTimes: [ 0, 5 ] })
|
||||||
await servers[1].views.simulateViewer({ id: vodVideoId, currentTimes: [ 0, 5 ] })
|
await servers[1].views.simulateViewer({ id: vodVideoId, currentTimes: [ 0, 5 ] })
|
||||||
|
@ -131,23 +141,51 @@ describe('Test video views/viewers counters', function () {
|
||||||
await servers[1].views.simulateViewer({ id: liveVideoId, currentTimes: [ 0, 35 ] })
|
await servers[1].views.simulateViewer({ id: liveVideoId, currentTimes: [ 0, 35 ] })
|
||||||
await servers[1].views.simulateViewer({ id: liveVideoId, currentTimes: [ 0, 35 ] })
|
await servers[1].views.simulateViewer({ id: liveVideoId, currentTimes: [ 0, 35 ] })
|
||||||
|
|
||||||
|
await wait(3000) // Throttled federation
|
||||||
await waitJobs(servers)
|
await waitJobs(servers)
|
||||||
|
|
||||||
await checkCounter('viewers', liveVideoId, 2)
|
await checkCounter('viewers', liveVideoId, 2)
|
||||||
await checkCounter('viewers', vodVideoId, 2)
|
await checkCounter('viewers', vodVideoId, 3)
|
||||||
|
|
||||||
await processViewsBuffer(servers)
|
await processViewsBuffer(servers)
|
||||||
|
|
||||||
await checkCounter('views', liveVideoId, 3)
|
await checkCounter('views', liveVideoId, 3)
|
||||||
await checkCounter('views', vodVideoId, 3)
|
await checkCounter('views', vodVideoId, 4)
|
||||||
})
|
})
|
||||||
|
|
||||||
after(async function () {
|
after(async function () {
|
||||||
await stopFfmpeg(command)
|
await stopFfmpeg(command)
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
describe('Federation V1', function () {
|
||||||
|
|
||||||
|
before(async function () {
|
||||||
|
this.timeout(120000)
|
||||||
|
|
||||||
|
servers = await prepareViewsServers({ viewExpiration: '5 seconds', viewersFederationV2: false })
|
||||||
|
})
|
||||||
|
|
||||||
|
runTests()
|
||||||
|
|
||||||
after(async function () {
|
after(async function () {
|
||||||
await cleanupTests(servers)
|
await cleanupTests(servers)
|
||||||
})
|
})
|
||||||
|
})
|
||||||
|
|
||||||
|
describe('Federation V2', function () {
|
||||||
|
|
||||||
|
before(async function () {
|
||||||
|
this.timeout(120000)
|
||||||
|
|
||||||
|
servers = await prepareViewsServers({ viewExpiration: '5 seconds', viewersFederationV2: true })
|
||||||
|
})
|
||||||
|
|
||||||
|
runTests()
|
||||||
|
|
||||||
|
after(async function () {
|
||||||
|
await cleanupTests(servers)
|
||||||
|
})
|
||||||
|
})
|
||||||
})
|
})
|
||||||
|
|
|
@ -30,8 +30,17 @@ async function processViewsBuffer (servers: PeerTubeServer[]) {
|
||||||
await waitJobs(servers)
|
await waitJobs(servers)
|
||||||
}
|
}
|
||||||
|
|
||||||
async function prepareViewsServers () {
|
async function prepareViewsServers (options: {
|
||||||
const servers = await createMultipleServers(2)
|
viewersFederationV2?: boolean
|
||||||
|
viewExpiration?: string // default 1 second
|
||||||
|
} = {}) {
|
||||||
|
const { viewExpiration = '1 second' } = options
|
||||||
|
|
||||||
|
const env = options?.viewersFederationV2 === true
|
||||||
|
? { USE_VIEWERS_FEDERATION_V2: 'true' }
|
||||||
|
: undefined
|
||||||
|
|
||||||
|
const servers = await createMultipleServers(2, { views: { videos: { ip_view_expiration: viewExpiration } } }, { env })
|
||||||
await setAccessTokensToServers(servers)
|
await setAccessTokensToServers(servers)
|
||||||
await setDefaultVideoChannel(servers)
|
await setDefaultVideoChannel(servers)
|
||||||
|
|
||||||
|
|
|
@ -196,11 +196,17 @@ const contextStore: { [ id in ContextType ]: (string | { [ id: string ]: string
|
||||||
uuid: 'sc:identifier'
|
uuid: 'sc:identifier'
|
||||||
}),
|
}),
|
||||||
|
|
||||||
|
View: buildContext({
|
||||||
|
WatchAction: 'sc:WatchAction',
|
||||||
|
InteractionCounter: 'sc:InteractionCounter',
|
||||||
|
interactionType: 'sc:interactionType',
|
||||||
|
userInteractionCount: 'sc:userInteractionCount'
|
||||||
|
}),
|
||||||
|
|
||||||
Collection: buildContext(),
|
Collection: buildContext(),
|
||||||
Follow: buildContext(),
|
Follow: buildContext(),
|
||||||
Reject: buildContext(),
|
Reject: buildContext(),
|
||||||
Accept: buildContext(),
|
Accept: buildContext(),
|
||||||
View: buildContext(),
|
|
||||||
Announce: buildContext(),
|
Announce: buildContext(),
|
||||||
Comment: buildContext(),
|
Comment: buildContext(),
|
||||||
Delete: buildContext(),
|
Delete: buildContext(),
|
||||||
|
|
|
@ -9,7 +9,6 @@ export function Debounce (config: { timeoutMS: number }) {
|
||||||
|
|
||||||
timeoutRef = setTimeout(() => {
|
timeoutRef = setTimeout(() => {
|
||||||
original.apply(this, args)
|
original.apply(this, args)
|
||||||
|
|
||||||
}, config.timeoutMS)
|
}, config.timeoutMS)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
@ -28,11 +28,15 @@ async function processCreateView (activity: ActivityView, byActor: MActorSignatu
|
||||||
allowRefresh: false
|
allowRefresh: false
|
||||||
})
|
})
|
||||||
|
|
||||||
const viewerExpires = activity.expires
|
await VideoViewsManager.Instance.processRemoteView({
|
||||||
? new Date(activity.expires)
|
video,
|
||||||
: undefined
|
viewerId: activity.id,
|
||||||
|
|
||||||
await VideoViewsManager.Instance.processRemoteView({ video, viewerId: activity.id, viewerExpires })
|
viewerExpires: activity.expires
|
||||||
|
? new Date(activity.expires)
|
||||||
|
: undefined,
|
||||||
|
viewerResultCounter: getViewerResultCounter(activity)
|
||||||
|
})
|
||||||
|
|
||||||
if (video.isOwned()) {
|
if (video.isOwned()) {
|
||||||
// Forward the view but don't resend the activity to the sender
|
// Forward the view but don't resend the activity to the sender
|
||||||
|
@ -40,3 +44,15 @@ async function processCreateView (activity: ActivityView, byActor: MActorSignatu
|
||||||
await forwardVideoRelatedActivity(activity, undefined, exceptions, video)
|
await forwardVideoRelatedActivity(activity, undefined, exceptions, video)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Viewer protocol V2
|
||||||
|
function getViewerResultCounter (activity: ActivityView) {
|
||||||
|
const result = activity.result
|
||||||
|
|
||||||
|
if (!activity.expires || result?.interactionType !== 'WatchAction' || result?.type !== 'InteractionCounter') return undefined
|
||||||
|
|
||||||
|
const counter = parseInt(result.userInteractionCount + '')
|
||||||
|
if (isNaN(counter)) return undefined
|
||||||
|
|
||||||
|
return counter
|
||||||
|
}
|
||||||
|
|
|
@ -6,24 +6,23 @@ import { logger } from '../../../helpers/logger.js'
|
||||||
import { audiencify, getAudience } from '../audience.js'
|
import { audiencify, getAudience } from '../audience.js'
|
||||||
import { getLocalVideoViewActivityPubUrl } from '../url.js'
|
import { getLocalVideoViewActivityPubUrl } from '../url.js'
|
||||||
import { sendVideoRelatedActivity } from './shared/send-utils.js'
|
import { sendVideoRelatedActivity } from './shared/send-utils.js'
|
||||||
|
import { isUsingViewersFederationV2 } from '@peertube/peertube-node-utils'
|
||||||
type ViewType = 'view' | 'viewer'
|
|
||||||
|
|
||||||
async function sendView (options: {
|
async function sendView (options: {
|
||||||
byActor: MActorLight
|
byActor: MActorLight
|
||||||
type: ViewType
|
|
||||||
video: MVideoImmutable
|
video: MVideoImmutable
|
||||||
viewerIdentifier: string
|
viewerIdentifier: string
|
||||||
|
viewersCount?: number
|
||||||
transaction?: Transaction
|
transaction?: Transaction
|
||||||
}) {
|
}) {
|
||||||
const { byActor, type, video, viewerIdentifier, transaction } = options
|
const { byActor, viewersCount, video, viewerIdentifier, transaction } = options
|
||||||
|
|
||||||
logger.info('Creating job to send %s of %s.', type, video.url)
|
logger.info('Creating job to send %s of %s.', viewersCount !== undefined ? 'viewer' : 'view', video.url)
|
||||||
|
|
||||||
const activityBuilder = (audience: ActivityAudience) => {
|
const activityBuilder = (audience: ActivityAudience) => {
|
||||||
const url = getLocalVideoViewActivityPubUrl(byActor, video, viewerIdentifier)
|
const url = getLocalVideoViewActivityPubUrl(byActor, video, viewerIdentifier)
|
||||||
|
|
||||||
return buildViewActivity({ url, byActor, video, audience, type })
|
return buildViewActivity({ url, byActor, video, audience, viewersCount })
|
||||||
}
|
}
|
||||||
|
|
||||||
return sendVideoRelatedActivity(activityBuilder, { byActor, video, transaction, contextType: 'View', parallelizable: true })
|
return sendVideoRelatedActivity(activityBuilder, { byActor, video, transaction, contextType: 'View', parallelizable: true })
|
||||||
|
@ -41,22 +40,33 @@ function buildViewActivity (options: {
|
||||||
url: string
|
url: string
|
||||||
byActor: MActorAudience
|
byActor: MActorAudience
|
||||||
video: MVideoUrl
|
video: MVideoUrl
|
||||||
type: ViewType
|
viewersCount?: number
|
||||||
audience?: ActivityAudience
|
audience?: ActivityAudience
|
||||||
}): ActivityView {
|
}): ActivityView {
|
||||||
const { url, byActor, type, video, audience = getAudience(byActor) } = options
|
const { url, byActor, viewersCount, video, audience = getAudience(byActor) } = options
|
||||||
|
|
||||||
return audiencify(
|
const base = {
|
||||||
{
|
|
||||||
id: url,
|
id: url,
|
||||||
type: 'View' as 'View',
|
type: 'View' as 'View',
|
||||||
actor: byActor.url,
|
actor: byActor.url,
|
||||||
object: video.url,
|
object: video.url
|
||||||
|
}
|
||||||
|
|
||||||
expires: type === 'viewer'
|
if (viewersCount === undefined) {
|
||||||
? new Date(VideoViewsManager.Instance.buildViewerExpireTime()).toISOString()
|
return audiencify(base, audience)
|
||||||
|
}
|
||||||
|
|
||||||
|
return audiencify({
|
||||||
|
...base,
|
||||||
|
|
||||||
|
expires: new Date(VideoViewsManager.Instance.buildViewerExpireTime()).toISOString(),
|
||||||
|
|
||||||
|
result: isUsingViewersFederationV2()
|
||||||
|
? {
|
||||||
|
interactionType: 'WatchAction',
|
||||||
|
type: 'InteractionCounter',
|
||||||
|
userInteractionCount: viewersCount
|
||||||
|
}
|
||||||
: undefined
|
: undefined
|
||||||
},
|
}, audience)
|
||||||
audience
|
|
||||||
)
|
|
||||||
}
|
}
|
||||||
|
|
|
@ -1,4 +1,5 @@
|
||||||
import { buildUUID, isTestOrDevInstance, sha256 } from '@peertube/peertube-node-utils'
|
import { buildUUID, isTestOrDevInstance, isUsingViewersFederationV2, sha256 } from '@peertube/peertube-node-utils'
|
||||||
|
import { exists } from '@server/helpers/custom-validators/misc.js'
|
||||||
import { logger, loggerTagsFactory } from '@server/helpers/logger.js'
|
import { logger, loggerTagsFactory } from '@server/helpers/logger.js'
|
||||||
import { VIEW_LIFETIME } from '@server/initializers/constants.js'
|
import { VIEW_LIFETIME } from '@server/initializers/constants.js'
|
||||||
import { sendView } from '@server/lib/activitypub/send/send-view.js'
|
import { sendView } from '@server/lib/activitypub/send/send-view.js'
|
||||||
|
@ -17,6 +18,7 @@ type Viewer = {
|
||||||
id: string
|
id: string
|
||||||
viewerScope: ViewerScope
|
viewerScope: ViewerScope
|
||||||
videoScope: VideoScope
|
videoScope: VideoScope
|
||||||
|
viewerCount: number
|
||||||
lastFederation?: number
|
lastFederation?: number
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@ -54,22 +56,48 @@ export class VideoViewerCounters {
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
||||||
const newViewer = await this.addViewerToVideo({ viewerId, video, viewerScope: 'local' })
|
const newViewer = this.addViewerToVideo({ viewerId, video, viewerScope: 'local', viewerCount: 1 })
|
||||||
await this.federateViewerIfNeeded(video, newViewer)
|
await this.federateViewerIfNeeded(video, newViewer)
|
||||||
|
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
|
|
||||||
async addRemoteViewer (options: {
|
addRemoteViewerOnLocalVideo (options: {
|
||||||
video: MVideo
|
video: MVideo
|
||||||
viewerId: string
|
viewerId: string
|
||||||
viewerExpires: Date
|
viewerExpires: Date
|
||||||
}) {
|
}) {
|
||||||
const { video, viewerExpires, viewerId } = options
|
const { video, viewerExpires, viewerId } = options
|
||||||
|
|
||||||
logger.debug('Adding remote viewer to video %s.', video.uuid, { ...lTags(video.uuid) })
|
logger.debug('Adding remote viewer to local video %s.', video.uuid, { viewerId, viewerExpires, ...lTags(video.uuid) })
|
||||||
|
|
||||||
await this.addViewerToVideo({ video, viewerExpires, viewerId, viewerScope: 'remote' })
|
this.addViewerToVideo({ video, viewerExpires, viewerId, viewerScope: 'remote', viewerCount: 1 })
|
||||||
|
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
addRemoteViewerOnRemoteVideo (options: {
|
||||||
|
video: MVideo
|
||||||
|
viewerId: string
|
||||||
|
viewerExpires: Date
|
||||||
|
viewerResultCounter?: number
|
||||||
|
}) {
|
||||||
|
const { video, viewerExpires, viewerId, viewerResultCounter } = options
|
||||||
|
|
||||||
|
logger.debug(
|
||||||
|
'Adding remote viewer to remote video %s.', video.uuid,
|
||||||
|
{ viewerId, viewerResultCounter, viewerExpires, ...lTags(video.uuid) }
|
||||||
|
)
|
||||||
|
|
||||||
|
this.addViewerToVideo({
|
||||||
|
video,
|
||||||
|
viewerExpires,
|
||||||
|
viewerId,
|
||||||
|
viewerScope: 'remote',
|
||||||
|
// The origin server sends a summary of all viewers, so we can replace our local copy
|
||||||
|
replaceCurrentViewers: exists(viewerResultCounter),
|
||||||
|
viewerCount: viewerResultCounter ?? 1
|
||||||
|
})
|
||||||
|
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
|
@ -83,17 +111,17 @@ export class VideoViewerCounters {
|
||||||
let total = 0
|
let total = 0
|
||||||
|
|
||||||
for (const viewers of this.viewersPerVideo.values()) {
|
for (const viewers of this.viewersPerVideo.values()) {
|
||||||
total += viewers.filter(v => v.viewerScope === options.viewerScope && v.videoScope === options.videoScope).length
|
total += viewers.filter(v => v.viewerScope === options.viewerScope && v.videoScope === options.videoScope)
|
||||||
|
.reduce((p, c) => p + c.viewerCount, 0)
|
||||||
}
|
}
|
||||||
|
|
||||||
return total
|
return total
|
||||||
}
|
}
|
||||||
|
|
||||||
getViewers (video: MVideo) {
|
getTotalViewersOf (video: MVideoImmutable) {
|
||||||
const viewers = this.viewersPerVideo.get(video.id)
|
const viewers = this.viewersPerVideo.get(video.id)
|
||||||
if (!viewers) return 0
|
|
||||||
|
|
||||||
return viewers.length
|
return viewers?.reduce((p, c) => p + c.viewerCount, 0) || 0
|
||||||
}
|
}
|
||||||
|
|
||||||
buildViewerExpireTime () {
|
buildViewerExpireTime () {
|
||||||
|
@ -102,17 +130,19 @@ export class VideoViewerCounters {
|
||||||
|
|
||||||
// ---------------------------------------------------------------------------
|
// ---------------------------------------------------------------------------
|
||||||
|
|
||||||
private async addViewerToVideo (options: {
|
private addViewerToVideo (options: {
|
||||||
video: MVideoImmutable
|
video: MVideoImmutable
|
||||||
viewerId: string
|
viewerId: string
|
||||||
viewerScope: ViewerScope
|
viewerScope: ViewerScope
|
||||||
|
viewerCount: number
|
||||||
|
replaceCurrentViewers?: boolean
|
||||||
viewerExpires?: Date
|
viewerExpires?: Date
|
||||||
}) {
|
}) {
|
||||||
const { video, viewerExpires, viewerId, viewerScope } = options
|
const { video, viewerExpires, viewerId, viewerScope, viewerCount, replaceCurrentViewers } = options
|
||||||
|
|
||||||
let watchers = this.viewersPerVideo.get(video.id)
|
let watchers = this.viewersPerVideo.get(video.id)
|
||||||
|
|
||||||
if (!watchers) {
|
if (!watchers || replaceCurrentViewers) {
|
||||||
watchers = []
|
watchers = []
|
||||||
this.viewersPerVideo.set(video.id, watchers)
|
this.viewersPerVideo.set(video.id, watchers)
|
||||||
}
|
}
|
||||||
|
@ -125,12 +155,12 @@ export class VideoViewerCounters {
|
||||||
? 'remote'
|
? 'remote'
|
||||||
: 'local'
|
: 'local'
|
||||||
|
|
||||||
const viewer = { id: viewerId, expires, videoScope, viewerScope }
|
const viewer = { id: viewerId, expires, videoScope, viewerScope, viewerCount }
|
||||||
watchers.push(viewer)
|
watchers.push(viewer)
|
||||||
|
|
||||||
this.idToViewer.set(viewerId, viewer)
|
this.idToViewer.set(viewerId, viewer)
|
||||||
|
|
||||||
await this.notifyClients(video.id, watchers.length)
|
this.notifyClients(video)
|
||||||
|
|
||||||
return viewer
|
return viewer
|
||||||
}
|
}
|
||||||
|
@ -162,7 +192,16 @@ export class VideoViewerCounters {
|
||||||
if (newViewers.length === 0) this.viewersPerVideo.delete(videoId)
|
if (newViewers.length === 0) this.viewersPerVideo.delete(videoId)
|
||||||
else this.viewersPerVideo.set(videoId, newViewers)
|
else this.viewersPerVideo.set(videoId, newViewers)
|
||||||
|
|
||||||
await this.notifyClients(videoId, newViewers.length)
|
const video = await VideoModel.loadImmutableAttributes(videoId)
|
||||||
|
|
||||||
|
if (video) {
|
||||||
|
this.notifyClients(video)
|
||||||
|
|
||||||
|
// Let total viewers expire on remote instances if there are no more viewers
|
||||||
|
if (video.remote === false && newViewers.length !== 0) {
|
||||||
|
await this.federateTotalViewers(video)
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
logger.error('Error in video clean viewers scheduler.', { err, ...lTags() })
|
logger.error('Error in video clean viewers scheduler.', { err, ...lTags() })
|
||||||
|
@ -171,13 +210,11 @@ export class VideoViewerCounters {
|
||||||
this.processingViewerCounters = false
|
this.processingViewerCounters = false
|
||||||
}
|
}
|
||||||
|
|
||||||
private async notifyClients (videoId: string | number, viewersLength: number) {
|
private notifyClients (video: MVideoImmutable) {
|
||||||
const video = await VideoModel.loadImmutableAttributes(videoId)
|
const totalViewers = this.getTotalViewersOf(video)
|
||||||
if (!video) return
|
PeerTubeSocket.Instance.sendVideoViewsUpdate(video, totalViewers)
|
||||||
|
|
||||||
PeerTubeSocket.Instance.sendVideoViewsUpdate(video, viewersLength)
|
logger.debug('Video viewers update for %s is %d.', video.url, totalViewers, lTags())
|
||||||
|
|
||||||
logger.debug('Video viewers update for %s is %d.', video.url, viewersLength, lTags())
|
|
||||||
}
|
}
|
||||||
|
|
||||||
private generateViewerId (ip: string, videoUUID: string) {
|
private generateViewerId (ip: string, videoUUID: string) {
|
||||||
|
@ -190,8 +227,26 @@ export class VideoViewerCounters {
|
||||||
const federationLimit = now - (VIEW_LIFETIME.VIEWER_COUNTER * 0.75)
|
const federationLimit = now - (VIEW_LIFETIME.VIEWER_COUNTER * 0.75)
|
||||||
|
|
||||||
if (viewer.lastFederation && viewer.lastFederation > federationLimit) return
|
if (viewer.lastFederation && viewer.lastFederation > federationLimit) return
|
||||||
|
if (video.remote === false && isUsingViewersFederationV2()) return
|
||||||
|
|
||||||
|
await sendView({
|
||||||
|
byActor: await getServerActor(),
|
||||||
|
video,
|
||||||
|
viewersCount: 1,
|
||||||
|
viewerIdentifier: viewer.id
|
||||||
|
})
|
||||||
|
|
||||||
await sendView({ byActor: await getServerActor(), video, type: 'viewer', viewerIdentifier: viewer.id })
|
|
||||||
viewer.lastFederation = now
|
viewer.lastFederation = now
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private async federateTotalViewers (video: MVideoImmutable) {
|
||||||
|
if (!isUsingViewersFederationV2()) return
|
||||||
|
|
||||||
|
await sendView({
|
||||||
|
byActor: await getServerActor(),
|
||||||
|
video,
|
||||||
|
viewersCount: this.getTotalViewersOf(video),
|
||||||
|
viewerIdentifier: video.uuid
|
||||||
|
})
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
@ -35,7 +35,7 @@ export class VideoViews {
|
||||||
|
|
||||||
await this.addView(video)
|
await this.addView(video)
|
||||||
|
|
||||||
await sendView({ byActor: await getServerActor(), video, type: 'view', viewerIdentifier: buildUUID() })
|
await sendView({ byActor: await getServerActor(), video, viewerIdentifier: buildUUID() })
|
||||||
|
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
|
|
|
@ -66,17 +66,29 @@ export class VideoViewsManager {
|
||||||
video: MVideo
|
video: MVideo
|
||||||
viewerId: string | null
|
viewerId: string | null
|
||||||
viewerExpires?: Date
|
viewerExpires?: Date
|
||||||
|
viewerResultCounter?: number
|
||||||
}) {
|
}) {
|
||||||
const { video, viewerId, viewerExpires } = options
|
const { video, viewerId, viewerExpires, viewerResultCounter } = options
|
||||||
|
|
||||||
logger.debug('Processing remote view for %s.', video.url, { viewerExpires, viewerId, ...lTags() })
|
logger.debug('Processing remote view for %s.', video.url, { viewerExpires, viewerId, ...lTags() })
|
||||||
|
|
||||||
if (viewerExpires) await this.videoViewerCounters.addRemoteViewer({ video, viewerId, viewerExpires })
|
// Viewer
|
||||||
else await this.videoViews.addRemoteView({ video })
|
if (viewerExpires) {
|
||||||
|
if (video.remote === false) {
|
||||||
|
this.videoViewerCounters.addRemoteViewerOnLocalVideo({ video, viewerId, viewerExpires })
|
||||||
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
getViewers (video: MVideo) {
|
this.videoViewerCounters.addRemoteViewerOnRemoteVideo({ video, viewerId, viewerExpires, viewerResultCounter })
|
||||||
return this.videoViewerCounters.getViewers(video)
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
// Just a view
|
||||||
|
await this.videoViews.addRemoteView({ video })
|
||||||
|
}
|
||||||
|
|
||||||
|
getTotalViewersOf (video: MVideo) {
|
||||||
|
return this.videoViewerCounters.getTotalViewersOf(video)
|
||||||
}
|
}
|
||||||
|
|
||||||
getTotalViewers (options: {
|
getTotalViewers (options: {
|
||||||
|
|
|
@ -90,7 +90,7 @@ export function videoModelToFormattedJSON (video: MVideoFormattable, options: Vi
|
||||||
duration: video.duration,
|
duration: video.duration,
|
||||||
|
|
||||||
views: video.views,
|
views: video.views,
|
||||||
viewers: VideoViewsManager.Instance.getViewers(video),
|
viewers: VideoViewsManager.Instance.getTotalViewersOf(video),
|
||||||
|
|
||||||
likes: video.likes,
|
likes: video.likes,
|
||||||
dislikes: video.dislikes,
|
dislikes: video.dislikes,
|
||||||
|
|
Loading…
Reference in a new issue