Upload files to "internal/storage"
This commit is contained in:
1 parent
bc445ad7be
commit
ab48d50988
3 files changed
+3274
No files matched your search
@@ -0,0 +1,451 @@
|
||||
/*
|
||||
* Copyright 2026 Safronov Grigorii
|
||||
*
|
||||
* Licensed under the CDDL, Version 1.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
*
|
||||
* You may obtain a copy of the License at
|
||||
* https://opensource.org/licenses/CDDL-1.0
|
||||
*/
|
||||
|
||||
// Файл: internal/storage/runtime_limits.go
|
||||
// Назначение: Ограничения на размер коллекции/документа в рантайме
|
||||
// Добавлен механизм eviction при нехватке памяти (OOM protection)
|
||||
// Eviction проверяет активные batch-операции перед удалением
|
||||
//
|
||||
// ИСПРАВЛЕНО (2026-10): после замены ACID-транзакций на batch-операции
|
||||
// функция isDocInActiveTransaction ссылалась на удалённые globalTxManager
|
||||
// и Transaction. Заменена на isDocInActiveBatch, который проверяет
|
||||
// реестр активных batch'ей в globalBatchManager.
|
||||
//
|
||||
// Логика та же: не вытеснять документы, которые в данный момент
|
||||
// участвуют в незакоммиченной batch-операции. Batch живёт коротко
|
||||
// (создаётся, коммитится и удаляется в рамках одной функции), но
|
||||
// между AddUpdate и Commit документ ещё не изменён, и eviction
|
||||
// мог бы его удалить — защита сохраняется.
|
||||
|
||||
package storage
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"runtime"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
)
|
||||
|
||||
// LoggerInterface определяет интерфейс для логирования
|
||||
type LoggerInterface interface {
|
||||
Debug(msg string)
|
||||
Info(msg string)
|
||||
Error(msg string)
|
||||
Warn(msg string)
|
||||
}
|
||||
|
||||
// EvictionPolicy определяет политику вытеснения
|
||||
type EvictionPolicy int
|
||||
|
||||
const (
|
||||
EvictionNone EvictionPolicy = iota
|
||||
EvictionLRU
|
||||
EvictionTTL
|
||||
EvictionOldest
|
||||
)
|
||||
|
||||
// RuntimeLimitsManager управляет runtime-ограничениями
|
||||
type RuntimeLimitsManager struct {
|
||||
mu sync.RWMutex
|
||||
globalMaxDocSize int64
|
||||
globalMaxCollSize int64
|
||||
globalMaxDocsPerColl int64
|
||||
globalMaxMemory int64
|
||||
collectionOverrides map[string]*CollectionLimits
|
||||
metrics *LimitMetrics
|
||||
logger LoggerInterface
|
||||
enabled bool
|
||||
evictionPolicy EvictionPolicy
|
||||
memoryThreshold float64
|
||||
evictionChan chan string
|
||||
stopChan chan struct{}
|
||||
wg sync.WaitGroup
|
||||
}
|
||||
|
||||
// CollectionLimits содержит лимиты для конкретной коллекции
|
||||
type CollectionLimits struct {
|
||||
MaxDocSize int64
|
||||
MaxCollectionSize int64
|
||||
MaxDocuments int64
|
||||
EvictionPolicy EvictionPolicy
|
||||
LastUpdated int64
|
||||
}
|
||||
|
||||
// LimitMetrics хранит метрики ограничений
|
||||
type LimitMetrics struct {
|
||||
RejectedBySize atomic.Uint64
|
||||
RejectedByDocCount atomic.Uint64
|
||||
RejectedByCollSize atomic.Uint64
|
||||
RejectedByMemory atomic.Uint64
|
||||
EvictedDocuments atomic.Uint64
|
||||
EvictedBytes atomic.Uint64
|
||||
LastCheckTime atomic.Int64
|
||||
LastEvictionTime atomic.Int64
|
||||
SkippedEviction atomic.Uint64 // Пропущено из-за активных batch-операций
|
||||
}
|
||||
|
||||
// RuntimeLimitsConfig содержит конфигурацию ограничений
|
||||
type RuntimeLimitsConfig struct {
|
||||
Enabled bool `json:"enabled"`
|
||||
GlobalMaxDocSizeMB int `json:"global_max_doc_size_mb"`
|
||||
GlobalMaxCollSizeMB int64 `json:"global_max_coll_size_mb"`
|
||||
GlobalMaxDocsPerColl int64 `json:"global_max_docs_per_coll"`
|
||||
GlobalMaxMemoryMB int64 `json:"global_max_memory_mb"`
|
||||
EvictionPolicy EvictionPolicy `json:"eviction_policy"`
|
||||
MemoryThreshold float64 `json:"memory_threshold"`
|
||||
}
|
||||
|
||||
// DefaultRuntimeLimitsConfig возвращает конфигурацию по умолчанию
|
||||
func DefaultRuntimeLimitsConfig() *RuntimeLimitsConfig {
|
||||
return &RuntimeLimitsConfig{
|
||||
Enabled: true,
|
||||
GlobalMaxDocSizeMB: 16,
|
||||
GlobalMaxCollSizeMB: 10240,
|
||||
GlobalMaxDocsPerColl: 10000000,
|
||||
GlobalMaxMemoryMB: 0,
|
||||
EvictionPolicy: EvictionLRU,
|
||||
MemoryThreshold: 0.85,
|
||||
}
|
||||
}
|
||||
|
||||
// NewRuntimeLimitsManager создаёт новый менеджер ограничений
|
||||
func NewRuntimeLimitsManager(cfg *RuntimeLimitsConfig, logger LoggerInterface) *RuntimeLimitsManager {
|
||||
if cfg == nil {
|
||||
cfg = DefaultRuntimeLimitsConfig()
|
||||
}
|
||||
|
||||
maxMemory := cfg.GlobalMaxMemoryMB * 1024 * 1024
|
||||
if maxMemory <= 0 {
|
||||
var memStats runtime.MemStats
|
||||
runtime.ReadMemStats(&memStats)
|
||||
maxMemory = int64(float64(memStats.Sys) * 0.8)
|
||||
}
|
||||
|
||||
rlm := &RuntimeLimitsManager{
|
||||
globalMaxDocSize: int64(cfg.GlobalMaxDocSizeMB) * 1024 * 1024,
|
||||
globalMaxCollSize: cfg.GlobalMaxCollSizeMB * 1024 * 1024,
|
||||
globalMaxDocsPerColl: cfg.GlobalMaxDocsPerColl,
|
||||
globalMaxMemory: maxMemory,
|
||||
collectionOverrides: make(map[string]*CollectionLimits),
|
||||
metrics: &LimitMetrics{},
|
||||
logger: logger,
|
||||
enabled: cfg.Enabled,
|
||||
evictionPolicy: cfg.EvictionPolicy,
|
||||
memoryThreshold: cfg.MemoryThreshold,
|
||||
evictionChan: make(chan string, 100),
|
||||
stopChan: make(chan struct{}),
|
||||
}
|
||||
|
||||
rlm.wg.Add(1)
|
||||
go rlm.memoryMonitorLoop()
|
||||
|
||||
if logger != nil {
|
||||
logger.Debug(fmt.Sprintf("Runtime limits manager initialized: maxDoc=%dMB, maxColl=%dMB, maxDocs=%d, maxMemory=%dMB, eviction=%d",
|
||||
cfg.GlobalMaxDocSizeMB, cfg.GlobalMaxCollSizeMB, cfg.GlobalMaxDocsPerColl, maxMemory/(1024*1024), cfg.EvictionPolicy))
|
||||
}
|
||||
|
||||
return rlm
|
||||
}
|
||||
|
||||
func (rlm *RuntimeLimitsManager) memoryMonitorLoop() {
|
||||
defer rlm.wg.Done()
|
||||
ticker := time.NewTicker(10 * time.Second)
|
||||
defer ticker.Stop()
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ticker.C:
|
||||
if !rlm.enabled {
|
||||
continue
|
||||
}
|
||||
var memStats runtime.MemStats
|
||||
runtime.ReadMemStats(&memStats)
|
||||
currentUsage := int64(memStats.Alloc)
|
||||
usageRatio := float64(currentUsage) / float64(rlm.globalMaxMemory)
|
||||
|
||||
if usageRatio > rlm.memoryThreshold {
|
||||
if rlm.logger != nil {
|
||||
rlm.logger.Warn(fmt.Sprintf("Memory usage %.2f%% exceeds threshold %.2f%%, triggering eviction",
|
||||
usageRatio*100, rlm.memoryThreshold*100))
|
||||
}
|
||||
rlm.triggerEviction()
|
||||
}
|
||||
case <-rlm.stopChan:
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (rlm *RuntimeLimitsManager) triggerEviction() {
|
||||
select {
|
||||
case rlm.evictionChan <- "global":
|
||||
default:
|
||||
}
|
||||
}
|
||||
|
||||
// isDocInActiveBatch проверяет, участвует ли документ в активной
|
||||
// batch-операции.
|
||||
//
|
||||
// ИСПРАВЛЕНО (2026-10): ранее метод назывался isDocInActiveTransaction
|
||||
// и обращался к globalTxManager и Transaction, которые были удалены при
|
||||
// переходе на batch-операции. Теперь проверяем реестр активных batch'ей
|
||||
// в globalBatchManager (см. transactions.go).
|
||||
//
|
||||
// Защита сохраняется: batch живёт коротко (создаётся и коммитится в
|
||||
// рамках одной функции), но между AddUpdate и Commit документ ещё не
|
||||
// изменён, и eviction мог бы его удалить.
|
||||
func (rlm *RuntimeLimitsManager) isDocInActiveBatch(docID, dbName, collName string) bool {
|
||||
bm := GetBatchManager()
|
||||
if bm == nil {
|
||||
return false
|
||||
}
|
||||
return bm.IsDocInActiveBatch(docID, dbName, collName)
|
||||
}
|
||||
|
||||
// EvictFromCollection выполняет вытеснение документов из коллекции.
|
||||
//
|
||||
// Проверка активных batch-операций перед удалением.
|
||||
func (rlm *RuntimeLimitsManager) EvictFromCollection(coll *Collection, targetBytes int64) (int64, error) {
|
||||
if !rlm.enabled {
|
||||
return 0, nil
|
||||
}
|
||||
|
||||
evictedBytes := int64(0)
|
||||
evictedCount := int64(0)
|
||||
skippedCount := int64(0)
|
||||
|
||||
docs := coll.GetAllDocumentsIncludingDeleted()
|
||||
|
||||
// Сначала вытесняем удалённые документы
|
||||
for _, doc := range docs {
|
||||
if evictedBytes >= targetBytes {
|
||||
break
|
||||
}
|
||||
|
||||
if doc.IsDeleted() {
|
||||
// Проверяем активные batch-операции
|
||||
if rlm.isDocInActiveBatch(doc.ID, coll.DBName(), coll.Name()) {
|
||||
skippedCount++
|
||||
continue
|
||||
}
|
||||
|
||||
size := doc.OriginalSize
|
||||
if size == 0 {
|
||||
size = 1024
|
||||
}
|
||||
|
||||
if err := coll.PermanentDelete(doc.ID); err == nil {
|
||||
evictedBytes += size
|
||||
evictedCount++
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Затем вытесняем самые старые
|
||||
if evictedBytes < targetBytes && rlm.evictionPolicy == EvictionLRU {
|
||||
for _, doc := range docs {
|
||||
if evictedBytes >= targetBytes {
|
||||
break
|
||||
}
|
||||
|
||||
if !doc.IsDeleted() {
|
||||
// Проверяем активные batch-операции
|
||||
if rlm.isDocInActiveBatch(doc.ID, coll.DBName(), coll.Name()) {
|
||||
skippedCount++
|
||||
continue
|
||||
}
|
||||
|
||||
size := doc.OriginalSize
|
||||
if size == 0 {
|
||||
size = 1024
|
||||
}
|
||||
|
||||
if coll.metadata.Settings.SoftDelete {
|
||||
if err := coll.Delete(doc.ID); err == nil {
|
||||
evictedBytes += size
|
||||
evictedCount++
|
||||
}
|
||||
} else {
|
||||
if err := coll.PermanentDelete(doc.ID); err == nil {
|
||||
evictedBytes += size
|
||||
evictedCount++
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
rlm.metrics.EvictedDocuments.Add(uint64(evictedCount))
|
||||
rlm.metrics.EvictedBytes.Add(uint64(evictedBytes))
|
||||
if skippedCount > 0 {
|
||||
rlm.metrics.SkippedEviction.Add(uint64(skippedCount))
|
||||
}
|
||||
rlm.metrics.LastEvictionTime.Store(time.Now().UnixMilli())
|
||||
|
||||
if rlm.logger != nil {
|
||||
rlm.logger.Info(fmt.Sprintf("Evicted %d documents (%d bytes) from collection %s.%s (skipped %d in active batches)",
|
||||
evictedCount, evictedBytes, coll.DBName(), coll.Name(), skippedCount))
|
||||
}
|
||||
|
||||
return evictedBytes, nil
|
||||
}
|
||||
|
||||
func (rlm *RuntimeLimitsManager) ValidateDocumentSize(dbName, collName string, docSize int64) error {
|
||||
if !rlm.enabled {
|
||||
return nil
|
||||
}
|
||||
|
||||
limit := rlm.getCollectionLimit(dbName, collName)
|
||||
maxSize := rlm.globalMaxDocSize
|
||||
if limit != nil && limit.MaxDocSize > 0 {
|
||||
maxSize = limit.MaxDocSize
|
||||
}
|
||||
|
||||
if docSize > maxSize {
|
||||
rlm.metrics.RejectedBySize.Add(1)
|
||||
return fmt.Errorf("document size %d bytes exceeds limit %d bytes", docSize, maxSize)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (rlm *RuntimeLimitsManager) ValidateCollectionSize(coll *Collection, newDocSize int64) error {
|
||||
if !rlm.enabled {
|
||||
return nil
|
||||
}
|
||||
|
||||
limit := rlm.getCollectionLimit(coll.DBName(), coll.Name())
|
||||
maxSize := rlm.globalMaxCollSize
|
||||
if limit != nil && limit.MaxCollectionSize > 0 {
|
||||
maxSize = limit.MaxCollectionSize
|
||||
}
|
||||
|
||||
currentSize := coll.Size()
|
||||
if currentSize+newDocSize > maxSize {
|
||||
neededBytes := currentSize + newDocSize - maxSize
|
||||
if evicted, err := rlm.EvictFromCollection(coll, neededBytes); err == nil && evicted >= neededBytes {
|
||||
return nil
|
||||
}
|
||||
|
||||
rlm.metrics.RejectedByCollSize.Add(1)
|
||||
return fmt.Errorf("collection size would exceed limit %d bytes (current: %d, new: %d)",
|
||||
maxSize, currentSize, newDocSize)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (rlm *RuntimeLimitsManager) ValidateDocumentCount(coll *Collection) error {
|
||||
if !rlm.enabled {
|
||||
return nil
|
||||
}
|
||||
|
||||
limit := rlm.getCollectionLimit(coll.DBName(), coll.Name())
|
||||
maxDocs := rlm.globalMaxDocsPerColl
|
||||
if limit != nil && limit.MaxDocuments > 0 {
|
||||
maxDocs = limit.MaxDocuments
|
||||
}
|
||||
|
||||
currentCount := coll.Count()
|
||||
if currentCount >= maxDocs {
|
||||
targetBytes := int64((currentCount - maxDocs + 1) * 1024)
|
||||
if evicted, err := rlm.EvictFromCollection(coll, targetBytes); err == nil && evicted > 0 {
|
||||
if coll.Count() < maxDocs {
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
rlm.metrics.RejectedByDocCount.Add(1)
|
||||
return fmt.Errorf("collection has reached maximum document count %d", maxDocs)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (rlm *RuntimeLimitsManager) CheckMemoryUsage() error {
|
||||
if !rlm.enabled {
|
||||
return nil
|
||||
}
|
||||
|
||||
var memStats runtime.MemStats
|
||||
runtime.ReadMemStats(&memStats)
|
||||
|
||||
currentUsage := int64(memStats.Alloc)
|
||||
if currentUsage > rlm.globalMaxMemory {
|
||||
rlm.metrics.RejectedByMemory.Add(1)
|
||||
return fmt.Errorf("memory usage %d bytes exceeds limit %d bytes", currentUsage, rlm.globalMaxMemory)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (rlm *RuntimeLimitsManager) getCollectionLimit(dbName, collName string) *CollectionLimits {
|
||||
rlm.mu.RLock()
|
||||
defer rlm.mu.RUnlock()
|
||||
|
||||
key := fmt.Sprintf("%s.%s", dbName, collName)
|
||||
if limits, ok := rlm.collectionOverrides[key]; ok {
|
||||
return limits
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (rlm *RuntimeLimitsManager) SetCollectionLimits(dbName, collName string, maxDocSizeMB int, maxCollSizeMB int64, maxDocuments int64) {
|
||||
rlm.mu.Lock()
|
||||
defer rlm.mu.Unlock()
|
||||
|
||||
key := fmt.Sprintf("%s.%s", dbName, collName)
|
||||
rlm.collectionOverrides[key] = &CollectionLimits{
|
||||
MaxDocSize: int64(maxDocSizeMB) * 1024 * 1024,
|
||||
MaxCollectionSize: maxCollSizeMB * 1024 * 1024,
|
||||
MaxDocuments: maxDocuments,
|
||||
EvictionPolicy: rlm.evictionPolicy,
|
||||
LastUpdated: time.Now().UnixMilli(),
|
||||
}
|
||||
|
||||
if rlm.logger != nil {
|
||||
rlm.logger.Info(fmt.Sprintf("Set limits for %s: maxDoc=%dMB, maxColl=%dMB, maxDocs=%d",
|
||||
key, maxDocSizeMB, maxCollSizeMB, maxDocuments))
|
||||
}
|
||||
}
|
||||
|
||||
func (rlm *RuntimeLimitsManager) RemoveCollectionLimits(dbName, collName string) {
|
||||
rlm.mu.Lock()
|
||||
defer rlm.mu.Unlock()
|
||||
|
||||
key := fmt.Sprintf("%s.%s", dbName, collName)
|
||||
delete(rlm.collectionOverrides, key)
|
||||
|
||||
if rlm.logger != nil {
|
||||
rlm.logger.Info(fmt.Sprintf("Removed limits override for %s", key))
|
||||
}
|
||||
}
|
||||
|
||||
func (rlm *RuntimeLimitsManager) GetMetrics() map[string]interface{} {
|
||||
return map[string]interface{}{
|
||||
"rejected_by_size": rlm.metrics.RejectedBySize.Load(),
|
||||
"rejected_by_doc_count": rlm.metrics.RejectedByDocCount.Load(),
|
||||
"rejected_by_coll_size": rlm.metrics.RejectedByCollSize.Load(),
|
||||
"rejected_by_memory": rlm.metrics.RejectedByMemory.Load(),
|
||||
"evicted_documents": rlm.metrics.EvictedDocuments.Load(),
|
||||
"evicted_bytes": rlm.metrics.EvictedBytes.Load(),
|
||||
"skipped_eviction": rlm.metrics.SkippedEviction.Load(),
|
||||
"last_eviction_time": rlm.metrics.LastEvictionTime.Load(),
|
||||
"global_max_doc_size_mb": rlm.globalMaxDocSize / (1024 * 1024),
|
||||
"global_max_coll_size_mb": rlm.globalMaxCollSize / (1024 * 1024),
|
||||
"global_max_docs_per_coll": rlm.globalMaxDocsPerColl,
|
||||
"global_max_memory_mb": rlm.globalMaxMemory / (1024 * 1024),
|
||||
"eviction_policy": rlm.evictionPolicy,
|
||||
"memory_threshold": rlm.memoryThreshold,
|
||||
"enabled": rlm.enabled,
|
||||
}
|
||||
}
|
||||
|
||||
func (rlm *RuntimeLimitsManager) Stop() {
|
||||
close(rlm.stopChan)
|
||||
rlm.wg.Wait()
|
||||
}
|
||||
@@ -0,0 +1,705 @@
|
||||
/*
|
||||
* Copyright 2026 Safronov Grigorii
|
||||
*
|
||||
* Licensed under the CDDL, Version 1.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
*
|
||||
* You may obtain a copy of the License at
|
||||
* https://opensource.org/licenses/CDDL-1.0
|
||||
*/
|
||||
|
||||
// Файл: internal/storage/skiplist.go
|
||||
// Назначение: Lock-free skip list (список с пропусками) для инвертированного
|
||||
// индекса коллекции. Заменяет sync.Map в Index.data.
|
||||
//
|
||||
// ИСПРАВЛЕНО (2026-10):
|
||||
// - Баг #6: рассинхронизация size при воскрешении узла (Insert после Delete).
|
||||
// Теперь проверяется wasDeleted и size увеличивается корректно.
|
||||
// - Баг #7: дубликат на верхних уровнях при конкурентной вставке одного
|
||||
// ключа с разными newLevel. Решение: двухфазная вставка — сначала
|
||||
// уровень 0 (гарантированно единственный узел с таким ключом), затем
|
||||
// верхние уровни (best-effort, дубликаты безопасны).
|
||||
//
|
||||
// КЛЮЧЕВОЙ ИНВАРИАНТ:
|
||||
// - На уровне 0 не может быть двух узлов с одинаковым ключом.
|
||||
// - На уровнях > 0 дубликаты допустимы (они лишь замедляют поиск, но не
|
||||
// ломают корректность, потому что Search всегда спускается до уровня 0
|
||||
// и там находит единственный узел с нужным ключом).
|
||||
//
|
||||
// ДРУГИЕ ИСПРАВЛЕНИЯ (2026):
|
||||
// - Устранён data race при обновлении value существующего узла:
|
||||
// value хранится в atomic.Value.
|
||||
// - Исправлен обход Range: спуск на базовый уровень теперь корректный.
|
||||
// - Исправлена вставка на уровнях выше текущего sl.level.
|
||||
// - Идемпотентный Delete.
|
||||
// - RangeWithKeyPrefix корректно останавливается при выходе за префикс.
|
||||
|
||||
package storage
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"math/rand"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
)
|
||||
|
||||
// =============================================================================
|
||||
// КОНСТАНТЫ
|
||||
// =============================================================================
|
||||
|
||||
const (
|
||||
// Максимальное количество уровней в skip list.
|
||||
maxSkipListLevel = 32
|
||||
|
||||
// Вероятность "подъёма" на следующий уровень.
|
||||
skipListProbability = 0.5
|
||||
|
||||
// Порог, после которого Range запускает физическую очистку
|
||||
// помеченных (tombstone) узлов.
|
||||
skipListSweepThreshold = 1024
|
||||
)
|
||||
|
||||
// =============================================================================
|
||||
// ТИПЫ
|
||||
// =============================================================================
|
||||
|
||||
// skipListNode представляет узел skip list.
|
||||
type skipListNode struct {
|
||||
// key — значение индексированного поля.
|
||||
key interface{}
|
||||
|
||||
// value — ID документа. Хранится в atomic.Value для защиты от
|
||||
// data race при обновлении существующего ключа.
|
||||
// atomic.Value не может хранить nil строку напрямую, поэтому
|
||||
// используется обёртка *string.
|
||||
value atomic.Value // *string
|
||||
|
||||
// deleted — tombstone-флаг.
|
||||
deleted atomic.Bool
|
||||
|
||||
// forward — массив атомарных указателей на следующий узел
|
||||
// на каждом уровне.
|
||||
forward []atomic.Value // *skipListNode
|
||||
}
|
||||
|
||||
// loadValue атомарно загружает value.
|
||||
func (n *skipListNode) loadValue() string {
|
||||
val := n.value.Load()
|
||||
if val == nil {
|
||||
return ""
|
||||
}
|
||||
ptr, ok := val.(*string)
|
||||
if !ok || ptr == nil {
|
||||
return ""
|
||||
}
|
||||
return *ptr
|
||||
}
|
||||
|
||||
// storeValue атомарно сохраняет value.
|
||||
func (n *skipListNode) storeValue(v string) {
|
||||
n.value.Store(&v)
|
||||
}
|
||||
|
||||
// loadForward атомарно загружает указатель на следующий узел на уровне.
|
||||
func (n *skipListNode) loadForward(level int) *skipListNode {
|
||||
if level < 0 || level >= len(n.forward) {
|
||||
return nil
|
||||
}
|
||||
val := n.forward[level].Load()
|
||||
if val == nil {
|
||||
return nil
|
||||
}
|
||||
ptr, ok := val.(*skipListNode)
|
||||
if !ok {
|
||||
return nil
|
||||
}
|
||||
return ptr
|
||||
}
|
||||
|
||||
// storeForward атомарно сохраняет указатель на следующий узел на уровне.
|
||||
func (n *skipListNode) storeForward(level int, next *skipListNode) {
|
||||
if level < 0 || level >= len(n.forward) {
|
||||
return
|
||||
}
|
||||
n.forward[level].Store(next)
|
||||
}
|
||||
|
||||
// casForward атомарно заменяет указатель, если он равен old.
|
||||
func (n *skipListNode) casForward(level int, old, new *skipListNode) bool {
|
||||
if level < 0 || level >= len(n.forward) {
|
||||
return false
|
||||
}
|
||||
return n.forward[level].CompareAndSwap(old, new)
|
||||
}
|
||||
|
||||
// SkipList — lock-free список с пропусками.
|
||||
type SkipList struct {
|
||||
head *skipListNode
|
||||
level atomic.Int32
|
||||
size atomic.Int64
|
||||
|
||||
rngMu sync.Mutex
|
||||
rng *rand.Rand
|
||||
}
|
||||
|
||||
// NewSkipList создаёт новый пустой skip list.
|
||||
func NewSkipList() *SkipList {
|
||||
head := &skipListNode{
|
||||
key: nil,
|
||||
forward: make([]atomic.Value, maxSkipListLevel),
|
||||
}
|
||||
// Инициализируем forward-указатели типизированным nil.
|
||||
for i := 0; i < maxSkipListLevel; i++ {
|
||||
head.forward[i].Store((*skipListNode)(nil))
|
||||
}
|
||||
// Инициализируем value head-а, чтобы atomic.Value имел
|
||||
// установленный concrete type.
|
||||
emptyStr := ""
|
||||
head.value.Store(&emptyStr)
|
||||
|
||||
sl := &SkipList{
|
||||
head: head,
|
||||
rng: rand.New(rand.NewSource(int64(0x5eed5eed))),
|
||||
}
|
||||
sl.level.Store(0)
|
||||
sl.size.Store(0)
|
||||
return sl
|
||||
}
|
||||
|
||||
// randomLevel генерирует случайный уровень для нового узла.
|
||||
func (sl *SkipList) randomLevel() int {
|
||||
sl.rngMu.Lock()
|
||||
defer sl.rngMu.Unlock()
|
||||
|
||||
level := 0
|
||||
for level < maxSkipListLevel-1 && sl.rng.Float64() < skipListProbability {
|
||||
level++
|
||||
}
|
||||
return level
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// ПОИСК ПРЕДШЕСТВЕННИКОВ (вспомогательный метод)
|
||||
// =============================================================================
|
||||
|
||||
// findPredecessors заполняет update[] указателями на узлы, после которых
|
||||
// нужно вставлять/искать узел с заданным ключом на каждом уровне.
|
||||
//
|
||||
// Возвращает:
|
||||
// - update[] — массив предшественников (длина maxSkipListLevel)
|
||||
// - existing — узел на уровне 0 с ключом == key, либо nil
|
||||
//
|
||||
// Этот метод используется и в Insert, и в Delete.
|
||||
func (sl *SkipList) findPredecessors(key interface{}) ([]*skipListNode, *skipListNode) {
|
||||
update := make([]*skipListNode, maxSkipListLevel)
|
||||
|
||||
currentLevel := int(sl.level.Load())
|
||||
|
||||
// Спускаемся с верхнего уровня до базового.
|
||||
x := sl.head
|
||||
for i := currentLevel; i >= 0; i-- {
|
||||
for {
|
||||
next := x.loadForward(i)
|
||||
if next == nil {
|
||||
break
|
||||
}
|
||||
if next.deleted.Load() {
|
||||
// Пытаемся физически вытолкнуть tombstone.
|
||||
if x.casForward(i, next, next.loadForward(i)) {
|
||||
continue
|
||||
}
|
||||
continue
|
||||
}
|
||||
if compareSkipListKeys(next.key, key) < 0 {
|
||||
x = next
|
||||
continue
|
||||
}
|
||||
break
|
||||
}
|
||||
update[i] = x
|
||||
}
|
||||
|
||||
// Для уровней выше currentLevel предшественник — head.
|
||||
// (Это может понадобиться, если новый уровень больше текущего.)
|
||||
for i := currentLevel + 1; i < maxSkipListLevel; i++ {
|
||||
update[i] = sl.head
|
||||
}
|
||||
|
||||
// Ищем существующий узел на уровне 0.
|
||||
existing := update[0].loadForward(0)
|
||||
for existing != nil && compareSkipListKeys(existing.key, key) < 0 {
|
||||
existing = existing.loadForward(0)
|
||||
}
|
||||
if existing != nil && compareSkipListKeys(existing.key, key) != 0 {
|
||||
existing = nil
|
||||
}
|
||||
|
||||
return update, existing
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// ВСТАВКА (ИСПРАВЛЕНО: баги #6 и #7)
|
||||
// =============================================================================
|
||||
|
||||
// Insert добавляет пару (key, value) в skip list.
|
||||
// Если ключ уже существует, значение перезаписывается.
|
||||
//
|
||||
// РЕАЛИЗАЦИЯ (двухфазная вставка):
|
||||
//
|
||||
// Фаза 1: вставка на уровень 0.
|
||||
// - Если ключ уже есть на уровне 0 — обновляем value (с корректной
|
||||
// обработкой wasDeleted для size).
|
||||
// - Иначе — CAS-вставка нового узла на уровень 0. Если CAS провалился
|
||||
// из-за конкурента — retry.
|
||||
// - После успешной вставки на уровень 0 ключ ГАРАНТИРОВАННО есть в списке
|
||||
// (в единственном экземпляре на уровне 0).
|
||||
//
|
||||
// Фаза 2: вставка на верхние уровни (best-effort).
|
||||
// - Для каждого уровня i от 1 до newLevel пытаемся CAS-вставить узел.
|
||||
// - Если CAS провалился — не страшно: узел уже есть на уровне 0,
|
||||
// а верхние уровни — только ускорители поиска. Дубликаты на верхних
|
||||
// уровнях безопасны, потому что Search всегда спускается до уровня 0.
|
||||
func (sl *SkipList) Insert(key interface{}, value string) {
|
||||
if key == nil {
|
||||
return
|
||||
}
|
||||
|
||||
// Внешний цикл — retry при провале CAS на уровне 0.
|
||||
for attempt := 0; attempt < 1000; attempt++ {
|
||||
newLevel := sl.randomLevel()
|
||||
|
||||
// Поднимаем максимальный уровень, если нужно.
|
||||
for {
|
||||
currentLevel := int(sl.level.Load())
|
||||
if newLevel <= currentLevel {
|
||||
break
|
||||
}
|
||||
if sl.level.CompareAndSwap(int32(currentLevel), int32(newLevel)) {
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
// Находим предшественников и существующий узел.
|
||||
update, existing := sl.findPredecessors(key)
|
||||
|
||||
// Случай 1: ключ уже есть на уровне 0.
|
||||
// Обновляем value и, если узел был tombstone, воскрешаем его.
|
||||
if existing != nil {
|
||||
wasDeleted := existing.deleted.Load()
|
||||
existing.storeValue(value)
|
||||
existing.deleted.Store(false)
|
||||
|
||||
// ИСПРАВЛЕНО (баг #6): если узел был tombstone,
|
||||
// увеличиваем size, потому что Delete его уже уменьшил.
|
||||
if wasDeleted {
|
||||
sl.size.Add(1)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
// Случай 2: ключа нет — создаём новый узел и вставляем
|
||||
// сначала на уровень 0.
|
||||
newNode := &skipListNode{
|
||||
key: key,
|
||||
forward: make([]atomic.Value, newLevel+1),
|
||||
}
|
||||
for i := 0; i <= newLevel; i++ {
|
||||
newNode.forward[i].Store((*skipListNode)(nil))
|
||||
}
|
||||
newNode.storeValue(value)
|
||||
|
||||
// Пытаемся вставить на уровень 0.
|
||||
next0 := update[0].loadForward(0)
|
||||
// Защита от ситуации, когда между findPredecessors и сейчас
|
||||
// кто-то вставил узел с таким же ключом.
|
||||
for next0 != nil && compareSkipListKeys(next0.key, key) < 0 {
|
||||
next0 = next0.loadForward(0)
|
||||
}
|
||||
|
||||
// Если между findPredecessors и сейчас появился узел с нашим
|
||||
// ключом — обновляем его value и выходим.
|
||||
if next0 != nil && compareSkipListKeys(next0.key, key) == 0 {
|
||||
wasDeleted := next0.deleted.Load()
|
||||
next0.storeValue(value)
|
||||
next0.deleted.Store(false)
|
||||
if wasDeleted {
|
||||
sl.size.Add(1)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
newNode.storeForward(0, next0)
|
||||
if !update[0].casForward(0, next0, newNode) {
|
||||
// CAS провалился — кто-то другой вставил узел между
|
||||
// findPredecessors и сейчас. Retry.
|
||||
continue
|
||||
}
|
||||
|
||||
// Узел успешно вставлен на уровень 0. Теперь увеличиваем size
|
||||
// (единственное место, где мы это делаем для новых узлов).
|
||||
sl.size.Add(1)
|
||||
|
||||
// ============================================================
|
||||
// ФАЗА 2: вставка на верхние уровни (best-effort).
|
||||
// ============================================================
|
||||
//
|
||||
// Дубликаты на верхних уровнях безопасны: Search всегда спускается
|
||||
// до уровня 0 и там находит единственный узел с нужным ключом.
|
||||
// Если CAS на верхнем уровне провалился — просто пропускаем уровень.
|
||||
for i := 1; i <= newLevel; i++ {
|
||||
if update[i] == nil {
|
||||
// Предшественник не определён — пропускаем уровень.
|
||||
continue
|
||||
}
|
||||
// Перечитываем next на этом уровне (мог измениться).
|
||||
next := update[i].loadForward(i)
|
||||
newNode.storeForward(i, next)
|
||||
// CAS best-effort: если провалился — не retry, просто пропускаем.
|
||||
update[i].casForward(i, next, newNode)
|
||||
}
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
// Если за 1000 попыток не удалось — что-то не так.
|
||||
// Молча выходим (не паникуем).
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// ПОИСК
|
||||
// =============================================================================
|
||||
|
||||
// Search ищет значение по ключу.
|
||||
func (sl *SkipList) Search(key interface{}) (string, bool) {
|
||||
if key == nil {
|
||||
return "", false
|
||||
}
|
||||
|
||||
x := sl.head
|
||||
for i := int(sl.level.Load()); i >= 0; i-- {
|
||||
for {
|
||||
next := x.loadForward(i)
|
||||
if next == nil {
|
||||
break
|
||||
}
|
||||
if next.deleted.Load() {
|
||||
// Пропускаем tombstone.
|
||||
if x.casForward(i, next, next.loadForward(i)) {
|
||||
continue
|
||||
}
|
||||
continue
|
||||
}
|
||||
if compareSkipListKeys(next.key, key) < 0 {
|
||||
x = next
|
||||
continue
|
||||
}
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
candidate := x.loadForward(0)
|
||||
if candidate == nil {
|
||||
return "", false
|
||||
}
|
||||
if compareSkipListKeys(candidate.key, key) == 0 {
|
||||
if candidate.deleted.Load() {
|
||||
return "", false
|
||||
}
|
||||
return candidate.loadValue(), true
|
||||
}
|
||||
return "", false
|
||||
}
|
||||
|
||||
// SearchAll возвращает все значения, ключ которых равен заданному.
|
||||
//
|
||||
// ВАЖНО: текущая реализация индекса хранит только один docID на ключ,
|
||||
// поэтому метод возвращает 0 или 1 значение. Для полноценной поддержки
|
||||
// неуникальных индексов потребуется изменить структуру узла.
|
||||
func (sl *SkipList) SearchAll(key interface{}) []string {
|
||||
val, ok := sl.Search(key)
|
||||
if !ok {
|
||||
return nil
|
||||
}
|
||||
return []string{val}
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// УДАЛЕНИЕ
|
||||
// =============================================================================
|
||||
|
||||
// Delete помечает узел с заданным ключом как удалённый (tombstone).
|
||||
// Идемпотентен: повторный вызов не уменьшает size.
|
||||
func (sl *SkipList) Delete(key interface{}) {
|
||||
if key == nil {
|
||||
return
|
||||
}
|
||||
|
||||
_, existing := sl.findPredecessors(key)
|
||||
if existing == nil {
|
||||
return
|
||||
}
|
||||
|
||||
// CAS как атомарная проверка+установка.
|
||||
// Уменьшаем size только если мы первыми поставили tombstone.
|
||||
if existing.deleted.CompareAndSwap(false, true) {
|
||||
sl.size.Add(-1)
|
||||
}
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// ОБХОД (RANGE)
|
||||
// =============================================================================
|
||||
|
||||
// Range вызывает fn для каждой пары (key, value) в порядке возрастания ключей.
|
||||
// Если fn возвращает false — обход прекращается.
|
||||
func (sl *SkipList) Range(fn func(key interface{}, value string) bool) {
|
||||
prev := sl.head
|
||||
tombstoneCount := 0
|
||||
|
||||
for {
|
||||
next := prev.loadForward(0)
|
||||
if next == nil {
|
||||
break
|
||||
}
|
||||
|
||||
if next.deleted.Load() {
|
||||
tombstoneCount++
|
||||
// Пытаемся физически удалить узел из базового уровня.
|
||||
if prev.casForward(0, next, next.loadForward(0)) {
|
||||
if tombstoneCount >= skipListSweepThreshold {
|
||||
sl.sweepUpperLevels()
|
||||
tombstoneCount = 0
|
||||
}
|
||||
continue
|
||||
}
|
||||
// CAS не удался — двигаемся дальше.
|
||||
prev = next
|
||||
continue
|
||||
}
|
||||
|
||||
if !fn(next.key, next.loadValue()) {
|
||||
return
|
||||
}
|
||||
|
||||
prev = next
|
||||
}
|
||||
}
|
||||
|
||||
// RangeWithKeyPrefix вызывает fn для каждой пары, ключ которой является
|
||||
// строкой с заданным префиксом.
|
||||
//
|
||||
// Благодаря сортированности skip list, обход можно прервать, как только
|
||||
// ключ перестал соответствовать префиксу.
|
||||
func (sl *SkipList) RangeWithKeyPrefix(prefix string, fn func(key interface{}, value string) bool) {
|
||||
sl.Range(func(key interface{}, value string) bool {
|
||||
keyStr, ok := key.(string)
|
||||
if !ok {
|
||||
return true
|
||||
}
|
||||
if strings.HasPrefix(keyStr, prefix) {
|
||||
return fn(key, value)
|
||||
}
|
||||
// Если ключ — строка и он больше префикса, дальнейшие ключи
|
||||
// (отсортированные) тоже не подойдут.
|
||||
if keyStr > prefix {
|
||||
return false
|
||||
}
|
||||
return true
|
||||
})
|
||||
}
|
||||
|
||||
// sweepUpperLevels проходит по верхним уровням и физически выталкивает
|
||||
// tombstone-узлы.
|
||||
func (sl *SkipList) sweepUpperLevels() {
|
||||
for i := 1; i < maxSkipListLevel; i++ {
|
||||
x := sl.head
|
||||
for {
|
||||
next := x.loadForward(i)
|
||||
if next == nil {
|
||||
break
|
||||
}
|
||||
if next.deleted.Load() {
|
||||
if x.casForward(i, next, next.loadForward(i)) {
|
||||
continue
|
||||
}
|
||||
}
|
||||
x = next
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// СЛУЖЕБНЫЕ МЕТОДЫ
|
||||
// =============================================================================
|
||||
|
||||
// Size возвращает приблизительное количество активных элементов.
|
||||
func (sl *SkipList) Size() int64 { return sl.size.Load() }
|
||||
|
||||
// Len — синоним Size.
|
||||
func (sl *SkipList) Len() int64 { return sl.size.Load() }
|
||||
|
||||
// IsEmpty возвращает true, если список пуст.
|
||||
func (sl *SkipList) IsEmpty() bool { return sl.size.Load() == 0 }
|
||||
|
||||
// Clear помечает все узлы как удалённые.
|
||||
func (sl *SkipList) Clear() {
|
||||
sl.Range(func(key interface{}, value string) bool {
|
||||
sl.Delete(key)
|
||||
return true
|
||||
})
|
||||
sl.size.Store(0)
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// СРАВНЕНИЕ КЛЮЧЕЙ
|
||||
// =============================================================================
|
||||
|
||||
// compareSkipListKeys сравнивает два ключа.
|
||||
func compareSkipListKeys(a, b interface{}) int {
|
||||
if a == nil && b == nil {
|
||||
return 0
|
||||
}
|
||||
if a == nil {
|
||||
return -1
|
||||
}
|
||||
if b == nil {
|
||||
return 1
|
||||
}
|
||||
|
||||
switch av := a.(type) {
|
||||
case string:
|
||||
if bv, ok := b.(string); ok {
|
||||
return strings.Compare(av, bv)
|
||||
}
|
||||
case int:
|
||||
if bv, ok := toInt64(b); ok {
|
||||
return compareInt64(int64(av), bv)
|
||||
}
|
||||
case int64:
|
||||
if bv, ok := toInt64(b); ok {
|
||||
return compareInt64(av, bv)
|
||||
}
|
||||
case float64:
|
||||
if bv, ok := toFloat64Comparable(b); ok {
|
||||
if av < bv {
|
||||
return -1
|
||||
}
|
||||
if av > bv {
|
||||
return 1
|
||||
}
|
||||
return 0
|
||||
}
|
||||
case bool:
|
||||
if bv, ok := b.(bool); ok {
|
||||
if av == bv {
|
||||
return 0
|
||||
}
|
||||
if !av {
|
||||
return -1
|
||||
}
|
||||
return 1
|
||||
}
|
||||
case []byte:
|
||||
if bv, ok := b.([]byte); ok {
|
||||
return strings.Compare(string(av), string(bv))
|
||||
}
|
||||
}
|
||||
|
||||
return strings.Compare(fmt.Sprintf("%v", a), fmt.Sprintf("%v", b))
|
||||
}
|
||||
|
||||
func toInt64(v interface{}) (int64, bool) {
|
||||
switch x := v.(type) {
|
||||
case int:
|
||||
return int64(x), true
|
||||
case int32:
|
||||
return int64(x), true
|
||||
case int64:
|
||||
return x, true
|
||||
case uint:
|
||||
return int64(x), true
|
||||
case uint32:
|
||||
return int64(x), true
|
||||
case uint64:
|
||||
return int64(x), true
|
||||
case float64:
|
||||
return int64(x), true
|
||||
case float32:
|
||||
return int64(x), true
|
||||
default:
|
||||
return 0, false
|
||||
}
|
||||
}
|
||||
|
||||
func toFloat64Comparable(v interface{}) (float64, bool) {
|
||||
switch x := v.(type) {
|
||||
case int:
|
||||
return float64(x), true
|
||||
case int32:
|
||||
return float64(x), true
|
||||
case int64:
|
||||
return float64(x), true
|
||||
case float32:
|
||||
return float64(x), true
|
||||
case float64:
|
||||
return x, true
|
||||
default:
|
||||
return 0, false
|
||||
}
|
||||
}
|
||||
|
||||
func compareInt64(a, b int64) int {
|
||||
if a < b {
|
||||
return -1
|
||||
}
|
||||
if a > b {
|
||||
return 1
|
||||
}
|
||||
return 0
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// ITERATOR
|
||||
// =============================================================================
|
||||
|
||||
// SkipListIterator — итератор по skip list.
|
||||
type SkipListIterator struct {
|
||||
sl *SkipList
|
||||
current *skipListNode
|
||||
started bool
|
||||
}
|
||||
|
||||
// NewIterator создаёт новый итератор.
|
||||
func (sl *SkipList) NewIterator() *SkipListIterator {
|
||||
return &SkipListIterator{sl: sl}
|
||||
}
|
||||
|
||||
// Next переходит к следующему элементу.
|
||||
func (it *SkipListIterator) Next() (interface{}, string, bool) {
|
||||
if it.sl == nil {
|
||||
return nil, "", false
|
||||
}
|
||||
|
||||
if !it.started {
|
||||
it.started = true
|
||||
it.current = it.sl.head
|
||||
}
|
||||
|
||||
for {
|
||||
next := it.current.loadForward(0)
|
||||
if next == nil {
|
||||
return nil, "", false
|
||||
}
|
||||
it.current = next
|
||||
if !next.deleted.Load() {
|
||||
return next.key, next.loadValue(), true
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Close освобождает ресурсы итератора.
|
||||
func (it *SkipListIterator) Close() {
|
||||
it.sl = nil
|
||||
it.current = nil
|
||||
}
|
||||
File diff suppressed because it is too large.
Load diff
Reference in new issue
Block a user