661 lines
23 KiB
Go
661 lines
23 KiB
Go
/*
|
||
* 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.
|
||
//
|
||
// ОСОБЕННОСТИ РЕАЛИЗАЦИИ:
|
||
// - Полностью lock-free для операций чтения и записи.
|
||
// - Использует atomic.Value и sync/atomic для совместимости с Go < 1.19
|
||
// (актуально для сборки на OpenIndiana/illumos, где часто стоит старый Go).
|
||
// - Не использует unsafe и generics — одинаково работает на Linux и illumos.
|
||
// - Удаление через mark-and-sweep: узел помечается как удалённый (tombstone),
|
||
// физическое удаление происходит лениво при обходе (Range).
|
||
// - Хранение: ключ — interface{} (значение индексированного поля),
|
||
// значение — string (ID документа).
|
||
//
|
||
// АЛГОРИТМ (упрощённый lock-free skip list, основанный на идеях
|
||
// William Pugh "Skip Lists: A Probabilistic Alternative to Balanced Trees"
|
||
// и подходах из работ Herlihy & Shavit "The Art of Multiprocessor Programming"):
|
||
//
|
||
// - Каждый узел имеет массив forward-указателей (уровни).
|
||
// - Вставка: находим позицию на каждом уровне, атомарно "вклеиваем"
|
||
// узел через Compare-And-Swap (CAS).
|
||
// - Поиск: спускаемся по уровням сверху вниз.
|
||
// - Удаление: помечаем узел как удалённый (tombstone), обходчик
|
||
// (Range) физически удаляет такие узлы лениво.
|
||
//
|
||
// ВАЖНО: это не строго "wait-free" реализация в академическом смысле —
|
||
// операции вставки используют CAS-цикл (lock-free, но с возможностью
|
||
// повторных попыток). Однако, в отличие от sync.Map, здесь нет
|
||
// глобального мьютекса и нет внутренних блокировок на уровне шардов,
|
||
// что даёт лучшую масштабируемость при конкурентных чтениях/записях.
|
||
|
||
package storage
|
||
|
||
import (
|
||
"fmt"
|
||
"math/rand"
|
||
"strings"
|
||
"sync"
|
||
"sync/atomic"
|
||
)
|
||
|
||
// =============================================================================
|
||
// КОНСТАНТЫ
|
||
// =============================================================================
|
||
|
||
const (
|
||
// Максимальное количество уровней в skip list.
|
||
// 32 уровня достаточно для хранения до 4 миллиардов элементов
|
||
// с вероятностью 1/2 на каждом уровне.
|
||
maxSkipListLevel = 32
|
||
|
||
// Вероятность "подъёма" на следующий уровень.
|
||
skipListProbability = 0.5
|
||
|
||
// Порог, после которого Range запускает физическую очистку
|
||
// помеченных (tombstone) узлов.
|
||
skipListSweepThreshold = 1024
|
||
)
|
||
|
||
// =============================================================================
|
||
// ТИПЫ
|
||
// =============================================================================
|
||
|
||
// skipListNode представляет узел skip list.
|
||
type skipListNode struct {
|
||
// key — значение индексированного поля (может быть любым comparable value).
|
||
key interface{}
|
||
|
||
// value — ID документа (строка).
|
||
value string
|
||
|
||
// deleted — tombstone-флаг. Если true, узел логически удалён
|
||
// и должен быть пропущен при обходе.
|
||
deleted atomic.Bool
|
||
|
||
// forward — массив атомарных указателей на следующий узел
|
||
// на каждом уровне. forward[0] — базовый уровень.
|
||
// Используем atomic.Value для совместимости с Go < 1.19.
|
||
forward []atomic.Value // хранит *skipListNode
|
||
}
|
||
|
||
// 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
|
||
}
|
||
return val.(*skipListNode)
|
||
}
|
||
|
||
// storeForward атомарно сохраняет указатель на следующий узел на уровне.
|
||
func (n *skipListNode) storeForward(level int, next *skipListNode) {
|
||
if level < 0 || level >= len(n.forward) {
|
||
return
|
||
}
|
||
// atomic.Value не может хранить nil — используем типизированный указатель.
|
||
if next == nil {
|
||
n.forward[level].Store((*skipListNode)(nil))
|
||
} else {
|
||
n.forward[level].Store(next)
|
||
}
|
||
}
|
||
|
||
// casForward атомарно заменяет указатель, если он равен old.
|
||
// Возвращает true, если замена произошла.
|
||
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 список с пропусками для инвертированного индекса.
|
||
//
|
||
// Поля:
|
||
// - head: фиктивный головной узел, key == nil.
|
||
// - level: текущий максимальный уровень (атомарный).
|
||
// - size: приблизительное количество элементов (атомарный счётчик).
|
||
// - rngMu: защита rand.Rand (используется только при вставке).
|
||
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,
|
||
value: "",
|
||
forward: make([]atomic.Value, maxSkipListLevel),
|
||
}
|
||
// Инициализируем forward-указатели nil-ами (atomic.Value требует
|
||
// типизированного значения, чтобы Load не падал).
|
||
for i := 0; i < maxSkipListLevel; i++ {
|
||
head.forward[i].Store((*skipListNode)(nil))
|
||
}
|
||
|
||
sl := &SkipList{
|
||
head: head,
|
||
rng: rand.New(rand.NewSource(int64(0x5eed5eed))),
|
||
}
|
||
sl.level.Store(0)
|
||
sl.size.Store(0)
|
||
return sl
|
||
}
|
||
|
||
// randomLevel генерирует случайный уровень для нового узла.
|
||
// Использует геометрическое распределение с p = 0.5.
|
||
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
|
||
}
|
||
|
||
// =============================================================================
|
||
// ВСТАВКА
|
||
// =============================================================================
|
||
|
||
// Insert добавляет пару (key, value) в skip list.
|
||
// Если ключ уже существует, значение перезаписывается.
|
||
//
|
||
// Реализация lock-free: используется CAS-цикл на каждом уровне.
|
||
// При неудаче CAS — повторная попытка.
|
||
func (sl *SkipList) Insert(key interface{}, value string) {
|
||
if key == nil {
|
||
return
|
||
}
|
||
|
||
newLevel := sl.randomLevel()
|
||
|
||
// Поднимаем максимальный уровень, если нужно.
|
||
for {
|
||
currentLevel := int(sl.level.Load())
|
||
if newLevel <= currentLevel {
|
||
break
|
||
}
|
||
if sl.level.CompareAndSwap(int32(currentLevel), int32(newLevel)) {
|
||
break
|
||
}
|
||
}
|
||
|
||
newNode := &skipListNode{
|
||
key: key,
|
||
value: value,
|
||
forward: make([]atomic.Value, newLevel+1),
|
||
}
|
||
for i := 0; i <= newLevel; i++ {
|
||
newNode.forward[i].Store((*skipListNode)(nil))
|
||
}
|
||
|
||
// update[i] — узел, после которого нужно вставить новый
|
||
// на уровне i.
|
||
update := make([]*skipListNode, maxSkipListLevel)
|
||
|
||
retry:
|
||
// Спускаемся с верхнего уровня до базового, запоминая предшественников.
|
||
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
|
||
}
|
||
// CAS не удался — перечитываем next.
|
||
continue
|
||
}
|
||
if compareSkipListKeys(next.key, key) < 0 {
|
||
x = next
|
||
continue
|
||
}
|
||
break
|
||
}
|
||
update[i] = x
|
||
}
|
||
|
||
// Теперь update[0].forward[0] — либо nil, либо узел с key >= newKey.
|
||
existing := update[0].loadForward(0)
|
||
for existing != nil && compareSkipListKeys(existing.key, key) < 0 {
|
||
existing = existing.loadForward(0)
|
||
}
|
||
|
||
// Если ключ уже есть — обновляем значение (и снимаем tombstone).
|
||
if existing != nil && compareSkipListKeys(existing.key, key) == 0 {
|
||
existing.value = value
|
||
existing.deleted.Store(false)
|
||
return
|
||
}
|
||
|
||
// Вставляем newNode на всех уровнях от 0 до newLevel.
|
||
for i := 0; i <= newLevel; i++ {
|
||
for {
|
||
next := update[i].loadForward(i)
|
||
newNode.forward[i].Store(next)
|
||
if update[i].casForward(i, next, newNode) {
|
||
break
|
||
}
|
||
// CAS не удался — кто-то другой вставил узел.
|
||
// Перечитываем update[i] и повторяем.
|
||
// Упрощение: перезапускаем всю вставку.
|
||
goto retry
|
||
}
|
||
}
|
||
|
||
sl.size.Add(1)
|
||
}
|
||
|
||
// =============================================================================
|
||
// ПОИСК
|
||
// =============================================================================
|
||
|
||
// Search ищет значение по ключу. Возвращает (value, true) при успехе.
|
||
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 compareSkipListKeys(next.key, key) < 0 {
|
||
x = next
|
||
continue
|
||
}
|
||
break
|
||
}
|
||
}
|
||
|
||
// x — предшественник искомого ключа на базовом уровне.
|
||
candidate := x.loadForward(0)
|
||
if candidate == nil {
|
||
return "", false
|
||
}
|
||
if compareSkipListKeys(candidate.key, key) == 0 {
|
||
if candidate.deleted.Load() {
|
||
return "", false
|
||
}
|
||
return candidate.value, true
|
||
}
|
||
return "", false
|
||
}
|
||
|
||
// SearchAll возвращает все значения, ключ которых равен заданному.
|
||
// В неуникальном индексе по одному ключу может быть несколько документов —
|
||
// в текущей реализации индекс хранит только последний docID для данного
|
||
// ключа (как и sync.Map). Для уникальных индексов это не критично.
|
||
// Для неуникальных индексов метод возвращает единственное значение.
|
||
func (sl *SkipList) SearchAll(key interface{}) []string {
|
||
val, ok := sl.Search(key)
|
||
if !ok {
|
||
return nil
|
||
}
|
||
return []string{val}
|
||
}
|
||
|
||
// =============================================================================
|
||
// УДАЛЕНИЕ
|
||
// =============================================================================
|
||
|
||
// Delete помечает узел с заданным ключом как удалённый (tombstone).
|
||
// Физическое удаление произойдёт лениво при следующем Range.
|
||
func (sl *SkipList) Delete(key interface{}) {
|
||
if key == nil {
|
||
return
|
||
}
|
||
|
||
x := sl.head
|
||
for i := int(sl.level.Load()); i >= 0; i-- {
|
||
for {
|
||
next := x.loadForward(i)
|
||
if next == nil {
|
||
break
|
||
}
|
||
if compareSkipListKeys(next.key, key) < 0 {
|
||
x = next
|
||
continue
|
||
}
|
||
break
|
||
}
|
||
}
|
||
|
||
candidate := x.loadForward(0)
|
||
if candidate == nil {
|
||
return
|
||
}
|
||
if compareSkipListKeys(candidate.key, key) == 0 {
|
||
if !candidate.deleted.Load() {
|
||
if candidate.deleted.CompareAndSwap(false, true) {
|
||
sl.size.Add(-1)
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
// =============================================================================
|
||
// ОБХОД (RANGE)
|
||
// =============================================================================
|
||
|
||
// Range вызывает fn для каждой пары (key, value) в порядке возрастания ключей.
|
||
// Если fn возвращает false — обход прекращается.
|
||
//
|
||
// Во время обхода выполняется ленивая физическая очистка tombstone-узлов:
|
||
// если подряд встретилось больше skipListSweepThreshold удалённых узлов,
|
||
// они физически выталкиваются из списка.
|
||
func (sl *SkipList) Range(fn func(key interface{}, value string) bool) {
|
||
x := sl.head
|
||
level0 := 0
|
||
|
||
// Спускаемся на базовый уровень.
|
||
for i := int(sl.level.Load()); i > 0; i-- {
|
||
for {
|
||
next := x.loadForward(i)
|
||
if next == nil {
|
||
break
|
||
}
|
||
x = next
|
||
}
|
||
}
|
||
_ = level0
|
||
|
||
prev := x
|
||
tombstoneCount := 0
|
||
|
||
for {
|
||
next := prev.loadForward(0)
|
||
if next == nil {
|
||
break
|
||
}
|
||
|
||
if next.deleted.Load() {
|
||
tombstoneCount++
|
||
// Пытаемся физически удалить узел из базового уровня.
|
||
if prev.casForward(0, next, next.loadForward(0)) {
|
||
// Успешно вытолкнули — не двигаем prev.
|
||
if tombstoneCount >= skipListSweepThreshold {
|
||
// Периодически чистим верхние уровни (lazy).
|
||
sl.sweepUpperLevels()
|
||
tombstoneCount = 0
|
||
}
|
||
continue
|
||
}
|
||
// CAS не удался — двигаемся дальше.
|
||
prev = next
|
||
continue
|
||
}
|
||
|
||
if !fn(next.key, next.value) {
|
||
return
|
||
}
|
||
|
||
prev = next
|
||
}
|
||
}
|
||
|
||
// RangeWithKeyPrefix вызывает fn для каждой пары, ключ которой является
|
||
// строкой с заданным префиксом.
|
||
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)
|
||
}
|
||
return true
|
||
})
|
||
}
|
||
|
||
// sweepUpperLevels проходит по верхним уровням и физически выталкивает
|
||
// tombstone-узлы. Вызывается лениво из Range.
|
||
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 для совместимости с sync.Map-стилем).
|
||
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 сравнивает два ключа.
|
||
// Поддерживает основные типы, используемые в инвертированном индексе:
|
||
// string, int, int64, float64, bool, []byte.
|
||
//
|
||
// Возвращает -1, 0 или 1.
|
||
//
|
||
// ВАЖНО: если типы ключей несовместимы, используется строковое сравнение
|
||
// через fmt.Sprintf, чтобы обеспечить хоть какой-то порядок (не падать).
|
||
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))
|
||
}
|
||
}
|
||
|
||
// Fallback: строковое сравнение.
|
||
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 переходит к следующему элементу. Возвращает (key, value, ok).
|
||
func (it *SkipListIterator) Next() (interface{}, string, bool) {
|
||
if it.sl == nil {
|
||
return nil, "", false
|
||
}
|
||
|
||
if !it.started {
|
||
it.started = true
|
||
// Спускаемся на базовый уровень.
|
||
x := it.sl.head
|
||
for i := int(it.sl.level.Load()); i > 0; i-- {
|
||
for {
|
||
next := x.loadForward(i)
|
||
if next == nil {
|
||
break
|
||
}
|
||
x = next
|
||
}
|
||
}
|
||
it.current = x
|
||
}
|
||
|
||
for {
|
||
next := it.current.loadForward(0)
|
||
if next == nil {
|
||
return nil, "", false
|
||
}
|
||
it.current = next
|
||
if !next.deleted.Load() {
|
||
return next.key, next.value, true
|
||
}
|
||
}
|
||
}
|
||
|
||
// Close освобождает ресурсы итератора (no-op, оставлен для совместимости).
|
||
func (it *SkipListIterator) Close() {
|
||
it.sl = nil
|
||
it.current = nil
|
||
}
|