From 204723f3c09357e4eb8463420d4178d2b6006403 Mon Sep 17 00:00:00 2001 From: binwiederhier Date: Sun, 31 May 2026 14:28:04 -0400 Subject: [PATCH] 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