Merge pull request #193 from Dvorinka/feat/event-coverage

feat: quota-notify and membership-unsubscribe activity events
pull/3587/head
Tomáš Dvořák 2 weeks ago committed by GitHub
commit 8060eee426
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194

@ -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