@@ -9,21 +9,18 @@ const Redis = require('ioredis');
99const log = require ( 'npmlog' ) ;
1010const counters = require ( './counters' ) ;
1111const errors = require ( './errors' ) ;
12+ const { SettingsHandler } = require ( './settings-handler' ) ;
1213const { publish, MARKED_SPAM , MARKED_HAM } = require ( './events' ) ;
1314
15+ const LEGACY_CHANNEL = 'wd_events' ;
1416const WORKER_CHANNEL_PREFIX = 'wd_events:worker:' ;
1517const USER_WORKERS_PREFIX = 'wd:imap:users:' ;
1618const USER_WORKERS_TTL = Math . max ( Number ( config ?. imap ?. notifyUserWorkersTtl ) || 120 , 30 ) ;
1719const USER_WORKERS_REFRESH = Math . max ( Number ( config ?. imap ?. notifyUserWorkersRefresh ) || Math . floor ( USER_WORKERS_TTL / 4 ) , 10 ) ;
1820const USER_WORKERS_TTL_MS = USER_WORKERS_TTL * 1000 ;
21+ const USE_WORKER_WD_EVENTS_SETTING = 'const:imap:use_wd_worker_channels' ;
22+ const USE_WORKER_WD_EVENTS_CACHE_TTL = 5000 ; // ms
1923const WORKER_ID = `${ config ?. dbs ?. workerId || os . hostname ( ) } :${ process . pid } ` ;
20- const USER_REGISTRY_STATE = {
21- counts : new Map ( ) ,
22- workerId : WORKER_ID ,
23- workerChannel : `${ WORKER_CHANNEL_PREFIX } ${ WORKER_ID } ` ,
24- timer : null ,
25- redis : null
26- } ;
2724
2825class ImapNotifier extends EventEmitter {
2926 constructor ( options ) {
@@ -33,9 +30,17 @@ class ImapNotifier extends EventEmitter {
3330 this . redis = options . redis || new Redis ( tools . redisConfig ( config . dbs . redis ) ) ;
3431 errors . registerRedisErrorLogger ( this . redis , { role : 'imap-notifier' } ) ; // if separate redis will add the appropriate role
3532 this . counters = counters ( this . redis ) ;
36- this . _userRegistryState = USER_REGISTRY_STATE ;
37- this . workerId = this . _userRegistryState . workerId ;
38- this . _workerChannel = this . _userRegistryState . workerChannel ;
33+ this . settingsHandler = options . settingsHandler || new SettingsHandler ( { db : this . database } ) ;
34+ this . _userRegistryState = {
35+ counts : new Map ( ) ,
36+ workerId : WORKER_ID ,
37+ workerChannel : `${ WORKER_CHANNEL_PREFIX } ${ WORKER_ID } ` ,
38+ useWorkerWdEvents : false ,
39+ useWorkerWdEventsUpdated : 0 ,
40+ useWorkerWdEventsPending : null ,
41+ timer : null ,
42+ redis : null
43+ } ;
3944
4045 this . logger = options . logger || {
4146 info : log . silly . bind ( log , 'IMAP' ) ,
@@ -97,46 +102,39 @@ class ImapNotifier extends EventEmitter {
97102 } ;
98103
99104 this . subscriber . on ( 'message' , ( channel , message ) => {
100- if ( channel === this . _workerChannel ) {
101- let data ;
102- // if e present at beginning, check if p also is present
103- // if no p -> no json parse
104- // if p -> json parse ONLY p
105- // if e not in beginning but p is -> json parse whole
106-
107- let needFullParse = true ;
108-
109- if ( message . length === 32 && message [ 2 ] === 'e' && message [ 5 ] === '"' && message [ 6 + 24 ] === '"' ) {
110- // there is only e, no p -> no need for full parse
111- needFullParse = false ;
112- }
105+ if ( channel !== this . _userRegistryState . workerChannel && channel !== LEGACY_CHANNEL ) {
106+ return ;
107+ }
113108
114- if ( ! needFullParse ) {
115- // get e and continue
116- data = { e : message . slice ( 6 , 6 + 24 ) } ;
117- } else {
118- // full parse
119- try {
120- data = JSON . parse ( message ) ;
121- } catch ( E ) {
122- return ;
123- }
109+ let data ;
110+ let needFullParse = ! ( message . length === 32 && message [ 2 ] === 'e' && message [ 5 ] === '"' && message [ 6 + 24 ] === '"' ) ;
111+
112+ if ( ! needFullParse ) {
113+ // get e and continue
114+ data = { e : message . slice ( 6 , 6 + 24 ) } ;
115+ } else {
116+ // full parse
117+ try {
118+ data = JSON . parse ( message ) ;
119+ } catch ( E ) {
120+ return ;
124121 }
122+ }
125123
126- if ( this . _listeners . _events [ data . e ] ?. length > 0 ) {
127- // do not schedule or fire/emit empty events
128- if ( data . e && ! data . p ) {
129- // events without payload are scheduled, these are notifications about changes in journal
130- scheduleDataEvent ( data . e ) ;
131- } else if ( data . e ) {
132- // events with payload are triggered immediately, these are actions for doing something
133- this . _listeners . emit ( data . e , data . p ) ;
134- }
124+ if ( this . _listeners . _events [ data . e ] ?. length > 0 ) {
125+ // do not schedule or fire/emit empty events
126+ if ( data . e && ! data . p ) {
127+ // events without payload are scheduled, these are notifications about changes in journal
128+ scheduleDataEvent ( data . e ) ;
129+ } else if ( data . e ) {
130+ // events with payload are triggered immediately, these are actions for doing something
131+ this . _listeners . emit ( data . e , data . p ) ;
135132 }
136133 }
137134 } ) ;
138135
139- this . subscriber . subscribe ( this . _workerChannel ) ;
136+ this . subscriber . subscribe ( this . _userRegistryState . workerChannel ) ;
137+ this . subscriber . subscribe ( LEGACY_CHANNEL ) ;
140138
141139 if ( ! this . _userRegistryState . timer ) {
142140 this . _userRegistryState . redis = this . redis ;
@@ -329,27 +327,65 @@ class ImapNotifier extends EventEmitter {
329327 return ;
330328 }
331329 setImmediate ( ( ) => {
332- let data = JSON . stringify ( {
333- e : user . toString ( ) ,
330+ const userId = user . toString ( ) ;
331+ const data = JSON . stringify ( {
332+ e : userId ,
334333 p : payload
335334 } ) ;
336- let key = this . _getUserWorkersKey ( user . toString ( ) ) ;
337- let minScore = Date . now ( ) - USER_WORKERS_TTL_MS ;
338- // send message only to workers that actually service this user
339- this . redis
340- . zrangebyscore ( key , minScore , '+inf' ) // get active workers
341- . then ( workers => {
342- if ( ! workers || ! workers . length ) {
343- return ;
335+ this . _usesWorkerWdEvents ( )
336+ . then ( useWorkerWdEvents => {
337+ if ( ! useWorkerWdEvents ) {
338+ return this . redis . publish ( LEGACY_CHANNEL , data ) ;
344339 }
345- let pipeline = this . redis . pipeline ( ) ;
346- workers . forEach ( workerId => pipeline . publish ( this . _getWorkerChannel ( workerId ) , data ) ) ;
347- return pipeline . exec ( ) ;
340+
341+ const key = this . _getUserWorkersKey ( userId ) ;
342+ const minScore = Date . now ( ) - USER_WORKERS_TTL_MS ;
343+ // send message only to workers that actually service this user
344+ return this . redis . zrangebyscore ( key , minScore , '+inf' ) . then ( workers => {
345+ if ( ! workers || ! workers . length ) {
346+ return ;
347+ }
348+
349+ const pipeline = this . redis . pipeline ( ) ;
350+ workers . forEach ( workerId => pipeline . publish ( this . _getWorkerChannel ( workerId ) , data ) ) ;
351+ return pipeline . exec ( ) ;
352+ } ) ;
348353 } )
349354 . catch ( ( ) => false ) ;
350355 } ) ;
351356 }
352357
358+ _usesWorkerWdEvents ( ) {
359+ // Use promise memoization
360+ if ( this . _userRegistryState . useWorkerWdEventsUpdated > Date . now ( ) - USE_WORKER_WD_EVENTS_CACHE_TTL ) {
361+ return Promise . resolve ( this . _userRegistryState . useWorkerWdEvents ) ;
362+ }
363+
364+ // return memoized promise
365+ if ( this . _userRegistryState . useWorkerWdEventsPending ) {
366+ return this . _userRegistryState . useWorkerWdEventsPending ;
367+ }
368+
369+ // Return settingsHandler async request promise
370+ this . _userRegistryState . useWorkerWdEventsPending = this . settingsHandler
371+ . get ( USE_WORKER_WD_EVENTS_SETTING , { default : false } )
372+ . then ( useWorkerWdEvents => {
373+ this . _userRegistryState . useWorkerWdEvents = ! ! useWorkerWdEvents ;
374+ this . _userRegistryState . useWorkerWdEventsUpdated = Date . now ( ) ;
375+ return this . _userRegistryState . useWorkerWdEvents ;
376+ } )
377+ . catch ( ( ) => {
378+ this . _userRegistryState . useWorkerWdEventsUpdated = Date . now ( ) ;
379+ return this . _userRegistryState . useWorkerWdEvents ;
380+ } )
381+ . finally ( ( ) => {
382+ this . _userRegistryState . useWorkerWdEventsPending = null ;
383+ } ) ;
384+
385+ // return promise
386+ return this . _userRegistryState . useWorkerWdEventsPending ;
387+ }
388+
353389 _getUserWorkersKey ( userId ) {
354390 return `${ USER_WORKERS_PREFIX } ${ userId } ` ;
355391 }
@@ -384,15 +420,15 @@ class ImapNotifier extends EventEmitter {
384420 let key = this . _getUserWorkersKey ( userId ) ;
385421 this . redis
386422 . pipeline ( )
387- . zadd ( key , Date . now ( ) , this . workerId )
423+ . zadd ( key , Date . now ( ) , this . _userRegistryState . workerId )
388424 . expire ( key , USER_WORKERS_TTL + 60 )
389425 . exec ( )
390426 . catch ( ( ) => false ) ;
391427 }
392428
393429 _unregisterUserWorker ( userId ) {
394430 let key = this . _getUserWorkersKey ( userId ) ;
395- this . redis . zrem ( key , this . workerId ) . catch ( ( ) => false ) ;
431+ this . redis . zrem ( key , this . _userRegistryState . workerId ) . catch ( ( ) => false ) ;
396432 }
397433
398434 _refreshUserWorkers ( ) {
@@ -408,7 +444,7 @@ class ImapNotifier extends EventEmitter {
408444 for ( let userId of counts . keys ( ) ) {
409445 let key = this . _getUserWorkersKey ( userId ) ;
410446 pipeline . zremrangebyscore ( key , 0 , minScore ) ;
411- pipeline . zadd ( key , now , this . workerId ) ;
447+ pipeline . zadd ( key , now , this . _userRegistryState . workerId ) ;
412448 pipeline . expire ( key , USER_WORKERS_TTL + 60 ) ;
413449 }
414450 pipeline . exec ( ) . catch ( ( ) => false ) ;
0 commit comments