Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 42 additions & 0 deletions docs/design/2026-09-02-refactor-hll.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
# Refactor ModeHLL layout (no contract change)

**Status:** approved (chat: start with HLL)
**Branch:** `feat/refactor-hll`
**Date:** 2026-09-02

## Problem

HLL store methods still live in `pkg/store/memory.go` (~2400 lines). TopK and CMS already have their own files. Engine HLL is one 190-line file mixing public API, owner apply, and cluster hops. `pkg/hllx.Merge` is exported but unused on the data path (install is LWW replace).

## Non-goals

- No API / proto / flag / version / fan-out / hintID change
- No `hllx` algorithm change
- No decoded cache / Get-Peek flush helper
- Not extracting other modes

## Contract

Unchanged: `HLLAdd` / `HLLCount` / `Delete(name)`, `FlagHLL` + inbox `FlagHLLAdd`, snapshot fan-out, version under store mutex, present-bit, empty-until-delete.

`pkg/hllx.Merge` becomes unexported (`merge`). Engine never called it. Package tests still cover per-register max.

## Approach

Standard Go file layout (one concern per file, same package):

| Before | After |
|--------|--------|
| `pkg/store/memory.go` HLL* | `pkg/store/hll.go` (same as `cms.go` / `topk.go`) |
| `pkg/engine/hll.go` (all) | `hll.go` public verbs; `hll_apply.go` owner write; `hll_cluster.go` inbox/fan-out/GetOrLoad |
| `pkg/hllx.Merge` | `merge` (unexported) |

Rejected: rewrite registers; split `hllx` (150 lines is one file); touch `ApplyPut` order.

## Tests (already exist; keep them)

Existing `pkg/hllx`, `pkg/store` HLL, `pkg/engine` HLL unit + cluster + hint-after-down. No new behavior tests.

## Bench risk

Hot path? No new Get/Peek logic. Extra file split must not add allocs. Local `BenchmarkEngineGetHit` / `BenchmarkStoreGetHit` allocs/op must stay 15 / 2.
1 change: 1 addition & 0 deletions docs/design/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ Do not implement from a draft.
| [2026-08-25-mode-hll.md](./2026-08-25-mode-hll.md) | `ModeHLL` |
| [2026-08-31-mode-topk.md](./2026-08-31-mode-topk.md) | `ModeTopK` |
| [2026-09-01-mode-cms.md](./2026-09-01-mode-cms.md) | `ModeCMS` |
| [2026-09-02-refactor-hll.md](./2026-09-02-refactor-hll.md) | ModeHLL file layout (no contract) |
| [2026-08-25-list-counter-version.md](./2026-08-25-list-counter-version.md) | List/Counter snapshot version |
| [2026-08-13-unify-grpc-error-map.md](./2026-08-13-unify-grpc-error-map.md) | grpcmap |

Expand Down
105 changes: 2 additions & 103 deletions pkg/engine/hll.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,10 +6,10 @@ import (

"github.com/Code0987/supercache/pkg/hllx"
"github.com/Code0987/supercache/pkg/keyspace"
"github.com/Code0987/supercache/pkg/store"
)

// HLLAdd hashes item into a ModeHLL sketch (Redis PFADD). ACK-only.
// HLLAdd hashes item into a ModeHLL sketch (ACK-only; no Redis changed-bool).
// Non-owners forward an inbox FlagHLLAdd; the owner applies and fans a snapshot.
func (e *Engine) HLLAdd(ctx context.Context, keyspaceName, name string, item []byte) error {
if err := ctx.Err(); err != nil {
return err
Expand Down Expand Up @@ -87,104 +87,3 @@ func (e *Engine) hllNeed(ks *ksRuntime) error {
}
return nil
}

func (e *Engine) hllMutViaOwner(ctx context.Context, ks *ksRuntime, name string, item []byte) error {
c := e.clusterSnapshot()
owner, _ := c.Ring.Owner(name)
ent := store.Entry{Value: append([]byte(nil), item...), Flags: store.FlagHLLAdd, Version: 1}
pctx, cancel := e.peerCtx(ctx, ks)
defer cancel()
applied, err := c.Transport.ApplyPut(pctx, owner.Addr, ks.cfg.Name, name, ent, c.Ring.Generation())
if err != nil {
return err
}
if !applied {
return fmt.Errorf("%w: hll add rejected", ErrInvalidArgument)
}
return nil
}

func (e *Engine) hllAddLocal(ks *ksRuntime, name string, item []byte) error {
expire := e.expireAt(ks.cfg.TTL)
max := e.maxValueSize
if ks.cfg.MaxValueSize > 0 {
max = ks.cfg.MaxValueSize
}
cur, _ := ks.store.PeekVersion(name)
gate := cur + 1
applied, tooLarge := ks.store.HLLAdd(name, item, gate, expire, max)
if tooLarge {
return ErrValueTooLarge
}
if !applied {
return fmt.Errorf("%w: hll add rejected", ErrInvalidArgument)
}
ver, _ := ks.store.PeekVersion(name)
ks.observeVersion(name, ver)
e.hllReplicateSnapshot(ks, name, ver, expire)
return nil
}

func (e *Engine) hllReplicateSnapshot(ks *ksRuntime, name string, ver uint64, expire int64) {
ent, ok := ks.store.Peek(name)
if !ok || !ent.IsHLL() {
return
}
e.replicate(ks.cfg.Name, name, store.Entry{
Value: ent.Value,
Version: ver,
ExpireAt: expire,
Flags: store.FlagHLL,
}, false)
}

func (e *Engine) hllFetchOwner(ctx context.Context, ks *ksRuntime, name string) (store.Entry, bool, error) {
c := e.clusterSnapshot()
if c == nil || c.Ring == nil || c.Transport == nil {
return store.Entry{}, false, nil
}
owner, ok := c.Ring.Owner(name)
if !ok || owner.ID == "" || owner.ID == c.SelfID || owner.Addr == "" {
return store.Entry{}, false, nil
}
pctx, cancel := e.peerCtx(ctx, ks)
defer cancel()
res, err := c.Transport.GetOrLoad(pctx, owner.Addr, ks.cfg.Name, name)
if err != nil || !res.Found || !res.Entry.IsHLL() {
return store.Entry{}, false, nil
}
if e.holdsReplica(c, ks, name) {
_ = ks.store.HLLInstall(name, res.Entry.Value, res.Entry.Version, res.Entry.ExpireAt)
}
return res.Entry, true, nil
}

func (e *Engine) applyHLLAdd(ks *ksRuntime, name string, item []byte, expireAt int64) bool {
if len(item) == 0 {
return false
}
if expireAt == 0 {
expireAt = e.expireAt(ks.cfg.TTL)
}
max := e.maxValueSize
if ks.cfg.MaxValueSize > 0 {
max = ks.cfg.MaxValueSize
}
cur, _ := ks.store.PeekVersion(name)
gate := cur + 1
ok, tooLarge := ks.store.HLLAdd(name, item, gate, expireAt, max)
if !ok || tooLarge {
return false
}
ver, _ := ks.store.PeekVersion(name)
ks.observeVersion(name, ver)
e.hllReplicateSnapshot(ks, name, ver, expireAt)
return true
}

func (e *Engine) applyHLLInstall(ks *ksRuntime, name string, blob []byte, version uint64, expireAt int64) bool {
if expireAt == 0 {
expireAt = e.expireAt(ks.cfg.TTL)
}
return ks.store.HLLInstall(name, blob, version, expireAt)
}
59 changes: 59 additions & 0 deletions pkg/engine/hll_apply.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,59 @@
package engine

import "fmt"

// hllAddLocal is the owner / single-node write path.
// The store assigns the stored version; we fan PeekVersion after the write.
func (e *Engine) hllAddLocal(ks *ksRuntime, name string, item []byte) error {
expire := e.expireAt(ks.cfg.TTL)
max := e.maxValueSize
if ks.cfg.MaxValueSize > 0 {
max = ks.cfg.MaxValueSize
}
cur, _ := ks.store.PeekVersion(name)
gate := cur + 1
applied, tooLarge := ks.store.HLLAdd(name, item, gate, expire, max)
if tooLarge {
return ErrValueTooLarge
}
if !applied {
return fmt.Errorf("%w: hll add rejected", ErrInvalidArgument)
}
ver, _ := ks.store.PeekVersion(name)
ks.observeVersion(name, ver)
e.hllReplicateSnapshot(ks, name, ver, expire)
return nil
}

// applyHLLAdd is the owner-inbox ApplyPut of FlagHLLAdd (raw item).
// Non-owners must not reach here (ApplyPut returns applied=false first).
func (e *Engine) applyHLLAdd(ks *ksRuntime, name string, item []byte, expireAt int64) bool {
if len(item) == 0 {
return false
}
if expireAt == 0 {
expireAt = e.expireAt(ks.cfg.TTL)
}
max := e.maxValueSize
if ks.cfg.MaxValueSize > 0 {
max = ks.cfg.MaxValueSize
}
cur, _ := ks.store.PeekVersion(name)
gate := cur + 1
ok, tooLarge := ks.store.HLLAdd(name, item, gate, expireAt, max)
if !ok || tooLarge {
return false
}
ver, _ := ks.store.PeekVersion(name)
ks.observeVersion(name, ver)
e.hllReplicateSnapshot(ks, name, ver, expireAt)
return true
}

// applyHLLInstall is replica / handoff ApplyPut of FlagHLL (LWW replace).
func (e *Engine) applyHLLInstall(ks *ksRuntime, name string, blob []byte, version uint64, expireAt int64) bool {
if expireAt == 0 {
expireAt = e.expireAt(ks.cfg.TTL)
}
return ks.store.HLLInstall(name, blob, version, expireAt)
}
65 changes: 65 additions & 0 deletions pkg/engine/hll_cluster.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,65 @@
package engine

import (
"context"
"fmt"

"github.com/Code0987/supercache/pkg/store"
)

// hllMutViaOwner sends an inbox FlagHLLAdd to the ring owner.
// ACK-only: a rejected apply surfaces as InvalidArgument (no return payload).
func (e *Engine) hllMutViaOwner(ctx context.Context, ks *ksRuntime, name string, item []byte) error {
c := e.clusterSnapshot()
owner, _ := c.Ring.Owner(name)
ent := store.Entry{Value: append([]byte(nil), item...), Flags: store.FlagHLLAdd, Version: 1}
pctx, cancel := e.peerCtx(ctx, ks)
defer cancel()
applied, err := c.Transport.ApplyPut(pctx, owner.Addr, ks.cfg.Name, name, ent, c.Ring.Generation())
if err != nil {
return err
}
if !applied {
return fmt.Errorf("%w: hll add rejected", ErrInvalidArgument)
}
return nil
}

// hllReplicateSnapshot fans the post-write FlagHLL blob to RF−1 replicas.
// hintID is (ks, name), so a later snapshot replaces a pending hint — both items stay.
func (e *Engine) hllReplicateSnapshot(ks *ksRuntime, name string, ver uint64, expire int64) {
ent, ok := ks.store.Peek(name)
if !ok || !ent.IsHLL() {
return
}
e.replicate(ks.cfg.Name, name, store.Entry{
Value: ent.Value,
Version: ver,
ExpireAt: expire,
Flags: store.FlagHLL,
}, false)
}

// hllFetchOwner loads a missing local name from the owner (GetOrLoad).
// RPC / !Found / wrong type → miss + nil error (do not return Unavailable).
// Replicas may install the snapshot; non-replicas do not.
func (e *Engine) hllFetchOwner(ctx context.Context, ks *ksRuntime, name string) (store.Entry, bool, error) {
c := e.clusterSnapshot()
if c == nil || c.Ring == nil || c.Transport == nil {
return store.Entry{}, false, nil
}
owner, ok := c.Ring.Owner(name)
if !ok || owner.ID == "" || owner.ID == c.SelfID || owner.Addr == "" {
return store.Entry{}, false, nil
}
pctx, cancel := e.peerCtx(ctx, ks)
defer cancel()
res, err := c.Transport.GetOrLoad(pctx, owner.Addr, ks.cfg.Name, name)
if err != nil || !res.Found || !res.Entry.IsHLL() {
return store.Entry{}, false, nil
}
if e.holdsReplica(c, ks, name) {
_ = ks.store.HLLInstall(name, res.Entry.Value, res.Entry.Version, res.Entry.ExpireAt)
}
return res.Entry, true, nil
}
5 changes: 3 additions & 2 deletions pkg/hllx/hllx.go
Original file line number Diff line number Diff line change
Expand Up @@ -103,8 +103,9 @@ func Count(regs []byte) uint64 {
return uint64(math.Floor(e + 0.5))
}

// Merge writes per-register max(dst, src) into dst.
func Merge(dst, src []byte) error {
// merge writes per-register max(dst, src) into dst.
// Not used on install (handoff is LWW replace, not max-merge).
func merge(dst, src []byte) error {
if len(dst) != DenseSize || len(src) != DenseSize {
return ErrSize
}
Expand Down
4 changes: 2 additions & 2 deletions pkg/hllx/hllx_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -92,7 +92,7 @@ func TestNewSize(t *testing.T) {
if len(New()) != DenseSize || DenseSize != 12288 {
t.Fatal(len(New()), DenseSize)
}
if err := Merge(New(), []byte{1}); err != ErrSize {
if err := merge(New(), []byte{1}); err != ErrSize {
t.Fatal(err)
}
}
Expand All @@ -101,7 +101,7 @@ func TestMergeMax(t *testing.T) {
a, b := New(), New()
Add(a, []byte("alice"))
Add(b, []byte("bob"))
if err := Merge(a, b); err != nil {
if err := merge(a, b); err != nil {
t.Fatal(err)
}
n := Count(a)
Expand Down
Loading
Loading