feat: hash-based duplicate detection and instant upload (#191)

* feat: hash-based duplicate detection and instant upload

Clients send a client-computed sha256 with the upload session request;
when upload_dedup_scope is enabled and a completed entity with the same
hash and size exists within scope, the file is materialized by linking
the existing entity (refcount++) instead of transferring bytes.

- entity.hash column + (hash,size) index; stored on entity creation
- FileClient.FindEntityByHash: completed version entities only
  (reference_count > 0, no upload session), owner-scoped unless global
- CreateFile LinkedEntityID branch: attach entity, bump refcount, set
  primary entity + logical size, charge owner quota via storage diff
- DBFS.PrepareUpload rapid path: dedup lookup after validation, commits
  file + emits upload event, skips session metadata/lock/credentials
- manager short-circuits rapid sessions before driver token + KV store
- postProcessRapidUpload queues media-meta / FTS tasks per file
- upload_dedup_scope setting (off/owner/global, default owner) surfaced
  in admin FileSystem parameters and basic site config
- Frontend streams SHA-256 via hash-wasm (no whole-file buffering),
  skips hashing for encrypted policies, short-circuits on rapid_uploaded
- Tests: FindEntityByHash scoping/isolation, linked CreateFile, and
  full PrepareUpload rapid path incl. scope=off

Implements upstream cloudreve/cloudreve#3044 (dedup detection);
addresses fork issues #66/#15.

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

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* docs: roadmap — dedup landed (#191), multi-file share scoped

---------

Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
pull/3587/head
Tomáš Dvořák 2 weeks ago committed by GitHub
parent 5e7f938f41
commit cf688634c4
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194

@ -174,6 +174,8 @@ Order = user-visible value first; each ships with backend + UI + tests.
- [x] File/dir ACL entity — `(subject_type ∈ user|group|anonymous|everyone, subject_id)` → R/C/U/D bitmask; Permissions dialog under More actions; enforced in share-navigator capability checks; group bit 15 gate (#181, fixes #3517)
- [x] Default shares — `setting.default_symbolics` + group `default_pinned` chip-input of share IDs; materialize as share-shortcut entries on fs init (#180)
- [x] Paid shares — `share.price_points` + `share_purchase` (buyer debit → owner income at `share_score_rate`) + resume ticket; download/thumb gated in share navigator, listing stays visible; `share_sell` (bit 29) gates price-setting, `share_free` (bit 8) bypasses paywall; `PaidShareGate` UI + restore via `purchase_ticket`
- [x] Hash dedup / instant upload — `entity.hash` + `(hash,size)` index; `FindEntityByHash` (completed version entities, owner scope default / `global` opt-in); `PrepareUpload` links existing entity (refcount++) → `rapid_uploaded` session, no transfer; media-meta/FTS queued per file; frontend streams SHA-256 via `hash-wasm`, skips encrypted policies; `upload_dedup_scope` admin setting (#191, fixes #3044)
- [ ] Multi-file share — one share covering N selected files (upstream #3032): share↔files M:N, navigator union, share page listing, create dialog multi-pick
2. **Storage policy advanced** — multiple policies per group, per-directory binding, load-balancer policy, file migration (fixes #3518, #2961, #2262). See §1.3a.
- [x] PR #175 — resumable admin relocation task (entities or whole-policy scope), encryption-aware re-wrap, admin UI + per-policy migrate action (#9, #125, #136)
- [x] Group→policies M:N (`allowed_policies` edge, empty = legacy single) + group-editor multi-select; per-directory `sys:preferred_policy` metadata marker with nearest-ancestor precedence (invalid marker cuts inheritance); user `preferred_policy` setting applied in own tree only; `load_balance` policy type with weighted children resolved before drivers (#182, fixes #2961)

@ -36,6 +36,8 @@ type Entity struct {
Size int64 `json:"size,omitempty"`
// ReferenceCount holds the value of the "reference_count" field.
ReferenceCount int `json:"reference_count,omitempty"`
// Hash holds the value of the "hash" field.
Hash string `json:"hash,omitempty"`
// StoragePolicyEntities holds the value of the "storage_policy_entities" field.
StoragePolicyEntities int `json:"storage_policy_entities,omitempty"`
// CreatedBy holds the value of the "created_by" field.
@ -109,7 +111,7 @@ func (*Entity) scanValues(columns []string) ([]any, error) {
values[i] = new([]byte)
case entity.FieldID, entity.FieldType, entity.FieldSize, entity.FieldReferenceCount, entity.FieldStoragePolicyEntities, entity.FieldCreatedBy:
values[i] = new(sql.NullInt64)
case entity.FieldSource:
case entity.FieldSource, entity.FieldHash:
values[i] = new(sql.NullString)
case entity.FieldCreatedAt, entity.FieldUpdatedAt, entity.FieldDeletedAt:
values[i] = new(sql.NullTime)
@ -177,6 +179,12 @@ func (e *Entity) assignValues(columns []string, values []any) error {
} else if value.Valid {
e.ReferenceCount = int(value.Int64)
}
case entity.FieldHash:
if value, ok := values[i].(*sql.NullString); !ok {
return fmt.Errorf("unexpected type %T for field hash", values[i])
} else if value.Valid {
e.Hash = value.String
}
case entity.FieldStoragePolicyEntities:
if value, ok := values[i].(*sql.NullInt64); !ok {
return fmt.Errorf("unexpected type %T for field storage_policy_entities", values[i])
@ -278,6 +286,9 @@ func (e *Entity) String() string {
builder.WriteString("reference_count=")
builder.WriteString(fmt.Sprintf("%v", e.ReferenceCount))
builder.WriteString(", ")
builder.WriteString("hash=")
builder.WriteString(e.Hash)
builder.WriteString(", ")
builder.WriteString("storage_policy_entities=")
builder.WriteString(fmt.Sprintf("%v", e.StoragePolicyEntities))
builder.WriteString(", ")

@ -29,6 +29,8 @@ const (
FieldSize = "size"
// FieldReferenceCount holds the string denoting the reference_count field in the database.
FieldReferenceCount = "reference_count"
// FieldHash holds the string denoting the hash field in the database.
FieldHash = "hash"
// FieldStoragePolicyEntities holds the string denoting the storage_policy_entities field in the database.
FieldStoragePolicyEntities = "storage_policy_entities"
// FieldCreatedBy holds the string denoting the created_by field in the database.
@ -76,6 +78,7 @@ var Columns = []string{
FieldSource,
FieldSize,
FieldReferenceCount,
FieldHash,
FieldStoragePolicyEntities,
FieldCreatedBy,
FieldUploadSessionID,
@ -114,6 +117,8 @@ var (
UpdateDefaultUpdatedAt func() time.Time
// DefaultReferenceCount holds the default value on creation for the "reference_count" field.
DefaultReferenceCount int
// HashValidator is a validator for the "hash" field. It is called by the builders before save.
HashValidator func(string) error
)
// OrderOption defines the ordering options for the Entity queries.
@ -159,6 +164,11 @@ func ByReferenceCount(opts ...sql.OrderTermOption) OrderOption {
return sql.OrderByField(FieldReferenceCount, opts...).ToFunc()
}
// ByHash orders the results by the hash field.
func ByHash(opts ...sql.OrderTermOption) OrderOption {
return sql.OrderByField(FieldHash, opts...).ToFunc()
}
// ByStoragePolicyEntities orders the results by the storage_policy_entities field.
func ByStoragePolicyEntities(opts ...sql.OrderTermOption) OrderOption {
return sql.OrderByField(FieldStoragePolicyEntities, opts...).ToFunc()

@ -91,6 +91,11 @@ func ReferenceCount(v int) predicate.Entity {
return predicate.Entity(sql.FieldEQ(FieldReferenceCount, v))
}
// Hash applies equality check predicate on the "hash" field. It's identical to HashEQ.
func Hash(v string) predicate.Entity {
return predicate.Entity(sql.FieldEQ(FieldHash, v))
}
// StoragePolicyEntities applies equality check predicate on the "storage_policy_entities" field. It's identical to StoragePolicyEntitiesEQ.
func StoragePolicyEntities(v int) predicate.Entity {
return predicate.Entity(sql.FieldEQ(FieldStoragePolicyEntities, v))
@ -421,6 +426,81 @@ func ReferenceCountLTE(v int) predicate.Entity {
return predicate.Entity(sql.FieldLTE(FieldReferenceCount, v))
}
// HashEQ applies the EQ predicate on the "hash" field.
func HashEQ(v string) predicate.Entity {
return predicate.Entity(sql.FieldEQ(FieldHash, v))
}
// HashNEQ applies the NEQ predicate on the "hash" field.
func HashNEQ(v string) predicate.Entity {
return predicate.Entity(sql.FieldNEQ(FieldHash, v))
}
// HashIn applies the In predicate on the "hash" field.
func HashIn(vs ...string) predicate.Entity {
return predicate.Entity(sql.FieldIn(FieldHash, vs...))
}
// HashNotIn applies the NotIn predicate on the "hash" field.
func HashNotIn(vs ...string) predicate.Entity {
return predicate.Entity(sql.FieldNotIn(FieldHash, vs...))
}
// HashGT applies the GT predicate on the "hash" field.
func HashGT(v string) predicate.Entity {
return predicate.Entity(sql.FieldGT(FieldHash, v))
}
// HashGTE applies the GTE predicate on the "hash" field.
func HashGTE(v string) predicate.Entity {
return predicate.Entity(sql.FieldGTE(FieldHash, v))
}
// HashLT applies the LT predicate on the "hash" field.
func HashLT(v string) predicate.Entity {
return predicate.Entity(sql.FieldLT(FieldHash, v))
}
// HashLTE applies the LTE predicate on the "hash" field.
func HashLTE(v string) predicate.Entity {
return predicate.Entity(sql.FieldLTE(FieldHash, v))
}
// HashContains applies the Contains predicate on the "hash" field.
func HashContains(v string) predicate.Entity {
return predicate.Entity(sql.FieldContains(FieldHash, v))
}
// HashHasPrefix applies the HasPrefix predicate on the "hash" field.
func HashHasPrefix(v string) predicate.Entity {
return predicate.Entity(sql.FieldHasPrefix(FieldHash, v))
}
// HashHasSuffix applies the HasSuffix predicate on the "hash" field.
func HashHasSuffix(v string) predicate.Entity {
return predicate.Entity(sql.FieldHasSuffix(FieldHash, v))
}
// HashIsNil applies the IsNil predicate on the "hash" field.
func HashIsNil() predicate.Entity {
return predicate.Entity(sql.FieldIsNull(FieldHash))
}
// HashNotNil applies the NotNil predicate on the "hash" field.
func HashNotNil() predicate.Entity {
return predicate.Entity(sql.FieldNotNull(FieldHash))
}
// HashEqualFold applies the EqualFold predicate on the "hash" field.
func HashEqualFold(v string) predicate.Entity {
return predicate.Entity(sql.FieldEqualFold(FieldHash, v))
}
// HashContainsFold applies the ContainsFold predicate on the "hash" field.
func HashContainsFold(v string) predicate.Entity {
return predicate.Entity(sql.FieldContainsFold(FieldHash, v))
}
// StoragePolicyEntitiesEQ applies the EQ predicate on the "storage_policy_entities" field.
func StoragePolicyEntitiesEQ(v int) predicate.Entity {
return predicate.Entity(sql.FieldEQ(FieldStoragePolicyEntities, v))

@ -101,6 +101,20 @@ func (ec *EntityCreate) SetNillableReferenceCount(i *int) *EntityCreate {
return ec
}
// SetHash sets the "hash" field.
func (ec *EntityCreate) SetHash(s string) *EntityCreate {
ec.mutation.SetHash(s)
return ec
}
// SetNillableHash sets the "hash" field if the given value is not nil.
func (ec *EntityCreate) SetNillableHash(s *string) *EntityCreate {
if s != nil {
ec.SetHash(*s)
}
return ec
}
// SetStoragePolicyEntities sets the "storage_policy_entities" field.
func (ec *EntityCreate) SetStoragePolicyEntities(i int) *EntityCreate {
ec.mutation.SetStoragePolicyEntities(i)
@ -264,6 +278,11 @@ func (ec *EntityCreate) check() error {
if _, ok := ec.mutation.ReferenceCount(); !ok {
return &ValidationError{Name: "reference_count", err: errors.New(`ent: missing required field "Entity.reference_count"`)}
}
if v, ok := ec.mutation.Hash(); ok {
if err := entity.HashValidator(v); err != nil {
return &ValidationError{Name: "hash", err: fmt.Errorf(`ent: validator failed for field "Entity.hash": %w`, err)}
}
}
if _, ok := ec.mutation.StoragePolicyEntities(); !ok {
return &ValidationError{Name: "storage_policy_entities", err: errors.New(`ent: missing required field "Entity.storage_policy_entities"`)}
}
@ -332,6 +351,10 @@ func (ec *EntityCreate) createSpec() (*Entity, *sqlgraph.CreateSpec) {
_spec.SetField(entity.FieldReferenceCount, field.TypeInt, value)
_node.ReferenceCount = value
}
if value, ok := ec.mutation.Hash(); ok {
_spec.SetField(entity.FieldHash, field.TypeString, value)
_node.Hash = value
}
if value, ok := ec.mutation.UploadSessionID(); ok {
_spec.SetField(entity.FieldUploadSessionID, field.TypeUUID, value)
_node.UploadSessionID = &value
@ -538,6 +561,24 @@ func (u *EntityUpsert) AddReferenceCount(v int) *EntityUpsert {
return u
}
// SetHash sets the "hash" field.
func (u *EntityUpsert) SetHash(v string) *EntityUpsert {
u.Set(entity.FieldHash, v)
return u
}
// UpdateHash sets the "hash" field to the value that was provided on create.
func (u *EntityUpsert) UpdateHash() *EntityUpsert {
u.SetExcluded(entity.FieldHash)
return u
}
// ClearHash clears the value of the "hash" field.
func (u *EntityUpsert) ClearHash() *EntityUpsert {
u.SetNull(entity.FieldHash)
return u
}
// SetStoragePolicyEntities sets the "storage_policy_entities" field.
func (u *EntityUpsert) SetStoragePolicyEntities(v int) *EntityUpsert {
u.Set(entity.FieldStoragePolicyEntities, v)
@ -761,6 +802,27 @@ func (u *EntityUpsertOne) UpdateReferenceCount() *EntityUpsertOne {
})
}
// SetHash sets the "hash" field.
func (u *EntityUpsertOne) SetHash(v string) *EntityUpsertOne {
return u.Update(func(s *EntityUpsert) {
s.SetHash(v)
})
}
// UpdateHash sets the "hash" field to the value that was provided on create.
func (u *EntityUpsertOne) UpdateHash() *EntityUpsertOne {
return u.Update(func(s *EntityUpsert) {
s.UpdateHash()
})
}
// ClearHash clears the value of the "hash" field.
func (u *EntityUpsertOne) ClearHash() *EntityUpsertOne {
return u.Update(func(s *EntityUpsert) {
s.ClearHash()
})
}
// SetStoragePolicyEntities sets the "storage_policy_entities" field.
func (u *EntityUpsertOne) SetStoragePolicyEntities(v int) *EntityUpsertOne {
return u.Update(func(s *EntityUpsert) {
@ -1166,6 +1228,27 @@ func (u *EntityUpsertBulk) UpdateReferenceCount() *EntityUpsertBulk {
})
}
// SetHash sets the "hash" field.
func (u *EntityUpsertBulk) SetHash(v string) *EntityUpsertBulk {
return u.Update(func(s *EntityUpsert) {
s.SetHash(v)
})
}
// UpdateHash sets the "hash" field to the value that was provided on create.
func (u *EntityUpsertBulk) UpdateHash() *EntityUpsertBulk {
return u.Update(func(s *EntityUpsert) {
s.UpdateHash()
})
}
// ClearHash clears the value of the "hash" field.
func (u *EntityUpsertBulk) ClearHash() *EntityUpsertBulk {
return u.Update(func(s *EntityUpsert) {
s.ClearHash()
})
}
// SetStoragePolicyEntities sets the "storage_policy_entities" field.
func (u *EntityUpsertBulk) SetStoragePolicyEntities(v int) *EntityUpsertBulk {
return u.Update(func(s *EntityUpsert) {

@ -136,6 +136,26 @@ func (eu *EntityUpdate) AddReferenceCount(i int) *EntityUpdate {
return eu
}
// SetHash sets the "hash" field.
func (eu *EntityUpdate) SetHash(s string) *EntityUpdate {
eu.mutation.SetHash(s)
return eu
}
// SetNillableHash sets the "hash" field if the given value is not nil.
func (eu *EntityUpdate) SetNillableHash(s *string) *EntityUpdate {
if s != nil {
eu.SetHash(*s)
}
return eu
}
// ClearHash clears the value of the "hash" field.
func (eu *EntityUpdate) ClearHash() *EntityUpdate {
eu.mutation.ClearHash()
return eu
}
// SetStoragePolicyEntities sets the "storage_policy_entities" field.
func (eu *EntityUpdate) SetStoragePolicyEntities(i int) *EntityUpdate {
eu.mutation.SetStoragePolicyEntities(i)
@ -329,6 +349,11 @@ func (eu *EntityUpdate) defaults() error {
// check runs all checks and user-defined validators on the builder.
func (eu *EntityUpdate) check() error {
if v, ok := eu.mutation.Hash(); ok {
if err := entity.HashValidator(v); err != nil {
return &ValidationError{Name: "hash", err: fmt.Errorf(`ent: validator failed for field "Entity.hash": %w`, err)}
}
}
if _, ok := eu.mutation.StoragePolicyID(); eu.mutation.StoragePolicyCleared() && !ok {
return errors.New(`ent: clearing a required unique edge "Entity.storage_policy"`)
}
@ -377,6 +402,12 @@ func (eu *EntityUpdate) sqlSave(ctx context.Context) (n int, err error) {
if value, ok := eu.mutation.AddedReferenceCount(); ok {
_spec.AddField(entity.FieldReferenceCount, field.TypeInt, value)
}
if value, ok := eu.mutation.Hash(); ok {
_spec.SetField(entity.FieldHash, field.TypeString, value)
}
if eu.mutation.HashCleared() {
_spec.ClearField(entity.FieldHash, field.TypeString)
}
if value, ok := eu.mutation.UploadSessionID(); ok {
_spec.SetField(entity.FieldUploadSessionID, field.TypeUUID, value)
}
@ -615,6 +646,26 @@ func (euo *EntityUpdateOne) AddReferenceCount(i int) *EntityUpdateOne {
return euo
}
// SetHash sets the "hash" field.
func (euo *EntityUpdateOne) SetHash(s string) *EntityUpdateOne {
euo.mutation.SetHash(s)
return euo
}
// SetNillableHash sets the "hash" field if the given value is not nil.
func (euo *EntityUpdateOne) SetNillableHash(s *string) *EntityUpdateOne {
if s != nil {
euo.SetHash(*s)
}
return euo
}
// ClearHash clears the value of the "hash" field.
func (euo *EntityUpdateOne) ClearHash() *EntityUpdateOne {
euo.mutation.ClearHash()
return euo
}
// SetStoragePolicyEntities sets the "storage_policy_entities" field.
func (euo *EntityUpdateOne) SetStoragePolicyEntities(i int) *EntityUpdateOne {
euo.mutation.SetStoragePolicyEntities(i)
@ -821,6 +872,11 @@ func (euo *EntityUpdateOne) defaults() error {
// check runs all checks and user-defined validators on the builder.
func (euo *EntityUpdateOne) check() error {
if v, ok := euo.mutation.Hash(); ok {
if err := entity.HashValidator(v); err != nil {
return &ValidationError{Name: "hash", err: fmt.Errorf(`ent: validator failed for field "Entity.hash": %w`, err)}
}
}
if _, ok := euo.mutation.StoragePolicyID(); euo.mutation.StoragePolicyCleared() && !ok {
return errors.New(`ent: clearing a required unique edge "Entity.storage_policy"`)
}
@ -886,6 +942,12 @@ func (euo *EntityUpdateOne) sqlSave(ctx context.Context) (_node *Entity, err err
if value, ok := euo.mutation.AddedReferenceCount(); ok {
_spec.AddField(entity.FieldReferenceCount, field.TypeInt, value)
}
if value, ok := euo.mutation.Hash(); ok {
_spec.SetField(entity.FieldHash, field.TypeString, value)
}
if euo.mutation.HashCleared() {
_spec.ClearField(entity.FieldHash, field.TypeString)
}
if value, ok := euo.mutation.UploadSessionID(); ok {
_spec.SetField(entity.FieldUploadSessionID, field.TypeUUID, value)
}

File diff suppressed because one or more lines are too long

@ -222,6 +222,7 @@ var (
{Name: "source", Type: field.TypeString, Size: 2147483647},
{Name: "size", Type: field.TypeInt64},
{Name: "reference_count", Type: field.TypeInt, Default: 1},
{Name: "hash", Type: field.TypeString, Nullable: true, Size: 64},
{Name: "upload_session_id", Type: field.TypeUUID, Nullable: true},
{Name: "recycle_options", Type: field.TypeJSON, Nullable: true},
{Name: "storage_policy_entities", Type: field.TypeInt},
@ -235,17 +236,24 @@ var (
ForeignKeys: []*schema.ForeignKey{
{
Symbol: "entities_storage_policies_entities",
Columns: []*schema.Column{EntitiesColumns[10]},
Columns: []*schema.Column{EntitiesColumns[11]},
RefColumns: []*schema.Column{StoragePoliciesColumns[0]},
OnDelete: schema.NoAction,
},
{
Symbol: "entities_users_entities",
Columns: []*schema.Column{EntitiesColumns[11]},
Columns: []*schema.Column{EntitiesColumns[12]},
RefColumns: []*schema.Column{UsersColumns[0]},
OnDelete: schema.SetNull,
},
},
Indexes: []*schema.Index{
{
Name: "entity_hash_size",
Unique: false,
Columns: []*schema.Column{EntitiesColumns[8], EntitiesColumns[6]},
},
},
}
// FilesColumns holds the columns for the "files" table.
FilesColumns = []*schema.Column{

@ -5537,6 +5537,7 @@ type EntityMutation struct {
addsize *int64
reference_count *int
addreference_count *int
hash *string
upload_session_id *uuid.UUID
props **types.EntityProps
clearedFields map[string]struct{}
@ -5975,6 +5976,55 @@ func (m *EntityMutation) ResetReferenceCount() {
m.addreference_count = nil
}
// SetHash sets the "hash" field.
func (m *EntityMutation) SetHash(s string) {
m.hash = &s
}
// Hash returns the value of the "hash" field in the mutation.
func (m *EntityMutation) Hash() (r string, exists bool) {
v := m.hash
if v == nil {
return
}
return *v, true
}
// OldHash returns the old "hash" field's value of the Entity entity.
// If the Entity object wasn't provided to the builder, the object is fetched from the database.
// An error is returned if the mutation operation is not UpdateOne, or the database query fails.
func (m *EntityMutation) OldHash(ctx context.Context) (v string, err error) {
if !m.op.Is(OpUpdateOne) {
return v, errors.New("OldHash is only allowed on UpdateOne operations")
}
if m.id == nil || m.oldValue == nil {
return v, errors.New("OldHash requires an ID field in the mutation")
}
oldValue, err := m.oldValue(ctx)
if err != nil {
return v, fmt.Errorf("querying old value for OldHash: %w", err)
}
return oldValue.Hash, nil
}
// ClearHash clears the value of the "hash" field.
func (m *EntityMutation) ClearHash() {
m.hash = nil
m.clearedFields[entity.FieldHash] = struct{}{}
}
// HashCleared returns if the "hash" field was cleared in this mutation.
func (m *EntityMutation) HashCleared() bool {
_, ok := m.clearedFields[entity.FieldHash]
return ok
}
// ResetHash resets all changes to the "hash" field.
func (m *EntityMutation) ResetHash() {
m.hash = nil
delete(m.clearedFields, entity.FieldHash)
}
// SetStoragePolicyEntities sets the "storage_policy_entities" field.
func (m *EntityMutation) SetStoragePolicyEntities(i int) {
m.storage_policy = &i
@ -6326,7 +6376,7 @@ func (m *EntityMutation) Type() string {
// order to get all numeric fields that were incremented/decremented, call
// AddedFields().
func (m *EntityMutation) Fields() []string {
fields := make([]string, 0, 11)
fields := make([]string, 0, 12)
if m.created_at != nil {
fields = append(fields, entity.FieldCreatedAt)
}
@ -6348,6 +6398,9 @@ func (m *EntityMutation) Fields() []string {
if m.reference_count != nil {
fields = append(fields, entity.FieldReferenceCount)
}
if m.hash != nil {
fields = append(fields, entity.FieldHash)
}
if m.storage_policy != nil {
fields = append(fields, entity.FieldStoragePolicyEntities)
}
@ -6382,6 +6435,8 @@ func (m *EntityMutation) Field(name string) (ent.Value, bool) {
return m.Size()
case entity.FieldReferenceCount:
return m.ReferenceCount()
case entity.FieldHash:
return m.Hash()
case entity.FieldStoragePolicyEntities:
return m.StoragePolicyEntities()
case entity.FieldCreatedBy:
@ -6413,6 +6468,8 @@ func (m *EntityMutation) OldField(ctx context.Context, name string) (ent.Value,
return m.OldSize(ctx)
case entity.FieldReferenceCount:
return m.OldReferenceCount(ctx)
case entity.FieldHash:
return m.OldHash(ctx)
case entity.FieldStoragePolicyEntities:
return m.OldStoragePolicyEntities(ctx)
case entity.FieldCreatedBy:
@ -6479,6 +6536,13 @@ func (m *EntityMutation) SetField(name string, value ent.Value) error {
}
m.SetReferenceCount(v)
return nil
case entity.FieldHash:
v, ok := value.(string)
if !ok {
return fmt.Errorf("unexpected type %T for field %s", value, name)
}
m.SetHash(v)
return nil
case entity.FieldStoragePolicyEntities:
v, ok := value.(int)
if !ok {
@ -6579,6 +6643,9 @@ func (m *EntityMutation) ClearedFields() []string {
if m.FieldCleared(entity.FieldDeletedAt) {
fields = append(fields, entity.FieldDeletedAt)
}
if m.FieldCleared(entity.FieldHash) {
fields = append(fields, entity.FieldHash)
}
if m.FieldCleared(entity.FieldCreatedBy) {
fields = append(fields, entity.FieldCreatedBy)
}
@ -6605,6 +6672,9 @@ func (m *EntityMutation) ClearField(name string) error {
case entity.FieldDeletedAt:
m.ClearDeletedAt()
return nil
case entity.FieldHash:
m.ClearHash()
return nil
case entity.FieldCreatedBy:
m.ClearCreatedBy()
return nil
@ -6643,6 +6713,9 @@ func (m *EntityMutation) ResetField(name string) error {
case entity.FieldReferenceCount:
m.ResetReferenceCount()
return nil
case entity.FieldHash:
m.ResetHash()
return nil
case entity.FieldStoragePolicyEntities:
m.ResetStoragePolicyEntities()
return nil

@ -180,6 +180,10 @@ func init() {
entityDescReferenceCount := entityFields[3].Descriptor()
// entity.DefaultReferenceCount holds the default value on creation for the reference_count field.
entity.DefaultReferenceCount = entityDescReferenceCount.Default.(int)
// entityDescHash is the schema descriptor for hash field.
entityDescHash := entityFields[4].Descriptor()
// entity.HashValidator is a validator for the "hash" field. It is called by the builders before save.
entity.HashValidator = entityDescHash.Validators[0].(func(string) error)
fileHooks := schema.File{}.Hooks()
file.Hooks[0] = fileHooks[0]
fileFields := schema.File{}.Fields()

@ -4,6 +4,7 @@ import (
"entgo.io/ent"
"entgo.io/ent/schema/edge"
"entgo.io/ent/schema/field"
"entgo.io/ent/schema/index"
"github.com/cloudreve/Cloudreve/v4/inventory/types"
"github.com/gofrs/uuid"
)
@ -20,6 +21,12 @@ func (Entity) Fields() []ent.Field {
field.Text("source"),
field.Int64("size"),
field.Int("reference_count").Default(1),
// Content hash asserted by the uploader (sha256 hex). Used for
// duplicate detection / instant upload; not verified server-side
// because remote-policy blobs never pass through this server.
field.String("hash").
Optional().
MaxLen(64),
field.Int("storage_policy_entities"),
field.Int("created_by").Optional(),
field.UUID("upload_session_id", uuid.Must(uuid.NewV4())).
@ -53,3 +60,9 @@ func (Entity) Mixin() []ent.Mixin {
CommonMixin{},
}
}
func (Entity) Indexes() []ent.Index {
return []ent.Index{
index.Fields("hash", "size"),
}
}

@ -38,6 +38,7 @@
"axios": "^1.12.2",
"dayjs": "^1.11.10",
"fuse.js": "^7.0.0",
"hash-wasm": "4.12.0",
"heic-to": "^1.1.14",
"hls.js": "^1.6.2",
"i18next": "^23.7.11",

@ -241,6 +241,14 @@
"slaveAPIExpirationDes": "The signature validity period used by the master node when accessing the slave node API.",
"uploadSessionTimeout": "Upload session TTL (seconds)",
"uploadSessionDes": "In a valid upload session period, for supported storage policies, users can resume unfinished tasks. The maximum value that can be set is limited by the rules of different storage policy providers.",
"uploadDedupScope": "Duplicate detection scope",
"uploadDedupScopeDes": "When a client sends a content hash, the server looks for an identical completed file and materializes the upload instantly without transferring bytes.",
"uploadDedupScope_off": "Off",
"uploadDedupScope_offDes": "No duplicate detection; every upload transfers bytes.",
"uploadDedupScope_owner": "Same user",
"uploadDedupScope_ownerDes": "Only deduplicate against files the uploading user already owns.",
"uploadDedupScope_global": "Global",
"uploadDedupScope_globalDes": "Deduplicate against all completed files on this site.",
"archiveTimeout": "Server-side batch download session TTL (seconds)",
"advanceOptions": "Advanced options",
"emojiOptions": "Emoji options",

@ -241,6 +241,14 @@
"slaveAPIExpirationDes": "主机访问从机 API 时使用的签名有效期。",
"uploadSessionTimeout": "上传会话有效期 (秒)",
"uploadSessionDes": "在上传会话有效期内,对于支持的存储策略,用户可以断点续传未完成的任务。最大可设定的值受限于不同存储策略服务商的规则。",
"uploadDedupScope": "重复检测范围",
"uploadDedupScopeDes": "客户端发送内容哈希时,服务端会查找内容一致的已完成文件并直接完成上传(秒传),无需传输数据。",
"uploadDedupScope_off": "关闭",
"uploadDedupScope_offDes": "不进行重复检测,每次上传都传输数据。",
"uploadDedupScope_owner": "仅本人",
"uploadDedupScope_ownerDes": "仅对上传者本人已有的文件进行去重。",
"uploadDedupScope_global": "全站",
"uploadDedupScope_globalDes": "对全站所有已完成文件进行去重。",
"archiveTimeout": "服务端打包下载会话有效期 (秒)",
"advanceOptions": "高级设置",
"emojiOptions": "Emoji 选项",

@ -556,6 +556,8 @@ export interface UploadSessionRequest {
};
mime_type?: string;
encryption_supported?: EncryptionCipher[];
// sha256 hex of the file content — enables server-side dedup / instant upload
hash?: string;
}
export interface EncryptMetadata {
@ -582,6 +584,8 @@ export interface UploadCredential {
mime_type?: string;
upload_policy?: string;
encrypt_metadata?: EncryptMetadata;
// File was materialized from an existing identical entity — no bytes to send
rapid_uploaded?: boolean;
}
export interface DeleteUploadSessionService {

@ -35,6 +35,7 @@ export interface SiteConfig {
qq_connect_enabled?: boolean;
download_cdn_routes?: { name: string; url: string }[];
abuse_captcha?: boolean;
upload_dedup?: boolean;
allow_select_node?: boolean;
task_nodes?: { id: string; name: string }[];
logo?: string;

@ -101,6 +101,7 @@ const FileSystem = () => {
"explorer_category_document_query",
"archive_timeout",
"upload_session_timeout",
"upload_dedup_scope",
"slave_api_timeout",
"folder_props_timeout",
"chunk_retries",

@ -1,5 +1,5 @@
import { DeleteOutline } from "@mui/icons-material";
import { Box, FormControl, FormControlLabel, Link, Switch, Typography } from "@mui/material";
import { Box, FormControl, FormControlLabel, Link, ListItemText, Switch, Typography } from "@mui/material";
import { useSnackbar } from "notistack";
import { useContext, useState } from "react";
import { Trans, useTranslation } from "react-i18next";
@ -7,7 +7,8 @@ import { sendClearBlobUrlCache } from "../../../../api/api.ts";
import { useAppDispatch } from "../../../../redux/hooks.ts";
import { isTrueVal } from "../../../../session/utils.ts";
import { DefaultCloseAction } from "../../../Common/Snackbar/snackbar.tsx";
import { DenseFilledTextField, SecondaryButton } from "../../../Common/StyledComponents.tsx";
import { DenseFilledTextField, DenseSelect, SecondaryButton } from "../../../Common/StyledComponents.tsx";
import { SquareMenuItem } from "../../../FileManager/ContextMenu/ContextMenu.tsx";
import SettingForm from "../../../Pages/Setting/SettingForm.tsx";
import { NoMarginHelperText, SettingSection, SettingSectionContent } from "../../Settings/Settings.tsx";
import { SettingContext } from "../../Settings/SettingWrapper.tsx";
@ -64,6 +65,41 @@ const AdvancedOptionsSection = () => {
/>
<NoMarginHelperText>{t("settings.uploadSessionDes")}</NoMarginHelperText>
</SettingForm>
<SettingForm title={t("settings.uploadDedupScope")} lgWidth={5}>
<FormControl>
<DenseSelect
renderValue={(v) => (
<ListItemText
slotProps={{
primary: { variant: "body2" },
}}
>
{t(`settings.uploadDedupScope_${v}`)}
</ListItemText>
)}
onChange={(e) =>
setSettings({
upload_dedup_scope: e.target.value as string,
})
}
value={values.upload_dedup_scope ?? "owner"}
>
{["off", "owner", "global"].map((scope) => (
<SquareMenuItem key={scope} value={scope}>
<Box sx={{ display: "flex", flexDirection: "column" }}>
<Typography variant={"body2"} fontWeight={600}>
{t(`settings.uploadDedupScope_${scope}`)}
</Typography>
<Typography variant={"body2"} color={"textSecondary"}>
{t(`settings.uploadDedupScope_${scope}Des`)}
</Typography>
</Box>
</SquareMenuItem>
))}
</DenseSelect>
<NoMarginHelperText>{t("settings.uploadDedupScopeDes")}</NoMarginHelperText>
</FormControl>
</SettingForm>
<SettingForm title={t("settings.slaveAPIExpiration")} lgWidth={5}>
<DenseFilledTextField
type="number"

@ -1,6 +1,7 @@
// 所有 Uploader 的基类
import axios, { CanceledError, CancelTokenSource } from "axios";
import { EncryptionCipher, PolicyType } from "../../../../api/explorer.ts";
import { store } from "../../../../redux/store.ts";
import CrUri from "../../../../util/uri.ts";
import { createUploadSession, deleteUploadSession } from "../api";
import { UploaderError } from "../errors";
@ -139,6 +140,18 @@ export default abstract class Base {
if (cachedInfo == null) {
const crUri = new CrUri(this.task.dst);
crUri.join(this.task.name);
// Content hash enables server-side dedup / instant upload. Skipped for
// encrypted policies: the blob is ciphertext with a random IV, so a
// plaintext hash would never match — and storing it would leak content
// identity.
let hash: string | undefined;
if (store.getState().siteConfig.basic.config.upload_dedup && !this.task.policy.encryption) {
try {
hash = await utils.sha256File(this.task.file);
} catch (e) {
this.logger.warn("Failed to hash file for dedup, uploading without hash:", e);
}
}
this.task.session = await createUploadSession(
{
uri: crUri.toString(),
@ -149,6 +162,7 @@ export default abstract class Base {
entity_type: this.task.overwrite ? "version" : undefined,
encryption_supported:
this.task.policy.encryption && "crypto" in window ? [EncryptionCipher.aes256ctr] : undefined,
hash,
},
this.cancelToken.token,
);
@ -160,6 +174,14 @@ export default abstract class Base {
this.logger.info("Resume upload from cached ctx:", cachedInfo);
}
if (this.task.session?.rapid_uploaded) {
// Instant upload: server materialized the file from an identical
// existing entity — nothing to transfer.
this.logger.info("Rapid upload: identical content already exists on server");
this.transit(Status.finished);
return;
}
if (this.task.session?.encrypt_metadata && !this.task.policy?.relay) {
// Check browser support for encryption
if (!("crypto" in window)) {
@ -236,7 +258,7 @@ export default abstract class Base {
protected cancelUploadSession = (): Promise<void> => {
return new Promise<void>((resolve) => {
utils.removeResumeCtx(this.task, this.logger);
if (this.task.session) {
if (this.task.session?.session_id) {
setTimeout(() => {
deleteUploadSession(this.task.session!?.session_id, this.task.session!?.uri)
.catch((e) => {

@ -1,9 +1,26 @@
import { createSHA256 } from "hash-wasm";
import CrUri from "../../../../util/uri";
import { UploaderError, UploaderErrorName } from "../errors";
import Logger from "../logger";
import { Task } from "../types";
import { ChunkProgress } from "../uploader/chunk";
// sha256File streams the file through WASM sha256 without buffering it
// whole — used for server-side dedup / instant upload.
const hashChunkSize = 8 * 1024 * 1024;
export async function sha256File(file: Blob): Promise<string> {
const hasher = await createSHA256();
hasher.init();
for (let offset = 0; offset < file.size; offset += hashChunkSize) {
const buf = await file.slice(offset, offset + hashChunkSize).arrayBuffer();
hasher.update(new Uint8Array(buf));
}
if (file.size === 0) {
hasher.update(new Uint8Array(0));
}
return hasher.digest("hex");
}
// 文件分块
export function getChunks(file: Blob, chunkByteSize: number | undefined): Blob[] {
// 如果 chunkByteSize 比文件大或为0,则直接取文件的大小

@ -6594,6 +6594,11 @@ has-tostringtag@^1.0.2:
dependencies:
has-symbols "^1.0.3"
hash-wasm@4.12.0:
version "4.12.0"
resolved "https://registry.yarnpkg.com/hash-wasm/-/hash-wasm-4.12.0.tgz#f9f1a9f9121e027a9acbf6db5d59452ace1ef9bb"
integrity sha512-+/2B2rYLb48I/evdOIhP+K/DD2ca2fgBjp6O+GBEnCDk2e4rpeXIK8GvIyRPjTezgmWn9gmKwkQjjx6BtqDHVQ==
hasown@^2.0.0:
version "2.0.0"
resolved "https://registry.npmjs.org/hasown/-/hasown-2.0.0.tgz"

@ -147,6 +147,12 @@ type (
UploadSessionID uuid.UUID
Importing bool
EncryptMetadata *types.EncryptMetadata
// Hash is the uploader-asserted content hash (sha256 hex) stored
// on the entity for later dedup lookups.
Hash string
// LinkedEntityID, when non-zero, makes CreateFile attach an
// existing entity (instant upload) instead of creating a new one.
LinkedEntityID int
}
RelocateEntityParameter struct {
@ -251,6 +257,11 @@ type FileClient interface {
Update(ctx context.Context, file *ent.File) (*ent.File, error)
// ListEntities lists entities
ListEntities(ctx context.Context, args *ListEntityParameters) (*ListEntityResult, error)
// FindEntityByHash returns the newest completed version entity matching
// the uploader-asserted content hash and size. When global is false the
// match is restricted to entities created by creatorID so one user
// cannot probe another user's blobs by hash.
FindEntityByHash(ctx context.Context, hash string, size int64, creatorID int, global bool) (*ent.Entity, error)
// UpdateProps updates props of a file
UpdateProps(ctx context.Context, file *ent.File, props *types.FileProps) (*ent.File, error)
// UpdateModifiedAt updates modified at of a file
@ -863,7 +874,22 @@ func (f *fileClient) CreateFile(ctx context.Context, root *ent.File, args *Creat
// Create default primary file entity if needed
var storageDiff StorageDiff
if args.EntityParameters != nil {
if args.EntityParameters != nil && args.EntityParameters.LinkedEntityID != 0 {
// Instant upload: attach an already-stored entity and bump its
// refcount. The file still consumes the owner's logical quota.
linked, err := f.client.Entity.Get(ctx, args.EntityParameters.LinkedEntityID)
if err != nil {
return nil, nil, nil, fmt.Errorf("failed to get linked entity: %v", err)
}
if err := f.client.Entity.UpdateOne(linked).AddFile(newFile).AddReferenceCount(1).Exec(ctx); err != nil {
return nil, nil, nil, fmt.Errorf("failed to link entity: %v", err)
}
if err := f.client.File.UpdateOne(newFile).SetPrimaryEntity(linked.ID).SetSize(linked.Size).Exec(ctx); err != nil {
return nil, nil, nil, fmt.Errorf("failed to set primary entity: %v", err)
}
defaultEntity = linked
storageDiff = map[int]int64{newFile.OwnerID: linked.Size}
} else if args.EntityParameters != nil {
args.EntityParameters.OwnerID = root.OwnerID
args.EntityParameters.StoragePolicyID = args.StoragePolicyID
defaultEntity, storageDiff, err = f.CreateEntity(ctx, newFile, args.EntityParameters)
@ -997,6 +1023,10 @@ func (f *fileClient) CreateEntity(ctx context.Context, file *ent.File, args *Ent
SetSize(args.Size).
SetStoragePolicyID(args.StoragePolicyID)
if args.Hash != "" {
stm.SetHash(args.Hash)
}
if opt != nil {
stm.SetProps(opt)
}
@ -1276,6 +1306,26 @@ func ParseReferenceCountFilter(expr string) (*ReferenceCountFilter, error) {
return &ReferenceCountFilter{Op: op, Value: v}, nil
}
func (f *fileClient) FindEntityByHash(ctx context.Context, hash string, size int64, creatorID int, global bool) (*ent.Entity, error) {
if hash == "" {
return nil, &ent.NotFoundError{}
}
query := f.client.Entity.Query().
Where(
entity.Hash(hash),
entity.Size(size),
entity.Type(int(types.EntityTypeVersion)),
entity.ReferenceCountGT(0),
entity.UploadSessionIDIsNil(),
)
if !global {
query = query.Where(entity.CreatedBy(creatorID))
}
return query.Order(ent.Desc(entity.FieldID)).First(ctx)
}
func (f *fileClient) ListEntities(ctx context.Context, args *ListEntityParameters) (*ListEntityResult, error) {
query := f.client.Entity.Query()
if args.EntityType != nil {

@ -0,0 +1,143 @@
package inventory
import (
"context"
"testing"
"github.com/cloudreve/Cloudreve/v4/ent"
"github.com/cloudreve/Cloudreve/v4/ent/enttest"
"github.com/cloudreve/Cloudreve/v4/ent/storagepolicy"
"github.com/cloudreve/Cloudreve/v4/inventory/types"
"github.com/cloudreve/Cloudreve/v4/pkg/boolset"
"github.com/cloudreve/Cloudreve/v4/pkg/conf"
"github.com/gofrs/uuid"
"github.com/stretchr/testify/require"
)
func dedupFixture(t *testing.T, client *ent.Client, suffix string) (*ent.User, *ent.File, *ent.StoragePolicy) {
ctx := context.Background()
group := client.Group.Create().SetName("g" + suffix).SetPermissions(&boolset.BooleanSet{}).SaveX(ctx)
u := client.User.Create().SetEmail("u" + suffix + "@example.com").SetNick("u" + suffix).SetGroup(group).SaveX(ctx)
root := client.File.Create().SetName(RootFolderName).SetType(int(types.FileTypeFolder)).SetOwner(u).SaveX(ctx)
p := client.StoragePolicy.Create().SetName("p" + suffix).SetType("local").
SetStatus(storagepolicy.StatusActive).SaveX(ctx)
return u, root, p
}
func mkHashedEntity(client *ent.Client, ctx context.Context, u *ent.User, p *ent.StoragePolicy, hash string, size int64) *ent.Entity {
return client.Entity.Create().
SetType(int(types.EntityTypeVersion)).
SetSource("cloudreve/" + hash).
SetSize(size).
SetHash(hash).
SetReferenceCount(1).
SetCreatedBy(u.ID).
SetStoragePolicyEntities(p.ID).
SaveX(ctx)
}
func TestFindEntityByHash(t *testing.T) {
ctx := context.Background()
client := enttest.Open(t, "sqlite3", "file:"+t.Name()+"?mode=memory&cache=shared")
t.Cleanup(func() { require.NoError(t, client.Close()) })
fc := NewFileClient(client, conf.SQLite3DB, nil)
u, _, p := dedupFixture(t, client, "a")
other, _, _ := dedupFixture(t, client, "b")
const hash = "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"
e := mkHashedEntity(client, ctx, u, p, hash, 100)
// Owner scope: match
found, err := fc.FindEntityByHash(ctx, hash, 100, u.ID, false)
require.NoError(t, err)
require.Equal(t, e.ID, found.ID)
// Owner scope: other user's entity must not leak
_, err = fc.FindEntityByHash(ctx, hash, 100, other.ID, false)
require.True(t, ent.IsNotFound(err))
// Global scope: cross-user match allowed
found, err = fc.FindEntityByHash(ctx, hash, 100, other.ID, true)
require.NoError(t, err)
require.Equal(t, e.ID, found.ID)
// Size mismatch: no match
_, err = fc.FindEntityByHash(ctx, hash, 101, u.ID, false)
require.True(t, ent.IsNotFound(err))
// Empty hash: no match
_, err = fc.FindEntityByHash(ctx, "", 100, u.ID, false)
require.True(t, ent.IsNotFound(err))
// Unfinished upload entity (session still attached): excluded
unfinished := client.Entity.Create().
SetType(int(types.EntityTypeVersion)).
SetSource("cloudreve/unfinished").
SetSize(50).
SetHash("aaaa").
SetReferenceCount(0).
SetCreatedBy(u.ID).
SetUploadSessionID(uuid.Must(uuid.NewV4())).
SetStoragePolicyEntities(p.ID).
SaveX(ctx)
_, err = fc.FindEntityByHash(ctx, "aaaa", 50, u.ID, false)
require.True(t, ent.IsNotFound(err))
_ = unfinished
// Stale entity (refcount 0, session cleared): excluded
client.Entity.Create().
SetType(int(types.EntityTypeVersion)).
SetSource("cloudreve/stale").
SetSize(60).
SetHash("bbbb").
SetReferenceCount(0).
SetCreatedBy(u.ID).
SetStoragePolicyEntities(p.ID).
SaveX(ctx)
_, err = fc.FindEntityByHash(ctx, "bbbb", 60, u.ID, false)
require.True(t, ent.IsNotFound(err))
}
func TestCreateFileLinkedEntity(t *testing.T) {
ctx := context.Background()
client := enttest.Open(t, "sqlite3", "file:"+t.Name()+"?mode=memory&cache=shared")
t.Cleanup(func() { require.NoError(t, client.Close()) })
fc := NewFileClient(client, conf.SQLite3DB, nil)
u, root, p := dedupFixture(t, client, "l")
const hash = "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"
existing := mkHashedEntity(client, ctx, u, p, hash, 512)
newFile, linked, diff, err := fc.CreateFile(ctx, root, &CreateFileParameters{
FileType: types.FileTypeFile,
Name: "copy.txt",
StoragePolicyID: p.ID,
EntityParameters: &EntityParameters{
LinkedEntityID: existing.ID,
},
})
require.NoError(t, err)
require.Equal(t, existing.ID, linked.ID)
// Refcount incremented
reloaded := client.Entity.GetX(ctx, existing.ID)
require.Equal(t, 2, reloaded.ReferenceCount)
// File wired to the shared entity (reload: primary/size are set via
// a follow-up update after the initial insert)
reloadedFile := client.File.GetX(ctx, newFile.ID)
require.Equal(t, int64(512), reloadedFile.Size)
require.Equal(t, existing.ID, reloadedFile.PrimaryEntity)
require.Equal(t, u.ID, reloadedFile.OwnerID)
// Logical quota charged to the new file's owner
require.Equal(t, int64(512), diff[u.ID])
// Entity->file edge attached
require.Len(t, reloaded.Edges.File, 0) // edges not eager-loaded
files := reloaded.QueryFile().AllX(ctx)
require.Len(t, files, 1)
require.Equal(t, newFile.ID, files[0].ID)
}

@ -570,6 +570,7 @@ var DefaultSettings = map[string]string{
"qq_connect_app_id": "",
"qq_connect_app_secret": "",
"qq_connect_register_enabled": "1",
"upload_dedup_scope": "owner",
"download_cdn_routes": "",
"email_filter_mode": "0",
"email_filter_list": "",

@ -705,8 +705,15 @@ func (f *DBFS) createFile(ctx context.Context, parent *File, name string, fileTy
UploadSessionID: uuid.FromStringOrNil(o.UploadRequest.Props.UploadSessionID),
Importing: o.UploadRequest.ImportFrom != nil,
EncryptMetadata: o.encryptMetadata,
Hash: o.UploadRequest.Props.Hash,
}
}
if o.existingEntity != nil {
if createFileArgs.EntityParameters == nil {
createFileArgs.EntityParameters = &inventory.EntityParameters{}
}
createFileArgs.EntityParameters.LinkedEntityID = o.existingEntity.ID
}
// Start transaction to create files
fc, tx, ctx, err := inventory.WithTx(ctx, f.fileClient)

@ -28,6 +28,7 @@ type dbfsOption struct {
ancestor *File
notRoot bool
encryptMetadata *types.EncryptMetadata
existingEntity *ent.Entity
}
func newDbfsOption() *dbfsOption {
@ -186,3 +187,11 @@ func WithAncestor(f *File) fs.Option {
o.ancestor = f
})
}
// WithExistingEntity links an already-stored entity to the new file
// instead of creating a fresh one (instant upload / hash dedup).
func WithExistingEntity(e *ent.Entity) fs.Option {
return optionFunc(func(o *dbfsOption) {
o.existingEntity = e
})
}

@ -192,6 +192,22 @@ func (f *DBFS) PrepareUpload(ctx context.Context, req *fs.UploadRequest, opts ..
}
}
// Rapid upload: when the uploader asserted a content hash and a
// completed identical entity already exists within the configured
// dedup scope, the file is materialized from that entity and no
// transfer happens.
var rapidEntity *ent.Entity
if !fileExisted && req.ImportFrom == nil && req.Props.Hash != "" {
if scope := f.settingClient.DBFS(ctx).DedupScope; scope != "off" {
candidate, err := f.fileClient.FindEntityByHash(ctx, req.Props.Hash, req.Props.Size, ancestor.Owner().ID, scope == "global")
if err != nil && !ent.IsNotFound(err) {
return nil, serializer.NewError(serializer.CodeDBError, "Failed to query dedup entity", err)
} else if err == nil {
rapidEntity = candidate
}
}
}
// Create upload placeholder
var (
fileId int
@ -237,6 +253,7 @@ func (f *DBFS) PrepareUpload(ctx context.Context, req *fs.UploadRequest, opts ..
WithErrorOnConflict(),
WithAncestor(ancestor),
WithEncryptMetadata(encryptMetadata),
WithExistingEntity(rapidEntity),
)
if err != nil {
_ = inventory.Rollback(dbTx)
@ -248,6 +265,27 @@ func (f *DBFS) PrepareUpload(ctx context.Context, req *fs.UploadRequest, opts ..
targetFile = uploadPlaceholder.(*File).Model
}
if rapidEntity != nil {
// Instant upload is complete at this point — no session metadata
// or transfer lock is needed.
if err := inventory.CommitWithStorageDiff(ctx, dbTx, f.l, f.userClient); err != nil {
return nil, serializer.NewError(serializer.CodeDBError, "Failed to commit rapid upload", err)
}
f.record(ctx, types.EventEntityUploaded, activity.File(fileId), activity.Extra(map[string]any{"uri": req.Props.Uri.String(), "size": req.Props.Size, "rapid": true}))
return &fs.UploadSession{
Props: &fs.UploadProps{
Uri: req.Props.Uri,
Size: req.Props.Size,
RapidUploaded: true,
},
FileID: fileId,
NewFileCreated: true,
EntityID: entityId,
UID: f.user.ID,
Policy: policy,
}, nil
}
if req.ImportFrom == nil {
// If not importing, we can keep the lock
lockToken = ls.Exclude(lr, f.user, f.hasher)

@ -0,0 +1,150 @@
package dbfs
import (
"context"
"testing"
"time"
"github.com/cloudreve/Cloudreve/v4/ent"
"github.com/cloudreve/Cloudreve/v4/ent/enttest"
entfile "github.com/cloudreve/Cloudreve/v4/ent/file"
"github.com/cloudreve/Cloudreve/v4/ent/storagepolicy"
"github.com/cloudreve/Cloudreve/v4/inventory"
"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/filemanager/eventhub"
"github.com/cloudreve/Cloudreve/v4/pkg/filemanager/fs"
"github.com/cloudreve/Cloudreve/v4/pkg/filemanager/lock"
"github.com/cloudreve/Cloudreve/v4/pkg/hashid"
"github.com/cloudreve/Cloudreve/v4/pkg/logging"
"github.com/cloudreve/Cloudreve/v4/pkg/setting"
"github.com/gofrs/uuid"
"github.com/stretchr/testify/require"
)
type dedupSettingProvider struct {
setting.Provider
scope string
}
func (p dedupSettingProvider) DBFS(context.Context) *setting.DBFS {
return &setting.DBFS{DedupScope: p.scope, MaxPageSize: 200}
}
type stubEventHub struct{}
func (stubEventHub) Subscribe(context.Context, int, string) (chan *eventhub.Event, bool, error) {
return nil, false, nil
}
func (stubEventHub) Unsubscribe(context.Context, int, string) {}
func (stubEventHub) GetSubscribers(context.Context, int) []eventhub.Subscriber {
return nil
}
func (stubEventHub) Close() {}
func dedupUploadFixture(t *testing.T, client *ent.Client, scope string) (*ent.User, *ent.StoragePolicy, *DBFS) {
t.Helper()
ctx := context.Background()
l := logging.NewConsoleLogger(logging.LevelError)
hasher, err := hashid.New("dedup-test-salt")
require.NoError(t, err)
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)
u := client.User.Create().SetEmail("u@example.com").SetNick("u").SetGroup(group).SaveX(ctx)
u.SetGroup(group)
client.File.Create().SetName(inventory.RootFolderName).
SetType(int(types.FileTypeFolder)).SetOwner(u).SaveX(ctx)
f := &DBFS{
user: u,
navigators: make(map[string]Navigator),
fileClient: inventory.NewFileClient(client, conf.SQLiteDB, hasher),
userClient: inventory.NewUserClient(client),
storagePolicyClient: inventory.NewStoragePolicyClient(client, nil),
settingClient: dedupSettingProvider{scope: scope},
hasher: hasher,
l: l,
ls: lock.NewMemLS(hasher, l),
eventHub: stubEventHub{},
}
return u, p, f
}
const dedupTestHash = "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"
func uploadReq(t *testing.T, hasher hashid.Encoder, u *ent.User, name string, size int64, hash string) *fs.UploadRequest {
t.Helper()
uri, err := fs.NewUriFromString(fs.NewMyUri(hashid.EncodeUserID(hasher, u.ID)) + "/" + name)
require.NoError(t, err)
return &fs.UploadRequest{
Props: &fs.UploadProps{
Uri: uri,
Size: size,
Hash: hash,
UploadSessionID: uuid.Must(uuid.NewV4()).String(),
ExpireAt: time.Now().Add(time.Hour),
},
}
}
func seedCompletedEntity(t *testing.T, client *ent.Client, u *ent.User, p *ent.StoragePolicy, hash string, size int64) *ent.Entity {
t.Helper()
return client.Entity.Create().
SetType(int(types.EntityTypeVersion)).
SetSource("cloudreve/data/" + hash).
SetSize(size).
SetHash(hash).
SetReferenceCount(1).
SetCreatedBy(u.ID).
SetStoragePolicyEntities(p.ID).
SaveX(context.Background())
}
func TestPrepareUploadRapid(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, p, f := dedupUploadFixture(t, client, "owner")
existing := seedCompletedEntity(t, client, u, p, dedupTestHash, 1024)
// Normal upload without hash -> transfer session
s, err := f.PrepareUpload(ctx, uploadReq(t, f.hasher, u, "first.txt", 2048, ""))
require.NoError(t, err)
require.False(t, s.Props.RapidUploaded)
require.NotEmpty(t, s.Props.UploadSessionID)
// Same hash + size -> rapid session, no transfer required
s, err = f.PrepareUpload(ctx, uploadReq(t, f.hasher, u, "copy.txt", 1024, dedupTestHash))
require.NoError(t, err)
require.True(t, s.Props.RapidUploaded)
require.Equal(t, existing.ID, s.EntityID)
// Entity refcount bumped, new file linked
require.Equal(t, 2, client.Entity.GetX(ctx, existing.ID).ReferenceCount)
newFile := client.File.Query().Where(entfile.Name("copy.txt")).OnlyX(ctx)
require.Equal(t, existing.ID, newFile.PrimaryEntity)
require.Equal(t, int64(1024), newFile.Size)
// Size mismatch -> normal session
s, err = f.PrepareUpload(ctx, uploadReq(t, f.hasher, u, "bigger.txt", 4096, dedupTestHash))
require.NoError(t, err)
require.False(t, s.Props.RapidUploaded)
}
func TestPrepareUploadRapidScopeOff(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, p, f := dedupUploadFixture(t, client, "off")
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)
}

@ -257,6 +257,9 @@ type (
MimeType string `json:"mime_type,omitempty"` // Expected mimetype
UploadPolicy string `json:"upload_policy,omitempty"` // Upyun upload policy
EncryptMetadata *types.EncryptMetadata `json:"encrypt_metadata,omitempty"`
// RapidUploaded marks that the file was materialized from an existing
// identical entity — the client must not transfer any bytes.
RapidUploaded bool `json:"rapid_uploaded,omitempty"`
}
// UploadSession stores the information of an upload session, used in server side.
@ -305,6 +308,12 @@ type (
ExpireAt time.Time
EncryptionSupported []types.Cipher
ClientSideEncrypted bool // Whether the file stream is already encrypted by client side.
// Hash is the client-asserted sha256 of the file content, used for
// duplicate detection and instant upload. Stored on the entity.
Hash string
// RapidUploaded marks a session whose file was materialized from an
// existing identical entity — no transfer is needed.
RapidUploaded bool
}
// FsOption options for underlying file system.

@ -16,6 +16,7 @@ import (
"github.com/cloudreve/Cloudreve/v4/pkg/cluster"
"github.com/cloudreve/Cloudreve/v4/pkg/filemanager/driver"
"github.com/cloudreve/Cloudreve/v4/pkg/filemanager/fs"
"github.com/cloudreve/Cloudreve/v4/pkg/filemanager/fs/dbfs"
"github.com/cloudreve/Cloudreve/v4/pkg/logging"
"github.com/cloudreve/Cloudreve/v4/pkg/queue"
"github.com/cloudreve/Cloudreve/v4/pkg/serializer"
@ -127,6 +128,22 @@ func (m *manager) CreateUploadSession(ctx context.Context, req *fs.UploadRequest
}
}
if uploadSession.Props != nil && uploadSession.Props.RapidUploaded {
// File was materialized from an existing identical entity — no
// storage credential or session cache is needed. Still queue the
// same post-upload processing (media meta, full-text index) a
// completed upload would get, since those are per-file.
if !m.stateless {
m.postProcessRapidUpload(ctx, uploadSession)
}
return &fs.UploadCredential{
Uri: uploadSession.Props.Uri.String(),
StoragePolicy: uploadSession.Policy,
Expires: req.Props.ExpireAt.Unix(),
RapidUploaded: true,
}, nil
}
d, err := m.GetStorageDriver(ctx, m.CastStoragePolicyOnSlave(ctx, uploadSession.Policy))
if err != nil {
m.OnUploadFailed(ctx, uploadSession)
@ -189,6 +206,30 @@ func (m *manager) CreateUploadSession(ctx context.Context, req *fs.UploadRequest
return credential, nil
}
// postProcessRapidUpload resolves the linked entity's real storage policy
// driver (it may differ from the file's upload policy when dedup matched a
// cross-policy entity) and queues media-meta / full-text tasks.
func (m *manager) postProcessRapidUpload(ctx context.Context, session *fs.UploadSession) {
file, err := m.fs.Get(ctx, session.Props.Uri, dbfs.WithFileEntities())
if err != nil {
m.l.Warning("Failed to load rapid-uploaded file for post-processing: %s", err)
return
}
entity, found := lo.Find(file.Entities(), func(e fs.Entity) bool {
return e.ID() == session.EntityID
})
if !found {
m.l.Warning("Rapid-uploaded file %d missing entity %d", session.FileID, session.EntityID)
return
}
_, d, err := m.getEntityPolicyDriver(ctx, entity, nil)
if err != nil {
m.l.Warning("Failed to resolve rapid-upload entity driver: %s", err)
return
}
m.onNewEntityUploaded(ctx, session, d, file.OwnerID())
}
func (m *manager) ConfirmUploadSession(ctx context.Context, session *fs.UploadSession, chunkIndex int) (fs.File, error) {
// Get placeholder file
file, err := m.fs.Get(ctx, session.Props.Uri)

@ -774,6 +774,7 @@ func (s *settingProvider) DBFS(ctx context.Context) *DBFS {
MaxPageSize: s.getInt(ctx, "max_page_size", 2000),
MaxRecursiveSearchedFolder: s.getInt(ctx, "max_recursive_searched_folder", 65535),
UseSSEForSearch: s.getBoolean(ctx, "use_sse_for_search", false),
DedupScope: s.getString(ctx, "upload_dedup_scope", "owner"),
}
}

@ -118,6 +118,10 @@ type DBFS struct {
MaxPageSize int
MaxRecursiveSearchedFolder int
UseSSEForSearch bool
// DedupScope controls hash-based duplicate detection on upload:
// "off" disables it, "owner" dedups against the uploader's own
// entities, "global" dedups across all users.
DedupScope string
}
type (

@ -61,6 +61,10 @@ type SiteConfig struct {
// AbuseCaptcha controls whether the report-abuse dialog shows captcha.
AbuseCaptcha bool `json:"abuse_captcha,omitempty"`
// UploadDedup signals clients to send content hashes with upload
// sessions, enabling server-side duplicate detection / instant upload.
UploadDedup bool `json:"upload_dedup,omitempty"`
// TaskNodes lists nodes the current user may target when creating tasks
// (remote download, archive ops); populated when the group allows node
// selection, filtered to the group's allowed pool.
@ -257,6 +261,7 @@ func (s *GetSettingService) GetSiteConfig(c *gin.Context) (*SiteConfig, error) {
DefaultShareLinksInProfile: string(shareDefaults.LinksInProfile),
DownloadCDNRoutes: settings.DownloadCDNRoutes(c),
AbuseCaptcha: settings.AbuseCaptchaEnabled(c),
UploadDedup: settings.DBFS(c).DedupScope != "off",
TaskNodes: taskNodes,
AllowSelectNode: allowSelect,
}, nil

@ -163,6 +163,9 @@ type UploadSessionResponse struct {
MimeType string `json:"mime_type,omitempty"`
UploadPolicy string `json:"upload_policy,omitempty"`
EncryptMetadata *types.EncryptMetadata `json:"encrypt_metadata,omitempty"`
// RapidUploaded marks that the file already exists identically on the
// server — no bytes need to be transferred.
RapidUploaded bool `json:"rapid_uploaded,omitempty"`
}
func BuildUploadSessionResponse(session *fs.UploadCredential, hasher hashid.Encoder) *UploadSessionResponse {
@ -180,6 +183,7 @@ func BuildUploadSessionResponse(session *fs.UploadCredential, hasher hashid.Enco
MimeType: session.MimeType,
UploadPolicy: session.UploadPolicy,
EncryptMetadata: session.EncryptMetadata,
RapidUploaded: session.RapidUploaded,
}
if session.EncryptMetadata != nil {

@ -4,6 +4,7 @@ import (
"context"
"fmt"
"strconv"
"strings"
"time"
"github.com/cloudreve/Cloudreve/v4/application/dependency"
@ -31,6 +32,9 @@ type (
EntityType string `json:"entity_type" binding:"eq=|eq=live_photo|eq=version"`
EncryptionSupported []types.Cipher `json:"encryption_supported"`
Previous string `json:"previous" form:"previous"`
// Hash is the client-computed sha256 of the content, enabling
// duplicate detection / instant upload when dedup is enabled.
Hash string `json:"hash" binding:"omitempty,hexadecimal,len=64"`
}
)
@ -77,6 +81,7 @@ func (service *CreateUploadSessionService) Create(c context.Context) (*UploadSes
PreferredStoragePolicy: policyId,
EncryptionSupported: service.EncryptionSupported,
ClientSideEncrypted: len(service.EncryptionSupported) > 0,
Hash: strings.ToLower(service.Hash),
},
}

Loading…
Cancel
Save