@@ -451,75 +451,78 @@ function loadJournalStream(db, res, user, lastEventId, done, onEntry) {
451451 let changedKeywords = new Set ( ) ;
452452 let flaggedChanged = false ;
453453
454- let emitFlaggedCounter = next => {
454+ let emitFlaggedCounter = async next => {
455455 if ( ! flaggedChanged ) {
456456 return next ( ) ;
457457 }
458458
459- Promise . all ( [ tools . getFlaggedCounter ( db , user ) , tools . getFlaggedCounter ( db , user , 'unseen' ) ] )
460- . then ( ( [ total , unseen ] ) => {
461- res . write (
462- formatJournalData ( {
463- command : 'FLAGGED_COUNTER' ,
464- _id : lastEventId ,
465- total,
466- unseen
467- } )
468- ) ;
469- next ( ) ;
470- } )
471- . catch ( ( ) => next ( ) ) ;
459+ try {
460+ const [ total , unseen ] = await Promise . all ( [ tools . getFlaggedCounter ( db , user ) , tools . getFlaggedCounter ( db , user , 'unseen' ) ] ) ;
461+ res . write (
462+ formatJournalData ( {
463+ command : 'FLAGGED_COUNTER' ,
464+ _id : lastEventId ,
465+ total,
466+ unseen
467+ } )
468+ ) ;
469+ } catch {
470+ // ignore
471+ }
472+ next ( ) ;
472473 } ;
473474
474- let emitKeywordCounters = next => {
475+ let emitKeywordCounters = async next => {
475476 if ( ! changedKeywords . size ) {
476477 return next ( ) ;
477478 }
478479
479- const userKey = user . toString ( ) ;
480- Promise . all (
481- [ ...changedKeywords ] . map ( async keyword => {
482- const isCached = await db . redis . exists ( `kw:total:${ userKey } :${ keyword } ` ) ;
483- return isCached ? keyword : null ;
484- } )
485- )
486- . then ( async results => {
487- const toEmit = results . filter ( Boolean ) ;
488- if ( ! toEmit . length ) {
489- return next ( ) ;
490- }
491-
492- const keywordResults = await Promise . all (
493- toEmit . map ( async keyword => {
494- let total , unseen ;
495- try {
496- total = await tools . getKeywordCounter ( db , user , keyword ) ;
497- } catch {
498- total = 0 ;
499- }
500- try {
501- unseen = await tools . getKeywordCounter ( db , user , keyword , 'unseen' ) ;
502- } catch {
503- unseen = 0 ;
504- }
505- return { keyword, total, unseen } ;
506- } )
507- ) ;
480+ try {
481+ const userKey = user . toString ( ) ;
482+ const cachedResults = await Promise . all (
483+ [ ...changedKeywords ] . map ( async keyword => {
484+ const isCached = await db . redis . exists ( `kw:total:${ userKey } :${ keyword } ` ) ;
485+ return isCached ? keyword : null ;
486+ } )
487+ ) ;
488+
489+ const toEmit = cachedResults . filter ( Boolean ) ;
490+ if ( ! toEmit . length ) {
491+ return next ( ) ;
492+ }
508493
509- for ( const { keyword, total, unseen } of keywordResults ) {
510- let keywordEntry = {
511- command : 'KEYWORD_COUNTERS' ,
512- _id : lastEventId ,
513- keyword,
514- total,
515- unseen
516- } ;
517- res . write ( formatJournalData ( keywordEntry ) ) ;
518- onEntry ( 'keyword-counters' , keywordEntry ) ;
519- }
520- next ( ) ;
521- } )
522- . catch ( ( ) => next ( ) ) ;
494+ const keywordResults = await Promise . all (
495+ toEmit . map ( async keyword => {
496+ let total , unseen ;
497+ try {
498+ total = await tools . getKeywordCounter ( db , user , keyword ) ;
499+ } catch {
500+ total = 0 ;
501+ }
502+ try {
503+ unseen = await tools . getKeywordCounter ( db , user , keyword , 'unseen' ) ;
504+ } catch {
505+ unseen = 0 ;
506+ }
507+ return { keyword, total, unseen } ;
508+ } )
509+ ) ;
510+
511+ for ( const { keyword, total, unseen } of keywordResults ) {
512+ let keywordEntry = {
513+ command : 'KEYWORD_COUNTERS' ,
514+ _id : lastEventId ,
515+ keyword,
516+ total,
517+ unseen
518+ } ;
519+ res . write ( formatJournalData ( keywordEntry ) ) ;
520+ onEntry ( 'keyword-counters' , keywordEntry ) ;
521+ }
522+ } catch {
523+ // ignore
524+ }
525+ next ( ) ;
523526 } ;
524527
525528 let cursor = db . database . collection ( 'journal' ) . find ( query ) . sort ( { _id : 1 } ) ;
@@ -612,6 +615,26 @@ function loadJournalStream(db, res, user, lastEventId, done, onEntry) {
612615 break ;
613616 }
614617
618+ let writeEntryAndContinue = ( ) => {
619+ try {
620+ let data = formatJournalData ( e ) ;
621+ res . write ( data ) ;
622+ onEntry ( 'replay' , e ) ;
623+ } catch ( err ) {
624+ log . error (
625+ 'API' ,
626+ 'action=updates-event-write-fail user=%s event=%s payload=%s error=%s' ,
627+ user . toString ( ) ,
628+ formatLogValue ( e . command ) ,
629+ stringifyJournalPayload ( e ) ,
630+ err . stack || err
631+ ) ;
632+ }
633+
634+ processed ++ ;
635+ return setImmediate ( processNext ) ;
636+ } ;
637+
615638 for ( const keyword of [ ...( e . keywords ?? [ ] ) , ...( e . addedKeywords ?? [ ] ) , ...( e . removedKeywords ?? [ ] ) ] ) {
616639 changedKeywords . add ( keyword ) ;
617640 }
@@ -624,23 +647,7 @@ function loadJournalStream(db, res, user, lastEventId, done, onEntry) {
624647 }
625648 }
626649
627- try {
628- let data = formatJournalData ( e ) ;
629- res . write ( data ) ;
630- onEntry ( 'replay' , e ) ;
631- } catch ( err ) {
632- log . error (
633- 'API' ,
634- 'action=updates-event-write-fail user=%s event=%s payload=%s error=%s' ,
635- user . toString ( ) ,
636- formatLogValue ( e . command ) ,
637- stringifyJournalPayload ( e ) ,
638- err . stack || err
639- ) ;
640- }
641-
642- processed ++ ;
643- return setImmediate ( processNext ) ;
650+ writeEntryAndContinue ( ) ;
644651 } ) ;
645652 } ;
646653
0 commit comments