@@ -34,10 +34,12 @@ const logReportError = throttled(logger, 'error', 30000);
3434
3535/** Reports Node.js runtime metrics via gRPC MeterReportService (Go/Python-compatible pipeline). */
3636export default class MeterSender implements BootService , GRPCChannelListener {
37- private reporterClient ! : MeterReportServiceClient ;
37+ private status = GRPCChannelStatus . DISCONNECT ;
38+ private reporterClient ?: MeterReportServiceClient ;
3839 private readonly buffer : RuntimeSnapshot [ ] = [ ] ;
3940 private collectTimer ?: NodeJS . Timeout ;
4041 private reportTimer ?: NodeJS . Timeout ;
42+ private reporting ?: Promise < void > ;
4143
4244 private collector ! : RuntimeMetricsCollector ;
4345
@@ -57,9 +59,8 @@ export default class MeterSender implements BootService, GRPCChannelListener {
5759 }
5860
5961 statusChanged ( status : GRPCChannelStatus ) : void {
60- if ( status === GRPCChannelStatus . CONNECTED ) {
61- this . reporterClient = this . createReporterClient ( ) ;
62- }
62+ this . status = status ;
63+ this . reporterClient = status === GRPCChannelStatus . CONNECTED ? this . createReporterClient ( ) : undefined ;
6364 }
6465
6566 private createReporterClient ( ) : MeterReportServiceClient {
@@ -70,20 +71,15 @@ export default class MeterSender implements BootService, GRPCChannelListener {
7071 ) ;
7172 }
7273
73- get isConnected ( ) : boolean {
74- return ServiceManager . INSTANCE . findService ( GRPCChannelManager ) ! . isConnected ( ) ;
75- }
76-
7774 private startTimers ( ) : void {
7875 this . collectTimer = setInterval (
7976 ( ) => this . collectSample ( ) ,
8077 config . runtimeMetricsCollectPeriod || 1000 ,
8178 ) as NodeJS . Timeout ;
8279 this . collectTimer . unref ( ) ;
83- this . reportTimer = setInterval (
84- ( ) => this . reportBufferedMetrics ( ) ,
85- config . runtimeMetricsReportPeriod || 1000 ,
86- ) as NodeJS . Timeout ;
80+ this . reportTimer = setInterval ( ( ) => {
81+ void this . reportBufferedMetrics ( ) ;
82+ } , config . runtimeMetricsReportPeriod || 1000 ) as NodeJS . Timeout ;
8783 this . reportTimer . unref ( ) ;
8884 }
8985
@@ -95,56 +91,68 @@ export default class MeterSender implements BootService, GRPCChannelListener {
9591 this . buffer . push ( this . collector . sample ( ) ) ;
9692 }
9793
98- private reportBufferedMetrics ( callback ?: ( ) => void ) : void {
99- try {
100- if ( this . buffer . length === 0 || ! this . isConnected || ! this . reporterClient ) {
101- if ( callback ) callback ( ) ;
102- return ;
103- }
94+ private reportBufferedMetrics ( ) : Promise < void > {
95+ if ( this . reporting ) {
96+ return this . reporting ;
97+ }
10498
105- const snapshots = this . buffer . splice ( 0 , this . buffer . length ) ;
106- const stream = this . reporterClient . collect (
107- new grpc . Metadata ( ) ,
108- { deadline : Date . now ( ) + ( config . traceTimeout || 10000 ) } ,
109- ( error : grpc . ServiceError | null ) => {
110- if ( error ) {
111- logReportError ( 'Failed to report runtime meter data' , error ) ;
112- ServiceManager . INSTANCE . findService ( GRPCChannelManager ) ! . reportError ( error ) ;
113- }
114- if ( callback ) callback ( ) ;
115- } ,
116- ) ;
99+ this . reporting = this . doReportBufferedMetrics ( ) . finally ( ( ) => {
100+ this . reporting = undefined ;
101+ } ) ;
102+
103+ return this . reporting ;
104+ }
117105
106+ private doReportBufferedMetrics ( ) : Promise < void > {
107+ return new Promise ( ( resolve ) => {
118108 try {
119- let metadataWritten = false ;
120- const timestamp = Date . now ( ) ;
121- for ( const snapshot of snapshots ) {
122- for ( const meterData of this . collector . toMeterData ( snapshot ) ) {
123- if ( ! metadataWritten ) {
124- meterData
125- . setService ( config . serviceName ! )
126- . setServiceinstance ( config . serviceInstance ! )
127- . setTimestamp ( timestamp ) ;
128- metadataWritten = true ;
109+ if ( this . buffer . length === 0 || this . status !== GRPCChannelStatus . CONNECTED || ! this . reporterClient ) {
110+ resolve ( ) ;
111+ return ;
112+ }
113+
114+ const snapshots = this . buffer . splice ( 0 , this . buffer . length ) ;
115+ const stream = this . reporterClient . collect (
116+ new grpc . Metadata ( ) ,
117+ { deadline : Date . now ( ) + ( config . traceTimeout || 10000 ) } ,
118+ ( error : grpc . ServiceError | null ) => {
119+ if ( error ) {
120+ logReportError ( 'Failed to report runtime meter data' , error ) ;
121+ ServiceManager . INSTANCE . findService ( GRPCChannelManager ) ! . reportError ( error ) ;
122+ }
123+ resolve ( ) ;
124+ } ,
125+ ) ;
126+
127+ try {
128+ let metadataWritten = false ;
129+ const timestamp = Date . now ( ) ;
130+ for ( const snapshot of snapshots ) {
131+ for ( const meterData of this . collector . toMeterData ( snapshot ) ) {
132+ if ( ! metadataWritten ) {
133+ meterData
134+ . setService ( config . serviceName ! )
135+ . setServiceinstance ( config . serviceInstance ! )
136+ . setTimestamp ( timestamp ) ;
137+ metadataWritten = true ;
138+ }
139+ stream . write ( meterData ) ;
129140 }
130- stream . write ( meterData ) ;
131141 }
142+ } finally {
143+ stream . end ( ) ;
132144 }
133- } finally {
134- stream . end ( ) ;
145+ } catch ( error ) {
146+ logReportError ( 'Failed to report runtime meter data' , error ) ;
147+ ServiceManager . INSTANCE . findService ( GRPCChannelManager ) ! . reportError ( error ) ;
148+ resolve ( ) ;
135149 }
136- } catch ( error ) {
137- logReportError ( 'Failed to report runtime meter data' , error ) ;
138- ServiceManager . INSTANCE . findService ( GRPCChannelManager ) ! . reportError ( error ) ;
139- if ( callback ) callback ( ) ;
140- }
150+ } ) ;
141151 }
142152
143153 flush ( ) : Promise < any > | null {
144- return new Promise ( ( resolve ) => {
145- this . collectSample ( ) ;
146- this . reportBufferedMetrics ( resolve ) ;
147- } ) ;
154+ this . collectSample ( ) ;
155+ return this . reportBufferedMetrics ( ) ;
148156 }
149157
150158 shutdown ( ) : void {
@@ -156,6 +164,8 @@ export default class MeterSender implements BootService, GRPCChannelListener {
156164 clearInterval ( this . reportTimer ) ;
157165 this . reportTimer = undefined ;
158166 }
167+ this . reporting = undefined ;
168+ this . reporterClient = undefined ;
159169 this . buffer . length = 0 ;
160170 this . collector . destroy ( ) ;
161171 logger . info ( 'MeterSender destroyed and resources cleaned up' ) ;
0 commit comments