diff --git a/integration/helpers.go b/integration/helpers.go index e7d237cc4fd..1bc8fcaa415 100644 --- a/integration/helpers.go +++ b/integration/helpers.go @@ -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 diff --git a/lib/cache/cache.go b/lib/cache/cache.go index 42fc0e99968..d836dbb0d21 100644 --- a/lib/cache/cache.go +++ b/lib/cache/cache.go @@ -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 diff --git a/lib/cache/cache_test.go b/lib/cache/cache_test.go index d894714d884..185504d5c96 100644 --- a/lib/cache/cache_test.go +++ b/lib/cache/cache_test.go @@ -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) { diff --git a/lib/config/configuration.go b/lib/config/configuration.go index 5a90f966cf8..1fd3befe15c 100644 --- a/lib/config/configuration.go +++ b/lib/config/configuration.go @@ -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 diff --git a/lib/config/configuration_test.go b/lib/config/configuration_test.go index 71d8ed81f4c..95c0334ebd9 100644 --- a/lib/config/configuration_test.go +++ b/lib/config/configuration_test.go @@ -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) diff --git a/lib/config/fileconf.go b/lib/config/fileconf.go index 5c2b719fa3c..48ce3fa9b5b 100644 --- a/lib/config/fileconf.go +++ b/lib/config/fileconf.go @@ -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 { diff --git a/lib/defaults/defaults.go b/lib/defaults/defaults.go index 88539fa023c..d9ef985b131 100644 --- a/lib/defaults/defaults.go +++ b/lib/defaults/defaults.go @@ -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 diff --git a/lib/service/cfg.go b/lib/service/cfg.go index 3cd7955fb55..b39d8ce3c6c 100644 --- a/lib/service/cfg.go +++ b/lib/service/cfg.go @@ -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 diff --git a/lib/service/service.go b/lib/service/service.go index bf882857ec8..d4e7d1cd2f8 100644 --- a/lib/service/service.go +++ b/lib/service/service.go @@ -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,