mirror of
https://github.com/gravitational/teleport.git
synced 2026-09-24 16:17:11 +08:00
Always use in-memory caches (#11386)
* Always use in-memory caches This also cleans up now-useless fields and constants related to on-disk caches. * Remove the cache tombstone mechanism As we're never reopening the same cache backend twice, this is no longer useful. * Warn if a cache directory exists on disk We can't remove it automatically because we might be in the middle of an upgrade with a old version of Teleport still running.
This commit is contained in:
@@ -88,7 +88,6 @@ func SetTestTimeouts(t time.Duration) {
|
||||
defaults.ResyncInterval = t
|
||||
defaults.SessionRefreshPeriod = t
|
||||
defaults.HeartbeatCheckPeriod = t
|
||||
defaults.CachePollPeriod = t
|
||||
}
|
||||
|
||||
// TeleInstance represents an in-memory instance of a teleport
|
||||
|
||||
Vendored
-51
@@ -61,10 +61,6 @@ var (
|
||||
cacheCollectors = []prometheus.Collector{cacheEventsReceived, cacheStaleEventsReceived}
|
||||
)
|
||||
|
||||
func tombstoneKey() []byte {
|
||||
return backend.Key("cache", teleport.Version, "tombstone", "ok")
|
||||
}
|
||||
|
||||
const cacheTargetAuth string = "auth"
|
||||
|
||||
// ForAuth sets up watch configuration for the auth server
|
||||
@@ -629,9 +625,6 @@ const (
|
||||
WatcherStarted = "watcher_started"
|
||||
// WatcherFailed is emitted when event watcher has failed
|
||||
WatcherFailed = "watcher_failed"
|
||||
// TombstoneWritten is emitted if cache is closed in a healthy
|
||||
// state and successfully writes its tombstone.
|
||||
TombstoneWritten = "tombstone_written"
|
||||
// Reloading is emitted when an error occurred watching events
|
||||
// and the cache is waiting to create a new watcher
|
||||
Reloading = "reloading_cache"
|
||||
@@ -700,28 +693,6 @@ func New(config Config) (*Cache, error) {
|
||||
}
|
||||
cs.collections = collections
|
||||
|
||||
// if the ok tombstone is present, set the initial read state of the cache
|
||||
// to ok. this tombstone's presence indicates that we are dealing with an
|
||||
// on-disk cache produced by the same teleport version which gracefully shutdown
|
||||
// while in an ok state. We delete the tombstone rather than check for its
|
||||
// presence to ensure self-healing in the event that the tombstone wasn't actually
|
||||
// valid. Note that setting the cache's read state to ok does not cause us to skip
|
||||
// our normal init logic, it just means that reads against the local cache will
|
||||
// be allowed in the event the init step fails before it starts applying.
|
||||
// Note also that we aren't setting our event fanout system to an initialized state
|
||||
// or incrementing the generation counter; this cache isn't so much "healthy" as it is
|
||||
// "slightly preferable to an unreachable auth server".
|
||||
err = cs.wrapper.Delete(ctx, tombstoneKey())
|
||||
switch {
|
||||
case err == nil:
|
||||
cs.setReadOK(true)
|
||||
case trace.IsNotFound(err):
|
||||
// do nothing
|
||||
default:
|
||||
cs.Close()
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
|
||||
retry, err := utils.NewLinear(utils.LinearConfig{
|
||||
First: utils.HalfJitter(cs.MaxRetryPeriod / 10),
|
||||
Step: cs.MaxRetryPeriod / 5,
|
||||
@@ -775,10 +746,6 @@ func (c *Cache) update(ctx context.Context, retry utils.Retry) {
|
||||
c.Debugf("Cache is closing, returning from update loop.")
|
||||
// ensure that close operations have been run
|
||||
c.Close()
|
||||
// run tombstone operations in an orphaned context
|
||||
tombCtx, cancel := context.WithTimeout(context.Background(), time.Second*10)
|
||||
defer cancel()
|
||||
c.writeTombstone(tombCtx)
|
||||
}()
|
||||
timer := time.NewTimer(c.Config.WatcherInitTimeout)
|
||||
for {
|
||||
@@ -811,24 +778,6 @@ func (c *Cache) update(ctx context.Context, retry utils.Retry) {
|
||||
}
|
||||
}
|
||||
|
||||
// writeTombstone writes the cache tombstone.
|
||||
func (c *Cache) writeTombstone(ctx context.Context) {
|
||||
if !c.getReadOK() || c.generation.Load() == 0 {
|
||||
// state is unhealthy or was loaded from a previously
|
||||
// entombed state; do nothing.
|
||||
return
|
||||
}
|
||||
item := backend.Item{
|
||||
Key: tombstoneKey(),
|
||||
Value: []byte("{}"),
|
||||
}
|
||||
if _, err := c.wrapper.Create(ctx, item); err != nil {
|
||||
c.Warningf("Failed to set tombstone: %v", err)
|
||||
} else {
|
||||
c.notify(ctx, Event{Type: TombstoneWritten})
|
||||
}
|
||||
}
|
||||
|
||||
func (c *Cache) notify(ctx context.Context, event Event) {
|
||||
if c.EventsC == nil {
|
||||
return
|
||||
|
||||
Vendored
-81
@@ -659,87 +659,6 @@ func TestCompletenessReset(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestTombstones verifies that healthy caches leave tombstones
|
||||
// on closure, giving new caches the ability to start from a known
|
||||
// good state if the origin state is unavailable.
|
||||
func TestTombstones(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
const caCount = 10
|
||||
p := newTestPackWithoutCache(t)
|
||||
t.Cleanup(p.Close)
|
||||
|
||||
// put lots of CAs in the backend
|
||||
for i := 0; i < caCount; i++ {
|
||||
ca := suite.NewTestCA(types.UserCA, fmt.Sprintf("%d.example.com", i))
|
||||
require.NoError(t, p.trustS.UpsertCertAuthority(ca))
|
||||
}
|
||||
|
||||
var err error
|
||||
p.cache, err = New(ForAuth(Config{
|
||||
Context: ctx,
|
||||
Backend: p.cacheBackend,
|
||||
Events: p.eventsS,
|
||||
ClusterConfig: p.clusterConfigS,
|
||||
Provisioner: p.provisionerS,
|
||||
Trust: p.trustS,
|
||||
Users: p.usersS,
|
||||
Access: p.accessS,
|
||||
DynamicAccess: p.dynamicAccessS,
|
||||
Presence: p.presenceS,
|
||||
AppSession: p.appSessionS,
|
||||
WebSession: p.webSessionS,
|
||||
WebToken: p.webTokenS,
|
||||
Restrictions: p.restrictions,
|
||||
Apps: p.apps,
|
||||
Databases: p.databases,
|
||||
WindowsDesktops: p.windowsDesktops,
|
||||
MaxRetryPeriod: 200 * time.Millisecond,
|
||||
EventsC: p.eventsC,
|
||||
}))
|
||||
require.NoError(t, err)
|
||||
|
||||
// verify that CAs are immediately available
|
||||
cas, err := p.cache.GetCertAuthorities(ctx, types.UserCA, false)
|
||||
require.NoError(t, err)
|
||||
require.Len(t, cas, caCount)
|
||||
|
||||
require.NoError(t, p.cache.Close())
|
||||
// wait for TombstoneWritten, ignoring all other event types
|
||||
expectEvent(t, p.eventsC, TombstoneWritten)
|
||||
// simulate bad connection to auth server
|
||||
p.backend.SetReadError(trace.ConnectionProblem(nil, "backend is unavailable"))
|
||||
p.eventsS.closeWatchers()
|
||||
|
||||
p.cache, err = New(ForAuth(Config{
|
||||
Context: ctx,
|
||||
Backend: p.cacheBackend,
|
||||
Events: p.eventsS,
|
||||
ClusterConfig: p.clusterConfigS,
|
||||
Provisioner: p.provisionerS,
|
||||
Trust: p.trustS,
|
||||
Users: p.usersS,
|
||||
Access: p.accessS,
|
||||
DynamicAccess: p.dynamicAccessS,
|
||||
Presence: p.presenceS,
|
||||
AppSession: p.appSessionS,
|
||||
WebSession: p.webSessionS,
|
||||
WebToken: p.webTokenS,
|
||||
Restrictions: p.restrictions,
|
||||
Apps: p.apps,
|
||||
Databases: p.databases,
|
||||
WindowsDesktops: p.windowsDesktops,
|
||||
MaxRetryPeriod: 200 * time.Millisecond,
|
||||
EventsC: p.eventsC,
|
||||
}))
|
||||
require.NoError(t, err)
|
||||
|
||||
// verify that CAs are immediately available despite the fact
|
||||
// that the origin state was never available.
|
||||
cas, err = p.cache.GetCertAuthorities(ctx, types.UserCA, false)
|
||||
require.NoError(t, err)
|
||||
require.Len(t, cas, caCount)
|
||||
}
|
||||
|
||||
// TestInitStrategy verifies that cache uses expected init strategy
|
||||
// of serving backend state when init is taking too long.
|
||||
func TestInitStrategy(t *testing.T) {
|
||||
|
||||
@@ -47,6 +47,7 @@ import (
|
||||
"github.com/gravitational/teleport/lib"
|
||||
"github.com/gravitational/teleport/lib/backend"
|
||||
"github.com/gravitational/teleport/lib/backend/lite"
|
||||
"github.com/gravitational/teleport/lib/backend/memory"
|
||||
"github.com/gravitational/teleport/lib/backend/postgres"
|
||||
"github.com/gravitational/teleport/lib/client"
|
||||
"github.com/gravitational/teleport/lib/defaults"
|
||||
@@ -304,7 +305,12 @@ func ApplyFileConfig(fc *FileConfig, cfg *service.Config) error {
|
||||
}
|
||||
|
||||
if fc.CachePolicy.TTL != "" {
|
||||
log.Warnf("cache.ttl config option is deprecated and will be ignored, caches no longer attempt to anticipate resource expiration.")
|
||||
log.Warn("cache.ttl config option is deprecated and will be ignored, caches no longer attempt to anticipate resource expiration.")
|
||||
}
|
||||
if fc.CachePolicy.Type == memory.GetName() {
|
||||
log.Debugf("cache.type config option is explicitly set to %v.", memory.GetName())
|
||||
} else if fc.CachePolicy.Type != "" {
|
||||
log.Warn("cache.type config option is deprecated and will be ignored, caches are always in memory in this version.")
|
||||
}
|
||||
|
||||
// apply cache policy for node and proxy
|
||||
|
||||
@@ -37,7 +37,6 @@ import (
|
||||
"github.com/gravitational/teleport/lib"
|
||||
"github.com/gravitational/teleport/lib/backend"
|
||||
"github.com/gravitational/teleport/lib/backend/lite"
|
||||
"github.com/gravitational/teleport/lib/backend/memory"
|
||||
"github.com/gravitational/teleport/lib/defaults"
|
||||
"github.com/gravitational/teleport/lib/fixtures"
|
||||
"github.com/gravitational/teleport/lib/limiter"
|
||||
@@ -919,12 +918,11 @@ func TestParseCachePolicy(t *testing.T) {
|
||||
out *service.CachePolicy
|
||||
err error
|
||||
}{
|
||||
{in: &CachePolicy{EnabledFlag: "yes", TTL: "never"}, out: &service.CachePolicy{Enabled: true, Type: lite.GetName()}},
|
||||
{in: &CachePolicy{EnabledFlag: "yes", TTL: "10h"}, out: &service.CachePolicy{Enabled: true, Type: lite.GetName()}},
|
||||
{in: &CachePolicy{Type: memory.GetName(), EnabledFlag: "false", TTL: "10h"}, out: &service.CachePolicy{Enabled: false, Type: memory.GetName()}},
|
||||
{in: &CachePolicy{Type: memory.GetName(), EnabledFlag: "yes", TTL: "never"}, out: &service.CachePolicy{Enabled: true, Type: memory.GetName()}},
|
||||
{in: &CachePolicy{EnabledFlag: "no"}, out: &service.CachePolicy{Type: lite.GetName(), Enabled: false}},
|
||||
{in: &CachePolicy{Type: "memsql"}, err: trace.BadParameter("unsupported backend")},
|
||||
{in: &CachePolicy{EnabledFlag: "yes", TTL: "never"}, out: &service.CachePolicy{Enabled: true}},
|
||||
{in: &CachePolicy{EnabledFlag: "true", TTL: "10h"}, out: &service.CachePolicy{Enabled: true}},
|
||||
{in: &CachePolicy{Type: "whatever", EnabledFlag: "false", TTL: "10h"}, out: &service.CachePolicy{Enabled: false}},
|
||||
{in: &CachePolicy{Type: "name", EnabledFlag: "yes", TTL: "never"}, out: &service.CachePolicy{Enabled: true}},
|
||||
{in: &CachePolicy{EnabledFlag: "no"}, out: &service.CachePolicy{Enabled: false}},
|
||||
}
|
||||
for i, tc := range tcs {
|
||||
comment := fmt.Sprintf("test case #%v", i)
|
||||
|
||||
@@ -441,7 +441,6 @@ func (c *CachePolicy) Enabled() bool {
|
||||
// Parse parses cache policy from Teleport config
|
||||
func (c *CachePolicy) Parse() (*service.CachePolicy, error) {
|
||||
out := service.CachePolicy{
|
||||
Type: c.Type,
|
||||
Enabled: c.Enabled(),
|
||||
}
|
||||
if err := out.CheckAndSetDefaults(); err != nil {
|
||||
|
||||
@@ -389,12 +389,6 @@ var (
|
||||
// TopRequestsCapacity sets up default top requests capacity
|
||||
TopRequestsCapacity = 128
|
||||
|
||||
// CachePollPeriod is a period for cache internal events polling,
|
||||
// used in cases when cache is being used to subscribe for events
|
||||
// and this parameter controls how often cache checks for new events
|
||||
// to arrive
|
||||
CachePollPeriod = 500 * time.Millisecond
|
||||
|
||||
// AuthQueueSize is auth service queue size
|
||||
AuthQueueSize = 8192
|
||||
|
||||
|
||||
+2
-13
@@ -42,7 +42,6 @@ import (
|
||||
"github.com/gravitational/teleport/lib/auth/keystore"
|
||||
"github.com/gravitational/teleport/lib/backend"
|
||||
"github.com/gravitational/teleport/lib/backend/lite"
|
||||
"github.com/gravitational/teleport/lib/backend/memory"
|
||||
"github.com/gravitational/teleport/lib/bpf"
|
||||
"github.com/gravitational/teleport/lib/defaults"
|
||||
"github.com/gravitational/teleport/lib/events"
|
||||
@@ -308,31 +307,21 @@ func (cfg *Config) DebugDumpToYAML() string {
|
||||
|
||||
// CachePolicy sets caching policy for proxies and nodes
|
||||
type CachePolicy struct {
|
||||
// Type sets the cache type
|
||||
Type string
|
||||
// Enabled enables or disables caching
|
||||
Enabled bool
|
||||
}
|
||||
|
||||
// CheckAndSetDefaults checks and sets default values
|
||||
func (c *CachePolicy) CheckAndSetDefaults() error {
|
||||
switch c.Type {
|
||||
case "", lite.GetName():
|
||||
c.Type = lite.GetName()
|
||||
case memory.GetName():
|
||||
default:
|
||||
return trace.BadParameter("unsupported cache type %q, supported values are %q and %q",
|
||||
c.Type, lite.GetName(), memory.GetName())
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// String returns human-friendly representation of the policy
|
||||
func (c CachePolicy) String() string {
|
||||
if !c.Enabled {
|
||||
return "no cache policy"
|
||||
return "no cache"
|
||||
}
|
||||
return fmt.Sprintf("%v cache will store frequently accessed items", c.Type)
|
||||
return "in-memory cache"
|
||||
}
|
||||
|
||||
// ProxyConfig specifies configuration for proxy service
|
||||
|
||||
+16
-40
@@ -636,6 +636,13 @@ func NewTeleport(cfg *Config) (*TeleportProcess, error) {
|
||||
}
|
||||
}
|
||||
|
||||
// TODO(espadolini): DELETE IN 11.0, replace with
|
||||
// os.RemoveAll(filepath.Join(cfg.DataDir, "cache")), because no stable v10
|
||||
// should ever use the cache directory, and 11 requires upgrading from 10
|
||||
if fi, err := os.Stat(filepath.Join(cfg.DataDir, "cache")); err == nil && fi.IsDir() {
|
||||
cfg.Log.Warnf("An old cache directory exists at %q. It can be safely deleted after ensuring that no other Teleport instance is running.", filepath.Join(cfg.DataDir, "cache"))
|
||||
}
|
||||
|
||||
if len(cfg.FileDescriptors) == 0 {
|
||||
cfg.FileDescriptors, err = importFileDescriptors(cfg.Log)
|
||||
if err != nil {
|
||||
@@ -1314,7 +1321,6 @@ func (process *TeleportProcess) initAuthService() error {
|
||||
services: authServer.Services,
|
||||
setup: cache.ForAuth,
|
||||
cacheName: []string{teleport.ComponentAuth},
|
||||
inMemory: true,
|
||||
events: true,
|
||||
})
|
||||
if err != nil {
|
||||
@@ -1534,13 +1540,8 @@ type accessCacheConfig struct {
|
||||
setup cache.SetupConfigFn
|
||||
// cacheName is a cache name
|
||||
cacheName []string
|
||||
// inMemory is true if cache
|
||||
// should use memory
|
||||
inMemory bool
|
||||
// events is true if cache should turn on events
|
||||
events bool
|
||||
// pollPeriod contains period for polling
|
||||
pollPeriod time.Duration
|
||||
}
|
||||
|
||||
func (c *accessCacheConfig) CheckAndSetDefaults() error {
|
||||
@@ -1553,9 +1554,6 @@ func (c *accessCacheConfig) CheckAndSetDefaults() error {
|
||||
if len(c.cacheName) == 0 {
|
||||
return trace.BadParameter("missing parameter cacheName")
|
||||
}
|
||||
if c.pollPeriod == 0 {
|
||||
c.pollPeriod = defaults.CachePollPeriod
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -1564,39 +1562,18 @@ func (process *TeleportProcess) newAccessCache(cfg accessCacheConfig) (*cache.Ca
|
||||
if err := cfg.CheckAndSetDefaults(); err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
var cacheBackend backend.Backend
|
||||
if cfg.inMemory {
|
||||
process.log.Debugf("Creating in-memory backend for %v.", cfg.cacheName)
|
||||
mem, err := memory.New(memory.Config{
|
||||
Context: process.ExitContext(),
|
||||
EventsOff: !cfg.events,
|
||||
Mirror: true,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
cacheBackend = mem
|
||||
} else {
|
||||
process.log.Debugf("Creating sqlite backend for %v.", cfg.cacheName)
|
||||
path := filepath.Join(append([]string{process.Config.DataDir, "cache"}, cfg.cacheName...)...)
|
||||
if err := os.MkdirAll(path, teleport.SharedDirMode); err != nil {
|
||||
return nil, trace.ConvertSystemError(err)
|
||||
}
|
||||
liteBackend, err := lite.NewWithConfig(process.ExitContext(),
|
||||
lite.Config{
|
||||
Path: path,
|
||||
EventsOff: !cfg.events,
|
||||
Mirror: true,
|
||||
PollStreamPeriod: 100 * time.Millisecond,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
cacheBackend = liteBackend
|
||||
process.log.Debugf("Creating in-memory backend for %v.", cfg.cacheName)
|
||||
mem, err := memory.New(memory.Config{
|
||||
Context: process.ExitContext(),
|
||||
EventsOff: !cfg.events,
|
||||
Mirror: true,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
reporter, err := backend.NewReporter(backend.ReporterConfig{
|
||||
Component: teleport.ComponentCache,
|
||||
Backend: cacheBackend,
|
||||
Backend: mem,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
@@ -1751,7 +1728,6 @@ func (process *TeleportProcess) newLocalCacheForWindowsDesktop(clt auth.ClientI,
|
||||
// newLocalCache returns new instance of access point
|
||||
func (process *TeleportProcess) newLocalCache(clt auth.ClientI, setupConfig cache.SetupConfigFn, cacheName []string) (*cache.Cache, error) {
|
||||
return process.newAccessCache(accessCacheConfig{
|
||||
inMemory: process.Config.CachePolicy.Type == memory.GetName(),
|
||||
services: clt,
|
||||
setup: setupConfig,
|
||||
cacheName: cacheName,
|
||||
|
||||
Reference in New Issue
Block a user