amit-jain commented on code in PR #2817: URL: https://github.com/apache/jackrabbit-oak/pull/2817#discussion_r4181364871
########## oak-search-lucene-ng/src/main/java/org/apache/jackrabbit/oak/plugins/index/luceneNg/internal/IndexSearcherHolder.java: ########## @@ -0,0 +1,130 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.jackrabbit.oak.plugins.index.luceneNg.internal; + +import org.apache.jackrabbit.oak.plugins.index.luceneNg.LuceneNgIndexDefinition; +import org.apache.jackrabbit.oak.plugins.index.luceneNg.LuceneNgIndexStorage; +import org.apache.jackrabbit.oak.plugins.index.luceneNg.directory.LuceneNgIndexCopier; +import org.apache.jackrabbit.oak.plugins.index.luceneNg.directory.OakDirectory; +import org.apache.jackrabbit.oak.spi.state.NodeState; +import org.apache.lucene.facet.sortedset.DefaultSortedSetDocValuesReaderState; +import org.apache.lucene.index.DirectoryReader; +import org.apache.lucene.search.IndexSearcher; +import org.apache.lucene.store.Directory; +import org.jetbrains.annotations.Nullable; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.io.Closeable; +import java.io.IOException; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; + +/** + * Manages IndexSearcher lifecycle for a Lucene 9 index. + * Opens the index from the {@link LuceneNgIndexStorage} node state passed in (typically the + * {@link LuceneNgIndexStorage#STORAGE_NODE_NAME} child under the index definition). + */ +public class IndexSearcherHolder implements Closeable { + + private static final Logger LOG = LoggerFactory.getLogger(IndexSearcherHolder.class); + + private final String indexName; + private DirectoryReader reader; + private IndexSearcher searcher; + private Directory directory; + private final ConcurrentMap<String, DefaultSortedSetDocValuesReaderState> facetStateCache = + new ConcurrentHashMap<>(); + + /** + * @param storageState {@link LuceneNgIndexStorage#storageState(NodeState)} for the index definition + * @param indexName the index name, used only for logging/error messages + */ + public IndexSearcherHolder(NodeState storageState, String indexName) throws IOException { + this(storageState, indexName, null, null); + } + + /** + * @param storageState {@link LuceneNgIndexStorage#storageState(NodeState)} for the index definition + * @param indexName the index name, used only for logging/error messages + * @param copier when non-null, wraps the remote {@link OakDirectory} with a local-disk + * cache (CopyOnRead) via {@link LuceneNgIndexCopier#wrapForRead} + * @param definition the index definition; required (non-null) when {@code copier} is non-null + */ + public IndexSearcherHolder(NodeState storageState, String indexName, + @Nullable LuceneNgIndexCopier copier, @Nullable LuceneNgIndexDefinition definition) throws IOException { + this.indexName = indexName; + OakDirectory oakDirectory = new OakDirectory(storageState.builder(), indexName, true); + Directory toOpen = oakDirectory; + if (copier != null) { + try { + toOpen = copier.wrapForRead(definition.getIndexPath(), definition, oakDirectory, LuceneNgIndexStorage.STORAGE_NODE_NAME); Review Comment: **CopyOnRead local cache gets wiped on every deactivate / index removal** `wrapForRead` is invoked even when the storage node doesn't exist (legacy `DefaultIndexReaderFactory` only wraps `if (data.exists())`). On `deactivate()` (`update(EMPTY_NODE)`) — and on runtime removal of a lucene9 index — `openIndex(path, EMPTY, MISSING_NODE)` runs: 1. no `:status` → `getUniqueId()` is null → `IndexRootDirectory.getIndexDir` falls back to the old-format `<sha256(path)>/0` dir; 2. `createLocalDirForIndexReader` sees a changed path and wraps it in `DeleteOldDirOnClose(old = live <name>-<uid>/lucene9)`; 3. `DirectoryReader.open` fails on the empty dir → `directory.close()` → **the live local cache is deleted**. Effect: every restart starts cold (the readiness-probe problem CoR is meant to fix), and on index removal the dir can be deleted under a reader still serving queries. Fix: only build the `OakDirectory` / call `wrapForRead` when `storageState.exists()`. Please add a test that the local dir survives `deactivate()` and an index removal. ########## oak-search-lucene-ng/src/main/java/org/apache/jackrabbit/oak/plugins/index/luceneNg/internal/LuceneNgIndexNode.java: ########## @@ -0,0 +1,202 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.jackrabbit.oak.plugins.index.luceneNg.internal; + +import org.apache.jackrabbit.oak.commons.PathUtils; +import org.apache.jackrabbit.oak.plugins.index.luceneNg.LuceneNgIndexDefinition; +import org.apache.jackrabbit.oak.plugins.index.luceneNg.LuceneNgIndexStorage; +import org.apache.jackrabbit.oak.plugins.index.luceneNg.directory.LuceneNgIndexCopier; +import org.apache.jackrabbit.oak.plugins.index.search.IndexNode; +import org.apache.jackrabbit.oak.plugins.index.search.IndexStatistics; +import org.apache.jackrabbit.oak.spi.state.NodeState; +import org.apache.lucene.facet.sortedset.DefaultSortedSetDocValuesReaderState; +import org.apache.lucene.search.IndexSearcher; +import org.jetbrains.annotations.NotNull; +import org.jetbrains.annotations.Nullable; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.io.IOException; +import java.util.concurrent.atomic.AtomicInteger; + +/** + * Represents a Lucene 9 index with its definition and a cached searcher. + * + * <p>One instance is built per generation of the index (whenever the tracker detects a + * definition or storage change) — it is never mutated or reopened in place. Wrapped by + * {@link LuceneNgIndexNodeManager}, whose inherited {@code IndexNodeManager} read/write + * lock is what makes {@link #release()} / {@link #closeResources()} safe: {@code close()} + * on the manager cannot return, and therefore {@link #closeResources()} cannot run, until + * every {@code acquire()}-holder has called {@link #release()}. Do not reintroduce + * per-call {@code IndexReader.tryIncRef()/decRef()} bookkeeping here — it is redundant + * with that lock and duplicating it reintroduces a concurrency race.</p> + */ +public class LuceneNgIndexNode implements IndexNode { + + private static final Logger LOG = LoggerFactory.getLogger(LuceneNgIndexNode.class); + private static final AtomicInteger ID_COUNTER = new AtomicInteger(); + + private final String indexPath; + /** Immutable snapshot of the index definition — used for definition change detection. */ + private final NodeState indexState; + /** + * Immutable snapshot of the storage node ({@link LuceneNgIndexStorage#STORAGE_NODE_NAME} child). + * Used together with {@link #indexState} to detect when data changes independently + * of the definition (which is the normal case during incremental indexing). + */ + private final NodeState storageState; + private final LuceneNgIndexDefinition definition; + /** Cached searcher; null when index has not been populated yet. */ + private final IndexSearcherHolder searcherHolder; + private final int indexNodeId = ID_COUNTER.incrementAndGet(); + + /** Set once by {@link LuceneNgIndexNodeManager}'s constructor. Package-private: + * only the owning manager binds itself, and only {@link #release()} reads it. */ + private LuceneNgIndexNodeManager owner; + + /** + * Creates a new index node, opening a cached {@link IndexSearcher} from + * {@link LuceneNgIndexStorage}. + * If the storage path does not exist yet the searcher is left null and + * {@link #getSearcher()} returns null. + * + * @param indexPath path to the index definition (e.g. "/oak:index/myIndex") + * @param root repository root state + * @param indexState index definition node state (immutable snapshot) + */ + public LuceneNgIndexNode(@NotNull String indexPath, + @NotNull NodeState root, + @NotNull NodeState indexState) { + this(indexPath, root, indexState, null); + } + + /** + * Creates a new index node, opening a cached {@link IndexSearcher} from + * {@link LuceneNgIndexStorage}. + * If the storage path does not exist yet the searcher is left null and + * {@link #getSearcher()} returns null. + * + * @param indexPath path to the index definition (e.g. "/oak:index/myIndex") + * @param root repository root state + * @param indexState index definition node state (immutable snapshot) + * @param copier when non-null, the {@link LuceneNgIndexCopier} used to wrap the + * opened directory with a local-disk cache (CopyOnRead) + */ + public LuceneNgIndexNode(@NotNull String indexPath, + @NotNull NodeState root, + @NotNull NodeState indexState, + @Nullable LuceneNgIndexCopier copier) { + this.indexPath = indexPath; + this.indexState = indexState; + this.definition = new LuceneNgIndexDefinition(root, indexState, indexPath); + + String indexName = PathUtils.getName(indexPath); + this.storageState = LuceneNgIndexStorage.storageState(indexState); + + IndexSearcherHolder holder = null; + try { + holder = new IndexSearcherHolder(storageState, indexName, copier, definition); + } catch (IOException e) { + LOG.debug("No index data for {} yet, searcher not opened: {}", indexPath, e.getMessage()); Review Comment: **Index open / copier failures are hidden and retried on every query** Every `IOException` here — including real failures from `wrapForRead` (disk full, permissions, sanity check, `FSDirectory` errors) — is logged at DEBUG as "No index data … yet". `openIndex` then returns null, so the index silently drops out of query planning, and `FulltextIndexTracker.findIndexNode` re-runs `openIndex` on every query (`mkdirs`, new `FSDirectory`, async close each time). Legacy propagates the exception so `BadIndexTracker` logs it and backs off. Fix: treat only "storage node missing" as the normal null case; on a copier failure log WARN and fall back to the bare `OakDirectory`; rethrow anything else so `BadIndexTracker` handles it. ########## oak-search-lucene-ng/src/main/java/org/apache/jackrabbit/oak/plugins/index/luceneNg/directory/CopyOnReadDirectory.java: ########## @@ -0,0 +1,401 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.jackrabbit.oak.plugins.index.luceneNg.directory; + +import java.io.File; +import java.io.IOException; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashSet; +import java.util.List; +import java.util.Objects; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.Executor; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.stream.Collectors; + +import org.apache.jackrabbit.oak.commons.PerfLogger; +import org.apache.jackrabbit.oak.commons.collections.SetUtils; +import org.apache.lucene.store.Directory; +import org.apache.lucene.store.FilterDirectory; +import org.apache.lucene.store.IOContext; +import org.apache.lucene.store.IndexInput; +import org.apache.lucene.store.IndexOutput; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import static java.util.Arrays.stream; +import static org.apache.jackrabbit.oak.commons.IOUtils.humanReadableByteCount; + +/** + * Directory implementation which lazily copies the index files from a + * remote directory in background. + * <p> + * Port of {@code oak-lucene}'s {@code CopyOnReadDirectory} for Lucene 9's {@link Directory} + * API. {@code remote} is narrowed to {@link OakDirectory} (this module's only remote + * implementation, see {@link LuceneNgIndexCopier#wrapForRead}) and a {@code localDir} + * {@link File} is threaded through so local-side existence checks can use direct + * {@link File} calls (via {@link LuceneNgIndexCopier#existsLocally(File, String)}) instead of + * the removed {@code Directory.fileExists}. Lucene 9 also removed the source-invoked + * {@code Directory.copy(...)} in favour of a destination-invoked {@code copyFrom(...)}. + */ +public class CopyOnReadDirectory extends FilterDirectory { + private static final Logger log = LoggerFactory.getLogger(CopyOnReadDirectory.class); + private static final PerfLogger PERF_LOGGER = new PerfLogger(LoggerFactory.getLogger(log.getName() + ".perf")); + + public static final String DELETE_MARGIN_MILLIS_NAME = "oak.lucene.delete.margin"; + public final long DELETE_MARGIN_MILLIS = Long.getLong(DELETE_MARGIN_MILLIS_NAME, TimeUnit.MINUTES.toMillis(5)); + + private final LuceneNgIndexCopier indexCopier; + private final OakDirectory remote; + private final Directory local; + private final File localDir; + private final boolean prefetch; + private final String indexPath; + private final Executor executor; + private final AtomicBoolean closed = new AtomicBoolean(); + + // exported as package private to be useful in tests + static final String WAIT_OTHER_COPY_SYSPROP_NAME = "cor.waitCopyMillis"; + + long waitOtherCopyTimeoutMillis = Long.getLong(WAIT_OTHER_COPY_SYSPROP_NAME, TimeUnit.SECONDS.toMillis(30)); + + private final ConcurrentMap<String, CORFileReference> files = new ConcurrentHashMap<>(); + + public CopyOnReadDirectory(LuceneNgIndexCopier indexCopier, OakDirectory remote, Directory local, File localDir, + boolean prefetch, String indexPath, Executor executor) throws IOException { + super(remote); + this.indexCopier = indexCopier; + this.executor = executor; + this.remote = remote; + this.local = local; + this.localDir = localDir; + this.prefetch = prefetch; + this.indexPath = indexPath; + + if (prefetch) { + prefetchIndexFiles(); + } + } + + @Override + public void deleteFile(String name) throws IOException { + throw new UnsupportedOperationException("Cannot delete in a ReadOnly directory"); + } + + @Override + public IndexOutput createOutput(String name, IOContext context) throws IOException { + throw new UnsupportedOperationException("Cannot write in a ReadOnly directory"); + } + + @Override + public IndexInput openInput(String name, IOContext context) throws IOException { + if (Objects.nonNull(name) && LuceneNgIndexCopier.REMOTE_ONLY.contains(name)) { + log.trace("[{}] opening remote only file {}", indexPath, name); + return remote.openInput(name, context); + } + + CORFileReference ref = files.get(name); + if (ref != null) { + if (ref.isLocalValid()) { + log.trace("[{}] opening existing local file {}", indexPath, name); + return files.get(name).openLocalInput(context); + } else { + indexCopier.readFromRemote(true); + logRemoteAccess( + "[{}] opening existing remote file as local version is not valid {}", + indexPath, name); + return remote.openInput(name, context); + } + } + + //If file does not exist then just delegate to remote and not + //schedule a copy task + if (!remote.fileExists(name)){ + if (log.isDebugEnabled()) { + log.debug("[{}] Looking for non existent file {}. Current known files {}", + indexPath, name, Arrays.toString(remote.listAll())); + } + return remote.openInput(name, context); + } + + CORFileReference toPut = new CORFileReference(name); + CORFileReference old = files.putIfAbsent(name, toPut); + if (old == null) { + log.trace("[{}] scheduled local copy for {}", indexPath, name); + copy(toPut); + } + + //If immediate executor is used the result would be ready right away + if (toPut.isLocalValid()) { + log.trace("[{}] opening new local file {}", indexPath, name); + return toPut.openLocalInput(context); + } + + logRemoteAccess("[{}] opening new remote file {}", indexPath, name); + indexCopier.readFromRemote(true); + return remote.openInput(name, context); + } + + public Directory getLocal() { + return local; + } + + private void copy(final CORFileReference reference) { + indexCopier.scheduledForCopy(); + executor.execute(new Runnable() { + @Override + public void run() { + indexCopier.copyDone(); + copyFilesToLocal(reference, true, true); + } + }); + } + + private void prefetchIndexFiles() throws IOException { + long start = PERF_LOGGER.start(); + long totalSize = 0; + int copyCount = 0; + List<String> copiedFileNames = new ArrayList<>(); + for (String name : remote.listAll()) { + if (LuceneNgIndexCopier.REMOTE_ONLY.contains(name)) { + continue; + } + CORFileReference fileRef = new CORFileReference(name); + files.putIfAbsent(name, fileRef); + long fileSize = copyFilesToLocal(fileRef, false, false); + if (fileSize > 0) { + copyCount++; + totalSize += fileSize; + copiedFileNames.add(name); + } + } + + local.sync(copiedFileNames); + PERF_LOGGER.end(start, -1, "[{}] Copied {} files totaling {}", indexPath, copyCount, humanReadableByteCount(totalSize)); + } + + private long copyFilesToLocal(CORFileReference reference, boolean sync, boolean logDuration) { + String name = reference.name; + boolean success = false; + boolean copyAttempted = false; + long fileSize = 0; + try { + if (!LuceneNgIndexCopier.existsLocally(localDir, name)) { + long perfStart = -1; + if (logDuration) { + perfStart = PERF_LOGGER.start(); + } + + fileSize = remote.fileLength(name); + LocalIndexFile file = new LocalIndexFile(local, name, fileSize, true); + long start = indexCopier.startCopy(file); + copyAttempted = true; + + local.copyFrom(remote, name, name, IOContext.READ); + reference.markValid(); Review Comment: **Concurrent copies of the same file can delete each other's result** Lucene 9 differs from 4 here: `Directory.copyFrom` opens the remote input *before* creating the local output, `FSIndexOutput` opens with `CREATE_NEW`, and a failed `copyFrom` deletes the destination. When two generations share a local dir and both queue a copy of the same file, both pass `existsLocally` (L204); the second `createOutput` fails with `FileAlreadyExistsException` and `copyFrom` deletes the **winner's** file. The winner has already called `markValid()` (L216) before `sync` (L219), so it serves a file that no longer exists, and `doneCopy` is skipped (stats drift). `ConcurrentCopyOnReadDirectoryTest` now blocks the local `createOutput` instead of the remote `openInput`, which recreates Lucene 4 ordering, so this window isn't covered. Fix: claim the file in a copier-wide map (losers wait for completion), or copy to a temp name + atomic move treating "already exists" as success; `markValid()` only after `sync`; `doneCopy` in `finally`. ########## oak-search-lucene-ng/src/main/java/org/apache/jackrabbit/oak/plugins/index/luceneNg/internal/editor/LuceneNgFulltextIndexWriterFactory.java: ########## @@ -0,0 +1,75 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.jackrabbit.oak.plugins.index.luceneNg.internal.editor; + +import org.apache.jackrabbit.oak.plugins.index.luceneNg.LuceneNgIndexDefinition; +import org.apache.jackrabbit.oak.plugins.index.luceneNg.LuceneNgIndexStorage; +import org.apache.jackrabbit.oak.plugins.index.luceneNg.directory.OakDirectory; +import org.apache.jackrabbit.oak.plugins.index.search.IndexDefinition; +import org.apache.jackrabbit.oak.plugins.index.search.spi.editor.FulltextIndexWriter; +import org.apache.jackrabbit.oak.plugins.index.search.spi.editor.FulltextIndexWriterFactory; +import org.apache.jackrabbit.oak.spi.commit.CommitInfo; +import org.apache.jackrabbit.oak.spi.state.NodeBuilder; +import org.apache.lucene.document.Document; +import org.apache.lucene.index.IndexWriter; +import org.apache.lucene.index.IndexWriterConfig; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.io.IOException; +import java.io.UncheckedIOException; + +/** + * Opens the same {@link OakDirectory}-backed Lucene {@link IndexWriter} that + * {@code LuceneNgIndexEditor}'s constructor previously opened directly, wrapped behind the + * {@link FulltextIndexWriterFactory} shape {@code FulltextIndexEditorContext} expects. + * + * <p>{@link #newInstance} does not declare a checked exception (per the + * {@link FulltextIndexWriterFactory} interface), so any {@link IOException} raised while + * opening the directory or writer is wrapped in an {@link UncheckedIOException}.</p> + */ +public class LuceneNgFulltextIndexWriterFactory implements FulltextIndexWriterFactory<Document> { + + private static final Logger LOG = LoggerFactory.getLogger(LuceneNgFulltextIndexWriterFactory.class); + + @Override + public FulltextIndexWriter<Document> newInstance(IndexDefinition definition, NodeBuilder definitionBuilder, + CommitInfo commitInfo, boolean reindex) { + LuceneNgIndexDefinition luceneNgDefinition = (LuceneNgIndexDefinition) definition; + String indexName = luceneNgDefinition.getIndexName(); + NodeBuilder storage = LuceneNgIndexStorage.getOrCreateStorageBuilder(definitionBuilder); + + try { + OakDirectory directory = new OakDirectory(storage, indexName, false); Review Comment: **CopyOnWrite / read-before-write is missing on the write path — this will hurt async-indexing performance on large indexes** The `IndexWriter` here is opened directly on the remote-backed `OakDirectory`. Legacy (`DefaultDirectoryFactory.newInstance`, L60-87) wraps the write directory in two ways, both enabled by default: 1. **Read-before-write prefetch** (`READ_BEFORE_WRITE`): the existing segments are synced to local disk before writing, *"to avoid having to stream it when merging."* 2. **CopyOnWrite** (`enableCopyOnWriteSupport`, default `true`): new segments are written locally and uploaded in the background, and reads of existing segments are served from the local copies. Without these, every async-indexing cycle on the indexing node: - **Reads existing segments from the remote blob store.** Applying the deletes from `updateDocument` reads term dictionaries, and segment merges read whole segments. On a large index (e.g. DAM-sized), one merge can stream GBs over the network, so async-indexing lag grows with index size. - **Transfers the same data twice.** Segments are uploaded directly, and then this node's CopyOnRead downloads them again. The writer also ignores local copies that CopyOnRead already has. - **Uploads inline on each file close**, so the upload isn't overlapped with indexing. This is a scalability gap for write-heavy and large indexes, which is where legacy needed these features. Suggestion: at minimum, port the read-before-write prefetch and route the writer's reads through the CopyOnRead local copies. That removes the merge-streaming cost without the full async-upload machinery of CopyOnWrite. ########## oak-search-lucene-ng/src/main/java/org/apache/jackrabbit/oak/plugins/index/luceneNg/LuceneNgIndex.java: ########## @@ -0,0 +1,967 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.jackrabbit.oak.plugins.index.luceneNg; + +import org.apache.jackrabbit.oak.api.Type; +import org.apache.jackrabbit.oak.commons.PathUtils; +import org.apache.jackrabbit.oak.plugins.index.IndexConstants; +import org.apache.jackrabbit.oak.plugins.index.cursor.Cursors; +import org.apache.jackrabbit.oak.plugins.index.search.FieldNames; +import org.apache.jackrabbit.oak.plugins.index.search.IndexDefinition; +import org.apache.jackrabbit.oak.plugins.index.search.IndexDefinition.IndexingRule; +import org.apache.jackrabbit.oak.plugins.index.search.IndexDefinition.SecureFacetConfiguration; +import org.apache.jackrabbit.oak.plugins.index.search.IndexNode; +import org.apache.jackrabbit.oak.plugins.index.search.SizeEstimator; +import org.apache.jackrabbit.oak.plugins.index.search.spi.query.FulltextIndex; +import org.apache.jackrabbit.oak.plugins.index.search.spi.query.FulltextIndexPlanner; +import org.apache.jackrabbit.oak.plugins.index.luceneNg.internal.LuceneNgCursor; +import org.apache.jackrabbit.oak.plugins.index.luceneNg.internal.LuceneNgIndexNode; +import org.apache.jackrabbit.oak.plugins.index.luceneNg.internal.LuceneNgSecureSortedSetDocValuesFacetCounts; +import org.apache.jackrabbit.oak.plugins.index.luceneNg.internal.LuceneNgStatisticalSortedSetDocValuesFacetCounts; +import org.apache.jackrabbit.oak.plugins.memory.PropertyValues; +import org.apache.jackrabbit.oak.spi.query.Cursor; +import org.apache.jackrabbit.oak.spi.query.Filter; +import org.apache.jackrabbit.oak.spi.query.QueryIndex; +import org.apache.jackrabbit.oak.spi.query.QueryIndex.OrderEntry; +import org.apache.jackrabbit.oak.spi.query.QueryIndex.IndexPlan; +import org.apache.jackrabbit.oak.spi.query.QueryConstants; +import org.apache.jackrabbit.oak.spi.query.fulltext.FullTextAnd; +import org.apache.jackrabbit.oak.spi.query.fulltext.FullTextContains; +import org.apache.jackrabbit.oak.spi.query.fulltext.FullTextExpression; +import org.apache.jackrabbit.oak.spi.query.fulltext.FullTextOr; +import org.apache.jackrabbit.oak.spi.query.fulltext.FullTextTerm; +import org.apache.jackrabbit.oak.spi.query.fulltext.FullTextVisitor; +import org.apache.jackrabbit.oak.spi.state.NodeState; +import org.apache.jackrabbit.util.ISO8601; +import org.apache.lucene.analysis.Analyzer; +import org.apache.lucene.analysis.TokenStream; +import org.apache.lucene.analysis.standard.StandardAnalyzer; +import org.apache.lucene.analysis.tokenattributes.CharTermAttribute; +import org.apache.lucene.document.DoublePoint; +import org.apache.lucene.document.LongPoint; +import org.apache.lucene.index.DocValuesType; +import org.apache.lucene.index.FieldInfo; +import org.apache.lucene.index.FieldInfos; +import org.apache.lucene.index.IndexReader; +import org.apache.lucene.index.Term; +import org.apache.lucene.facet.Facets; +import org.apache.lucene.facet.FacetsCollector; +import org.apache.lucene.facet.sortedset.DefaultSortedSetDocValuesReaderState; +import org.apache.lucene.facet.sortedset.SortedSetDocValuesFacetCounts; +import org.apache.lucene.search.BooleanClause.Occur; +import org.apache.lucene.search.BooleanQuery; +import org.apache.lucene.search.IndexSearcher; +import org.apache.lucene.search.MatchAllDocsQuery; +import org.apache.lucene.search.PhraseQuery; +import org.apache.lucene.search.Query; +import org.apache.lucene.search.Sort; +import org.apache.lucene.search.SortField; +import org.apache.lucene.search.SortedSetSortField; +import org.apache.lucene.search.TermQuery; +import org.apache.lucene.search.PrefixQuery; +import org.apache.lucene.search.TermRangeQuery; +import org.apache.lucene.search.TotalHitCountCollector; +import org.apache.lucene.search.BoostQuery; +import org.apache.lucene.search.WildcardQuery; +import org.apache.lucene.util.BytesRef; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import javax.jcr.PropertyType; +import java.io.IOException; +import java.io.StringReader; +import java.util.ArrayList; +import java.util.Locale; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.atomic.AtomicReference; +import java.util.function.Predicate; +import java.util.stream.Collectors; + +/** + * Lucene 9 query index implementation. + * Executes queries against Lucene 9 indexes. + */ +public class LuceneNgIndex extends FulltextIndex { + + private static final Logger LOG = LoggerFactory.getLogger(LuceneNgIndex.class); + // Must equal FulltextIndexPlanner.ATTR_FACET_FIELDS — the inherited FulltextIndexPlanner + // sets facet fields on the plan under this key; query(IndexPlan) reads them back. + private static final String ATTR_FACET_FIELDS = "oak.facet.fields"; + + private final LuceneNgIndexTracker tracker; + private final String indexPath; + + public LuceneNgIndex(LuceneNgIndexTracker tracker, String indexPath) { + this.tracker = tracker; + this.indexPath = indexPath; + } + + // ===== FulltextIndex abstract hooks ===== + // Cost estimation and plan building come from the inherited FulltextIndexPlanner, which + // only offers a plan for properties this index actually declares, matching + // LucenePropertyIndex and ElasticIndex. getCost(Filter,...), getPlan(Filter,...) and + // query(Filter,...) are unsupported here (inherited default throws). + + @Override + protected LuceneNgIndexNode acquireIndexNode(String indexPath) { + return tracker.acquireIndexNode(indexPath); + } + + @Override + protected LuceneNgIndexNode acquireIndexNode(IndexPlan plan) { + return (LuceneNgIndexNode) super.acquireIndexNode(plan); + } + + @Override + protected String getType() { + return LuceneNgIndexConstants.TYPE_LUCENE9; + } + + @Override + public String getIndexName() { + return LuceneNgIndexConstants.TYPE_LUCENE9; + } + + @Override + protected SizeEstimator getSizeEstimator(IndexPlan plan) { + // Port of LucenePropertyIndex.getSizeEstimator: a bounded count-only search over the + // plan's built query. Builds the query via buildQuery(plan.getFilter(), getPlanResult(plan)), + // the same PlanResult-driven construction the executed query uses. Note: LuceneNg's + // query(IndexPlan,...) returns its own LuceneNgCursor, which supplies its own size, so this + // estimator is not on the hot path today — but the hook is abstract and must be implemented + // correctly. + return () -> { + LuceneNgIndexNode indexNode = acquireIndexNode(plan); + if (indexNode == null) { + return -1L; + } + try { + IndexSearcher searcher = indexNode.getSearcher(); + if (searcher == null) { + return -1L; + } + Query query = buildQuery(plan.getFilter(), getPlanResult(plan)); + TotalHitCountCollector collector = new TotalHitCountCollector(); + searcher.search(query, collector); + int totalHits = collector.getTotalHits(); + LOG.debug("Estimated size for query {} is {}", query, totalHits); + return (long) totalHits; + } catch (IOException e) { + LOG.warn("Size-estimate query failed on index {}", indexPath, e); + return -1L; + } finally { + indexNode.release(); + } + }; + } + + @Override + protected Predicate<NodeState> getIndexDefinitionPredicate() { + return state -> LuceneNgIndexConstants.TYPE_LUCENE9.equals( + state.getString(IndexConstants.TYPE_PROPERTY_NAME)); + } + + @Override + protected String getFulltextRequestString(IndexPlan plan, IndexNode indexNode, NodeState rootState) { + // The diagnostic representation of the query this plan would run — the same Lucene + // Query buildQuery(...) constructs for execution. + return buildQuery(plan.getFilter(), getPlanResult(plan)).toString(); + } + + @Override + protected boolean filterReplacedIndexes() { + return false; // matches this module's current behavior — no blue/green mount-info concept yet + } + + @Override + protected boolean runIsActiveIndexCheck() { + return false; // matches ElasticIndex's choice; LuceneNg has no active-index-check concept yet + } + + private Query buildQuery(Filter filter, FulltextIndexPlanner.PlanResult planResult) { + FullTextExpression ft = filter.getFullTextConstraint(); + + // Strip rep:facet pseudo-restrictions and function restrictions we don't index. + // Function restrictions (e.g. "function*@:localname") are paired with their dedicated + // equivalents (e.g. ":localname") and are handled by createPropertyQuery(); including + // them as separate clauses would produce a term query on a non-existent field. + // + // A property restriction only becomes a Lucene clause when the planner validated it — + // matching LucenePropertyIndex.addNonFullTextConstraints, which skips any restriction + // whose planResult.getPropDefn(pr) is null (undeclared/unindexed property) and leaves it + // for the query engine to post-filter instead. The localname() pseudo-restriction has no + // declared PropertyDefinition, so it is gated on evaluateNodeNameRestriction() instead, + // exactly as legacy does. + List<Filter.PropertyRestriction> propRestrictions = filter.getPropertyRestrictions() + .stream() + .filter(pr -> !QueryConstants.REP_FACET.equals(pr.propertyName)) + .filter(pr -> pr.propertyName == null + || !pr.propertyName.startsWith(QueryConstants.FUNCTION_RESTRICTION_PREFIX)) + .filter(pr -> isPlannerValidated(pr, planResult)) + .collect(Collectors.toList()); + + Query pathQuery = buildPathQuery(filter); + + // Build content query (fulltext and/or property constraints) + Query contentQuery; + if (ft == null && propRestrictions.isEmpty()) { + contentQuery = new MatchAllDocsQuery(); + } else if (ft != null) { + try (Analyzer analyzer = new StandardAnalyzer()) { + Query ftQuery = getFullTextQuery(ft, analyzer); + LOG.debug("Building full-text query: {}", ftQuery); + if (!propRestrictions.isEmpty()) { + BooleanQuery.Builder bq = new BooleanQuery.Builder(); + bq.add(ftQuery, Occur.MUST); + for (Filter.PropertyRestriction pr : propRestrictions) { + Query propQuery = createPropertyQuery(pr); + if (propQuery != null) { + bq.add(propQuery, Occur.MUST); + } + } + contentQuery = bq.build(); + } else { + contentQuery = ftQuery; + } + } + } else if (propRestrictions.size() == 1) { + Query q = createPropertyQuery(propRestrictions.get(0)); + contentQuery = q != null ? q : new MatchAllDocsQuery(); + } else { + BooleanQuery.Builder bq = new BooleanQuery.Builder(); + for (Filter.PropertyRestriction pr : propRestrictions) { + Query propQuery = createPropertyQuery(pr); + if (propQuery != null) { + bq.add(propQuery, Occur.MUST); + } + } + contentQuery = bq.build(); + } + + if (pathQuery == null) { + return contentQuery; + } + BooleanQuery.Builder combined = new BooleanQuery.Builder(); + combined.add(contentQuery, Occur.MUST); + combined.add(pathQuery, Occur.FILTER); + return combined.build(); + } + + /** + * Decides whether a property restriction may be turned into a Lucene query clause, driven by + * the {@link FulltextIndexPlanner.PlanResult} the planner already built and attached to the plan + * (rather than re-deciding from the raw {@link Filter}). Mirrors + * {@code LucenePropertyIndex.addNonFullTextConstraints}: + * <ul> + * <li>the {@code localname()} pseudo-restriction is retained only when the planner marked the + * node-name restriction as evaluable ({@link FulltextIndexPlanner.PlanResult#evaluateNodeNameRestriction()});</li> + * <li>every other property restriction is retained only when the planner matched it to a + * declared/indexed property ({@link FulltextIndexPlanner.PlanResult#getPropDefn} is non-null) — + * restrictions on undeclared properties are dropped here and left for the query engine to + * post-filter, exactly as legacy does.</li> + * </ul> + */ + private static boolean isPlannerValidated(Filter.PropertyRestriction pr, + FulltextIndexPlanner.PlanResult planResult) { + // In real query execution the plan is always built by the inherited FulltextIndexPlanner, + // so getPlanResult(plan) is non-null (the same assumption LucenePropertyIndex makes). A null + // PlanResult only arises for lower-level building-block callers that construct a plan without + // going through the planner (e.g. mock-plan unit tests). With no planner decision to be + // consistent with, there is nothing to gate on, so we retain the restriction — i.e. fall back + // to the pre-D3 "derive every constraint from the raw Filter" behavior. + if (planResult == null) { + return true; + } + if (QueryConstants.RESTRICTION_LOCAL_NAME.equals(pr.propertyName)) { + return planResult.evaluateNodeNameRestriction(); + } + return planResult.getPropDefn(pr) != null; + } + + /** + * Translates the Oak PathRestriction to a Lucene query clause, + * or returns null for NO_RESTRICTION (no clause added). + */ + @org.jetbrains.annotations.Nullable + private Query buildPathQuery(Filter filter) { + Filter.PathRestriction restriction = filter.getPathRestriction(); + if (restriction == null) { + return null; + } + String path = filter.getPath(); + switch (restriction) { + case ALL_CHILDREN: + if ("/".equals(path)) { + return null; // matches everything + } + return new PrefixQuery(new Term(FieldNames.PATH, path + "/")); + case DIRECT_CHILDREN: + return new TermQuery(new Term(LuceneNgIndexConstants.FIELD_PARENT_PATH, path)); + case EXACT: + return new TermQuery(new Term(FieldNames.PATH, path)); + case PARENT: + if ("/".equals(path)) { + // root has no parent — match nothing + return new TermQuery(new Term(FieldNames.PATH, "\u0000")); + } + int lastSlash = path.lastIndexOf('/'); + String parentPath = lastSlash == 0 ? "/" : path.substring(0, lastSlash); + return new TermQuery(new Term(FieldNames.PATH, parentPath)); + case NO_RESTRICTION: + default: + return null; + } + } + + /** + * Creates a Lucene Query for a property restriction. + * Handles equality, range, NOT NULL, NULL, NOT, and IN queries. + * Based on legacy LuceneIndex pattern. + */ + private Query createPropertyQuery(Filter.PropertyRestriction pr) { + String propertyName = pr.propertyName; + + // localname() restriction — maps to the NODE_NAME StringField + if (QueryConstants.RESTRICTION_LOCAL_NAME.equals(propertyName)) { + return createLocalNameQuery(pr); + } + + // Function restrictions (e.g. "function*@:localname", "function*lower*@name") are + // only supported when the index has an explicit function property definition. + // We don't support that yet, so skip these to avoid false negatives. + if (propertyName.startsWith(QueryConstants.FUNCTION_RESTRICTION_PREFIX)) { + return null; + } + + // Skip special properties (rep:facet etc.) + if (propertyName.startsWith("rep:") || propertyName.startsWith("oak:")) { + return null; + } + + // Handle IS NOT NULL: matches all documents that have the property indexed + if (pr.isNotNullRestriction()) { Review Comment: **`IS NOT NULL` and typed `IN` return 0 rows on Long/Double/Date properties** Typed properties are indexed only as `LongPoint`/`DoublePoint` (no terms), but: - **IS NOT NULL** (here) builds `TermRangeQuery(prop, null, null)` — matches nothing on a point-only field. The planner lets NOT NULL through (`FulltextIndexPlanner` L281 only skips NULL), and `indexNotNullProperty` is a no-op, so the result is silently empty. Legacy uses `:notNullProps` when `notNullCheckEnabled`, otherwise a full numeric range per type (`LucenePropertyIndex` L1208, L1237-1239, L1273-1275). - **Type dispatch** (`determinePropertyType(pr)`, L370/L386) only looks at `first/last/not` — it ignores `pr.list`, `pr.propertyType` and the definition's declared type. So `[count] IN (1,2,3)` on a Long property is treated as STRING and builds `TermQuery`s on a point field → 0 rows (the numeric `IN` branch `np.set`, L453, is unreachable). Same for `[count] = '5'`. Fix: use the inherited `FulltextIndex.determinePropertyType(defn, pr)` like legacy, and handle NOT NULL per type (or via `:notNullProps`). There are no tests for typed `IN` / `IS NOT NULL` today. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
