2020-08-11 18:24:59 +02:00
// @ts-check
2017-06-26 04:49:39 +02:00
const os = require ( 'os' ) ;
const throng = require ( 'throng' ) ;
const dotenv = require ( 'dotenv' ) ;
const express = require ( 'express' ) ;
const http = require ( 'http' ) ;
const redis = require ( 'redis' ) ;
const pg = require ( 'pg' ) ;
const log = require ( 'npmlog' ) ;
const url = require ( 'url' ) ;
const uuid = require ( 'uuid' ) ;
2018-08-24 18:16:53 +02:00
const fs = require ( 'fs' ) ;
2021-03-24 09:37:41 +01:00
const WebSocket = require ( 'ws' ) ;
2017-05-20 17:31:47 +02:00
const env = process . env . NODE _ENV || 'development' ;
2020-08-11 18:24:59 +02:00
const alwaysRequireAuth = process . env . LIMITED _FEDERATION _MODE === 'true' || process . env . WHITELIST _MODE === 'true' || process . env . AUTHORIZED _FETCH === 'true' ;
2017-02-02 16:11:36 +01:00
dotenv . config ( {
2017-05-20 17:31:47 +02:00
path : env === 'production' ? '.env.production' : '.env' ,
} ) ;
2017-02-02 01:31:09 +01:00
2017-05-28 16:25:26 +02:00
log . level = process . env . LOG _LEVEL || 'verbose' ;
2020-08-11 18:24:59 +02:00
/ * *
* @ param { string } dbUrl
* @ return { Object . < string , any > }
* /
2017-05-03 23:18:13 +02:00
const dbUrlToConfig = ( dbUrl ) => {
if ( ! dbUrl ) {
2017-05-20 17:31:47 +02:00
return { } ;
2017-05-03 23:18:13 +02:00
}
2019-03-10 16:00:54 +01:00
const params = url . parse ( dbUrl , true ) ;
2017-05-20 17:31:47 +02:00
const config = { } ;
2017-05-04 15:53:44 +02:00
if ( params . auth ) {
2017-05-20 17:31:47 +02:00
[ config . user , config . password ] = params . auth . split ( ':' ) ;
2017-05-04 15:53:44 +02:00
}
if ( params . hostname ) {
2017-05-20 17:31:47 +02:00
config . host = params . hostname ;
2017-05-04 15:53:44 +02:00
}
if ( params . port ) {
2017-05-20 17:31:47 +02:00
config . port = params . port ;
2017-05-03 23:18:13 +02:00
}
2017-05-04 15:53:44 +02:00
if ( params . pathname ) {
2017-05-20 17:31:47 +02:00
config . database = params . pathname . split ( '/' ) [ 1 ] ;
2017-05-04 15:53:44 +02:00
}
2017-05-20 17:31:47 +02:00
const ssl = params . query && params . query . ssl ;
2017-05-20 21:06:09 +02:00
2019-03-10 16:00:54 +01:00
if ( ssl && ssl === 'true' || ssl === '1' ) {
config . ssl = true ;
2017-05-04 15:53:44 +02:00
}
2017-05-20 17:31:47 +02:00
return config ;
} ;
2017-05-03 23:18:13 +02:00
2020-08-11 18:24:59 +02:00
/ * *
* @ param { Object . < string , any > } defaultConfig
* @ param { string } redisUrl
* /
2021-12-25 22:55:06 +01:00
const redisUrlToClient = async ( defaultConfig , redisUrl ) => {
2017-05-20 21:06:09 +02:00
const config = defaultConfig ;
2021-12-25 22:55:06 +01:00
let client ;
2017-05-20 21:06:09 +02:00
if ( ! redisUrl ) {
2021-12-25 22:55:06 +01:00
client = redis . createClient ( config ) ;
} else if ( redisUrl . startsWith ( 'unix://' ) ) {
client = redis . createClient ( Object . assign ( config , {
socket : {
path : redisUrl . slice ( 7 ) ,
} ,
} ) ) ;
} else {
client = redis . createClient ( Object . assign ( config , {
url : redisUrl ,
} ) ) ;
2017-05-20 21:06:09 +02:00
}
2021-12-25 22:55:06 +01:00
client . on ( 'error' , ( err ) => log . error ( 'Redis Client Error!' , err ) ) ;
await client . connect ( ) ;
2017-05-20 21:06:09 +02:00
2021-12-25 22:55:06 +01:00
return client ;
2017-05-20 21:06:09 +02:00
} ;
2017-05-28 16:25:26 +02:00
const numWorkers = + process . env . STREAMING _CLUSTER _NUM || ( env === 'development' ? 1 : Math . max ( os . cpus ( ) . length - 1 , 1 ) ) ;
2017-05-03 23:18:13 +02:00
2020-09-22 15:30:41 +02:00
/ * *
2023-06-09 19:29:16 +02:00
* Attempts to safely parse a string as JSON , used when both receiving a message
* from redis and when receiving a message from a client over a websocket
* connection , this is why it accepts a ` req ` argument .
2020-09-22 15:30:41 +02:00
* @ param { string } json
2023-06-09 19:29:16 +02:00
* @ param { any ? } req
* @ returns { Object . < string , any > | null }
2020-09-22 15:30:41 +02:00
* /
2022-02-16 14:37:26 +01:00
const parseJSON = ( json , req ) => {
2020-09-22 15:30:41 +02:00
try {
return JSON . parse ( json ) ;
} catch ( err ) {
2023-06-09 19:29:16 +02:00
/ * F I X M E : T h i s l o g g i n g i s n ' t g r e a t , a n d s h o u l d p r o b a b l y b e d o n e a t t h e
* call - site of parseJSON , not in the method , but this would require changing
* the signature of parseJSON to return something akin to a Result type :
* [ Error | null , null | Object < string , any } ] , and then handling the error
* scenarios .
* /
if ( req ) {
if ( req . accountId ) {
log . warn ( req . requestId , ` Error parsing message from user ${ req . accountId } : ${ err } ` ) ;
} else {
log . silly ( req . requestId , ` Error parsing message from ${ req . remoteAddress } : ${ err } ` ) ;
}
2022-02-16 14:37:26 +01:00
} else {
2023-06-09 19:29:16 +02:00
log . warn ( ` Error parsing message from redis: ${ err } ` ) ;
2022-02-16 14:37:26 +01:00
}
2020-09-22 15:30:41 +02:00
return null ;
}
} ;
2017-05-28 16:25:26 +02:00
const startMaster = ( ) => {
2018-08-24 18:16:53 +02:00
if ( ! process . env . SOCKET && process . env . PORT && isNaN ( + process . env . PORT ) ) {
log . warn ( 'UNIX domain socket is now supported by using SOCKET. Please migrate from PORT hack.' ) ;
}
2018-10-20 02:25:25 +02:00
2021-05-01 23:19:18 +02:00
log . warn ( ` Starting streaming API server master with ${ numWorkers } workers ` ) ;
2017-05-28 16:25:26 +02:00
} ;
2017-05-03 23:18:13 +02:00
2021-12-25 22:55:06 +01:00
const startWorker = async ( workerId ) => {
2021-05-01 23:19:18 +02:00
log . warn ( ` Starting worker ${ workerId } ` ) ;
2017-04-17 04:32:30 +02:00
const pgConfigs = {
development : {
2017-06-25 18:13:31 +02:00
user : process . env . DB _USER || pg . defaults . user ,
password : process . env . DB _PASS || pg . defaults . password ,
2017-10-17 11:45:37 +02:00
database : process . env . DB _NAME || 'mastodon_development' ,
2017-06-25 18:13:31 +02:00
host : process . env . DB _HOST || pg . defaults . host ,
port : process . env . DB _PORT || pg . defaults . port ,
2017-05-20 17:31:47 +02:00
max : 10 ,
2017-04-17 04:32:30 +02:00
} ,
production : {
user : process . env . DB _USER || 'mastodon' ,
password : process . env . DB _PASS || '' ,
database : process . env . DB _NAME || 'mastodon_production' ,
host : process . env . DB _HOST || 'localhost' ,
port : process . env . DB _PORT || 5432 ,
2017-05-20 17:31:47 +02:00
max : 10 ,
} ,
} ;
2017-02-02 01:31:09 +01:00
2019-03-11 00:51:23 +01:00
if ( ! ! process . env . DB _SSLMODE && process . env . DB _SSLMODE !== 'disable' ) {
pgConfigs . development . ssl = true ;
2021-12-25 22:55:06 +01:00
pgConfigs . production . ssl = true ;
2019-03-11 00:51:23 +01:00
}
const app = express ( ) ;
2020-08-11 18:24:59 +02:00
2022-04-19 09:11:58 +02:00
app . set ( 'trust proxy' , process . env . TRUSTED _PROXY _IP ? process . env . TRUSTED _PROXY _IP . split ( /(?:\s*,\s*|\s+)/ ) : 'loopback,uniquelocal' ) ;
2017-12-12 15:13:24 +01:00
2017-05-20 17:31:47 +02:00
const pgPool = new pg . Pool ( Object . assign ( pgConfigs [ env ] , dbUrlToConfig ( process . env . DATABASE _URL ) ) ) ;
const server = http . createServer ( app ) ;
const redisNamespace = process . env . REDIS _NAMESPACE || null ;
2017-02-07 14:37:12 +01:00
2017-05-07 19:42:32 +02:00
const redisParams = {
2021-12-25 22:55:06 +01:00
socket : {
host : process . env . REDIS _HOST || '127.0.0.1' ,
port : process . env . REDIS _PORT || 6379 ,
} ,
database : process . env . REDIS _DB || 0 ,
2020-06-24 22:25:23 +02:00
password : process . env . REDIS _PASSWORD || undefined ,
2017-05-20 17:31:47 +02:00
} ;
2017-05-07 19:42:32 +02:00
if ( redisNamespace ) {
2017-05-20 17:31:47 +02:00
redisParams . namespace = redisNamespace ;
2017-05-07 19:42:32 +02:00
}
2017-05-20 21:06:09 +02:00
2017-05-20 17:31:47 +02:00
const redisPrefix = redisNamespace ? ` ${ redisNamespace } : ` : '' ;
2017-05-07 19:42:32 +02:00
2022-03-21 19:08:29 +01:00
/ * *
2023-06-09 19:29:16 +02:00
* @ type { Object . < string , Array . < function ( Object < string , any > ) : void >> }
2022-03-21 19:08:29 +01:00
* /
const subs = { } ;
2021-12-25 22:55:06 +01:00
const redisSubscribeClient = await redisUrlToClient ( redisParams , process . env . REDIS _URL ) ;
const redisClient = await redisUrlToClient ( redisParams , process . env . REDIS _URL ) ;
2017-02-07 14:37:12 +01:00
2020-08-11 18:24:59 +02:00
/ * *
* @ param { string [ ] } channels
* @ return { function ( ) : void }
* /
2020-06-02 19:24:53 +02:00
const subscriptionHeartbeat = channels => {
const interval = 6 * 60 ;
2017-06-03 20:50:53 +02:00
const tellSubscribed = ( ) => {
2020-06-02 19:24:53 +02:00
channels . forEach ( channel => redisClient . set ( ` ${ redisPrefix } subscribed: ${ channel } ` , '1' , 'EX' , interval * 3 ) ) ;
2017-06-03 20:50:53 +02:00
} ;
2020-06-02 19:24:53 +02:00
2017-06-03 20:50:53 +02:00
tellSubscribed ( ) ;
2020-06-02 19:24:53 +02:00
const heartbeat = setInterval ( tellSubscribed , interval * 1000 ) ;
2017-06-03 20:50:53 +02:00
return ( ) => {
clearInterval ( heartbeat ) ;
} ;
} ;
2017-02-07 14:37:12 +01:00
2022-03-21 19:08:29 +01:00
/ * *
* @ param { string } message
* @ param { string } channel
* /
const onRedisMessage = ( message , channel ) => {
const callbacks = subs [ channel ] ;
log . silly ( ` New message on channel ${ channel } ` ) ;
if ( ! callbacks ) {
return ;
}
2023-06-09 19:29:16 +02:00
const json = parseJSON ( message , null ) ;
if ( ! json ) return ;
callbacks . forEach ( callback => callback ( json ) ) ;
2022-03-21 19:08:29 +01:00
} ;
2020-08-11 18:24:59 +02:00
/ * *
* @ param { string } channel
* @ param { function ( string ) : void } callback
* /
2017-04-17 04:32:30 +02:00
const subscribe = ( channel , callback ) => {
2017-05-20 17:31:47 +02:00
log . silly ( ` Adding listener for ${ channel } ` ) ;
2020-08-11 18:24:59 +02:00
2022-03-21 19:08:29 +01:00
subs [ channel ] = subs [ channel ] || [ ] ;
if ( subs [ channel ] . length === 0 ) {
log . verbose ( ` Subscribe ${ channel } ` ) ;
redisSubscribeClient . subscribe ( channel , onRedisMessage ) ;
}
subs [ channel ] . push ( callback ) ;
2017-05-20 17:31:47 +02:00
} ;
2017-02-03 18:27:42 +01:00
2020-08-11 18:24:59 +02:00
/ * *
* @ param { string } channel
2023-06-09 19:29:16 +02:00
* @ param { function ( Object < string , any > ) : void } callback
2020-08-11 18:24:59 +02:00
* /
2022-01-07 19:50:12 +01:00
const unsubscribe = ( channel , callback ) => {
log . silly ( ` Removing listener for ${ channel } ` ) ;
2020-08-11 18:24:59 +02:00
2022-03-21 19:08:29 +01:00
if ( ! subs [ channel ] ) {
return ;
}
subs [ channel ] = subs [ channel ] . filter ( item => item !== callback ) ;
if ( subs [ channel ] . length === 0 ) {
log . verbose ( ` Unsubscribe ${ channel } ` ) ;
redisSubscribeClient . unsubscribe ( channel ) ;
delete subs [ channel ] ;
}
2017-05-20 17:31:47 +02:00
} ;
2017-02-03 18:27:42 +01:00
2020-08-11 18:24:59 +02:00
const FALSE _VALUES = [
false ,
0 ,
2020-11-23 17:35:14 +01:00
'0' ,
'f' ,
'F' ,
'false' ,
'FALSE' ,
'off' ,
'OFF' ,
2020-08-11 18:24:59 +02:00
] ;
/ * *
* @ param { any } value
* @ return { boolean }
* /
const isTruthy = value =>
value && ! FALSE _VALUES . includes ( value ) ;
/ * *
* @ param { any } req
* @ param { any } res
* @ param { function ( Error = ) : void }
* /
2017-04-17 04:32:30 +02:00
const allowCrossDomain = ( req , res , next ) => {
2017-05-20 17:31:47 +02:00
res . header ( 'Access-Control-Allow-Origin' , '*' ) ;
res . header ( 'Access-Control-Allow-Headers' , 'Authorization, Accept, Cache-Control' ) ;
res . header ( 'Access-Control-Allow-Methods' , 'GET, OPTIONS' ) ;
2017-02-05 23:37:25 +01:00
2017-05-20 17:31:47 +02:00
next ( ) ;
} ;
2017-02-05 23:37:25 +01:00
2020-08-11 18:24:59 +02:00
/ * *
* @ param { any } req
* @ param { any } res
* @ param { function ( Error = ) : void }
* /
2017-04-17 04:32:30 +02:00
const setRequestId = ( req , res , next ) => {
2017-05-20 17:31:47 +02:00
req . requestId = uuid . v4 ( ) ;
res . header ( 'X-Request-Id' , req . requestId ) ;
2017-02-02 01:31:09 +01:00
2017-05-20 17:31:47 +02:00
next ( ) ;
} ;
2017-02-02 01:31:09 +01:00
2020-08-11 18:24:59 +02:00
/ * *
* @ param { any } req
* @ param { any } res
* @ param { function ( Error = ) : void }
* /
2017-12-12 15:13:24 +01:00
const setRemoteAddress = ( req , res , next ) => {
req . remoteAddress = req . connection . remoteAddress ;
next ( ) ;
} ;
2021-09-26 13:23:28 +02:00
/ * *
* @ param { any } req
* @ param { string [ ] } necessaryScopes
* @ return { boolean }
* /
const isInScope = ( req , necessaryScopes ) =>
req . scopes . some ( scope => necessaryScopes . includes ( scope ) ) ;
2020-08-11 18:24:59 +02:00
/ * *
* @ param { string } token
* @ param { any } req
* @ return { Promise . < void > }
* /
const accountFromToken = ( token , req ) => new Promise ( ( resolve , reject ) => {
2017-04-17 04:32:30 +02:00
pgPool . connect ( ( err , client , done ) => {
2017-02-02 01:31:09 +01:00
if ( err ) {
2020-08-11 18:24:59 +02:00
reject ( err ) ;
2017-05-20 17:31:47 +02:00
return ;
2017-02-02 01:31:09 +01:00
}
2020-11-12 23:05:24 +01:00
client . query ( 'SELECT oauth_access_tokens.id, oauth_access_tokens.resource_owner_id, users.account_id, users.chosen_languages, oauth_access_tokens.scopes, devices.device_id FROM oauth_access_tokens INNER JOIN users ON oauth_access_tokens.resource_owner_id = users.id LEFT OUTER JOIN devices ON oauth_access_tokens.id = devices.access_token_id WHERE oauth_access_tokens.token = $1 AND oauth_access_tokens.revoked_at IS NULL LIMIT 1' , [ token ] , ( err , result ) => {
2017-05-20 17:31:47 +02:00
done ( ) ;
2017-02-02 01:31:09 +01:00
2017-04-17 04:32:30 +02:00
if ( err ) {
2020-08-11 18:24:59 +02:00
reject ( err ) ;
2017-05-20 17:31:47 +02:00
return ;
2017-04-17 04:32:30 +02:00
}
2017-02-02 01:31:09 +01:00
2017-04-17 04:32:30 +02:00
if ( result . rows . length === 0 ) {
2017-05-20 17:31:47 +02:00
err = new Error ( 'Invalid access token' ) ;
2020-08-11 18:24:59 +02:00
err . status = 401 ;
2019-05-24 15:21:42 +02:00
2020-08-11 18:24:59 +02:00
reject ( err ) ;
2019-05-24 15:21:42 +02:00
return ;
}
2020-11-12 23:05:24 +01:00
req . accessTokenId = result . rows [ 0 ] . id ;
2020-08-11 18:24:59 +02:00
req . scopes = result . rows [ 0 ] . scopes . split ( ' ' ) ;
2017-05-20 17:31:47 +02:00
req . accountId = result . rows [ 0 ] . account _id ;
2018-07-14 03:59:31 +02:00
req . chosenLanguages = result . rows [ 0 ] . chosen _languages ;
2020-06-02 19:24:53 +02:00
req . deviceId = result . rows [ 0 ] . device _id ;
2017-02-04 00:34:31 +01:00
2020-08-11 18:24:59 +02:00
resolve ( ) ;
2017-05-20 17:31:47 +02:00
} ) ;
} ) ;
2020-08-11 18:24:59 +02:00
} ) ;
2017-02-04 00:34:31 +01:00
2020-08-11 18:24:59 +02:00
/ * *
* @ param { any } req
* @ param { boolean = } required
* @ return { Promise . < void > }
* /
const accountFromRequest = ( req , required = true ) => new Promise ( ( resolve , reject ) => {
2017-05-29 18:20:53 +02:00
const authorization = req . headers . authorization ;
2020-08-11 18:24:59 +02:00
const location = url . parse ( req . url , true ) ;
const accessToken = location . query . access _token || req . headers [ 'sec-websocket-protocol' ] ;
2017-02-04 00:34:31 +01:00
2017-05-21 21:13:11 +02:00
if ( ! authorization && ! accessToken ) {
2017-12-12 15:13:24 +01:00
if ( required ) {
const err = new Error ( 'Missing access token' ) ;
2020-08-11 18:24:59 +02:00
err . status = 401 ;
2017-02-02 01:31:09 +01:00
2020-08-11 18:24:59 +02:00
reject ( err ) ;
2017-12-12 15:13:24 +01:00
return ;
} else {
2020-08-11 18:24:59 +02:00
resolve ( ) ;
2017-12-12 15:13:24 +01:00
return ;
}
2017-04-17 04:32:30 +02:00
}
2017-02-02 17:10:59 +01:00
2017-05-21 21:13:11 +02:00
const token = authorization ? authorization . replace ( /^Bearer / , '' ) : accessToken ;
2017-02-02 13:56:14 +01:00
2020-08-11 18:24:59 +02:00
resolve ( accountFromToken ( token , req ) ) ;
} ) ;
/ * *
* @ param { any } req
2023-06-09 19:29:16 +02:00
* @ returns { string | undefined }
2020-08-11 18:24:59 +02:00
* /
const channelNameFromPath = req => {
const { path , query } = req ;
const onlyMedia = isTruthy ( query . only _media ) ;
2021-12-25 22:55:06 +01:00
switch ( path ) {
2020-08-11 18:24:59 +02:00
case '/api/v1/streaming/user' :
return 'user' ;
case '/api/v1/streaming/user/notification' :
return 'user:notification' ;
case '/api/v1/streaming/public' :
return onlyMedia ? 'public:media' : 'public' ;
case '/api/v1/streaming/public/local' :
return onlyMedia ? 'public:local:media' : 'public:local' ;
case '/api/v1/streaming/public/remote' :
return onlyMedia ? 'public:remote:media' : 'public:remote' ;
case '/api/v1/streaming/hashtag' :
return 'hashtag' ;
case '/api/v1/streaming/hashtag/local' :
return 'hashtag:local' ;
case '/api/v1/streaming/direct' :
return 'direct' ;
case '/api/v1/streaming/list' :
return 'list' ;
2020-11-23 17:35:14 +01:00
default :
return undefined ;
2020-08-11 18:24:59 +02:00
}
2017-05-20 17:31:47 +02:00
} ;
2017-02-05 03:19:04 +01:00
2020-08-11 18:24:59 +02:00
const PUBLIC _CHANNELS = [
2017-12-12 15:13:24 +01:00
'public' ,
2018-05-21 12:43:38 +02:00
'public:media' ,
2017-12-12 15:13:24 +01:00
'public:local' ,
2018-05-21 12:43:38 +02:00
'public:local:media' ,
2020-05-10 10:36:18 +02:00
'public:remote' ,
'public:remote:media' ,
2017-12-12 15:13:24 +01:00
'hashtag' ,
'hashtag:local' ,
] ;
2020-08-11 18:24:59 +02:00
/ * *
* @ param { any } req
* @ param { string } channelName
* @ return { Promise . < void > }
* /
const checkScopes = ( req , channelName ) => new Promise ( ( resolve , reject ) => {
log . silly ( req . requestId , ` Checking OAuth scopes for ${ channelName } ` ) ;
// When accessing public channels, no scopes are needed
if ( PUBLIC _CHANNELS . includes ( channelName ) ) {
resolve ( ) ;
return ;
}
2019-05-24 15:21:42 +02:00
2020-08-11 18:24:59 +02:00
// The `read` scope has the highest priority, if the token has it
// then it can access all streams
const requiredScopes = [ 'read' ] ;
// When accessing specifically the notifications stream,
// we need a read:notifications, while in all other cases,
// we can allow access with read:statuses. Mind that the
// user stream will not contain notifications unless
// the token has either read or read:notifications scope
// as well, this is handled separately.
if ( channelName === 'user:notification' ) {
requiredScopes . push ( 'read:notifications' ) ;
} else {
requiredScopes . push ( 'read:statuses' ) ;
2019-05-24 15:21:42 +02:00
}
2017-12-12 15:13:24 +01:00
2021-10-13 05:02:55 +02:00
if ( req . scopes && requiredScopes . some ( requiredScope => req . scopes . includes ( requiredScope ) ) ) {
2020-08-11 18:24:59 +02:00
resolve ( ) ;
return ;
}
2017-05-29 18:20:53 +02:00
2020-08-11 18:24:59 +02:00
const err = new Error ( 'Access token does not cover required scopes' ) ;
err . status = 401 ;
reject ( err ) ;
} ) ;
2017-12-12 15:13:24 +01:00
2020-08-11 18:24:59 +02:00
/ * *
* @ param { any } info
* @ param { function ( boolean , number , string ) : void } callback
* /
const wsVerifyClient = ( info , callback ) => {
// When verifying the websockets connection, we no longer pre-emptively
// check OAuth scopes and drop the connection if they're missing. We only
// drop the connection if access without token is not allowed by environment
// variables. OAuth scope checks are moved to the point of subscription
// to a specific stream.
accountFromRequest ( info . req , alwaysRequireAuth ) . then ( ( ) => {
callback ( true , undefined , undefined ) ;
} ) . catch ( err => {
log . error ( info . req . requestId , err . toString ( ) ) ;
callback ( false , 401 , 'Unauthorized' ) ;
} ) ;
} ;
2020-11-12 23:05:24 +01:00
/ * *
* @ typedef SystemMessageHandlers
* @ property { function ( ) : void } onKill
* /
/ * *
* @ param { any } req
* @ param { SystemMessageHandlers } eventHandlers
2023-06-09 19:29:16 +02:00
* @ returns { function ( object ) : void }
2020-11-12 23:05:24 +01:00
* /
const createSystemMessageListener = ( req , eventHandlers ) => {
return message => {
2023-06-09 19:29:16 +02:00
const { event } = message ;
2020-11-12 23:05:24 +01:00
log . silly ( req . requestId , ` System message for ${ req . accountId } : ${ event } ` ) ;
if ( event === 'kill' ) {
log . verbose ( req . requestId , ` Closing connection for ${ req . accountId } due to expired access token ` ) ;
eventHandlers . onKill ( ) ;
}
2020-11-23 17:35:14 +01:00
} ;
2020-11-12 23:05:24 +01:00
} ;
/ * *
* @ param { any } req
* @ param { any } res
* /
const subscribeHttpToSystemChannel = ( req , res ) => {
const systemChannelId = ` timeline:access_token: ${ req . accessTokenId } ` ;
const listener = createSystemMessageListener ( req , {
2021-12-25 22:55:06 +01:00
onKill ( ) {
2020-11-12 23:05:24 +01:00
res . end ( ) ;
} ,
} ) ;
res . on ( 'close' , ( ) => {
unsubscribe ( ` ${ redisPrefix } ${ systemChannelId } ` , listener ) ;
} ) ;
subscribe ( ` ${ redisPrefix } ${ systemChannelId } ` , listener ) ;
} ;
2020-08-11 18:24:59 +02:00
/ * *
* @ param { any } req
* @ param { any } res
* @ param { function ( Error = ) : void } next
* /
2017-05-29 18:20:53 +02:00
const authenticationMiddleware = ( req , res , next ) => {
if ( req . method === 'OPTIONS' ) {
next ( ) ;
return ;
}
2020-08-11 18:24:59 +02:00
accountFromRequest ( req , alwaysRequireAuth ) . then ( ( ) => checkScopes ( req , channelNameFromPath ( req ) ) ) . then ( ( ) => {
2020-11-12 23:05:24 +01:00
subscribeHttpToSystemChannel ( req , res ) ;
} ) . then ( ( ) => {
2020-08-11 18:24:59 +02:00
next ( ) ;
} ) . catch ( err => {
next ( err ) ;
} ) ;
2017-05-29 18:20:53 +02:00
} ;
2020-08-11 18:24:59 +02:00
/ * *
* @ param { Error } err
* @ param { any } req
* @ param { any } res
* @ param { function ( Error = ) : void } next
* /
const errorMiddleware = ( err , req , res , next ) => {
2017-05-28 16:25:26 +02:00
log . error ( req . requestId , err . toString ( ) ) ;
2020-08-11 18:24:59 +02:00
if ( res . headersSent ) {
2020-11-23 17:35:14 +01:00
next ( err ) ;
return ;
2020-08-11 18:24:59 +02:00
}
res . writeHead ( err . status || 500 , { 'Content-Type' : 'application/json' } ) ;
res . end ( JSON . stringify ( { error : err . status ? err . toString ( ) : 'An unexpected error occurred' } ) ) ;
2017-05-20 17:31:47 +02:00
} ;
2017-02-05 03:19:04 +01:00
2020-08-11 18:24:59 +02:00
/ * *
2021-12-25 22:55:06 +01:00
* @ param { array } arr
2020-08-11 18:24:59 +02:00
* @ param { number = } shift
* @ return { string }
* /
2017-04-17 04:32:30 +02:00
const placeholders = ( arr , shift = 0 ) => arr . map ( ( _ , i ) => ` $ ${ i + 1 + shift } ` ) . join ( ', ' ) ;
2017-02-02 01:31:09 +01:00
2020-08-11 18:24:59 +02:00
/ * *
* @ param { string } listId
* @ param { any } req
* @ return { Promise . < void > }
* /
const authorizeListAccess = ( listId , req ) => new Promise ( ( resolve , reject ) => {
const { accountId } = req ;
2017-11-18 00:16:48 +01:00
pgPool . connect ( ( err , client , done ) => {
if ( err ) {
2020-08-11 18:24:59 +02:00
reject ( ) ;
2017-11-18 00:16:48 +01:00
return ;
}
2020-08-11 18:24:59 +02:00
client . query ( 'SELECT id, account_id FROM lists WHERE id = $1 LIMIT 1' , [ listId ] , ( err , result ) => {
2017-11-18 00:16:48 +01:00
done ( ) ;
2020-08-11 18:24:59 +02:00
if ( err || result . rows . length === 0 || result . rows [ 0 ] . account _id !== accountId ) {
reject ( ) ;
2017-11-18 00:16:48 +01:00
return ;
}
2020-08-11 18:24:59 +02:00
resolve ( ) ;
2017-11-18 00:16:48 +01:00
} ) ;
} ) ;
2020-08-11 18:24:59 +02:00
} ) ;
2017-11-18 00:16:48 +01:00
2020-08-11 18:24:59 +02:00
/ * *
* @ param { string [ ] } ids
* @ param { any } req
* @ param { function ( string , string ) : void } output
* @ param { function ( string [ ] , function ( string ) : void ) : void } attachCloseHandler
* @ param { boolean = } needsFiltering
2023-06-09 19:29:16 +02:00
* @ returns { function ( object ) : void }
2020-08-11 18:24:59 +02:00
* /
2021-09-26 13:23:28 +02:00
const streamFrom = ( ids , req , output , attachCloseHandler , needsFiltering = false ) => {
2021-12-25 22:55:06 +01:00
const accountId = req . accountId || req . remoteAddress ;
2020-06-02 19:24:53 +02:00
2021-09-26 13:23:28 +02:00
log . verbose ( req . requestId , ` Starting stream from ${ ids . join ( ', ' ) } for ${ accountId } ` ) ;
2017-04-17 04:32:30 +02:00
2023-06-09 19:29:16 +02:00
// Currently message is of type string, soon it'll be Record<string, any>
2017-04-17 04:32:30 +02:00
const listener = message => {
2023-06-09 19:29:16 +02:00
const { event , payload , queued _at } = message ;
2017-02-02 13:56:14 +01:00
2017-04-17 04:32:30 +02:00
const transmit = ( ) => {
2021-12-25 22:55:06 +01:00
const now = new Date ( ) . getTime ( ) ;
const delta = now - queued _at ;
2017-09-24 15:31:03 +02:00
const encodedPayload = typeof payload === 'object' ? JSON . stringify ( payload ) : payload ;
2017-02-02 13:56:14 +01:00
2017-12-12 15:13:24 +01:00
log . silly ( req . requestId , ` Transmitting for ${ accountId } : ${ event } ${ encodedPayload } Delay: ${ delta } ms ` ) ;
2017-07-07 16:56:52 +02:00
output ( event , encodedPayload ) ;
2017-05-20 17:31:47 +02:00
} ;
2017-02-02 13:56:14 +01:00
2017-04-17 04:32:30 +02:00
// Only messages that may require filtering are statuses, since notifications
// are already personalized and deletes do not matter
2018-04-17 13:49:09 +02:00
if ( ! needsFiltering || event !== 'update' ) {
transmit ( ) ;
return ;
}
2017-02-02 13:56:14 +01:00
2021-12-25 22:55:06 +01:00
const unpackedPayload = payload ;
2018-04-17 13:49:09 +02:00
const targetAccountIds = [ unpackedPayload . account . id ] . concat ( unpackedPayload . mentions . map ( item => item . id ) ) ;
2021-12-25 22:55:06 +01:00
const accountDomain = unpackedPayload . account . acct . split ( '@' ) [ 1 ] ;
2017-04-17 04:32:30 +02:00
2018-07-14 03:59:31 +02:00
if ( Array . isArray ( req . chosenLanguages ) && unpackedPayload . language !== null && req . chosenLanguages . indexOf ( unpackedPayload . language ) === - 1 ) {
2018-04-17 13:49:09 +02:00
log . silly ( req . requestId , ` Message ${ unpackedPayload . id } filtered by language ( ${ unpackedPayload . language } ) ` ) ;
return ;
}
// When the account is not logged in, it is not necessary to confirm the block or mute
if ( ! req . accountId ) {
transmit ( ) ;
return ;
}
pgPool . connect ( ( err , client , done ) => {
if ( err ) {
log . error ( err ) ;
return ;
}
const queries = [
2021-12-25 22:55:06 +01:00
client . query ( ` SELECT 1
FROM blocks
WHERE ( account _id = $1 AND target _account _id IN ( $ { placeholders ( targetAccountIds , 2 ) } ) )
OR ( account _id = $2 AND target _account _id = $1 )
UNION
SELECT 1
FROM mutes
WHERE account _id = $1
AND target _account _id IN ( $ { placeholders ( targetAccountIds , 2 ) } ) ` , [req.accountId, unpackedPayload.account.id].concat(targetAccountIds)),
2018-04-17 13:49:09 +02:00
] ;
if ( accountDomain ) {
queries . push ( client . query ( 'SELECT 1 FROM account_domain_blocks WHERE account_id = $1 AND domain = $2' , [ req . accountId , accountDomain ] ) ) ;
}
Promise . all ( queries ) . then ( values => {
done ( ) ;
if ( values [ 0 ] . rows . length > 0 || ( values . length > 1 && values [ 1 ] . rows . length > 0 ) ) {
2017-05-27 00:53:48 +02:00
return ;
}
2018-04-17 13:49:09 +02:00
transmit ( ) ;
} ) . catch ( err => {
done ( ) ;
log . error ( err ) ;
2017-05-20 17:31:47 +02:00
} ) ;
2018-04-17 13:49:09 +02:00
} ) ;
2017-05-20 17:31:47 +02:00
} ;
2017-04-17 04:32:30 +02:00
2020-06-02 19:24:53 +02:00
ids . forEach ( id => {
subscribe ( ` ${ redisPrefix } ${ id } ` , listener ) ;
} ) ;
2020-08-11 18:24:59 +02:00
if ( attachCloseHandler ) {
attachCloseHandler ( ids . map ( id => ` ${ redisPrefix } ${ id } ` ) , listener ) ;
}
return listener ;
2017-05-20 17:31:47 +02:00
} ;
2017-02-02 01:31:09 +01:00
2020-08-11 18:24:59 +02:00
/ * *
* @ param { any } req
* @ param { any } res
* @ return { function ( string , string ) : void }
* /
2017-04-17 04:32:30 +02:00
const streamToHttp = ( req , res ) => {
2017-12-12 15:13:24 +01:00
const accountId = req . accountId || req . remoteAddress ;
2017-05-20 17:31:47 +02:00
res . setHeader ( 'Content-Type' , 'text/event-stream' ) ;
2020-01-24 20:51:33 +01:00
res . setHeader ( 'Cache-Control' , 'no-store' ) ;
2017-05-20 17:31:47 +02:00
res . setHeader ( 'Transfer-Encoding' , 'chunked' ) ;
2017-02-04 00:34:31 +01:00
2020-01-24 20:51:33 +01:00
res . write ( ':)\n' ) ;
2017-05-20 17:31:47 +02:00
const heartbeat = setInterval ( ( ) => res . write ( ':thump\n' ) , 15000 ) ;
2017-02-04 00:34:31 +01:00
2017-04-17 04:32:30 +02:00
req . on ( 'close' , ( ) => {
2017-12-12 15:13:24 +01:00
log . verbose ( req . requestId , ` Ending stream for ${ accountId } ` ) ;
2017-05-20 17:31:47 +02:00
clearInterval ( heartbeat ) ;
} ) ;
2017-02-02 15:20:31 +01:00
2017-04-17 04:32:30 +02:00
return ( event , payload ) => {
2017-05-20 17:31:47 +02:00
res . write ( ` event: ${ event } \n ` ) ;
res . write ( ` data: ${ payload } \n \n ` ) ;
} ;
} ;
2017-02-02 01:31:09 +01:00
2020-08-11 18:24:59 +02:00
/ * *
* @ param { any } req
* @ param { function ( ) : void } [ closeHandler ]
2021-12-25 22:55:06 +01:00
* @ return { function ( string [ ] ) : void }
2020-08-11 18:24:59 +02:00
* /
2021-12-25 22:55:06 +01:00
const streamHttpEnd = ( req , closeHandler = undefined ) => ( ids ) => {
2017-04-17 04:32:30 +02:00
req . on ( 'close' , ( ) => {
2020-06-02 19:24:53 +02:00
ids . forEach ( id => {
2021-12-25 22:55:06 +01:00
unsubscribe ( id ) ;
2020-06-02 19:24:53 +02:00
} ) ;
2017-06-03 20:50:53 +02:00
if ( closeHandler ) {
closeHandler ( ) ;
}
2017-05-20 17:31:47 +02:00
} ) ;
} ;
2017-02-04 00:34:31 +01:00
2020-08-11 18:24:59 +02:00
/ * *
* @ param { any } req
* @ param { any } ws
* @ param { string [ ] } streamName
* @ return { function ( string , string ) : void }
* /
const streamToWs = ( req , ws , streamName ) => ( event , payload ) => {
2017-05-28 16:25:26 +02:00
if ( ws . readyState !== ws . OPEN ) {
log . error ( req . requestId , 'Tried writing to closed socket' ) ;
return ;
}
2017-02-04 00:34:31 +01:00
2020-08-11 18:24:59 +02:00
ws . send ( JSON . stringify ( { stream : streamName , event , payload } ) ) ;
2017-05-20 17:31:47 +02:00
} ;
2017-02-04 00:34:31 +01:00
2020-08-11 18:24:59 +02:00
/ * *
* @ param { any } res
* /
2018-10-11 19:24:43 +02:00
const httpNotFound = res => {
res . writeHead ( 404 , { 'Content-Type' : 'application/json' } ) ;
res . end ( JSON . stringify ( { error : 'Not found' } ) ) ;
} ;
2017-05-20 17:31:47 +02:00
app . use ( setRequestId ) ;
2017-12-12 15:13:24 +01:00
app . use ( setRemoteAddress ) ;
2017-05-20 17:31:47 +02:00
app . use ( allowCrossDomain ) ;
2018-08-26 11:54:25 +02:00
app . get ( '/api/v1/streaming/health' , ( req , res ) => {
res . writeHead ( 200 , { 'Content-Type' : 'text/plain' } ) ;
res . end ( 'OK' ) ;
} ) ;
2017-05-20 17:31:47 +02:00
app . use ( authenticationMiddleware ) ;
app . use ( errorMiddleware ) ;
2017-02-02 01:31:09 +01:00
2020-08-11 18:24:59 +02:00
app . get ( '/api/v1/streaming/*' , ( req , res ) => {
channelNameToIds ( req , channelNameFromPath ( req ) , req . query ) . then ( ( { channelIds , options } ) => {
const onSend = streamToHttp ( req , res ) ;
2021-12-25 22:55:06 +01:00
const onEnd = streamHttpEnd ( req , subscriptionHeartbeat ( channelIds ) ) ;
2018-05-21 12:43:38 +02:00
2021-09-26 13:23:28 +02:00
streamFrom ( channelIds , req , onSend , onEnd , options . needsFiltering ) ;
2020-08-11 18:24:59 +02:00
} ) . catch ( err => {
log . verbose ( req . requestId , 'Subscription error:' , err . toString ( ) ) ;
2018-10-11 19:24:43 +02:00
httpNotFound ( res ) ;
2017-11-18 00:16:48 +01:00
} ) ;
} ) ;
2021-03-24 09:37:41 +01:00
const wss = new WebSocket . Server ( { server , verifyClient : wsVerifyClient } ) ;
2017-05-29 18:20:53 +02:00
2020-08-11 18:24:59 +02:00
/ * *
* @ typedef StreamParams
* @ property { string } [ tag ]
* @ property { string } [ list ]
* @ property { string } [ only _media ]
* /
2021-09-26 13:23:28 +02:00
/ * *
* @ param { any } req
* @ return { string [ ] }
* /
const channelsForUserStream = req => {
const arr = [ ` timeline: ${ req . accountId } ` ] ;
if ( isInScope ( req , [ 'crypto' ] ) && req . deviceId ) {
arr . push ( ` timeline: ${ req . accountId } : ${ req . deviceId } ` ) ;
}
if ( isInScope ( req , [ 'read' , 'read:notifications' ] ) ) {
arr . push ( ` timeline: ${ req . accountId } :notifications ` ) ;
}
return arr ;
} ;
2020-08-11 18:24:59 +02:00
/ * *
* @ param { any } req
* @ param { string } name
* @ param { StreamParams } params
2021-09-26 13:23:28 +02:00
* @ return { Promise . < { channelIds : string [ ] , options : { needsFiltering : boolean } } > }
2020-08-11 18:24:59 +02:00
* /
const channelNameToIds = ( req , name , params ) => new Promise ( ( resolve , reject ) => {
2021-12-25 22:55:06 +01:00
switch ( name ) {
2017-05-29 18:20:53 +02:00
case 'user' :
2020-08-11 18:24:59 +02:00
resolve ( {
2021-09-26 13:23:28 +02:00
channelIds : channelsForUserStream ( req ) ,
options : { needsFiltering : false } ,
2020-08-11 18:24:59 +02:00
} ) ;
2020-06-02 19:24:53 +02:00
2017-06-03 20:50:53 +02:00
break ;
case 'user:notification' :
2020-08-11 18:24:59 +02:00
resolve ( {
2021-09-26 13:23:28 +02:00
channelIds : [ ` timeline: ${ req . accountId } :notifications ` ] ,
options : { needsFiltering : false } ,
2020-08-11 18:24:59 +02:00
} ) ;
2017-05-29 18:20:53 +02:00
break ;
case 'public' :
2020-08-11 18:24:59 +02:00
resolve ( {
channelIds : [ 'timeline:public' ] ,
2021-09-26 13:23:28 +02:00
options : { needsFiltering : true } ,
2020-08-11 18:24:59 +02:00
} ) ;
2017-05-29 18:20:53 +02:00
break ;
case 'public:local' :
2020-08-11 18:24:59 +02:00
resolve ( {
channelIds : [ 'timeline:public:local' ] ,
2021-09-26 13:23:28 +02:00
options : { needsFiltering : true } ,
2020-08-11 18:24:59 +02:00
} ) ;
2017-05-29 18:20:53 +02:00
break ;
2020-05-10 10:36:18 +02:00
case 'public:remote' :
2020-08-11 18:24:59 +02:00
resolve ( {
channelIds : [ 'timeline:public:remote' ] ,
2021-09-26 13:23:28 +02:00
options : { needsFiltering : true } ,
2020-08-11 18:24:59 +02:00
} ) ;
2020-05-10 10:36:18 +02:00
break ;
2018-05-21 12:43:38 +02:00
case 'public:media' :
2020-08-11 18:24:59 +02:00
resolve ( {
channelIds : [ 'timeline:public:media' ] ,
2021-09-26 13:23:28 +02:00
options : { needsFiltering : true } ,
2020-08-11 18:24:59 +02:00
} ) ;
2018-05-21 12:43:38 +02:00
break ;
case 'public:local:media' :
2020-08-11 18:24:59 +02:00
resolve ( {
channelIds : [ 'timeline:public:local:media' ] ,
2021-09-26 13:23:28 +02:00
options : { needsFiltering : true } ,
2020-08-11 18:24:59 +02:00
} ) ;
2018-05-21 12:43:38 +02:00
break ;
2020-05-10 10:36:18 +02:00
case 'public:remote:media' :
2020-08-11 18:24:59 +02:00
resolve ( {
channelIds : [ 'timeline:public:remote:media' ] ,
2021-09-26 13:23:28 +02:00
options : { needsFiltering : true } ,
2020-08-11 18:24:59 +02:00
} ) ;
2020-05-10 10:36:18 +02:00
break ;
2018-04-18 13:09:06 +02:00
case 'direct' :
2020-08-11 18:24:59 +02:00
resolve ( {
channelIds : [ ` timeline:direct: ${ req . accountId } ` ] ,
2021-09-26 13:23:28 +02:00
options : { needsFiltering : false } ,
2020-08-11 18:24:59 +02:00
} ) ;
2018-04-18 13:09:06 +02:00
break ;
2017-05-29 18:20:53 +02:00
case 'hashtag' :
2020-08-11 18:24:59 +02:00
if ( ! params . tag || params . tag . length === 0 ) {
reject ( 'No tag for stream provided' ) ;
} else {
resolve ( {
channelIds : [ ` timeline:hashtag: ${ params . tag . toLowerCase ( ) } ` ] ,
2021-09-26 13:23:28 +02:00
options : { needsFiltering : true } ,
2020-08-11 18:24:59 +02:00
} ) ;
2018-10-11 19:24:43 +02:00
}
2017-05-29 18:20:53 +02:00
break ;
case 'hashtag:local' :
2020-08-11 18:24:59 +02:00
if ( ! params . tag || params . tag . length === 0 ) {
reject ( 'No tag for stream provided' ) ;
} else {
resolve ( {
channelIds : [ ` timeline:hashtag: ${ params . tag . toLowerCase ( ) } :local ` ] ,
2021-09-26 13:23:28 +02:00
options : { needsFiltering : true } ,
2020-08-11 18:24:59 +02:00
} ) ;
2018-10-11 19:24:43 +02:00
}
2017-05-29 18:20:53 +02:00
break ;
2017-11-18 00:16:48 +01:00
case 'list' :
2020-08-11 18:24:59 +02:00
authorizeListAccess ( params . list , req ) . then ( ( ) => {
resolve ( {
channelIds : [ ` timeline:list: ${ params . list } ` ] ,
2021-09-26 13:23:28 +02:00
options : { needsFiltering : false } ,
2020-08-11 18:24:59 +02:00
} ) ;
} ) . catch ( ( ) => {
reject ( 'Not authorized to stream this list' ) ;
2017-11-18 00:16:48 +01:00
} ) ;
2020-08-11 18:24:59 +02:00
2017-11-18 00:16:48 +01:00
break ;
2017-05-29 18:20:53 +02:00
default :
2020-08-11 18:24:59 +02:00
reject ( 'Unknown stream type' ) ;
}
} ) ;
/ * *
* @ param { string } channelName
* @ param { StreamParams } params
* @ return { string [ ] }
* /
const streamNameFromChannelName = ( channelName , params ) => {
if ( channelName === 'list' ) {
return [ channelName , params . list ] ;
} else if ( [ 'hashtag' , 'hashtag:local' ] . includes ( channelName ) ) {
return [ channelName , params . tag ] ;
} else {
return [ channelName ] ;
}
} ;
/ * *
* @ typedef WebSocketSession
* @ property { any } socket
* @ property { any } request
* @ property { Object . < string , { listener : function ( string ) : void , stopHeartbeat : function ( ) : void } > } subscriptions
* /
/ * *
* @ param { WebSocketSession } session
* @ param { string } channelName
* @ param { StreamParams } params
* /
const subscribeWebsocketToChannel = ( { socket , request , subscriptions } , channelName , params ) =>
2021-12-25 22:55:06 +01:00
checkScopes ( request , channelName ) . then ( ( ) => channelNameToIds ( request , channelName , params ) ) . then ( ( {
channelIds ,
options ,
} ) => {
2020-08-11 18:24:59 +02:00
if ( subscriptions [ channelIds . join ( ';' ) ] ) {
return ;
}
2021-12-25 22:55:06 +01:00
const onSend = streamToWs ( request , socket , streamNameFromChannelName ( channelName , params ) ) ;
2020-08-11 18:24:59 +02:00
const stopHeartbeat = subscriptionHeartbeat ( channelIds ) ;
2021-12-25 22:55:06 +01:00
const listener = streamFrom ( channelIds , request , onSend , undefined , options . needsFiltering ) ;
2020-08-11 18:24:59 +02:00
subscriptions [ channelIds . join ( ';' ) ] = {
listener ,
stopHeartbeat ,
} ;
} ) . catch ( err => {
log . verbose ( request . requestId , 'Subscription error:' , err . toString ( ) ) ;
socket . send ( JSON . stringify ( { error : err . toString ( ) } ) ) ;
} ) ;
/ * *
* @ param { WebSocketSession } session
* @ param { string } channelName
* @ param { StreamParams } params
* /
const unsubscribeWebsocketFromChannel = ( { socket , request , subscriptions } , channelName , params ) =>
channelNameToIds ( request , channelName , params ) . then ( ( { channelIds } ) => {
log . verbose ( request . requestId , ` Ending stream from ${ channelIds . join ( ', ' ) } for ${ request . accountId } ` ) ;
2020-08-12 15:36:07 +02:00
const subscription = subscriptions [ channelIds . join ( ';' ) ] ;
2020-08-11 18:24:59 +02:00
2020-08-12 15:36:07 +02:00
if ( ! subscription ) {
2020-08-11 18:24:59 +02:00
return ;
}
2020-08-12 15:36:07 +02:00
const { listener , stopHeartbeat } = subscription ;
2020-08-11 18:24:59 +02:00
channelIds . forEach ( channelId => {
unsubscribe ( ` ${ redisPrefix } ${ channelId } ` , listener ) ;
} ) ;
stopHeartbeat ( ) ;
2020-08-12 15:36:07 +02:00
delete subscriptions [ channelIds . join ( ';' ) ] ;
2020-08-11 18:24:59 +02:00
} ) . catch ( err => {
log . verbose ( request . requestId , 'Unsubscription error:' , err ) ;
socket . send ( JSON . stringify ( { error : err . toString ( ) } ) ) ;
} ) ;
2020-11-12 23:05:24 +01:00
/ * *
* @ param { WebSocketSession } session
* /
const subscribeWebsocketToSystemChannel = ( { socket , request , subscriptions } ) => {
const systemChannelId = ` timeline:access_token: ${ request . accessTokenId } ` ;
const listener = createSystemMessageListener ( request , {
2021-12-25 22:55:06 +01:00
onKill ( ) {
2020-11-12 23:05:24 +01:00
socket . close ( ) ;
} ,
} ) ;
subscribe ( ` ${ redisPrefix } ${ systemChannelId } ` , listener ) ;
subscriptions [ systemChannelId ] = {
listener ,
2021-12-25 22:55:06 +01:00
stopHeartbeat : ( ) => {
} ,
2020-11-12 23:05:24 +01:00
} ;
} ;
2020-08-11 18:24:59 +02:00
/ * *
* @ param { string | string [ ] } arrayOrString
* @ return { string }
* /
const firstParam = arrayOrString => {
if ( Array . isArray ( arrayOrString ) ) {
return arrayOrString [ 0 ] ;
} else {
return arrayOrString ;
}
} ;
wss . on ( 'connection' , ( ws , req ) => {
const location = url . parse ( req . url , true ) ;
2021-12-25 22:55:06 +01:00
req . requestId = uuid . v4 ( ) ;
2020-08-11 18:24:59 +02:00
req . remoteAddress = ws . _socket . remoteAddress ;
2021-03-24 09:37:41 +01:00
ws . isAlive = true ;
ws . on ( 'pong' , ( ) => {
ws . isAlive = true ;
} ) ;
2020-08-11 18:24:59 +02:00
/ * *
* @ type { WebSocketSession }
* /
const session = {
socket : ws ,
request : req ,
subscriptions : { } ,
} ;
const onEnd = ( ) => {
const keys = Object . keys ( session . subscriptions ) ;
keys . forEach ( channelIds => {
const { listener , stopHeartbeat } = session . subscriptions [ channelIds ] ;
channelIds . split ( ';' ) . forEach ( channelId => {
unsubscribe ( ` ${ redisPrefix } ${ channelId } ` , listener ) ;
} ) ;
stopHeartbeat ( ) ;
} ) ;
} ;
ws . on ( 'close' , onEnd ) ;
ws . on ( 'error' , onEnd ) ;
2023-06-09 19:29:16 +02:00
ws . on ( 'message' , ( data , isBinary ) => {
if ( isBinary ) {
log . debug ( 'Received binary data, closing connection' ) ;
ws . close ( 1003 , 'The mastodon streaming server does not support binary messages' ) ;
return ;
}
const message = data . toString ( 'utf8' ) ;
const json = parseJSON ( message , session . request ) ;
2020-11-12 23:05:24 +01:00
2020-09-22 15:30:41 +02:00
if ( ! json ) return ;
2020-11-12 23:05:24 +01:00
2020-09-22 15:30:41 +02:00
const { type , stream , ... params } = json ;
2020-08-11 18:24:59 +02:00
if ( type === 'subscribe' ) {
subscribeWebsocketToChannel ( session , firstParam ( stream ) , params ) ;
} else if ( type === 'unsubscribe' ) {
2020-11-23 17:35:14 +01:00
unsubscribeWebsocketFromChannel ( session , firstParam ( stream ) , params ) ;
2020-08-11 18:24:59 +02:00
} else {
// Unknown action type
}
} ) ;
2020-11-12 23:05:24 +01:00
subscribeWebsocketToSystemChannel ( session ) ;
2020-08-11 18:24:59 +02:00
if ( location . query . stream ) {
subscribeWebsocketToChannel ( session , firstParam ( location . query . stream ) , location . query ) ;
2017-05-29 18:20:53 +02:00
}
2017-05-20 17:31:47 +02:00
} ) ;
2017-02-04 00:34:31 +01:00
2021-03-24 09:37:41 +01:00
setInterval ( ( ) => {
wss . clients . forEach ( ws => {
if ( ws . isAlive === false ) {
ws . terminate ( ) ;
return ;
}
ws . isAlive = false ;
2021-05-02 14:30:26 +02:00
ws . ping ( '' , false ) ;
2021-03-24 09:37:41 +01:00
} ) ;
} , 30000 ) ;
2017-05-28 16:25:26 +02:00
2018-10-20 02:25:25 +02:00
attachServerWithConfig ( server , address => {
2021-05-01 23:19:18 +02:00
log . warn ( ` Worker ${ workerId } now listening on ${ address } ` ) ;
2018-10-20 02:25:25 +02:00
} ) ;
2017-04-21 19:24:31 +02:00
2017-05-28 16:25:26 +02:00
const onExit = ( ) => {
2021-05-01 23:19:18 +02:00
log . warn ( ` Worker ${ workerId } exiting ` ) ;
2017-05-20 17:31:47 +02:00
server . close ( ) ;
2017-07-07 20:01:00 +02:00
process . exit ( 0 ) ;
2017-05-28 16:25:26 +02:00
} ;
const onError = ( err ) => {
log . error ( err ) ;
2017-12-12 20:19:33 +01:00
server . close ( ) ;
process . exit ( 0 ) ;
2017-05-28 16:25:26 +02:00
} ;
process . on ( 'SIGINT' , onExit ) ;
process . on ( 'SIGTERM' , onExit ) ;
process . on ( 'exit' , onExit ) ;
2017-12-12 20:19:33 +01:00
process . on ( 'uncaughtException' , onError ) ;
2017-05-28 16:25:26 +02:00
} ;
2020-08-11 18:24:59 +02:00
/ * *
* @ param { any } server
* @ param { function ( string ) : void } [ onSuccess ]
* /
2018-10-20 02:25:25 +02:00
const attachServerWithConfig = ( server , onSuccess ) => {
if ( process . env . SOCKET || process . env . PORT && isNaN ( + process . env . PORT ) ) {
server . listen ( process . env . SOCKET || process . env . PORT , ( ) => {
if ( onSuccess ) {
2018-10-21 16:41:33 +02:00
fs . chmodSync ( server . address ( ) , 0o666 ) ;
2018-10-20 02:25:25 +02:00
onSuccess ( server . address ( ) ) ;
}
} ) ;
} else {
2019-07-15 05:56:35 +02:00
server . listen ( + process . env . PORT || 4000 , process . env . BIND || '127.0.0.1' , ( ) => {
2018-10-20 02:25:25 +02:00
if ( onSuccess ) {
onSuccess ( ` ${ server . address ( ) . address } : ${ server . address ( ) . port } ` ) ;
}
} ) ;
}
} ;
2020-08-11 18:24:59 +02:00
/ * *
* @ param { function ( Error = ) : void } onSuccess
* /
2018-10-20 02:25:25 +02:00
const onPortAvailable = onSuccess => {
const testServer = http . createServer ( ) ;
testServer . once ( 'error' , err => {
onSuccess ( err ) ;
} ) ;
testServer . once ( 'listening' , ( ) => {
testServer . once ( 'close' , ( ) => onSuccess ( ) ) ;
testServer . close ( ) ;
} ) ;
attachServerWithConfig ( testServer ) ;
} ;
onPortAvailable ( err => {
if ( err ) {
log . error ( 'Could not start server, the port or socket is in use' ) ;
return ;
}
throng ( {
workers : numWorkers ,
lifetime : Infinity ,
start : startWorker ,
master : startMaster ,
} ) ;
2017-05-28 16:25:26 +02:00
} ) ;