MakeIngesterClient makes a new IngesterClient
(addr string, cfg Config, useStreamConnection bool)
| 140 | |
| 141 | // MakeIngesterClient makes a new IngesterClient |
| 142 | func MakeIngesterClient(addr string, cfg Config, useStreamConnection bool) (HealthAndIngesterClient, error) { |
| 143 | unaryClientInterceptor, streamClientInterceptor := grpcclient.Instrument(ingesterClientRequestDuration) |
| 144 | if useStreamConnection { |
| 145 | unaryClientInterceptor, streamClientInterceptor = grpcclient.InstrumentReusableStream(ingesterClientRequestDuration) |
| 146 | } |
| 147 | dialOpts, err := cfg.GRPCClientConfig.DialOption(unaryClientInterceptor, streamClientInterceptor) |
| 148 | if err != nil { |
| 149 | return nil, err |
| 150 | } |
| 151 | conn, err := grpc.NewClient(addr, dialOpts...) |
| 152 | if err != nil { |
| 153 | return nil, err |
| 154 | } |
| 155 | c := &closableHealthAndIngesterClient{ |
| 156 | IngesterClient: NewIngesterClient(conn), |
| 157 | HealthClient: grpc_health_v1.NewHealthClient(conn), |
| 158 | conn: conn, |
| 159 | addr: addr, |
| 160 | maxInflightPushRequests: cfg.MaxInflightPushRequests, |
| 161 | inflightPushRequests: ingesterClientInflightPushRequests, |
| 162 | } |
| 163 | if useStreamConnection { |
| 164 | streamCtx, streamCancel := context.WithCancel(context.Background()) |
| 165 | err = c.Run(make(chan *streamWriteJob, INGESTER_CLIENT_STREAM_WORKER_COUNT), streamCtx, streamCancel) |
| 166 | if err != nil { |
| 167 | return nil, err |
| 168 | } |
| 169 | } |
| 170 | return c, nil |
| 171 | } |
| 172 | |
| 173 | func (c *closableHealthAndIngesterClient) Close() error { |
| 174 | c.inflightPushRequests.DeleteLabelValues(c.addr) |