From f68fd88e28d21d39b6885384132c1d178fb78659 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=D0=93=D1=80=D0=B8=D0=B3=D0=BE=D1=80=D0=B8=D0=B9=20=D0=A1?= =?UTF-8?q?=D0=B0=D1=84=D1=80=D0=BE=D0=BD=D0=BE=D0=B2?= Date: Mon, 28 Sep 2026 21:26:04 +0000 Subject: [PATCH] Upload files to "internal/storage" --- internal/storage/skiplist.go | 660 +++++++++++++++++++++++++++++++++++ 1 file changed, 660 insertions(+) create mode 100644 internal/storage/skiplist.go diff --git a/internal/storage/skiplist.go b/internal/storage/skiplist.go new file mode 100644 index 0000000..ae0238b --- /dev/null +++ b/internal/storage/skiplist.go @@ -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 +}