From 521bfda63c534db94de59b6e663d2ef7adaf6fb1 Mon Sep 17 00:00:00 2001 From: Tomas Dvorak Date: Sun, 20 Sep 2026 04:47:06 +0200 Subject: [PATCH] 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> --- ROADMAP.md | 3 +- inventory/vas.go | 24 +++--- inventory/vas_test.go | 5 +- pkg/filemanager/fs/dbfs/manage.go | 8 ++ pkg/filemanager/fs/dbfs/upload.go | 3 + pkg/filemanager/fs/dbfs/upload_dedup_test.go | 30 +++++++- pkg/filemanager/fs/dbfs/validator.go | 5 ++ service/user/vas_cron.go | 25 ++++++- service/user/vas_cron_test.go | 78 ++++++++++++++++++++ 9 files changed, 165 insertions(+), 16 deletions(-) create mode 100644 service/user/vas_cron_test.go diff --git a/ROADMAP.md b/ROADMAP.md index cf1b3928..bfa76cd0 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -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 diff --git a/inventory/vas.go b/inventory/vas.go index d875e664..3ccd861f 100644 --- a/inventory/vas.go +++ b/inventory/vas.go @@ -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) { diff --git a/inventory/vas_test.go b/inventory/vas_test.go index 837a09b6..1c584717 100644 --- a/inventory/vas_test.go +++ b/inventory/vas_test.go @@ -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)) } diff --git a/pkg/filemanager/fs/dbfs/manage.go b/pkg/filemanager/fs/dbfs/manage.go index ce322e9b..19c6b360 100644 --- a/pkg/filemanager/fs/dbfs/manage.go +++ b/pkg/filemanager/fs/dbfs/manage.go @@ -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 } diff --git a/pkg/filemanager/fs/dbfs/upload.go b/pkg/filemanager/fs/dbfs/upload.go index 32fdb7af..22ae21d1 100644 --- a/pkg/filemanager/fs/dbfs/upload.go +++ b/pkg/filemanager/fs/dbfs/upload.go @@ -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) diff --git a/pkg/filemanager/fs/dbfs/upload_dedup_test.go b/pkg/filemanager/fs/dbfs/upload_dedup_test.go index aaad4161..4afec6d4 100644 --- a/pkg/filemanager/fs/dbfs/upload_dedup_test.go +++ b/pkg/filemanager/fs/dbfs/upload_dedup_test.go @@ -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"]) +} diff --git a/pkg/filemanager/fs/dbfs/validator.go b/pkg/filemanager/fs/dbfs/validator.go index f38b57bf..dbe9582e 100644 --- a/pkg/filemanager/fs/dbfs/validator.go +++ b/pkg/filemanager/fs/dbfs/validator.go @@ -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 diff --git a/service/user/vas_cron.go b/service/user/vas_cron.go index 8ad0cf6a..45e87af3 100644 --- a/service/user/vas_cron.go +++ b/service/user/vas_cron.go @@ -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, + })) + } +} diff --git a/service/user/vas_cron_test.go b/service/user/vas_cron_test.go new file mode 100644 index 00000000..b09018f0 --- /dev/null +++ b/service/user/vas_cron_test.go @@ -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"]) +}