Further review

This commit is contained in:
binwiederhier
2026-05-31 15:42:12 -04:00
parent 2f4afbdae5
commit 0e2c459d6b
4 changed files with 28 additions and 20 deletions
+25 -17
View File
@@ -62,7 +62,7 @@ type Manager struct {
queries queries queries queries
statsQueue map[string]*Stats // "Queue" to asynchronously write user stats to the database (UserID -> Stats) 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) 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 quit chan struct{} // Closed by Close() to signal background goroutines to stop
mu sync.Mutex mu sync.Mutex
} }
@@ -76,7 +76,7 @@ func newManager(d *db.DB, queries queries, config *Config) (*Manager, error) {
if config.QueueWriterInterval.Seconds() <= 0 { if config.QueueWriterInterval.Seconds() <= 0 {
config.QueueWriterInterval = DefaultUserStatsQueueWriterInterval config.QueueWriterInterval = DefaultUserStatsQueueWriterInterval
} }
if config.AccessCacheReloadInterval == 0 { if config.AccessCacheReloadInterval <= 0 {
config.AccessCacheReloadInterval = DefaultAccessCacheReloadInterval config.AccessCacheReloadInterval = DefaultAccessCacheReloadInterval
} }
manager := &Manager{ manager := &Manager{
@@ -87,13 +87,11 @@ func newManager(d *db.DB, queries queries, config *Config) (*Manager, error) {
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
} }
if manager.accessCache != nil { if config.AccessCacheEnabled {
manager.accessCache = newAccessCache()
if err := manager.maybeReloadAccessCache(); err != nil { if err := manager.maybeReloadAccessCache(); err != nil {
return nil, err return nil, err
} }
@@ -114,10 +112,14 @@ func (a *Manager) maybeReloadAccessCache(usernames ...string) error {
if len(usernames) == 0 { if len(usernames) == 0 {
return a.accessCache.reload(a.db, a.queries.selectAccessCacheAll) 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) { func (a *Manager) asyncAccessCacheReloader(interval time.Duration) {
ticker := time.NewTicker(interval) ticker := time.NewTicker(interval)
defer ticker.Stop() defer ticker.Stop()
@@ -430,12 +432,18 @@ func (a *Manager) EnqueueUserStats(userID string, stats *Stats) {
func (a *Manager) asyncQueueWriter(interval time.Duration) { func (a *Manager) asyncQueueWriter(interval time.Duration) {
ticker := time.NewTicker(interval) ticker := time.NewTicker(interval)
for range ticker.C { defer ticker.Stop()
if err := a.writeUserStatsQueue(); err != nil { for {
log.Tag(tag).Err(err).Warn("Writing user stats queue failed") select {
} case <-a.quit:
if err := a.writeTokenUpdateQueue(); err != nil { return
log.Tag(tag).Err(err).Warn("Writing token update queue failed") 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 { if err != nil {
return err 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 // Otherwise refresh the named user plus Everyone, since resetUserAccessTx
// and deleteTopicAccess both touch rows owned by the user (typically // and deleteTopicAccess both touch rows owned by the user (typically the
// Everyone rows from their reservations). // Everyone row from their reservations).
if username == "" { if username == "" {
return a.maybeReloadAccessCache() return a.maybeReloadAccessCache()
} }
+1 -1
View File
@@ -269,7 +269,7 @@ var postgresQueries = queries{
deleteUsersProvisioned: postgresDeleteUsersProvisionedQuery, deleteUsersProvisioned: postgresDeleteUsersProvisionedQuery,
selectTopicPerms: postgresSelectTopicPermsQuery, selectTopicPerms: postgresSelectTopicPermsQuery,
selectAccessCacheAll: postgresSelectAccessCacheAllQuery, selectAccessCacheAll: postgresSelectAccessCacheAllQuery,
selectAccessCacheUsersFn: postgresSelectAccessCacheUsersQuery, selectAccessCacheUsers: postgresSelectAccessCacheUsersQuery,
selectUserAllAccess: postgresSelectUserAllAccessQuery, selectUserAllAccess: postgresSelectUserAllAccessQuery,
selectUserAccess: postgresSelectUserAccessQuery, selectUserAccess: postgresSelectUserAccessQuery,
selectUserReservations: postgresSelectUserReservationsQuery, selectUserReservations: postgresSelectUserReservationsQuery,
+1 -1
View File
@@ -260,7 +260,7 @@ var sqliteQueries = queries{
deleteUsersProvisioned: sqliteDeleteUsersProvisionedQuery, deleteUsersProvisioned: sqliteDeleteUsersProvisionedQuery,
selectTopicPerms: sqliteSelectTopicPermsQuery, selectTopicPerms: sqliteSelectTopicPermsQuery,
selectAccessCacheAll: sqliteSelectAccessCacheAllQuery, selectAccessCacheAll: sqliteSelectAccessCacheAllQuery,
selectAccessCacheUsersFn: sqliteSelectAccessCacheUsersQuery, selectAccessCacheUsers: sqliteSelectAccessCacheUsersQuery,
selectUserAllAccess: sqliteSelectUserAllAccessQuery, selectUserAllAccess: sqliteSelectUserAllAccessQuery,
selectUserAccess: sqliteSelectUserAccessQuery, selectUserAccess: sqliteSelectUserAccessQuery,
selectUserReservations: sqliteSelectUserReservationsQuery, selectUserReservations: sqliteSelectUserReservationsQuery,
+1 -1
View File
@@ -320,7 +320,7 @@ type queries struct {
// Access queries // Access queries
selectTopicPerms string // Direct-DB authorizeTopicAccess query; used when the in-memory cache is disabled 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 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 selectUserAllAccess string
selectUserAccess string selectUserAccess string
selectUserReservations string selectUserReservations string