diff --git a/CHANGELOG.md b/CHANGELOG.md index 44886800..4b512426 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,6 +1,19 @@ # Changelog All notable changes to this project will be documented in this file. +## Unreleased — next minor release + +### Added + +* Opt-in `ParallelOptions.AutoOverlap` protects matches across all chunk boundaries, including keywords longer than a chunk, using one engine snapshot per call. +* `CacheStats.PresetReloadFailures` and `PresetPollFailures` report cumulative refresh failures, excluding cancellation. + +### Fixed + +* Preset reload waiters respond independently to cancellation. Shared jobs stop when all waiters leave or the instance closes. Snapshot fetching and engine building run outside the state lock; generation checks reject obsolete reloads after writes or invalidation. +* Preset polling reads only the existing version field. Failed refreshes still return search errors and retry on later requests; polling is disabled by default and provides no outage freshness bound. +* Parallel options are copied before normalization. Existing behavior remains the default. No Redis schema change or data migration is required. + ## [v1.5.2](https://github.com/skyoo2003/acor/releases/tag/v1.5.2) - 2026-08-09 ### Changed diff --git a/api/v1-audit.txt b/api/v1-audit.txt index 93df7e79..f7e20a11 100644 --- a/api/v1-audit.txt +++ b/api/v1-audit.txt @@ -24,7 +24,7 @@ field AhoCorasickArgs.DB int fixed acor.go:262 gave '0-15' as if checked; nothin field AhoCorasickArgs.Debug bool ok acor.go:271; newLogger switches the default logger to stdout at acor.go:448 field AhoCorasickArgs.DialTimeout time.Duration ok acor.go:280; carried into every topology through universalOptions (client.go:86) and the hand-built ring (client.go:100), so the shared 'all topologies' preamble holds for it field AhoCorasickArgs.EnableCache bool ok acor.go:283; both documented rejections fire at acor.go:437,503 -field AhoCorasickArgs.InvalidationPollInterval time.Duration ok acor.go:354; read only at redis_backed.go:91 and the poller starts only when > 0 (redis_backed.go:117), so 'disabled by default, Preset mode only' is accurate +field AhoCorasickArgs.InvalidationPollInterval time.Duration ok redis_backed.go:372 reads only version and marks stale; the next search reloads. Disabled by default, Preset only; failures mean the interval cannot bound freshness. field AhoCorasickArgs.Logger Logger ok acor.go:307; a non-nil Logger wins over the default at acor.go:494 field AhoCorasickArgs.MasterName string ok acor.go:254; client.go:27-28 selects the failover client on a non-blank MasterName and client.go:55-57 requires Addrs with it, exactly as documented field AhoCorasickArgs.MaxRetries int ok acor.go:286; client.go:89,104. -1 disabling retries is go-redis's contract, not this package's @@ -51,8 +51,10 @@ field BatchResult.Skipped []string fixed options.go:110 gave only "duplicates in field CacheStats.Hits uint64 ok stats.go:21; re-verdicted after #206, which landed the one-read-per-call behavior the sentence now describes. FindParallelContext, FindIndexParallelContext and FindManyContext each call loadEngine exactly once (context_ops.go:143,188,111), and hit/miss are recorded only inside loadEngine (v2_ops.go:274,282, redis_backed.go:230,239, engine_memo.go:43,46), so writes, Suggest and Info record nothing. TestCacheStatsCountsOneReadPerCall (stats_test.go:61) pins it field CacheStats.LastInvalidationLag time.Duration ok stats.go:74; recordInvalidationLag drops only negatives (stats.go:143) per TestCacheStatsDiscardsNegativeLag, and the listener-only modes match invalidation.go:192 field CacheStats.Misses uint64 fixed stats.go:35 claimed a failed Redis fetch is always a miss; true for preset (redis_backed.go:239) and cached V2 (v2_ops.go:282), false for default V2, which fetches at v2_ops.go:260 and only then reaches the counter. Sentence now names the split; TestCacheStatsFailedFetchByMode pins all three modes -field CacheStats.RebuildDuration time.Duration ok stats.go:62; timeRebuild (stats.go:165) wraps build alone — the Redis fetch happens before it at v2_ops.go:260 and the lock is taken before it at engine_memo.go:40, matching both exclusions -field CacheStats.Rebuilds uint64 ok stats.go:50; starts at 1 in Preset per TestCacheStatsPreset (stats_test.go:214), and coalesced misses share one build at engine_memo.go:39-47, which is the documented Misses-Rebuilds gap +field CacheStats.PresetPollFailures uint64 ok redis_backed.go:372 records failed version reads; TestPresetPollVersionOnlyAndRecovery verifies version-only polling, errors and recovery. +field CacheStats.PresetReloadFailures uint64 ok redis_backed.go:297 records once per shared failed job, excluding cancellation; TestPresetReloadFailureSharedAndRetry verifies one failure for two waiting requests. +field CacheStats.RebuildDuration time.Duration ok redis_backed.go:177 times only engine construction in Preset, excluding Redis reads, materialization and publication. Other modes retain their existing timed callbacks. +field CacheStats.Rebuilds uint64 ok redis_backed.go:177 counts every Preset engine build, including rejected generations. Construction starts at one; shared misses coalesce, local writes also rebuild. field KeywordError.Error error ok options.go:97; carries ErrEmptyKeyword or the write error, batch.go:69,126,141 field KeywordError.Keyword string ok options.go:95; set from the offending keyword, batch.go:68,126,140 field Match.End int ok matches.go:25; exclusive, indexed at matches.go:327 with an m.End >= len bound @@ -82,17 +84,18 @@ field OperationError.Err error ok errors.go:77; the wrapped cause, returned by U field OperationError.Keyword string ok errors.go:73; left empty by newOperationError, which is why Error has the two-branch format at errors.go:83 field OperationError.Op string ok errors.go:71; set from the op argument, errors.go:112 field OperationError.Schema int ok errors.go:75; carries the schema constant passed in, e.g. SchemaV2 at v2_ops.go:114 +field ParallelOptions.AutoOverlap bool ok parallel.go:243 extends right by max(Overlap, longest rune length minus one); context_ops.go:170 filters owned starts and preserves chunk order. TestAutoOverlapSerialParity verifies Unicode and long keywords. field ParallelOptions.Boundary ChunkBoundary ok options.go:68; the zero value is ChunkBoundaryWord (options.go:37), so unset does split on whitespace field ParallelOptions.ChunkSize int fixed options.go:61 promised a DefaultChunkSize fallback; normalizeParallelOptions (parallel.go:89-104) never sets it and context_ops.go:127,177 reject <= 0 with ErrInvalidChunkSize, so the documented default was an error return. TestParallelOptionsHaveNoImpliedDefaults pins it -field ParallelOptions.Overlap int fixed options.go:71 promised DefaultOverlap; parallel.go:96-102 only clamps negatives, leaving an unset Overlap at 0 and silently missing boundary-straddling keywords. Same test pins it -field ParallelOptions.Workers int ok options.go:58; parallel.go:93-95 substitutes runtime.NumCPU() when <= 0, exactly as documented +field ParallelOptions.Overlap int ok parallel.go:90 copies options and clamps negative overlap to zero. Legacy splitting is unchanged; AutoOverlap treats it as a minimum right extension. +field ParallelOptions.Workers int ok parallel.go:90 copies options and replaces nonpositive Workers with runtime.NumCPU without mutating the caller. field RedisError.Err error ok errors.go:99; the client error, returned by Unwrap at errors.go:109 field RedisError.Key string ok errors.go:97; the key involved, v2_ops.go:108 field RedisError.Op string ok errors.go:95; the Redis verb, e.g. "HGETALL" at v2_ops.go:108 func Create(args *AhoCorasickArgs) (*AhoCorasick, error) ok acor.go:403 delegates to CreateContext with context.Background, and the documented error cases are the guards at acor.go:420-437 and client.go:47-68 func CreateContext(ctx context.Context, args *AhoCorasickArgs) (*AhoCorasick, error) ok acor.go:418; ctx governs setup only, and the background listener runs on an internal context per acor.go:476 func DefaultMigrationOptions() *MigrationOptions ok schema.go:64 names DryRun=false, KeepOldKeys=false, Progress=nil; the body returns the zero value at schema.go:67, which is exactly those three -func DefaultParallelOptions() *ParallelOptions ok options.go:83 returns exactly the four documented values, and is the only source of them +func DefaultParallelOptions() *ParallelOptions ok options.go:88 supplies CPU workers, 1000 rune chunks, word boundary, overlap 50 and AutoOverlap false. method (*AhoCorasick) Add(keyword string) (int, error) fixed acor.go:654 listed only "added" and "already exists" for a 0 return; an empty keyword also returns (0, nil) at redis_backed_ops.go:20 and v2_ops.go:81. Case added; TestEmptyKeywordIsNotAnErrorOutsideBatch pins it method (*AhoCorasick) AddContext(ctx context.Context, keyword string) (int, error) ok context_ops.go:7 forwards to the same ops.add that Add uses (acor.go:701), so ctx reaches Redis in V2 and preset mode; cross-reference to Add's return values added method (*AhoCorasick) AddMany(keywords []string, opts *BatchOptions) (*BatchResult, error) fixed batch.go:45 said 'duplicate keywords are skipped', which reads as exact duplicates; screening is on the normalized form (batch.go:110-114), so 'Foo' and 'foo' are one keyword on a case-insensitive collection. Also added that a transactional failure returns a nil *BatchResult (batch.go:204,214). TestBatchDuplicatesAreJudgedNormalized pins both @@ -132,8 +135,8 @@ method (*AhoCorasick) RollbackToV1() error fixed migration.go:349 named only the method (*AhoCorasick) SchemaVersion() int ok acor.go:562 returns the stored version with no Redis I/O method (*AhoCorasick) Suggest(input string) ([]string, error) ok acor.go:710 delegates to ops.suggest; preset mode returns ErrSuggestRequiresRedis at redis_backed_ops.go:160 method (*AhoCorasick) SuggestContext(ctx context.Context, input string) ([]string, error) fixed context_ops.go:49 omitted that preset mode cannot serve it at all - redis_backed_ops.go:160 returns ErrSuggestRequiresRedis, since the local automaton holds no prefix index. Added -method (*AhoCorasick) SuggestIndex(input string) (map[string][]int, error) ok acor.go:716 delegates to ops.suggestIndex; preset mode returns ErrSuggestRequiresRedis at redis_backed_ops.go:164 -method (*AhoCorasick) SuggestIndexContext(ctx context.Context, input string) (map[string][]int, error) fixed context_ops.go:57; same omission and same sentinel at redis_backed_ops.go:164 +method (*AhoCorasick) SuggestIndex(input string) (map[string][]int, error) ok acor.go:716 delegates to ops.suggestIndex; preset mode returns ErrSuggestRequiresRedis at redis_backed_ops.go:159 +method (*AhoCorasick) SuggestIndexContext(ctx context.Context, input string) (map[string][]int, error) fixed context_ops.go:57; same omission and same sentinel at redis_backed_ops.go:159 method (*MigrationResult) Stats() map[string]interface{} fixed schema.go:127 offered 'migration statistics'; it returns 6 of the 13 fields (schema.go:128-135), omitting every outcome field, so a caller cannot tell success from a dry run or a failure by reading the map. Now documented as a projection with the six named method (*OperationError) Error() string ok errors.go:82; includes op, schema and cause, and adds the keyword only when set method (*OperationError) Unwrap() error ok errors.go:90 returns Err, so errors.Is and errors.As reach the cause as documented diff --git a/api/v1.txt b/api/v1.txt index 2bec1cfb..c7923865 100644 --- a/api/v1.txt +++ b/api/v1.txt @@ -51,6 +51,8 @@ field BatchResult.Skipped []string field CacheStats.Hits uint64 field CacheStats.LastInvalidationLag time.Duration field CacheStats.Misses uint64 +field CacheStats.PresetPollFailures uint64 +field CacheStats.PresetReloadFailures uint64 field CacheStats.RebuildDuration time.Duration field CacheStats.Rebuilds uint64 field KeywordError.Error error @@ -82,6 +84,7 @@ field OperationError.Err error field OperationError.Keyword string field OperationError.Op string field OperationError.Schema int +field ParallelOptions.AutoOverlap bool field ParallelOptions.Boundary ChunkBoundary field ParallelOptions.ChunkSize int field ParallelOptions.Overlap int diff --git a/changes/unreleased/20260905-search-refresh.yaml b/changes/unreleased/20260905-search-refresh.yaml new file mode 100644 index 00000000..74b75030 --- /dev/null +++ b/changes/unreleased/20260905-search-refresh.yaml @@ -0,0 +1,5 @@ +kind: Added +body: "Add opt-in automatic parallel boundary protection and cancellation-independent Preset reloads, version-only polling, and refresh failure counters. No data migration is required." +time: 2026-09-05T19:00:00+09:00 +custom: + Issue: "243" diff --git a/docs/content/guides/parallel-matching.md b/docs/content/guides/parallel-matching.md index e84fc2f7..b799bf4e 100644 --- a/docs/content/guides/parallel-matching.md +++ b/docs/content/guides/parallel-matching.md @@ -5,7 +5,7 @@ weight: 2 # Parallel Matching -For large texts, use parallel matching to leverage multiple goroutines. +For large texts, use parallel matching to scan with multiple goroutines. ## Overview @@ -13,19 +13,24 @@ Parallel matching splits text into chunks and processes them concurrently, signi ## Basic Usage + ```go matches, err := ac.FindParallel(largeText, &acor.ParallelOptions{ Workers: 4, + ChunkSize: 1000, + AutoOverlap: true, Boundary: acor.ChunkBoundaryWord, }) if err != nil { panic(err) } +_ = matches ``` ## Chunk Boundaries -Chunk boundaries ensure matches aren't split across chunks: +Boundary selection chooses where base chunks end. Enable `AutoOverlap` to protect +matches across every boundary; boundary selection alone cannot guarantee this. ### ChunkBoundaryWord (default) @@ -34,6 +39,8 @@ Splits at word boundaries, ideal for natural language text: ```go opts := &acor.ParallelOptions{ Workers: 4, + ChunkSize: 1000, + AutoOverlap: true, Boundary: acor.ChunkBoundaryWord, } ``` @@ -45,6 +52,8 @@ Splits at line breaks, ideal for log files: ```go opts := &acor.ParallelOptions{ Workers: 4, + ChunkSize: 1000, + AutoOverlap: true, Boundary: acor.ChunkBoundaryLine, } ``` @@ -56,10 +65,30 @@ Splits at sentence endings, ideal for document processing: ```go opts := &acor.ParallelOptions{ Workers: 4, + ChunkSize: 1000, + AutoOverlap: true, Boundary: acor.ChunkBoundarySentence, } ``` +## Automatic Boundary Protection + +`AutoOverlap: true` uses the longest keyword in the engine loaded for this call. +Each base chunk owns its starting positions and extends its right search range by +`max(Overlap, longest keyword rune length - 1)`. Matches starting in the extension +belong to the next base chunk and are excluded from the current chunk. This also +finds keywords longer than `ChunkSize`, including Korean and emoji keywords. +No extra Redis query is needed, and all workers use the same dictionary snapshot. + +`FindParallel` returns each keyword once, in chunk order and then first-match scan +order within each chunk. That order need not equal serial scan order. +`FindIndexParallel` returns sorted unique rune positions in the original text. + +The default `AutoOverlap` is `false`, preserving legacy overlapping chunks. +In that mode, an insufficient `Overlap` can miss boundary matches, and a keyword +longer than a chunk may not fit in any chunk. Options are copied before normalization. +`ChunkSize` must be positive; empty input performs no Redis reads. + ## Performance Tuning ### Worker Count @@ -83,7 +112,7 @@ Control chunk size with the `ChunkSize` option: ```go opts := &acor.ParallelOptions{ Workers: 4, - ChunkSize: 10000, // 10KB chunks + ChunkSize: 10000, // 10,000 runes per base chunk } ``` diff --git a/docs/content/guides/preset-engine.md b/docs/content/guides/preset-engine.md index f791ecd4..c56f0fa9 100644 --- a/docs/content/guides/preset-engine.md +++ b/docs/content/guides/preset-engine.md @@ -110,6 +110,14 @@ info, err := ac.Info() // (*AhoCorasickInfo, error) ac.Flush() ``` +## Refresh and Cancellation + +Stale reads share a reload but respond independently to request cancellation. +A failed reload returns a search error and is retried by a later request. Optional +version polling helps recover from missed Pub/Sub messages; its interval does not +guarantee freshness during failures. See [invalidation safety](../redis-backed-engine/#invalidation-safety) +for the refresh lifecycle and failure counters. + ## Next Steps - [Redis-Backed Engine](../redis-backed-engine/) - Redis persistence details diff --git a/docs/content/guides/redis-backed-engine.md b/docs/content/guides/redis-backed-engine.md index 69afdb21..ca4f4f21 100644 --- a/docs/content/guides/redis-backed-engine.md +++ b/docs/content/guides/redis-backed-engine.md @@ -35,13 +35,13 @@ Instance A ──Find()──▶ local engine (0 RTT) - **Writes**: V2 Lua scripts with optimistic locking (up to 3 retries with backoff) - **Reads**: Local preset-optimized automaton — no Redis I/O - **Invalidation**: Redis Pub/Sub notifies all instances on mutation -- **Degraded mode**: If reload fails, the last-good engine continues serving reads +- **Reload failures**: The previous engine is retained, but the waiting search returns an error. A later search retries the reload. ## Invalidation Safety Redis Pub/Sub is best effort. In a multi-instance deployment, set -`InvalidationPollInterval` to bound how long a dropped invalidation can leave a -local preset engine stale: +`InvalidationPollInterval` to recover from dropped invalidations after a successful +version poll and reload: ```go @@ -55,7 +55,23 @@ _ = args ``` The zero value disables polling. Polling only applies to Preset mode; normal -invalidation still uses Pub/Sub. +invalidation still uses Pub/Sub. Each poll reads only the existing `version` hash +field; a changed version marks the engine stale, and the next search reads the +full snapshot. The interval is not a freshness upper bound: Redis failures, query +latency, and rebuild time can delay recovery. This mode provides eventual refresh, +not strong consistency or a guarantee of fresh results during an outage. + +Concurrent stale reads share a reload, while each request independently observes +its own context cancellation. Canceling one waiter leaves the others running. +The instance lifetime owns the job; all waiters leaving or `Close` cancels it. +Redis reads and engine builds run outside the state lock. A generation check +rejects and retries snapshots overtaken by local writes or invalidations. + +`CacheStats().PresetReloadFailures` counts failed shared jobs once per job; +`PresetPollFailures` counts failed version polls. Cancellation is excluded. +Both counters are cumulative per instance and remain zero outside Preset mode. +Monitor these counters alongside application search errors; a retained previous +engine does not turn a failed reload into a successful response. ## Quick Start diff --git a/docs/content/operations/monitoring.md b/docs/content/operations/monitoring.md index 2bf59033..453fed92 100644 --- a/docs/content/operations/monitoring.md +++ b/docs/content/operations/monitoring.md @@ -33,9 +33,8 @@ you run ACOR as a service, not to instrument your own process. ## Core library: cache statistics -`CacheStats()` answers the three questions that decide whether ACOR is behaving: how -often reads avoid Redis, how much a peer's write costs every reader, and how fast -invalidations propagate. It does no Redis I/O, so scraping it on a timer is cheap. +`CacheStats()` reports how often reads avoid Redis, rebuild cost, invalidation lag, +and Preset refresh failures. It does no Redis I/O, so scraping it on a timer is cheap. ```go @@ -61,6 +60,24 @@ Wire those three into whatever you already run — a Prometheus collector, an OT meter, a log line. ACOR deliberately depends on no metrics library, so the choice stays yours. +### Preset refresh failures + +Track increases in `PresetReloadFailures` and `PresetPollFailures` alongside search +errors. A shared failed reload increments the first counter once, even when many +requests receive the error. A failed version poll increments the second counter; +its next tick retries. Request cancellation is excluded from both counters. +These counters remain zero outside Preset mode, and polling is disabled by default. + +A failed reload retains the previous engine but returns an error to the search; +it does not serve that engine as a fallback. Each waiting request responds to its +own context cancellation without canceling other waiters. All waiters leaving or +`Close` cancels the shared job. Redis reads and engine builds run outside the +state lock, and snapshots overtaken by local state changes are rejected. + +Polling reads only the version field. It detects missed invalidations after a +successful poll; the next search must then fetch and build the full dictionary. +The configured interval is not an upper bound on staleness during failures. + ### Reading the numbers - **The counters are per instance and per process.** Nothing is aggregated through @@ -68,7 +85,7 @@ stays yours. - **`Rebuilds` will not equal `Misses`.** Concurrent misses coalesce onto one build, so `Misses - Rebuilds` is what that coalescing saved; local writes rebuild off the read path and push the count the other way. In `Preset` mode `Rebuilds` starts at 1, from - the build during `Create`. Both counters are `uint64`, so check `Misses > Rebuilds` + the build during `Create`; builds discarded after a generation conflict also count. Both counters are `uint64`, so check `Misses > Rebuilds` before subtracting — a write-heavy instance is routinely the other way round, and the difference wraps to roughly 1.8e19 rather than going negative. - **One scanning call is one read, whatever it scans over.** `FindParallel`, diff --git a/docs/content/operations/troubleshooting.md b/docs/content/operations/troubleshooting.md index 578a0c7a..faf634a1 100644 --- a/docs/content/operations/troubleshooting.md +++ b/docs/content/operations/troubleshooting.md @@ -146,13 +146,17 @@ _ = args disconnected subscriber can miss an invalidation. **Solution:** In multi-instance deployments, set -`InvalidationPollInterval` to the maximum acceptable staleness window: +`InvalidationPollInterval` to retry version checks periodically: ```go args.InvalidationPollInterval = 30 * time.Second ``` -The option is disabled by default and ignored outside Preset mode. +The option is disabled by default and ignored outside Preset mode. The interval +is not a freshness bound: recovery requires a successful version poll and a +successful reload on the next search. Inspect `CacheStats().PresetPollFailures` +and `PresetReloadFailures` when updates remain invisible. Reload errors are +returned to searches; the retained engine is not automatically served as fallback. ## Performance Issues diff --git a/docs/content/reference/api.md b/docs/content/reference/api.md index 305b7b51..9157cedf 100644 --- a/docs/content/reference/api.md +++ b/docs/content/reference/api.md @@ -462,3 +462,14 @@ const ( ChunkBoundaryLine // Split at newlines ) ``` + +## Parallel Boundary Protection and Preset Refresh + +`ParallelOptions.AutoOverlap` (default `false`) enables dictionary-aware right +extensions for each base chunk. It protects keywords longer than `ChunkSize` +without additional Redis reads. Results use the existing parallel order and rune +positions. See [parallel matching](../../guides/parallel-matching/). + +`CacheStats.PresetReloadFailures` and `CacheStats.PresetPollFailures` are cumulative +`uint64` counters. Shared reload failures count once per job; cancellations are +excluded. They remain zero in other modes. See [Preset refresh](../../guides/redis-backed-engine/#invalidation-safety). diff --git a/examples/parallel/main.go b/examples/parallel/main.go index 337d58ec..42e9ada7 100644 --- a/examples/parallel/main.go +++ b/examples/parallel/main.go @@ -30,9 +30,10 @@ func main() { largeText := "foo bar baz " matches, err := ac.FindParallel(largeText, &acor.ParallelOptions{ - Workers: 4, - Boundary: acor.ChunkBoundaryWord, - ChunkSize: 1000, + Workers: 4, + Boundary: acor.ChunkBoundaryWord, + ChunkSize: 1000, + AutoOverlap: true, }) if err != nil { fmt.Fprintf(os.Stderr, "failed to find parallel: %v\n", err) diff --git a/internal/engine/engine_handle.go b/internal/engine/engine_handle.go index afdad5f8..e046cf00 100644 --- a/internal/engine/engine_handle.go +++ b/internal/engine/engine_handle.go @@ -2,6 +2,8 @@ package engine +import "unicode/utf8" + // container is the optional specialization behind Engine.Contains. Routing a // presence check through matchString was the most expensive of the cheap // operations: it missed the byte-scan fast path on ASCII text and reported a match @@ -23,7 +25,8 @@ const findResultHint = 8 // package can build and query the automaton without depending on the concrete // engine types (which stay unexported). type Engine struct { - impl matchEngine + impl matchEngine + maxKeywordRunes int } // New returns an Engine backed by the implementation selected for preset. @@ -33,6 +36,10 @@ func New(preset Preset) *Engine { // Build (re)constructs the automaton from the given keyword set. func (e *Engine) Build(keywords map[string]struct{}) { + e.maxKeywordRunes = 0 + for keyword := range keywords { + e.maxKeywordRunes = max(e.maxKeywordRunes, utf8.RuneCountInString(keyword)) + } e.impl.buildFromKeywords(keywords) } @@ -97,3 +104,6 @@ func (e *Engine) Stream(next func() (rune, bool), emit func(keyword string, star func (e *Engine) Info() *InMemoryInfo { return e.impl.info() } + +// MaxKeywordRunes returns the longest keyword length in runes. +func (e *Engine) MaxKeywordRunes() int { return e.maxKeywordRunes } diff --git a/pkg/acor/acor.go b/pkg/acor/acor.go index b654c40b..d88e9d3e 100644 --- a/pkg/acor/acor.go +++ b/pkg/acor/acor.go @@ -360,10 +360,12 @@ type AhoCorasickArgs struct { // InvalidationPollInterval enables a background safety net for the Preset // engine: every interval it compares the collection's stored version against - // the local one and reloads if they differ. Cross-instance invalidation is + // the local one and marks the engine stale if they differ; the next read + // reloads it. Only the version field is fetched by a poll. Invalidation is // normally driven by best-effort Redis Pub/Sub, which has no delivery // guarantee — a dropped message leaves a node serving stale data until the - // next local write. This poll bounds that staleness to the interval. + // next local write. Recovery requires a successful poll and reload, so the + // interval is not a freshness bound, especially during Redis failures. // // Disabled by default (zero). Recommended for multi-instance deployments // (e.g. 30 * time.Second). Only applies to Preset mode; ignored otherwise. diff --git a/pkg/acor/batch_atomic.go b/pkg/acor/batch_atomic.go index 22e98324..f9cf83ef 100644 --- a/pkg/acor/batch_atomic.go +++ b/pkg/acor/batch_atomic.go @@ -102,9 +102,21 @@ func (ac *redisBackedAC) removeManyAtomic(ctx context.Context, keywords []string // already folded this write into snap.Keywords. func (ac *redisBackedAC) applyCommittedWrite(snap *trieSnapshot, newVersion int64) { ac.mu.Lock() - ac.applyReload(snap) - ac.localVersion = newVersion + ac.generation++ + generation := ac.generation + ac.stale = true ac.mu.Unlock() + keywords, engine := ac.prepareSnapshot(snap) + ac.mu.Lock() + defer ac.mu.Unlock() + // A later write, invalidation, or completed reload owns the newer state. + if ac.generation != generation { + return + } + ac.keywordSet = keywords + ac.engine = engine + ac.localVersion = newVersion + ac.stale = false } // --- V2 mode --- diff --git a/pkg/acor/context_ops.go b/pkg/acor/context_ops.go index 14769b3a..9ec84716 100644 --- a/pkg/acor/context_ops.go +++ b/pkg/acor/context_ops.go @@ -156,8 +156,6 @@ func (ac *AhoCorasick) FindParallelContext(ctx context.Context, text string, opt return []string{}, nil } - chunks := splitChunks(text, opts) - // One engine for the whole call, not one per chunk. Each ops.find reloaded it, // so an N-chunk text cost N round trips where serial Find costs one at any // input size — the fixed read cost V2 is built around (#205). Loading once also @@ -168,6 +166,13 @@ func (ac *AhoCorasick) FindParallelContext(ctx context.Context, text string, opt return nil, err } + var chunks []chunk + if opts.AutoOverlap { + chunks = splitAutoChunks(text, opts, eng.MaxKeywordRunes()) + } else { + chunks = splitChunks(text, opts) + } + perChunk, err := scanChunks(ctx, chunks, opts.Workers, func(ctx context.Context, c chunk) ([]string, error) { // Per chunk, not once above: the in-memory scan is not ctx-threaded, and in // Preset mode loadEngine touches no Redis, so this is the only place a @@ -181,6 +186,20 @@ func (ac *AhoCorasick) FindParallelContext(ctx context.Context, text string, opt // dedupPreservingOrder below. On match-dense text that per-occurrence slice // is most of the scan's allocation, and it is accumulated across every chunk // before the dedup runs. + if opts.AutoOverlap { + matches := make([]string, 0) + seen := make(map[string]struct{}) + eng.MatchString(normalizeText(c.text, ac.caseSensitive), func(keyword string, start, _ int) bool { + if start < c.ownedRunes { + if _, exists := seen[keyword]; !exists { + seen[keyword] = struct{}{} + matches = append(matches, keyword) + } + } + return true + }) + return matches, nil + } return eng.FindSet(normalizeText(c.text, ac.caseSensitive)), nil }) if err != nil { @@ -205,18 +224,33 @@ func (ac *AhoCorasick) FindIndexParallelContext(ctx context.Context, text string return map[string][]int{}, nil } - chunks := splitChunks(text, opts) - // One engine for the whole call; see FindParallelContext. eng, err := ac.ops.loadEngine(ctx) if err != nil { return nil, err } + var chunks []chunk + if opts.AutoOverlap { + chunks = splitAutoChunks(text, opts, eng.MaxKeywordRunes()) + } else { + chunks = splitChunks(text, opts) + } + perChunk, err := scanChunks(ctx, chunks, opts.Workers, func(ctx context.Context, c chunk) (map[string][]int, error) { if ctxErr := ctx.Err(); ctxErr != nil { return nil, ctxErr } + if opts.AutoOverlap { + matches := make(map[string][]int) + eng.MatchString(normalizeText(c.text, ac.caseSensitive), func(keyword string, start, _ int) bool { + if start < c.ownedRunes { + matches[keyword] = append(matches[keyword], start) + } + return true + }) + return matches, nil + } return eng.FindIndex(normalizeText(c.text, ac.caseSensitive)), nil }) if err != nil { diff --git a/pkg/acor/options.go b/pkg/acor/options.go index df51ba93..23118e2c 100644 --- a/pkg/acor/options.go +++ b/pkg/acor/options.go @@ -76,6 +76,11 @@ type ParallelOptions struct { // so set it to at least your longest keyword or start from DefaultParallelOptions. // A negative value is clamped to zero. Overlap int + // AutoOverlap protects every boundary using the loaded dictionary’s longest + // keyword. Base chunks do not overlap; their right search ranges extend by + // max(Overlap, longest keyword rune length - 1). Default false preserves + // legacy splitting. Keywords longer than ChunkSize are supported. + AutoOverlap bool } // DefaultParallelOptions returns parallel processing options with sensible defaults: diff --git a/pkg/acor/parallel.go b/pkg/acor/parallel.go index 3f301b5e..c0313273 100644 --- a/pkg/acor/parallel.go +++ b/pkg/acor/parallel.go @@ -14,6 +14,7 @@ import ( type chunk struct { text string textOffset int + ownedRunes int } func splitChunks(text string, opts *ParallelOptions) []chunk { @@ -90,6 +91,8 @@ func normalizeParallelOptions(opts *ParallelOptions) *ParallelOptions { if opts == nil { return DefaultParallelOptions() } + normalized := *opts + opts = &normalized if opts.Workers <= 0 { opts.Workers = runtime.NumCPU() } @@ -157,10 +160,13 @@ func scanChunks[T any](ctx context.Context, chunks []chunk, workers int, scan fu // most once regardless of how many times or in how many chunks it occurs. (Find // reports every occurrence; FindParallel reports a set.) // -// Limitation: a keyword longer than opts.Overlap that straddles a chunk boundary -// can be missed, since it fits in no single chunk. Set Overlap to at least your -// longest expected keyword length, or use FindStream (which never splits a match) -// when this matters. +// With AutoOverlap enabled, every boundary is protected using the loaded +// dictionary, including keywords longer than ChunkSize. Results retain chunk +// order, then first-match scan order within each chunk. +// +// Limitation with AutoOverlap disabled: a keyword longer than opts.Overlap that straddles a chunk boundary +// can be missed, since it fits in no single chunk. Enable AutoOverlap or use +// FindStream (which never splits a match) when this matters. // // Example: // @@ -183,9 +189,11 @@ func (ac *AhoCorasick) FindParallel(text string, opts *ParallelOptions) ([]strin // Due to chunk overlap, matches at chunk boundaries may have duplicate indices // that are automatically deduplicated. // -// Limitation: like FindParallel, a keyword longer than opts.Overlap that straddles -// a chunk boundary can be missed. Set Overlap to at least your longest expected -// keyword length. +// Enable AutoOverlap to protect all boundaries, including keywords longer than +// ChunkSize. Positions are rune offsets in the original text. +// +// Limitation with AutoOverlap disabled: a keyword longer than opts.Overlap that straddles +// a chunk boundary can be missed. Enable AutoOverlap to protect these matches. // // Example: // @@ -228,3 +236,23 @@ func mergeIndexResults(chunks []chunk, perChunk []map[string][]int) map[string][ } return result } + +// splitAutoChunks assigns each start position to exactly one base chunk. +func splitAutoChunks(text string, opts *ParallelOptions, longest int) []chunk { + runes := []rune(text) + extension := max(opts.Overlap, longest-1) + chunks := make([]chunk, 0) + for start := 0; start < len(runes); { + end := start + min(opts.ChunkSize, len(runes)-start) + if end < len(runes) { + boundary := findBoundary(runes, end, opts.Boundary, opts.ChunkSize/defaultMaxBacktrackDivisor) + if boundary > start { + end = boundary + } + } + searchEnd := end + min(extension, len(runes)-end) + chunks = append(chunks, chunk{text: string(runes[start:searchEnd]), textOffset: start, ownedRunes: end - start}) + start = end + } + return chunks +} diff --git a/pkg/acor/preset_poll_test.go b/pkg/acor/preset_poll_test.go new file mode 100644 index 00000000..16747e8e --- /dev/null +++ b/pkg/acor/preset_poll_test.go @@ -0,0 +1,130 @@ +// SPDX-License-Identifier: Apache-2.0 + +package acor + +import ( + "context" + "errors" + "testing" + + "github.com/redis/go-redis/v9" +) + +type versionProbeStorage struct { + kvStorage + gets int + snapshots int + failure error +} + +func (s *versionProbeStorage) HGet(ctx context.Context, key, field string) (string, error) { + s.gets++ + if field != fieldVersion { + return "", errors.New("unexpected field") + } + if s.failure != nil { + return "", s.failure + } + return s.kvStorage.HGet(ctx, key, field) +} +func (s *versionProbeStorage) HGetAll(ctx context.Context, key string) (map[string]string, error) { + s.snapshots++ + return s.kvStorage.HGetAll(ctx, key) +} +func TestPresetPollVersionOnlyAndRecovery(t *testing.T) { + ac := newTestPresetRedis(t, PresetBalanced) + testPresetPollRecovery(t, ac) +} +func TestIntegrationPresetPollVersionOnlyAndRecovery(t *testing.T) { + ac := newIntegrationAC(t, "acor-it-preset-poll", &AhoCorasickArgs{Preset: PresetBalanced}) + testPresetPollRecovery(t, ac) +} +func testPresetPollRecovery(t *testing.T, ac *AhoCorasick) { + t.Helper() + rb := ac.ops.(*redisBackedAC) + probe := &versionProbeStorage{kvStorage: rb.storage} + rb.storage = probe + rb.pollVersion() + if probe.gets != 1 || probe.snapshots != 0 { + t.Fatalf("poll commands HGET=%d HGETALL=%d", probe.gets, probe.snapshots) + } + // Redis mutation with no PUBLISH deliberately drops the notification. + if err := rb.redisClient.HSet(context.Background(), trieKey(rb.name), fieldKeywords, `["new"]`, fieldVersion, "123").Err(); err != nil { + t.Fatal(err) + } + probe.failure = errors.New("Redis unavailable") + rb.pollVersion() + if ac.CacheStats().PresetPollFailures != 1 { + t.Fatal(ac.CacheStats()) + } + probe.failure = nil + rb.pollVersion() + found, err := ac.Find("new") + if err != nil || len(found) != 1 { + t.Fatalf("recovery: %v %v", found, err) + } + if probe.snapshots != 1 { + t.Fatalf("reloads %d", probe.snapshots) + } + rb.pollVersion() + if probe.snapshots != 1 { + t.Fatal("unchanged poll read full dictionary") + } + if _, err := ac.Find("new"); err != nil { + t.Fatal(err) + } + if probe.snapshots != 1 { + t.Fatal("warm read touched Redis") + } +} + +func TestTrieVersionCompatibility(t *testing.T) { + ac := newTestPresetRedis(t, PresetBalanced) + rb := ac.ops.(*redisBackedAC) + for _, value := range []string{"0", "123", "9223372036854775807", "-23", "null", "bad", "1.5", "9223372036854775808"} { + if err := rb.redisClient.HSet(context.Background(), trieKey(rb.name), fieldVersion, value).Err(); err != nil { + t.Fatal(err) + } + snap, err := readTrieSnapshot(context.Background(), rb.storage, rb.name) + if err != nil { + t.Fatal(err) + } + version, err := readTrieVersion(context.Background(), rb.storage, rb.name) + if err != nil || version != snap.Version { + t.Fatalf("%s: %d %v", value, version, err) + } + } + if err := rb.redisClient.HDel(context.Background(), trieKey(rb.name), fieldVersion).Err(); err != nil { + t.Fatal(err) + } + if version, err := readTrieVersion(context.Background(), rb.storage, rb.name); err != nil || version != 0 { + t.Fatalf("missing: %d %v", version, err) + } + if version, err := readTrieVersion(context.Background(), rb.storage, "absent"); err != nil || version != 0 { + t.Fatalf("absent: %d %v", version, err) + } + probe := &versionProbeStorage{kvStorage: rb.storage, failure: redis.Nil} + if version, err := readTrieVersion(context.Background(), probe, rb.name); err != nil || version != 0 { + t.Fatalf("nil: %d %v", version, err) + } +} + +// Backend timeouts count as failures when the job/request itself is not canceled. +func TestPresetRefreshTimeoutFailures(t *testing.T) { + ac := newTestPresetRedis(t, PresetBalanced) + rb, held := holdPreset(t, ac) + held.failure = context.DeadlineExceeded + read := startPresetRead(ac, context.Background()) + <-held.started + close(held.release) + if err := awaitPresetError(t, read); !errors.Is(err, context.DeadlineExceeded) { + t.Fatal(err) + } + probe := &versionProbeStorage{kvStorage: rb.storage, failure: context.DeadlineExceeded} + rb.storage = probe + rb.pollVersion() + stats := ac.CacheStats() + if stats.PresetReloadFailures != 1 || stats.PresetPollFailures != 1 { + t.Fatal(stats) + } +} diff --git a/pkg/acor/redis_backed.go b/pkg/acor/redis_backed.go index f35d7a03..84aa46e6 100644 --- a/pkg/acor/redis_backed.go +++ b/pkg/acor/redis_backed.go @@ -11,7 +11,6 @@ import ( "time" redis "github.com/redis/go-redis/v9" - "golang.org/x/sync/singleflight" matchengine "github.com/skyoo2003/acor/internal/engine" ) @@ -40,14 +39,23 @@ type redisBackedAC struct { stats *cacheStats - selfSkip selfSkipSet - reloadGroup singleflight.Group - pubsub subscription - stopCh chan struct{} - ctx context.Context - cancel context.CancelFunc - closeOnce sync.Once - closed int32 + selfSkip selfSkipSet + reload *presetReload + generation uint64 + pubsub subscription + stopCh chan struct{} + ctx context.Context + cancel context.CancelFunc + closeOnce sync.Once + closed int32 +} + +// presetReload is protected by mu; closing done publishes err to its waiters. +type presetReload struct { + done chan struct{} + cancel context.CancelFunc + waiters int + err error } // newRedisBacked creates a Redis-backed Aho-Corasick engine. It loads the @@ -164,31 +172,16 @@ func buildEngine(preset Preset, keywordSet map[string]struct{}) *matchengine.Eng return e } -// rebuildEngine replaces the engine from the current keyword set and records the -// build. Every build in this mode routes through here, not just the reload path: a -// single-keyword write and a flush rebuild directly, and a build the counters miss -// inflates the mean rebuild cost that a write is supposed to explain. -// -// The timing covers the build alone. The Redis round trip that fetched the keywords -// happened in the caller, and folding it in would report the network as build time. -// -// Caller holds ac.mu. -func (ac *redisBackedAC) rebuildEngine() { - start := time.Now() - engine := buildEngine(ac.preset, ac.keywordSet) - ac.stats.recordRebuild(time.Since(start)) - ac.engine = engine -} - -func (ac *redisBackedAC) applyReload(snap *trieSnapshot) { - keywordSet := make(map[string]struct{}, len(snap.Keywords)) +// prepareSnapshot builds an immutable engine without holding the state lock. +func (ac *redisBackedAC) prepareSnapshot(snap *trieSnapshot) (map[string]struct{}, *matchengine.Engine) { + keywords := make(map[string]struct{}, len(snap.Keywords)) for _, kw := range snap.Keywords { - keywordSet[kw] = struct{}{} + keywords[kw] = struct{}{} } - ac.keywordSet = keywordSet - ac.rebuildEngine() - ac.localVersion = snap.Version - ac.stale = false + start := time.Now() + engine := buildEngine(ac.preset, keywords) + ac.stats.recordRebuild(time.Since(start)) + return keywords, engine } // loadEngine returns an immutable engine snapshot for the current keyword set, @@ -205,55 +198,114 @@ func (ac *redisBackedAC) loadEngine(ctx context.Context) (*matchengine.Engine, e return e, nil } +// reloadFromRedis fetches and builds outside the state lock. A local write or +// invalidation during either phase makes the snapshot ineligible for installation. func (ac *redisBackedAC) reloadFromRedis(ctx context.Context) error { - snap, err := readTrieSnapshot(ctx, ac.storage, ac.name) - if err != nil { - return err + for { + if err := ctx.Err(); err != nil { + return err + } + ac.mu.RLock() + generation := ac.generation + ac.mu.RUnlock() + snap, err := readTrieSnapshot(ctx, ac.storage, ac.name) + if err != nil { + return err + } + if err := ctx.Err(); err != nil { + return err + } + keywords, engine := ac.prepareSnapshot(snap) + ac.mu.Lock() + if err := ctx.Err(); err != nil { + ac.mu.Unlock() + return err + } + if ac.generation != generation { + ac.mu.Unlock() + continue + } + ac.engine = engine + ac.keywordSet = keywords + ac.localVersion = snap.Version + ac.stale = false + ac.generation++ + ac.mu.Unlock() + return nil } - - ac.mu.Lock() - defer ac.mu.Unlock() - ac.applyReload(snap) - return nil } func (ac *redisBackedAC) markStale() { ac.mu.Lock() ac.stale = true + ac.generation++ ac.mu.Unlock() } func (ac *redisBackedAC) ensureValid(ctx context.Context) error { + if err := ctx.Err(); err != nil { + return err + } ac.mu.RLock() + valid := !ac.stale + ac.mu.RUnlock() + if valid { + ac.stats.hit() + return nil + } + ac.mu.Lock() if !ac.stale { - ac.mu.RUnlock() + ac.mu.Unlock() ac.stats.hit() return nil } - ac.mu.RUnlock() - - // Counted before the singleflight, not inside it: this read found stale state and - // waits for a rebuild whether or not it performs one. Readers coalesced onto - // another goroutine's reload therefore still count as misses, which is what makes - // Misses-Rebuilds the work the coalescing saved. ac.stats.miss() + work := ac.reload + if work == nil { + workCtx, cancel := context.WithCancel(ac.ctx) //nolint:gosec // Ownership transfers to work; completion or the final waiter calls cancel. + work = &presetReload{done: make(chan struct{}), cancel: cancel} + ac.reload = work + go ac.runReload(workCtx, work) + } + work.waiters++ + ac.mu.Unlock() - _, err, _ := ac.reloadGroup.Do("reload", func() (interface{}, error) { - ac.mu.Lock() - defer ac.mu.Unlock() - - if !ac.stale { - return nil, nil + select { + case <-ctx.Done(): + case <-ac.ctx.Done(): + case <-work.done: + } + ac.mu.Lock() + work.waiters-- + if work.waiters == 0 { + work.cancel() + if ac.reload == work { + ac.reload = nil } + } + ac.mu.Unlock() + if err := ctx.Err(); err != nil { + return err + } + if err := ac.ctx.Err(); err != nil { + return err + } + return work.err +} - snap, err := readTrieSnapshot(ctx, ac.storage, ac.name) - if err != nil { - return nil, err - } - ac.applyReload(snap) - return nil, nil - }) - return err +func (ac *redisBackedAC) runReload(ctx context.Context, work *presetReload) { + err := ac.reloadFromRedis(ctx) + ac.mu.Lock() + defer ac.mu.Unlock() + if err != nil && ctx.Err() == nil { + ac.stats.reloadFailure() + } + work.err = err + if ac.reload == work { + ac.reload = nil + } + close(work.done) + work.cancel() } // --- Pub/Sub --- @@ -295,7 +347,8 @@ func (ac *redisBackedAC) handleInvalidation(stats *cacheStats, payload string) { // startPoller runs a background safety net for missed Pub/Sub invalidations: // every pollInterval it compares the stored collection version against the local // one and marks the engine stale on any difference, so a dropped invalidation -// self-heals within one interval instead of persisting until the next local write. +// can recover on a subsequent successful poll and reload. The interval is not +// a freshness bound during failures. func (ac *redisBackedAC) startPoller() { go func() { ticker := time.NewTicker(ac.pollInterval) @@ -316,15 +369,20 @@ func (ac *redisBackedAC) startPoller() { // pollVersion marks the engine stale if Redis holds a version other than the one // last loaded locally. A transient read error is ignored; the next tick retries. func (ac *redisBackedAC) pollVersion() { - snap, err := readTrieSnapshot(ac.ctx, ac.storage, ac.name) + version, err := readTrieVersion(ac.ctx, ac.storage, ac.name) if err != nil { + if ac.ctx.Err() == nil { + ac.stats.pollFailure() + } return } - ac.mu.RLock() - changed := snap.Version != ac.localVersion - ac.mu.RUnlock() - if changed { - ac.markStale() + ac.mu.Lock() + defer ac.mu.Unlock() + // Repeated polls must not invalidate a reload already in progress. Otherwise + // a build slower than the poll interval could be retried forever. + if version != ac.localVersion && !ac.stale { + ac.stale = true + ac.generation++ } } diff --git a/pkg/acor/redis_backed_ops.go b/pkg/acor/redis_backed_ops.go index 69a28476..9f49e904 100644 --- a/pkg/acor/redis_backed_ops.go +++ b/pkg/acor/redis_backed_ops.go @@ -130,11 +130,7 @@ func (ac *redisBackedAC) flush(ctx context.Context) error { return err } - ac.mu.Lock() - ac.keywordSet = make(map[string]struct{}) - ac.rebuildEngine() - ac.stale = false - ac.mu.Unlock() + ac.applyCommittedWrite(&trieSnapshot{}, 0) ac.publishInvalidate(ctx) return nil diff --git a/pkg/acor/redis_storage.go b/pkg/acor/redis_storage.go index 8a7ce6ee..2fcd2f39 100644 --- a/pkg/acor/redis_storage.go +++ b/pkg/acor/redis_storage.go @@ -197,3 +197,7 @@ func (p *redisPipeliner) Exec(ctx context.Context) error { _, err := p.pipe.Exec(ctx) return err } + +func (s *redisStorage) HGet(ctx context.Context, key, field string) (string, error) { + return s.client.HGet(ctx, key, field).Result() +} diff --git a/pkg/acor/release_reliability_test.go b/pkg/acor/release_reliability_test.go new file mode 100644 index 00000000..57baf9ea --- /dev/null +++ b/pkg/acor/release_reliability_test.go @@ -0,0 +1,345 @@ +// SPDX-License-Identifier: Apache-2.0 + +package acor + +import ( + "context" + "errors" + "fmt" + "reflect" + "slices" + "sync/atomic" + "testing" + "time" +) + +//nolint:gocyclo // Cross-product parity matrix covers all modes and boundary options. +func TestAutoOverlapSerialParity(t *testing.T) { + for _, preset := range []Preset{PresetNone, PresetSpeed, PresetBalanced, PresetMemoryEfficient} { + t.Run(preset.String(), func(t *testing.T) { + var ac *AhoCorasick + if preset == PresetNone { + var err error + mr := createTestRedisServer(t) + t.Cleanup(mr.Close) + ac, err = Create(&AhoCorasickArgs{Addr: mr.Addr(), Name: t.Name()}) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = ac.Close() }) + } else { + ac = newTestPresetRedis(t, preset) + } + keywords := []string{"abcdefghijk", "hijk", "ijk", "한글😀키워드", "😀키워드", "키워드", "a. b\nc", "b\nc"} + if _, err := ac.AddMany(keywords, nil); err != nil { + t.Fatal(err) + } + text := "xxabcdefghijk 한글😀키워드. a. b\nc! abcdefghijk 한글😀키워드" + want, err := ac.FindIndex(text) + if err != nil { + t.Fatal(err) + } + wantSet := make([]string, 0, len(want)) + for k := range want { + wantSet = append(wantSet, k) + } + slices.Sort(wantSet) + counter := countRTT(t, ac) + for _, boundary := range []ChunkBoundary{ChunkBoundaryWord, ChunkBoundaryLine, ChunkBoundarySentence, ChunkBoundary(99)} { + for _, size := range []int{1, 2, 5, 11, 1000} { + for _, overlap := range []int{-1, 0, 3, 100} { + opts := &ParallelOptions{ChunkSize: size, Boundary: boundary, Overlap: overlap, AutoOverlap: true} + before := *opts + counter.reset() + got, err := ac.FindIndexParallel(text, opts) + if err != nil { + t.Fatal(err) + } + if !reflect.DeepEqual(got, want) { + t.Fatalf("boundary %v size %d overlap %d: %v != %v", boundary, size, overlap, got, want) + } + expectedRTT := 0 + if preset == PresetNone { + expectedRTT = 1 + } + if counter.count() != expectedRTT { + t.Fatalf("RTT = %d", counter.count()) + } + found, err := ac.FindParallel(text, opts) + if err != nil { + t.Fatal(err) + } + slices.Sort(found) + if !reflect.DeepEqual(found, wantSet) { + t.Fatalf("set = %v", found) + } + if *opts != before { + t.Fatal("caller options mutated") + } + } + } + } + counter.reset() + if _, err := ac.FindParallel("", &ParallelOptions{ChunkSize: 1, AutoOverlap: true}); err != nil { + t.Fatal(err) + } + if _, err := ac.FindIndexParallel("", &ParallelOptions{AutoOverlap: true}); !errors.Is(err, ErrInvalidChunkSize) { + t.Fatal(err) + } + if counter.count() != 0 { + t.Fatal("empty input read Redis") + } + }) + } +} + +func TestAutoOverlapOrder(t *testing.T) { + ac := newTestPresetRedis(t, PresetBalanced) + if _, err := ac.AddMany([]string{"abcdef", "bc", "de", "f"}, nil); err != nil { + t.Fatal(err) + } + got, err := ac.FindParallel("abcdefabcdef", &ParallelOptions{ChunkSize: 2, AutoOverlap: true}) + if err != nil { + t.Fatal(err) + } + // Chunk order, then engine scan order within each owned chunk. + if !reflect.DeepEqual(got, []string{"bc", "abcdef", "de", "f"}) { + t.Fatal(got) + } +} + +// Hold the first snapshot after reading it, so writes can commit while the +// reloader has an old snapshot. Later calls (including writes) pass through. +type heldPresetStorage struct { + kvStorage + calls atomic.Int64 + started chan context.Context + release chan struct{} + failure error +} + +func (s *heldPresetStorage) HGetAll(ctx context.Context, key string) (map[string]string, error) { + data, err := s.kvStorage.HGetAll(ctx, key) + if s.calls.Add(1) == 1 { + s.started <- ctx + select { + case <-ctx.Done(): + return nil, ctx.Err() + case <-s.release: + } + if s.failure != nil { + return nil, s.failure + } + } + return data, err +} +func holdPreset(t *testing.T, ac *AhoCorasick) (*redisBackedAC, *heldPresetStorage) { + t.Helper() + rb := ac.ops.(*redisBackedAC) + s := &heldPresetStorage{kvStorage: rb.storage, started: make(chan context.Context, 1), release: make(chan struct{})} + rb.storage = s + rb.markStale() + return rb, s +} +func awaitPresetError(t *testing.T, ch <-chan error) error { + t.Helper() + select { + case err := <-ch: + return err + case <-time.After(3 * time.Second): + t.Fatal("request did not complete") + return nil + } +} +func startPresetRead(ac *AhoCorasick, ctx context.Context) <-chan error { + ch := make(chan error, 1) + go func() { _, err := ac.FindContext(ctx, "old new"); ch <- err }() + return ch +} +func awaitPresetWaiters(t *testing.T, rb *redisBackedAC, n int) { + t.Helper() + deadline := time.Now().Add(3 * time.Second) + for time.Now().Before(deadline) { + rb.mu.RLock() + ok := rb.reload != nil && rb.reload.waiters == n + rb.mu.RUnlock() + if ok { + return + } + time.Sleep(time.Millisecond) + } + t.Fatal("waiters did not join") +} +func TestPresetReloadIndependentCancellation(t *testing.T) { + ac := newTestPresetRedis(t, PresetBalanced) + rb, s := holdPreset(t, ac) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + first := startPresetRead(ac, ctx) + workCtx := <-s.started + second := startPresetRead(ac, context.Background()) + awaitPresetWaiters(t, rb, 2) + cancel() + if err := awaitPresetError(t, first); !errors.Is(err, context.Canceled) { + t.Fatal(err) + } + if workCtx.Err() != nil { + t.Fatal("one cancellation canceled shared work") + } + close(s.release) + if err := awaitPresetError(t, second); err != nil { + t.Fatal(err) + } + if s.calls.Load() != 1 { + t.Fatal("reload was not shared") + } + if ac.CacheStats().PresetReloadFailures != 0 { + t.Fatal("cancellation counted as failure") + } +} +func TestPresetReloadAllCancelAndClose(t *testing.T) { + for _, closeInstance := range []bool{false, true} { + t.Run(fmt.Sprint(closeInstance), func(t *testing.T) { + ac := newTestPresetRedis(t, PresetBalanced) + rb, s := holdPreset(t, ac) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + first := startPresetRead(ac, ctx) + workCtx := <-s.started + second := startPresetRead(ac, ctx) + awaitPresetWaiters(t, rb, 2) + if closeInstance { + if err := ac.Close(); err != nil { + t.Fatal(err) + } + } else { + cancel() + } + for _, ch := range []<-chan error{first, second} { + if err := awaitPresetError(t, ch); !errors.Is(err, context.Canceled) { + t.Fatal(err) + } + } + select { + case <-workCtx.Done(): + case <-time.After(time.Second): + t.Fatal("shared work not canceled") + } + if !closeInstance { + if _, err := ac.Find("new"); err != nil { + t.Fatal(err) + } + } + if ac.CacheStats().PresetReloadFailures != 0 { + t.Fatal("cancellation counted") + } + }) + } +} +func TestPresetReloadFailureSharedAndRetry(t *testing.T) { + ac := newTestPresetRedis(t, PresetBalanced) + if _, err := ac.Add("old"); err != nil { + t.Fatal(err) + } + rb, s := holdPreset(t, ac) + s.failure = errors.New("snapshot unavailable") + first := startPresetRead(ac, context.Background()) + <-s.started + second := startPresetRead(ac, context.Background()) + awaitPresetWaiters(t, rb, 2) + close(s.release) + for _, ch := range []<-chan error{first, second} { + if err := awaitPresetError(t, ch); !errors.Is(err, s.failure) { + t.Fatal(err) + } + } + if ac.CacheStats().PresetReloadFailures != 1 { + t.Fatal(ac.CacheStats()) + } + rb.mu.RLock() + old := rb.engine.Find("old") + stale := rb.stale + rb.mu.RUnlock() + if len(old) != 1 || !stale { + t.Fatal("failed reload discarded previous state") + } + if got, err := ac.Find("old"); err != nil || len(got) != 1 { + t.Fatalf("retry %v %v", got, err) + } +} +func TestPresetReloadRejectsObsoleteSnapshot(t *testing.T) { + for _, action := range []string{"add", "remove", "flush", "invalidate"} { + t.Run(action, func(t *testing.T) { + ac := newTestPresetRedis(t, PresetBalanced) + if _, err := ac.Add("old"); err != nil { + t.Fatal(err) + } + rb, s := holdPreset(t, ac) + read := startPresetRead(ac, context.Background()) + <-s.started + done := make(chan error, 1) + go func() { + var err error + switch action { + case "add": + _, err = ac.Add("new") + case "remove": + _, err = ac.Remove("old") + case "flush": + err = ac.Flush() + case "invalidate": + rb.markStale() + } + done <- err + }() + if err := awaitPresetError(t, done); err != nil { + t.Fatal(err) + } + close(s.release) + if err := awaitPresetError(t, read); err != nil { + t.Fatal(err) + } + got, err := ac.Find("old new") + if err != nil { + t.Fatal(err) + } + want := []string{} + switch action { + case "add": + want = []string{"old", "new"} + case "invalidate": + want = []string{"old"} + } + if !reflect.DeepEqual(got, want) { + t.Fatalf("got %v want %v", got, want) + } + if s.calls.Load() < 2 { + t.Fatal("obsolete reload was not retried") + } + }) + } +} + +func TestPresetPollDoesNotRestartSlowReload(t *testing.T) { + ac := newTestPresetRedis(t, PresetBalanced) + rb := ac.ops.(*redisBackedAC) + if err := rb.redisClient.HSet(context.Background(), trieKey(rb.name), fieldKeywords, `["new"]`, fieldVersion, "123").Err(); err != nil { + t.Fatal(err) + } + _, held := holdPreset(t, ac) + read := startPresetRead(ac, context.Background()) + <-held.started + for range 3 { + rb.pollVersion() + } + close(held.release) + if err := awaitPresetError(t, read); err != nil { + t.Fatal(err) + } + if held.calls.Load() != 1 { + t.Fatal("polling restarted an already stale reload") + } + if got, err := ac.Find("new"); err != nil || len(got) != 1 { + t.Fatalf("got %v, err %v", got, err) + } +} diff --git a/pkg/acor/rtt_counter_test.go b/pkg/acor/rtt_counter_test.go index f16c163d..ad18dc17 100644 --- a/pkg/acor/rtt_counter_test.go +++ b/pkg/acor/rtt_counter_test.go @@ -220,3 +220,8 @@ var ( _ kvStorage = (*countingStorage)(nil) _ pipeliner = (*countingPipeliner)(nil) ) + +func (s *countingStorage) HGet(ctx context.Context, key, field string) (string, error) { + s.c.add() + return s.inner.HGet(ctx, key, field) +} diff --git a/pkg/acor/stats.go b/pkg/acor/stats.go index 13701bc0..957dc1b6 100644 --- a/pkg/acor/stats.go +++ b/pkg/acor/stats.go @@ -18,6 +18,12 @@ import ( // from or written to Redis, so a peer instance's activity is invisible: scrape each // instance separately rather than expecting a fleet-wide total. type CacheStats struct { + // PresetReloadFailures counts failed shared reload jobs, once per job. + // Request cancellation is excluded. Zero outside Preset mode. + PresetReloadFailures uint64 + // PresetPollFailures counts failed version polls, excluding cancellation. + // Zero outside Preset mode or when polling is disabled. + PresetPollFailures uint64 // Hits is the number of reads served from the local automaton without rebuilding // it. // @@ -48,7 +54,8 @@ type CacheStats struct { // means. Misses uint64 // Rebuilds is the number of automaton builds. It starts at 1 in Preset mode, which - // builds once during Create before any read. + // builds once during Create before any read. Preset builds discarded by a + // generation conflict or cancellation are included: their build cost was paid. // // It is deliberately not equal to Misses in either direction. Concurrent misses // coalesce onto one build, so Misses-Rebuilds is the work that coalescing saved; @@ -60,14 +67,15 @@ type CacheStats struct { // wraps to something near 2^64 instead of going negative. Rebuilds uint64 // RebuildDuration is the total time spent building automatons, excluding the Redis - // fetch and the wait for the lock that serializes rebuilds. The brief cache write - // lock a finished build takes to publish itself is included. + // fetch and synchronization waits before building. Preset times the engine build + // alone, excluding its later publication lock; cache modes may also include + // the brief publication step performed inside their timed build callback. // RebuildDuration/Rebuilds is the mean cost of one rebuild — what a write to a // large collection makes every reader pay. // // Where decoding the fetched payload lands differs by mode. A default V2 instance // parses inside the memoized build, so the decode counts here; Preset materializes - // its keyword set in applyReload before rebuildEngine starts timing, so there it + // its keyword set in prepareSnapshot before the build timer starts, so there it // does not. Read the mean against itself over time rather than across two // differently configured instances. RebuildDuration time.Duration @@ -103,11 +111,13 @@ type cacheStats struct { // types guarantee 8-byte alignment on 386/arm (both released by goreleaser, where a // misaligned 64-bit atomic panics) instead of leaving it to field order, for the // same reason selfSkipSet.publishCount uses one. - hits atomic.Uint64 - misses atomic.Uint64 - rebuilds atomic.Uint64 - rebuildNanos atomic.Int64 - lastLagNanos atomic.Int64 + presetReloadFailures atomic.Uint64 + presetPollFailures atomic.Uint64 + hits atomic.Uint64 + misses atomic.Uint64 + rebuilds atomic.Uint64 + rebuildNanos atomic.Int64 + lastLagNanos atomic.Int64 } func (s *cacheStats) hit() { @@ -151,11 +161,13 @@ func (s *cacheStats) snapshot() CacheStats { return CacheStats{} } return CacheStats{ - Hits: s.hits.Load(), - Misses: s.misses.Load(), - Rebuilds: s.rebuilds.Load(), - RebuildDuration: time.Duration(s.rebuildNanos.Load()), - LastInvalidationLag: time.Duration(s.lastLagNanos.Load()), + PresetReloadFailures: s.presetReloadFailures.Load(), + PresetPollFailures: s.presetPollFailures.Load(), + Hits: s.hits.Load(), + Misses: s.misses.Load(), + Rebuilds: s.rebuilds.Load(), + RebuildDuration: time.Duration(s.rebuildNanos.Load()), + LastInvalidationLag: time.Duration(s.lastLagNanos.Load()), } } @@ -170,3 +182,14 @@ func timeRebuild[T any](s *cacheStats, build func() (T, error)) (T, error) { } return out, err } + +func (s *cacheStats) reloadFailure() { + if s != nil { + s.presetReloadFailures.Add(1) + } +} +func (s *cacheStats) pollFailure() { + if s != nil { + s.presetPollFailures.Add(1) + } +} diff --git a/pkg/acor/storage.go b/pkg/acor/storage.go index 77ffaf36..0660df9d 100644 --- a/pkg/acor/storage.go +++ b/pkg/acor/storage.go @@ -27,6 +27,8 @@ type zMember struct { // // All operations accept a context for cancellation and timeout support. type kvStorage interface { + // HGet retrieves one hash field. A missing field returns redis.Nil. + HGet(ctx context.Context, key, field string) (string, error) // HGetAll retrieves all field-value pairs from a hash. HGetAll(ctx context.Context, key string) (map[string]string, error) // HSet sets multiple field-value pairs in a hash. diff --git a/pkg/acor/v2_transaction.go b/pkg/acor/v2_transaction.go index a4d6786d..9d84fd5a 100644 --- a/pkg/acor/v2_transaction.go +++ b/pkg/acor/v2_transaction.go @@ -72,11 +72,7 @@ func readTrieSnapshot(ctx context.Context, storage kvStorage, name string) (*tri return nil, newOperationError("unmarshal", SchemaV2, err) } } - if v, ok := trieData[fieldVersion]; ok { - if err := json.Unmarshal([]byte(v), &snap.Version); err != nil { - snap.Version = 0 - } - } + snap.Version = parseTrieVersion(trieData[fieldVersion]) return snap, nil } @@ -248,3 +244,23 @@ func (o *v2Operations) tryRemoveV2(ctx context.Context, keyword string) (int, er o.publishInvalidate(ctx) return 1, nil } + +// parseTrieVersion preserves the snapshot semantics: absent or malformed is zero. +func parseTrieVersion(value string) int64 { + var version int64 + if err := json.Unmarshal([]byte(value), &version); err != nil { + return 0 + } + return version +} + +func readTrieVersion(ctx context.Context, storage kvStorage, name string) (int64, error) { + value, err := storage.HGet(ctx, trieKey(name), fieldVersion) + if errors.Is(err, redis.Nil) { + return 0, nil + } + if err != nil { + return 0, newRedisError("HGET", trieKey(name), err) + } + return parseTrieVersion(value), nil +}