From b474dbf55f5e47de24fec697130ecd3d97acee26 Mon Sep 17 00:00:00 2001 From: Tonis Tiigi Date: Mon, 10 Aug 2020 17:11:24 -0700 Subject: [PATCH] resolver: clean up unused resolver pool Signed-off-by: Tonis Tiigi --- cache/remotecache/local/local.go | 2 +- exporter/local/export.go | 2 +- exporter/oci/export.go | 2 +- exporter/tar/export.go | 2 +- session/content/content_test.go | 2 +- session/filesync/filesync_test.go | 2 +- session/group.go | 2 +- session/manager.go | 8 +++++-- util/resolver/authorizer.go | 9 +++++++- util/resolver/pool.go | 37 +++++++++++++++++++++++++++++-- 10 files changed, 56 insertions(+), 12 deletions(-) diff --git a/cache/remotecache/local/local.go b/cache/remotecache/local/local.go index dabb23564..1e99ebbca 100644 --- a/cache/remotecache/local/local.go +++ b/cache/remotecache/local/local.go @@ -76,7 +76,7 @@ func getContentStore(ctx context.Context, sm *session.Manager, g session.Group, timeoutCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() - caller, err := sm.Get(timeoutCtx, sessionID) + caller, err := sm.Get(timeoutCtx, sessionID, false) if err != nil { return nil, err } diff --git a/exporter/local/export.go b/exporter/local/export.go index 28e32204f..c50100e99 100644 --- a/exporter/local/export.go +++ b/exporter/local/export.go @@ -51,7 +51,7 @@ func (e *localExporterInstance) Export(ctx context.Context, inp exporter.Source, timeoutCtx, cancel := context.WithTimeout(ctx, 5*time.Second) defer cancel() - caller, err := e.opt.SessionManager.Get(timeoutCtx, sessionID) + caller, err := e.opt.SessionManager.Get(timeoutCtx, sessionID, false) if err != nil { return nil, err } diff --git a/exporter/oci/export.go b/exporter/oci/export.go index 3874a3cfe..bd5387e51 100644 --- a/exporter/oci/export.go +++ b/exporter/oci/export.go @@ -165,7 +165,7 @@ func (e *imageExporterInstance) Export(ctx context.Context, src exporter.Source, timeoutCtx, cancel := context.WithTimeout(ctx, 5*time.Second) defer cancel() - caller, err := e.opt.SessionManager.Get(timeoutCtx, sessionID) + caller, err := e.opt.SessionManager.Get(timeoutCtx, sessionID, false) if err != nil { return nil, err } diff --git a/exporter/tar/export.go b/exporter/tar/export.go index 0f635fc1b..77e3f6888 100644 --- a/exporter/tar/export.go +++ b/exporter/tar/export.go @@ -135,7 +135,7 @@ func (e *localExporterInstance) Export(ctx context.Context, inp exporter.Source, timeoutCtx, cancel := context.WithTimeout(ctx, 5*time.Second) defer cancel() - caller, err := e.opt.SessionManager.Get(timeoutCtx, sessionID) + caller, err := e.opt.SessionManager.Get(timeoutCtx, sessionID, false) if err != nil { return nil, err } diff --git a/session/content/content_test.go b/session/content/content_test.go index 89f4ccd84..d52c60790 100644 --- a/session/content/content_test.go +++ b/session/content/content_test.go @@ -63,7 +63,7 @@ func TestContentAttachable(t *testing.T) { }) g.Go(func() error { - c, err := m.Get(ctx, s.ID()) + c, err := m.Get(ctx, s.ID(), false) if err != nil { return err } diff --git a/session/filesync/filesync_test.go b/session/filesync/filesync_test.go index 39a0125eb..b569d1731 100644 --- a/session/filesync/filesync_test.go +++ b/session/filesync/filesync_test.go @@ -46,7 +46,7 @@ func TestFileSyncIncludePatterns(t *testing.T) { }) g.Go(func() (reterr error) { - c, err := m.Get(ctx, s.ID()) + c, err := m.Get(ctx, s.ID(), false) if err != nil { return err } diff --git a/session/group.go b/session/group.go index 88409bf8b..4b9ba221f 100644 --- a/session/group.go +++ b/session/group.go @@ -74,7 +74,7 @@ func (sm *Manager) Any(ctx context.Context, g Group, f func(context.Context, str timeoutCtx, cancel := context.WithTimeout(ctx, 5*time.Second) defer cancel() - c, err := sm.Get(timeoutCtx, id) + c, err := sm.Get(timeoutCtx, id, false) if err != nil { lastErr = err continue diff --git a/session/manager.go b/session/manager.go index e01b047e2..edac93063 100644 --- a/session/manager.go +++ b/session/manager.go @@ -149,7 +149,7 @@ func (sm *Manager) handleConn(ctx context.Context, conn net.Conn, opts map[strin } // Get returns a session by ID -func (sm *Manager) Get(ctx context.Context, id string) (Caller, error) { +func (sm *Manager) Get(ctx context.Context, id string, noWait bool) (Caller, error) { // session prefix is used to identify vertexes with different contexts so // they would not collide, but for lookup we don't need the prefix if p := strings.SplitN(id, ":", 2); len(p) == 2 && len(p[1]) > 0 { @@ -180,7 +180,7 @@ func (sm *Manager) Get(ctx context.Context, id string) (Caller, error) { } var ok bool c, ok = sm.sessions[id] - if !ok || c.closed() { + if (!ok || c.closed()) && !noWait { sm.updateCondition.Wait() continue } @@ -188,6 +188,10 @@ func (sm *Manager) Get(ctx context.Context, id string) (Caller, error) { break } + if c == nil { + return nil, nil + } + return c, nil } diff --git a/util/resolver/authorizer.go b/util/resolver/authorizer.go index ec0e9c9d0..b32c10fb7 100644 --- a/util/resolver/authorizer.go +++ b/util/resolver/authorizer.go @@ -27,13 +27,15 @@ type authHandlerNS struct { mu sync.Mutex handlers map[string]*authHandler hosts map[string][]docker.RegistryHost + sm *session.Manager g flightcontrol.Group } -func newAuthHandlerNS() *authHandlerNS { +func newAuthHandlerNS(sm *session.Manager) *authHandlerNS { return &authHandlerNS{ handlers: map[string]*authHandler{}, hosts: map[string][]docker.RegistryHost{}, + sm: sm, } } @@ -54,6 +56,7 @@ func (a *authHandlerNS) get(host string, sm *session.Manager, g session.Group) * } h, ok := a.handlers[path.Join(host, id)] if ok { + h.lastUsed = time.Now() return h } } @@ -69,6 +72,7 @@ func (a *authHandlerNS) get(host string, sm *session.Manager, g session.Group) * if err == nil { if username == h.common.Username && password == h.common.Secret { a.handlers[path.Join(host, session)] = h + h.lastUsed = time.Now() return h } } @@ -214,6 +218,8 @@ type authHandler struct { // scopedTokens caches token indexed by scopes, which used in // bearer auth case scopedTokens map[string]*authResult + + lastUsed time.Time } func newAuthHandler(client *http.Client, scheme auth.AuthenticationScheme, opts auth.TokenOptions) *authHandler { @@ -222,6 +228,7 @@ func newAuthHandler(client *http.Client, scheme auth.AuthenticationScheme, opts scheme: scheme, common: opts, scopedTokens: map[string]*authResult{}, + lastUsed: time.Now(), } } diff --git a/util/resolver/pool.go b/util/resolver/pool.go index d7dc8da9f..55e932606 100644 --- a/util/resolver/pool.go +++ b/util/resolver/pool.go @@ -3,8 +3,10 @@ package resolver import ( "context" "fmt" + "strings" "sync" "sync/atomic" + "time" "github.com/containerd/containerd/images" "github.com/containerd/containerd/remotes" @@ -23,9 +25,40 @@ type Pool struct { } func NewPool() *Pool { - return &Pool{ + p := &Pool{ m: map[string]*authHandlerNS{}, } + time.AfterFunc(5*time.Minute, p.gc) + return p +} + +func (p *Pool) gc() { + p.mu.Lock() + defer p.mu.Unlock() + + for k, ns := range p.m { + ns.mu.Lock() + for key, h := range ns.handlers { + if time.Since(h.lastUsed) < 10*time.Minute { + continue + } + parts := strings.SplitN(key, "/", 2) + if len(parts) != 2 { + delete(ns.handlers, key) + continue + } + c, err := ns.sm.Get(context.TODO(), parts[1], true) + if c == nil || err != nil { + delete(ns.handlers, key) + } + } + if len(ns.handlers) == 0 { + delete(p.m, k) + } + ns.mu.Unlock() + } + + time.AfterFunc(5*time.Minute, p.gc) } func (p *Pool) Clear() { @@ -47,7 +80,7 @@ func (p *Pool) GetResolver(hosts docker.RegistryHosts, ref, scope string, sm *se defer p.mu.Unlock() h, ok := p.m[key] if !ok { - h = newAuthHandlerNS() + h = newAuthHandlerNS(sm) p.m[key] = h } return newResolver(hosts, h, sm, g)