2023-05-10 15:05:08 +09:00
|
|
|
import { setTimeout } from 'node:timers/promises';
|
|
|
|
import { Inject, Module, OnApplicationShutdown } from '@nestjs/common';
|
2022-09-18 03:27:08 +09:00
|
|
|
import Bull from 'bull';
|
|
|
|
import { DI } from '@/di-symbols.js';
|
|
|
|
import type { Config } from '@/config.js';
|
|
|
|
import type { Provider } from '@nestjs/common';
|
2023-04-12 09:13:58 +09:00
|
|
|
import type { DeliverJobData, InboxJobData, DbJobData, ObjectStorageJobData, EndedPollNotificationJobData, WebhookDeliverJobData, RelationshipJobData, DbJobMap } from '../queue/types.js';
|
2022-09-18 03:27:08 +09:00
|
|
|
|
|
|
|
function q<T>(config: Config, name: string, limitPerSec = -1) {
|
|
|
|
return new Bull<T>(name, {
|
|
|
|
redis: {
|
2023-04-07 11:27:01 +09:00
|
|
|
port: config.redisForJobQueue.port,
|
|
|
|
host: config.redisForJobQueue.host,
|
|
|
|
family: config.redisForJobQueue.family == null ? 0 : config.redisForJobQueue.family,
|
|
|
|
password: config.redisForJobQueue.pass,
|
|
|
|
db: config.redisForJobQueue.db ?? 0,
|
2022-09-18 03:27:08 +09:00
|
|
|
},
|
2023-04-07 11:27:01 +09:00
|
|
|
prefix: config.redisForJobQueue.prefix ? `${config.redisForJobQueue.prefix}:queue` : 'queue',
|
2022-09-18 03:27:08 +09:00
|
|
|
limiter: limitPerSec > 0 ? {
|
|
|
|
max: limitPerSec,
|
|
|
|
duration: 1000,
|
|
|
|
} : undefined,
|
|
|
|
settings: {
|
|
|
|
backoffStrategies: {
|
|
|
|
apBackoff,
|
|
|
|
},
|
|
|
|
},
|
|
|
|
});
|
|
|
|
}
|
|
|
|
|
|
|
|
// ref. https://github.com/misskey-dev/misskey/pull/7635#issue-971097019
|
|
|
|
function apBackoff(attemptsMade: number, err: Error) {
|
|
|
|
const baseDelay = 60 * 1000; // 1min
|
|
|
|
const maxBackoff = 8 * 60 * 60 * 1000; // 8hours
|
|
|
|
let backoff = (Math.pow(2, attemptsMade) - 1) * baseDelay;
|
|
|
|
backoff = Math.min(backoff, maxBackoff);
|
|
|
|
backoff += Math.round(backoff * Math.random() * 0.2);
|
|
|
|
return backoff;
|
|
|
|
}
|
|
|
|
|
|
|
|
export type SystemQueue = Bull.Queue<Record<string, unknown>>;
|
|
|
|
export type EndedPollNotificationQueue = Bull.Queue<EndedPollNotificationJobData>;
|
|
|
|
export type DeliverQueue = Bull.Queue<DeliverJobData>;
|
|
|
|
export type InboxQueue = Bull.Queue<InboxJobData>;
|
2023-05-10 15:05:08 +09:00
|
|
|
export type DbQueue = Bull.Queue;
|
2023-04-12 09:13:58 +09:00
|
|
|
export type RelationshipQueue = Bull.Queue<RelationshipJobData>;
|
2023-05-10 15:05:08 +09:00
|
|
|
export type ObjectStorageQueue = Bull.Queue;
|
2022-09-18 03:27:08 +09:00
|
|
|
export type WebhookDeliverQueue = Bull.Queue<WebhookDeliverJobData>;
|
|
|
|
|
|
|
|
const $system: Provider = {
|
|
|
|
provide: 'queue:system',
|
|
|
|
useFactory: (config: Config) => q(config, 'system'),
|
|
|
|
inject: [DI.config],
|
|
|
|
};
|
|
|
|
|
|
|
|
const $endedPollNotification: Provider = {
|
|
|
|
provide: 'queue:endedPollNotification',
|
|
|
|
useFactory: (config: Config) => q(config, 'endedPollNotification'),
|
|
|
|
inject: [DI.config],
|
|
|
|
};
|
|
|
|
|
|
|
|
const $deliver: Provider = {
|
|
|
|
provide: 'queue:deliver',
|
|
|
|
useFactory: (config: Config) => q(config, 'deliver', config.deliverJobPerSec ?? 128),
|
|
|
|
inject: [DI.config],
|
|
|
|
};
|
|
|
|
|
|
|
|
const $inbox: Provider = {
|
|
|
|
provide: 'queue:inbox',
|
|
|
|
useFactory: (config: Config) => q(config, 'inbox', config.inboxJobPerSec ?? 16),
|
|
|
|
inject: [DI.config],
|
|
|
|
};
|
|
|
|
|
|
|
|
const $db: Provider = {
|
|
|
|
provide: 'queue:db',
|
|
|
|
useFactory: (config: Config) => q(config, 'db'),
|
|
|
|
inject: [DI.config],
|
|
|
|
};
|
|
|
|
|
2023-04-12 09:13:58 +09:00
|
|
|
const $relationship: Provider = {
|
|
|
|
provide: 'queue:relationship',
|
enhance: account migration (#10592)
* copy block and mute then create follow and unfollow jobs
* copy block and mute and update lists when detecting an account has moved
* no need to care promise orders
* refactor updating actor and target
* automatically accept if a locked account had accepted an old account
* fix exception format
* prevent the old account from calling some endpoints
* do not unfollow when moving
* adjust following and follower counts
* check movedToUri when receiving a follow request
* skip if no need to adjust
* Revert "disable account migration"
This reverts commit 2321214c98591bcfe1385c1ab5bf0ff7b471ae1d.
* fix translation specifier
* fix checking alsoKnownAs and uri
* fix updating account
* fix refollowing locked account
* decrease followersCount if followed by the old account
* adjust following and followers counts when unfollowing
* fix copying mutings
* prohibit moved account from moving again
* fix move service
* allow app creation after moving
* fix lint
* remove unnecessary field
* fix cache update
* add e2e test
* add e2e test of accepting the new account automatically
* force follow if any error happens
* remove unnecessary joins
* use Array.map instead of for const of
* ユーザーリストの移行は追加のみを行う
* nanka iroiro
* fix misskey-js?
* :v:
* 移行を行ったアカウントからのフォローリクエストの自動許可を調整
* newUriを外に出す
* newUriを外に出す2
* clean up
* fix newUri
* prevent moving if the destination account has already moved
* set alsoKnownAs via /i/update
* fix database initialization
* add return type
* prohibit updating alsoKnownAs after moving
* skip to add to alsoKnownAs if toUrl is known
* skip adding to the list if it already has
* use Acct.parse instead
* rename error code
* :art:
* 制限を5から10に緩和
* movedTo(Uri), alsoKnownAsはユーザーidを返すように
* test api res
* fix
* 元アカウントはミュートし続ける
* :art:
* unfollow
* fix
* getUserUriをUserEntityServiceに
* ?
* job!
* :art:
* instance => server
* accountMovedShort, forbiddenBecauseYouAreMigrated
* accountMovedShort
* fix test
* import, pin禁止
* 実績を凍結する
* clean up
* :v:
* change message
* ブロック, フォロー, ミュート, リストのインポートファイルの制限を32MiBに
* Revert "ブロック, フォロー, ミュート, リストのインポートファイルの制限を32MiBに"
This reverts commit 3bd7be35d8aa455cb01ae58f8172a71a50485db1.
* validateAlsoKnownAs
* 移行後2時間以内はインポート可能なファイルサイズを拡大
* clean up
* どうせactorをupdatePersonで更新するならupdatePersonしか移行処理を発行しないことにする
* handle error?
* リモートからの移行処理の条件を是正
* log, port
* fix
* fix
* enhance(dev): non-production環境でhttpサーバー間でもユーザー、ノートの連合が可能なように
* refactor (use checkHttps)
* MISSKEY_WEBFINGER_USE_HTTP
* Environment Variable readme
* NEVER USE IN PRODUCTION
* fix punyHost
* fix indent
* fix
* experimental
---------
Co-authored-by: tamaina <tamaina@hotmail.co.jp>
Co-authored-by: syuilo <Syuilotan@yahoo.co.jp>
2023-04-30 00:09:29 +09:00
|
|
|
useFactory: (config: Config) => q(config, 'relationship', config.relashionshipJobPerSec ?? 64),
|
2023-04-12 09:13:58 +09:00
|
|
|
inject: [DI.config],
|
|
|
|
};
|
|
|
|
|
2022-09-18 03:27:08 +09:00
|
|
|
const $objectStorage: Provider = {
|
|
|
|
provide: 'queue:objectStorage',
|
|
|
|
useFactory: (config: Config) => q(config, 'objectStorage'),
|
|
|
|
inject: [DI.config],
|
|
|
|
};
|
|
|
|
|
|
|
|
const $webhookDeliver: Provider = {
|
|
|
|
provide: 'queue:webhookDeliver',
|
|
|
|
useFactory: (config: Config) => q(config, 'webhookDeliver', 64),
|
|
|
|
inject: [DI.config],
|
|
|
|
};
|
|
|
|
|
|
|
|
@Module({
|
|
|
|
imports: [
|
|
|
|
],
|
|
|
|
providers: [
|
|
|
|
$system,
|
|
|
|
$endedPollNotification,
|
|
|
|
$deliver,
|
|
|
|
$inbox,
|
|
|
|
$db,
|
2023-04-12 09:13:58 +09:00
|
|
|
$relationship,
|
2022-09-18 03:27:08 +09:00
|
|
|
$objectStorage,
|
|
|
|
$webhookDeliver,
|
|
|
|
],
|
|
|
|
exports: [
|
|
|
|
$system,
|
|
|
|
$endedPollNotification,
|
|
|
|
$deliver,
|
|
|
|
$inbox,
|
|
|
|
$db,
|
2023-04-12 09:13:58 +09:00
|
|
|
$relationship,
|
2022-09-18 03:27:08 +09:00
|
|
|
$objectStorage,
|
|
|
|
$webhookDeliver,
|
|
|
|
],
|
|
|
|
})
|
2023-05-10 15:05:08 +09:00
|
|
|
export class QueueModule implements OnApplicationShutdown {
|
|
|
|
constructor(
|
|
|
|
@Inject('queue:system') public systemQueue: SystemQueue,
|
|
|
|
@Inject('queue:endedPollNotification') public endedPollNotificationQueue: EndedPollNotificationQueue,
|
|
|
|
@Inject('queue:deliver') public deliverQueue: DeliverQueue,
|
|
|
|
@Inject('queue:inbox') public inboxQueue: InboxQueue,
|
|
|
|
@Inject('queue:db') public dbQueue: DbQueue,
|
|
|
|
@Inject('queue:relationship') public relationshipQueue: RelationshipQueue,
|
|
|
|
@Inject('queue:objectStorage') public objectStorageQueue: ObjectStorageQueue,
|
|
|
|
@Inject('queue:webhookDeliver') public webhookDeliverQueue: WebhookDeliverQueue,
|
|
|
|
) {}
|
|
|
|
|
|
|
|
async onApplicationShutdown(signal: string): Promise<void> {
|
|
|
|
if (process.env.NODE_ENV === 'test') {
|
|
|
|
// XXX:
|
|
|
|
// Shutting down the existing connections causes errors on Jest as
|
|
|
|
// Misskey has asynchronous postgres/redis connections that are not
|
|
|
|
// awaited.
|
|
|
|
// Let's wait for some random time for them to finish.
|
|
|
|
await setTimeout(5000);
|
|
|
|
}
|
|
|
|
await Promise.all([
|
|
|
|
this.systemQueue.close(),
|
|
|
|
this.endedPollNotificationQueue.close(),
|
|
|
|
this.deliverQueue.close(),
|
|
|
|
this.inboxQueue.close(),
|
|
|
|
this.dbQueue.close(),
|
|
|
|
this.relationshipQueue.close(),
|
|
|
|
this.objectStorageQueue.close(),
|
|
|
|
this.webhookDeliverQueue.close(),
|
|
|
|
]);
|
|
|
|
}
|
|
|
|
}
|