-
Notifications
You must be signed in to change notification settings - Fork 2.1k
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Merge branch 'main' into logger_duration_ms
- Loading branch information
Showing
10 changed files
with
437 additions
and
19 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -4,7 +4,7 @@ | |
|-----------------------|------------------------|--------------------------|----------------------------------------------|---------------| | ||
| Bartłomiej Płotka | [email protected] | `@bwplotka` | [@bwplotka](https://github.com/bwplotka) | Red Hat | | ||
| Frederic Branczyk | [email protected] | `@brancz` | [@brancz](https://github.com/brancz) | Polar Signals | | ||
| Giedrius Statkevičius | [email protected] | `@Giedrius Statkevičius` | [@GiedriusS](https://github.com/GiedriusS) | AdForm | | ||
| Giedrius Statkevičius | [email protected] | `@Giedrius Statkevičius` | [@GiedriusS](https://github.com/GiedriusS) | Vinted | | ||
| Kemal Akkoyun | [email protected] | `@kakkoyun` | [@kakkoyun](https://github.com/kakkoyun) | Polar Signals | | ||
| Lucas Servén Marín | [email protected] | `@squat` | [@squat](https://github.com/squat) | Red Hat | | ||
| Prem Saraswat | [email protected] | `@Prem Saraswat` | [@onprem](https://github.com/onprem) | Red Hat | | ||
|
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,114 @@ | ||
// Copyright (c) The Thanos Authors. | ||
// Licensed under the Apache License 2.0. | ||
|
||
package memcache | ||
|
||
import ( | ||
"context" | ||
"fmt" | ||
"sync" | ||
"time" | ||
|
||
"github.com/go-kit/kit/log" | ||
"github.com/go-kit/kit/log/level" | ||
"github.com/prometheus/client_golang/prometheus" | ||
"github.com/prometheus/client_golang/prometheus/promauto" | ||
"github.com/thanos-io/thanos/pkg/errutil" | ||
"github.com/thanos-io/thanos/pkg/extprom" | ||
) | ||
|
||
// Provider is a stateful cache for asynchronous memcached auto-discovery resolution. It provides a way to resolve | ||
// addresses and obtain them. | ||
type Provider struct { | ||
sync.RWMutex | ||
resolver Resolver | ||
clusterConfigs map[string]*clusterConfig | ||
logger log.Logger | ||
|
||
configVersion *extprom.TxGaugeVec | ||
resolvedAddresses *extprom.TxGaugeVec | ||
resolverFailuresCount prometheus.Counter | ||
resolverLookupsCount prometheus.Counter | ||
} | ||
|
||
func NewProvider(logger log.Logger, reg prometheus.Registerer, dialTimeout time.Duration) *Provider { | ||
p := &Provider{ | ||
resolver: &memcachedAutoDiscovery{dialTimeout: dialTimeout}, | ||
clusterConfigs: map[string]*clusterConfig{}, | ||
configVersion: extprom.NewTxGaugeVec(reg, prometheus.GaugeOpts{ | ||
Name: "auto_discovery_config_version", | ||
Help: "The current auto discovery config version", | ||
}, []string{"addr"}), | ||
resolvedAddresses: extprom.NewTxGaugeVec(reg, prometheus.GaugeOpts{ | ||
Name: "auto_discovery_resolved_addresses", | ||
Help: "The number of memcached nodes found via auto discovery", | ||
}, []string{"addr"}), | ||
resolverLookupsCount: promauto.With(reg).NewCounter(prometheus.CounterOpts{ | ||
Name: "auto_discovery_total", | ||
Help: "The number of memcache auto discovery attempts", | ||
}), | ||
resolverFailuresCount: promauto.With(reg).NewCounter(prometheus.CounterOpts{ | ||
Name: "auto_discovery_failures_total", | ||
Help: "The number of memcache auto discovery failures", | ||
}), | ||
logger: logger, | ||
} | ||
return p | ||
} | ||
|
||
// Resolve stores a list of nodes auto-discovered from the provided addresses. | ||
func (p *Provider) Resolve(ctx context.Context, addresses []string) error { | ||
clusterConfigs := map[string]*clusterConfig{} | ||
errs := errutil.MultiError{} | ||
|
||
for _, address := range addresses { | ||
clusterConfig, err := p.resolver.Resolve(ctx, address) | ||
p.resolverLookupsCount.Inc() | ||
|
||
if err != nil { | ||
level.Warn(p.logger).Log( | ||
"msg", "failed to perform auto-discovery for memcached", | ||
"address", address, | ||
) | ||
errs.Add(err) | ||
p.resolverFailuresCount.Inc() | ||
|
||
// Use cached values. | ||
p.RLock() | ||
clusterConfigs[address] = p.clusterConfigs[address] | ||
p.RUnlock() | ||
} else { | ||
clusterConfigs[address] = clusterConfig | ||
} | ||
} | ||
|
||
p.Lock() | ||
defer p.Unlock() | ||
|
||
p.resolvedAddresses.ResetTx() | ||
p.configVersion.ResetTx() | ||
for address, config := range clusterConfigs { | ||
p.resolvedAddresses.WithLabelValues(address).Set(float64(len(config.nodes))) | ||
p.configVersion.WithLabelValues(address).Set(float64(config.version)) | ||
} | ||
p.resolvedAddresses.Submit() | ||
p.configVersion.Submit() | ||
|
||
p.clusterConfigs = clusterConfigs | ||
|
||
return errs.Err() | ||
} | ||
|
||
// Addresses returns the latest addresses present in the Provider. | ||
func (p *Provider) Addresses() []string { | ||
p.RLock() | ||
defer p.RUnlock() | ||
|
||
var result []string | ||
for _, config := range p.clusterConfigs { | ||
for _, node := range config.nodes { | ||
result = append(result, fmt.Sprintf("%s:%d", node.dns, node.port)) | ||
} | ||
} | ||
return result | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,81 @@ | ||
// Copyright (c) The Thanos Authors. | ||
// Licensed under the Apache License 2.0. | ||
|
||
package memcache | ||
|
||
import ( | ||
"context" | ||
"sort" | ||
"testing" | ||
"time" | ||
|
||
"github.com/pkg/errors" | ||
|
||
"github.com/go-kit/kit/log" | ||
"github.com/thanos-io/thanos/pkg/testutil" | ||
) | ||
|
||
func TestProviderUpdatesAddresses(t *testing.T) { | ||
ctx := context.TODO() | ||
clusters := []string{"memcached-cluster-1", "memcached-cluster-2"} | ||
provider := NewProvider(log.NewNopLogger(), nil, 5*time.Second) | ||
resolver := mockResolver{ | ||
configs: map[string]*clusterConfig{ | ||
"memcached-cluster-1": {nodes: []node{{dns: "dns-1", ip: "ip-1", port: 11211}}}, | ||
"memcached-cluster-2": {nodes: []node{{dns: "dns-2", ip: "ip-2", port: 8080}}}, | ||
}, | ||
} | ||
provider.resolver = &resolver | ||
|
||
testutil.Ok(t, provider.Resolve(ctx, clusters)) | ||
addresses := provider.Addresses() | ||
testutil.Equals(t, []string{"dns-1:11211", "dns-2:8080"}, addresses) | ||
|
||
resolver.configs = map[string]*clusterConfig{ | ||
"memcached-cluster-1": {nodes: []node{{dns: "dns-1", ip: "ip-1", port: 11211}, {dns: "dns-3", ip: "ip-3", port: 11211}}}, | ||
"memcached-cluster-2": {nodes: []node{{dns: "dns-2", ip: "ip-2", port: 8080}}}, | ||
} | ||
|
||
testutil.Ok(t, provider.Resolve(ctx, clusters)) | ||
addresses = provider.Addresses() | ||
sort.Strings(addresses) | ||
testutil.Equals(t, []string{"dns-1:11211", "dns-2:8080", "dns-3:11211"}, addresses) | ||
} | ||
|
||
func TestProviderDoesNotUpdateAddressIfFailed(t *testing.T) { | ||
ctx := context.TODO() | ||
clusters := []string{"memcached-cluster-1", "memcached-cluster-2"} | ||
provider := NewProvider(log.NewNopLogger(), nil, 5*time.Second) | ||
resolver := mockResolver{ | ||
configs: map[string]*clusterConfig{ | ||
"memcached-cluster-1": {nodes: []node{{dns: "dns-1", ip: "ip-1", port: 11211}}}, | ||
"memcached-cluster-2": {nodes: []node{{dns: "dns-2", ip: "ip-2", port: 8080}}}, | ||
}, | ||
} | ||
provider.resolver = &resolver | ||
|
||
testutil.Ok(t, provider.Resolve(ctx, clusters)) | ||
addresses := provider.Addresses() | ||
sort.Strings(addresses) | ||
testutil.Equals(t, []string{"dns-1:11211", "dns-2:8080"}, addresses) | ||
|
||
resolver.configs = nil | ||
resolver.err = errors.New("oops") | ||
|
||
testutil.NotOk(t, provider.Resolve(ctx, clusters)) | ||
addresses = provider.Addresses() | ||
sort.Strings(addresses) | ||
testutil.Equals(t, []string{"dns-1:11211", "dns-2:8080"}, addresses) | ||
} | ||
|
||
type mockResolver struct { | ||
configs map[string]*clusterConfig | ||
err error | ||
} | ||
|
||
func (r *mockResolver) Resolve(_ context.Context, address string) (*clusterConfig, error) { | ||
if r.err != nil { | ||
return nil, r.err | ||
} | ||
return r.configs[address], nil | ||
} |
Oops, something went wrong.