From 25520c450527dae6aa73b9ad1df99e527b278708 Mon Sep 17 00:00:00 2001 From: binwiederhier Date: Thu, 28 May 2026 17:13:14 -0400 Subject: [PATCH 01/18] WIP: Access cache --- cmd/access_test.go | 7 ++ cmd/user.go | 5 + cmd/user_test.go | 5 + server/config.go | 2 + server/server.go | 21 ++-- user/access_cache.go | 180 ++++++++++++++++++++++++++ user/access_cache_test.go | 259 ++++++++++++++++++++++++++++++++++++++ user/manager.go | 158 ++++++++++++++++------- user/manager_postgres.go | 8 +- user/manager_sqlite.go | 8 +- user/types.go | 12 +- 11 files changed, 601 insertions(+), 64 deletions(-) create mode 100644 user/access_cache.go create mode 100644 user/access_cache_test.go diff --git a/cmd/access_test.go b/cmd/access_test.go index 8810b6b3..f280a9e9 100644 --- a/cmd/access_test.go +++ b/cmd/access_test.go @@ -7,6 +7,7 @@ import ( "heckel.io/ntfy/v2/server" "heckel.io/ntfy/v2/test" "testing" + "time" ) func TestCLI_Access_Show(t *testing.T) { @@ -43,6 +44,12 @@ user * (role: anonymous, tier: none) ` require.Equal(t, expected, stdout.String()) + // The CLI commands above ran against a separate Manager instance (their own + // process-equivalent), so the server's ACL cache hasn't seen the new grants + // yet. Wait for the server's background reloader (interval set in + // newTestServerWithAuth) to pick them up. + time.Sleep(150 * time.Millisecond) + // See if access permissions match app, _, _, _ = newTestApp() require.Error(t, app.Run([]string{ diff --git a/cmd/user.go b/cmd/user.go index cd6cf795..9bb0c5b0 100644 --- a/cmd/user.go +++ b/cmd/user.go @@ -378,6 +378,11 @@ func createUserManager(c *cli.Context) (*user.Manager, error) { ProvisionEnabled: false, // Hack: Do not re-provision users on manager initialization BcryptCost: user.DefaultUserPasswordBcryptCost, QueueWriterInterval: user.DefaultUserStatsQueueWriterInterval, + // CLI Managers are short-lived; the background ACL cache poller would only + // spam "database is closed" warnings after the subcommand returns. Mutations + // still refresh the local cache synchronously; the running server (if any) + // picks them up via its own poller. + AccessCacheReloadInterval: -1, } if databaseURL != "" { host, dbErr := pg.Open(databaseURL) diff --git a/cmd/user_test.go b/cmd/user_test.go index ed6f5de4..5b5bff73 100644 --- a/cmd/user_test.go +++ b/cmd/user_test.go @@ -9,6 +9,7 @@ import ( "os" "path/filepath" "testing" + "time" ) func TestCLI_User_Add(t *testing.T) { @@ -128,6 +129,10 @@ func newTestServerWithAuth(t *testing.T) (s *server.Server, conf *server.Config, conf.File = configFile conf.AuthFile = filepath.Join(t.TempDir(), "user.db") conf.AuthDefault = user.PermissionDenyAll + // Tight interval so cross-process writes from the `ntfy access`/`ntfy user` + // CLI commands (which run via a separate Manager) propagate to the server's + // ACL cache within tens of ms instead of the default 5s. + conf.AuthAccessCacheReloadInterval = 25 * time.Millisecond s, port = test.StartServerWithConfig(t, conf) return } diff --git a/server/config.go b/server/config.go index 1cfed3fc..f6f13eca 100644 --- a/server/config.go +++ b/server/config.go @@ -116,6 +116,7 @@ type Config struct { AuthTokens map[string][]*user.Token AuthBcryptCost int AuthStatsQueueWriterInterval time.Duration + AuthAccessCacheReloadInterval time.Duration AttachmentCacheDir string AttachmentTotalSizeLimit int64 AttachmentFileSizeLimit int64 @@ -223,6 +224,7 @@ func NewConfig() *Config { AuthDefault: user.PermissionReadWrite, AuthBcryptCost: user.DefaultUserPasswordBcryptCost, AuthStatsQueueWriterInterval: user.DefaultUserStatsQueueWriterInterval, + AuthAccessCacheReloadInterval: user.DefaultAccessCacheReloadInterval, AttachmentCacheDir: "", AttachmentTotalSizeLimit: DefaultAttachmentTotalSizeLimit, AttachmentFileSizeLimit: DefaultAttachmentFileSizeLimit, diff --git a/server/server.go b/server/server.go index 7ca0b4e7..7bcbbb09 100644 --- a/server/server.go +++ b/server/server.go @@ -247,16 +247,17 @@ func New(conf *Config) (*Server, error) { var userManager *user.Manager if conf.AuthFile != "" || pool != nil { authConfig := &user.Config{ - Filename: conf.AuthFile, - DatabaseURL: conf.DatabaseURL, - StartupQueries: conf.AuthStartupQueries, - DefaultAccess: conf.AuthDefault, - ProvisionEnabled: true, // Enable provisioning of users and access - Users: conf.AuthUsers, - Access: conf.AuthAccess, - Tokens: conf.AuthTokens, - BcryptCost: conf.AuthBcryptCost, - QueueWriterInterval: conf.AuthStatsQueueWriterInterval, + Filename: conf.AuthFile, + DatabaseURL: conf.DatabaseURL, + StartupQueries: conf.AuthStartupQueries, + DefaultAccess: conf.AuthDefault, + ProvisionEnabled: true, // Enable provisioning of users and access + Users: conf.AuthUsers, + Access: conf.AuthAccess, + Tokens: conf.AuthTokens, + BcryptCost: conf.AuthBcryptCost, + QueueWriterInterval: conf.AuthStatsQueueWriterInterval, + AccessCacheReloadInterval: conf.AuthAccessCacheReloadInterval, } if pool != nil { userManager, err = user.NewPostgresManager(pool, authConfig) diff --git a/user/access_cache.go b/user/access_cache.go new file mode 100644 index 00000000..98942d05 --- /dev/null +++ b/user/access_cache.go @@ -0,0 +1,180 @@ +package user + +import ( + "regexp" + "strings" + "sync/atomic" + + "heckel.io/ntfy/v2/db" +) + +// aclEntry mirrors one user_access row in the in-memory snapshot. +// +// topic is the raw stored value: it may contain \_ escapes (for literal underscores) +// and % wildcards (translated from user-supplied *). For exact-match entries (no %) +// matcher is nil and the entry is keyed by topic in aclSnapshot.exact. For wildcard +// entries (with %) matcher is the pre-compiled regex equivalent of the LIKE pattern. +type aclEntry struct { + topic string + read bool + write bool + matcher *regexp.Regexp +} + +// aclSnapshot is an immutable indexed form of the entire user_access table. +// +// exact[userName][escapedTopic] returns the matching entry in O(1) for the common +// case where the requested topic appears verbatim in some rule. The key is the +// stored form of the topic (i.e. with \_ escapes), so callers must pass topics +// through escapeUnderscore before probing. +// +// wildcards[userName] is the linear scan list of %-bearing rules for that user. +// Walked per request; trivially small in practice. Wildcards are NOT u_everyone- +// only -- any user can create them. +type aclSnapshot struct { + exact map[string]map[string]aclEntry + wildcards map[string][]aclEntry +} + +// aclCache holds the current snapshot behind an atomic pointer so that the hot +// path (Lookup) is lock-free. reload builds a fresh snapshot off the request +// path and atomically swaps the pointer; the old snapshot is GC'd once in-flight +// Lookups release their references. +// +// A nil receiver behaves as if no snapshot were loaded -- Lookup returns +// found=false, which the caller then resolves via DefaultAccess. This keeps +// tests and edge cases (e.g. early-startup) safe. +type aclCache struct { + snap atomic.Pointer[aclSnapshot] +} + +func newAccessCache() *aclCache { + return &aclCache{} +} + +// reload runs the bulk-load query against the primary and atomically swaps in +// a fresh snapshot. The primary is used (not ReadOnly) so a reload immediately +// after an ACL mutation sees the freshly-written rows without replica lag. +func (c *aclCache) reload(d *db.DB, query string) error { + rows, err := d.Query(query) + if err != nil { + return err + } + defer rows.Close() + snap := &aclSnapshot{ + exact: make(map[string]map[string]aclEntry), + wildcards: make(map[string][]aclEntry), + } + for rows.Next() { + var userName string + var entry aclEntry + if err := rows.Scan(&userName, &entry.topic, &entry.read, &entry.write); err != nil { + return err + } + if strings.Contains(entry.topic, "%") { + entry.matcher = compileLikeToRegex(entry.topic) + snap.wildcards[userName] = append(snap.wildcards[userName], entry) + } else { + if snap.exact[userName] == nil { + snap.exact[userName] = make(map[string]aclEntry) + } + snap.exact[userName][entry.topic] = entry + } + } + if err := rows.Err(); err != nil { + return err + } + c.snap.Store(snap) + return nil +} + +// Lookup returns the effective (read, write, found) permission for the given +// (username, topic), preserving the priority ordering of the original SQL query: +// 1. specific user beats Everyone +// 2. longer pattern beats shorter (more specific wins) +// 3. write beats read at equal length (write is "stronger") +func (c *aclCache) Lookup(usernameOrEveryone, topic string) (read, write, found bool) { + if c == nil { + return false, false, false + } + snap := c.snap.Load() + if snap == nil { + return false, false, false + } + // Pre-compute the escaped form once: exact-match keys in the snapshot are + // stored as toSQLWildcard would emit them (literal _ -> \_), so the + // incoming topic must be escaped the same way before map lookup. + escaped := escapeUnderscore(topic) + + // Specific user takes priority over Everyone. Skip the first lookup when + // the request is already anonymous to avoid scanning the same map twice. + if usernameOrEveryone != Everyone { + if e, ok := pickBest(snap, usernameOrEveryone, topic, escaped); ok { + return e.read, e.write, true + } + } + if e, ok := pickBest(snap, Everyone, topic, escaped); ok { + return e.read, e.write, true + } + return false, false, false +} + +// pickBest returns the highest-priority entry for a single user, combining the +// exact-match O(1) probe with a linear scan over the (usually empty or tiny) +// wildcard list. Priority within a user: longer pattern wins; write wins ties. +func pickBest(snap *aclSnapshot, userName, topic, escaped string) (aclEntry, bool) { + var best aclEntry + var found bool + if m, ok := snap.exact[userName]; ok { + if e, ok := m[escaped]; ok { + best, found = e, true + } + } + for _, w := range snap.wildcards[userName] { + if !w.matcher.MatchString(topic) { + continue + } + if !found || better(w, best) { + best, found = w, true + } + } + return best, found +} + +// better implements the (length DESC, write DESC) tie-break used by the original +// query's ORDER BY for entries owned by the same user. +func better(a, b aclEntry) bool { + if len(a.topic) != len(b.topic) { + return len(a.topic) > len(b.topic) + } + if a.write != b.write { + return a.write + } + return false +} + +// compileLikeToRegex converts a stored ntfy LIKE pattern into an equivalent Go +// regexp. In ntfy's stored form, % is the only wildcard (translated from *) and +// \_ is a literal underscore; no other backslashes occur. Topics themselves are +// restricted to [A-Za-z0-9_-] (see AllowedTopic), so neither % nor stray +// backslashes appear in user-supplied input. +func compileLikeToRegex(pattern string) *regexp.Regexp { + var sb strings.Builder + sb.WriteString("^") + i := 0 + for i < len(pattern) { + switch { + case pattern[i] == '\\' && i+1 < len(pattern) && pattern[i+1] == '_': + sb.WriteString(regexp.QuoteMeta("_")) + i += 2 + case pattern[i] == '%': + sb.WriteString(".*") + i++ + default: + sb.WriteString(regexp.QuoteMeta(string(pattern[i]))) + i++ + } + } + sb.WriteString("$") + return regexp.MustCompile(sb.String()) +} diff --git a/user/access_cache_test.go b/user/access_cache_test.go new file mode 100644 index 00000000..d8f11d71 --- /dev/null +++ b/user/access_cache_test.go @@ -0,0 +1,259 @@ +package user + +import ( + "sync" + "sync/atomic" + "testing" + + "github.com/stretchr/testify/require" +) + +// Cache-only unit tests. Integration with the Manager (loading from the DB, +// reload-after-mutation, end-to-end Authorize behavior) is covered by the +// existing TestStoreAuthorizeTopicAccess* tests in manager_test.go via +// forEachStoreBackend. + +func TestCompileLikeToRegex_Exact(t *testing.T) { + r := compileLikeToRegex("foo") + require.True(t, r.MatchString("foo")) + require.False(t, r.MatchString("foox")) + require.False(t, r.MatchString("xfoo")) +} + +func TestCompileLikeToRegex_TrailingPercent(t *testing.T) { + r := compileLikeToRegex("up%") + require.True(t, r.MatchString("up")) + require.True(t, r.MatchString("up123")) + require.False(t, r.MatchString("xup")) +} + +func TestCompileLikeToRegex_LeadingAndEmbeddedPercent(t *testing.T) { + r := compileLikeToRegex("%test%") + require.True(t, r.MatchString("test")) + require.True(t, r.MatchString("mytest")) + require.True(t, r.MatchString("testxxx")) + require.True(t, r.MatchString("xtestx")) + require.False(t, r.MatchString("nope")) +} + +func TestCompileLikeToRegex_EscapedUnderscore(t *testing.T) { + // "my\_topic" is the stored form of a literal "my_topic" -- the underscore + // must match itself, NOT act as a SQL one-character wildcard. + r := compileLikeToRegex(`my\_topic`) + require.True(t, r.MatchString("my_topic")) + require.False(t, r.MatchString("myXtopic")) + require.False(t, r.MatchString("mytopic")) +} + +func TestCompileLikeToRegex_EscapedUnderscoreAdjacentToPercent(t *testing.T) { + // "nz\_vip\_%" is the stored form of "nz_vip_*" -- literal "nz_vip_" prefix + // followed by any suffix. + r := compileLikeToRegex(`nz\_vip\_%`) + require.True(t, r.MatchString("nz_vip_")) + require.True(t, r.MatchString("nz_vip_alpha")) + require.False(t, r.MatchString("nz_vipX")) + require.False(t, r.MatchString("nzvip_alpha")) +} + +func TestCompileLikeToRegex_RegexMetaCharsInTopic(t *testing.T) { + // Topics in ntfy can include '-', which is benign, but make sure + // regex metacharacters in the pattern are escaped properly anyway. + r := compileLikeToRegex("foo-bar") + require.True(t, r.MatchString("foo-bar")) + require.False(t, r.MatchString("foo.bar")) // would match if '-' leaked into a character class +} + +func TestACLCache_LookupOnNilReceiverSafe(t *testing.T) { + var c *aclCache + read, write, found := c.Lookup("phil", "mytopic") + require.False(t, found) + require.False(t, read) + require.False(t, write) +} + +func TestACLCache_LookupBeforeReload(t *testing.T) { + // Before reload the snapshot pointer is nil. The cache treats this as + // "no rule found", which the caller resolves via DefaultAccess. + c := newAccessCache() + read, write, found := c.Lookup("phil", "mytopic") + require.False(t, found) + require.False(t, read) + require.False(t, write) +} + +func TestACLCache_ExactMatchHit(t *testing.T) { + c := newAccessCache() + c.snap.Store(buildSnapshot(t, []rawACLRow{ + {user: "phil", topic: "mytopic", read: true, write: true}, + })) + read, write, found := c.Lookup("phil", "mytopic") + require.True(t, found) + require.True(t, read) + require.True(t, write) +} + +func TestACLCache_ExactMatchMiss(t *testing.T) { + c := newAccessCache() + c.snap.Store(buildSnapshot(t, []rawACLRow{ + {user: "phil", topic: "mytopic", read: true, write: true}, + })) + _, _, found := c.Lookup("phil", "othertopic") + require.False(t, found) +} + +func TestACLCache_LiteralUnderscoreExactMatch(t *testing.T) { + // Stored as "my\_topic" (toSQLWildcard of "my_topic"). A literal underscore + // in the requested topic must match, while any other single char must not. + c := newAccessCache() + c.snap.Store(buildSnapshot(t, []rawACLRow{ + {user: "phil", topic: `my\_topic`, read: true, write: false}, + })) + read, write, found := c.Lookup("phil", "my_topic") + require.True(t, found) + require.True(t, read) + require.False(t, write) + + _, _, found = c.Lookup("phil", "myXtopic") + require.False(t, found) +} + +func TestACLCache_WildcardMatch(t *testing.T) { + c := newAccessCache() + c.snap.Store(buildSnapshot(t, []rawACLRow{ + {user: Everyone, topic: "up%", read: false, write: true}, + })) + read, write, found := c.Lookup("phil", "up42") + require.True(t, found) + require.False(t, read) + require.True(t, write) +} + +func TestACLCache_SpecificUserBeatsEveryone(t *testing.T) { + c := newAccessCache() + c.snap.Store(buildSnapshot(t, []rawACLRow{ + {user: Everyone, topic: "mytopic", read: true, write: false}, + {user: "phil", topic: "mytopic", read: false, write: false}, // deny-all for phil + })) + read, write, found := c.Lookup("phil", "mytopic") + require.True(t, found) + require.False(t, read) + require.False(t, write) +} + +func TestACLCache_AnonymousReadsEveryone(t *testing.T) { + c := newAccessCache() + c.snap.Store(buildSnapshot(t, []rawACLRow{ + {user: Everyone, topic: "announcements", read: true, write: false}, + })) + read, write, found := c.Lookup(Everyone, "announcements") + require.True(t, found) + require.True(t, read) + require.False(t, write) +} + +func TestACLCache_LongerPatternWinsForSameUser(t *testing.T) { + // Both rules belong to the same user (Everyone). The more specific (longer) + // "mytopic%" should beat the catch-all "%". + c := newAccessCache() + c.snap.Store(buildSnapshot(t, []rawACLRow{ + {user: Everyone, topic: "%", read: true, write: false}, + {user: Everyone, topic: "mytopic%", read: true, write: true}, + })) + read, write, found := c.Lookup(Everyone, "mytopicX") + require.True(t, found) + require.True(t, read) + require.True(t, write) +} + +func TestACLCache_WriteBeatsReadAtEqualLength(t *testing.T) { + // Two wildcard rules of identical length for the same user. The write rule + // should win the tie-break. + c := newAccessCache() + c.snap.Store(buildSnapshot(t, []rawACLRow{ + {user: Everyone, topic: "ab%", read: true, write: false}, + {user: Everyone, topic: "ab%", read: false, write: true}, // synthesized; impossible via real upsert but exercises the tie-break + })) + // One of the two will be the surviving exact-key entry (map collision keeps last); + // but the wildcard slice is what we want to exercise. Inject two wildcard entries + // directly to force the tie-break path. + c.snap.Store(&aclSnapshot{ + exact: map[string]map[string]aclEntry{}, + wildcards: map[string][]aclEntry{ + Everyone: { + {topic: "ab%", read: true, write: false, matcher: compileLikeToRegex("ab%")}, + {topic: "ab%", read: false, write: true, matcher: compileLikeToRegex("ab%")}, + }, + }, + }) + _, write, found := c.Lookup(Everyone, "abc") + require.True(t, found) + require.True(t, write) +} + +func TestACLCache_ConcurrentLookupAndReload(t *testing.T) { + // Atomic-pointer swap must be safe under concurrent reads. The race detector + // catches any unsafe shared mutation. + c := newAccessCache() + c.snap.Store(buildSnapshot(t, []rawACLRow{ + {user: Everyone, topic: "mytopic", read: true, write: true}, + })) + + var stop atomic.Bool + var wg sync.WaitGroup + wg.Add(2) + go func() { + defer wg.Done() + for !stop.Load() { + _, _, _ = c.Lookup(Everyone, "mytopic") + } + }() + go func() { + defer wg.Done() + for i := 0; i < 100; i++ { + c.snap.Store(buildSnapshot(t, []rawACLRow{ + {user: Everyone, topic: "mytopic", read: i%2 == 0, write: i%2 == 1}, + })) + } + stop.Store(true) + }() + wg.Wait() +} + +// rawACLRow + buildSnapshot mirror the rows that reload would Scan from the DB +// but avoid actually opening a DB for these unit tests. +type rawACLRow struct { + user string + topic string + read bool + write bool +} + +func buildSnapshot(t *testing.T, rows []rawACLRow) *aclSnapshot { + t.Helper() + snap := &aclSnapshot{ + exact: make(map[string]map[string]aclEntry), + wildcards: make(map[string][]aclEntry), + } + for _, r := range rows { + e := aclEntry{topic: r.topic, read: r.read, write: r.write} + if containsPercent(r.topic) { + e.matcher = compileLikeToRegex(r.topic) + snap.wildcards[r.user] = append(snap.wildcards[r.user], e) + } else { + if snap.exact[r.user] == nil { + snap.exact[r.user] = make(map[string]aclEntry) + } + snap.exact[r.user][r.topic] = e + } + } + return snap +} + +func containsPercent(s string) bool { + for i := 0; i < len(s); i++ { + if s[i] == '%' { + return true + } + } + return false +} diff --git a/user/manager.go b/user/manager.go index 303c7a49..52217e29 100644 --- a/user/manager.go +++ b/user/manager.go @@ -38,6 +38,11 @@ const ( const ( DefaultUserStatsQueueWriterInterval = 33 * time.Second DefaultUserPasswordBcryptCost = 10 + // DefaultAccessCacheReloadInterval bounds how stale the in-memory ACL snapshot + // can be relative to writes made by *other* processes (e.g. a separate `ntfy + // access` CLI invocation modifying the same database). Mutations performed + // by this Manager refresh the cache synchronously and do not depend on this. + DefaultAccessCacheReloadInterval = 5 * time.Second ) var ( @@ -48,12 +53,14 @@ var ( // Manager handles user authentication, authorization, and management type Manager struct { - config *Config - db *db.DB - queries queries - statsQueue map[string]*Stats // "Queue" to asynchronously write user stats to the database (UserID -> Stats) - tokenQueue map[string]*TokenUpdate // "Queue" to asynchronously write token access stats to the database (Token ID -> TokenUpdate) - mu sync.Mutex + config *Config + db *db.DB + queries queries + statsQueue map[string]*Stats // "Queue" to asynchronously write user stats to the database (UserID -> Stats) + tokenQueue map[string]*TokenUpdate // "Queue" to asynchronously write token access stats to the database (Token ID -> TokenUpdate) + accessCache *aclCache // In-memory snapshot of user_access; rebuilt after every ACL mutation + quit chan struct{} // Closed by Close() to signal background goroutines to stop + mu sync.Mutex } var _ Auther = (*Manager)(nil) @@ -65,20 +72,60 @@ func newManager(d *db.DB, queries queries, config *Config) (*Manager, error) { if config.QueueWriterInterval.Seconds() <= 0 { config.QueueWriterInterval = DefaultUserStatsQueueWriterInterval } + if config.AccessCacheReloadInterval == 0 { + config.AccessCacheReloadInterval = DefaultAccessCacheReloadInterval + } manager := &Manager{ - config: config, - db: d, - statsQueue: make(map[string]*Stats), - tokenQueue: make(map[string]*TokenUpdate), - queries: queries, + config: config, + db: d, + statsQueue: make(map[string]*Stats), + tokenQueue: make(map[string]*TokenUpdate), + accessCache: newAccessCache(), + quit: make(chan struct{}), + queries: queries, } if err := manager.maybeProvisionUsersAccessAndTokens(); err != nil { return nil, err } + // Populate the ACL cache after provisioning so the initial snapshot includes + // any provisioned access rules. Subsequent mutations call reloadAccessCache. + if err := manager.reloadAccessCache(); err != nil { + return nil, err + } go manager.asyncQueueWriter(manager.config.QueueWriterInterval) + if manager.config.AccessCacheReloadInterval > 0 { + go manager.asyncAccessCacheReloader(manager.config.AccessCacheReloadInterval) + } return manager, nil } +// reloadAccessCache rebuilds the in-memory ACL snapshot from the primary +// database. Called once at startup (after provisioning) and after every method +// that mutates user_access (directly or by cascade from "user" deletion). +func (a *Manager) reloadAccessCache() error { + return a.accessCache.reload(a.db, a.queries.selectAllAccessForCache) +} + +// asyncAccessCacheReloader periodically refreshes the ACL snapshot so that +// writes made by other processes against the same database (most notably the +// `ntfy access` CLI subcommand running while a server holds the cache) become +// visible within the configured interval. This Manager's own writes do not +// depend on the poller -- they refresh the cache synchronously. +func (a *Manager) asyncAccessCacheReloader(interval time.Duration) { + ticker := time.NewTicker(interval) + defer ticker.Stop() + for { + select { + case <-a.quit: + return + case <-ticker.C: + if err := a.reloadAccessCache(); err != nil { + log.Tag(tag).Err(err).Warn("Reloading ACL cache failed") + } + } + } +} + // Authenticate checks username and password and returns a User if correct, and the user has not been // marked as deleted. The method returns in constant-ish time, regardless of whether the user exists or // the password is correct or incorrect. @@ -151,9 +198,13 @@ func (a *Manager) RemoveUser(username string) error { if err := a.CanChangeUser(username); err != nil { return err } - return db.ExecTx(a.db, func(tx *sql.Tx) error { + if err := db.ExecTx(a.db, func(tx *sql.Tx) error { return a.removeUserTx(tx, username) - }) + }); err != nil { + return err + } + // user_access rows are cascade-deleted along with the user; refresh the snapshot. + return a.reloadAccessCache() } // removeUserTx deletes the user with the given username @@ -174,7 +225,7 @@ func (a *Manager) MarkUserRemoved(user *User) error { if !AllowedUsername(user.Name) { return ErrInvalidArgument } - return db.ExecTx(a.db, func(tx *sql.Tx) error { + if err := db.ExecTx(a.db, func(tx *sql.Tx) error { if err := a.resetUserAccessTx(tx, user.Name); err != nil { return err } @@ -186,7 +237,11 @@ func (a *Manager) MarkUserRemoved(user *User) error { return err } return nil - }) + }); err != nil { + return err + } + // resetUserAccessTx wiped this user's user_access rows; refresh the snapshot. + return a.reloadAccessCache() } // RemoveDeletedUsers deletes all users that have been marked deleted @@ -194,7 +249,8 @@ func (a *Manager) RemoveDeletedUsers() error { if _, err := a.db.Exec(a.queries.deleteUsersMarked, time.Now().Unix()); err != nil { return err } - return nil + // user_access rows are cascade-deleted with the users; refresh the snapshot. + return a.reloadAccessCache() } // ChangePassword changes a user's password @@ -225,9 +281,15 @@ func (a *Manager) ChangeRole(username string, role Role) error { if err := a.CanChangeUser(username); err != nil { return err } - return db.ExecTx(a.db, func(tx *sql.Tx) error { + if err := db.ExecTx(a.db, func(tx *sql.Tx) error { return a.changeRoleTx(tx, username, role) - }) + }); err != nil { + return err + } + // Promotion to admin clears user_access rows for the user; refresh the snapshot. + // Other role changes are no-ops for the cache but reloading is cheap and keeps + // the code path uniform. + return a.reloadAccessCache() } // changeRoleTx changes a user's role @@ -588,9 +650,12 @@ func (a *Manager) resolvePerms(base, perm Permission) error { // read/write access to a topic. The parameter topicPattern may include wildcards (*). The ACL entry // owner may either be a user (username), or the system (empty). func (a *Manager) AllowAccess(username string, topicPattern string, permission Permission) error { - return db.ExecTx(a.db, func(tx *sql.Tx) error { + if err := db.ExecTx(a.db, func(tx *sql.Tx) error { return a.allowAccessTx(tx, username, topicPattern, permission, false) - }) + }); err != nil { + return err + } + return a.reloadAccessCache() } func (a *Manager) allowAccessTx(tx *sql.Tx, username string, topicPattern string, permission Permission, provisioned bool) error { @@ -606,9 +671,12 @@ func (a *Manager) allowAccessTx(tx *sql.Tx, username string, topicPattern string // ResetAccess removes an access control list entry for a specific username/topic, or (if topic is // empty) for an entire user. The parameter topicPattern may include wildcards (*). func (a *Manager) ResetAccess(username string, topicPattern string) error { - return db.ExecTx(a.db, func(tx *sql.Tx) error { + if err := db.ExecTx(a.db, func(tx *sql.Tx) error { return a.resetAccessTx(tx, username, topicPattern) - }) + }); err != nil { + return err + } + return a.reloadAccessCache() } func (a *Manager) resetAccessTx(tx *sql.Tx, username string, topicPattern string) error { @@ -650,24 +718,15 @@ func (a *Manager) AllowReservation(username string, topic string) error { // authorizeTopicAccess returns the read/write permissions for the given username and topic. // The found return value indicates whether an ACL entry was found at all. // -// - The query may return two rows (one for everyone, and one for the user), but prioritizes the user. -// - Furthermore, the query prioritizes more specific permissions (longer!) over more generic ones, e.g. "test*" > "*" +// - The cache may contain two matching entries (one for everyone, and one for the user), but prioritizes the user. +// - Furthermore, the lookup prioritizes more specific permissions (longer!) over more generic ones, e.g. "test*" > "*" // - It also prioritizes write permissions over read permissions +// +// The lookup is served entirely from the in-memory snapshot maintained by accessCache, +// so this is on the hot path of every authenticatable HTTP request and must stay allocation-free. func (a *Manager) authorizeTopicAccess(usernameOrEveryone, topic string) (read, write, found bool, err error) { - rows, err := a.db.ReadOnly().Query(a.queries.selectTopicPerms, Everyone, usernameOrEveryone, topic) - if err != nil { - return false, false, false, err - } - defer rows.Close() - if !rows.Next() { - return false, false, false, nil - } - if err := rows.Scan(&read, &write); err != nil { - return false, false, false, err - } else if err := rows.Err(); err != nil { - return false, false, false, err - } - return read, write, true, nil + read, write, found = a.accessCache.Lookup(usernameOrEveryone, topic) + return read, write, found, nil } // AllGrants returns all user-specific access control entries, mapped to their respective user IDs @@ -731,7 +790,7 @@ func (a *Manager) AddReservation(username string, topic string, everyone Permiss if !AllowedUsername(username) || username == Everyone || !AllowedTopic(topic) { return ErrInvalidArgument } - return db.ExecTx(a.db, func(tx *sql.Tx) error { + if err := db.ExecTx(a.db, func(tx *sql.Tx) error { if limit > 0 { hasReservation, err := a.hasReservationTx(tx, username, topic) if err != nil { @@ -754,7 +813,10 @@ func (a *Manager) AddReservation(username string, topic string, everyone Permiss return err } return nil - }) + }); err != nil { + return err + } + return a.reloadAccessCache() } // RemoveReservations deletes the access control entries associated with the given username/topic, @@ -769,14 +831,17 @@ func (a *Manager) RemoveReservations(username string, topics ...string) error { return ErrInvalidArgument } } - return db.ExecTx(a.db, func(tx *sql.Tx) error { + if err := db.ExecTx(a.db, func(tx *sql.Tx) error { for _, topic := range topics { if err := a.removeReservationAccessTx(tx, username, topic); err != nil { return err } } return nil - }) + }); err != nil { + return err + } + return a.reloadAccessCache() } // Reservations returns all user-owned topics, and the associated everyone-access @@ -1515,8 +1580,15 @@ func (a *Manager) maybeProvisionTokens(tx *sql.Tx, provisionUsernames []string, return nil } -// Close closes the underlying database +// Close stops background goroutines and closes the underlying database. +// Safe to call multiple times. func (a *Manager) Close() error { + select { + case <-a.quit: + // already closed + default: + close(a.quit) + } return a.db.Close() } diff --git a/user/manager_postgres.go b/user/manager_postgres.go index 02cffd84..d3bb6f5f 100644 --- a/user/manager_postgres.go +++ b/user/manager_postgres.go @@ -70,12 +70,10 @@ const ( postgresDeleteUsersProvisionedQuery = `DELETE FROM "user" WHERE provisioned = true` // Access queries - postgresSelectTopicPermsQuery = ` - SELECT read, write + postgresSelectAllAccessForCacheQuery = ` + SELECT u.user_name, a.topic, a.read, a.write FROM user_access a JOIN "user" u ON u.id = a.user_id - WHERE (u.user_name = $1 OR u.user_name = $2) AND $3 LIKE a.topic ESCAPE '\' - ORDER BY u.user_name DESC, LENGTH(a.topic) DESC, CASE WHEN a.write THEN 1 ELSE 0 END DESC ` postgresSelectUserAllAccessQuery = ` SELECT user_id, topic, read, write, provisioned @@ -244,7 +242,7 @@ var postgresQueries = queries{ deleteUserTier: postgresDeleteUserTierQuery, deleteUsersMarked: postgresDeleteUsersMarkedQuery, deleteUsersProvisioned: postgresDeleteUsersProvisionedQuery, - selectTopicPerms: postgresSelectTopicPermsQuery, + selectAllAccessForCache: postgresSelectAllAccessForCacheQuery, selectUserAllAccess: postgresSelectUserAllAccessQuery, selectUserAccess: postgresSelectUserAccessQuery, selectUserReservations: postgresSelectUserReservationsQuery, diff --git a/user/manager_sqlite.go b/user/manager_sqlite.go index 0f1a9227..4d498ce2 100644 --- a/user/manager_sqlite.go +++ b/user/manager_sqlite.go @@ -76,12 +76,10 @@ const ( sqliteDeleteUsersProvisionedQuery = `DELETE FROM user WHERE provisioned = 1` // Access queries - sqliteSelectTopicPermsQuery = ` - SELECT read, write + sqliteSelectAllAccessForCacheQuery = ` + SELECT u.user, a.topic, a.read, a.write FROM user_access a JOIN user u ON u.id = a.user_id - WHERE (u.user = ? OR u.user = ?) AND ? LIKE a.topic ESCAPE '\' - ORDER BY u.user DESC, LENGTH(a.topic) DESC, a.write DESC ` sqliteSelectUserAllAccessQuery = ` SELECT user_id, topic, read, write, provisioned @@ -242,7 +240,7 @@ var sqliteQueries = queries{ deleteUserTier: sqliteDeleteUserTierQuery, deleteUsersMarked: sqliteDeleteUsersMarkedQuery, deleteUsersProvisioned: sqliteDeleteUsersProvisionedQuery, - selectTopicPerms: sqliteSelectTopicPermsQuery, + selectAllAccessForCache: sqliteSelectAllAccessForCacheQuery, selectUserAllAccess: sqliteSelectUserAllAccessQuery, selectUserAccess: sqliteSelectUserAccessQuery, selectUserReservations: sqliteSelectUserReservationsQuery, diff --git a/user/types.go b/user/types.go index d0d40e33..6d2650ed 100644 --- a/user/types.go +++ b/user/types.go @@ -255,6 +255,16 @@ type Config struct { Tokens map[string][]*Token // Predefined users to create on startup (username -> []*Token) QueueWriterInterval time.Duration // Interval for the async queue writer to flush stats and token updates to the database BcryptCost int // Cost of generated passwords; lowering makes testing faster + + // AccessCacheReloadInterval bounds the staleness of the in-memory ACL cache + // relative to writes from other processes (e.g. `ntfy access` CLI against a + // running server). + // 0 -> use DefaultAccessCacheReloadInterval + // negative -> disable the background poller; cache only refreshes on this + // Manager's own ACL mutations. Use this for short-lived Managers + // (e.g. the CLI subcommands) where polling is wasted work. + // positive -> poll at the given interval + AccessCacheReloadInterval time.Duration } // Error constants used by the package @@ -303,7 +313,7 @@ type queries struct { deleteUsersProvisioned string // Access queries - selectTopicPerms string + selectAllAccessForCache string // Bulk load: (user_name, topic, read, write) for the in-memory ACL cache selectUserAllAccess string selectUserAccess string selectUserReservations string From 6310e3a96f4baaf7585e1c67427c8bf70bd3f565 Mon Sep 17 00:00:00 2001 From: binwiederhier Date: Sun, 31 May 2026 10:45:10 -0400 Subject: [PATCH 02/18] Replace atomic.Pointer for readbility --- user/access_cache.go | 128 ++++++++++++++++++-------------------- user/access_cache_test.go | 119 ++++++++++++++++++----------------- 2 files changed, 125 insertions(+), 122 deletions(-) diff --git a/user/access_cache.go b/user/access_cache.go index 98942d05..d73f48a1 100644 --- a/user/access_cache.go +++ b/user/access_cache.go @@ -3,105 +3,101 @@ package user import ( "regexp" "strings" - "sync/atomic" + "sync" "heckel.io/ntfy/v2/db" ) -// aclEntry mirrors one user_access row in the in-memory snapshot. +// aclCache is an in-memory index over the entire user_access table. +// +// exact[username][escapedTopic] returns the matching entry in O(1) for the common +// case where the requested topic appears verbatim in some rule. The key is the +// stored form of the topic (i.e. with \_ escapes), so Lookup escapes incoming +// topics through escapeUnderscore before probing. +// +// wildcard[username] is the linear-scan list of %-bearing rules for that user. +// Walked per request; trivially small in practice. Wildcards are NOT u_everyone- +// only -- any user can create them. +type aclCache struct { + exact map[string]map[string]aclEntry + wildcard map[string][]aclEntry + mu sync.RWMutex // Protect exact and wildcard +} + +// aclEntry mirrors one user_access row in the in-memory cache. // // topic is the raw stored value: it may contain \_ escapes (for literal underscores) // and % wildcards (translated from user-supplied *). For exact-match entries (no %) -// matcher is nil and the entry is keyed by topic in aclSnapshot.exact. For wildcard +// matcher is nil and the entry is keyed by topic in aclCache.exact. For wildcard // entries (with %) matcher is the pre-compiled regex equivalent of the LIKE pattern. type aclEntry struct { topic string read bool write bool - matcher *regexp.Regexp -} - -// aclSnapshot is an immutable indexed form of the entire user_access table. -// -// exact[userName][escapedTopic] returns the matching entry in O(1) for the common -// case where the requested topic appears verbatim in some rule. The key is the -// stored form of the topic (i.e. with \_ escapes), so callers must pass topics -// through escapeUnderscore before probing. -// -// wildcards[userName] is the linear scan list of %-bearing rules for that user. -// Walked per request; trivially small in practice. Wildcards are NOT u_everyone- -// only -- any user can create them. -type aclSnapshot struct { - exact map[string]map[string]aclEntry - wildcards map[string][]aclEntry -} - -// aclCache holds the current snapshot behind an atomic pointer so that the hot -// path (Lookup) is lock-free. reload builds a fresh snapshot off the request -// path and atomically swaps the pointer; the old snapshot is GC'd once in-flight -// Lookups release their references. -// -// A nil receiver behaves as if no snapshot were loaded -- Lookup returns -// found=false, which the caller then resolves via DefaultAccess. This keeps -// tests and edge cases (e.g. early-startup) safe. -type aclCache struct { - snap atomic.Pointer[aclSnapshot] + matcher *regexp.Regexp // Nil for exact entries } func newAccessCache() *aclCache { - return &aclCache{} + return &aclCache{ + exact: make(map[string]map[string]aclEntry), + wildcard: make(map[string][]aclEntry), + } } -// reload runs the bulk-load query against the primary and atomically swaps in -// a fresh snapshot. The primary is used (not ReadOnly) so a reload immediately -// after an ACL mutation sees the freshly-written rows without replica lag. +// reload runs the bulk-load query against the primary and swaps in freshly-built +// exact and wildcard maps under the write lock. The primary is used (not +// ReadOnly) so a reload immediately after an ACL mutation sees the freshly- +// written rows without replica lag. func (c *aclCache) reload(d *db.DB, query string) error { rows, err := d.Query(query) if err != nil { return err } defer rows.Close() - snap := &aclSnapshot{ - exact: make(map[string]map[string]aclEntry), - wildcards: make(map[string][]aclEntry), - } + exact := make(map[string]map[string]aclEntry) + wildcards := make(map[string][]aclEntry) for rows.Next() { - var userName string + var username string var entry aclEntry - if err := rows.Scan(&userName, &entry.topic, &entry.read, &entry.write); err != nil { + if err := rows.Scan(&username, &entry.topic, &entry.read, &entry.write); err != nil { return err } if strings.Contains(entry.topic, "%") { - entry.matcher = compileLikeToRegex(entry.topic) - snap.wildcards[userName] = append(snap.wildcards[userName], entry) - } else { - if snap.exact[userName] == nil { - snap.exact[userName] = make(map[string]aclEntry) + re, err := compileLikeToRegex(entry.topic) + if err != nil { + return err } - snap.exact[userName][entry.topic] = entry + entry.matcher = re + wildcards[username] = append(wildcards[username], entry) + } else { + if exact[username] == nil { + exact[username] = make(map[string]aclEntry) + } + exact[username][entry.topic] = entry } } if err := rows.Err(); err != nil { return err } - c.snap.Store(snap) + c.mu.Lock() + c.exact = exact + c.wildcard = wildcards + c.mu.Unlock() return nil } // Lookup returns the effective (read, write, found) permission for the given // (username, topic), preserving the priority ordering of the original SQL query: -// 1. specific user beats Everyone -// 2. longer pattern beats shorter (more specific wins) -// 3. write beats read at equal length (write is "stronger") +// 1. specific user beats Everyone +// 2. longer pattern beats shorter (more specific wins) +// 3. write beats read at equal length (write is "stronger") func (c *aclCache) Lookup(usernameOrEveryone, topic string) (read, write, found bool) { if c == nil { return false, false, false } - snap := c.snap.Load() - if snap == nil { - return false, false, false - } - // Pre-compute the escaped form once: exact-match keys in the snapshot are + c.mu.RLock() + defer c.mu.RUnlock() + // Pre-compute the escaped form once: exact-match keys in the cache are // stored as toSQLWildcard would emit them (literal _ -> \_), so the // incoming topic must be escaped the same way before map lookup. escaped := escapeUnderscore(topic) @@ -109,28 +105,28 @@ func (c *aclCache) Lookup(usernameOrEveryone, topic string) (read, write, found // Specific user takes priority over Everyone. Skip the first lookup when // the request is already anonymous to avoid scanning the same map twice. if usernameOrEveryone != Everyone { - if e, ok := pickBest(snap, usernameOrEveryone, topic, escaped); ok { + if e, ok := c.pickBestLocked(usernameOrEveryone, topic, escaped); ok { return e.read, e.write, true } } - if e, ok := pickBest(snap, Everyone, topic, escaped); ok { + if e, ok := c.pickBestLocked(Everyone, topic, escaped); ok { return e.read, e.write, true } return false, false, false } -// pickBest returns the highest-priority entry for a single user, combining the -// exact-match O(1) probe with a linear scan over the (usually empty or tiny) -// wildcard list. Priority within a user: longer pattern wins; write wins ties. -func pickBest(snap *aclSnapshot, userName, topic, escaped string) (aclEntry, bool) { +// pickBestLocked returns the highest-priority entry for a single user, combining +// the exact-match O(1) probe with a linear scan over the (usually empty or tiny) +// wildcard list. Caller must hold c.mu (RLock is sufficient). +func (c *aclCache) pickBestLocked(username, topic, escaped string) (aclEntry, bool) { var best aclEntry var found bool - if m, ok := snap.exact[userName]; ok { + if m, ok := c.exact[username]; ok { if e, ok := m[escaped]; ok { best, found = e, true } } - for _, w := range snap.wildcards[userName] { + for _, w := range c.wildcard[username] { if !w.matcher.MatchString(topic) { continue } @@ -158,7 +154,7 @@ func better(a, b aclEntry) bool { // \_ is a literal underscore; no other backslashes occur. Topics themselves are // restricted to [A-Za-z0-9_-] (see AllowedTopic), so neither % nor stray // backslashes appear in user-supplied input. -func compileLikeToRegex(pattern string) *regexp.Regexp { +func compileLikeToRegex(pattern string) (*regexp.Regexp, error) { var sb strings.Builder sb.WriteString("^") i := 0 @@ -176,5 +172,5 @@ func compileLikeToRegex(pattern string) *regexp.Regexp { } } sb.WriteString("$") - return regexp.MustCompile(sb.String()) + return regexp.Compile(sb.String()) } diff --git a/user/access_cache_test.go b/user/access_cache_test.go index d8f11d71..6c617a0c 100644 --- a/user/access_cache_test.go +++ b/user/access_cache_test.go @@ -1,6 +1,7 @@ package user import ( + "regexp" "sync" "sync/atomic" "testing" @@ -14,21 +15,21 @@ import ( // forEachStoreBackend. func TestCompileLikeToRegex_Exact(t *testing.T) { - r := compileLikeToRegex("foo") + r := mustCompileLikeToRegex(t, "foo") require.True(t, r.MatchString("foo")) require.False(t, r.MatchString("foox")) require.False(t, r.MatchString("xfoo")) } func TestCompileLikeToRegex_TrailingPercent(t *testing.T) { - r := compileLikeToRegex("up%") + r := mustCompileLikeToRegex(t, "up%") require.True(t, r.MatchString("up")) require.True(t, r.MatchString("up123")) require.False(t, r.MatchString("xup")) } func TestCompileLikeToRegex_LeadingAndEmbeddedPercent(t *testing.T) { - r := compileLikeToRegex("%test%") + r := mustCompileLikeToRegex(t, "%test%") require.True(t, r.MatchString("test")) require.True(t, r.MatchString("mytest")) require.True(t, r.MatchString("testxxx")) @@ -39,7 +40,7 @@ func TestCompileLikeToRegex_LeadingAndEmbeddedPercent(t *testing.T) { func TestCompileLikeToRegex_EscapedUnderscore(t *testing.T) { // "my\_topic" is the stored form of a literal "my_topic" -- the underscore // must match itself, NOT act as a SQL one-character wildcard. - r := compileLikeToRegex(`my\_topic`) + r := mustCompileLikeToRegex(t, `my\_topic`) require.True(t, r.MatchString("my_topic")) require.False(t, r.MatchString("myXtopic")) require.False(t, r.MatchString("mytopic")) @@ -48,7 +49,7 @@ func TestCompileLikeToRegex_EscapedUnderscore(t *testing.T) { func TestCompileLikeToRegex_EscapedUnderscoreAdjacentToPercent(t *testing.T) { // "nz\_vip\_%" is the stored form of "nz_vip_*" -- literal "nz_vip_" prefix // followed by any suffix. - r := compileLikeToRegex(`nz\_vip\_%`) + r := mustCompileLikeToRegex(t, `nz\_vip\_%`) require.True(t, r.MatchString("nz_vip_")) require.True(t, r.MatchString("nz_vip_alpha")) require.False(t, r.MatchString("nz_vipX")) @@ -58,7 +59,7 @@ func TestCompileLikeToRegex_EscapedUnderscoreAdjacentToPercent(t *testing.T) { func TestCompileLikeToRegex_RegexMetaCharsInTopic(t *testing.T) { // Topics in ntfy can include '-', which is benign, but make sure // regex metacharacters in the pattern are escaped properly anyway. - r := compileLikeToRegex("foo-bar") + r := mustCompileLikeToRegex(t, "foo-bar") require.True(t, r.MatchString("foo-bar")) require.False(t, r.MatchString("foo.bar")) // would match if '-' leaked into a character class } @@ -72,8 +73,9 @@ func TestACLCache_LookupOnNilReceiverSafe(t *testing.T) { } func TestACLCache_LookupBeforeReload(t *testing.T) { - // Before reload the snapshot pointer is nil. The cache treats this as - // "no rule found", which the caller resolves via DefaultAccess. + // A freshly-constructed cache has empty exact and wildcards maps. The + // cache treats this as "no rule found", which the caller resolves via + // DefaultAccess. c := newAccessCache() read, write, found := c.Lookup("phil", "mytopic") require.False(t, found) @@ -83,9 +85,9 @@ func TestACLCache_LookupBeforeReload(t *testing.T) { func TestACLCache_ExactMatchHit(t *testing.T) { c := newAccessCache() - c.snap.Store(buildSnapshot(t, []rawACLRow{ + loadCache(t, c, []rawACLRow{ {user: "phil", topic: "mytopic", read: true, write: true}, - })) + }) read, write, found := c.Lookup("phil", "mytopic") require.True(t, found) require.True(t, read) @@ -94,9 +96,9 @@ func TestACLCache_ExactMatchHit(t *testing.T) { func TestACLCache_ExactMatchMiss(t *testing.T) { c := newAccessCache() - c.snap.Store(buildSnapshot(t, []rawACLRow{ + loadCache(t, c, []rawACLRow{ {user: "phil", topic: "mytopic", read: true, write: true}, - })) + }) _, _, found := c.Lookup("phil", "othertopic") require.False(t, found) } @@ -105,9 +107,9 @@ func TestACLCache_LiteralUnderscoreExactMatch(t *testing.T) { // Stored as "my\_topic" (toSQLWildcard of "my_topic"). A literal underscore // in the requested topic must match, while any other single char must not. c := newAccessCache() - c.snap.Store(buildSnapshot(t, []rawACLRow{ + loadCache(t, c, []rawACLRow{ {user: "phil", topic: `my\_topic`, read: true, write: false}, - })) + }) read, write, found := c.Lookup("phil", "my_topic") require.True(t, found) require.True(t, read) @@ -119,9 +121,9 @@ func TestACLCache_LiteralUnderscoreExactMatch(t *testing.T) { func TestACLCache_WildcardMatch(t *testing.T) { c := newAccessCache() - c.snap.Store(buildSnapshot(t, []rawACLRow{ + loadCache(t, c, []rawACLRow{ {user: Everyone, topic: "up%", read: false, write: true}, - })) + }) read, write, found := c.Lookup("phil", "up42") require.True(t, found) require.False(t, read) @@ -130,10 +132,10 @@ func TestACLCache_WildcardMatch(t *testing.T) { func TestACLCache_SpecificUserBeatsEveryone(t *testing.T) { c := newAccessCache() - c.snap.Store(buildSnapshot(t, []rawACLRow{ + loadCache(t, c, []rawACLRow{ {user: Everyone, topic: "mytopic", read: true, write: false}, {user: "phil", topic: "mytopic", read: false, write: false}, // deny-all for phil - })) + }) read, write, found := c.Lookup("phil", "mytopic") require.True(t, found) require.False(t, read) @@ -142,9 +144,9 @@ func TestACLCache_SpecificUserBeatsEveryone(t *testing.T) { func TestACLCache_AnonymousReadsEveryone(t *testing.T) { c := newAccessCache() - c.snap.Store(buildSnapshot(t, []rawACLRow{ + loadCache(t, c, []rawACLRow{ {user: Everyone, topic: "announcements", read: true, write: false}, - })) + }) read, write, found := c.Lookup(Everyone, "announcements") require.True(t, found) require.True(t, read) @@ -155,10 +157,10 @@ func TestACLCache_LongerPatternWinsForSameUser(t *testing.T) { // Both rules belong to the same user (Everyone). The more specific (longer) // "mytopic%" should beat the catch-all "%". c := newAccessCache() - c.snap.Store(buildSnapshot(t, []rawACLRow{ + loadCache(t, c, []rawACLRow{ {user: Everyone, topic: "%", read: true, write: false}, {user: Everyone, topic: "mytopic%", read: true, write: true}, - })) + }) read, write, found := c.Lookup(Everyone, "mytopicX") require.True(t, found) require.True(t, read) @@ -167,36 +169,31 @@ func TestACLCache_LongerPatternWinsForSameUser(t *testing.T) { func TestACLCache_WriteBeatsReadAtEqualLength(t *testing.T) { // Two wildcard rules of identical length for the same user. The write rule - // should win the tie-break. + // should win the tie-break. The two-rows-with-same-topic shape is + // impossible via real upsert (pkey would conflict), so we inject the entries + // directly into the cache's wildcard slice. c := newAccessCache() - c.snap.Store(buildSnapshot(t, []rawACLRow{ - {user: Everyone, topic: "ab%", read: true, write: false}, - {user: Everyone, topic: "ab%", read: false, write: true}, // synthesized; impossible via real upsert but exercises the tie-break - })) - // One of the two will be the surviving exact-key entry (map collision keeps last); - // but the wildcard slice is what we want to exercise. Inject two wildcard entries - // directly to force the tie-break path. - c.snap.Store(&aclSnapshot{ - exact: map[string]map[string]aclEntry{}, - wildcards: map[string][]aclEntry{ - Everyone: { - {topic: "ab%", read: true, write: false, matcher: compileLikeToRegex("ab%")}, - {topic: "ab%", read: false, write: true, matcher: compileLikeToRegex("ab%")}, - }, + c.mu.Lock() + c.exact = map[string]map[string]aclEntry{} + c.wildcard = map[string][]aclEntry{ + Everyone: { + {topic: "ab%", read: true, write: false, matcher: mustCompileLikeToRegex(t, "ab%")}, + {topic: "ab%", read: false, write: true, matcher: mustCompileLikeToRegex(t, "ab%")}, }, - }) + } + c.mu.Unlock() _, write, found := c.Lookup(Everyone, "abc") require.True(t, found) require.True(t, write) } func TestACLCache_ConcurrentLookupAndReload(t *testing.T) { - // Atomic-pointer swap must be safe under concurrent reads. The race detector + // Lock-based swap must be safe under concurrent reads. The race detector // catches any unsafe shared mutation. c := newAccessCache() - c.snap.Store(buildSnapshot(t, []rawACLRow{ + loadCache(t, c, []rawACLRow{ {user: Everyone, topic: "mytopic", read: true, write: true}, - })) + }) var stop atomic.Bool var wg sync.WaitGroup @@ -210,17 +207,17 @@ func TestACLCache_ConcurrentLookupAndReload(t *testing.T) { go func() { defer wg.Done() for i := 0; i < 100; i++ { - c.snap.Store(buildSnapshot(t, []rawACLRow{ + loadCache(t, c, []rawACLRow{ {user: Everyone, topic: "mytopic", read: i%2 == 0, write: i%2 == 1}, - })) + }) } stop.Store(true) }() wg.Wait() } -// rawACLRow + buildSnapshot mirror the rows that reload would Scan from the DB -// but avoid actually opening a DB for these unit tests. +// rawACLRow models the rows that reload would Scan from the DB but avoids +// actually opening a DB for these unit tests. type rawACLRow struct { user string topic string @@ -228,25 +225,35 @@ type rawACLRow struct { write bool } -func buildSnapshot(t *testing.T, rows []rawACLRow) *aclSnapshot { +// loadCache writes the given rows into the cache under its write lock, +// preserving the same exact/wildcard partitioning that reload would produce. +func loadCache(t *testing.T, c *aclCache, rows []rawACLRow) { t.Helper() - snap := &aclSnapshot{ - exact: make(map[string]map[string]aclEntry), - wildcards: make(map[string][]aclEntry), - } + exact := make(map[string]map[string]aclEntry) + wildcards := make(map[string][]aclEntry) for _, r := range rows { e := aclEntry{topic: r.topic, read: r.read, write: r.write} if containsPercent(r.topic) { - e.matcher = compileLikeToRegex(r.topic) - snap.wildcards[r.user] = append(snap.wildcards[r.user], e) + e.matcher = mustCompileLikeToRegex(t, r.topic) + wildcards[r.user] = append(wildcards[r.user], e) } else { - if snap.exact[r.user] == nil { - snap.exact[r.user] = make(map[string]aclEntry) + if exact[r.user] == nil { + exact[r.user] = make(map[string]aclEntry) } - snap.exact[r.user][r.topic] = e + exact[r.user][r.topic] = e } } - return snap + c.mu.Lock() + c.exact = exact + c.wildcard = wildcards + c.mu.Unlock() +} + +func mustCompileLikeToRegex(t *testing.T, pattern string) *regexp.Regexp { + t.Helper() + r, err := compileLikeToRegex(pattern) + require.NoError(t, err) + return r } func containsPercent(s string) bool { From c841caa3b335884c3f13dd10e0cf46cb17efddaf Mon Sep 17 00:00:00 2001 From: binwiederhier Date: Sun, 31 May 2026 11:05:27 -0400 Subject: [PATCH 03/18] Review, rename to pattern --- user/access_cache.go | 105 ++++++++++++++++++-------------------- user/access_cache_test.go | 12 ++--- 2 files changed, 57 insertions(+), 60 deletions(-) diff --git a/user/access_cache.go b/user/access_cache.go index d73f48a1..adcf6e02 100644 --- a/user/access_cache.go +++ b/user/access_cache.go @@ -15,37 +15,43 @@ import ( // stored form of the topic (i.e. with \_ escapes), so Lookup escapes incoming // topics through escapeUnderscore before probing. // -// wildcard[username] is the linear-scan list of %-bearing rules for that user. +// pattern[username] is the linear-scan list of %-bearing rules for that user. // Walked per request; trivially small in practice. Wildcards are NOT u_everyone- // only -- any user can create them. type aclCache struct { - exact map[string]map[string]aclEntry - wildcard map[string][]aclEntry - mu sync.RWMutex // Protect exact and wildcard + exact map[string]map[string]aclEntry + pattern map[string][]aclEntry + mu sync.RWMutex // Protect exact and pattern } // aclEntry mirrors one user_access row in the in-memory cache. // -// topic is the raw stored value: it may contain \_ escapes (for literal underscores) -// and % wildcards (translated from user-supplied *). For exact-match entries (no %) -// matcher is nil and the entry is keyed by topic in aclCache.exact. For wildcard -// entries (with %) matcher is the pre-compiled regex equivalent of the LIKE pattern. +// length is the length of the original stored value (topic for exact rows, +// SQL LIKE pattern for wildcard rows). It is only used by better() to +// implement the "longer pattern beats shorter" tie-break from the original +// SQL ORDER BY. The string itself is intentionally not stored on the entry: +// the exact map already keys on it, and surfacing it would invite misuse +// (wildcard "topics" are actually SQL patterns like "up%"). +// +// pattern is the pre-compiled regex equivalent of the stored LIKE pattern. +// For exact-match entries (no % in the stored value) pattern is nil and the +// entry is reachable only through aclCache.exact[username][topic]. type aclEntry struct { - topic string + length int // len() of the original stored topic/pattern + pattern *regexp.Regexp // nil for exact entries read bool write bool - matcher *regexp.Regexp // Nil for exact entries } func newAccessCache() *aclCache { return &aclCache{ - exact: make(map[string]map[string]aclEntry), - wildcard: make(map[string][]aclEntry), + exact: make(map[string]map[string]aclEntry), + pattern: make(map[string][]aclEntry), } } // reload runs the bulk-load query against the primary and swaps in freshly-built -// exact and wildcard maps under the write lock. The primary is used (not +// exact and pattern maps under the write lock. The primary is used (not // ReadOnly) so a reload immediately after an ACL mutation sees the freshly- // written rows without replica lag. func (c *aclCache) reload(d *db.DB, query string) error { @@ -54,34 +60,35 @@ func (c *aclCache) reload(d *db.DB, query string) error { return err } defer rows.Close() - exact := make(map[string]map[string]aclEntry) - wildcards := make(map[string][]aclEntry) + exacts := make(map[string]map[string]aclEntry) + patterns := make(map[string][]aclEntry) for rows.Next() { - var username string + var username, topic string var entry aclEntry - if err := rows.Scan(&username, &entry.topic, &entry.read, &entry.write); err != nil { + if err := rows.Scan(&username, &topic, &entry.read, &entry.write); err != nil { return err } - if strings.Contains(entry.topic, "%") { - re, err := compileLikeToRegex(entry.topic) + entry.length = len(topic) + if strings.Contains(topic, "%") { + re, err := compileLikeToRegex(topic) if err != nil { return err } - entry.matcher = re - wildcards[username] = append(wildcards[username], entry) + entry.pattern = re + patterns[username] = append(patterns[username], entry) } else { - if exact[username] == nil { - exact[username] = make(map[string]aclEntry) + if exacts[username] == nil { + exacts[username] = make(map[string]aclEntry) } - exact[username][entry.topic] = entry + exacts[username][topic] = entry } } if err := rows.Err(); err != nil { return err } c.mu.Lock() - c.exact = exact - c.wildcard = wildcards + c.exact = exacts + c.pattern = patterns c.mu.Unlock() return nil } @@ -92,56 +99,46 @@ func (c *aclCache) reload(d *db.DB, query string) error { // 2. longer pattern beats shorter (more specific wins) // 3. write beats read at equal length (write is "stronger") func (c *aclCache) Lookup(usernameOrEveryone, topic string) (read, write, found bool) { - if c == nil { - return false, false, false - } + escapedTopic := escapeUnderscore(topic) c.mu.RLock() defer c.mu.RUnlock() - // Pre-compute the escaped form once: exact-match keys in the cache are - // stored as toSQLWildcard would emit them (literal _ -> \_), so the - // incoming topic must be escaped the same way before map lookup. - escaped := escapeUnderscore(topic) - - // Specific user takes priority over Everyone. Skip the first lookup when - // the request is already anonymous to avoid scanning the same map twice. if usernameOrEveryone != Everyone { - if e, ok := c.pickBestLocked(usernameOrEveryone, topic, escaped); ok { - return e.read, e.write, true + if entry, ok := c.pickBestNoLock(usernameOrEveryone, topic, escapedTopic); ok { + return entry.read, entry.write, true } } - if e, ok := c.pickBestLocked(Everyone, topic, escaped); ok { - return e.read, e.write, true + if entry, ok := c.pickBestNoLock(Everyone, topic, escapedTopic); ok { + return entry.read, entry.write, true } return false, false, false } -// pickBestLocked returns the highest-priority entry for a single user, combining +// pickBestNoLock returns the highest-priority entry for a single user, combining // the exact-match O(1) probe with a linear scan over the (usually empty or tiny) -// wildcard list. Caller must hold c.mu (RLock is sufficient). -func (c *aclCache) pickBestLocked(username, topic, escaped string) (aclEntry, bool) { +// pattern list. Caller must hold c.mu (RLock is sufficient). +func (c *aclCache) pickBestNoLock(username, topic, escapedTopic string) (*aclEntry, bool) { var best aclEntry var found bool - if m, ok := c.exact[username]; ok { - if e, ok := m[escaped]; ok { - best, found = e, true + if m, exists := c.exact[username]; exists { + if entry, exists := m[escapedTopic]; exists { + best, found = entry, true } } - for _, w := range c.wildcard[username] { - if !w.matcher.MatchString(topic) { + for _, pattern := range c.pattern[username] { + if !pattern.pattern.MatchString(topic) { continue - } - if !found || better(w, best) { - best, found = w, true + } else if !found || better(pattern, best) { + best, found = pattern, true } } - return best, found + return &best, found } // better implements the (length DESC, write DESC) tie-break used by the original // query's ORDER BY for entries owned by the same user. func better(a, b aclEntry) bool { - if len(a.topic) != len(b.topic) { - return len(a.topic) > len(b.topic) + if a.length != b.length { + return a.length > b.length } if a.write != b.write { return a.write diff --git a/user/access_cache_test.go b/user/access_cache_test.go index 6c617a0c..54d8d071 100644 --- a/user/access_cache_test.go +++ b/user/access_cache_test.go @@ -175,10 +175,10 @@ func TestACLCache_WriteBeatsReadAtEqualLength(t *testing.T) { c := newAccessCache() c.mu.Lock() c.exact = map[string]map[string]aclEntry{} - c.wildcard = map[string][]aclEntry{ + c.pattern = map[string][]aclEntry{ Everyone: { - {topic: "ab%", read: true, write: false, matcher: mustCompileLikeToRegex(t, "ab%")}, - {topic: "ab%", read: false, write: true, matcher: mustCompileLikeToRegex(t, "ab%")}, + {length: len("ab%"), read: true, write: false, pattern: mustCompileLikeToRegex(t, "ab%")}, + {length: len("ab%"), read: false, write: true, pattern: mustCompileLikeToRegex(t, "ab%")}, }, } c.mu.Unlock() @@ -232,9 +232,9 @@ func loadCache(t *testing.T, c *aclCache, rows []rawACLRow) { exact := make(map[string]map[string]aclEntry) wildcards := make(map[string][]aclEntry) for _, r := range rows { - e := aclEntry{topic: r.topic, read: r.read, write: r.write} + e := aclEntry{length: len(r.topic), read: r.read, write: r.write} if containsPercent(r.topic) { - e.matcher = mustCompileLikeToRegex(t, r.topic) + e.pattern = mustCompileLikeToRegex(t, r.topic) wildcards[r.user] = append(wildcards[r.user], e) } else { if exact[r.user] == nil { @@ -245,7 +245,7 @@ func loadCache(t *testing.T, c *aclCache, rows []rawACLRow) { } c.mu.Lock() c.exact = exact - c.wildcard = wildcards + c.pattern = wildcards c.mu.Unlock() } From d987796243e732232658ec89ed63344d595084db Mon Sep 17 00:00:00 2001 From: binwiederhier Date: Sun, 31 May 2026 11:08:30 -0400 Subject: [PATCH 04/18] Rename --- user/access_cache.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/user/access_cache.go b/user/access_cache.go index adcf6e02..5c5359fd 100644 --- a/user/access_cache.go +++ b/user/access_cache.go @@ -119,8 +119,8 @@ func (c *aclCache) Lookup(usernameOrEveryone, topic string) (read, write, found func (c *aclCache) pickBestNoLock(username, topic, escapedTopic string) (*aclEntry, bool) { var best aclEntry var found bool - if m, exists := c.exact[username]; exists { - if entry, exists := m[escapedTopic]; exists { + if exact, exists := c.exact[username]; exists { + if entry, exists := exact[escapedTopic]; exists { best, found = entry, true } } From 4b87a2732633e3a6e7285e0378b4bc1719eb2c06 Mon Sep 17 00:00:00 2001 From: binwiederhier Date: Sun, 31 May 2026 11:22:30 -0400 Subject: [PATCH 05/18] Docblock, more tests --- user/access_cache.go | 17 +++++++--- user/access_cache_test.go | 70 ++++++++++++++++++++++++++++++++++----- 2 files changed, 74 insertions(+), 13 deletions(-) diff --git a/user/access_cache.go b/user/access_cache.go index 5c5359fd..a49327d2 100644 --- a/user/access_cache.go +++ b/user/access_cache.go @@ -113,9 +113,17 @@ func (c *aclCache) Lookup(usernameOrEveryone, topic string) (read, write, found return false, false, false } -// pickBestNoLock returns the highest-priority entry for a single user, combining -// the exact-match O(1) probe with a linear scan over the (usually empty or tiny) -// pattern list. Caller must hold c.mu (RLock is sufficient). +// pickBestNoLock returns the highest-priority entry for a single user. When +// more than one of that user's rules matches the requested topic, the winner +// is chosen by: +// +// 1. longer stored pattern beats shorter (a more specific rule wins over a +// more general one) +// 2. at equal length, write beats read (a stronger permission wins the tie) +// +// Exact and wildcard rules are ranked together under the same criteria, so +// an exact "foo" (length 3) beats a wildcard "f%" (length 2), but a wildcard +// "foo%" (length 4) beats an exact "foo" (length 3). func (c *aclCache) pickBestNoLock(username, topic, escapedTopic string) (*aclEntry, bool) { var best aclEntry var found bool @@ -139,8 +147,7 @@ func (c *aclCache) pickBestNoLock(username, topic, escapedTopic string) (*aclEnt func better(a, b aclEntry) bool { if a.length != b.length { return a.length > b.length - } - if a.write != b.write { + } else if a.write != b.write { return a.write } return false diff --git a/user/access_cache_test.go b/user/access_cache_test.go index 54d8d071..90839872 100644 --- a/user/access_cache_test.go +++ b/user/access_cache_test.go @@ -64,14 +64,6 @@ func TestCompileLikeToRegex_RegexMetaCharsInTopic(t *testing.T) { require.False(t, r.MatchString("foo.bar")) // would match if '-' leaked into a character class } -func TestACLCache_LookupOnNilReceiverSafe(t *testing.T) { - var c *aclCache - read, write, found := c.Lookup("phil", "mytopic") - require.False(t, found) - require.False(t, read) - require.False(t, write) -} - func TestACLCache_LookupBeforeReload(t *testing.T) { // A freshly-constructed cache has empty exact and wildcards maps. The // cache treats this as "no rule found", which the caller resolves via @@ -142,6 +134,36 @@ func TestACLCache_SpecificUserBeatsEveryone(t *testing.T) { require.False(t, write) } +func TestACLCache_SpecificUserBeatsEveryoneEvenWhenShorter(t *testing.T) { + // The SQL's "user_name DESC" sort key takes precedence over LENGTH(topic). + // Concretely: a specific user with a shorter matching rule still wins over + // Everyone with a longer matching rule. + c := newAccessCache() + loadCache(t, c, []rawACLRow{ + {user: Everyone, topic: "foo", read: true, write: true}, // exact, length 3 + {user: "phil", topic: "f%", read: false, write: false}, // wildcard, length 2, deny-all + }) + read, write, found := c.Lookup("phil", "foo") + require.True(t, found) + require.False(t, read) + require.False(t, write) +} + +func TestACLCache_SpecificUserBeatsEveryoneRegardlessOfWrite(t *testing.T) { + // Same-length rules but conflicting permissions across user boundary: the + // specific user always wins, even if its permission set is weaker (or + // stronger, in either direction). + c := newAccessCache() + loadCache(t, c, []rawACLRow{ + {user: Everyone, topic: "mytopic", read: true, write: true}, // wide-open + {user: "phil", topic: "mytopic", read: true, write: false}, // read-only for phil + }) + read, write, found := c.Lookup("phil", "mytopic") + require.True(t, found) + require.True(t, read) + require.False(t, write) +} + func TestACLCache_AnonymousReadsEveryone(t *testing.T) { c := newAccessCache() loadCache(t, c, []rawACLRow{ @@ -167,6 +189,38 @@ func TestACLCache_LongerPatternWinsForSameUser(t *testing.T) { require.True(t, write) } +func TestACLCache_ExactBeatsShorterWildcardSameUser(t *testing.T) { + // Same user, two matching rules: exact "foo" (length 3) and wildcard "f%" + // (length 2). The longer one wins, which is the exact rule -- mirroring + // the SQL's "LENGTH(topic) DESC" tie-break. Crucially, the cache must seed + // "best" from the exact map probe before walking wildcards, otherwise a + // shorter wildcard could overwrite a longer exact. + c := newAccessCache() + loadCache(t, c, []rawACLRow{ + {user: "phil", topic: "foo", read: true, write: true}, // exact, length 3 + {user: "phil", topic: "f%", read: false, write: false}, // wildcard, length 2, deny-all + }) + read, write, found := c.Lookup("phil", "foo") + require.True(t, found) + require.True(t, read) + require.True(t, write) +} + +func TestACLCache_LongerWildcardBeatsExactSameUser(t *testing.T) { + // Same user, two matching rules: exact "foo" (length 3) and wildcard "foo%" + // (length 4). The wildcard wins on length DESC. Exercises the "swap best + // to wildcard when better() returns true" path. + c := newAccessCache() + loadCache(t, c, []rawACLRow{ + {user: "phil", topic: "foo", read: false, write: false}, // exact, length 3, deny-all + {user: "phil", topic: "foo%", read: true, write: true}, // wildcard, length 4 + }) + read, write, found := c.Lookup("phil", "foo") + require.True(t, found) + require.True(t, read) + require.True(t, write) +} + func TestACLCache_WriteBeatsReadAtEqualLength(t *testing.T) { // Two wildcard rules of identical length for the same user. The write rule // should win the tie-break. The two-rows-with-same-topic shape is From 03d405ed807f37e3cac836cb62b5ff394ac051d3 Mon Sep 17 00:00:00 2001 From: binwiederhier Date: Sun, 31 May 2026 11:42:54 -0400 Subject: [PATCH 06/18] Rename acl cahce --- user/access_cache.go | 16 ++++++++-------- user/access_cache_test.go | 18 +++++++++--------- user/manager.go | 2 +- 3 files changed, 18 insertions(+), 18 deletions(-) diff --git a/user/access_cache.go b/user/access_cache.go index a49327d2..cfd983b3 100644 --- a/user/access_cache.go +++ b/user/access_cache.go @@ -8,7 +8,7 @@ import ( "heckel.io/ntfy/v2/db" ) -// aclCache is an in-memory index over the entire user_access table. +// accessCache is an in-memory index over the entire user_access table. // // exact[username][escapedTopic] returns the matching entry in O(1) for the common // case where the requested topic appears verbatim in some rule. The key is the @@ -18,7 +18,7 @@ import ( // pattern[username] is the linear-scan list of %-bearing rules for that user. // Walked per request; trivially small in practice. Wildcards are NOT u_everyone- // only -- any user can create them. -type aclCache struct { +type accessCache struct { exact map[string]map[string]aclEntry pattern map[string][]aclEntry mu sync.RWMutex // Protect exact and pattern @@ -35,7 +35,7 @@ type aclCache struct { // // pattern is the pre-compiled regex equivalent of the stored LIKE pattern. // For exact-match entries (no % in the stored value) pattern is nil and the -// entry is reachable only through aclCache.exact[username][topic]. +// entry is reachable only through accessCache.exact[username][topic]. type aclEntry struct { length int // len() of the original stored topic/pattern pattern *regexp.Regexp // nil for exact entries @@ -43,8 +43,8 @@ type aclEntry struct { write bool } -func newAccessCache() *aclCache { - return &aclCache{ +func newAccessCache() *accessCache { + return &accessCache{ exact: make(map[string]map[string]aclEntry), pattern: make(map[string][]aclEntry), } @@ -54,7 +54,7 @@ func newAccessCache() *aclCache { // exact and pattern maps under the write lock. The primary is used (not // ReadOnly) so a reload immediately after an ACL mutation sees the freshly- // written rows without replica lag. -func (c *aclCache) reload(d *db.DB, query string) error { +func (c *accessCache) reload(d *db.DB, query string) error { rows, err := d.Query(query) if err != nil { return err @@ -98,7 +98,7 @@ func (c *aclCache) reload(d *db.DB, query string) error { // 1. specific user beats Everyone // 2. longer pattern beats shorter (more specific wins) // 3. write beats read at equal length (write is "stronger") -func (c *aclCache) Lookup(usernameOrEveryone, topic string) (read, write, found bool) { +func (c *accessCache) Lookup(usernameOrEveryone, topic string) (read, write, found bool) { escapedTopic := escapeUnderscore(topic) c.mu.RLock() defer c.mu.RUnlock() @@ -124,7 +124,7 @@ func (c *aclCache) Lookup(usernameOrEveryone, topic string) (read, write, found // Exact and wildcard rules are ranked together under the same criteria, so // an exact "foo" (length 3) beats a wildcard "f%" (length 2), but a wildcard // "foo%" (length 4) beats an exact "foo" (length 3). -func (c *aclCache) pickBestNoLock(username, topic, escapedTopic string) (*aclEntry, bool) { +func (c *accessCache) pickBestNoLock(username, topic, escapedTopic string) (*aclEntry, bool) { var best aclEntry var found bool if exact, exists := c.exact[username]; exists { diff --git a/user/access_cache_test.go b/user/access_cache_test.go index 90839872..65c109cb 100644 --- a/user/access_cache_test.go +++ b/user/access_cache_test.go @@ -140,8 +140,8 @@ func TestACLCache_SpecificUserBeatsEveryoneEvenWhenShorter(t *testing.T) { // Everyone with a longer matching rule. c := newAccessCache() loadCache(t, c, []rawACLRow{ - {user: Everyone, topic: "foo", read: true, write: true}, // exact, length 3 - {user: "phil", topic: "f%", read: false, write: false}, // wildcard, length 2, deny-all + {user: Everyone, topic: "foo", read: true, write: true}, // exact, length 3 + {user: "phil", topic: "f%", read: false, write: false}, // wildcard, length 2, deny-all }) read, write, found := c.Lookup("phil", "foo") require.True(t, found) @@ -155,8 +155,8 @@ func TestACLCache_SpecificUserBeatsEveryoneRegardlessOfWrite(t *testing.T) { // stronger, in either direction). c := newAccessCache() loadCache(t, c, []rawACLRow{ - {user: Everyone, topic: "mytopic", read: true, write: true}, // wide-open - {user: "phil", topic: "mytopic", read: true, write: false}, // read-only for phil + {user: Everyone, topic: "mytopic", read: true, write: true}, // wide-open + {user: "phil", topic: "mytopic", read: true, write: false}, // read-only for phil }) read, write, found := c.Lookup("phil", "mytopic") require.True(t, found) @@ -197,8 +197,8 @@ func TestACLCache_ExactBeatsShorterWildcardSameUser(t *testing.T) { // shorter wildcard could overwrite a longer exact. c := newAccessCache() loadCache(t, c, []rawACLRow{ - {user: "phil", topic: "foo", read: true, write: true}, // exact, length 3 - {user: "phil", topic: "f%", read: false, write: false}, // wildcard, length 2, deny-all + {user: "phil", topic: "foo", read: true, write: true}, // exact, length 3 + {user: "phil", topic: "f%", read: false, write: false}, // wildcard, length 2, deny-all }) read, write, found := c.Lookup("phil", "foo") require.True(t, found) @@ -212,8 +212,8 @@ func TestACLCache_LongerWildcardBeatsExactSameUser(t *testing.T) { // to wildcard when better() returns true" path. c := newAccessCache() loadCache(t, c, []rawACLRow{ - {user: "phil", topic: "foo", read: false, write: false}, // exact, length 3, deny-all - {user: "phil", topic: "foo%", read: true, write: true}, // wildcard, length 4 + {user: "phil", topic: "foo", read: false, write: false}, // exact, length 3, deny-all + {user: "phil", topic: "foo%", read: true, write: true}, // wildcard, length 4 }) read, write, found := c.Lookup("phil", "foo") require.True(t, found) @@ -281,7 +281,7 @@ type rawACLRow struct { // loadCache writes the given rows into the cache under its write lock, // preserving the same exact/wildcard partitioning that reload would produce. -func loadCache(t *testing.T, c *aclCache, rows []rawACLRow) { +func loadCache(t *testing.T, c *accessCache, rows []rawACLRow) { t.Helper() exact := make(map[string]map[string]aclEntry) wildcards := make(map[string][]aclEntry) diff --git a/user/manager.go b/user/manager.go index 52217e29..f9995e84 100644 --- a/user/manager.go +++ b/user/manager.go @@ -58,7 +58,7 @@ type Manager struct { queries queries statsQueue map[string]*Stats // "Queue" to asynchronously write user stats to the database (UserID -> Stats) tokenQueue map[string]*TokenUpdate // "Queue" to asynchronously write token access stats to the database (Token ID -> TokenUpdate) - accessCache *aclCache // In-memory snapshot of user_access; rebuilt after every ACL mutation + accessCache *accessCache // In-memory snapshot of user_access; rebuilt after every ACL mutation quit chan struct{} // Closed by Close() to signal background goroutines to stop mu sync.Mutex } From 301be79f0a752b24550621fecbc367a0bf8f2749 Mon Sep 17 00:00:00 2001 From: binwiederhier Date: Sun, 31 May 2026 14:11:55 -0400 Subject: [PATCH 07/18] Per-user reload --- user/access_cache.go | 87 +++++++++++++++++++++++++++++---- user/manager.go | 102 ++++++++++++++++++++++++--------------- user/manager_postgres.go | 7 +++ user/manager_sqlite.go | 7 +++ user/types.go | 1 + 5 files changed, 156 insertions(+), 48 deletions(-) diff --git a/user/access_cache.go b/user/access_cache.go index cfd983b3..a1518058 100644 --- a/user/access_cache.go +++ b/user/access_cache.go @@ -64,17 +64,15 @@ func (c *accessCache) reload(d *db.DB, query string) error { patterns := make(map[string][]aclEntry) for rows.Next() { var username, topic string - var entry aclEntry - if err := rows.Scan(&username, &topic, &entry.read, &entry.write); err != nil { + var read, write bool + if err := rows.Scan(&username, &topic, &read, &write); err != nil { return err } - entry.length = len(topic) - if strings.Contains(topic, "%") { - re, err := compileLikeToRegex(topic) - if err != nil { - return err - } - entry.pattern = re + entry, isWildcard, err := newACLEntry(topic, read, write) + if err != nil { + return err + } + if isWildcard { patterns[username] = append(patterns[username], entry) } else { if exacts[username] == nil { @@ -93,6 +91,77 @@ func (c *accessCache) reload(d *db.DB, query string) error { return nil } +// reloadUser refreshes just one user's rules from the database and swaps them +// into the cache. Used as a cheaper, owner-cascade-safe alternative to the +// full bulk reload after a mutation that affects a known set of users. The +// caller is expected to invoke reloadUser for every user whose row set may +// have changed -- typically the targeted user plus Everyone, since most +// reservation flows touch both. +// +// The query must return (topic, read, write) for every row whose user_id +// matches the given username. An empty result set is treated as "this user +// has no rules": the user's entries are removed from both maps so the inner +// maps don't grow unbounded under churn. +func (c *accessCache) reloadUser(d *db.DB, query, username string) error { + rows, err := d.Query(query, username) + if err != nil { + return err + } + defer rows.Close() + exact := make(map[string]aclEntry) + var pattern []aclEntry + for rows.Next() { + var topic string + var read, write bool + if err := rows.Scan(&topic, &read, &write); err != nil { + return err + } + entry, isWildcard, err := newACLEntry(topic, read, write) + if err != nil { + return err + } + if isWildcard { + pattern = append(pattern, entry) + } else { + exact[topic] = entry + } + } + if err := rows.Err(); err != nil { + return err + } + c.mu.Lock() + if len(exact) == 0 { + delete(c.exact, username) + } else { + c.exact[username] = exact + } + if len(pattern) == 0 { + delete(c.pattern, username) + } else { + c.pattern[username] = pattern + } + c.mu.Unlock() + return nil +} + +// newACLEntry builds an aclEntry from one user_access row's values. The +// isWildcard return tells the caller which storage slot the entry belongs in: +// the per-user wildcard slice if true, the per-user exact map if false. +// Wildcards have their LIKE pattern pre-compiled into entry.pattern; exact +// entries leave entry.pattern nil. +func newACLEntry(topic string, read, write bool) (entry aclEntry, isWildcard bool, err error) { + entry = aclEntry{length: len(topic), read: read, write: write} + if !strings.Contains(topic, "%") { + return entry, false, nil + } + re, err := compileLikeToRegex(topic) + if err != nil { + return entry, true, err + } + entry.pattern = re + return entry, true, nil +} + // Lookup returns the effective (read, write, found) permission for the given // (username, topic), preserving the priority ordering of the original SQL query: // 1. specific user beats Everyone diff --git a/user/manager.go b/user/manager.go index f9995e84..4626d8e6 100644 --- a/user/manager.go +++ b/user/manager.go @@ -40,9 +40,8 @@ const ( DefaultUserPasswordBcryptCost = 10 // DefaultAccessCacheReloadInterval bounds how stale the in-memory ACL snapshot // can be relative to writes made by *other* processes (e.g. a separate `ntfy - // access` CLI invocation modifying the same database). Mutations performed - // by this Manager refresh the cache synchronously and do not depend on this. - DefaultAccessCacheReloadInterval = 5 * time.Second + // access` CLI invocation modifying the same database) + DefaultAccessCacheReloadInterval = 60 * time.Second ) var ( @@ -87,8 +86,6 @@ func newManager(d *db.DB, queries queries, config *Config) (*Manager, error) { if err := manager.maybeProvisionUsersAccessAndTokens(); err != nil { return nil, err } - // Populate the ACL cache after provisioning so the initial snapshot includes - // any provisioned access rules. Subsequent mutations call reloadAccessCache. if err := manager.reloadAccessCache(); err != nil { return nil, err } @@ -99,18 +96,22 @@ func newManager(d *db.DB, queries queries, config *Config) (*Manager, error) { return manager, nil } -// reloadAccessCache rebuilds the in-memory ACL snapshot from the primary -// database. Called once at startup (after provisioning) and after every method -// that mutates user_access (directly or by cascade from "user" deletion). +// reloadAccessCache rebuilds the in-memory access cache from the primary database func (a *Manager) reloadAccessCache() error { return a.accessCache.reload(a.db, a.queries.selectAllAccessForCache) } -// asyncAccessCacheReloader periodically refreshes the ACL snapshot so that -// writes made by other processes against the same database (most notably the -// `ntfy access` CLI subcommand running while a server holds the cache) become -// visible within the configured interval. This Manager's own writes do not -// depend on the poller -- they refresh the cache synchronously. +// reloadAccessCacheUsers refreshes the cache slices for the given usernames +func (a *Manager) reloadAccessCacheUsers(usernames ...string) error { + for _, username := range usernames { + if err := a.accessCache.reloadUser(a.db, a.queries.selectAccessForCacheByUser, username); err != nil { + return err + } + } + return nil +} + +// asyncAccessCacheReloader periodically refreshes the access cache func (a *Manager) asyncAccessCacheReloader(interval time.Duration) { ticker := time.NewTicker(interval) defer ticker.Stop() @@ -198,13 +199,16 @@ func (a *Manager) RemoveUser(username string) error { if err := a.CanChangeUser(username); err != nil { return err } - if err := db.ExecTx(a.db, func(tx *sql.Tx) error { + err := db.ExecTx(a.db, func(tx *sql.Tx) error { return a.removeUserTx(tx, username) - }); err != nil { + }) + if err != nil { return err } - // user_access rows are cascade-deleted along with the user; refresh the snapshot. - return a.reloadAccessCache() + // user_access rows are cascade-deleted along with the user (both by user_id + // and by owner_user_id). Refresh this user's own slice (now empty) and + // Everyone's slice, since reservations owned by this user landed there too. + return a.reloadAccessCacheUsers(username, Everyone) } // removeUserTx deletes the user with the given username @@ -225,7 +229,7 @@ func (a *Manager) MarkUserRemoved(user *User) error { if !AllowedUsername(user.Name) { return ErrInvalidArgument } - if err := db.ExecTx(a.db, func(tx *sql.Tx) error { + err := db.ExecTx(a.db, func(tx *sql.Tx) error { if err := a.resetUserAccessTx(tx, user.Name); err != nil { return err } @@ -237,11 +241,14 @@ func (a *Manager) MarkUserRemoved(user *User) error { return err } return nil - }); err != nil { + }) + if err != nil { return err } - // resetUserAccessTx wiped this user's user_access rows; refresh the snapshot. - return a.reloadAccessCache() + // resetUserAccessTx deleted this user's rows AND any row owned by this user + // (typically the matching Everyone rows from their reservations). Refresh + // both slices to mirror the DB exactly. + return a.reloadAccessCacheUsers(user.Name, Everyone) } // RemoveDeletedUsers deletes all users that have been marked deleted @@ -281,15 +288,17 @@ func (a *Manager) ChangeRole(username string, role Role) error { if err := a.CanChangeUser(username); err != nil { return err } - if err := db.ExecTx(a.db, func(tx *sql.Tx) error { + err := db.ExecTx(a.db, func(tx *sql.Tx) error { return a.changeRoleTx(tx, username, role) - }); err != nil { + }) + if err != nil { return err } - // Promotion to admin clears user_access rows for the user; refresh the snapshot. - // Other role changes are no-ops for the cache but reloading is cheap and keeps - // the code path uniform. - return a.reloadAccessCache() + // Promotion to admin clears user_access rows for this user AND rows owned + // by this user (Everyone rows from their reservations). Other role changes + // are no-ops for the cache but reloading the two affected slices is cheap + // and keeps the code path uniform. + return a.reloadAccessCacheUsers(username, Everyone) } // changeRoleTx changes a user's role @@ -650,12 +659,14 @@ func (a *Manager) resolvePerms(base, perm Permission) error { // read/write access to a topic. The parameter topicPattern may include wildcards (*). The ACL entry // owner may either be a user (username), or the system (empty). func (a *Manager) AllowAccess(username string, topicPattern string, permission Permission) error { - if err := db.ExecTx(a.db, func(tx *sql.Tx) error { + err := db.ExecTx(a.db, func(tx *sql.Tx) error { return a.allowAccessTx(tx, username, topicPattern, permission, false) - }); err != nil { + }) + if err != nil { return err } - return a.reloadAccessCache() + // Only this user's row set changed; refresh their slice only. + return a.reloadAccessCacheUsers(username) } func (a *Manager) allowAccessTx(tx *sql.Tx, username string, topicPattern string, permission Permission, provisioned bool) error { @@ -671,12 +682,20 @@ func (a *Manager) allowAccessTx(tx *sql.Tx, username string, topicPattern string // ResetAccess removes an access control list entry for a specific username/topic, or (if topic is // empty) for an entire user. The parameter topicPattern may include wildcards (*). func (a *Manager) ResetAccess(username string, topicPattern string) error { - if err := db.ExecTx(a.db, func(tx *sql.Tx) error { + err := db.ExecTx(a.db, func(tx *sql.Tx) error { return a.resetAccessTx(tx, username, topicPattern) - }); err != nil { + }) + if err != nil { return err } - return a.reloadAccessCache() + // "Delete all access" affects every user; do the bulk reload. + // Otherwise refresh the named user plus Everyone, since resetUserAccessTx + // and deleteTopicAccess both touch rows owned by the user (typically + // Everyone rows from their reservations). + if username == "" { + return a.reloadAccessCache() + } + return a.reloadAccessCacheUsers(username, Everyone) } func (a *Manager) resetAccessTx(tx *sql.Tx, username string, topicPattern string) error { @@ -790,7 +809,7 @@ func (a *Manager) AddReservation(username string, topic string, everyone Permiss if !AllowedUsername(username) || username == Everyone || !AllowedTopic(topic) { return ErrInvalidArgument } - if err := db.ExecTx(a.db, func(tx *sql.Tx) error { + err := db.ExecTx(a.db, func(tx *sql.Tx) error { if limit > 0 { hasReservation, err := a.hasReservationTx(tx, username, topic) if err != nil { @@ -813,10 +832,12 @@ func (a *Manager) AddReservation(username string, topic string, everyone Permiss return err } return nil - }); err != nil { + }) + if err != nil { return err } - return a.reloadAccessCache() + // Both user's and Everyone's rows changed. + return a.reloadAccessCacheUsers(username, Everyone) } // RemoveReservations deletes the access control entries associated with the given username/topic, @@ -831,17 +852,20 @@ func (a *Manager) RemoveReservations(username string, topics ...string) error { return ErrInvalidArgument } } - if err := db.ExecTx(a.db, func(tx *sql.Tx) error { + err := db.ExecTx(a.db, func(tx *sql.Tx) error { for _, topic := range topics { if err := a.removeReservationAccessTx(tx, username, topic); err != nil { return err } } return nil - }); err != nil { + }) + if err != nil { return err } - return a.reloadAccessCache() + // Mirror the DB: rows for this user and any Everyone rows owned by this + // user are gone. Refresh both slices. + return a.reloadAccessCacheUsers(username, Everyone) } // Reservations returns all user-owned topics, and the associated everyone-access diff --git a/user/manager_postgres.go b/user/manager_postgres.go index d3bb6f5f..8dec4c7d 100644 --- a/user/manager_postgres.go +++ b/user/manager_postgres.go @@ -75,6 +75,12 @@ const ( FROM user_access a JOIN "user" u ON u.id = a.user_id ` + postgresSelectAccessForCacheByUserQuery = ` + SELECT a.topic, a.read, a.write + FROM user_access a + JOIN "user" u ON u.id = a.user_id + WHERE u.user_name = $1 + ` postgresSelectUserAllAccessQuery = ` SELECT user_id, topic, read, write, provisioned FROM user_access @@ -243,6 +249,7 @@ var postgresQueries = queries{ deleteUsersMarked: postgresDeleteUsersMarkedQuery, deleteUsersProvisioned: postgresDeleteUsersProvisionedQuery, selectAllAccessForCache: postgresSelectAllAccessForCacheQuery, + selectAccessForCacheByUser: postgresSelectAccessForCacheByUserQuery, selectUserAllAccess: postgresSelectUserAllAccessQuery, selectUserAccess: postgresSelectUserAccessQuery, selectUserReservations: postgresSelectUserReservationsQuery, diff --git a/user/manager_sqlite.go b/user/manager_sqlite.go index 4d498ce2..26e22b36 100644 --- a/user/manager_sqlite.go +++ b/user/manager_sqlite.go @@ -81,6 +81,12 @@ const ( FROM user_access a JOIN user u ON u.id = a.user_id ` + sqliteSelectAccessForCacheByUserQuery = ` + SELECT a.topic, a.read, a.write + FROM user_access a + JOIN user u ON u.id = a.user_id + WHERE u.user = ? + ` sqliteSelectUserAllAccessQuery = ` SELECT user_id, topic, read, write, provisioned FROM user_access @@ -241,6 +247,7 @@ var sqliteQueries = queries{ deleteUsersMarked: sqliteDeleteUsersMarkedQuery, deleteUsersProvisioned: sqliteDeleteUsersProvisionedQuery, selectAllAccessForCache: sqliteSelectAllAccessForCacheQuery, + selectAccessForCacheByUser: sqliteSelectAccessForCacheByUserQuery, selectUserAllAccess: sqliteSelectUserAllAccessQuery, selectUserAccess: sqliteSelectUserAccessQuery, selectUserReservations: sqliteSelectUserReservationsQuery, diff --git a/user/types.go b/user/types.go index 6d2650ed..3af01f5d 100644 --- a/user/types.go +++ b/user/types.go @@ -314,6 +314,7 @@ type queries struct { // Access queries selectAllAccessForCache string // Bulk load: (user_name, topic, read, write) for the in-memory ACL cache + selectAccessForCacheByUser string // Per-user load: (topic, read, write) for one username; used to refresh just one user's slice of the cache after mutation selectUserAllAccess string selectUserAccess string selectUserReservations string From 204723f3c09357e4eb8463420d4178d2b6006403 Mon Sep 17 00:00:00 2001 From: binwiederhier Date: Sun, 31 May 2026 14:28:04 -0400 Subject: [PATCH 08/18] Make opt-in flag --- cmd/access_test.go | 7 -- cmd/user.go | 10 +-- cmd/user_test.go | 7 +- server/config.go | 2 + server/server.go | 1 + user/manager.go | 72 +++++++++++++++----- user/manager_postgres.go | 8 +++ user/manager_sqlite.go | 8 +++ user/manager_test.go | 142 +++++++++++++++++++++++++++++++++++++++ user/types.go | 22 +++--- 10 files changed, 236 insertions(+), 43 deletions(-) diff --git a/cmd/access_test.go b/cmd/access_test.go index f280a9e9..8810b6b3 100644 --- a/cmd/access_test.go +++ b/cmd/access_test.go @@ -7,7 +7,6 @@ import ( "heckel.io/ntfy/v2/server" "heckel.io/ntfy/v2/test" "testing" - "time" ) func TestCLI_Access_Show(t *testing.T) { @@ -44,12 +43,6 @@ user * (role: anonymous, tier: none) ` require.Equal(t, expected, stdout.String()) - // The CLI commands above ran against a separate Manager instance (their own - // process-equivalent), so the server's ACL cache hasn't seen the new grants - // yet. Wait for the server's background reloader (interval set in - // newTestServerWithAuth) to pick them up. - time.Sleep(150 * time.Millisecond) - // See if access permissions match app, _, _, _ = newTestApp() require.Error(t, app.Run([]string{ diff --git a/cmd/user.go b/cmd/user.go index 9bb0c5b0..8eca5ce5 100644 --- a/cmd/user.go +++ b/cmd/user.go @@ -378,11 +378,11 @@ func createUserManager(c *cli.Context) (*user.Manager, error) { ProvisionEnabled: false, // Hack: Do not re-provision users on manager initialization BcryptCost: user.DefaultUserPasswordBcryptCost, QueueWriterInterval: user.DefaultUserStatsQueueWriterInterval, - // CLI Managers are short-lived; the background ACL cache poller would only - // spam "database is closed" warnings after the subcommand returns. Mutations - // still refresh the local cache synchronously; the running server (if any) - // picks them up via its own poller. - AccessCacheReloadInterval: -1, + // CLI subcommands never serve authorizeTopicAccess and are short-lived, + // so the cache (and its background poller) would be wasted work. Mutations + // hit the DB directly; the running server, if any, picks them up via its + // own poller when the cache is enabled there. + AccessCacheEnabled: false, } if databaseURL != "" { host, dbErr := pg.Open(databaseURL) diff --git a/cmd/user_test.go b/cmd/user_test.go index 5b5bff73..91373694 100644 --- a/cmd/user_test.go +++ b/cmd/user_test.go @@ -9,7 +9,6 @@ import ( "os" "path/filepath" "testing" - "time" ) func TestCLI_User_Add(t *testing.T) { @@ -129,10 +128,8 @@ func newTestServerWithAuth(t *testing.T) (s *server.Server, conf *server.Config, conf.File = configFile conf.AuthFile = filepath.Join(t.TempDir(), "user.db") conf.AuthDefault = user.PermissionDenyAll - // Tight interval so cross-process writes from the `ntfy access`/`ntfy user` - // CLI commands (which run via a separate Manager) propagate to the server's - // ACL cache within tens of ms instead of the default 5s. - conf.AuthAccessCacheReloadInterval = 25 * time.Millisecond + // Cache is off by default (matches self-hoster setup), so the server reads + // authorizations directly from the DB and sees CLI mutations immediately. s, port = test.StartServerWithConfig(t, conf) return } diff --git a/server/config.go b/server/config.go index f6f13eca..a1ba4d40 100644 --- a/server/config.go +++ b/server/config.go @@ -116,6 +116,7 @@ type Config struct { AuthTokens map[string][]*user.Token AuthBcryptCost int AuthStatsQueueWriterInterval time.Duration + AuthAccessCacheEnabled bool AuthAccessCacheReloadInterval time.Duration AttachmentCacheDir string AttachmentTotalSizeLimit int64 @@ -224,6 +225,7 @@ func NewConfig() *Config { AuthDefault: user.PermissionReadWrite, AuthBcryptCost: user.DefaultUserPasswordBcryptCost, AuthStatsQueueWriterInterval: user.DefaultUserStatsQueueWriterInterval, + AuthAccessCacheEnabled: user.DefaultAccessCacheEnabled, // Opt-in (e.g. ntfy.sh) via server.yml AuthAccessCacheReloadInterval: user.DefaultAccessCacheReloadInterval, AttachmentCacheDir: "", AttachmentTotalSizeLimit: DefaultAttachmentTotalSizeLimit, diff --git a/server/server.go b/server/server.go index 7bcbbb09..c380bd26 100644 --- a/server/server.go +++ b/server/server.go @@ -257,6 +257,7 @@ func New(conf *Config) (*Server, error) { Tokens: conf.AuthTokens, BcryptCost: conf.AuthBcryptCost, QueueWriterInterval: conf.AuthStatsQueueWriterInterval, + AccessCacheEnabled: conf.AuthAccessCacheEnabled, AccessCacheReloadInterval: conf.AuthAccessCacheReloadInterval, } if pool != nil { diff --git a/user/manager.go b/user/manager.go index 4626d8e6..292d59b2 100644 --- a/user/manager.go +++ b/user/manager.go @@ -38,9 +38,14 @@ const ( const ( DefaultUserStatsQueueWriterInterval = 33 * time.Second DefaultUserPasswordBcryptCost = 10 + // DefaultAccessCacheEnabled is the default for Config.AccessCacheEnabled. + // Off by default so self-hosters keep the direct-DB authorizeTopicAccess + // path; ntfy.sh opts in via server config. + DefaultAccessCacheEnabled = false // DefaultAccessCacheReloadInterval bounds how stale the in-memory ACL snapshot // can be relative to writes made by *other* processes (e.g. a separate `ntfy - // access` CLI invocation modifying the same database) + // access` CLI invocation modifying the same database). Only honored when the + // cache is enabled. DefaultAccessCacheReloadInterval = 60 * time.Second ) @@ -75,34 +80,46 @@ func newManager(d *db.DB, queries queries, config *Config) (*Manager, error) { config.AccessCacheReloadInterval = DefaultAccessCacheReloadInterval } manager := &Manager{ - config: config, - db: d, - statsQueue: make(map[string]*Stats), - tokenQueue: make(map[string]*TokenUpdate), - accessCache: newAccessCache(), - quit: make(chan struct{}), - queries: queries, + config: config, + db: d, + statsQueue: make(map[string]*Stats), + tokenQueue: make(map[string]*TokenUpdate), + quit: make(chan struct{}), + queries: queries, + } + if config.AccessCacheEnabled { + manager.accessCache = newAccessCache() } if err := manager.maybeProvisionUsersAccessAndTokens(); err != nil { return nil, err } + // Populate the cache after provisioning so the initial snapshot includes + // any provisioned access rules. No-op when the cache is disabled. if err := manager.reloadAccessCache(); err != nil { return nil, err } go manager.asyncQueueWriter(manager.config.QueueWriterInterval) - if manager.config.AccessCacheReloadInterval > 0 { + if manager.accessCache != nil && manager.config.AccessCacheReloadInterval > 0 { go manager.asyncAccessCacheReloader(manager.config.AccessCacheReloadInterval) } return manager, nil } -// reloadAccessCache rebuilds the in-memory access cache from the primary database +// reloadAccessCache rebuilds the in-memory access cache from the primary +// database. No-op when the cache is disabled. func (a *Manager) reloadAccessCache() error { + if a.accessCache == nil { + return nil + } return a.accessCache.reload(a.db, a.queries.selectAllAccessForCache) } -// reloadAccessCacheUsers refreshes the cache slices for the given usernames +// reloadAccessCacheUsers refreshes the cache slices for the given usernames. +// No-op when the cache is disabled. func (a *Manager) reloadAccessCacheUsers(usernames ...string) error { + if a.accessCache == nil { + return nil + } for _, username := range usernames { if err := a.accessCache.reloadUser(a.db, a.queries.selectAccessForCacheByUser, username); err != nil { return err @@ -737,15 +754,34 @@ func (a *Manager) AllowReservation(username string, topic string) error { // authorizeTopicAccess returns the read/write permissions for the given username and topic. // The found return value indicates whether an ACL entry was found at all. // -// - The cache may contain two matching entries (one for everyone, and one for the user), but prioritizes the user. -// - Furthermore, the lookup prioritizes more specific permissions (longer!) over more generic ones, e.g. "test*" > "*" -// - It also prioritizes write permissions over read permissions +// Priority: +// - specific user beats Everyone +// - longer pattern beats shorter (a more specific rule beats a more general one, +// e.g. "test*" > "*") +// - write beats read at equal length // -// The lookup is served entirely from the in-memory snapshot maintained by accessCache, -// so this is on the hot path of every authenticatable HTTP request and must stay allocation-free. +// When AccessCacheEnabled is true (config), the lookup is served entirely from +// the in-memory snapshot maintained by accessCache. Otherwise the original SQL +// query is executed against the database on every call. func (a *Manager) authorizeTopicAccess(usernameOrEveryone, topic string) (read, write, found bool, err error) { - read, write, found = a.accessCache.Lookup(usernameOrEveryone, topic) - return read, write, found, nil + if a.accessCache != nil { + read, write, found = a.accessCache.Lookup(usernameOrEveryone, topic) + return read, write, found, nil + } + rows, err := a.db.ReadOnly().Query(a.queries.selectTopicPerms, Everyone, usernameOrEveryone, topic) + if err != nil { + return false, false, false, err + } + defer rows.Close() + if !rows.Next() { + return false, false, false, nil + } + if err := rows.Scan(&read, &write); err != nil { + return false, false, false, err + } else if err := rows.Err(); err != nil { + return false, false, false, err + } + return read, write, true, nil } // AllGrants returns all user-specific access control entries, mapped to their respective user IDs diff --git a/user/manager_postgres.go b/user/manager_postgres.go index 8dec4c7d..8f0079b6 100644 --- a/user/manager_postgres.go +++ b/user/manager_postgres.go @@ -70,6 +70,13 @@ const ( postgresDeleteUsersProvisionedQuery = `DELETE FROM "user" WHERE provisioned = true` // Access queries + postgresSelectTopicPermsQuery = ` + SELECT read, write + FROM user_access a + JOIN "user" u ON u.id = a.user_id + WHERE (u.user_name = $1 OR u.user_name = $2) AND $3 LIKE a.topic ESCAPE '\' + ORDER BY u.user_name DESC, LENGTH(a.topic) DESC, CASE WHEN a.write THEN 1 ELSE 0 END DESC + ` postgresSelectAllAccessForCacheQuery = ` SELECT u.user_name, a.topic, a.read, a.write FROM user_access a @@ -248,6 +255,7 @@ var postgresQueries = queries{ deleteUserTier: postgresDeleteUserTierQuery, deleteUsersMarked: postgresDeleteUsersMarkedQuery, deleteUsersProvisioned: postgresDeleteUsersProvisionedQuery, + selectTopicPerms: postgresSelectTopicPermsQuery, selectAllAccessForCache: postgresSelectAllAccessForCacheQuery, selectAccessForCacheByUser: postgresSelectAccessForCacheByUserQuery, selectUserAllAccess: postgresSelectUserAllAccessQuery, diff --git a/user/manager_sqlite.go b/user/manager_sqlite.go index 26e22b36..a1ca4700 100644 --- a/user/manager_sqlite.go +++ b/user/manager_sqlite.go @@ -76,6 +76,13 @@ const ( sqliteDeleteUsersProvisionedQuery = `DELETE FROM user WHERE provisioned = 1` // Access queries + sqliteSelectTopicPermsQuery = ` + SELECT read, write + FROM user_access a + JOIN user u ON u.id = a.user_id + WHERE (u.user = ? OR u.user = ?) AND ? LIKE a.topic ESCAPE '\' + ORDER BY u.user DESC, LENGTH(a.topic) DESC, a.write DESC + ` sqliteSelectAllAccessForCacheQuery = ` SELECT u.user, a.topic, a.read, a.write FROM user_access a @@ -246,6 +253,7 @@ var sqliteQueries = queries{ deleteUserTier: sqliteDeleteUserTierQuery, deleteUsersMarked: sqliteDeleteUsersMarkedQuery, deleteUsersProvisioned: sqliteDeleteUsersProvisionedQuery, + selectTopicPerms: sqliteSelectTopicPermsQuery, selectAllAccessForCache: sqliteSelectAllAccessForCacheQuery, selectAccessForCacheByUser: sqliteSelectAccessForCacheByUserQuery, selectUserAllAccess: sqliteSelectUserAllAccessQuery, diff --git a/user/manager_test.go b/user/manager_test.go index 3bdb15b2..5d6944f2 100644 --- a/user/manager_test.go +++ b/user/manager_test.go @@ -2169,6 +2169,148 @@ func TestStoreAuthorizeTopicAccessDenyAll(t *testing.T) { }) } +// TestAuthorizeTopicAccess_CacheAndDirectDBAgree wires up two Managers on the +// same backend storage -- one with AccessCacheEnabled=true (in-memory cache +// path) and one with AccessCacheEnabled=false (direct SQL path) -- then runs +// an identical battery of authorizeTopicAccess queries against both and +// asserts byte-identical (read, write, found) responses for every query. +// This protects the in-memory implementation from drifting away from the +// SQL behavior it is meant to mirror. +func TestAuthorizeTopicAccess_CacheAndDirectDBAgree(t *testing.T) { + forEachBackend(t, func(t *testing.T, newManager newManagerFunc) { + // Seed via a Manager with the cache enabled. Writes go to the shared + // backend; both Managers will see them after the writes commit. + writer := newManager(&Config{ + DefaultAccess: PermissionDenyAll, + BcryptCost: bcrypt.MinCost, + AccessCacheEnabled: true, + }) + t.Cleanup(func() { writer.Close() }) + + require.Nil(t, writer.AddUser("phil", "mypass", RoleAdmin, false)) + require.Nil(t, writer.AddUser("ben", "mypass", RoleUser, false)) + require.Nil(t, writer.AddUser("alice", "mypass", RoleUser, false)) + + // A mix that exercises every branch of the priority logic: + // - exact and wildcard rules for the same user + // - exact and wildcard rules under Everyone + // - Everyone rules that are longer than the matching user rule + // - literal underscores (stored as "\_") + // - deny-all permissions + require.Nil(t, writer.AllowAccess("ben", "mytopic", PermissionReadWrite)) + require.Nil(t, writer.AllowAccess("ben", "readme", PermissionRead)) + require.Nil(t, writer.AllowAccess("ben", "writeme", PermissionWrite)) + require.Nil(t, writer.AllowAccess("ben", "ben_topic", PermissionReadWrite)) + require.Nil(t, writer.AllowAccess("ben", "mytopic*", PermissionRead)) + require.Nil(t, writer.AllowAccess("alice", "alice_*", PermissionWrite)) + require.Nil(t, writer.AllowAccess("alice", "secret", PermissionDenyAll)) + require.Nil(t, writer.AllowAccess(Everyone, "announcements", PermissionRead)) + require.Nil(t, writer.AllowAccess(Everyone, "up*", PermissionWrite)) + require.Nil(t, writer.AllowAccess(Everyone, "mytopic", PermissionDenyAll)) + + // Build a reader Manager with the cache OFF, pointing at the same backend. + reader := newManager(&Config{ + DefaultAccess: PermissionDenyAll, + BcryptCost: bcrypt.MinCost, + AccessCacheEnabled: false, + }) + t.Cleanup(func() { reader.Close() }) + + // Probe matrix: every (user, topic) pair that exercises some branch. + cases := []struct { + user, topic string + }{ + // Anonymous reads. + {Everyone, "announcements"}, + {Everyone, "up42"}, + {Everyone, "up"}, + {Everyone, "downstream"}, + {Everyone, "mytopic"}, + {Everyone, "nope"}, + // Specific user, only-user rules. + {"ben", "mytopic"}, + {"ben", "readme"}, + {"ben", "writeme"}, + {"ben", "ben_topic"}, + {"ben", "benXtopic"}, // underscore in rule means "X" must NOT match + // Specific user falls through to Everyone. + {"ben", "announcements"}, + {"ben", "up5"}, + {"alice", "announcements"}, + // Wildcards with literal underscores. + {"alice", "alice_anything"}, + {"alice", "alice_"}, + {"alice", "aliceX"}, // does NOT match alice_* + // Exact-vs-wildcard overlap for the same user (ben has both + // "mytopic" exact and "mytopic*" wildcard). + {"ben", "mytopic"}, // exact wins on length + {"ben", "mytopicX"}, // only wildcard matches + {"ben", "mytopicYZ"}, // only wildcard matches + // Deny-all override. + {"alice", "secret"}, + // No matching rule anywhere. + {"ben", "completely_unmatched"}, + {"alice", "completely_unmatched"}, + {Everyone, "completely_unmatched"}, + } + + // Sanity: the two Managers must agree on every probe. + for _, tc := range cases { + cRead, cWrite, cFound, cErr := writer.authorizeTopicAccess(tc.user, tc.topic) + dRead, dWrite, dFound, dErr := reader.authorizeTopicAccess(tc.user, tc.topic) + require.Nil(t, cErr, "cache path errored for (%s, %s)", tc.user, tc.topic) + require.Nil(t, dErr, "direct-DB path errored for (%s, %s)", tc.user, tc.topic) + require.Equal(t, dFound, cFound, "found mismatch for (%s, %s)", tc.user, tc.topic) + require.Equal(t, dRead, cRead, "read mismatch for (%s, %s)", tc.user, tc.topic) + require.Equal(t, dWrite, cWrite, "write mismatch for (%s, %s)", tc.user, tc.topic) + } + }) +} + +// TestAccessCacheReloadInterval_PicksUpExternalWrite proves that the +// background reloader actually closes the cross-process coherence gap: a +// write made through a *different* Manager on the same backend becomes +// visible to a cache-enabled Manager within roughly one reload interval, +// without that Manager being told about the write. +func TestAccessCacheReloadInterval_PicksUpExternalWrite(t *testing.T) { + const interval = 25 * time.Millisecond + forEachBackend(t, func(t *testing.T, newManager newManagerFunc) { + // reader holds the cache and polls; writer plays the role of an + // out-of-band process (e.g. `ntfy access` CLI) writing to the same + // backend. + reader := newManager(&Config{ + DefaultAccess: PermissionDenyAll, + BcryptCost: bcrypt.MinCost, + AccessCacheEnabled: true, + AccessCacheReloadInterval: interval, + }) + t.Cleanup(func() { reader.Close() }) + + writer := newManager(&Config{ + DefaultAccess: PermissionDenyAll, + BcryptCost: bcrypt.MinCost, + AccessCacheEnabled: false, + }) + t.Cleanup(func() { writer.Close() }) + + require.Nil(t, writer.AddUser("phil", "mypass", RoleUser, false)) + // Sanity: before the write, the reader sees no rule for this topic. + _, _, found, err := reader.authorizeTopicAccess("phil", "via-poller") + require.Nil(t, err) + require.False(t, found) + + // Write through the second Manager. reader's cache is unaware. + require.Nil(t, writer.AllowAccess("phil", "via-poller", PermissionReadWrite)) + + // Wait for the poller to catch up. The interval is 25ms; allow a + // generous multiple to keep this test from flaking on slow CI. + require.Eventually(t, func() bool { + read, write, found, err := reader.authorizeTopicAccess("phil", "via-poller") + return err == nil && found && read && write + }, 2*time.Second, 10*time.Millisecond, "reader's cache never observed the external write") + }) +} + func TestStoreReservations(t *testing.T) { forEachStoreBackend(t, func(t *testing.T, manager *Manager) { require.Nil(t, manager.AddUser("phil", "mypass", RoleUser, false)) diff --git a/user/types.go b/user/types.go index 3af01f5d..65032aa3 100644 --- a/user/types.go +++ b/user/types.go @@ -256,14 +256,19 @@ type Config struct { QueueWriterInterval time.Duration // Interval for the async queue writer to flush stats and token updates to the database BcryptCost int // Cost of generated passwords; lowering makes testing faster - // AccessCacheReloadInterval bounds the staleness of the in-memory ACL cache - // relative to writes from other processes (e.g. `ntfy access` CLI against a - // running server). - // 0 -> use DefaultAccessCacheReloadInterval - // negative -> disable the background poller; cache only refreshes on this - // Manager's own ACL mutations. Use this for short-lived Managers - // (e.g. the CLI subcommands) where polling is wasted work. - // positive -> poll at the given interval + // AccessCacheEnabled gates the in-memory ACL cache. When false (the + // default), authorizeTopicAccess runs the direct SQL query against the + // database on every call -- mutations and authorization are unaffected by + // any cache logic. When true, the Manager keeps an in-memory snapshot of + // user_access and serves authorizeTopicAccess from it; mutations refresh + // the affected slices, and a background poller picks up cross-process + // writes at AccessCacheReloadInterval. + AccessCacheEnabled bool + + // AccessCacheReloadInterval bounds the staleness of the in-memory ACL + // cache relative to writes from other processes (e.g. `ntfy access` CLI + // against a running server). Only honored when AccessCacheEnabled is true. + // Zero falls back to DefaultAccessCacheReloadInterval. AccessCacheReloadInterval time.Duration } @@ -313,6 +318,7 @@ type queries struct { deleteUsersProvisioned string // Access queries + selectTopicPerms string // Direct-DB authorizeTopicAccess query; used when the in-memory cache is disabled selectAllAccessForCache string // Bulk load: (user_name, topic, read, write) for the in-memory ACL cache selectAccessForCacheByUser string // Per-user load: (topic, read, write) for one username; used to refresh just one user's slice of the cache after mutation selectUserAllAccess string From 2f4afbdae5690b16f4c539be820212dc4d138281 Mon Sep 17 00:00:00 2001 From: binwiederhier Date: Sun, 31 May 2026 15:26:28 -0400 Subject: [PATCH 09/18] Manual refinements --- cmd/user.go | 6 +- user/access_cache.go | 213 ++++++++++++++++----------------------- user/manager.go | 57 +++++------ user/manager_postgres.go | 30 ++++-- user/manager_sqlite.go | 23 +++-- user/types.go | 6 +- 6 files changed, 149 insertions(+), 186 deletions(-) diff --git a/cmd/user.go b/cmd/user.go index 8eca5ce5..2e5af5f4 100644 --- a/cmd/user.go +++ b/cmd/user.go @@ -378,11 +378,7 @@ func createUserManager(c *cli.Context) (*user.Manager, error) { ProvisionEnabled: false, // Hack: Do not re-provision users on manager initialization BcryptCost: user.DefaultUserPasswordBcryptCost, QueueWriterInterval: user.DefaultUserStatsQueueWriterInterval, - // CLI subcommands never serve authorizeTopicAccess and are short-lived, - // so the cache (and its background poller) would be wasted work. Mutations - // hit the DB directly; the running server, if any, picks them up via its - // own poller when the cache is enabled there. - AccessCacheEnabled: false, + AccessCacheEnabled: false, // Do not cache for CLI commands } if databaseURL != "" { host, dbErr := pg.Open(databaseURL) diff --git a/user/access_cache.go b/user/access_cache.go index a1518058..1318b01d 100644 --- a/user/access_cache.go +++ b/user/access_cache.go @@ -24,21 +24,14 @@ type accessCache struct { mu sync.RWMutex // Protect exact and pattern } -// aclEntry mirrors one user_access row in the in-memory cache. -// -// length is the length of the original stored value (topic for exact rows, -// SQL LIKE pattern for wildcard rows). It is only used by better() to -// implement the "longer pattern beats shorter" tie-break from the original -// SQL ORDER BY. The string itself is intentionally not stored on the entry: -// the exact map already keys on it, and surfacing it would invite misuse -// (wildcard "topics" are actually SQL patterns like "up%"). -// -// pattern is the pre-compiled regex equivalent of the stored LIKE pattern. -// For exact-match entries (no % in the stored value) pattern is nil and the -// entry is reachable only through accessCache.exact[username][topic]. +// aclEntry mirrors one user_access row. length feeds better()'s "longer +// pattern wins" tie-break; the stored topic/pattern string itself is not kept +// on the entry (the exact map already keys on it; surfacing wildcard "topics" +// like "up%" alongside real ones would invite misuse). pattern is the +// compiled regex form of the LIKE pattern; nil for exact entries. type aclEntry struct { - length int // len() of the original stored topic/pattern - pattern *regexp.Regexp // nil for exact entries + length int + pattern *regexp.Regexp read bool write bool } @@ -50,118 +43,6 @@ func newAccessCache() *accessCache { } } -// reload runs the bulk-load query against the primary and swaps in freshly-built -// exact and pattern maps under the write lock. The primary is used (not -// ReadOnly) so a reload immediately after an ACL mutation sees the freshly- -// written rows without replica lag. -func (c *accessCache) reload(d *db.DB, query string) error { - rows, err := d.Query(query) - if err != nil { - return err - } - defer rows.Close() - exacts := make(map[string]map[string]aclEntry) - patterns := make(map[string][]aclEntry) - for rows.Next() { - var username, topic string - var read, write bool - if err := rows.Scan(&username, &topic, &read, &write); err != nil { - return err - } - entry, isWildcard, err := newACLEntry(topic, read, write) - if err != nil { - return err - } - if isWildcard { - patterns[username] = append(patterns[username], entry) - } else { - if exacts[username] == nil { - exacts[username] = make(map[string]aclEntry) - } - exacts[username][topic] = entry - } - } - if err := rows.Err(); err != nil { - return err - } - c.mu.Lock() - c.exact = exacts - c.pattern = patterns - c.mu.Unlock() - return nil -} - -// reloadUser refreshes just one user's rules from the database and swaps them -// into the cache. Used as a cheaper, owner-cascade-safe alternative to the -// full bulk reload after a mutation that affects a known set of users. The -// caller is expected to invoke reloadUser for every user whose row set may -// have changed -- typically the targeted user plus Everyone, since most -// reservation flows touch both. -// -// The query must return (topic, read, write) for every row whose user_id -// matches the given username. An empty result set is treated as "this user -// has no rules": the user's entries are removed from both maps so the inner -// maps don't grow unbounded under churn. -func (c *accessCache) reloadUser(d *db.DB, query, username string) error { - rows, err := d.Query(query, username) - if err != nil { - return err - } - defer rows.Close() - exact := make(map[string]aclEntry) - var pattern []aclEntry - for rows.Next() { - var topic string - var read, write bool - if err := rows.Scan(&topic, &read, &write); err != nil { - return err - } - entry, isWildcard, err := newACLEntry(topic, read, write) - if err != nil { - return err - } - if isWildcard { - pattern = append(pattern, entry) - } else { - exact[topic] = entry - } - } - if err := rows.Err(); err != nil { - return err - } - c.mu.Lock() - if len(exact) == 0 { - delete(c.exact, username) - } else { - c.exact[username] = exact - } - if len(pattern) == 0 { - delete(c.pattern, username) - } else { - c.pattern[username] = pattern - } - c.mu.Unlock() - return nil -} - -// newACLEntry builds an aclEntry from one user_access row's values. The -// isWildcard return tells the caller which storage slot the entry belongs in: -// the per-user wildcard slice if true, the per-user exact map if false. -// Wildcards have their LIKE pattern pre-compiled into entry.pattern; exact -// entries leave entry.pattern nil. -func newACLEntry(topic string, read, write bool) (entry aclEntry, isWildcard bool, err error) { - entry = aclEntry{length: len(topic), read: read, write: write} - if !strings.Contains(topic, "%") { - return entry, false, nil - } - re, err := compileLikeToRegex(topic) - if err != nil { - return entry, true, err - } - entry.pattern = re - return entry, true, nil -} - // Lookup returns the effective (read, write, found) permission for the given // (username, topic), preserving the priority ordering of the original SQL query: // 1. specific user beats Everyone @@ -182,6 +63,68 @@ func (c *accessCache) Lookup(usernameOrEveryone, topic string) (read, write, fou return false, false, false } +// reload scans (user_name, topic, read, write) rows and merges them into the +// cache. With no usernames the cache is replaced wholesale; otherwise the +// query is invoked with those usernames as positional args and only the +// listed users' slices are touched (a username absent from the result drops +// them from both maps). Runs against the primary so a reload after a +// mutation sees the just-written rows. +func (c *accessCache) reload(d *db.DB, query string, usernames ...string) error { + args := make([]any, len(usernames)) + for i, u := range usernames { + args[i] = u + } + rows, err := d.Query(query, args...) + if err != nil { + return err + } + defer rows.Close() + exacts := make(map[string]map[string]aclEntry) + patterns := make(map[string][]aclEntry) + for rows.Next() { + var u, topic string + var read, write bool + if err := rows.Scan(&u, &topic, &read, &write); err != nil { + return err + } + entry, isPattern, err := toACLEntry(topic, read, write) + if err != nil { + return err + } + if isPattern { + patterns[u] = append(patterns[u], entry) + } else { + if exacts[u] == nil { + exacts[u] = make(map[string]aclEntry) + } + exacts[u][topic] = entry + } + } + if err := rows.Err(); err != nil { + return err + } + c.mu.Lock() + defer c.mu.Unlock() + if len(usernames) == 0 { + c.exact = exacts + c.pattern = patterns + return nil + } + for _, u := range usernames { + if e, ok := exacts[u]; ok { + c.exact[u] = e + } else { + delete(c.exact, u) + } + if p, ok := patterns[u]; ok { + c.pattern[u] = p + } else { + delete(c.pattern, u) + } + } + return nil +} + // pickBestNoLock returns the highest-priority entry for a single user. When // more than one of that user's rules matches the requested topic, the winner // is chosen by: @@ -211,6 +154,24 @@ func (c *accessCache) pickBestNoLock(username, topic, escapedTopic string) (*acl return &best, found } +// toACLEntry builds an aclEntry from one user_access row's values. The +// isWildcard return tells the caller which storage slot the entry belongs in: +// the per-user wildcard slice if true, the per-user exact map if false. +// Wildcards have their LIKE pattern pre-compiled into entry.pattern; exact +// entries leave entry.pattern nil. +func toACLEntry(topic string, read, write bool) (entry aclEntry, isWildcard bool, err error) { + entry = aclEntry{length: len(topic), read: read, write: write} + if !strings.Contains(topic, "%") { + return entry, false, nil + } + pattern, err := compileLikeToRegex(topic) + if err != nil { + return entry, true, err + } + entry.pattern = pattern + return entry, true, nil +} + // better implements the (length DESC, write DESC) tie-break used by the original // query's ORDER BY for entries owned by the same user. func better(a, b aclEntry) bool { diff --git a/user/manager.go b/user/manager.go index 292d59b2..587ba8f2 100644 --- a/user/manager.go +++ b/user/manager.go @@ -93,39 +93,28 @@ func newManager(d *db.DB, queries queries, config *Config) (*Manager, error) { if err := manager.maybeProvisionUsersAccessAndTokens(); err != nil { return nil, err } - // Populate the cache after provisioning so the initial snapshot includes - // any provisioned access rules. No-op when the cache is disabled. - if err := manager.reloadAccessCache(); err != nil { - return nil, err - } - go manager.asyncQueueWriter(manager.config.QueueWriterInterval) - if manager.accessCache != nil && manager.config.AccessCacheReloadInterval > 0 { + if manager.accessCache != nil { + if err := manager.maybeReloadAccessCache(); err != nil { + return nil, err + } go manager.asyncAccessCacheReloader(manager.config.AccessCacheReloadInterval) } + go manager.asyncQueueWriter(manager.config.QueueWriterInterval) return manager, nil } -// reloadAccessCache rebuilds the in-memory access cache from the primary -// database. No-op when the cache is disabled. -func (a *Manager) reloadAccessCache() error { +// maybeReloadAccessCache refreshes the in-memory access cache from the +// primary database. No-op when the cache is disabled. With no usernames it +// does a full bulk reload; with one or more it refreshes only those users' +// slices in a single DB round-trip via an IN clause. +func (a *Manager) maybeReloadAccessCache(usernames ...string) error { if a.accessCache == nil { return nil } - return a.accessCache.reload(a.db, a.queries.selectAllAccessForCache) -} - -// reloadAccessCacheUsers refreshes the cache slices for the given usernames. -// No-op when the cache is disabled. -func (a *Manager) reloadAccessCacheUsers(usernames ...string) error { - if a.accessCache == nil { - return nil + if len(usernames) == 0 { + return a.accessCache.reload(a.db, a.queries.selectAccessCacheAll) } - for _, username := range usernames { - if err := a.accessCache.reloadUser(a.db, a.queries.selectAccessForCacheByUser, username); err != nil { - return err - } - } - return nil + return a.accessCache.reload(a.db, a.queries.selectAccessCacheUsersFn(len(usernames)), usernames...) } // asyncAccessCacheReloader periodically refreshes the access cache @@ -137,7 +126,7 @@ func (a *Manager) asyncAccessCacheReloader(interval time.Duration) { case <-a.quit: return case <-ticker.C: - if err := a.reloadAccessCache(); err != nil { + if err := a.maybeReloadAccessCache(); err != nil { log.Tag(tag).Err(err).Warn("Reloading ACL cache failed") } } @@ -225,7 +214,7 @@ func (a *Manager) RemoveUser(username string) error { // user_access rows are cascade-deleted along with the user (both by user_id // and by owner_user_id). Refresh this user's own slice (now empty) and // Everyone's slice, since reservations owned by this user landed there too. - return a.reloadAccessCacheUsers(username, Everyone) + return a.maybeReloadAccessCache(username, Everyone) } // removeUserTx deletes the user with the given username @@ -265,7 +254,7 @@ func (a *Manager) MarkUserRemoved(user *User) error { // resetUserAccessTx deleted this user's rows AND any row owned by this user // (typically the matching Everyone rows from their reservations). Refresh // both slices to mirror the DB exactly. - return a.reloadAccessCacheUsers(user.Name, Everyone) + return a.maybeReloadAccessCache(user.Name, Everyone) } // RemoveDeletedUsers deletes all users that have been marked deleted @@ -274,7 +263,7 @@ func (a *Manager) RemoveDeletedUsers() error { return err } // user_access rows are cascade-deleted with the users; refresh the snapshot. - return a.reloadAccessCache() + return a.maybeReloadAccessCache() } // ChangePassword changes a user's password @@ -315,7 +304,7 @@ func (a *Manager) ChangeRole(username string, role Role) error { // by this user (Everyone rows from their reservations). Other role changes // are no-ops for the cache but reloading the two affected slices is cheap // and keeps the code path uniform. - return a.reloadAccessCacheUsers(username, Everyone) + return a.maybeReloadAccessCache(username, Everyone) } // changeRoleTx changes a user's role @@ -683,7 +672,7 @@ func (a *Manager) AllowAccess(username string, topicPattern string, permission P return err } // Only this user's row set changed; refresh their slice only. - return a.reloadAccessCacheUsers(username) + return a.maybeReloadAccessCache(username) } func (a *Manager) allowAccessTx(tx *sql.Tx, username string, topicPattern string, permission Permission, provisioned bool) error { @@ -710,9 +699,9 @@ func (a *Manager) ResetAccess(username string, topicPattern string) error { // and deleteTopicAccess both touch rows owned by the user (typically // Everyone rows from their reservations). if username == "" { - return a.reloadAccessCache() + return a.maybeReloadAccessCache() } - return a.reloadAccessCacheUsers(username, Everyone) + return a.maybeReloadAccessCache(username, Everyone) } func (a *Manager) resetAccessTx(tx *sql.Tx, username string, topicPattern string) error { @@ -873,7 +862,7 @@ func (a *Manager) AddReservation(username string, topic string, everyone Permiss return err } // Both user's and Everyone's rows changed. - return a.reloadAccessCacheUsers(username, Everyone) + return a.maybeReloadAccessCache(username, Everyone) } // RemoveReservations deletes the access control entries associated with the given username/topic, @@ -901,7 +890,7 @@ func (a *Manager) RemoveReservations(username string, topics ...string) error { } // Mirror the DB: rows for this user and any Everyone rows owned by this // user are gone. Refresh both slices. - return a.reloadAccessCacheUsers(username, Everyone) + return a.maybeReloadAccessCache(username, Everyone) } // Reservations returns all user-owned topics, and the associated everyone-access diff --git a/user/manager_postgres.go b/user/manager_postgres.go index 8f0079b6..49d294b3 100644 --- a/user/manager_postgres.go +++ b/user/manager_postgres.go @@ -1,6 +1,9 @@ package user import ( + "fmt" + "strings" + "heckel.io/ntfy/v2/db" ) @@ -77,17 +80,11 @@ const ( WHERE (u.user_name = $1 OR u.user_name = $2) AND $3 LIKE a.topic ESCAPE '\' ORDER BY u.user_name DESC, LENGTH(a.topic) DESC, CASE WHEN a.write THEN 1 ELSE 0 END DESC ` - postgresSelectAllAccessForCacheQuery = ` + postgresSelectAccessCacheAllQuery = ` SELECT u.user_name, a.topic, a.read, a.write FROM user_access a JOIN "user" u ON u.id = a.user_id ` - postgresSelectAccessForCacheByUserQuery = ` - SELECT a.topic, a.read, a.write - FROM user_access a - JOIN "user" u ON u.id = a.user_id - WHERE u.user_name = $1 - ` postgresSelectUserAllAccessQuery = ` SELECT user_id, topic, read, write, provisioned FROM user_access @@ -232,6 +229,21 @@ const ( ` ) +// postgresSelectAccessCacheUsersQuery builds the per-users cache-load query +// with a "$1, $2, ..." IN clause sized for n usernames. +func postgresSelectAccessCacheUsersQuery(n int) string { + var sb strings.Builder + sb.WriteString(`SELECT u.user_name, a.topic, a.read, a.write FROM user_access a JOIN "user" u ON u.id = a.user_id WHERE u.user_name IN (`) + for i := 0; i < n; i++ { + if i > 0 { + sb.WriteString(",") + } + fmt.Fprintf(&sb, "$%d", i+1) + } + sb.WriteString(")") + return sb.String() +} + // NewPostgresManager creates a new Manager backed by a PostgreSQL database using an existing connection pool. var postgresQueries = queries{ selectUserByID: postgresSelectUserByIDQuery, @@ -256,8 +268,8 @@ var postgresQueries = queries{ deleteUsersMarked: postgresDeleteUsersMarkedQuery, deleteUsersProvisioned: postgresDeleteUsersProvisionedQuery, selectTopicPerms: postgresSelectTopicPermsQuery, - selectAllAccessForCache: postgresSelectAllAccessForCacheQuery, - selectAccessForCacheByUser: postgresSelectAccessForCacheByUserQuery, + selectAccessCacheAll: postgresSelectAccessCacheAllQuery, + selectAccessCacheUsersFn: postgresSelectAccessCacheUsersQuery, selectUserAllAccess: postgresSelectUserAllAccessQuery, selectUserAccess: postgresSelectUserAccessQuery, selectUserReservations: postgresSelectUserReservationsQuery, diff --git a/user/manager_sqlite.go b/user/manager_sqlite.go index a1ca4700..7c6507f2 100644 --- a/user/manager_sqlite.go +++ b/user/manager_sqlite.go @@ -4,6 +4,7 @@ import ( "database/sql" "fmt" "path/filepath" + "strings" _ "github.com/mattn/go-sqlite3" // SQLite driver @@ -83,17 +84,11 @@ const ( WHERE (u.user = ? OR u.user = ?) AND ? LIKE a.topic ESCAPE '\' ORDER BY u.user DESC, LENGTH(a.topic) DESC, a.write DESC ` - sqliteSelectAllAccessForCacheQuery = ` + sqliteSelectAccessCacheAllQuery = ` SELECT u.user, a.topic, a.read, a.write FROM user_access a JOIN user u ON u.id = a.user_id ` - sqliteSelectAccessForCacheByUserQuery = ` - SELECT a.topic, a.read, a.write - FROM user_access a - JOIN user u ON u.id = a.user_id - WHERE u.user = ? - ` sqliteSelectUserAllAccessQuery = ` SELECT user_id, topic, read, write, provisioned FROM user_access @@ -231,6 +226,16 @@ const ( ` ) +// sqliteSelectAccessCacheUsersQuery builds the per-users cache-load query +// with a "?, ?, ..." IN clause sized for n usernames. +func sqliteSelectAccessCacheUsersQuery(n int) string { + placeholders := strings.Repeat(",?", n) + if n > 0 { + placeholders = placeholders[1:] // drop the leading comma + } + return `SELECT u.user, a.topic, a.read, a.write FROM user_access a JOIN user u ON u.id = a.user_id WHERE u.user IN (` + placeholders + `)` +} + var sqliteQueries = queries{ selectUserByID: sqliteSelectUserByIDQuery, selectUserByName: sqliteSelectUserByNameQuery, @@ -254,8 +259,8 @@ var sqliteQueries = queries{ deleteUsersMarked: sqliteDeleteUsersMarkedQuery, deleteUsersProvisioned: sqliteDeleteUsersProvisionedQuery, selectTopicPerms: sqliteSelectTopicPermsQuery, - selectAllAccessForCache: sqliteSelectAllAccessForCacheQuery, - selectAccessForCacheByUser: sqliteSelectAccessForCacheByUserQuery, + selectAccessCacheAll: sqliteSelectAccessCacheAllQuery, + selectAccessCacheUsersFn: sqliteSelectAccessCacheUsersQuery, selectUserAllAccess: sqliteSelectUserAllAccessQuery, selectUserAccess: sqliteSelectUserAccessQuery, selectUserReservations: sqliteSelectUserReservationsQuery, diff --git a/user/types.go b/user/types.go index 65032aa3..d2746c77 100644 --- a/user/types.go +++ b/user/types.go @@ -318,9 +318,9 @@ type queries struct { deleteUsersProvisioned string // Access queries - selectTopicPerms string // Direct-DB authorizeTopicAccess query; used when the in-memory cache is disabled - selectAllAccessForCache string // Bulk load: (user_name, topic, read, write) for the in-memory ACL cache - selectAccessForCacheByUser string // Per-user load: (topic, read, write) for one username; used to refresh just one user's slice of the cache after mutation + selectTopicPerms string // Direct-DB authorizeTopicAccess query; used when the in-memory cache is disabled + selectAccessCacheAll string // Bulk load: (user_name, topic, read, write) for the in-memory ACL cache + selectAccessCacheUsersFn func(n int) string // Returns a per-users load query whose IN clause is sized for n usernames selectUserAllAccess string selectUserAccess string selectUserReservations string From 0e2c459d6b601d5b50f6d801d58f966bff257200 Mon Sep 17 00:00:00 2001 From: binwiederhier Date: Sun, 31 May 2026 15:42:12 -0400 Subject: [PATCH 10/18] Further review --- user/manager.go | 42 ++++++++++++++++++++++++---------------- user/manager_postgres.go | 2 +- user/manager_sqlite.go | 2 +- user/types.go | 2 +- 4 files changed, 28 insertions(+), 20 deletions(-) diff --git a/user/manager.go b/user/manager.go index 587ba8f2..58fd3b2d 100644 --- a/user/manager.go +++ b/user/manager.go @@ -62,7 +62,7 @@ type Manager struct { queries queries statsQueue map[string]*Stats // "Queue" to asynchronously write user stats to the database (UserID -> Stats) tokenQueue map[string]*TokenUpdate // "Queue" to asynchronously write token access stats to the database (Token ID -> TokenUpdate) - accessCache *accessCache // In-memory snapshot of user_access; rebuilt after every ACL mutation + accessCache *accessCache // In-memory snapshot of user_access; refreshed by maybeReloadAccessCache after every ACL mutation quit chan struct{} // Closed by Close() to signal background goroutines to stop mu sync.Mutex } @@ -76,7 +76,7 @@ func newManager(d *db.DB, queries queries, config *Config) (*Manager, error) { if config.QueueWriterInterval.Seconds() <= 0 { config.QueueWriterInterval = DefaultUserStatsQueueWriterInterval } - if config.AccessCacheReloadInterval == 0 { + if config.AccessCacheReloadInterval <= 0 { config.AccessCacheReloadInterval = DefaultAccessCacheReloadInterval } manager := &Manager{ @@ -87,13 +87,11 @@ func newManager(d *db.DB, queries queries, config *Config) (*Manager, error) { quit: make(chan struct{}), queries: queries, } - if config.AccessCacheEnabled { - manager.accessCache = newAccessCache() - } if err := manager.maybeProvisionUsersAccessAndTokens(); err != nil { return nil, err } - if manager.accessCache != nil { + if config.AccessCacheEnabled { + manager.accessCache = newAccessCache() if err := manager.maybeReloadAccessCache(); err != nil { return nil, err } @@ -114,10 +112,14 @@ func (a *Manager) maybeReloadAccessCache(usernames ...string) error { if len(usernames) == 0 { return a.accessCache.reload(a.db, a.queries.selectAccessCacheAll) } - return a.accessCache.reload(a.db, a.queries.selectAccessCacheUsersFn(len(usernames)), usernames...) + return a.accessCache.reload(a.db, a.queries.selectAccessCacheUsers(len(usernames)), usernames...) } -// asyncAccessCacheReloader periodically refreshes the access cache +// asyncAccessCacheReloader periodically bulk-reloads the access cache so that +// writes made by other processes against the same database (most notably the +// `ntfy access` CLI subcommand running while a server holds the cache) become +// visible within the configured interval. This Manager's own mutations do +// not depend on the poller -- they refresh affected users synchronously. func (a *Manager) asyncAccessCacheReloader(interval time.Duration) { ticker := time.NewTicker(interval) defer ticker.Stop() @@ -430,12 +432,18 @@ func (a *Manager) EnqueueUserStats(userID string, stats *Stats) { func (a *Manager) asyncQueueWriter(interval time.Duration) { ticker := time.NewTicker(interval) - for range ticker.C { - if err := a.writeUserStatsQueue(); err != nil { - log.Tag(tag).Err(err).Warn("Writing user stats queue failed") - } - if err := a.writeTokenUpdateQueue(); err != nil { - log.Tag(tag).Err(err).Warn("Writing token update queue failed") + defer ticker.Stop() + for { + select { + case <-a.quit: + return + case <-ticker.C: + if err := a.writeUserStatsQueue(); err != nil { + log.Tag(tag).Err(err).Warn("Writing user stats queue failed") + } + if err := a.writeTokenUpdateQueue(); err != nil { + log.Tag(tag).Err(err).Warn("Writing token update queue failed") + } } } } @@ -694,10 +702,10 @@ func (a *Manager) ResetAccess(username string, topicPattern string) error { if err != nil { return err } - // "Delete all access" affects every user; do the bulk reload. + // Empty username -> deleteAllAccess affected every user, bulk reload. // Otherwise refresh the named user plus Everyone, since resetUserAccessTx - // and deleteTopicAccess both touch rows owned by the user (typically - // Everyone rows from their reservations). + // and deleteTopicAccess both touch rows owned by the user (typically the + // Everyone row from their reservations). if username == "" { return a.maybeReloadAccessCache() } diff --git a/user/manager_postgres.go b/user/manager_postgres.go index 49d294b3..0395baae 100644 --- a/user/manager_postgres.go +++ b/user/manager_postgres.go @@ -269,7 +269,7 @@ var postgresQueries = queries{ deleteUsersProvisioned: postgresDeleteUsersProvisionedQuery, selectTopicPerms: postgresSelectTopicPermsQuery, selectAccessCacheAll: postgresSelectAccessCacheAllQuery, - selectAccessCacheUsersFn: postgresSelectAccessCacheUsersQuery, + selectAccessCacheUsers: postgresSelectAccessCacheUsersQuery, selectUserAllAccess: postgresSelectUserAllAccessQuery, selectUserAccess: postgresSelectUserAccessQuery, selectUserReservations: postgresSelectUserReservationsQuery, diff --git a/user/manager_sqlite.go b/user/manager_sqlite.go index 7c6507f2..ebe81d8d 100644 --- a/user/manager_sqlite.go +++ b/user/manager_sqlite.go @@ -260,7 +260,7 @@ var sqliteQueries = queries{ deleteUsersProvisioned: sqliteDeleteUsersProvisionedQuery, selectTopicPerms: sqliteSelectTopicPermsQuery, selectAccessCacheAll: sqliteSelectAccessCacheAllQuery, - selectAccessCacheUsersFn: sqliteSelectAccessCacheUsersQuery, + selectAccessCacheUsers: sqliteSelectAccessCacheUsersQuery, selectUserAllAccess: sqliteSelectUserAllAccessQuery, selectUserAccess: sqliteSelectUserAccessQuery, selectUserReservations: sqliteSelectUserReservationsQuery, diff --git a/user/types.go b/user/types.go index d2746c77..93aaa4d2 100644 --- a/user/types.go +++ b/user/types.go @@ -320,7 +320,7 @@ type queries struct { // Access queries selectTopicPerms string // Direct-DB authorizeTopicAccess query; used when the in-memory cache is disabled selectAccessCacheAll string // Bulk load: (user_name, topic, read, write) for the in-memory ACL cache - selectAccessCacheUsersFn func(n int) string // Returns a per-users load query whose IN clause is sized for n usernames + selectAccessCacheUsers func(n int) string // Returns a per-users load query whose IN clause is sized for n usernames selectUserAllAccess string selectUserAccess string selectUserReservations string From a2dc290f31a4c969e02ca6b348e33fd96ef590c9 Mon Sep 17 00:00:00 2001 From: binwiederhier Date: Sun, 31 May 2026 21:25:53 -0400 Subject: [PATCH 11/18] Review --- server/config.go | 6 +++--- user/access_cache_test.go | 12 ++---------- user/manager.go | 31 ++++++++++--------------------- user/types.go | 37 ++++++++++++------------------------- 4 files changed, 27 insertions(+), 59 deletions(-) diff --git a/server/config.go b/server/config.go index a1ba4d40..b7dadddf 100644 --- a/server/config.go +++ b/server/config.go @@ -116,8 +116,8 @@ type Config struct { AuthTokens map[string][]*user.Token AuthBcryptCost int AuthStatsQueueWriterInterval time.Duration - AuthAccessCacheEnabled bool - AuthAccessCacheReloadInterval time.Duration + AuthAccessCacheEnabled bool // Enables the in-memory ACL cache (high volume servers only) + AuthAccessCacheReloadInterval time.Duration // Reload interval for access cache, relevant for ACL writes from CLI AttachmentCacheDir string AttachmentTotalSizeLimit int64 AttachmentFileSizeLimit int64 @@ -225,7 +225,7 @@ func NewConfig() *Config { AuthDefault: user.PermissionReadWrite, AuthBcryptCost: user.DefaultUserPasswordBcryptCost, AuthStatsQueueWriterInterval: user.DefaultUserStatsQueueWriterInterval, - AuthAccessCacheEnabled: user.DefaultAccessCacheEnabled, // Opt-in (e.g. ntfy.sh) via server.yml + AuthAccessCacheEnabled: user.DefaultAccessCacheEnabled, AuthAccessCacheReloadInterval: user.DefaultAccessCacheReloadInterval, AttachmentCacheDir: "", AttachmentTotalSizeLimit: DefaultAttachmentTotalSizeLimit, diff --git a/user/access_cache_test.go b/user/access_cache_test.go index 65c109cb..17fd2c2c 100644 --- a/user/access_cache_test.go +++ b/user/access_cache_test.go @@ -2,6 +2,7 @@ package user import ( "regexp" + "strings" "sync" "sync/atomic" "testing" @@ -287,7 +288,7 @@ func loadCache(t *testing.T, c *accessCache, rows []rawACLRow) { wildcards := make(map[string][]aclEntry) for _, r := range rows { e := aclEntry{length: len(r.topic), read: r.read, write: r.write} - if containsPercent(r.topic) { + if strings.Contains(r.topic, "%") { e.pattern = mustCompileLikeToRegex(t, r.topic) wildcards[r.user] = append(wildcards[r.user], e) } else { @@ -309,12 +310,3 @@ func mustCompileLikeToRegex(t *testing.T, pattern string) *regexp.Regexp { require.NoError(t, err) return r } - -func containsPercent(s string) bool { - for i := 0; i < len(s); i++ { - if s[i] == '%' { - return true - } - } - return false -} diff --git a/user/manager.go b/user/manager.go index 58fd3b2d..6cc7e8b7 100644 --- a/user/manager.go +++ b/user/manager.go @@ -38,15 +38,8 @@ const ( const ( DefaultUserStatsQueueWriterInterval = 33 * time.Second DefaultUserPasswordBcryptCost = 10 - // DefaultAccessCacheEnabled is the default for Config.AccessCacheEnabled. - // Off by default so self-hosters keep the direct-DB authorizeTopicAccess - // path; ntfy.sh opts in via server config. - DefaultAccessCacheEnabled = false - // DefaultAccessCacheReloadInterval bounds how stale the in-memory ACL snapshot - // can be relative to writes made by *other* processes (e.g. a separate `ntfy - // access` CLI invocation modifying the same database). Only honored when the - // cache is enabled. - DefaultAccessCacheReloadInterval = 60 * time.Second + DefaultAccessCacheEnabled = false + DefaultAccessCacheReloadInterval = 60 * time.Second ) var ( @@ -95,9 +88,9 @@ func newManager(d *db.DB, queries queries, config *Config) (*Manager, error) { if err := manager.maybeReloadAccessCache(); err != nil { return nil, err } - go manager.asyncAccessCacheReloader(manager.config.AccessCacheReloadInterval) + go manager.asyncAccessCacheReloadLoop(manager.config.AccessCacheReloadInterval) } - go manager.asyncQueueWriter(manager.config.QueueWriterInterval) + go manager.asyncQueueWriteLoop(manager.config.QueueWriterInterval) return manager, nil } @@ -115,12 +108,12 @@ func (a *Manager) maybeReloadAccessCache(usernames ...string) error { return a.accessCache.reload(a.db, a.queries.selectAccessCacheUsers(len(usernames)), usernames...) } -// asyncAccessCacheReloader periodically bulk-reloads the access cache so that +// asyncAccessCacheReloadLoop periodically bulk-reloads the access cache so that // writes made by other processes against the same database (most notably the // `ntfy access` CLI subcommand running while a server holds the cache) become // visible within the configured interval. This Manager's own mutations do // not depend on the poller -- they refresh affected users synchronously. -func (a *Manager) asyncAccessCacheReloader(interval time.Duration) { +func (a *Manager) asyncAccessCacheReloadLoop(interval time.Duration) { ticker := time.NewTicker(interval) defer ticker.Stop() for { @@ -213,9 +206,7 @@ func (a *Manager) RemoveUser(username string) error { if err != nil { return err } - // user_access rows are cascade-deleted along with the user (both by user_id - // and by owner_user_id). Refresh this user's own slice (now empty) and - // Everyone's slice, since reservations owned by this user landed there too. + // Reload user-specific parts of the access cache return a.maybeReloadAccessCache(username, Everyone) } @@ -253,9 +244,7 @@ func (a *Manager) MarkUserRemoved(user *User) error { if err != nil { return err } - // resetUserAccessTx deleted this user's rows AND any row owned by this user - // (typically the matching Everyone rows from their reservations). Refresh - // both slices to mirror the DB exactly. + // Reload user-specific parts of the access cache return a.maybeReloadAccessCache(user.Name, Everyone) } @@ -264,7 +253,7 @@ func (a *Manager) RemoveDeletedUsers() error { if _, err := a.db.Exec(a.queries.deleteUsersMarked, time.Now().Unix()); err != nil { return err } - // user_access rows are cascade-deleted with the users; refresh the snapshot. + // Full cache reload, because we don't know what the query affects return a.maybeReloadAccessCache() } @@ -430,7 +419,7 @@ func (a *Manager) EnqueueUserStats(userID string, stats *Stats) { a.statsQueue[userID] = stats } -func (a *Manager) asyncQueueWriter(interval time.Duration) { +func (a *Manager) asyncQueueWriteLoop(interval time.Duration) { ticker := time.NewTicker(interval) defer ticker.Stop() for { diff --git a/user/types.go b/user/types.go index 93aaa4d2..c400e48f 100644 --- a/user/types.go +++ b/user/types.go @@ -245,31 +245,18 @@ const ( // Config holds the configuration for the user Manager type Config struct { - Filename string // Database filename, e.g. "/var/lib/ntfy/user.db" (SQLite) - DatabaseURL string // Database connection string (PostgreSQL) - StartupQueries string // Queries to run on startup, e.g. to create initial users or tiers (SQLite only) - DefaultAccess Permission // Default permission if no ACL matches - ProvisionEnabled bool // Hack: Enable auto-provisioning of users and access grants, disabled for "ntfy user" commands - Users []*User // Predefined users to create on startup - Access map[string][]*Grant // Predefined access grants to create on startup (username -> []*Grant) - Tokens map[string][]*Token // Predefined users to create on startup (username -> []*Token) - QueueWriterInterval time.Duration // Interval for the async queue writer to flush stats and token updates to the database - BcryptCost int // Cost of generated passwords; lowering makes testing faster - - // AccessCacheEnabled gates the in-memory ACL cache. When false (the - // default), authorizeTopicAccess runs the direct SQL query against the - // database on every call -- mutations and authorization are unaffected by - // any cache logic. When true, the Manager keeps an in-memory snapshot of - // user_access and serves authorizeTopicAccess from it; mutations refresh - // the affected slices, and a background poller picks up cross-process - // writes at AccessCacheReloadInterval. - AccessCacheEnabled bool - - // AccessCacheReloadInterval bounds the staleness of the in-memory ACL - // cache relative to writes from other processes (e.g. `ntfy access` CLI - // against a running server). Only honored when AccessCacheEnabled is true. - // Zero falls back to DefaultAccessCacheReloadInterval. - AccessCacheReloadInterval time.Duration + Filename string // Database filename, e.g. "/var/lib/ntfy/user.db" (SQLite) + DatabaseURL string // Database connection string (PostgreSQL) + StartupQueries string // Queries to run on startup, e.g. to create initial users or tiers (SQLite only) + DefaultAccess Permission // Default permission if no ACL matches + ProvisionEnabled bool // Hack: Enable auto-provisioning of users and access grants, disabled for "ntfy user" commands + Users []*User // Predefined users to create on startup + Access map[string][]*Grant // Predefined access grants to create on startup (username -> []*Grant) + Tokens map[string][]*Token // Predefined users to create on startup (username -> []*Token) + QueueWriterInterval time.Duration // Interval for the async queue writer to flush stats and token updates to the database + BcryptCost int // Cost of generated passwords; lowering makes testing faster + AccessCacheEnabled bool // Enables the in-memory ACL cache (high volume servers only) + AccessCacheReloadInterval time.Duration // Reload interval for access cache, relevant for ACL writes from CLI } // Error constants used by the package From 4196e6444c680ca28054bf633e195d3ee847fec4 Mon Sep 17 00:00:00 2001 From: binwiederhier Date: Sun, 31 May 2026 21:31:23 -0400 Subject: [PATCH 12/18] Stringbuilder --- user/manager_sqlite.go | 13 +++++++++---- 1 file changed, 9 insertions(+), 4 deletions(-) diff --git a/user/manager_sqlite.go b/user/manager_sqlite.go index ebe81d8d..18cb4028 100644 --- a/user/manager_sqlite.go +++ b/user/manager_sqlite.go @@ -229,11 +229,16 @@ const ( // sqliteSelectAccessCacheUsersQuery builds the per-users cache-load query // with a "?, ?, ..." IN clause sized for n usernames. func sqliteSelectAccessCacheUsersQuery(n int) string { - placeholders := strings.Repeat(",?", n) - if n > 0 { - placeholders = placeholders[1:] // drop the leading comma + var sb strings.Builder + sb.WriteString(`SELECT u.user, a.topic, a.read, a.write FROM user_access a JOIN user u ON u.id = a.user_id WHERE u.user IN (`) + for i := 0; i < n; i++ { + if i > 0 { + sb.WriteString(",") + } + sb.WriteString("?") } - return `SELECT u.user, a.topic, a.read, a.write FROM user_access a JOIN user u ON u.id = a.user_id WHERE u.user IN (` + placeholders + `)` + sb.WriteString(")") + return sb.String() } var sqliteQueries = queries{ From 7614405332d3a7cc70aeec8fb4d4d932906fe2a1 Mon Sep 17 00:00:00 2001 From: binwiederhier Date: Sun, 31 May 2026 21:50:57 -0400 Subject: [PATCH 13/18] Wire up server.yml --- cmd/serve.go | 3 +++ docs/config.md | 2 ++ docs/releases.md | 4 ++++ server/server.yml | 3 +++ 4 files changed, 12 insertions(+) diff --git a/cmd/serve.go b/cmd/serve.go index 0cfc1cbc..9712f94f 100644 --- a/cmd/serve.go +++ b/cmd/serve.go @@ -52,6 +52,7 @@ var flagsServe = append( altsrc.NewStringSliceFlag(&cli.StringSliceFlag{Name: "auth-users", Aliases: []string{"auth_users"}, EnvVars: []string{"NTFY_AUTH_USERS"}, Usage: "pre-provisioned declarative users"}), altsrc.NewStringSliceFlag(&cli.StringSliceFlag{Name: "auth-access", Aliases: []string{"auth_access"}, EnvVars: []string{"NTFY_AUTH_ACCESS"}, Usage: "pre-provisioned declarative access control entries"}), altsrc.NewStringSliceFlag(&cli.StringSliceFlag{Name: "auth-tokens", Aliases: []string{"auth_tokens"}, EnvVars: []string{"NTFY_AUTH_TOKENS"}, Usage: "pre-provisioned declarative access tokens"}), + altsrc.NewBoolFlag(&cli.BoolFlag{Name: "auth-access-cache", Aliases: []string{"auth_access_cache"}, EnvVars: []string{"NTFY_AUTH_ACCESS_CACHE"}, Value: user.DefaultAccessCacheEnabled, Usage: "enables the in-memory ACL cache (high-volume servers only)"}), altsrc.NewStringFlag(&cli.StringFlag{Name: "attachment-cache-dir", Aliases: []string{"attachment_cache_dir"}, EnvVars: []string{"NTFY_ATTACHMENT_CACHE_DIR"}, Usage: "cache directory for attached files, or S3 URL (s3://ACCESS_KEY:SECRET_KEY@BUCKET[/PREFIX]?region=REGION[&endpoint=ENDPOINT])"}), altsrc.NewStringFlag(&cli.StringFlag{Name: "attachment-total-size-limit", Aliases: []string{"attachment_total_size_limit", "A"}, EnvVars: []string{"NTFY_ATTACHMENT_TOTAL_SIZE_LIMIT"}, Value: util.FormatSize(server.DefaultAttachmentTotalSizeLimit), Usage: "limit of the on-disk attachment cache"}), altsrc.NewStringFlag(&cli.StringFlag{Name: "attachment-file-size-limit", Aliases: []string{"attachment_file_size_limit", "Y"}, EnvVars: []string{"NTFY_ATTACHMENT_FILE_SIZE_LIMIT"}, Value: util.FormatSize(server.DefaultAttachmentFileSizeLimit), Usage: "per-file attachment size limit (e.g. 300k, 2M, 100M)"}), @@ -168,6 +169,7 @@ func execServe(c *cli.Context) error { authUsersRaw := c.StringSlice("auth-users") authAccessRaw := c.StringSlice("auth-access") authTokensRaw := c.StringSlice("auth-tokens") + authAccessCacheEnabled := c.Bool("auth-access-cache") attachmentCacheDir := c.String("attachment-cache-dir") attachmentTotalSizeLimitStr := c.String("attachment-total-size-limit") attachmentFileSizeLimitStr := c.String("attachment-file-size-limit") @@ -468,6 +470,7 @@ func execServe(c *cli.Context) error { conf.AuthUsers = authUsers conf.AuthAccess = authAccess conf.AuthTokens = authTokens + conf.AuthAccessCacheEnabled = authAccessCacheEnabled conf.AttachmentCacheDir = attachmentCacheDir conf.AttachmentTotalSizeLimit = attachmentTotalSizeLimit conf.AttachmentFileSizeLimit = attachmentFileSizeLimit diff --git a/docs/config.md b/docs/config.md index c934143a..af6b9ec9 100644 --- a/docs/config.md +++ b/docs/config.md @@ -2284,6 +2284,7 @@ variable before running the `ntfy` command (e.g. `export NTFY_LISTEN_HTTP=:80`). | `cache-batch-timeout` | `NTFY_CACHE_BATCH_TIMEOUT` | *duration* | 0s | Timeout for batched async writes to the message cache (if zero, writes are synchronous) | | `auth-file` | `NTFY_AUTH_FILE` | *filename* | - | Auth database file used for access control (SQLite). If set, enables authentication and access control. Not required if `database-url` is set. See [access control](#access-control). | | `auth-default-access` | `NTFY_AUTH_DEFAULT_ACCESS` | `read-write`, `read-only`, `write-only`, `deny-all` | `read-write` | Default permissions if no matching entries in the auth database are found. Default is `read-write`. | +| `auth-access-cache` | `NTFY_AUTH_ACCESS_CACHE` | *bool* | false | Enables an in-memory snapshot of the access control list so authorization checks no longer hit the database. Off by default; only worth enabling on high-volume servers. ACL changes from a separate `ntfy access` CLI invocation against the same database become visible within ~60s; the server's own changes are immediate. | | `behind-proxy` | `NTFY_BEHIND_PROXY` | *bool* | false | If set, use forwarded header (e.g. X-Forwarded-For, X-Client-IP) to determine visitor IP address (for rate limiting) | | `proxy-forwarded-header` | `NTFY_PROXY_FORWARDED_HEADER` | *string* | `X-Forwarded-For` | Use specified header to determine visitor IP address (for rate limiting) | | `proxy-trusted-hosts` | `NTFY_PROXY_TRUSTED_HOSTS` | *comma-separated host/IP/CIDR list* | - | Comma-separated list of trusted IP addresses, hosts, or CIDRs to remove from forwarded header | @@ -2392,6 +2393,7 @@ OPTIONS: --auth-file value, --auth_file value, -H value auth database file used for access control [$NTFY_AUTH_FILE] --auth-startup-queries value, --auth_startup_queries value queries run when the auth database is initialized [$NTFY_AUTH_STARTUP_QUERIES] --auth-default-access value, --auth_default_access value, -p value default permissions if no matching entries in the auth database are found (default: "read-write") [$NTFY_AUTH_DEFAULT_ACCESS] + --auth-access-cache, --auth_access_cache enables the in-memory ACL cache (high-volume servers only) (default: false) [$NTFY_AUTH_ACCESS_CACHE] --attachment-cache-dir value, --attachment_cache_dir value cache directory for attached files, or S3 URL (s3://ACCESS_KEY:SECRET_KEY@BUCKET[/PREFIX]?region=REGION[&endpoint=ENDPOINT][&disable_http2=true]) [$NTFY_ATTACHMENT_CACHE_DIR] --attachment-total-size-limit value, --attachment_total_size_limit value, -A value limit of the on-disk attachment cache (default: "5G") [$NTFY_ATTACHMENT_TOTAL_SIZE_LIMIT] --attachment-file-size-limit value, --attachment_file_size_limit value, -Y value per-file attachment size limit (e.g. 300k, 2M, 100M) (default: "15M") [$NTFY_ATTACHMENT_FILE_SIZE_LIMIT] diff --git a/docs/releases.md b/docs/releases.md index df18add5..2853f25d 100644 --- a/docs/releases.md +++ b/docs/releases.md @@ -1926,6 +1926,10 @@ and the [ntfy Android app](https://github.com/binwiederhier/ntfy-android/release ### ntfy server v2.24.0 (UNRELEASED) +**Features:** + +* Add opt-in in-memory ACL cache (`auth-access-cache`) that serves topic authorization without a database round-trip; off by default, intended for high-volume servers + **Bug fixes + maintenance:** * Extend account token automatically from the PWA service worker, so installed PWAs don't get logged out ([#1669](https://github.com/binwiederhier/ntfy/pull/1669), [#1203](https://github.com/binwiederhier/ntfy/issues/1203), [#1533](https://github.com/binwiederhier/ntfy/issues/1533), thanks to [@nihalgonsalves](https://github.com/nihalgonsalves) for the contribution) diff --git a/server/server.yml b/server/server.yml index 08161dc2..ee23dd8b 100644 --- a/server/server.yml +++ b/server/server.yml @@ -116,6 +116,8 @@ # - auth-tokens is a list of access tokens that are automatically created when the server starts. # Each entry is in the format ":[: