1515 along with this program. If not, see <https://www.gnu.org/licenses/>.
1616*/
1717
18- import './init.js' ;
19-
2018import { metrics } from '@balena/node-metrics-gatherer' ;
2119import cluster from 'cluster' ;
2220import express from 'express' ;
@@ -35,15 +33,18 @@ import {
3533 VPN_VERBOSE_LOGS ,
3634} from './utils/config.js' ;
3735
38- import proxyWorker from './proxy-worker.js' ;
39- import vpnWorker from './vpn-worker.js' ;
40- import { intVar } from '@balena/env-parsing' ;
41- import { describeMetrics , Metrics } from './utils/metrics.js' ;
36+ import { describePrimaryMetrics , Metrics } from './utils/metrics.js' ;
4237import { service } from './utils/service.js' ;
4338
39+ if ( ! cluster . isPrimary ) {
40+ throw new Error (
41+ 'init-primary should only be imported by the primary cluster process' ,
42+ ) ;
43+ }
44+
4445const masterLogger = getLogger ( 'master' ) ;
4546
46- describeMetrics ( ) ;
47+ describePrimaryMetrics ( ) ;
4748
4849export interface BitrateMessage {
4950 type : 'bitrate' ;
@@ -54,90 +55,91 @@ export interface BitrateMessage {
5455 } ;
5556}
5657
57- if ( cluster . isPrimary ) {
58- interface WorkerMetric {
59- rxBitrate : Array < BitrateMessage [ 'data' ] [ 'rxBitrate' ] > ;
60- txBitrate : Array < BitrateMessage [ 'data' ] [ 'txBitrate' ] > ;
61- }
62- const workerMetrics = new Map < string , WorkerMetric > ( ) ;
63-
64- let verbose = VPN_VERBOSE_LOGS ;
58+ interface WorkerMetric {
59+ rxBitrate : Array < BitrateMessage [ 'data' ] [ 'rxBitrate' ] > ;
60+ txBitrate : Array < BitrateMessage [ 'data' ] [ 'txBitrate' ] > ;
61+ }
62+ const workerMetrics = new Map < string , WorkerMetric > ( ) ;
6563
66- type WorkerState = {
67- instanceId : number ;
68- finished : boolean ;
69- } ;
70- const workerStates : { [ instanceId : number ] : WorkerState } = { } ;
64+ let verbose = VPN_VERBOSE_LOGS ;
7165
72- process . on ( 'SIGUSR2' , ( ) => {
73- masterLogger . notice ( 'caught SIGUSR2, toggling log verbosity' ) ;
74- verbose = ! verbose ;
75- for ( const clusterWorker of Object . values ( cluster . workers ?? { } ) ) {
76- clusterWorker ?. send ( 'toggleVerbosity' ) ;
77- }
78- } ) ;
66+ type WorkerState = {
67+ instanceId : number ;
68+ finished : boolean ;
69+ } ;
70+ const workerStates : { [ instanceId : number ] : WorkerState } = { } ;
7971
80- process . on ( 'SIGTERM' , ( ) => {
81- masterLogger . notice ( 'received SIGTERM' ) ;
82- for ( const clusterWorker of Object . values ( cluster . workers ?? { } ) ) {
83- clusterWorker ?. send ( 'prepareShutdown' ) ;
84- }
85- masterLogger . notice (
86- `waiting ${ DEFAULT_SIGTERM_TIMEOUT } ms for workers to finish` ,
87- ) ;
88- } ) ;
72+ process . on ( 'SIGUSR2' , ( ) => {
73+ masterLogger . notice ( 'caught SIGUSR2, toggling log verbosity' ) ;
74+ verbose = ! verbose ;
75+ for ( const clusterWorker of Object . values ( cluster . workers ?? { } ) ) {
76+ clusterWorker ?. send ( 'toggleVerbosity' ) ;
77+ }
78+ } ) ;
8979
90- cluster . on ( 'message' , ( _worker , { data, type } : BitrateMessage ) => {
91- if ( type === 'bitrate' ) {
92- let workerMetric = workerMetrics . get ( data . uuid ) ;
93- if ( workerMetric == null ) {
94- workerMetric = {
95- rxBitrate : [ ] ,
96- txBitrate : [ ] ,
97- } ;
98- workerMetrics . set ( data . uuid , workerMetric ) ;
99- }
100- workerMetric . rxBitrate . push ( data . rxBitrate ) ;
101- workerMetric . txBitrate . push ( data . txBitrate ) ;
80+ process . on ( 'SIGTERM' , ( ) => {
81+ masterLogger . notice ( 'received SIGTERM' ) ;
82+ for ( const clusterWorker of Object . values ( cluster . workers ?? { } ) ) {
83+ clusterWorker ?. send ( 'prepareShutdown' ) ;
84+ }
85+ masterLogger . notice (
86+ `waiting ${ DEFAULT_SIGTERM_TIMEOUT } ms for workers to finish` ,
87+ ) ;
88+ } ) ;
89+
90+ cluster . on ( 'message' , ( _worker , { data, type } : BitrateMessage ) => {
91+ if ( type === 'bitrate' ) {
92+ let workerMetric = workerMetrics . get ( data . uuid ) ;
93+ if ( workerMetric == null ) {
94+ workerMetric = {
95+ rxBitrate : [ ] ,
96+ txBitrate : [ ] ,
97+ } ;
98+ workerMetrics . set ( data . uuid , workerMetric ) ;
10299 }
103- } ) ;
104-
105- cluster . on ( 'message' , ( _worker , msg : { type : string ; data : WorkerState } ) => {
106- const { data, type } = msg ;
107-
108- // worker finished connection draining
109- if ( type === 'drain' ) {
110- try {
111- workerStates [ data . instanceId ] = data ;
112- const drainCount = Object . keys ( workerStates ) . length ;
113- masterLogger . notice (
114- `total: ${ VPN_INSTANCE_COUNT } drained: ${ drainCount } ` ,
115- ) ;
116- for ( const key in workerStates ) {
117- if ( key != null ) {
118- const value : WorkerState = workerStates [ key ] ;
119- const workerState = value . finished ;
120- masterLogger . notice ( `instanceId:${ key } finished:${ workerState } ` ) ;
121- }
122- }
123- if ( drainCount >= VPN_INSTANCE_COUNT ) {
124- masterLogger . notice ( `all ${ drainCount } worker(s) drained` ) ;
125- process . exit ( 0 ) ;
100+ workerMetric . rxBitrate . push ( data . rxBitrate ) ;
101+ workerMetric . txBitrate . push ( data . txBitrate ) ;
102+ }
103+ } ) ;
104+
105+ cluster . on ( 'message' , ( _worker , msg : { type : string ; data : WorkerState } ) => {
106+ const { data, type } = msg ;
107+
108+ // worker finished connection draining
109+ if ( type === 'drain' ) {
110+ try {
111+ workerStates [ data . instanceId ] = data ;
112+ const drainCount = Object . keys ( workerStates ) . length ;
113+ masterLogger . notice (
114+ `total: ${ VPN_INSTANCE_COUNT } drained: ${ drainCount } ` ,
115+ ) ;
116+ for ( const key in workerStates ) {
117+ if ( key != null ) {
118+ const value : WorkerState = workerStates [ key ] ;
119+ const workerState = value . finished ;
120+ masterLogger . notice ( `instanceId:${ key } finished:${ workerState } ` ) ;
126121 }
127- } catch ( err ) {
128- masterLogger . warning ( `${ err } handling message from worker` ) ;
129122 }
130- } else {
131- return ;
123+ if ( drainCount >= VPN_INSTANCE_COUNT ) {
124+ masterLogger . notice ( `all ${ drainCount } worker(s) drained` ) ;
125+ process . exit ( 0 ) ;
126+ }
127+ } catch ( err ) {
128+ masterLogger . warning ( `${ err } handling message from worker` ) ;
132129 }
133- } ) ;
134-
135- masterLogger . notice (
136- `open-balena-vpn@${ VERSION } process started with pid=${ process . pid } ` ,
137- ) ;
138- masterLogger . debug ( 'registering as service instance...' ) ;
139- service
140- . wrap ( { ipAddress : VPN_SERVICE_ADDRESS } , async ( serviceInstance ) => {
130+ } else {
131+ return ;
132+ }
133+ } ) ;
134+
135+ masterLogger . notice (
136+ `open-balena-vpn@${ VERSION } process started with pid=${ process . pid } ` ,
137+ ) ;
138+ masterLogger . debug ( 'registering as service instance...' ) ;
139+ try {
140+ await service . wrap (
141+ { ipAddress : VPN_SERVICE_ADDRESS } ,
142+ async ( serviceInstance ) => {
141143 const serviceLogger = getLogger ( 'master' , serviceInstance . getId ( ) ) ;
142144 serviceLogger . info (
143145 `registered as service instance with id=${ serviceInstance . getId ( ) } ipAddress=${ VPN_SERVICE_ADDRESS } ` ,
@@ -232,27 +234,9 @@ if (cluster.isPrimary) {
232234 . listen ( 8080 ) ;
233235
234236 return [ app , metrics ] ;
235- } )
236- . catch ( ( err ) => {
237- console . error ( 'Error starting master:' , err ) ;
238- process . exit ( 1 ) ;
239- } ) ;
240- }
241-
242- if ( cluster . isWorker ) {
243- // Ensure the prom-client worker listener is registered by instantiating the class
244- // tslint:disable-next-line:no-unused-expression-chai
245- new prometheus . AggregatorRegistry ( ) ;
246-
247- const instanceId = intVar ( 'WORKER_ID' ) ;
248- const serviceId = intVar ( 'SERVICE_ID' ) ;
249- getLogger ( 'worker' , serviceId , instanceId ) . notice (
250- `process started with pid=${ process . pid } ` ,
237+ } ,
251238 ) ;
252- vpnWorker ( instanceId , serviceId )
253- . then ( ( ) => proxyWorker ( instanceId , serviceId ) )
254- . catch ( ( err ) => {
255- console . error ( 'Error starting worker:' , err ) ;
256- process . exit ( 1 ) ;
257- } ) ;
239+ } catch ( err ) {
240+ console . error ( 'Error starting master:' , err ) ;
241+ process . exit ( 1 ) ;
258242}
0 commit comments