diff options
Diffstat (limited to 'server/lib/job-queue/handlers')
-rw-r--r-- | server/lib/job-queue/handlers/activitypub-http-broadcast.ts | 54 |
1 files changed, 29 insertions, 25 deletions
diff --git a/server/lib/job-queue/handlers/activitypub-http-broadcast.ts b/server/lib/job-queue/handlers/activitypub-http-broadcast.ts index 13eff5211..733c1378a 100644 --- a/server/lib/job-queue/handlers/activitypub-http-broadcast.ts +++ b/server/lib/job-queue/handlers/activitypub-http-broadcast.ts | |||
@@ -1,39 +1,28 @@ | |||
1 | import { map } from 'bluebird' | ||
2 | import { Job } from 'bullmq' | 1 | import { Job } from 'bullmq' |
3 | import { buildGlobalHeaders, buildSignedRequestOptions, computeBody } from '@server/lib/activitypub/send' | 2 | import { buildGlobalHeaders, buildSignedRequestOptions, computeBody } from '@server/lib/activitypub/send' |
4 | import { ActorFollowHealthCache } from '@server/lib/actor-follow-health-cache' | 3 | import { ActorFollowHealthCache } from '@server/lib/actor-follow-health-cache' |
4 | import { sequentialHTTPBroadcastFromWorker } from '@server/lib/worker/parent-process' | ||
5 | import { ActivitypubHttpBroadcastPayload } from '@shared/models' | 5 | import { ActivitypubHttpBroadcastPayload } from '@shared/models' |
6 | import { logger } from '../../../helpers/logger' | 6 | import { logger } from '../../../helpers/logger' |
7 | import { doRequest } from '../../../helpers/requests' | ||
8 | import { BROADCAST_CONCURRENCY } from '../../../initializers/constants' | ||
9 | 7 | ||
10 | async function processActivityPubHttpBroadcast (job: Job) { | 8 | // Prefer using a worker thread for HTTP requests because on high load we may have to sign many requests, which can be CPU intensive |
9 | |||
10 | async function processActivityPubHttpSequentialBroadcast (job: Job<ActivitypubHttpBroadcastPayload>) { | ||
11 | logger.info('Processing ActivityPub broadcast in job %s.', job.id) | 11 | logger.info('Processing ActivityPub broadcast in job %s.', job.id) |
12 | 12 | ||
13 | const payload = job.data as ActivitypubHttpBroadcastPayload | 13 | const requestOptions = await buildRequestOptions(job.data) |
14 | 14 | ||
15 | const body = await computeBody(payload) | 15 | const { badUrls, goodUrls } = await sequentialHTTPBroadcastFromWorker({ uris: job.data.uris, requestOptions }) |
16 | const httpSignatureOptions = await buildSignedRequestOptions(payload) | ||
17 | 16 | ||
18 | const options = { | 17 | return ActorFollowHealthCache.Instance.updateActorFollowsHealth(goodUrls, badUrls) |
19 | method: 'POST' as 'POST', | 18 | } |
20 | json: body, | 19 | |
21 | httpSignature: httpSignatureOptions, | 20 | async function processActivityPubParallelHttpBroadcast (job: Job<ActivitypubHttpBroadcastPayload>) { |
22 | headers: buildGlobalHeaders(body) | 21 | logger.info('Processing ActivityPub parallel broadcast in job %s.', job.id) |
23 | } | ||
24 | 22 | ||
25 | const badUrls: string[] = [] | 23 | const requestOptions = await buildRequestOptions(job.data) |
26 | const goodUrls: string[] = [] | ||
27 | 24 | ||
28 | await map(payload.uris, async uri => { | 25 | const { badUrls, goodUrls } = await sequentialHTTPBroadcastFromWorker({ uris: job.data.uris, requestOptions }) |
29 | try { | ||
30 | await doRequest(uri, options) | ||
31 | goodUrls.push(uri) | ||
32 | } catch (err) { | ||
33 | logger.debug('HTTP broadcast to %s failed.', uri, { err }) | ||
34 | badUrls.push(uri) | ||
35 | } | ||
36 | }, { concurrency: BROADCAST_CONCURRENCY }) | ||
37 | 26 | ||
38 | return ActorFollowHealthCache.Instance.updateActorFollowsHealth(goodUrls, badUrls) | 27 | return ActorFollowHealthCache.Instance.updateActorFollowsHealth(goodUrls, badUrls) |
39 | } | 28 | } |
@@ -41,5 +30,20 @@ async function processActivityPubHttpBroadcast (job: Job) { | |||
41 | // --------------------------------------------------------------------------- | 30 | // --------------------------------------------------------------------------- |
42 | 31 | ||
43 | export { | 32 | export { |
44 | processActivityPubHttpBroadcast | 33 | processActivityPubHttpSequentialBroadcast, |
34 | processActivityPubParallelHttpBroadcast | ||
35 | } | ||
36 | |||
37 | // --------------------------------------------------------------------------- | ||
38 | |||
39 | async function buildRequestOptions (payload: ActivitypubHttpBroadcastPayload) { | ||
40 | const body = await computeBody(payload) | ||
41 | const httpSignatureOptions = await buildSignedRequestOptions(payload) | ||
42 | |||
43 | return { | ||
44 | method: 'POST' as 'POST', | ||
45 | json: body, | ||
46 | httpSignature: httpSignatureOptions, | ||
47 | headers: buildGlobalHeaders(body) | ||
48 | } | ||
45 | } | 49 | } |