-
Notifications
You must be signed in to change notification settings - Fork 11
Expand file tree
/
Copy pathsetup.js
More file actions
128 lines (112 loc) · 5.18 KB
/
Copy pathsetup.js
File metadata and controls
128 lines (112 loc) · 5.18 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
/*****
License
--------------
Copyright © 2020-2025 Mojaloop Foundation
The Mojaloop files are made available by the Mojaloop Foundation under the Apache License, Version 2.0 (the "License") and you may not use these files except in compliance with the License. You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, the Mojaloop files are distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions and limitations under the License.
Contributors
--------------
This is the official list of the Mojaloop project contributors for this file.
Names of the original copyright holders (individuals or organizations)
should be listed with a '*' in the first column. People who have
contributed from an organization can be listed under the organization
that actually holds the copyright for their contributions (see the
Mojaloop Foundation for an example). Those individuals should have
their names indented and be marked with a '-'. Email address can be added
optionally within square brackets <email>.
* Mojaloop Foundation
- Name Surname <name.surname@mojaloop.io>
* Valentin Genev <valentin.genev@modusbox.com>
* Deon Botha <deon.botha@modusbox.com>
--------------
******/
'use strict'
/**
* @module src/setup
*/
const Kafka = require('@mojaloop/central-services-shared').Util.Kafka
const Util = require('@mojaloop/central-services-stream').Util
const Consumer = Util.Consumer
const Enum = require('@mojaloop/central-services-shared').Enum
const Logger = require('@mojaloop/central-services-logger')
const Rx = require('rxjs')
const { share, filter, flatMap, catchError } = require('rxjs/operators')
const Config = require('./lib/config')
const HealthCheck = require('@mojaloop/central-services-shared').HealthCheck.HealthCheck
const { createHealthCheckServer, defaultHealthHandler } = require('@mojaloop/central-services-health')
const packageJson = require('../package.json')
const { getSubServiceHealthBroker } = require('./lib/healthCheck/subServiceHealth')
const Observables = require('./observables')
const { initializeCache } = Observables.TraceObservable
const setup = async () => {
Rx.config.onUnhandledError = (err) => {
console.warn(err)
}
Rx.config.useDeprecatedNextContext = true
await registerEventHandler()
await initializeCache(Config.CACHE_CONFIG)
const topicName = Kafka.transformGeneralTopicName(Config.KAFKA_CONFIG.TOPIC_TEMPLATES.GENERAL_TOPIC_TEMPLATE.TEMPLATE, Enum.Events.Event.Action.EVENT)
const consumer = Consumer.getConsumer(topicName)
const healthCheck = new HealthCheck(packageJson, [
getSubServiceHealthBroker
])
await createHealthCheckServer(Config.PORT, defaultHealthHandler(healthCheck))
const topicObservable = Rx.Observable.create((observer) => {
consumer.on('message', async (message) => {
Logger.debug(`Central-Event-Processor :: Topic ${topicName} :: Payload: \n${JSON.stringify(message.value, null, 2)}`)
observer.next({ message })
if (!Consumer.isConsumerAutoCommitEnabled(topicName)) {
consumer.commitMessageSync(message)
}
})
})
const sharedMessageObservable = topicObservable.pipe(share(), catchError(e => { return Rx.onErrorResumeNext(sharedMessageObservable) }))
sharedMessageObservable.subscribe(async props => {
Observables.elasticsearchClientObservable(props).subscribe({
next: v => Logger.debug(v),
error: (e) => Logger.error(e.stack),
completed: () => Logger.debug('elastic API log completed')
})
})
const tracingObservable = sharedMessageObservable.pipe(
filter(props => props.message.value.metadata.event.type === 'trace'),
flatMap(Observables.TraceObservable.extractContextObservable),
flatMap(Observables.TraceObservable.cacheSpanContextObservable),
flatMap(Observables.TraceObservable.findLastSpanObservable),
flatMap(Observables.TraceObservable.recreateTraceObservable),
flatMap(Observables.TraceObservable.sendTraceToApmObservable),
catchError(e => { return Rx.onErrorResumeNext(tracingObservable) }))
tracingObservable.subscribe({
next: traceId => {
Logger.debug(`traceId ${traceId} sent to APM`)
},
error: (e) => Logger.error(e.stack),
completed: () => Logger.debug('trace info sent')
})
}
/**
* @function registerEventHandler
*
* @description This is used to register the handler for the Event topic according to a specified Kafka congfiguration
*
* @returns true
* @throws {Error} - if handler failed to create
*/
const registerEventHandler = async () => {
try {
const EventHandler = {
topicName: Kafka.transformGeneralTopicName(Config.KAFKA_CONFIG.TOPIC_TEMPLATES.GENERAL_TOPIC_TEMPLATE.TEMPLATE, Enum.Events.Event.Action.EVENT),
config: Kafka.getKafkaConfig(Config.KAFKA_CONFIG, Enum.Kafka.Config.CONSUMER, Enum.Events.Event.Type.EVENT.toUpperCase())
}
EventHandler.config.rdkafkaConf['client.id'] = EventHandler.topicName
await Consumer.createHandler(EventHandler.topicName, EventHandler.config)
return true
} catch (e) {
Logger.error(e.stack)
throw e
}
}
module.exports = {
setup
}