Update internal/storage/document.go
This commit is contained in:
1 parent
64231bbc67
commit
3b5458305b
1 file changed
+152
-52
+152
-52
@@ -14,10 +14,13 @@
|
|||||||
// Документ является основной единицей хранения в СУБД futriis.
|
// Документ является основной единицей хранения в СУБД futriis.
|
||||||
// Lock-free: Document.fields переведён на atomic.Value для wait-free доступа.
|
// Lock-free: Document.fields переведён на atomic.Value для wait-free доступа.
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
package storage
|
package storage
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"math/rand"
|
||||||
"strings"
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
"sync/atomic"
|
"sync/atomic"
|
||||||
@@ -28,10 +31,33 @@ import (
|
|||||||
"github.com/google/uuid"
|
"github.com/google/uuid"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// =============================================================================
|
||||||
|
// КОНСТАНТЫ ДЛЯ ЗАЩИТЫ ОТ LIVELOCK
|
||||||
|
// =============================================================================
|
||||||
|
|
||||||
|
const (
|
||||||
|
// Максимальное количество попыток CAS перед применением backoff
|
||||||
|
maxCASRetriesBeforeBackoff = 10
|
||||||
|
// Базовая задержка backoff в наносекундах
|
||||||
|
baseBackoffNanos = 100
|
||||||
|
// Максимальная задержка backoff в наносекундах
|
||||||
|
maxBackoffNanos = 10000
|
||||||
|
)
|
||||||
|
|
||||||
|
// =============================================================================
|
||||||
|
// FIELD VERSION - ЗАЩИТА ОТ ABA
|
||||||
|
// =============================================================================
|
||||||
|
|
||||||
|
// fieldVersionWrapper оборачивает map полей с версией для защиты от ABA
|
||||||
|
type fieldVersionWrapper struct {
|
||||||
|
fields map[string]interface{}
|
||||||
|
version uint64
|
||||||
|
}
|
||||||
|
|
||||||
// Document представляет документ в коллекции (аналог строки в реляционной СУБД)
|
// Document представляет документ в коллекции (аналог строки в реляционной СУБД)
|
||||||
type Document struct {
|
type Document struct {
|
||||||
ID string `msgpack:"_id"` // Уникальный идентификатор документа
|
ID string `msgpack:"_id"` // Уникальный идентификатор документа
|
||||||
fieldsPtr atomic.Value // map[string]interface{} - lock-free хранилище полей
|
fieldsPtr atomic.Value // *fieldVersionWrapper - lock-free хранилище полей с версией
|
||||||
CreatedAt int64 `msgpack:"created_at"` // Время создания (Unix миллисекунды)
|
CreatedAt int64 `msgpack:"created_at"` // Время создания (Unix миллисекунды)
|
||||||
UpdatedAt int64 `msgpack:"updated_at"` // Время последнего обновления
|
UpdatedAt int64 `msgpack:"updated_at"` // Время последнего обновления
|
||||||
DeletedAt int64 `msgpack:"deleted_at"` // Время удаления (Unix миллисекунды, 0 = не удалён)
|
DeletedAt int64 `msgpack:"deleted_at"` // Время удаления (Unix миллисекунды, 0 = не удалён)
|
||||||
@@ -80,7 +106,10 @@ func NewDocument() *Document {
|
|||||||
Compressed: false,
|
Compressed: false,
|
||||||
OriginalSize: 0,
|
OriginalSize: 0,
|
||||||
}
|
}
|
||||||
d.fieldsPtr.Store(make(map[string]interface{}))
|
d.fieldsPtr.Store(&fieldVersionWrapper{
|
||||||
|
fields: make(map[string]interface{}),
|
||||||
|
version: 0,
|
||||||
|
})
|
||||||
return d
|
return d
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -96,49 +125,132 @@ func NewDocumentWithID(id string) *Document {
|
|||||||
Compressed: false,
|
Compressed: false,
|
||||||
OriginalSize: 0,
|
OriginalSize: 0,
|
||||||
}
|
}
|
||||||
d.fieldsPtr.Store(make(map[string]interface{}))
|
d.fieldsPtr.Store(&fieldVersionWrapper{
|
||||||
|
fields: make(map[string]interface{}),
|
||||||
|
version: 0,
|
||||||
|
})
|
||||||
return d
|
return d
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// loadFieldsWrapper загружает обёртку полей с версией (lock-free)
|
||||||
|
func (d *Document) loadFieldsWrapper() *fieldVersionWrapper {
|
||||||
|
val := d.fieldsPtr.Load()
|
||||||
|
if val == nil {
|
||||||
|
return &fieldVersionWrapper{
|
||||||
|
fields: make(map[string]interface{}),
|
||||||
|
version: 0,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return val.(*fieldVersionWrapper)
|
||||||
|
}
|
||||||
|
|
||||||
// loadFields загружает карту полей (lock-free)
|
// loadFields загружает карту полей (lock-free)
|
||||||
func (d *Document) loadFields() map[string]interface{} {
|
func (d *Document) loadFields() map[string]interface{} {
|
||||||
val := d.fieldsPtr.Load()
|
wrapper := d.loadFieldsWrapper()
|
||||||
if val == nil {
|
if wrapper == nil || wrapper.fields == nil {
|
||||||
return make(map[string]interface{})
|
return make(map[string]interface{})
|
||||||
}
|
}
|
||||||
return val.(map[string]interface{})
|
return wrapper.fields
|
||||||
}
|
}
|
||||||
|
|
||||||
// storeFields сохраняет карту полей (lock-free)
|
// storeFields сохраняет карту полей (lock-free)
|
||||||
func (d *Document) storeFields(newMap map[string]interface{}) {
|
func (d *Document) storeFields(newMap map[string]interface{}) {
|
||||||
d.fieldsPtr.Store(newMap)
|
oldWrapper := d.loadFieldsWrapper()
|
||||||
|
newVersion := uint64(0)
|
||||||
|
if oldWrapper != nil {
|
||||||
|
newVersion = oldWrapper.version + 1
|
||||||
|
}
|
||||||
|
d.fieldsPtr.Store(&fieldVersionWrapper{
|
||||||
|
fields: newMap,
|
||||||
|
version: newVersion,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
// casFields выполняет CAS операцию для полей с защитой от ABA (lock-free)
|
||||||
|
// ИСПРАВЛЕНО: Использует версионирование для защиты от ABA-проблемы
|
||||||
|
func (d *Document) casFields(oldWrapper, newWrapper *fieldVersionWrapper) bool {
|
||||||
|
if oldWrapper == nil || newWrapper == nil {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
|
// Проверяем, что версия увеличилась (защита от ABA)
|
||||||
|
if newWrapper.version <= oldWrapper.version {
|
||||||
|
newWrapper.version = oldWrapper.version + 1
|
||||||
|
}
|
||||||
|
|
||||||
|
return d.fieldsPtr.CompareAndSwap(oldWrapper, newWrapper)
|
||||||
|
}
|
||||||
|
|
||||||
|
// casFieldsWithBackoff выполняет CAS с backoff для предотвращения livelock
|
||||||
|
// ИСПРАВЛЕНО: Добавлен экспоненциальный backoff для предотвращения livelock
|
||||||
|
// ИСПРАВЛЕНО: Исправлены типы int/int64 для устранения ошибок компилятора
|
||||||
|
func (d *Document) casFieldsWithBackoff(updateFunc func(map[string]interface{}) map[string]interface{}) bool {
|
||||||
|
attempt := 0
|
||||||
|
for {
|
||||||
|
oldWrapper := d.loadFieldsWrapper()
|
||||||
|
|
||||||
|
// Создаём копию полей для модификации
|
||||||
|
newFields := make(map[string]interface{})
|
||||||
|
for k, v := range oldWrapper.fields {
|
||||||
|
newFields[k] = v
|
||||||
|
}
|
||||||
|
|
||||||
|
// Применяем обновление
|
||||||
|
newFields = updateFunc(newFields)
|
||||||
|
|
||||||
|
newWrapper := &fieldVersionWrapper{
|
||||||
|
fields: newFields,
|
||||||
|
version: oldWrapper.version + 1,
|
||||||
|
}
|
||||||
|
|
||||||
|
if d.casFields(oldWrapper, newWrapper) {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
// Backoff при неудаче для предотвращения livelock
|
||||||
|
attempt++
|
||||||
|
if attempt > maxCASRetriesBeforeBackoff {
|
||||||
|
// Экспоненциальный backoff с jitter
|
||||||
|
backoff := baseBackoffNanos * (1 << min(attempt-maxCASRetriesBeforeBackoff, 10))
|
||||||
|
if backoff > maxBackoffNanos {
|
||||||
|
backoff = maxBackoffNanos
|
||||||
|
}
|
||||||
|
// ИСПРАВЛЕНО: Приводим типы для совместимости с rand.Int63n (int64)
|
||||||
|
// Добавляем jitter для избежания синхронизации потоков
|
||||||
|
jitter := rand.Int63n(int64(backoff)/2 + 1)
|
||||||
|
totalBackoff := backoff + int(jitter)
|
||||||
|
time.Sleep(time.Duration(totalBackoff))
|
||||||
|
}
|
||||||
|
|
||||||
|
// Ограничиваем максимальное количество попыток
|
||||||
|
if attempt > 1000 {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// min возвращает минимальное из двух чисел
|
||||||
|
func min(a, b int) int {
|
||||||
|
if a < b {
|
||||||
|
return a
|
||||||
|
}
|
||||||
|
return b
|
||||||
}
|
}
|
||||||
|
|
||||||
// SetField устанавливает значение поля документа (lock-free)
|
// SetField устанавливает значение поля документа (lock-free)
|
||||||
|
// ИСПРАВЛЕНО: Использует casFieldsWithBackoff для защиты от ABA и livelock
|
||||||
func (d *Document) SetField(name string, value interface{}) {
|
func (d *Document) SetField(name string, value interface{}) {
|
||||||
for {
|
d.casFieldsWithBackoff(func(fields map[string]interface{}) map[string]interface{} {
|
||||||
oldFields := d.loadFields()
|
fields[name] = value
|
||||||
newFields := make(map[string]interface{})
|
return fields
|
||||||
for k, v := range oldFields {
|
})
|
||||||
newFields[k] = v
|
|
||||||
}
|
|
||||||
newFields[name] = value
|
|
||||||
|
|
||||||
if d.compareAndSwapFields(oldFields, newFields) {
|
|
||||||
d.UpdatedAt = time.Now().UnixMilli()
|
d.UpdatedAt = time.Now().UnixMilli()
|
||||||
d.Version++
|
d.Version++
|
||||||
d.Compressed = false
|
d.Compressed = false
|
||||||
|
|
||||||
// Аудит изменения поля
|
// Аудит изменения поля
|
||||||
AuditFieldOperation("UPDATE", "", "", d.ID, name, value)
|
AuditFieldOperation("UPDATE", "", "", d.ID, name, value)
|
||||||
return
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// compareAndSwapFields выполняет CAS операцию для полей (lock-free)
|
|
||||||
func (d *Document) compareAndSwapFields(old, new map[string]interface{}) bool {
|
|
||||||
return d.fieldsPtr.CompareAndSwap(old, new)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// GetField возвращает значение поля документа (lock-free)
|
// GetField возвращает значение поля документа (lock-free)
|
||||||
@@ -151,30 +263,19 @@ func (d *Document) GetField(name string) (interface{}, error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// DeleteField удаляет поле из документа (lock-free)
|
// DeleteField удаляет поле из документа (lock-free)
|
||||||
|
// ИСПРАВЛЕНО: Использует casFieldsWithBackoff для защиты от ABA и livelock
|
||||||
func (d *Document) DeleteField(name string) {
|
func (d *Document) DeleteField(name string) {
|
||||||
for {
|
d.casFieldsWithBackoff(func(fields map[string]interface{}) map[string]interface{} {
|
||||||
oldFields := d.loadFields()
|
delete(fields, name)
|
||||||
if _, exists := oldFields[name]; !exists {
|
return fields
|
||||||
return
|
})
|
||||||
}
|
|
||||||
|
|
||||||
newFields := make(map[string]interface{})
|
|
||||||
for k, v := range oldFields {
|
|
||||||
if k != name {
|
|
||||||
newFields[k] = v
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
if d.compareAndSwapFields(oldFields, newFields) {
|
|
||||||
d.UpdatedAt = time.Now().UnixMilli()
|
d.UpdatedAt = time.Now().UnixMilli()
|
||||||
d.Version++
|
d.Version++
|
||||||
d.Compressed = false
|
d.Compressed = false
|
||||||
|
|
||||||
// Аудит удаления поля
|
// Аудит удаления поля
|
||||||
AuditFieldOperation("DELETE", "", "", d.ID, name, nil)
|
AuditFieldOperation("DELETE", "", "", d.ID, name, nil)
|
||||||
return
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// HasField проверяет наличие поля в документе (lock-free)
|
// HasField проверяет наличие поля в документе (lock-free)
|
||||||
@@ -246,7 +347,7 @@ func (d *Document) Serialize() ([]byte, error) {
|
|||||||
Compressed: d.Compressed,
|
Compressed: d.Compressed,
|
||||||
OriginalSize: d.OriginalSize,
|
OriginalSize: d.OriginalSize,
|
||||||
}
|
}
|
||||||
docCopy.fieldsPtr.Store(fields)
|
docCopy.storeFields(fields)
|
||||||
|
|
||||||
data, err := serializer.Marshal(docCopy)
|
data, err := serializer.Marshal(docCopy)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -285,7 +386,7 @@ func (d *Document) Deserialize(data []byte) error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
d.ID = doc.ID
|
d.ID = doc.ID
|
||||||
d.fieldsPtr.Store(doc.loadFields())
|
d.storeFields(doc.loadFields())
|
||||||
d.CreatedAt = doc.CreatedAt
|
d.CreatedAt = doc.CreatedAt
|
||||||
d.UpdatedAt = doc.UpdatedAt
|
d.UpdatedAt = doc.UpdatedAt
|
||||||
d.DeletedAt = doc.DeletedAt
|
d.DeletedAt = doc.DeletedAt
|
||||||
@@ -298,7 +399,7 @@ func (d *Document) Deserialize(data []byte) error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
d.ID = doc.ID
|
d.ID = doc.ID
|
||||||
d.fieldsPtr.Store(doc.loadFields())
|
d.storeFields(doc.loadFields())
|
||||||
d.CreatedAt = doc.CreatedAt
|
d.CreatedAt = doc.CreatedAt
|
||||||
d.UpdatedAt = doc.UpdatedAt
|
d.UpdatedAt = doc.UpdatedAt
|
||||||
d.DeletedAt = doc.DeletedAt
|
d.DeletedAt = doc.DeletedAt
|
||||||
@@ -329,30 +430,29 @@ func (d *Document) Clone() *Document {
|
|||||||
for k, v := range fields {
|
for k, v := range fields {
|
||||||
copiedFields[k] = deepCopyValue(v)
|
copiedFields[k] = deepCopyValue(v)
|
||||||
}
|
}
|
||||||
clone.fieldsPtr.Store(copiedFields)
|
clone.storeFields(copiedFields)
|
||||||
|
|
||||||
return clone
|
return clone
|
||||||
}
|
}
|
||||||
|
|
||||||
// Update применяет обновление к документу (атомарно, lock-free)
|
// Update применяет обновление к документу (атомарно, lock-free)
|
||||||
|
// ИСПРАВЛЕНО: Использует casFieldsWithBackoff для защиты от ABA и livelock
|
||||||
func (d *Document) Update(updates map[string]interface{}) error {
|
func (d *Document) Update(updates map[string]interface{}) error {
|
||||||
for {
|
success := d.casFieldsWithBackoff(func(fields map[string]interface{}) map[string]interface{} {
|
||||||
oldFields := d.loadFields()
|
|
||||||
newFields := make(map[string]interface{})
|
|
||||||
for k, v := range oldFields {
|
|
||||||
newFields[k] = v
|
|
||||||
}
|
|
||||||
for k, v := range updates {
|
for k, v := range updates {
|
||||||
newFields[k] = v
|
fields[k] = v
|
||||||
}
|
}
|
||||||
|
return fields
|
||||||
|
})
|
||||||
|
|
||||||
if d.compareAndSwapFields(oldFields, newFields) {
|
if success {
|
||||||
d.UpdatedAt = time.Now().UnixMilli()
|
d.UpdatedAt = time.Now().UnixMilli()
|
||||||
d.Version++
|
d.Version++
|
||||||
d.Compressed = false
|
d.Compressed = false
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
}
|
|
||||||
|
return fmt.Errorf("failed to update document after multiple attempts (possible high contention)")
|
||||||
}
|
}
|
||||||
|
|
||||||
// SoftDelete мягко удаляет документ (устанавливает метку времени удаления)
|
// SoftDelete мягко удаляет документ (устанавливает метку времени удаления)
|
||||||
|
|||||||
Reference in new issue
Block a user