Files
futriix/internal/storage/skiplist.go
T

661 lines
23 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
/*
* 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
}