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"]) +}