mirror of
https://github.com/multipleof4/ntfy.git
synced 2026-10-08 21:05:21 +00:00
Make opt-in flag
This commit is contained in:
@@ -7,7 +7,6 @@ import (
|
|||||||
"heckel.io/ntfy/v2/server"
|
"heckel.io/ntfy/v2/server"
|
||||||
"heckel.io/ntfy/v2/test"
|
"heckel.io/ntfy/v2/test"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
func TestCLI_Access_Show(t *testing.T) {
|
func TestCLI_Access_Show(t *testing.T) {
|
||||||
@@ -44,12 +43,6 @@ user * (role: anonymous, tier: none)
|
|||||||
`
|
`
|
||||||
require.Equal(t, expected, stdout.String())
|
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
|
// See if access permissions match
|
||||||
app, _, _, _ = newTestApp()
|
app, _, _, _ = newTestApp()
|
||||||
require.Error(t, app.Run([]string{
|
require.Error(t, app.Run([]string{
|
||||||
|
|||||||
+5
-5
@@ -378,11 +378,11 @@ func createUserManager(c *cli.Context) (*user.Manager, error) {
|
|||||||
ProvisionEnabled: false, // Hack: Do not re-provision users on manager initialization
|
ProvisionEnabled: false, // Hack: Do not re-provision users on manager initialization
|
||||||
BcryptCost: user.DefaultUserPasswordBcryptCost,
|
BcryptCost: user.DefaultUserPasswordBcryptCost,
|
||||||
QueueWriterInterval: user.DefaultUserStatsQueueWriterInterval,
|
QueueWriterInterval: user.DefaultUserStatsQueueWriterInterval,
|
||||||
// CLI Managers are short-lived; the background ACL cache poller would only
|
// CLI subcommands never serve authorizeTopicAccess and are short-lived,
|
||||||
// spam "database is closed" warnings after the subcommand returns. Mutations
|
// so the cache (and its background poller) would be wasted work. Mutations
|
||||||
// still refresh the local cache synchronously; the running server (if any)
|
// hit the DB directly; the running server, if any, picks them up via its
|
||||||
// picks them up via its own poller.
|
// own poller when the cache is enabled there.
|
||||||
AccessCacheReloadInterval: -1,
|
AccessCacheEnabled: false,
|
||||||
}
|
}
|
||||||
if databaseURL != "" {
|
if databaseURL != "" {
|
||||||
host, dbErr := pg.Open(databaseURL)
|
host, dbErr := pg.Open(databaseURL)
|
||||||
|
|||||||
+2
-5
@@ -9,7 +9,6 @@ import (
|
|||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
func TestCLI_User_Add(t *testing.T) {
|
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.File = configFile
|
||||||
conf.AuthFile = filepath.Join(t.TempDir(), "user.db")
|
conf.AuthFile = filepath.Join(t.TempDir(), "user.db")
|
||||||
conf.AuthDefault = user.PermissionDenyAll
|
conf.AuthDefault = user.PermissionDenyAll
|
||||||
// Tight interval so cross-process writes from the `ntfy access`/`ntfy user`
|
// Cache is off by default (matches self-hoster setup), so the server reads
|
||||||
// CLI commands (which run via a separate Manager) propagate to the server's
|
// authorizations directly from the DB and sees CLI mutations immediately.
|
||||||
// ACL cache within tens of ms instead of the default 5s.
|
|
||||||
conf.AuthAccessCacheReloadInterval = 25 * time.Millisecond
|
|
||||||
s, port = test.StartServerWithConfig(t, conf)
|
s, port = test.StartServerWithConfig(t, conf)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -116,6 +116,7 @@ type Config struct {
|
|||||||
AuthTokens map[string][]*user.Token
|
AuthTokens map[string][]*user.Token
|
||||||
AuthBcryptCost int
|
AuthBcryptCost int
|
||||||
AuthStatsQueueWriterInterval time.Duration
|
AuthStatsQueueWriterInterval time.Duration
|
||||||
|
AuthAccessCacheEnabled bool
|
||||||
AuthAccessCacheReloadInterval time.Duration
|
AuthAccessCacheReloadInterval time.Duration
|
||||||
AttachmentCacheDir string
|
AttachmentCacheDir string
|
||||||
AttachmentTotalSizeLimit int64
|
AttachmentTotalSizeLimit int64
|
||||||
@@ -224,6 +225,7 @@ func NewConfig() *Config {
|
|||||||
AuthDefault: user.PermissionReadWrite,
|
AuthDefault: user.PermissionReadWrite,
|
||||||
AuthBcryptCost: user.DefaultUserPasswordBcryptCost,
|
AuthBcryptCost: user.DefaultUserPasswordBcryptCost,
|
||||||
AuthStatsQueueWriterInterval: user.DefaultUserStatsQueueWriterInterval,
|
AuthStatsQueueWriterInterval: user.DefaultUserStatsQueueWriterInterval,
|
||||||
|
AuthAccessCacheEnabled: user.DefaultAccessCacheEnabled, // Opt-in (e.g. ntfy.sh) via server.yml
|
||||||
AuthAccessCacheReloadInterval: user.DefaultAccessCacheReloadInterval,
|
AuthAccessCacheReloadInterval: user.DefaultAccessCacheReloadInterval,
|
||||||
AttachmentCacheDir: "",
|
AttachmentCacheDir: "",
|
||||||
AttachmentTotalSizeLimit: DefaultAttachmentTotalSizeLimit,
|
AttachmentTotalSizeLimit: DefaultAttachmentTotalSizeLimit,
|
||||||
|
|||||||
@@ -257,6 +257,7 @@ func New(conf *Config) (*Server, error) {
|
|||||||
Tokens: conf.AuthTokens,
|
Tokens: conf.AuthTokens,
|
||||||
BcryptCost: conf.AuthBcryptCost,
|
BcryptCost: conf.AuthBcryptCost,
|
||||||
QueueWriterInterval: conf.AuthStatsQueueWriterInterval,
|
QueueWriterInterval: conf.AuthStatsQueueWriterInterval,
|
||||||
|
AccessCacheEnabled: conf.AuthAccessCacheEnabled,
|
||||||
AccessCacheReloadInterval: conf.AuthAccessCacheReloadInterval,
|
AccessCacheReloadInterval: conf.AuthAccessCacheReloadInterval,
|
||||||
}
|
}
|
||||||
if pool != nil {
|
if pool != nil {
|
||||||
|
|||||||
+54
-18
@@ -38,9 +38,14 @@ const (
|
|||||||
const (
|
const (
|
||||||
DefaultUserStatsQueueWriterInterval = 33 * time.Second
|
DefaultUserStatsQueueWriterInterval = 33 * time.Second
|
||||||
DefaultUserPasswordBcryptCost = 10
|
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
|
// DefaultAccessCacheReloadInterval bounds how stale the in-memory ACL snapshot
|
||||||
// can be relative to writes made by *other* processes (e.g. a separate `ntfy
|
// 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
|
DefaultAccessCacheReloadInterval = 60 * time.Second
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -75,34 +80,46 @@ func newManager(d *db.DB, queries queries, config *Config) (*Manager, error) {
|
|||||||
config.AccessCacheReloadInterval = DefaultAccessCacheReloadInterval
|
config.AccessCacheReloadInterval = DefaultAccessCacheReloadInterval
|
||||||
}
|
}
|
||||||
manager := &Manager{
|
manager := &Manager{
|
||||||
config: config,
|
config: config,
|
||||||
db: d,
|
db: d,
|
||||||
statsQueue: make(map[string]*Stats),
|
statsQueue: make(map[string]*Stats),
|
||||||
tokenQueue: make(map[string]*TokenUpdate),
|
tokenQueue: make(map[string]*TokenUpdate),
|
||||||
accessCache: newAccessCache(),
|
quit: make(chan struct{}),
|
||||||
quit: make(chan struct{}),
|
queries: queries,
|
||||||
queries: queries,
|
}
|
||||||
|
if config.AccessCacheEnabled {
|
||||||
|
manager.accessCache = newAccessCache()
|
||||||
}
|
}
|
||||||
if err := manager.maybeProvisionUsersAccessAndTokens(); err != nil {
|
if err := manager.maybeProvisionUsersAccessAndTokens(); err != nil {
|
||||||
return nil, err
|
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 {
|
if err := manager.reloadAccessCache(); err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
go manager.asyncQueueWriter(manager.config.QueueWriterInterval)
|
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)
|
go manager.asyncAccessCacheReloader(manager.config.AccessCacheReloadInterval)
|
||||||
}
|
}
|
||||||
return manager, nil
|
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 {
|
func (a *Manager) reloadAccessCache() error {
|
||||||
|
if a.accessCache == nil {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
return a.accessCache.reload(a.db, a.queries.selectAllAccessForCache)
|
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 {
|
func (a *Manager) reloadAccessCacheUsers(usernames ...string) error {
|
||||||
|
if a.accessCache == nil {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
for _, username := range usernames {
|
for _, username := range usernames {
|
||||||
if err := a.accessCache.reloadUser(a.db, a.queries.selectAccessForCacheByUser, username); err != nil {
|
if err := a.accessCache.reloadUser(a.db, a.queries.selectAccessForCacheByUser, username); err != nil {
|
||||||
return err
|
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.
|
// 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 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.
|
// Priority:
|
||||||
// - Furthermore, the lookup prioritizes more specific permissions (longer!) over more generic ones, e.g. "test*" > "*"
|
// - specific user beats Everyone
|
||||||
// - It also prioritizes write permissions over read permissions
|
// - 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,
|
// When AccessCacheEnabled is true (config), the lookup is served entirely from
|
||||||
// so this is on the hot path of every authenticatable HTTP request and must stay allocation-free.
|
// 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) {
|
func (a *Manager) authorizeTopicAccess(usernameOrEveryone, topic string) (read, write, found bool, err error) {
|
||||||
read, write, found = a.accessCache.Lookup(usernameOrEveryone, topic)
|
if a.accessCache != nil {
|
||||||
return read, write, found, 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
|
// AllGrants returns all user-specific access control entries, mapped to their respective user IDs
|
||||||
|
|||||||
@@ -70,6 +70,13 @@ const (
|
|||||||
postgresDeleteUsersProvisionedQuery = `DELETE FROM "user" WHERE provisioned = true`
|
postgresDeleteUsersProvisionedQuery = `DELETE FROM "user" WHERE provisioned = true`
|
||||||
|
|
||||||
// Access queries
|
// 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 = `
|
postgresSelectAllAccessForCacheQuery = `
|
||||||
SELECT u.user_name, a.topic, a.read, a.write
|
SELECT u.user_name, a.topic, a.read, a.write
|
||||||
FROM user_access a
|
FROM user_access a
|
||||||
@@ -248,6 +255,7 @@ var postgresQueries = queries{
|
|||||||
deleteUserTier: postgresDeleteUserTierQuery,
|
deleteUserTier: postgresDeleteUserTierQuery,
|
||||||
deleteUsersMarked: postgresDeleteUsersMarkedQuery,
|
deleteUsersMarked: postgresDeleteUsersMarkedQuery,
|
||||||
deleteUsersProvisioned: postgresDeleteUsersProvisionedQuery,
|
deleteUsersProvisioned: postgresDeleteUsersProvisionedQuery,
|
||||||
|
selectTopicPerms: postgresSelectTopicPermsQuery,
|
||||||
selectAllAccessForCache: postgresSelectAllAccessForCacheQuery,
|
selectAllAccessForCache: postgresSelectAllAccessForCacheQuery,
|
||||||
selectAccessForCacheByUser: postgresSelectAccessForCacheByUserQuery,
|
selectAccessForCacheByUser: postgresSelectAccessForCacheByUserQuery,
|
||||||
selectUserAllAccess: postgresSelectUserAllAccessQuery,
|
selectUserAllAccess: postgresSelectUserAllAccessQuery,
|
||||||
|
|||||||
@@ -76,6 +76,13 @@ const (
|
|||||||
sqliteDeleteUsersProvisionedQuery = `DELETE FROM user WHERE provisioned = 1`
|
sqliteDeleteUsersProvisionedQuery = `DELETE FROM user WHERE provisioned = 1`
|
||||||
|
|
||||||
// Access queries
|
// 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 = `
|
sqliteSelectAllAccessForCacheQuery = `
|
||||||
SELECT u.user, a.topic, a.read, a.write
|
SELECT u.user, a.topic, a.read, a.write
|
||||||
FROM user_access a
|
FROM user_access a
|
||||||
@@ -246,6 +253,7 @@ var sqliteQueries = queries{
|
|||||||
deleteUserTier: sqliteDeleteUserTierQuery,
|
deleteUserTier: sqliteDeleteUserTierQuery,
|
||||||
deleteUsersMarked: sqliteDeleteUsersMarkedQuery,
|
deleteUsersMarked: sqliteDeleteUsersMarkedQuery,
|
||||||
deleteUsersProvisioned: sqliteDeleteUsersProvisionedQuery,
|
deleteUsersProvisioned: sqliteDeleteUsersProvisionedQuery,
|
||||||
|
selectTopicPerms: sqliteSelectTopicPermsQuery,
|
||||||
selectAllAccessForCache: sqliteSelectAllAccessForCacheQuery,
|
selectAllAccessForCache: sqliteSelectAllAccessForCacheQuery,
|
||||||
selectAccessForCacheByUser: sqliteSelectAccessForCacheByUserQuery,
|
selectAccessForCacheByUser: sqliteSelectAccessForCacheByUserQuery,
|
||||||
selectUserAllAccess: sqliteSelectUserAllAccessQuery,
|
selectUserAllAccess: sqliteSelectUserAllAccessQuery,
|
||||||
|
|||||||
@@ -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) {
|
func TestStoreReservations(t *testing.T) {
|
||||||
forEachStoreBackend(t, func(t *testing.T, manager *Manager) {
|
forEachStoreBackend(t, func(t *testing.T, manager *Manager) {
|
||||||
require.Nil(t, manager.AddUser("phil", "mypass", RoleUser, false))
|
require.Nil(t, manager.AddUser("phil", "mypass", RoleUser, false))
|
||||||
|
|||||||
+14
-8
@@ -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
|
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
|
BcryptCost int // Cost of generated passwords; lowering makes testing faster
|
||||||
|
|
||||||
// AccessCacheReloadInterval bounds the staleness of the in-memory ACL cache
|
// AccessCacheEnabled gates the in-memory ACL cache. When false (the
|
||||||
// relative to writes from other processes (e.g. `ntfy access` CLI against a
|
// default), authorizeTopicAccess runs the direct SQL query against the
|
||||||
// running server).
|
// database on every call -- mutations and authorization are unaffected by
|
||||||
// 0 -> use DefaultAccessCacheReloadInterval
|
// any cache logic. When true, the Manager keeps an in-memory snapshot of
|
||||||
// negative -> disable the background poller; cache only refreshes on this
|
// user_access and serves authorizeTopicAccess from it; mutations refresh
|
||||||
// Manager's own ACL mutations. Use this for short-lived Managers
|
// the affected slices, and a background poller picks up cross-process
|
||||||
// (e.g. the CLI subcommands) where polling is wasted work.
|
// writes at AccessCacheReloadInterval.
|
||||||
// positive -> poll at the given interval
|
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
|
AccessCacheReloadInterval time.Duration
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -313,6 +318,7 @@ type queries struct {
|
|||||||
deleteUsersProvisioned string
|
deleteUsersProvisioned string
|
||||||
|
|
||||||
// Access queries
|
// 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
|
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
|
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
|
selectUserAllAccess string
|
||||||
|
|||||||
Reference in New Issue
Block a user