From 6c8ea32bc25cdec1da8ea457bd80235b775d4bdc Mon Sep 17 00:00:00 2001 From: Julia Ogris Date: Thu, 23 Apr 2026 12:44:47 +1000 Subject: [PATCH] app: Add active sessions Prometheus gauge (#65008) * fncache: Pass non-cancellable context to `OnExpiry` callbacks Call `OnExpiry` when `get()` replaces an expired entry on reload. Previously, entries that expired between cleanup intervals were silently dropped when a new request triggered a reload, skipping the `OnExpiry` callback entirely. Use `context.WithoutCancel(c.cfg.Context)` at both the `removeExpiredLocked` and `get` call sites so that `OnExpiry` work completes even after shutdown begins. `Shutdown` cancels `c.cfg.Context` via `c.cancel()` before processing entries, but these two call sites passed `c.cfg.Context` directly, meaning goroutines spawned by `OnExpiry` could run with an already-cancelled context. Update the `OnExpiry` doc comment to warn that the cache mutex may be held when the callback is invoked. * app: Add active sessions Prometheus gauge Add a `teleport_app_active_sessions` gauge labeled by app name that tracks HTTP app sessions on each agent. The gauge increments when a session chunk is created and decrements after the session chunk finishes closing, so it reflects sessions still holding resources (audit streams, disk I/O) rather than just sessions accepting new requests. TCP and MCP sessions are excluded because they bypass the session chunk cache. * fncache: Fix `OnExpires` typo in `Shutdown` doc comment Correct the stale field name `OnExpires` to `OnExpiry` in the `Shutdown` method's doc comment to match the actual field name in `FnCacheConfig`. --- docs/pages/includes/metrics.mdx | 6 ++ lib/srv/app/metrics.go | 41 ++++++++++ lib/srv/app/metrics_test.go | 137 ++++++++++++++++++++++++++++++++ lib/srv/app/session.go | 9 ++- lib/utils/fncache.go | 18 ++++- lib/utils/fncache_test.go | 58 ++++++++++++++ 6 files changed, 264 insertions(+), 5 deletions(-) create mode 100644 lib/srv/app/metrics.go create mode 100644 lib/srv/app/metrics_test.go diff --git a/docs/pages/includes/metrics.mdx b/docs/pages/includes/metrics.mdx index 7ad960ffddd..269d0e73ebc 100644 --- a/docs/pages/includes/metrics.mdx +++ b/docs/pages/includes/metrics.mdx @@ -200,6 +200,12 @@ The following table identifies all metrics available for incoming connections. | `teleport_kubernetes_server_join_sessions_total` | counter | Teleport Kubernetes Proxy | Total number of joining sessions. | +## Application Service + +| Name | Type | Component | Description | +|------|------|-----------|-------------| +| `teleport_app_active_sessions` | gauge | Teleport Application Service | Number of active HTTP app sessions on this agent, labeled by `app` name. Does not include TCP or MCP sessions. | + ## Teleport SSH Service | Name | Type | Component | Description | diff --git a/lib/srv/app/metrics.go b/lib/srv/app/metrics.go new file mode 100644 index 00000000000..70f9a11d356 --- /dev/null +++ b/lib/srv/app/metrics.go @@ -0,0 +1,41 @@ +/* + * Teleport + * Copyright (C) 2026 Gravitational, Inc. + * + * This program is free software: you can redistribute it and/or modify + * it under the terms of the GNU Affero General Public License as published by + * the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU Affero General Public License for more details. + * + * You should have received a copy of the GNU Affero General Public License + * along with this program. If not, see . + */ + +package app + +import ( + "github.com/prometheus/client_golang/prometheus" + + "github.com/gravitational/teleport" + "github.com/gravitational/teleport/lib/observability/metrics" +) + +func init() { + _ = metrics.RegisterPrometheusCollectors(activeSessions) +} + +// activeSessions tracks the number of active HTTP app sessions on this +// agent. Each session maps to one cached session chunk in the +// ConnectionsHandler. TCP and MCP sessions are not included because they +// bypass the session chunk cache. +var activeSessions = prometheus.NewGaugeVec(prometheus.GaugeOpts{ + Namespace: teleport.MetricNamespace, + Subsystem: "app", + Name: "active_sessions", + Help: "Number of active HTTP app sessions on this agent. Does not include TCP or MCP sessions.", +}, []string{"app"}) diff --git a/lib/srv/app/metrics_test.go b/lib/srv/app/metrics_test.go new file mode 100644 index 00000000000..8ae19bbd42b --- /dev/null +++ b/lib/srv/app/metrics_test.go @@ -0,0 +1,137 @@ +/* + * Teleport + * Copyright (C) 2026 Gravitational, Inc. + * + * This program is free software: you can redistribute it and/or modify + * it under the terms of the GNU Affero General Public License as published by + * the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU Affero General Public License for more details. + * + * You should have received a copy of the GNU Affero General Public License + * along with this program. If not, see . + */ + +package app + +import ( + "context" + "log/slog" + "testing" + "time" + + "github.com/jonboulle/clockwork" + dto "github.com/prometheus/client_model/go" + "github.com/stretchr/testify/require" + + "github.com/gravitational/teleport/api/client/proto" + "github.com/gravitational/teleport/api/types" + apievents "github.com/gravitational/teleport/api/types/events" + "github.com/gravitational/teleport/lib/auth/authclient" + "github.com/gravitational/teleport/lib/events" + "github.com/gravitational/teleport/lib/session" + "github.com/gravitational/teleport/lib/tlsca" +) + +// gaugeValue reads the current value of the activeSessions gauge for the +// given app name. +func gaugeValue(t *testing.T, appName string) float64 { + t.Helper() + var m dto.Metric + require.NoError(t, activeSessions.WithLabelValues(appName).Write(&m)) + return m.GetGauge().GetValue() +} + +func TestActiveSessionsGauge(t *testing.T) { + // Reset the gauge so previous test runs do not interfere. + activeSessions.Reset() + + c := newConnectionsHandler(t) + t.Cleanup(c.cacheCloseWg.Wait) + + identity := &tlsca.Identity{Username: "alice"} + appA, err := types.NewAppV3(types.Metadata{Name: "app-a"}, types.AppSpecV3{URI: "http://localhost"}) + require.NoError(t, err) + appB, err := types.NewAppV3(types.Metadata{Name: "app-b"}, types.AppSpecV3{URI: "http://localhost"}) + require.NoError(t, err) + + // Create two sessions for different apps. + now := time.Now() + sessA, err := c.newSessionChunk(t.Context(), identity, appA, now) + require.NoError(t, err) + sessB, err := c.newSessionChunk(t.Context(), identity, appB, now) + require.NoError(t, err) + require.Equal(t, 1.0, gaugeValue(t, "app-a")) //nolint:testifylint // Inc/Dec produce exact float64 integers + require.Equal(t, 1.0, gaugeValue(t, "app-b")) //nolint:testifylint // Inc/Dec produce exact float64 integers + + // Expire the first session via the cache eviction callback. The gauge is + // decremented asynchronously after the session finishes closing, so wait + // for the background goroutine to complete before checking. + c.onSessionExpired(t.Context(), "key1", sessA) + c.cacheCloseWg.Wait() + require.Equal(t, 0.0, gaugeValue(t, "app-a")) //nolint:testifylint // Inc/Dec produce exact float64 integers + require.Equal(t, 1.0, gaugeValue(t, "app-b")) //nolint:testifylint // Inc/Dec produce exact float64 integers + + // Expire the second session. + c.onSessionExpired(t.Context(), "key2", sessB) + c.cacheCloseWg.Wait() + require.Equal(t, 0.0, gaugeValue(t, "app-b")) //nolint:testifylint // Inc/Dec produce exact float64 integers +} + +// newConnectionsHandler returns a ConnectionsHandler wired with minimal fakes +// so that newSessionChunk can be called without a full auth server. +func newConnectionsHandler(t *testing.T) *ConnectionsHandler { + t.Helper() + return &ConnectionsHandler{ + cfg: &ConnectionsHandlerConfig{ + Clock: clockwork.NewRealClock(), + DataDir: t.TempDir(), + Emitter: events.NewDiscardEmitter(), + AuthClient: metricsAuthClient{}, + AccessPoint: metricsAccessPoint{}, + }, + closeContext: t.Context(), + log: slog.Default(), + } +} + +// metricsAuthClient is a minimal fake that satisfies authclient.ClientI for +// the calls made by newSessionChunk: session tracker creation, tracker state +// updates, and audit stream creation. +type metricsAuthClient struct { + authclient.ClientI +} + +func (metricsAuthClient) CreateSessionTracker(_ context.Context, st types.SessionTracker) (types.SessionTracker, error) { + return st, nil +} + +func (metricsAuthClient) UpdateSessionTracker(_ context.Context, _ *proto.UpdateSessionTrackerRequest) error { + return nil +} + +func (metricsAuthClient) CreateAuditStream(_ context.Context, _ session.ID) (apievents.Stream, error) { + return events.NewDiscardRecorder(), nil +} + +func (metricsAuthClient) ResumeAuditStream(_ context.Context, _ session.ID, _ string) (apievents.Stream, error) { + return events.NewDiscardRecorder(), nil +} + +// metricsAccessPoint is a minimal fake that satisfies +// [authclient.AppsAccessPoint] for the two calls made by newSessionRecorder. +type metricsAccessPoint struct { + authclient.AppsAccessPoint +} + +func (metricsAccessPoint) GetSessionRecordingConfig(_ context.Context) (types.SessionRecordingConfig, error) { + return types.DefaultSessionRecordingConfig(), nil +} + +func (metricsAccessPoint) GetClusterName(_ context.Context) (types.ClusterName, error) { + return types.NewClusterName(types.ClusterNameSpecV2{ClusterName: "test", ClusterID: "test"}) +} diff --git a/lib/srv/app/session.go b/lib/srv/app/session.go index 395ddfeba81..40935f4d94d 100644 --- a/lib/srv/app/session.go +++ b/lib/srv/app/session.go @@ -61,6 +61,9 @@ type sessionChunk struct { closeC chan struct{} // id is the session chunk's uuid, which is used as the id of its session upload. id string + // appName is the name of the app this session chunk belongs to, used as + // a Prometheus label on the active sessions gauge. + appName string // streamCloser closes the session chunk stream. streamCloser utils.WriteContextCloser // audit is the session chunk audit logger. @@ -94,6 +97,7 @@ type sessionOpt func(context.Context, *sessionChunk, *tlsca.Identity, types.Appl func (c *ConnectionsHandler) newSessionChunk(ctx context.Context, identity *tlsca.Identity, app types.Application, startTime time.Time, opts ...sessionOpt) (*sessionChunk, error) { sess := &sessionChunk{ id: uuid.New().String(), + appName: app.GetName(), closeC: make(chan struct{}), inflightCond: sync.NewCond(&sync.Mutex{}), closeTimeout: sessionChunkCloseTimeout, @@ -137,6 +141,7 @@ func (c *ConnectionsHandler) newSessionChunk(ctx context.Context, identity *tlsc return nil, trace.Wrap(err) } + activeSessions.WithLabelValues(sess.appName).Inc() sess.log.DebugContext(ctx, "Created app session chunk", "session_id", sess.id) return sess, nil } @@ -262,10 +267,12 @@ func (c *ConnectionsHandler) onSessionExpired(ctx context.Context, key, expired // Closing the session stream writer may trigger a flush operation which could // be time-consuming. Launch in another goroutine to prevent interfering with - // cache operations. + // cache operations. The gauge is decremented after close completes so that + // it reflects sessions still holding resources (audit streams, IOPS). c.cacheCloseWg.Add(1) go func() { defer c.cacheCloseWg.Done() + defer activeSessions.WithLabelValues(sess.appName).Dec() if err := sess.close(ctx); err != nil { c.log.DebugContext(ctx, "Error closing session", "session_id", sess.id, "error", err) } diff --git a/lib/utils/fncache.go b/lib/utils/fncache.go index 5600310f762..9b58f94f96f 100644 --- a/lib/utils/fncache.go +++ b/lib/utils/fncache.go @@ -76,8 +76,11 @@ type FnCacheConfig struct { // caches where keys are unlikely to become orphaned. Shorter cleanup // intervals should be used when keys regularly become orphaned. CleanupInterval time.Duration - // OnExpiry is an optional callback that will be executed any time - // an item is expired and removed from the cache. + // OnExpiry is an optional callback that will be executed any time an + // item is expired and removed from the cache, or replaced by a reload + // in get() when the entry's TTL has elapsed. The callback must not call + // any method on the same FnCache instance because the cache mutex may + // be held when the callback is invoked. OnExpiry func(ctx context.Context, key any, value any) } @@ -128,7 +131,7 @@ type fnCacheEntry struct { loaded chan struct{} } -// Shutdown expires all items in the cache. If the OnExpires +// Shutdown expires all items in the cache. If the OnExpiry // callback was set in the FnCacheConfig it will be called once // per item in the cache. func (c *FnCache) Shutdown(ctx context.Context) { @@ -248,7 +251,7 @@ func (c *FnCache) removeExpiredLocked(now time.Time) { case <-entry.loaded: if now.After(entry.t.Add(entry.ttl)) { if c.cfg.OnExpiry != nil && entry.e == nil { - c.cfg.OnExpiry(c.cfg.Context, key, entry.v) + c.cfg.OnExpiry(context.WithoutCancel(c.cfg.Context), key, entry.v) } delete(c.entries, key) @@ -333,6 +336,13 @@ func (c *FnCache) get(ctx context.Context, key any, ttl time.Duration, loadfn fu } if needsReload { + // If we are replacing a loaded, successful entry, call OnExpiry so + // the old value is properly cleaned up. Without this, entries that + // expire between cleanup intervals are silently dropped when a new + // request triggers a reload, skipping the OnExpiry callback. + if entry != nil && entry.e == nil && c.cfg.OnExpiry != nil { + c.cfg.OnExpiry(context.WithoutCancel(c.cfg.Context), key, entry.v) + } // Insert a new entry with a new loaded channel. This channel will // block subsequent reads, and serve as a memory barrier for the results. entry = &fnCacheEntry{ diff --git a/lib/utils/fncache_test.go b/lib/utils/fncache_test.go index 4de7a7a944c..39a0e40ec0a 100644 --- a/lib/utils/fncache_test.go +++ b/lib/utils/fncache_test.go @@ -532,6 +532,64 @@ func TestFnCacheEviction(t *testing.T) { require.ErrorIs(t, err, ErrFnCacheClosed) } +func TestFnCacheOnExpiryReloadReplace(t *testing.T) { + t.Parallel() + + ctx := t.Context() + clock := clockwork.NewFakeClock() + + type item struct { + k any + v any + } + expiredC := make(chan item, 5) + cache, err := NewFnCache(FnCacheConfig{ + TTL: time.Hour, + Context: ctx, + Clock: clock, + CleanupInterval: 24 * time.Hour, + OnExpiry: func(ctx context.Context, key, expired any) { + expiredC <- item{k: key, v: expired} + }, + }) + require.NoError(t, err) + + // Populate the cache. + val, err := FnCacheGet(ctx, cache, "key1", func(ctx context.Context) (string, error) { + return "old-value", nil + }) + require.NoError(t, err) + require.Equal(t, "old-value", val) + + // Advance past TTL but do NOT call RemoveExpired. The long CleanupInterval + // ensures the lazy cleanup in get() does not run either. The next + // FnCacheGet for the same key triggers a reload, which should call + // OnExpiry for the old entry. + clock.Advance(2 * time.Hour) + + val, err = FnCacheGet(ctx, cache, "key1", func(ctx context.Context) (string, error) { + return "new-value", nil + }) + require.NoError(t, err) + require.Equal(t, "new-value", val) + + // Verify OnExpiry was called with the old value. + select { + case expired := <-expiredC: + require.Equal(t, "key1", expired.k) + require.Equal(t, "old-value", expired.v) + case <-time.After(5 * time.Second): + t.Fatal("timed out waiting for OnExpiry callback") + } + + // Verify no extra OnExpiry calls. + select { + case extra := <-expiredC: + t.Fatalf("unexpected extra OnExpiry call for key %v", extra.k) + default: + } +} + func TestFnCacheRemove(t *testing.T) { t.Parallel()