From 0f02047d2ebfc105d25d2babcdcc5189e1b1ad41 Mon Sep 17 00:00:00 2001 From: rosstimothy <39066650+rosstimothy@users.noreply.github.com> Date: Tue, 15 Oct 2024 15:42:25 +0000 Subject: [PATCH] Convert resource watchers to use slog (#47560) The logrus logger was left in place for now and will be removed once teleport.e has been converted. --- lib/kube/proxy/watcher.go | 5 +- lib/proxy/peer/client.go | 3 +- lib/reversetunnel/localsite_test.go | 4 +- lib/reversetunnel/srv.go | 17 +++--- lib/service/db.go | 2 +- lib/service/desktop.go | 2 +- lib/service/kubernetes.go | 2 +- lib/service/service.go | 18 +++--- lib/services/watcher.go | 92 +++++++++++++++-------------- lib/srv/app/watcher.go | 5 +- lib/srv/db/watcher.go | 2 +- lib/srv/discovery/discovery.go | 5 +- lib/web/apiserver_test.go | 1 - 13 files changed, 84 insertions(+), 74 deletions(-) diff --git a/lib/kube/proxy/watcher.go b/lib/kube/proxy/watcher.go index 8b94e62dbde..188d16822a3 100644 --- a/lib/kube/proxy/watcher.go +++ b/lib/kube/proxy/watcher.go @@ -96,8 +96,9 @@ func (s *TLSServer) startKubeClusterResourceWatcher(ctx context.Context) (*servi watcher, err := services.NewKubeClusterWatcher(ctx, services.KubeClusterWatcherConfig{ ResourceWatcherConfig: services.ResourceWatcherConfig{ Component: s.Component, - Log: s.log, - Client: s.AccessPoint, + // TODO(tross): update this once converted to use slog + // Logger: s.log, + Client: s.AccessPoint, }, }) if err != nil { diff --git a/lib/proxy/peer/client.go b/lib/proxy/peer/client.go index f6a9e0aff42..227e6077261 100644 --- a/lib/proxy/peer/client.go +++ b/lib/proxy/peer/client.go @@ -414,7 +414,8 @@ func (c *Client) sync() { ResourceWatcherConfig: services.ResourceWatcherConfig{ Component: teleport.Component(teleport.ComponentProxyPeer), Client: c.config.AccessPoint, - Log: c.config.Log, + // TODO(tross): use the configured logger after updating peering to use slog + // Logger: c.config.Logger, }, ProxyDiffer: func(old, new types.Server) bool { return old.GetPeerAddr() != new.GetPeerAddr() diff --git a/lib/reversetunnel/localsite_test.go b/lib/reversetunnel/localsite_test.go index 3c1d8f3cfe3..3397aed7636 100644 --- a/lib/reversetunnel/localsite_test.go +++ b/lib/reversetunnel/localsite_test.go @@ -61,7 +61,7 @@ func TestRemoteConnCleanup(t *testing.T) { watcher, err := services.NewProxyWatcher(ctx, services.ProxyWatcherConfig{ ResourceWatcherConfig: services.ResourceWatcherConfig{ Component: "test", - Log: utils.NewLoggerForTests(), + Logger: utils.NewSlogLoggerForTests(), Clock: clock, Client: &mockLocalSiteClient{}, }, @@ -253,7 +253,7 @@ func TestProxyResync(t *testing.T) { watcher, err := services.NewProxyWatcher(ctx, services.ProxyWatcherConfig{ ResourceWatcherConfig: services.ResourceWatcherConfig{ Component: "test", - Log: utils.NewLoggerForTests(), + Logger: utils.NewSlogLoggerForTests(), Clock: clock, Client: &mockLocalSiteClient{ proxies: []types.Server{proxy1, proxy2}, diff --git a/lib/reversetunnel/srv.go b/lib/reversetunnel/srv.go index d47ff6c4c9f..19dfd9e2d43 100644 --- a/lib/reversetunnel/srv.go +++ b/lib/reversetunnel/srv.go @@ -302,7 +302,8 @@ func NewServer(cfg Config) (reversetunnelclient.Server, error) { ResourceWatcherConfig: services.ResourceWatcherConfig{ Component: cfg.Component, Client: cfg.LocalAccessPoint, - Log: cfg.Log, + // TODO(tross): update this after converting to slog here + // Logger: cfg.Log, }, ProxiesC: make(chan []types.Server, 10), ProxyGetter: cfg.LocalAccessPoint, @@ -1211,9 +1212,10 @@ func newRemoteSite(srv *server, domainName string, sconn ssh.Conn) (*remoteSite, remoteSite.remoteAccessPoint = accessPoint nodeWatcher, err := services.NewNodeWatcher(closeContext, services.NodeWatcherConfig{ ResourceWatcherConfig: services.ResourceWatcherConfig{ - Component: srv.Component, - Client: accessPoint, - Log: srv.Log, + Component: srv.Component, + Client: accessPoint, + // TODO(tross) update this after converting to use slog + // Logger: srv.Log, MaxStaleness: time.Minute, }, NodesGetter: accessPoint, @@ -1246,9 +1248,10 @@ func newRemoteSite(srv *server, domainName string, sconn ssh.Conn) (*remoteSite, remoteWatcher, err := services.NewCertAuthorityWatcher(srv.ctx, services.CertAuthorityWatcherConfig{ ResourceWatcherConfig: services.ResourceWatcherConfig{ Component: teleport.ComponentProxy, - Log: srv.log, - Clock: srv.Clock, - Client: remoteSite.remoteAccessPoint, + // TODO(tross): update this after converting to slog + // Logger: srv.log, + Clock: srv.Clock, + Client: remoteSite.remoteAccessPoint, }, Types: []types.CertAuthType{types.HostCA}, }) diff --git a/lib/service/db.go b/lib/service/db.go index 1e83b547ad9..80258e66582 100644 --- a/lib/service/db.go +++ b/lib/service/db.go @@ -90,7 +90,7 @@ func (process *TeleportProcess) initDatabaseService() (retErr error) { lockWatcher, err := services.NewLockWatcher(process.ExitContext(), services.LockWatcherConfig{ ResourceWatcherConfig: services.ResourceWatcherConfig{ Component: teleport.ComponentDatabase, - Log: process.log.WithField(teleport.ComponentKey, teleport.Component(teleport.ComponentDatabase, process.id)), + Logger: process.logger.With(teleport.ComponentKey, teleport.Component(teleport.ComponentDatabase, process.id)), Client: conn.Client, }, }) diff --git a/lib/service/desktop.go b/lib/service/desktop.go index adc3f19de92..c4950a882a8 100644 --- a/lib/service/desktop.go +++ b/lib/service/desktop.go @@ -144,7 +144,7 @@ func (process *TeleportProcess) initWindowsDesktopServiceRegistered(logger *slog lockWatcher, err := services.NewLockWatcher(process.ExitContext(), services.LockWatcherConfig{ ResourceWatcherConfig: services.ResourceWatcherConfig{ Component: teleport.ComponentWindowsDesktop, - Log: process.log.WithField(teleport.ComponentKey, teleport.Component(teleport.ComponentWindowsDesktop, process.id)), + Logger: process.logger.With(teleport.ComponentKey, teleport.Component(teleport.ComponentWindowsDesktop, process.id)), Clock: cfg.Clock, Client: conn.Client, }, diff --git a/lib/service/kubernetes.go b/lib/service/kubernetes.go index 935d327f6b8..e09569c4ba3 100644 --- a/lib/service/kubernetes.go +++ b/lib/service/kubernetes.go @@ -173,7 +173,7 @@ func (process *TeleportProcess) initKubernetesService(logger *slog.Logger, conn lockWatcher, err := services.NewLockWatcher(process.ExitContext(), services.LockWatcherConfig{ ResourceWatcherConfig: services.ResourceWatcherConfig{ Component: teleport.ComponentKube, - Log: process.log.WithField(teleport.ComponentKey, teleport.Component(teleport.ComponentKube, process.id)), + Logger: process.logger.With(teleport.ComponentKey, teleport.Component(teleport.ComponentKube, process.id)), Client: conn.Client, }, }) diff --git a/lib/service/service.go b/lib/service/service.go index 8b01eacf6ad..58c63b5671a 100644 --- a/lib/service/service.go +++ b/lib/service/service.go @@ -2117,7 +2117,7 @@ func (process *TeleportProcess) initAuthService() error { lockWatcher, err := services.NewLockWatcher(process.ExitContext(), services.LockWatcherConfig{ ResourceWatcherConfig: services.ResourceWatcherConfig{ Component: teleport.ComponentAuth, - Log: process.log.WithField(teleport.ComponentKey, teleport.Component(teleport.ComponentAuth, process.id)), + Logger: process.logger.With(teleport.ComponentKey, teleport.Component(teleport.ComponentAuth, process.id)), Client: authServer.Services, }, }) @@ -2134,7 +2134,7 @@ func (process *TeleportProcess) initAuthService() error { ResourceWatcherConfig: services.ResourceWatcherConfig{ QueueSize: defaults.UnifiedResourcesQueueSize, Component: teleport.ComponentUnifiedResource, - Log: process.log.WithField(teleport.ComponentKey, teleport.ComponentUnifiedResource), + Logger: process.logger.With(teleport.ComponentKey, teleport.ComponentUnifiedResource), Client: authServer, MaxStaleness: time.Minute, }, @@ -2923,7 +2923,7 @@ func (process *TeleportProcess) initSSH() error { lockWatcher, err := services.NewLockWatcher(process.ExitContext(), services.LockWatcherConfig{ ResourceWatcherConfig: services.ResourceWatcherConfig{ Component: teleport.ComponentNode, - Log: process.log.WithField(teleport.ComponentKey, teleport.Component(teleport.ComponentNode, process.id)), + Logger: process.logger.With(teleport.ComponentKey, teleport.Component(teleport.ComponentNode, process.id)), Client: conn.Client, }, }) @@ -4251,7 +4251,7 @@ func (process *TeleportProcess) initProxyEndpoint(conn *Connector) error { lockWatcher, err := services.NewLockWatcher(process.ExitContext(), services.LockWatcherConfig{ ResourceWatcherConfig: services.ResourceWatcherConfig{ Component: teleport.ComponentProxy, - Log: process.log.WithField(teleport.ComponentKey, teleport.ComponentProxy), + Logger: process.logger.With(teleport.ComponentKey, teleport.ComponentProxy), Client: conn.Client, }, }) @@ -4262,7 +4262,7 @@ func (process *TeleportProcess) initProxyEndpoint(conn *Connector) error { nodeWatcher, err := services.NewNodeWatcher(process.ExitContext(), services.NodeWatcherConfig{ ResourceWatcherConfig: services.ResourceWatcherConfig{ Component: teleport.ComponentProxy, - Log: process.log.WithField(teleport.ComponentKey, teleport.ComponentProxy), + Logger: process.logger.With(teleport.ComponentKey, teleport.ComponentProxy), Client: accessPoint, MaxStaleness: time.Minute, }, @@ -4275,7 +4275,7 @@ func (process *TeleportProcess) initProxyEndpoint(conn *Connector) error { caWatcher, err := services.NewCertAuthorityWatcher(process.ExitContext(), services.CertAuthorityWatcherConfig{ ResourceWatcherConfig: services.ResourceWatcherConfig{ Component: teleport.ComponentProxy, - Log: process.log.WithField(teleport.ComponentKey, teleport.ComponentProxy), + Logger: process.logger.With(teleport.ComponentKey, teleport.ComponentProxy), Client: accessPoint, }, AuthorityGetter: accessPoint, @@ -4511,7 +4511,7 @@ func (process *TeleportProcess) initProxyEndpoint(conn *Connector) error { lockWatcher, err := services.NewLockWatcher(process.GracefulExitContext(), services.LockWatcherConfig{ ResourceWatcherConfig: services.ResourceWatcherConfig{ Component: teleport.ComponentWebProxy, - Log: process.log, + Logger: process.logger, Client: conn.Client, Clock: process.Clock, }, @@ -5020,7 +5020,7 @@ func (process *TeleportProcess) initProxyEndpoint(conn *Connector) error { kubeServerWatcher, err := services.NewKubeServerWatcher(process.ExitContext(), services.KubeServerWatcherConfig{ ResourceWatcherConfig: services.ResourceWatcherConfig{ Component: component, - Log: process.log.WithField(teleport.ComponentKey, teleport.Component(teleport.ComponentReverseTunnelServer, process.id)), + Logger: process.logger.With(teleport.ComponentKey, teleport.Component(teleport.ComponentReverseTunnelServer, process.id)), Client: accessPoint, }, }) @@ -5955,7 +5955,7 @@ func (process *TeleportProcess) initApps() { lockWatcher, err := services.NewLockWatcher(process.ExitContext(), services.LockWatcherConfig{ ResourceWatcherConfig: services.ResourceWatcherConfig{ Component: teleport.ComponentApp, - Log: process.log.WithField(teleport.ComponentKey, component), + Logger: process.logger.With(teleport.ComponentKey, component), Client: conn.Client, }, }) diff --git a/lib/services/watcher.go b/lib/services/watcher.go index a36868e5a99..63a4af59ae7 100644 --- a/lib/services/watcher.go +++ b/lib/services/watcher.go @@ -20,6 +20,7 @@ package services import ( "context" + "log/slog" "strings" "sync" "sync/atomic" @@ -35,6 +36,7 @@ import ( "github.com/gravitational/teleport/api/utils/retryutils" "github.com/gravitational/teleport/lib/defaults" "github.com/gravitational/teleport/lib/utils" + logutils "github.com/gravitational/teleport/lib/utils/log" ) const ( @@ -88,8 +90,10 @@ func watchKindsString(kinds []types.WatchKind) string { type ResourceWatcherConfig struct { // Component is a component used in logs. Component string - // Log is a logger. + // TODO(tross): remove this once e has been updated. Log logrus.FieldLogger + // Logger emits log messages. + Logger *slog.Logger // MaxRetryPeriod is the maximum retry period on failed watchers. MaxRetryPeriod time.Duration // Clock is used to control time. @@ -110,8 +114,8 @@ func (cfg *ResourceWatcherConfig) CheckAndSetDefaults() error { if cfg.Component == "" { return trace.BadParameter("missing parameter Component") } - if cfg.Log == nil { - cfg.Log = logrus.StandardLogger() + if cfg.Logger == nil { + cfg.Logger = slog.Default() } if cfg.MaxRetryPeriod == 0 { cfg.MaxRetryPeriod = defaults.MaxWatcherBackoff @@ -146,7 +150,7 @@ func newResourceWatcher(ctx context.Context, collector resourceCollector, cfg Re return nil, trace.Wrap(err) } - cfg.Log = cfg.Log.WithField("resource-kind", watchKindsString(collector.resourceKinds())) + cfg.Logger = cfg.Logger.With("resource_kinds", watchKindsString(collector.resourceKinds())) ctx, cancel := context.WithCancel(ctx) p := &resourceWatcher{ ResourceWatcherConfig: cfg, @@ -220,7 +224,7 @@ func (p *resourceWatcher) WaitInitialization() error { case <-p.collector.initializationChan(): return nil case <-t.C: - p.Log.Debug("ResourceWatcher is not yet initialized.") + p.Logger.DebugContext(p.ctx, "ResourceWatcher is not yet initialized.") case <-p.ctx.Done(): return trace.BadParameter("ResourceWatcher %s failed to initialize.", watchKindsString(p.collector.resourceKinds())) } @@ -246,7 +250,7 @@ func (p *resourceWatcher) hasStaleView() bool { // runWatchLoop runs a watch loop. func (p *resourceWatcher) runWatchLoop() { for { - p.Log.Debug("Starting watch.") + p.Logger.DebugContext(p.ctx, "Starting watch.") err := p.watch() select { @@ -261,7 +265,7 @@ func (p *resourceWatcher) runWatchLoop() { p.failureStartedAt = p.Clock.Now() } if p.hasStaleView() { - p.Log.Warningf("Maximum staleness of %v exceeded, failure started at %v.", p.MaxStaleness, p.failureStartedAt) + p.Logger.WarnContext(p.ctx, "Maximum staleness of period exceeded.", "max_staleness", p.MaxStaleness, "failure_started", p.failureStartedAt) p.collector.notifyStale() } @@ -275,19 +279,19 @@ func (p *resourceWatcher) runWatchLoop() { startedWaiting := p.Clock.Now() select { case t := <-p.retry.After(): - p.Log.Debugf("Attempting to restart watch after waiting %v.", t.Sub(startedWaiting)) + p.Logger.DebugContext(p.ctx, "Attempting to restart watch after waiting", "waited", t.Sub(startedWaiting)) p.retry.Inc() case <-p.ctx.Done(): - p.Log.Debug("Closed, returning from watch loop.") + p.Logger.DebugContext(p.ctx, "Closed, returning from watch loop.") return case <-p.StaleC: // Used for testing that the watch routine is waiting for the // next restart attempt. We don't want to wait for the full // retry period in tests so we trigger the restart immediately. - p.Log.Debug("Stale view, continue watch loop.") + p.Logger.DebugContext(p.ctx, "Stale view, continue watch loop.") } if err != nil { - p.Log.Warningf("Restart watch on error: %v.", err) + p.Logger.WarnContext(p.ctx, "Restart watch on error", "error", err) } } } @@ -495,7 +499,7 @@ func (p *proxyCollector) processEventsAndUpdateCurrent(ctx context.Context, even for _, event := range events { if event.Resource == nil || event.Resource.GetKind() != types.KindProxy { - p.Log.Warningf("Unexpected event: %v.", event) + p.Logger.WarnContext(ctx, "Received unexpected event", "event", logutils.StringerAttr(event)) continue } @@ -507,7 +511,7 @@ func (p *proxyCollector) processEventsAndUpdateCurrent(ctx context.Context, even case types.OpPut: server, ok := event.Resource.(types.Server) if !ok { - p.Log.Warningf("Unexpected type %T.", event.Resource) + p.Logger.WarnContext(ctx, "Received unexpected type", "resource", event.Resource.GetKind()) continue } current, exists := p.current[server.GetName()] @@ -516,7 +520,7 @@ func (p *proxyCollector) processEventsAndUpdateCurrent(ctx context.Context, even updated = true } default: - p.Log.Warningf("Skipping unsupported event type %s.", event.Type) + p.Logger.WarnContext(ctx, "Skipping unsupported event type", "event_type", event.Type) } } @@ -531,7 +535,7 @@ func (p *proxyCollector) broadcastUpdate(ctx context.Context) { for k := range p.current { names = append(names, k) } - p.Log.Debugf("List of known proxies updated: %q.", names) + p.Logger.DebugContext(ctx, "List of known proxies updated", "proxies", names) select { case p.ProxiesC <- serverMapValues(p.current): @@ -747,7 +751,7 @@ func (p *lockCollector) processEventsAndUpdateCurrent(ctx context.Context, event eventsToEmit := events[:0] for _, event := range events { if event.Resource == nil || event.Resource.GetKind() != types.KindLock { - p.Log.Warningf("Unexpected event: %v.", event) + p.Logger.WarnContext(ctx, "Received unexpected event", "event", logutils.StringerAttr(event)) continue } @@ -758,7 +762,7 @@ func (p *lockCollector) processEventsAndUpdateCurrent(ctx context.Context, event case types.OpPut: lock, ok := event.Resource.(types.Lock) if !ok { - p.Log.Warningf("Unexpected resource type %T.", event.Resource) + p.Logger.WarnContext(ctx, "Unexpected resource type", "resource", event.Resource.GetKind()) continue } if lock.IsInForce(p.Clock.Now()) { @@ -768,7 +772,7 @@ func (p *lockCollector) processEventsAndUpdateCurrent(ctx context.Context, event delete(p.current, lock.GetName()) } default: - p.Log.Warningf("Skipping unsupported event type %s.", event.Type) + p.Logger.WarnContext(ctx, "Skipping unsupported event type", "event_type", event.Type) } } p.fanout.Emit(eventsToEmit...) @@ -929,7 +933,7 @@ func (p *databaseCollector) processEventsAndUpdateCurrent(ctx context.Context, e var updated bool for _, event := range events { if event.Resource == nil || event.Resource.GetKind() != types.KindDatabase { - p.Log.Warnf("Unexpected event: %v.", event) + p.Logger.WarnContext(ctx, "Received unexpected event", "event", logutils.StringerAttr(event)) continue } switch event.Type { @@ -939,13 +943,13 @@ func (p *databaseCollector) processEventsAndUpdateCurrent(ctx context.Context, e case types.OpPut: database, ok := event.Resource.(types.Database) if !ok { - p.Log.Warnf("Unexpected resource type %T.", event.Resource) + p.Logger.WarnContext(ctx, "Received unexpected resource type", "resource", event.Resource.GetKind()) continue } p.current[database.GetName()] = database updated = true default: - p.Log.Warnf("Unsupported event type %s.", event.Type) + p.Logger.WarnContext(ctx, "Received unsupported event type", "event_type", event.Type) } } @@ -1068,7 +1072,7 @@ func (p *appCollector) processEventsAndUpdateCurrent(ctx context.Context, events defer p.lock.Unlock() for _, event := range events { if event.Resource == nil || event.Resource.GetKind() != types.KindApp { - p.Log.Warnf("Unexpected event: %v.", event) + p.Logger.WarnContext(ctx, "Received unexpected event", "event", logutils.StringerAttr(event)) continue } switch event.Type { @@ -1084,7 +1088,7 @@ func (p *appCollector) processEventsAndUpdateCurrent(ctx context.Context, events case types.OpPut: app, ok := event.Resource.(types.Application) if !ok { - p.Log.Warnf("Unexpected resource type %T.", event.Resource) + p.Logger.WarnContext(ctx, "Received unexpected resource type", "resource", event.Resource.GetKind()) continue } p.current[app.GetName()] = app @@ -1094,7 +1098,7 @@ func (p *appCollector) processEventsAndUpdateCurrent(ctx context.Context, events case p.AppsC <- resourcesToSlice(p.current): } default: - p.Log.Warnf("Unsupported event type %s.", event.Type) + p.Logger.WarnContext(ctx, "Received unsupported event type", "event_type", event.Type) } } } @@ -1220,7 +1224,7 @@ func (k *kubeCollector) processEventsAndUpdateCurrent(ctx context.Context, event defer k.lock.Unlock() for _, event := range events { if event.Resource == nil || event.Resource.GetKind() != types.KindKubernetesCluster { - k.Log.Warnf("Unexpected event: %v.", event) + k.Logger.WarnContext(ctx, "Received unexpected event", "event", logutils.StringerAttr(event)) continue } switch event.Type { @@ -1236,7 +1240,7 @@ func (k *kubeCollector) processEventsAndUpdateCurrent(ctx context.Context, event case types.OpPut: cluster, ok := event.Resource.(types.KubeCluster) if !ok { - k.Log.Warnf("Unexpected resource type %T.", event.Resource) + k.Logger.WarnContext(ctx, "Received unexpected resource type", "resource", event.Resource.GetKind()) continue } k.current[cluster.GetName()] = cluster @@ -1246,7 +1250,7 @@ func (k *kubeCollector) processEventsAndUpdateCurrent(ctx context.Context, event case k.KubeClustersC <- resourcesToSlice(k.current): } default: - k.Log.Warnf("Unsupported event type %s.", event.Type) + k.Logger.WarnContext(ctx, "Received unsupported event type", "event_type", event.Type) } } } @@ -1425,7 +1429,7 @@ func (k *kubeServerCollector) processEventsAndUpdateCurrent(ctx context.Context, for _, event := range events { if event.Resource == nil || event.Resource.GetKind() != types.KindKubeServer { - k.Log.Warnf("Unexpected event: %v.", event) + k.Logger.WarnContext(ctx, "Received unexpected event", "event", logutils.StringerAttr(event)) continue } @@ -1440,7 +1444,7 @@ func (k *kubeServerCollector) processEventsAndUpdateCurrent(ctx context.Context, case types.OpPut: server, ok := event.Resource.(types.KubeServer) if !ok { - k.Log.Warnf("Unexpected resource type %T.", event.Resource) + k.Logger.WarnContext(ctx, "Received unexpected resource type", "resource", event.Resource.GetKind()) continue } @@ -1450,7 +1454,7 @@ func (k *kubeServerCollector) processEventsAndUpdateCurrent(ctx context.Context, } k.current[key] = server default: - k.Log.Warnf("Unsupported event type %s.", event.Type) + k.Logger.WarnContext(ctx, "Received unsupported event type", "event_type", event.Type) } } } @@ -1653,7 +1657,7 @@ func (c *caCollector) processEventsAndUpdateCurrent(ctx context.Context, events for _, event := range events { if event.Resource == nil || event.Resource.GetKind() != types.KindCertAuthority { - c.Log.Warnf("Unexpected event: %v.", event) + c.Logger.WarnContext(ctx, "Received unexpected event", "event", logutils.StringerAttr(event)) continue } switch event.Type { @@ -1668,7 +1672,7 @@ func (c *caCollector) processEventsAndUpdateCurrent(ctx context.Context, events case types.OpPut: ca, ok := event.Resource.(types.CertAuthority) if !ok { - c.Log.Warnf("Unexpected resource type %T.", event.Resource) + c.Logger.WarnContext(ctx, "Received unexpected resource type", "resource", event.Resource.GetKind()) continue } @@ -1684,7 +1688,7 @@ func (c *caCollector) processEventsAndUpdateCurrent(ctx context.Context, events c.cas[ca.GetType()][ca.GetName()] = ca eventsToEmit = append(eventsToEmit, event) default: - c.Log.Warnf("Unsupported event type %s.", event.Type) + c.Logger.WarnContext(ctx, "Received unsupported event type", "event_type", event.Type) } } @@ -1940,7 +1944,7 @@ func (n *nodeCollector) processEventsAndUpdateCurrent(ctx context.Context, event for _, event := range events { if event.Resource == nil || event.Resource.GetKind() != types.KindNode { - n.Log.Warningf("Unexpected event: %v.", event) + n.Logger.WarnContext(ctx, "Received unexpected event", "event", logutils.StringerAttr(event)) continue } @@ -1950,13 +1954,13 @@ func (n *nodeCollector) processEventsAndUpdateCurrent(ctx context.Context, event case types.OpPut: server, ok := event.Resource.(types.Server) if !ok { - n.Log.Warningf("Unexpected type %T.", event.Resource) + n.Logger.WarnContext(ctx, "Received unexpected type", "resource", event.Resource.GetKind()) continue } n.current[server.GetName()] = server default: - n.Log.Warningf("Skipping unsupported event type %s.", event.Type) + n.Logger.WarnContext(ctx, "Skipping unsupported event type", "event_type", event.Type) } } } @@ -2085,7 +2089,7 @@ func (p *accessRequestCollector) processEventsAndUpdateCurrent(ctx context.Conte for _, event := range events { if event.Resource == nil || event.Resource.GetKind() != types.KindAccessRequest { - p.Log.Warnf("Unexpected event: %v.", event) + p.Logger.WarnContext(ctx, "Received unexpected event", "event", logutils.StringerAttr(event)) continue } switch event.Type { @@ -2098,7 +2102,7 @@ func (p *accessRequestCollector) processEventsAndUpdateCurrent(ctx context.Conte case types.OpPut: accessRequest, ok := event.Resource.(types.AccessRequest) if !ok { - p.Log.Warnf("Unexpected resource type %T.", event.Resource) + p.Logger.WarnContext(ctx, "Received unexpected resource type", "resource", event.Resource.GetKind()) continue } p.current[accessRequest.GetName()] = accessRequest @@ -2108,7 +2112,7 @@ func (p *accessRequestCollector) processEventsAndUpdateCurrent(ctx context.Conte } default: - p.Log.Warnf("Unsupported event type %s.", event.Type) + p.Logger.WarnContext(ctx, "Received unsupported event type", "event_type", event.Type) } } } @@ -2152,7 +2156,7 @@ func NewOktaAssignmentWatcher(ctx context.Context, cfg OktaAssignmentWatcherConf return nil, trace.Wrap(err) } collector := &oktaAssignmentCollector{ - log: cfg.RWCfg.Log, + logger: cfg.RWCfg.Logger, cfg: cfg, initializationC: make(chan struct{}), } @@ -2189,7 +2193,7 @@ func (o *OktaAssignmentWatcher) Done() <-chan struct{} { // oktaAssignmentCollector accompanies resourceWatcher when monitoring Okta assignment resources. type oktaAssignmentCollector struct { - log logrus.FieldLogger + logger *slog.Logger // OktaAssignmentWatcherConfig is the watcher configuration. cfg OktaAssignmentWatcherConfig // mu guards "current" @@ -2260,7 +2264,7 @@ func (c *oktaAssignmentCollector) processEventsAndUpdateCurrent(ctx context.Cont for _, event := range events { if event.Resource == nil || event.Resource.GetKind() != types.KindOktaAssignment { - c.log.Warnf("Unexpected event: %v.", event) + c.logger.WarnContext(ctx, "Received unexpected event", "event", logutils.StringerAttr(event)) continue } switch event.Type { @@ -2274,7 +2278,7 @@ func (c *oktaAssignmentCollector) processEventsAndUpdateCurrent(ctx context.Cont case types.OpPut: oktaAssignment, ok := event.Resource.(types.OktaAssignment) if !ok { - c.log.Warnf("Unexpected resource type %T.", event.Resource) + c.logger.WarnContext(ctx, "Received unexpected resource type", "resource", event.Resource.GetKind()) continue } c.current[oktaAssignment.GetName()] = oktaAssignment @@ -2286,7 +2290,7 @@ func (c *oktaAssignmentCollector) processEventsAndUpdateCurrent(ctx context.Cont } default: - c.log.Warnf("Unsupported event type %s.", event.Type) + c.logger.WarnContext(ctx, "Received unsupported event type", "event_type", event.Type) } } } diff --git a/lib/srv/app/watcher.go b/lib/srv/app/watcher.go index 280a37882ab..fb0acc2bfad 100644 --- a/lib/srv/app/watcher.go +++ b/lib/srv/app/watcher.go @@ -74,8 +74,9 @@ func (s *Server) startResourceWatcher(ctx context.Context) (*services.AppWatcher watcher, err := services.NewAppWatcher(ctx, services.AppWatcherConfig{ ResourceWatcherConfig: services.ResourceWatcherConfig{ Component: teleport.ComponentApp, - Log: s.log, - Client: s.c.AccessPoint, + // TODO(tross): update this once converted to use slog + // Log: s.log, + Client: s.c.AccessPoint, }, }) if err != nil { diff --git a/lib/srv/db/watcher.go b/lib/srv/db/watcher.go index a5455c4d1fd..a3313a90792 100644 --- a/lib/srv/db/watcher.go +++ b/lib/srv/db/watcher.go @@ -78,7 +78,7 @@ func (s *Server) startResourceWatcher(ctx context.Context) (*services.DatabaseWa watcher, err := services.NewDatabaseWatcher(ctx, services.DatabaseWatcherConfig{ ResourceWatcherConfig: services.ResourceWatcherConfig{ Component: teleport.ComponentDatabase, - Log: s.logrusLogger, + Logger: s.log, Client: s.cfg.AccessPoint, }, }) diff --git a/lib/srv/discovery/discovery.go b/lib/srv/discovery/discovery.go index 6ce416212c7..c2c16eeb43e 100644 --- a/lib/srv/discovery/discovery.go +++ b/lib/srv/discovery/discovery.go @@ -1692,8 +1692,9 @@ func (s *Server) getAzureSubscriptions(ctx context.Context, subs []string) ([]st func (s *Server) initTeleportNodeWatcher() (err error) { s.nodeWatcher, err = services.NewNodeWatcher(s.ctx, services.NodeWatcherConfig{ ResourceWatcherConfig: services.ResourceWatcherConfig{ - Component: teleport.ComponentDiscovery, - Log: s.Log, + Component: teleport.ComponentDiscovery, + // TODO(tross): update this after converting logging to use slog + // Logger: s.Logger, Client: s.AccessPoint, MaxStaleness: time.Minute, }, diff --git a/lib/web/apiserver_test.go b/lib/web/apiserver_test.go index 5254cff6e4f..06b42f394e6 100644 --- a/lib/web/apiserver_test.go +++ b/lib/web/apiserver_test.go @@ -9073,7 +9073,6 @@ func startKubeWithoutCleanup(ctx context.Context, t *testing.T, cfg startKubeOpt watcher, err := services.NewKubeServerWatcher(ctx, services.KubeServerWatcherConfig{ ResourceWatcherConfig: services.ResourceWatcherConfig{ Component: component, - Log: log, Client: client, Clock: clock, },