feat: quota-notify and membership-unsubscribe activity events

EventUserExceedQuotaNotified now fires at every quota rejection:
validateUserCapacityRaw covers the upload pre-check, batch validation,
and copy walk; the copy-path caller re-records post-rollback since the
in-tx record is discarded, and the atomic ReserveStorage race gate
records after its rollback for the same reason.

ExpireGrants returns the grants whose group reversion actually applied
(still only when the user sits on the granted group), and the
grant_expire cron emits EventMembershipUnsubscribe for each with
from_group/to_group metadata.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
pull/3587/head
Tomas Dvorak 2 weeks ago
parent 0582fa8ab7
commit 521bfda63c

@ -194,7 +194,8 @@ Order = user-visible value first; each ships with backend + UI + tests.
- [x] `activity_event` entity (immutable, tx-aware, actor+IP+CID) + per-file Activity dialog + admin `/admin/event` feed + per-type enablement + retention cron (#184)
- [x] Coverage wave 2: email/user-activated/token-refresh/share-viewed/version/metadata/view/thumb/live-photo/copy-from/webdav/profile+security/oauth/admin-ops/import (1bbaddf)
- [x] link/unlink_account events wired via `sso_binding` flows (#190)
- [ ] Event coverage remainder (needs unbuilt features): payment_*, membership_unsubscribe, mount, quota-notify
- [x] Event coverage wave 3: `membership_unsubscribe` emitted by grant-expiry cron for reverted group grants; `user_exceed_quota_notified` at every quota rejection (upload pre-check, atomic reserve, copy)
- [ ] Event coverage remainder (needs unbuilt features): payment_* (no payment processor), mount (policy-mount feature absent)
- [x] site announcement: `announcement` setting (markdown) + post-login modal + per-user dismissal re-triggering on content change (#184)
- [x] `abuse_report` entity + public `POST /abuse/report` (IP rate-limit + `abuse_captcha` gate) + admin `/admin/abuse` queue (resolve/dismiss + reversible share block) + share-menu Report entry (#185)
- [x] group `allowed_nodes` pool + `allow_select_node` + task `target_node` dispatch (persisted in task state, weighted LB within pool); group admin multi-select + task-dialog node picker

@ -51,8 +51,9 @@ type (
ListGrants(ctx context.Context, userID int) ([]*ent.UserGrant, error)
// ExpireGrants deletes expired storage grants and reverts group
// upgrades whose users still sit on the granted group. Users without a
// recorded previous group fall back to defaultGroupID.
ExpireGrants(ctx context.Context, defaultGroupID int) error
// recorded previous group fall back to defaultGroupID. It returns the
// grants whose group reversion actually applied.
ExpireGrants(ctx context.Context, defaultGroupID int) ([]*ent.UserGrant, error)
// ListSkus returns products ordered by weight desc then id. When
// onlyEnabled is set, disabled products are excluded.
ListSkus(ctx context.Context, onlyEnabled bool) ([]*ent.Sku, error)
@ -319,15 +320,16 @@ func (c *vasClient) ListGrants(ctx context.Context, userID int) ([]*ent.UserGran
All(ctx)
}
func (c *vasClient) ExpireGrants(ctx context.Context, defaultGroupID int) error {
func (c *vasClient) ExpireGrants(ctx context.Context, defaultGroupID int) ([]*ent.UserGrant, error) {
now := time.Now()
expired, err := c.client.UserGrant.Query().
Where(usergrant.ExpiresAtLT(now)).
All(ctx)
if err != nil {
return err
return nil, err
}
var reverted []*ent.UserGrant
for _, g := range expired {
if g.Type == usergrant.TypeGroup {
revertTo := g.PrevGroupID
@ -336,20 +338,24 @@ func (c *vasClient) ExpireGrants(ctx context.Context, defaultGroupID int) error
}
// Revert only when the user still sits on the granted group —
// intervening group changes win.
if _, err := c.client.User.Update().
affected, err := c.client.User.Update().
Where(user.ID(g.UserID), user.GroupUsers(int(g.Amount))).
SetGroupUsers(revertTo).
Save(ctx); err != nil {
return err
Save(ctx)
if err != nil {
return nil, err
}
if affected > 0 {
reverted = append(reverted, g)
}
}
if _, err := c.client.UserGrant.Delete().
Where(usergrant.ID(g.ID)).
Exec(schema.SkipSoftDelete(ctx)); err != nil {
return err
return nil, err
}
}
return nil
return reverted, nil
}
func (c *vasClient) ListSkus(ctx context.Context, onlyEnabled bool) ([]*ent.Sku, error) {

@ -165,7 +165,10 @@ func TestExpireGrants(t *testing.T) {
require.NoError(t, err)
require.Equal(t, int64(512), bonus)
require.NoError(t, c.ExpireGrants(ctx, 2))
reverted, err := c.ExpireGrants(ctx, 2)
require.NoError(t, err)
require.Len(t, reverted, 1)
require.Equal(t, u.ID, reverted[0].UserID)
require.Equal(t, group.ID, client.User.GetX(ctx, u.ID).GroupUsers)
require.Equal(t, 1, client.UserGrant.Query().CountX(ctx))
}

@ -2,6 +2,7 @@ package dbfs
import (
"context"
"errors"
"fmt"
"path/filepath"
"strconv"
@ -697,6 +698,13 @@ func (f *DBFS) MoveOrCopy(ctx context.Context, path []*fs.URI, dst *fs.URI, isCo
if err != nil {
_ = inventory.Rollback(tx)
if errors.Is(err, fs.ErrInsufficientCapacity) {
extra := map[string]any{"dst": destination.Uri(true).String()}
if capacity, cerr := f.Capacity(ctx, destination.Owner()); cerr == nil {
extra["capacity"] = capacity.Total
}
f.record(ctx, types.EventUserExceedQuotaNotified, activity.Extra(extra))
}
return nil, err
}

@ -223,6 +223,9 @@ func (f *DBFS) PrepareUpload(ctx context.Context, req *fs.UploadRequest, opts ..
if err := dbTx.ReserveStorage(ctx, f.userClient, owner.ID, req.Props.Size, capacity.Total); err != nil {
_ = inventory.Rollback(dbTx)
if errors.Is(err, inventory.ErrInsufficientCapacity) {
f.record(ctx, types.EventUserExceedQuotaNotified, activity.Extra(map[string]any{
"size": req.Props.Size, "capacity": capacity.Total,
}))
return nil, fs.ErrInsufficientCapacity
}
return nil, serializer.NewError(serializer.CodeDBError, "Failed to reserve storage capacity", err)

@ -6,6 +6,7 @@ import (
"time"
"github.com/cloudreve/Cloudreve/v4/ent"
"github.com/cloudreve/Cloudreve/v4/ent/activityevent"
"github.com/cloudreve/Cloudreve/v4/ent/enttest"
entfile "github.com/cloudreve/Cloudreve/v4/ent/file"
"github.com/cloudreve/Cloudreve/v4/ent/storagepolicy"
@ -32,6 +33,10 @@ func (p dedupSettingProvider) DBFS(context.Context) *setting.DBFS {
return &setting.DBFS{DedupScope: p.scope, MaxPageSize: 200}
}
func (p dedupSettingProvider) AuditLogEnabled(context.Context, int) bool {
return true
}
type stubEventHub struct{}
func (stubEventHub) Subscribe(context.Context, int, string) (chan *eventhub.Event, bool, error) {
@ -43,7 +48,7 @@ func (stubEventHub) GetSubscribers(context.Context, int) []eventhub.Subscriber {
}
func (stubEventHub) Close() {}
func dedupUploadFixture(t *testing.T, client *ent.Client, scope string) (*ent.User, *ent.StoragePolicy, *DBFS) {
func dedupUploadFixture(t *testing.T, client *ent.Client, scope string, maxStorage int64) (*ent.User, *ent.StoragePolicy, *DBFS) {
t.Helper()
ctx := context.Background()
l := logging.NewConsoleLogger(logging.LevelError)
@ -53,7 +58,7 @@ func dedupUploadFixture(t *testing.T, client *ent.Client, scope string) (*ent.Us
p := client.StoragePolicy.Create().SetName("local").SetType("local").
SetStatus(storagepolicy.StatusActive).SetSettings(&types.PolicySetting{}).SaveX(ctx)
group := client.Group.Create().SetName("g").SetPermissions(&boolset.BooleanSet{}).
SetMaxStorage(1 << 40).SetStoragePolicies(p).SaveX(ctx)
SetMaxStorage(maxStorage).SetStoragePolicies(p).SaveX(ctx)
u := client.User.Create().SetEmail("u@example.com").SetNick("u").SetGroup(group).SaveX(ctx)
u.SetGroup(group)
client.File.Create().SetName(inventory.RootFolderName).
@ -65,6 +70,7 @@ func dedupUploadFixture(t *testing.T, client *ent.Client, scope string) (*ent.Us
fileClient: inventory.NewFileClient(client, conf.SQLiteDB, hasher),
userClient: inventory.NewUserClient(client),
storagePolicyClient: inventory.NewStoragePolicyClient(client, nil),
activityClient: inventory.NewActivityClient(client, conf.SQLiteDB),
settingClient: dedupSettingProvider{scope: scope},
hasher: hasher,
l: l,
@ -109,7 +115,7 @@ func TestPrepareUploadRapid(t *testing.T) {
t.Cleanup(func() { require.NoError(t, client.Close()) })
ctx := context.Background()
u, p, f := dedupUploadFixture(t, client, "owner")
u, p, f := dedupUploadFixture(t, client, "owner", 1<<40)
existing := seedCompletedEntity(t, client, u, p, dedupTestHash, 1024)
// Normal upload without hash -> transfer session
@ -141,10 +147,26 @@ func TestPrepareUploadRapidScopeOff(t *testing.T) {
t.Cleanup(func() { require.NoError(t, client.Close()) })
ctx := context.Background()
u, p, f := dedupUploadFixture(t, client, "off")
u, p, f := dedupUploadFixture(t, client, "off", 1<<40)
seedCompletedEntity(t, client, u, p, dedupTestHash, 1024)
s, err := f.PrepareUpload(ctx, uploadReq(t, f.hasher, u, "copy.txt", 1024, dedupTestHash))
require.NoError(t, err)
require.False(t, s.Props.RapidUploaded)
}
func TestPrepareUploadQuotaEvent(t *testing.T) {
client := enttest.Open(t, "sqlite3", "file:"+t.Name()+"?mode=memory&cache=shared")
t.Cleanup(func() { require.NoError(t, client.Close()) })
ctx := context.Background()
u, _, f := dedupUploadFixture(t, client, "owner", 1024)
_, err := f.PrepareUpload(ctx, uploadReq(t, f.hasher, u, "big.bin", 2048, ""))
require.ErrorIs(t, err, fs.ErrInsufficientCapacity)
ev := client.ActivityEvent.Query().
Where(activityevent.TypeEQ(types.EventUserExceedQuotaNotified)).OnlyX(ctx)
require.Equal(t, float64(2048), ev.Extra["size"])
require.Equal(t, float64(1024), ev.Extra["capacity"])
}

@ -7,6 +7,8 @@ import (
"strings"
"github.com/cloudreve/Cloudreve/v4/ent"
"github.com/cloudreve/Cloudreve/v4/inventory/types"
"github.com/cloudreve/Cloudreve/v4/pkg/activity"
"github.com/cloudreve/Cloudreve/v4/pkg/filemanager/fs"
"github.com/cloudreve/Cloudreve/v4/pkg/util"
)
@ -112,6 +114,9 @@ func (f *DBFS) validateUserCapacity(ctx context.Context, size int64, u *ent.User
// validateUserCapacityRaw validates the user capacity, but does not fetch the capacity.
func (f *DBFS) validateUserCapacityRaw(ctx context.Context, size int64, capacity *fs.Capacity) error {
if capacity.Used+size > capacity.Total {
f.record(ctx, types.EventUserExceedQuotaNotified, activity.Extra(map[string]any{
"size": size, "capacity": capacity.Total,
}))
return fs.ErrInsufficientCapacity
}
return nil

@ -4,6 +4,9 @@ import (
"context"
"github.com/cloudreve/Cloudreve/v4/application/dependency"
"github.com/cloudreve/Cloudreve/v4/ent"
"github.com/cloudreve/Cloudreve/v4/inventory/types"
"github.com/cloudreve/Cloudreve/v4/pkg/activity"
"github.com/cloudreve/Cloudreve/v4/pkg/crontab"
"github.com/cloudreve/Cloudreve/v4/pkg/setting"
)
@ -11,8 +14,28 @@ import (
func init() {
crontab.Register(setting.CronTypeGrantExpire, func(ctx context.Context) {
dep := dependency.FromContext(ctx)
if err := dep.VasClient().ExpireGrants(ctx, dep.SettingProvider().DefaultGroup(ctx)); err != nil {
defaultGroup := dep.SettingProvider().DefaultGroup(ctx)
reverted, err := dep.VasClient().ExpireGrants(ctx, defaultGroup)
if err != nil {
dep.Logger().Error("Failed to expire user grants: %s", err)
return
}
recordMembershipUnsubscribes(ctx, dep, reverted, defaultGroup)
})
}
// recordMembershipUnsubscribes emits an unsubscribe event for each grant
// whose group assignment was reverted by ExpireGrants.
func recordMembershipUnsubscribes(ctx context.Context, dep dependency.Dep, reverted []*ent.UserGrant, defaultGroup int) {
for _, g := range reverted {
revertTo := g.PrevGroupID
if revertTo <= 0 {
revertTo = defaultGroup
}
activity.Record(ctx, dep.SettingProvider(), dep.ActivityClient(), types.EventMembershipUnsubscribe,
activity.Actor(g.UserID), activity.Extra(map[string]any{
"from_group": g.Amount, "to_group": revertTo,
}))
}
}

@ -0,0 +1,78 @@
package user
import (
"context"
"testing"
"time"
"github.com/cloudreve/Cloudreve/v4/application/dependency"
"github.com/cloudreve/Cloudreve/v4/ent/activityevent"
"github.com/cloudreve/Cloudreve/v4/ent/enttest"
"github.com/cloudreve/Cloudreve/v4/ent/usergrant"
"github.com/cloudreve/Cloudreve/v4/inventory/types"
"github.com/cloudreve/Cloudreve/v4/pkg/boolset"
"github.com/cloudreve/Cloudreve/v4/pkg/conf"
"github.com/cloudreve/Cloudreve/v4/pkg/logging"
"github.com/cloudreve/Cloudreve/v4/pkg/setting"
"github.com/stretchr/testify/require"
)
type cronSettingProvider struct {
setting.Provider
defaultGroup int
}
func (p cronSettingProvider) DefaultGroup(context.Context) int {
return p.defaultGroup
}
func (p cronSettingProvider) AuditLogEnabled(context.Context, int) bool {
return true
}
func TestGrantExpireCronRecordsUnsubscribe(t *testing.T) {
client := enttest.Open(t, "sqlite3", "file:"+t.Name()+"?mode=memory&cache=shared")
t.Cleanup(func() { require.NoError(t, client.Close()) })
ctx := context.Background()
logger := logging.NewConsoleLogger(logging.LevelError)
cfg, err := conf.NewIniConfigProvider(t.TempDir()+"/conf.ini", logger)
require.NoError(t, err)
base := client.Group.Create().SetName("base").SetPermissions(&boolset.BooleanSet{}).SaveX(ctx)
vip := client.Group.Create().SetName("vip").SetPermissions(&boolset.BooleanSet{}).SaveX(ctx)
u := client.User.Create().SetEmail("c@example.com").SetNick("c").SetGroup(vip).SaveX(ctx)
// Expired group grant, user still on the granted group -> reverted.
client.UserGrant.Create().SetUserID(u.ID).SetType(usergrant.TypeGroup).
SetAmount(int64(vip.ID)).SetPrevGroupID(base.ID).
SetExpiresAt(time.Now().Add(-time.Minute)).SaveX(ctx)
// Expired grant for a user who moved on independently -> not reverted,
// no event.
u2 := client.User.Create().SetEmail("c2@example.com").SetNick("c2").SetGroup(base).SaveX(ctx)
client.UserGrant.Create().SetUserID(u2.ID).SetType(usergrant.TypeGroup).
SetAmount(int64(vip.ID)).SetPrevGroupID(base.ID).
SetExpiresAt(time.Now().Add(-time.Minute)).SaveX(ctx)
dep := dependency.NewDependency(
dependency.WithDbClient(client),
dependency.WithConfigProvider(cfg),
dependency.WithLogger(logger),
dependency.WithSettingProvider(cronSettingProvider{defaultGroup: base.ID}),
)
ctx = context.WithValue(ctx, dependency.DepCtx{}, dep)
reverted, err := dep.VasClient().ExpireGrants(ctx, base.ID)
require.NoError(t, err)
require.Len(t, reverted, 1)
require.Equal(t, u.ID, reverted[0].UserID)
require.Equal(t, base.ID, client.User.GetX(ctx, u.ID).GroupUsers)
require.Equal(t, base.ID, client.User.GetX(ctx, u2.ID).GroupUsers)
recordMembershipUnsubscribes(ctx, dep, reverted, base.ID)
ev := client.ActivityEvent.Query().
Where(activityevent.TypeEQ(types.EventMembershipUnsubscribe)).OnlyX(ctx)
require.Equal(t, u.ID, ev.ActorID)
require.Equal(t, float64(vip.ID), ev.Extra["from_group"])
require.Equal(t, float64(base.ID), ev.Extra["to_group"])
}
Loading…
Cancel
Save