Upload files to "internal/storage"

This commit is contained in:
gvsafronov committed 2026-09-28 21:26:04 +00:00
1 parent 754ed3ed94
commit f68fd88e28
1 file changed
+660
+660
View File
@@ -0,0 +1,660 @@
/*
* 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
}