From 105389490b0a8fe2647ead67226c45d34f4091fc Mon Sep 17 00:00:00 2001 From: Vladislav Pyatkov Date: Thu, 8 Oct 2026 22:06:53 +0300 Subject: [PATCH 1/2] IGNITE-28897 Return detailed result from transactional lock internal API --- .../calcite/exec/ExecutionServiceImpl.java | 2 +- .../processors/cache/GridCacheAdapter.java | 272 +++++++--------- .../processors/cache/GridCacheProxyImpl.java | 4 +- .../processors/cache/IgniteInternalCache.java | 27 +- .../distributed/dht/GridDhtCacheEntry.java | 29 ++ .../distributed/dht/GridDhtLockFuture.java | 103 +++++- .../dht/GridDhtTransactionalCacheAdapter.java | 42 ++- .../dht/GridDhtTxLocalAdapter.java | 18 +- .../colocated/GridDhtColocatedLockFuture.java | 51 ++- .../distributed/near/GridNearLockFuture.java | 63 +++- .../distributed/near/GridNearLockRequest.java | 24 ++ .../near/GridNearLockResponse.java | 28 +- .../distributed/near/GridNearTxLocal.java | 37 ++- .../cache/transactions/IgniteTxEntry.java | 35 ++ .../feature/CoreFeatureRegistry.java | 3 + ...heVersionedEntryTransactionalLockTest.java | 302 +++++++++++++++++- 16 files changed, 846 insertions(+), 194 deletions(-) diff --git a/modules/calcite/src/main/java/org/apache/ignite/internal/processors/query/calcite/exec/ExecutionServiceImpl.java b/modules/calcite/src/main/java/org/apache/ignite/internal/processors/query/calcite/exec/ExecutionServiceImpl.java index 3824416e78be7..280cd061e9635 100644 --- a/modules/calcite/src/main/java/org/apache/ignite/internal/processors/query/calcite/exec/ExecutionServiceImpl.java +++ b/modules/calcite/src/main/java/org/apache/ignite/internal/processors/query/calcite/exec/ExecutionServiceImpl.java @@ -946,7 +946,7 @@ private static boolean tryAcquireLocks( batchWaitMs = -1L; } - if (!batch.getKey().lockTxEntries(batch.getValue().values(), batchWaitMs)) { + if (batch.getKey().lockTxEntries(batch.getValue().values(), batchWaitMs).containsValue(false)) { locked = false; break; } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheAdapter.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheAdapter.java index 6b32f9408fb02..b7587b46ea1a0 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheAdapter.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheAdapter.java @@ -30,6 +30,7 @@ import java.util.HashMap; import java.util.HashSet; import java.util.Iterator; +import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.NoSuchElementException; @@ -124,6 +125,7 @@ import org.apache.ignite.internal.util.future.GridFinishedFuture; import org.apache.ignite.internal.util.future.GridFutureAdapter; import org.apache.ignite.internal.util.lang.GridCloseableIterator; +import org.apache.ignite.internal.util.lang.GridClosureException; import org.apache.ignite.internal.util.lang.GridPlainCallable; import org.apache.ignite.internal.util.lang.GridPlainRunnable; import org.apache.ignite.internal.util.tostring.GridToStringExclude; @@ -165,6 +167,7 @@ import static org.apache.ignite.internal.processors.dr.GridDrType.DR_LOAD; import static org.apache.ignite.internal.processors.dr.GridDrType.DR_NONE; import static org.apache.ignite.internal.processors.metric.impl.MetricUtils.cacheMetricsRegistryName; +import static org.apache.ignite.internal.processors.rollingupgrade.feature.CoreFeatureRegistry.VERSIONED_TX_LOCK_FEATURE; import static org.apache.ignite.internal.processors.task.TaskExecutionOptions.options; import static org.apache.ignite.internal.thread.pool.IgniteThreadPoolExecutor.newFixedThreadPool; import static org.apache.ignite.transactions.TransactionConcurrency.OPTIMISTIC; @@ -3040,10 +3043,10 @@ public CacheMetricsImpl metrics0() { } /** {@inheritDoc} */ - @Override public boolean lockTxEntries(Collection> entries, long waitTimeout) - throws IgniteCheckedException { - A.notNull(entries, "entries"); - + @Override public Map, Boolean> lockTxEntries( + Collection> entries, + long waitTimeout + ) throws IgniteCheckedException { return lockTxEntriesAsync(entries, waitTimeout).get(); } @@ -3051,11 +3054,18 @@ public CacheMetricsImpl metrics0() { @Override public IgniteInternalFuture lockTxEntryAsync(CacheEntry entry, long waitTimeout) { A.notNull(entry, "entry"); - return lockTxEntriesAsync(Collections.singleton(entry), waitTimeout); + return lockTxEntriesAsync(Collections.singleton(entry), waitTimeout).chain(fut -> { + try { + return fut.get().get(entry); + } + catch (IgniteCheckedException e) { + throw new GridClosureException(e); + } + }); } /** {@inheritDoc} */ - @Override public IgniteInternalFuture lockTxEntriesAsync( + @Override public IgniteInternalFuture, Boolean>> lockTxEntriesAsync( Collection> entries, long waitTimeout ) { @@ -3071,193 +3081,151 @@ public CacheMetricsImpl metrics0() { return new GridFinishedFuture<>( new IgniteCheckedException("Failed to acquire transactional lock in optimistic transaction.")); - // Wait for previous per-transaction async operations to finish. + if (!ctx.kernalContext().rollingUpgrade().features().isActive(VERSIONED_TX_LOCK_FEATURE)) + return new GridFinishedFuture<>(new IgniteCheckedException( + "Primary-side versioned transactional locks require rolling upgrade to be finalized.")); + tx.txState().awaitLastFuture(); if (!tx.init()) return new GridFinishedFuture<>(new IgniteTxRollbackCheckedException( "Failed to acquire transactional lock because transaction has been completed: " + tx)); - if (entries.isEmpty()) - return new GridFinishedFuture<>(true); + List> inputs = new ArrayList<>(entries); + Map, Boolean> results = new HashMap<>(); + Map enlisted = new LinkedHashMap<>(); + Map snapshots = new HashMap<>(); + CacheOperationContext opCtx = ctx.operationContextPerCall(); + + for (CacheEntry entry : inputs) + A.notNull(entry, "entry"); try { tx.addActiveCache(ctx, false); - } - catch (IgniteCheckedException e) { - return new GridFinishedFuture<>(e); - } - Collection keys = new ArrayList<>(entries.size()); - List txEntries = new ArrayList<>(entries.size()); - List expVers = new ArrayList<>(entries.size()); - Set txKeys = new HashSet<>(entries.size()); + for (CacheEntry entry : inputs) { + KeyCacheObject key = ctx.toCacheKeyObject(entry.getKey()); - CacheOperationContext opCtx = ctx.operationContextPerCall(); + if (enlisted.containsKey(key)) { + throw new IgniteCheckedException("Failed to acquire transactional lock because " + + "entry is duplicated [key=" + key + ", entry=" + entry + ']'); + } - for (CacheEntry entry : entries) { - A.notNull(entry, "entry"); + IgniteTxEntry oldEntry = tx.entry(ctx.txKey(key)); - KeyCacheObject key = ctx.toCacheKeyObject(entry.getKey()); - IgniteTxKey txKey = ctx.txKey(key); + if (oldEntry != null && (oldEntry.op() != READ || oldEntry.locked())) { + results.put(entry, true); - if (!txKeys.add(txKey)) - continue; + continue; + } - IgniteTxEntry lockedTxEntry = tx.entry(txKey); + if (!(entry.version() instanceof GridCacheVersion)) { + throw new IgniteCheckedException( + "Failed to acquire transactional lock for entry with unsupported version type: " + entry); + } - if (lockedTxEntry != null && (lockedTxEntry.op() != READ || lockedTxEntry.locked())) - continue; + snapshots.put(key, oldEntry == null ? null : oldEntry.copy()); - if (!(entry.version() instanceof GridCacheVersion)) { - tx.removeAndUnlockTxEntries(txEntries); + GridCacheEntryEx cached = ctx.isColocated() + ? ctx.colocated().entryExx(key, tx.topologyVersion(), true) : entryEx(key); + IgniteTxEntry txEntry = tx.addEntry( + READ, + ctx.toCacheObject(entry.getValue()), + null, + null, + cached, + null, + null, + true, + -1L, + -1L, + null, + opCtx != null && opCtx.skipStore(), + opCtx != null && opCtx.skipReadThrough(), + opCtx != null && opCtx.keepBinaryInInterceptor(), + opCtx != null && opCtx.isKeepBinary(), + CU.isNearEnabled(ctx) + ); - return new GridFinishedFuture<>(new IgniteCheckedException("Failed to acquire transactional lock for entry with " + - "unsupported version type [entry=" + entry + ", version=" + entry.version() + ']')); + txEntry.expectedLockVersion((GridCacheVersion)entry.version()); + txEntry.versionedLockResult(null); + enlisted.put(key, txEntry); } + } + catch (IgniteCheckedException e) { + for (Map.Entry entry : enlisted.entrySet()) + tx.removeFailedLockEntry(entry.getValue(), snapshots.get(entry.getKey())); - CacheObject val = ctx.toCacheObject(entry.getValue()); - GridCacheEntryEx entryEx = ctx.isColocated() ? ctx.colocated().entryExx(key, tx.topologyVersion(), true) : entryEx(key); + return new GridFinishedFuture<>(e); + } - IgniteTxEntry txEntry = tx.addEntry( - READ, - val, - null, - null, - entryEx, - null, - null, - true, - -1L, - -1L, - null, - opCtx != null && opCtx.skipStore(), - opCtx != null && opCtx.skipReadThrough(), - opCtx != null && opCtx.keepBinaryInInterceptor(), - opCtx != null && opCtx.isKeepBinary(), - CU.isNearEnabled(ctx) - ); + if (enlisted.isEmpty()) + return new GridFinishedFuture<>(results); - keys.add(key); - txEntries.add(txEntry); - expVers.add((GridCacheVersion)entry.version()); - } + GridFutureAdapter, Boolean>> res = new GridFutureAdapter<>(); + GridCacheAdapter.FutureHolder holder = tx.txState().lastAsyncFuture(); - if (keys.isEmpty()) - return new GridFinishedFuture<>(true); + if (holder != null) { + holder.lock(); - // Acquire transactional lock future from concrete cache implementation. Use txLockAsync which - // delegates to cache-specific lockAllAsync implementations for distributed caches. - long lockWaitStartTime = U.currentTimeMillis(); - long timeout = tx.remainingTime(); - long effectiveWaitTimeout = CU.isWaitTimeoutExpiresFirst(waitTimeout, timeout) ? waitTimeout : timeout; - long lockWaitEndTime = effectiveWaitTimeout > 0 - ? lockWaitStartTime + effectiveWaitTimeout - : effectiveWaitTimeout < 0 ? lockWaitStartTime : 0L; + try { + holder.saveFuture(res); + } + finally { + holder.unlock(); + } + } - IgniteInternalFuture lockFut = txLockAsync(keys, - timeout, + // Keep the existing distributed mapping and lock order: one request per consecutive primary group, + // rather than one request per entry. Individual outcomes travel in the primary's batch response. + IgniteInternalFuture lockFut = txLockAsync( + enlisted.keySet(), + tx.remainingTime(), waitTimeout, tx, - /*isRead*/true, - /*retval*/false, + true, + false, tx.isolation(), - /*invalidate*/false, - /*createTtl*/0L, - /*accessTtl*/0L); - - IgniteInternalFuture res = new GridEmbeddedFuture<>( - lockFut, - (locked, ex) -> { - if (ex != null) - return new GridFinishedFuture<>(ex); - - if (!locked) { - tx.removeAndUnlockTxEntries(txEntries); - - return new GridFinishedFuture<>(false); - } - - try { - for (int i = 0; i < txEntries.size(); i++) { - IgniteTxEntry txEntry = txEntries.get(i); - EntryGetResult getRes; - int retryCnt = 0; - - while (true) { - try { - GridCacheEntryEx cached = txEntry.cached(); - - getRes = cached.innerGetVersioned( - null, - tx, - /*update-metrics*/false, - /*event*/false, - null, - tx.resolveTaskName(), - null, - false, - null); - - break; - } - catch (GridCacheEntryRemovedException ignored) { - // NOWAIT still permits one immediate retry because renewing an obsolete entry does not - // wait for a lock. Other modes stop when their effective timeout expires. - boolean allowNowaitRetry = waitTimeout < 0 && retryCnt == 0; - - if (!allowNowaitRetry && lockWaitEndTime != 0 - && U.currentTimeMillis() >= lockWaitEndTime) { - getRes = null; + false, + 0L, + 0L + ); - break; - } + lockFut.listen(fut -> { + try { + boolean completed = fut.get(); - if (log.isDebugEnabled()) - log.debug("Got removed exception in lockTxEntries postLock (will retry): " - + txEntry.cached()); + if (tx.remainingTime() == -1) + throw tx.timeoutException(); - KeyCacheObject key = txEntry.key(); - GridCacheEntryEx cached = ctx.isColocated() - ? ctx.colocated().entryExx(key, tx.topologyVersion(), true) - : entryEx(key, tx.topologyVersion()); + if (tx.isRollbackOnly()) + throw tx.rollbackException(); - txEntry.cached(cached); - retryCnt++; - } - } + for (CacheEntry entry : inputs) { + if (!results.containsKey(entry)) { + IgniteTxEntry txEntry = enlisted.get(ctx.toCacheKeyObject(entry.getKey())); - if (getRes == null || !expVers.get(i).equals(getRes.version())) { - tx.removeAndUnlockTxEntries(txEntries); + assert !completed || txEntry.versionedLockResult() != null; - return new GridFinishedFuture<>(false); - } + results.put(entry, completed && Boolean.TRUE.equals(txEntry.versionedLockResult())); } - - return new GridFinishedFuture<>(true); } - catch (IgniteCheckedException e) { - tx.removeAndUnlockTxEntries(txEntries); - return new GridFinishedFuture<>(e); - } - } - ); - - // Register this future in transaction's async-holder so that subsequent operations - // that call tx.txState().awaitLastFuture() will wait for it. - GridCacheAdapter.FutureHolder holder = tx.txState().lastAsyncFuture(); + for (Map.Entry entry : enlisted.entrySet()) { + IgniteTxEntry txEntry = entry.getValue(); - if (holder != null) { - holder.lock(); + if (completed && Boolean.TRUE.equals(txEntry.versionedLockResult())) + txEntry.markLocked(); + else + tx.removeFailedLockEntry(txEntry, snapshots.get(entry.getKey())); + } - try { - holder.saveFuture(res); + res.onDone(results); } - finally { - holder.unlock(); + catch (IgniteCheckedException | RuntimeException e) { + res.onDone(e); } - } + }); return res; } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheProxyImpl.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheProxyImpl.java index 21c1fa361fe4a..1ce56e8003e70 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheProxyImpl.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheProxyImpl.java @@ -1336,7 +1336,7 @@ public IgniteInternalCache delegate() { } /** {@inheritDoc} */ - @Override public boolean lockTxEntries(Collection> entries, long waitTimeout) + @Override public Map, Boolean> lockTxEntries(Collection> entries, long waitTimeout) throws IgniteCheckedException { CacheOperationContext prev = gate.enter(opCtx); @@ -1349,7 +1349,7 @@ public IgniteInternalCache delegate() { } /** {@inheritDoc} */ - @Override public IgniteInternalFuture lockTxEntriesAsync( + @Override public IgniteInternalFuture, Boolean>> lockTxEntriesAsync( Collection> entries, long waitTimeout ) { diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/IgniteInternalCache.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/IgniteInternalCache.java index f98849c6761fa..0668d902ddabf 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/IgniteInternalCache.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/IgniteInternalCache.java @@ -1360,19 +1360,24 @@ public boolean lock(K key, long timeout) public boolean lockTxEntry(CacheEntry entry, long waitTimeout) throws IgniteCheckedException; /** - * Acquires transactional locks for the cached objects represented by the given entries if all current cached - * versions match the corresponding entry versions. This method works only in a - * {@link TransactionConcurrency#PESSIMISTIC} transaction. + * Attempts to acquire a transactional lock for each entry, validating its data version on the primary while + * owning the transactional lock and holding the entry mutex. This method works only in a + * {@link TransactionConcurrency#PESSIMISTIC} transaction. A conflicting, missing or changed entry is reported + * as {@code false}; other entries are still attempted. Successful locks are retained until transaction end + * (or an explicit rollback to a savepoint), even when other entries fail. Already owned entries are reused. + * Infrastructure failures complete the operation exceptionally and are not per-entry rejections. * * @param entries Entries whose keys, values and versions should be used. * @param waitTimeout Timeout in milliseconds to wait for locks to be acquired * ({@code 0} to use the transaction timeout, {@code -1} for immediate failure if * locks cannot be acquired immediately). - * @return {@code True} if all locks were acquired with the same entry versions. + * A positive limit is shared by all entries; once it expires the remaining attempts do not wait. + * @return Lock result for each supplied entry, in input iteration order. * @throws IgniteCheckedException If lock acquisition resulted in an error. * @throws NullPointerException If entries is {@code null}. */ - public boolean lockTxEntries(Collection> entries, long waitTimeout) throws IgniteCheckedException; + public Map, Boolean> lockTxEntries(Collection> entries, long waitTimeout) + throws IgniteCheckedException; /** * Asynchronously acquires a transactional lock for the cached object represented by the given entry if the current @@ -1390,19 +1395,19 @@ public boolean lock(K key, long timeout) public IgniteInternalFuture lockTxEntryAsync(CacheEntry entry, long waitTimeout); /** - * Asynchronously acquires transactional locks for the cached objects represented by the given entries if all - * current cached versions match the corresponding entry versions. This method works only in a - * {@link TransactionConcurrency#PESSIMISTIC} transaction. + * Asynchronous counterpart of {@link #lockTxEntries(Collection, long)} with the same partial-success semantics. * * @param entries Entries whose keys, values and versions should be used. * @param waitTimeout Timeout in milliseconds to wait for locks to be acquired * ({@code 0} to use the transaction timeout, {@code -1} for immediate failure if * locks cannot be acquired immediately). - * @return Future that resolves to {@code true} if all locks were acquired and all versions matched, or to - * {@code false} otherwise. + * @return Future containing a lock result for each supplied entry. * @throws NullPointerException If entries is {@code null}. */ - public IgniteInternalFuture lockTxEntriesAsync(Collection> entries, long waitTimeout); + public IgniteInternalFuture, Boolean>> lockTxEntriesAsync( + Collection> entries, + long waitTimeout + ); /** * Checks if current thread owns a lock on this key. diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridDhtCacheEntry.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridDhtCacheEntry.java index 05de436568e24..1e0e35232c46d 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridDhtCacheEntry.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridDhtCacheEntry.java @@ -49,6 +49,7 @@ import org.apache.ignite.internal.util.typedef.F; import org.apache.ignite.internal.util.typedef.internal.CU; import org.apache.ignite.internal.util.typedef.internal.S; +import org.apache.ignite.internal.util.typedef.internal.U; import org.apache.ignite.lang.IgniteBiTuple; import org.apache.ignite.lang.IgniteClosure; import org.jetbrains.annotations.Nullable; @@ -316,6 +317,34 @@ public GridDhtCacheEntry( return cand; } + /** + * Checks a conditional lock after this transaction became the MVCC owner and before acknowledging the lock. + * The entry mutex protects both the owner test and the data-version comparison. After it is released, the + * transaction's MVCC ownership prevents another transactional writer from changing the version until unlock. + * Merely registering an MVCC candidate is not sufficient: a preceding owner may commit while it is waiting. + * + * @param lockVer Transaction lock version (not the expected data version). + * @param expVer Expected data version. + * @return Whether this transaction owns a present entry with the expected version. + * @throws GridCacheEntryRemovedException If the entry is obsolete. + */ + boolean checkLockVersion(GridCacheVersion lockVer, GridCacheVersion expVer) + throws GridCacheEntryRemovedException { + lockEntry(); + + try { + checkObsolete(); + + long expireTime = expireTimeExtras(); + + return lockedBy(lockVer) && hasValueUnlocked() && expVer.equals(ver) + && (expireTime == 0 || expireTime > U.currentTimeMillis()); + } + finally { + unlockEntry(); + } + } + /** {@inheritDoc} */ @Override public boolean tmLock(IgniteInternalTx tx, long timeout, diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridDhtLockFuture.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridDhtLockFuture.java index 44d8e56abc97f..7ef3183012422 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridDhtLockFuture.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridDhtLockFuture.java @@ -169,6 +169,9 @@ public final class GridDhtLockFuture extends GridCacheCompoundIdentityFuture pendingLocks; + /** Entries enlisted by a conditional lock, including an entry rejected before candidate creation. */ + private final Map versionedEntries = new LinkedHashMap<>(); + /** TTL for create operation. */ private final long createTtl; @@ -348,7 +351,24 @@ void addInvalidPartition(GridCacheContext cacheCtx, int invalidPart) { * @return Entries. */ public Collection entries() { - return F.view(entries, F.notNull()); + return F.view(entries, entry -> entry != null && !rejected(entry)); + } + + /** @return Whether this entry was rejected by the current conditional batch. */ + private boolean rejected(GridCacheEntryEx entry) { + IgniteTxEntry txEntry = tx == null ? null : tx.entry(entry.txKey()); + + return versionedEntries.containsKey(entry) && txEntry != null + && Boolean.FALSE.equals(txEntry.versionedLockResult()); + } + + /** Rejects only this entry, preserving the other locks and the transaction. */ + private synchronized void rejectEntry(GridDhtCacheEntry entry) throws GridCacheEntryRemovedException { + tx.entry(entry.txKey()).versionedLockResult(false); + pendingLocks.remove(entry.key()); + + if (entry.candidate(lockVer) != null) + entry.removeLock(lockVer); } /** @@ -427,6 +447,13 @@ private boolean isInvalidate() { if (timedOut) return null; + if (tx != null) { + IgniteTxEntry txEntry = tx.entry(entry.txKey()); + + if (txEntry != null && txEntry.versionedLockPending()) + versionedEntries.put(entry, txEntry.expectedLockVersion()); + } + // Add local lock first, as it may throw GridCacheEntryRemovedException. GridCacheMvccCandidate c = entry.addDhtLocal( nearNodeId, @@ -446,7 +473,9 @@ private boolean isInvalidate() { if (log.isDebugEnabled()) log.debug("Failed to acquire lock with negative timeout: " + entry); - if (CU.isWaitTimeoutExpiresFirst(waitTimeout, timeout)) + if (versionedEntries.containsKey(entry)) + rejectEntry(entry); + else if (CU.isWaitTimeoutExpiresFirst(waitTimeout, timeout)) onComplete(false, false, false, false); else onFailed(); @@ -631,11 +660,20 @@ private void readyLocks() { if (entry == null) break; // While. + if (rejected(entry)) + break; + try { CacheLockCandidates owners = entry.readyLock(lockVer); if (lockTimeout() < 0) { if (owners == null || !owners.hasCandidate(lockVer)) { + if (versionedEntries.containsKey(entry)) { + rejectEntry((GridDhtCacheEntry)entry); + + break; + } + // We did not send any requests yet. if (CU.isWaitTimeoutExpiresFirst(waitTimeout, timeout)) onComplete(false, false, false, false); @@ -803,7 +841,7 @@ private synchronized boolean onComplete(boolean success, boolean stopping, boole } try { - if (err == null && !stopping) + if (err == null && !stopping && (success || versionedEntries.isEmpty())) loadMissingFromStore(); } finally { @@ -839,6 +877,10 @@ public void map() { readyLocks(); + // Removing rejected candidates does not notify this future of a new owner. + if (!versionedEntries.isEmpty() && checkLocks()) + map(entries()); + if (lockTimeout() > 0 && !isDone()) { // Prevent memory leak if future is completed by call to readyLocks. timeoutObj = new LockTimeoutObject(); @@ -872,6 +914,31 @@ private void map(Iterable entries) { mapped = true; } + // All local candidates are owners now. Validate on the primary before sending any backup requests. + // Ownership remains held across the comparison and replication, so another transaction cannot write + // between this check and the successful reply, including when this future previously waited for a lock. + for (Map.Entry entry : versionedEntries.entrySet()) { + if (rejected(entry.getKey())) + continue; + + try { + entry.getKey().unswap(false); + + if (entry.getKey().checkLockVersion(lockVer, entry.getValue())) + tx.entry(entry.getKey().txKey()).versionedLockResult(true); + else + rejectEntry(entry.getKey()); + } + catch (GridCacheEntryRemovedException e) { + tx.entry(entry.getKey().txKey()).versionedLockResult(false); + } + catch (IgniteCheckedException e) { + onError(e); + + return; + } + } + try { if (log.isDebugEnabled()) log.debug("Mapping entry for DHT lock future: " + this); @@ -1110,7 +1177,7 @@ private void loadMissingFromStore() { final GridCacheVersion ver = version(); - for (GridDhtCacheEntry entry : entries) { + for (GridDhtCacheEntry entry : entries()) { try { entry.unswap(false); @@ -1202,6 +1269,34 @@ private class LockTimeoutObject extends GridTimeoutObjectAdapter { long longOpsDumpTimeout = cctx.tm().longOperationsDumpTimeout(); synchronized (GridDhtLockFuture.this) { + // A conditional wait limit applies to acquiring primary ownership. Once mapping starts we + // already own every entry; replication is governed by the transaction deadline. Rejecting here + // would race with outgoing backup requests and require a distributed partial-lock cancellation. + if (!versionedEntries.isEmpty() && mapped && CU.isWaitTimeoutExpiresFirst(waitTimeout, timeout)) + return; + + if (!versionedEntries.isEmpty() && CU.isWaitTimeoutExpiresFirst(waitTimeout, timeout)) { + for (GridDhtCacheEntry entry : versionedEntries.keySet()) { + if (!pendingLocks.contains(entry.key())) + continue; + + try { + if (entry.lockedBy(lockVer)) + pendingLocks.remove(entry.key()); + else + rejectEntry(entry); + } + catch (GridCacheEntryRemovedException e) { + tx.entry(entry.txKey()).versionedLockResult(false); + pendingLocks.remove(entry.key()); + } + } + + map(entries()); + + return; + } + if (log.isDebugEnabled() || timeout >= longOpsDumpTimeout) { String msg = dumpPendingLocks(); diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridDhtTransactionalCacheAdapter.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridDhtTransactionalCacheAdapter.java index 0994df6c74fc3..cb15c88abc36d 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridDhtTransactionalCacheAdapter.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridDhtTransactionalCacheAdapter.java @@ -1065,7 +1065,8 @@ public IgniteInternalFuture lockAllAsync( req.keepBinaryInInterceptor(), req.keepBinary(), req.waitTimeout(), - req.nearCache()); + req.nearCache(), + req.expectedVersions()); final GridDhtTxLocal t = tx; @@ -1081,7 +1082,7 @@ public IgniteInternalFuture lockAllAsync( boolean lockAcquired = e == null && o != null && o.success() && !t.empty(); // Create response while holding locks. - final GridNearLockResponse resp = createLockReply(nearNode, + GridNearLockResponse resp = createLockReply(nearNode, entries, req, t, @@ -1089,6 +1090,24 @@ public IgniteInternalFuture lockAllAsync( e, lockAcquired); + if (e == null && req.expectedVersions() != null) { + for (GridCacheEntryEx entry : entries) { + IgniteTxEntry txEntry = t.entry(entry.txKey()); + + if (txEntry != null && Boolean.FALSE.equals(txEntry.versionedLockResult())) + t.clearEntry(entry.txKey()); + } + + try { + // The initiating side will keep mappings only for successful keys. + if (t.empty()) + t.rollbackDhtLocal(); + } + catch (IgniteCheckedException ex) { + resp = createLockReply(nearNode, entries, req, t, t.xidVersion(), ex, false); + } + } + assert !t.implicit() : t; assert !t.onePhaseCommit() : t; @@ -1243,6 +1262,18 @@ private GridNearLockResponse createLockReply( res.lockAcquired(lockAcquired); + if (err == null && req.expectedVersions() != null && lockAcquired) { + boolean[] results = new boolean[entries.size()]; + + for (int i = 0; i < entries.size(); i++) { + IgniteTxEntry txEntry = tx.entry(entries.get(i).txKey()); + + results[i] = txEntry != null && Boolean.TRUE.equals(txEntry.versionedLockResult()); + } + + res.lockResults(results); + } + if (err == null) { res.pending(localDhtPendingVersions(entries, mappedVer)); @@ -1262,6 +1293,13 @@ private GridNearLockResponse createLockReply( assert e != null; + if (!res.lockResult(i)) { + res.addValueBytes(null, false, null, null); + i++; + + continue; + } + while (true) { try { // Don't return anything for invalid partitions. diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridDhtTxLocalAdapter.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridDhtTxLocalAdapter.java index 74140402ca968..11be1db03d697 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridDhtTxLocalAdapter.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridDhtTxLocalAdapter.java @@ -565,6 +565,7 @@ private void addMapping( * @param keepBinaryInInterceptor Handle binary in interceptor operation flag. * @param keepBinary Keep binary flag. * @param nearCache {@code True} if near cache enabled on originating node. + * @param expectedVers Expected data versions for conditional locks, or {@code null}. * @return Lock future. */ @SuppressWarnings("ForLoopReplaceableByForEach") @@ -581,7 +582,8 @@ IgniteInternalFuture lockAllAsync( boolean keepBinaryInInterceptor, boolean keepBinary, long waitTimeout, - boolean nearCache + boolean nearCache, + @Nullable GridCacheVersion[] expectedVers ) { try { checkValid(); @@ -662,6 +664,9 @@ IgniteInternalFuture lockAllAsync( txEntry.cached(cached); + if (expectedVers != null) + txEntry.expectedLockVersion(expectedVers[i]); + addReader(msgId, cached, txEntry, topVer); } else { @@ -754,6 +759,9 @@ private IgniteInternalFuture obtainLockAsync( if (isRollbackOnly()) return new GridFinishedFuture<>(rollbackException()); + boolean versionedLock = passedKeys.stream() + .anyMatch(key -> entry(cacheCtx.txKey(key)).versionedLockPending()); + IgniteInternalFuture fut = dhtCache.lockAllAsyncInternal(passedKeys, timeout, waitTimeout, @@ -771,12 +779,15 @@ private IgniteInternalFuture obtainLockAsync( return new GridEmbeddedFuture<>( fut, - new PLC1(ret, true, !CU.isWaitTimeoutExpiresFirst(waitTimeout, timeout)) { + new PLC1(ret, true, !versionedLock && !CU.isWaitTimeoutExpiresFirst(waitTimeout, timeout)) { @Override protected GridCacheReturn postLock(GridCacheReturn ret) throws IgniteCheckedException { assert fut.error() == null : "Lock future completed with an error: " + fut.error(); boolean success = Boolean.TRUE.equals(fut.get()); + if (!success) + checkValid(); + ret.success(success); if (log.isDebugEnabled()) { @@ -788,7 +799,8 @@ private IgniteInternalFuture obtainLockAsync( if (ret.success()) { postLockWrite(cacheCtx, - passedKeys, + versionedLock ? F.view(passedKeys, key -> + !Boolean.FALSE.equals(entry(cacheCtx.txKey(key)).versionedLockResult())) : passedKeys, ret, /*remove*/false, /*retval*/false, diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/colocated/GridDhtColocatedLockFuture.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/colocated/GridDhtColocatedLockFuture.java index 3b7158b4e5dd6..9f4f9efa93d8c 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/colocated/GridDhtColocatedLockFuture.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/colocated/GridDhtColocatedLockFuture.java @@ -140,6 +140,12 @@ public final class GridDhtColocatedLockFuture extends GridCacheCompoundIdentityF /** Lock wait timeout. */ private final long waitTimeout; + /** Whether this operation conditionally locks a data version. */ + private final boolean versionedLock; + + /** Shared wait deadline for conditional primary batches. */ + private final long lockWaitEndTime; + /** Transaction. */ @GridToStringExclude private final GridNearTxLocal tx; @@ -227,6 +233,9 @@ public GridDhtColocatedLockFuture( this.retval = retval; this.timeout = timeout; this.waitTimeout = waitTimeout; + versionedLock = tx != null && keys.stream().anyMatch(key -> + tx.entry(cctx.txKey(key)) != null && tx.entry(cctx.txKey(key)).versionedLockPending()); + lockWaitEndTime = waitTimeout > 0 ? U.currentTimeMillis() + waitTimeout : 0; this.createTtl = createTtl; this.accessTtl = accessTtl; this.skipStore = skipStore; @@ -623,7 +632,7 @@ private synchronized void onError(Throwable t) { if (err != null) success = false; - if (!success && err == null && CU.isWaitTimeoutExpiresFirst(waitTimeout, timeout)) + if (!success && err == null && (versionedLock || CU.isWaitTimeoutExpiresFirst(waitTimeout, timeout))) return onComplete(false, true, false); return onComplete(success, true); @@ -1118,6 +1127,9 @@ private synchronized void map0( key, retval, dhtVer); // Include DHT version to match remote DHT entry. + + if (tx != null && tx.entry(txKey).versionedLockPending()) + req.expectedVersion(tx.entry(txKey).expectedLockVersion()); } explicit = inTx() && cand == null; @@ -1217,6 +1229,9 @@ private void proceedMapping0() final Collection mappedKeys = map.distributedKeys(); final ClusterNode node = map.node(); + if (versionedLock) + req.waitTimeout(remainingWaitTimeout()); + if (node.isLocal()) lockLocally(mappedKeys, req.topologyVersion()); else { @@ -1264,7 +1279,7 @@ private void lockLocally( read, retval, timeout, - waitTimeout, + remainingWaitTimeout(), createTtl, accessTtl, skipStore, @@ -1324,8 +1339,12 @@ private void lockLocally( /** @param keys Locally locked keys. */ private void markLocalDhtLocksAcquired(Collection keys) { if (inTx()) { - for (KeyCacheObject key : keys) - tx.entry(cctx.txKey(key)).markLocked(); + for (KeyCacheObject key : keys) { + IgniteTxEntry entry = tx.entry(cctx.txKey(key)); + + if (!Boolean.FALSE.equals(entry.versionedLockResult())) + entry.markLocked(); + } } else { for (KeyCacheObject key : keys) @@ -1488,9 +1507,23 @@ private boolean errorOrTimeoutOnTopologyVersion(IgniteCheckedException e, boolea * @return Timeout value for this lock future. */ private long lockTimeout() { + // Wait for the primary's conditional rejection and cleanup; only the transaction deadline is local. + if (versionedLock) + return timeout; + return CU.isWaitTimeoutExpiresFirst(waitTimeout, timeout) ? waitTimeout : timeout; } + /** @return Remaining shared wait budget for the next primary batch. */ + private long remainingWaitTimeout() { + if (!versionedLock || waitTimeout <= 0) + return waitTimeout; + + long remaining = lockWaitEndTime - U.currentTimeMillis(); + + return remaining > 0 ? remaining : -1; + } + /** * Lock request timeout object. */ @@ -1736,6 +1769,16 @@ void onResult(GridNearLockResponse res) { int i = 0; for (KeyCacheObject k : keys) { + if (res.hasLockResults()) { + tx.entry(cctx.txKey(k)).versionedLockResult(res.lockResult(i)); + + if (!res.lockResult(i)) { + i++; + + continue; + } + } + IgniteBiTuple oldValTup = valMap.get(k); CacheObject newVal = res.value(i); diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearLockFuture.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearLockFuture.java index 20923899b3201..edf00bd429d9e 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearLockFuture.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearLockFuture.java @@ -137,6 +137,12 @@ public final class GridNearLockFuture extends GridCacheCompoundIdentityFuture + tx.entry(cctx.txKey(key)) != null && tx.entry(cctx.txKey(key)).versionedLockPending()); + lockWaitEndTime = waitTimeout > 0 ? U.currentTimeMillis() + waitTimeout : 0; this.createTtl = createTtl; this.accessTtl = accessTtl; this.skipStore = skipStore; @@ -657,6 +666,9 @@ private boolean checkLocks() { while (true) { GridCacheEntryEx cached = entries.get(i); + if (inTx() && Boolean.FALSE.equals(tx.entry(cached.txKey()).versionedLockResult())) + break; + try { if (!locked(cached)) { if (log.isDebugEnabled()) @@ -723,7 +735,7 @@ private boolean checkLocks() { if (err != null) success = false; - if (!success && err == null && CU.isWaitTimeoutExpiresFirst(waitTimeout, timeout)) + if (!success && err == null && (versionedLock || CU.isWaitTimeoutExpiresFirst(waitTimeout, timeout))) return onComplete(false, true, false); return onComplete(success, true); @@ -1126,6 +1138,9 @@ private void map(Iterable keys, boolean remap, boolean topLocked retval && dhtVer == null, dhtVer); // Include DHT version to match remote DHT entry. + if (tx != null && tx.entry(txKey).versionedLockPending()) + req.expectedVersion(tx.entry(txKey).expectedLockVersion()); + } if (cand.reentry()) @@ -1216,6 +1231,12 @@ private void proceedMapping0() final Collection mappedKeys = map.distributedKeys(); final ClusterNode node = map.node(); + if (versionedLock && waitTimeout > 0) { + long remaining = lockWaitEndTime - U.currentTimeMillis(); + + req.waitTimeout(remaining > 0 ? remaining : -1); + } + if (node.isLocal()) { req.miniId(-1); @@ -1261,6 +1282,14 @@ private void proceedMapping0() int i = 0; for (KeyCacheObject k : mappedKeys) { + if (res.hasLockResults()) { + if (!processLockResult(k, res.lockResult(i))) { + i++; + + continue; + } + } + while (true) { GridNearCacheEntry entry = cctx.near().entryExx(k, req.topologyVersion()); @@ -1434,9 +1463,33 @@ private ClusterTopologyCheckedException newTopologyException(@Nullable Throwable * @return Timeout value for this lock future. */ private long lockTimeout() { + // For conditional locks the primary owns the wait deadline and acknowledges rejection only after + // cleaning up its candidate. An earlier near-side rejection would race with a retry of the same key. + if (versionedLock) + return timeout; + return CU.isWaitTimeoutExpiresFirst(waitTimeout, timeout) ? waitTimeout : timeout; } + /** Applies one primary outcome, removing a rejected local candidate before proceeding to the next batch. */ + private boolean processLockResult(KeyCacheObject key, boolean locked) { + tx.entry(cctx.txKey(key)).versionedLockResult(locked); + + if (!locked) { + GridCacheEntryEx entry = cctx.near().peekEx(key); + + try { + if (entry != null && entry.hasLockCandidate(lockVer)) + entry.removeLock(lockVer); + } + catch (GridCacheEntryRemovedException ignored) { + // The obsolete entry no longer has a local candidate. + } + } + + return locked; + } + /** * Lock request timeout object. */ @@ -1692,6 +1745,14 @@ void onResult(GridNearLockResponse res) { AffinityTopologyVersion topVer = GridNearLockFuture.this.topVer; for (KeyCacheObject k : keys) { + if (res.hasLockResults()) { + if (!processLockResult(k, res.lockResult(i))) { + i++; + + continue; + } + } + while (true) { GridNearCacheEntry entry = cctx.near().entryExx(k, topVer); diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearLockRequest.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearLockRequest.java index 46fae249ebd1f..9a8215a5c2aa8 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearLockRequest.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearLockRequest.java @@ -83,6 +83,30 @@ public class GridNearLockRequest extends GridDistributedLockRequest { @Order(8) long waitTimeout; + /** Expected data versions for conditional locks, indexed in the same order as keys. */ + @Order(value = 9, introducedBy = "VERSIONED_TX_LOCK_FEATURE") + GridCacheVersion[] expectedVers; + + /** @param ver Expected data version for the last added key. */ + public void expectedVersion(@Nullable GridCacheVersion ver) { + if (ver != null) { + if (expectedVers == null) + expectedVers = new GridCacheVersion[dhtVers.length]; + + expectedVers[idx - 1] = ver; + } + } + + /** @param timeout Remaining wait budget when this primary batch is sent. */ + public void waitTimeout(long timeout) { + waitTimeout = timeout; + } + + /** @return Expected data versions, or {@code null} for unconditional locks. */ + @Nullable public GridCacheVersion[] expectedVersions() { + return expectedVers; + } + /** * Empty constructor. */ diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearLockResponse.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearLockResponse.java index b179d094d8c09..dd387e3888100 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearLockResponse.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearLockResponse.java @@ -68,10 +68,32 @@ public class GridNearLockResponse extends GridDistributedLockResponse { @Order(6) boolean compatibleRemapVer; - /** {@code True} if requested locks were acquired. */ + /** Requested locks were acquired, or a conditional batch completed with individual {@link #lockResults}. */ @Order(7) boolean lockAcquired = true; + /** Conditional lock outcomes in request-key order. The common success flag describes batch completion. */ + @Order(value = 8, introducedBy = "VERSIONED_TX_LOCK_FEATURE") + boolean[] lockResults; + + /** @param results Conditional lock outcomes in request-key order. */ + public void lockResults(boolean[] results) { + lockResults = results; + } + + /** @return Whether this response contains individual conditional lock outcomes. */ + public boolean hasLockResults() { + return lockResults != null; + } + + /** + * @param idx Request key index. + * @return Whether the key was locked. + */ + public boolean lockResult(int idx) { + return lockResults == null ? lockAcquired : lockResults[idx]; + } + /** * Empty constructor. */ @@ -129,14 +151,14 @@ public boolean compatibleRemapVersion() { } /** - * @return {@code True} if requested locks were acquired. + * @return Whether lock acquisition completed; consult {@link #lockResult(int)} for conditional outcomes. */ public boolean lockAcquired() { return lockAcquired; } /** - * @param lockAcquired {@code True} if requested locks were acquired. + * @param lockAcquired Whether lock acquisition completed, possibly with individual conditional outcomes. */ public void lockAcquired(boolean lockAcquired) { this.lockAcquired = lockAcquired; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearTxLocal.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearTxLocal.java index ad3dbc70824d4..8b3834a03a01a 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearTxLocal.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearTxLocal.java @@ -3225,20 +3225,31 @@ private void removeEntryMappings(IgniteTxEntry entry) { } /** - * Removes transaction entries and releases their acquired transactional locks. + * Removes the initiating-side state of a conditional lock rejected and cleaned up by the primary. + * No distributed unlock is needed: the primary has not replicated or acknowledged this lock. * - * @param entries Entries to remove and unlock. + * @param entry Entry enlisted by the failed attempt. + * @param snapshot Previous read entry state, if the entry was already enlisted before this attempt. */ - public void removeAndUnlockTxEntries(Collection entries) { - if (F.isEmpty(entries)) - return; + public void removeFailedLockEntry(IgniteTxEntry entry, @Nullable IgniteTxEntry snapshot) { + removeEntryMappings(entry); - for (IgniteTxEntry entry : entries) { - txState().removeEntry(entry.txKey()); - removeEntryMappings(entry); + if (entry.context().isNear()) { + try { + GridCacheEntryEx cached = entry.cached(); + + if (cached.hasLockCandidate(xidVersion())) + cached.removeLock(xidVersion()); + } + catch (GridCacheEntryRemovedException ignored) { + // An obsolete near entry has no live candidate. + } } - unlockTxEntries(entries); + if (snapshot == null) + txState().removeEntry(entry.txKey()); + else + entry.restoreFrom(snapshot); } /** @@ -4278,6 +4289,9 @@ public IgniteInternalFuture lockAllAsync(GridCacheContext c if (timeout == -1) return new GridFinishedFuture<>(timeoutException()); + boolean versionedLock = keys.stream().anyMatch(key -> + entry(cacheCtx.txKey((KeyCacheObject)key)).versionedLockPending()); + IgniteInternalFuture fut = cacheCtx.colocated().lockAllAsyncInternal(keys, timeout, waitTimeout, @@ -4295,12 +4309,15 @@ public IgniteInternalFuture lockAllAsync(GridCacheContext c return new GridEmbeddedFuture<>( fut, - new PLC1(ret, false, !CU.isWaitTimeoutExpiresFirst(waitTimeout, timeout)) { + new PLC1(ret, false, !versionedLock && !CU.isWaitTimeoutExpiresFirst(waitTimeout, timeout)) { @Override protected GridCacheReturn postLock(GridCacheReturn ret) throws IgniteCheckedException { assert fut.error() == null : "Lock future completed with an error: " + fut.error(); boolean success = Boolean.TRUE.equals(fut.get()); + if (!success) + checkValid(); + ret.success(success); if (log.isDebugEnabled()) { diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/transactions/IgniteTxEntry.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/transactions/IgniteTxEntry.java index a2d36d76d8c73..9b3c5a1cd7e5d 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/transactions/IgniteTxEntry.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/transactions/IgniteTxEntry.java @@ -216,6 +216,37 @@ public class IgniteTxEntry implements GridPeerDeployAware, MarshallableMessage, @Order(12) GridCacheVersion serReadVer; + /** Expected data version for a conditional pessimistic lock; sent in the lock request, not in prepare. */ + private GridCacheVersion expectedLockVer; + + /** Per-key outcome of the current conditional lock request; not part of prepare messages. */ + private Boolean versionedLockRes; + + /** @return Per-key conditional lock outcome, or {@code null} if not completed. */ + @Nullable public Boolean versionedLockResult() { + return versionedLockRes; + } + + /** @param res Per-key conditional lock outcome. */ + public void versionedLockResult(@Nullable Boolean res) { + versionedLockRes = res; + } + + /** @return Whether the current operation still needs a conditional lock for this entry. */ + public boolean versionedLockPending() { + return expectedLockVer != null && versionedLockRes == null; + } + + /** @return Expected data version, or {@code null} for an unconditional lock. */ + @Nullable public GridCacheVersion expectedLockVersion() { + return expectedLockVer; + } + + /** @param ver Expected data version, or {@code null} for an unconditional lock. */ + public void expectedLockVersion(@Nullable GridCacheVersion ver) { + expectedLockVer = ver; + } + /** * Empty constructor. */ @@ -423,6 +454,8 @@ public IgniteTxEntry copy() { cp.flags = flags; cp.partUpdateCntr = partUpdateCntr; cp.serReadVer = serReadVer; + cp.expectedLockVer = expectedLockVer; + cp.versionedLockRes = versionedLockRes; return cp; } @@ -468,6 +501,8 @@ public void restoreFrom(IgniteTxEntry snapshot) { flags = snapshot.flags; partUpdateCntr = snapshot.partUpdateCntr; serReadVer = snapshot.serReadVer; + expectedLockVer = snapshot.expectedLockVer; + versionedLockRes = snapshot.versionedLockRes; } /** diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/rollingupgrade/feature/CoreFeatureRegistry.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/rollingupgrade/feature/CoreFeatureRegistry.java index 92b7eeb1a8c5b..9ebbae1672c61 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/rollingupgrade/feature/CoreFeatureRegistry.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/rollingupgrade/feature/CoreFeatureRegistry.java @@ -93,4 +93,7 @@ public class CoreFeatureRegistry { /** */ public static final IgniteFeature ROLLING_UPGRADE_FEATURE = new IgniteCoreFeature(0); + + /** Primary-side data version validation for transactional locks. */ + public static final IgniteFeature VERSIONED_TX_LOCK_FEATURE = new IgniteCoreFeature(1); } diff --git a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/CacheVersionedEntryTransactionalLockTest.java b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/CacheVersionedEntryTransactionalLockTest.java index 41acd7a472c33..d8dc7ae555e18 100644 --- a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/CacheVersionedEntryTransactionalLockTest.java +++ b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/CacheVersionedEntryTransactionalLockTest.java @@ -17,8 +17,10 @@ package org.apache.ignite.internal.processors.cache; +import java.util.ArrayList; import java.util.Collection; import java.util.List; +import java.util.Map; import org.apache.ignite.Ignite; import org.apache.ignite.IgniteCache; import org.apache.ignite.IgniteCheckedException; @@ -29,7 +31,14 @@ import org.apache.ignite.configuration.CacheConfiguration; import org.apache.ignite.configuration.IgniteConfiguration; import org.apache.ignite.configuration.NearCacheConfiguration; +import org.apache.ignite.internal.IgniteInternalFuture; +import org.apache.ignite.internal.TestRecordingCommunicationSpi; +import org.apache.ignite.internal.processors.cache.distributed.GridNearUnlockRequest; +import org.apache.ignite.internal.processors.cache.distributed.near.GridNearLockRequest; +import org.apache.ignite.internal.processors.cache.distributed.near.GridNearLockResponse; +import org.apache.ignite.internal.transactions.IgniteTxTimeoutCheckedException; import org.apache.ignite.internal.util.typedef.internal.U; +import org.apache.ignite.testframework.GridTestUtils; import org.apache.ignite.testframework.junits.common.GridCommonAbstractTest; import org.apache.ignite.transactions.Transaction; import org.apache.ignite.transactions.TransactionIsolation; @@ -139,6 +148,7 @@ public static Collection testData() { /** {@inheritDoc} */ @Override protected IgniteConfiguration getConfiguration(String igniteInstanceName) throws Exception { return super.getConfiguration(igniteInstanceName) + .setCommunicationSpi(new TestRecordingCommunicationSpi()) .setConsistentId(igniteInstanceName); } @@ -356,6 +366,296 @@ public void testVersionedEntryLockCanBeRetriedAfterWaitTimeout() throws Exceptio assertEquals(1, cache.get(key).intValue()); } + /** A batch retains successful locks and reports changed, deleted and contended entries individually. */ + @Test + public void testPerEntryResults() throws Exception { + transactionalCache(ignite0); + + Ignite initiator = batch ? client : ignite0; + IgniteCache cache = initiator.cache(DEFAULT_CACHE_NAME); + List keys = primaryKeys(ignite0.cache(DEFAULT_CACHE_NAME), 5); + + for (int key : keys) + cache.put(key, 0); + + List> entries = List.of( + cache.getEntry(keys.get(0)), cache.getEntry(keys.get(1)), cache.getEntry(keys.get(2)), + cache.getEntry(keys.get(3)), cache.getEntry(keys.get(4))); + + cache.put(keys.get(1), 1); + cache.remove(keys.get(2)); + + IgniteCache holderCache = ignite1.cache(DEFAULT_CACHE_NAME); + TestRecordingCommunicationSpi spi = TestRecordingCommunicationSpi.spi(initiator); + + try (Transaction holder = ignite1.transactions().txStart(PESSIMISTIC, READ_COMMITTED)) { + holderCache.put(keys.get(3), 1); + spi.record(GridNearUnlockRequest.class); + + try (Transaction tx = initiator.transactions().txStart(PESSIMISTIC, READ_COMMITTED)) { + Map, Boolean> res = batch + ? internalCache(cache).lockTxEntriesAsync(entries, -1).get() + : internalCache(cache).lockTxEntries(entries, -1); + + assertEquals(5, res.size()); + assertEquals(Boolean.TRUE, res.get(entries.get(0))); + assertEquals(Boolean.FALSE, res.get(entries.get(1))); + assertEquals(Boolean.FALSE, res.get(entries.get(2))); + assertEquals(Boolean.FALSE, res.get(entries.get(3))); + assertEquals(Boolean.TRUE, res.get(entries.get(4))); + assertTrue("Rejected primary locks must not require an initiating-side unlock", + spi.recordedMessages(false).isEmpty()); + + IgniteCache other = grid(2).cache(DEFAULT_CACHE_NAME); + + try (Transaction competitor = grid(2).transactions().txStart(PESSIMISTIC, READ_COMMITTED)) { + assertFalse(acquireLockForEntry(other, entries.get(0), -1)); + assertFalse(acquireLockForEntry(other, entries.get(4), -1)); + assertTrue(acquireLockForEntry(other, other.getEntry(keys.get(1)), -1)); + } + + CacheEntry current = cache.getEntry(keys.get(1)); + Map, Boolean> repeated = internalCache(cache) + .lockTxEntries(List.of(current, entries.get(1)), -1); + + assertEquals(Boolean.TRUE, repeated.get(current)); + assertEquals(Boolean.FALSE, repeated.get(entries.get(1))); + + // Failed entries must not invalidate the transaction or release successful locks. + cache.put(keys.get(0), 2); + cache.put(keys.get(4), 2); + tx.commit(); + } + finally { + spi.recordedMessages(true); + } + } + + assertEquals(Integer.valueOf(2), cache.get(keys.get(0))); + assertEquals(Integer.valueOf(2), cache.get(keys.get(4))); + } + + /** Each primary receives one batch with individual outcomes, including partial failure and wait expiry. */ + @Test + public void testRequestsAreBatchedByPrimary() throws Exception { + transactionalCache(ignite0); + + IgniteCache cache = client.cache(DEFAULT_CACHE_NAME); + List> entries = new ArrayList<>(); + List contended = new ArrayList<>(); + + for (int node = 0; node < 3; node++) { + List keys = primaryKeys(grid(node).cache(DEFAULT_CACHE_NAME), 3); + + for (Integer key : keys) { + cache.put(key, 0); + entries.add(cache.getEntry(key)); + } + + cache.put(keys.get(1), 1); + + // The first primary rejects its entire batch; the remaining primaries must still be contacted. + if (node == 0) + cache.put(keys.get(0), 1); + + contended.add(keys.get(2)); + } + + TestRecordingCommunicationSpi spi = TestRecordingCommunicationSpi.spi(client); + + try (Transaction holder = ignite0.transactions().txStart(PESSIMISTIC, READ_COMMITTED)) { + IgniteCache holderCache = ignite0.cache(DEFAULT_CACHE_NAME); + + for (Integer key : contended) + holderCache.put(key, 1); + + spi.record(GridNearLockRequest.class); + + try (Transaction tx = client.transactions().txStart(PESSIMISTIC, READ_COMMITTED)) { + Map, Boolean> results = internalCache(cache) + .lockTxEntries(entries, batch ? 200 : -1); + + assertEquals(entries.size(), results.size()); + + for (int i = 0; i < entries.size(); i++) + assertEquals(i >= 3 && i % 3 == 0, results.get(entries.get(i)).booleanValue()); + + List requests = spi.recordedMessages(true); + + assertEquals("One batch per primary, rather than one request per key", 3, requests.size()); + + for (Object msg : requests) { + GridNearLockRequest req = (GridNearLockRequest)msg; + + assertEquals(3, req.keys().size()); + assertEquals(3, req.expectedVersions().length); + } + + for (Map.Entry, Boolean> result : results.entrySet()) { + if (result.getValue()) + cache.put(result.getKey().getKey(), 2); + } + + tx.commit(); + } + finally { + spi.recordedMessages(true); + } + } + } + + /** A version changed by the preceding owner must be rejected after waiting, with the transaction still usable. */ + @Test + public void testVersionChangedWhileWaiting() throws Exception { + checkVersionAfterWaiting(true); + } + + /** Waiting for a rollback succeeds because the protected data version has not changed. */ + @Test + public void testVersionUnchangedAfterWaiting() throws Exception { + checkVersionAfterWaiting(false); + } + + /** A transaction timeout must remain an error, rather than a normal per-entry rejection. */ + @Test + public void testTransactionTimeoutIsNotRejection() throws Exception { + IgniteCache cache = transactionalCache(ignite0); + int key = primaryKey(cache); + + cache.put(key, 0); + + IgniteCache remote = client.cache(DEFAULT_CACHE_NAME); + CacheEntry entry = remote.getEntry(key); + + try (Transaction holder = ignite0.transactions().txStart(PESSIMISTIC, READ_COMMITTED)) { + cache.put(key, 1); + + try (Transaction tx = client.transactions().txStart(PESSIMISTIC, READ_COMMITTED, 200, 0)) { + GridTestUtils.assertThrows(log, () -> internalCache(remote).lockTxEntries(List.of(entry), 0), + IgniteTxTimeoutCheckedException.class, null); + } + } + } + + /** + * @param commit Whether the preceding owner commits its update. + * @throws Exception If failed. + */ + private void checkVersionAfterWaiting(boolean commit) throws Exception { + IgniteCache cache = transactionalCache(ignite0); + int key = primaryKey(cache); + + cache.put(key, 0); + + CacheEntry entry = client.cache(DEFAULT_CACHE_NAME).getEntry(key); + IgniteInternalFuture waiter; + + try (Transaction holder = ignite0.transactions().txStart(PESSIMISTIC, READ_COMMITTED)) { + cache.put(key, 1); + + waiter = GridTestUtils.runAsync(() -> { + IgniteCache remote = client.cache(DEFAULT_CACHE_NAME); + + try (Transaction tx = client.transactions().txStart(PESSIMISTIC, READ_COMMITTED, 10_000, 0)) { + // Exercise both an explicit wait limit and the transaction's remaining timeout. + long wait = batch ? 5_000 : 0; + + assertEquals(!commit, acquireLockForEntry(remote, entry, wait)); + + if (commit) + assertTrue(acquireLockForEntry(remote, remote.getEntry(key), -1)); + + remote.put(key, 2); + tx.commit(); + } + + return null; + }); + + awaitBlockedLock(cache, key); + + if (commit) + holder.commit(); + else + holder.rollback(); + } + + waiter.get(10_000); + + assertEquals(Integer.valueOf(2), cache.get(key)); + } + + /** A successful primary check protects the version even while its reply has not reached the initiator. */ + @Test + public void testVersionProtectedBeforeReply() throws Exception { + IgniteCache cache = transactionalCache(ignite0); + int key = primaryKey(cache); + + cache.put(key, 0); + + CacheEntry entry = client.cache(DEFAULT_CACHE_NAME).getEntry(key); + TestRecordingCommunicationSpi spi = TestRecordingCommunicationSpi.spi(ignite0); + + spi.blockMessages((node, msg) -> node.id().equals(client.cluster().localNode().id()) + && msg instanceof GridNearLockResponse && ((GridNearLockResponse)msg).lockAcquired()); + + IgniteInternalFuture locker = GridTestUtils.runAsync(() -> { + IgniteCache remote = client.cache(DEFAULT_CACHE_NAME); + + try (Transaction tx = client.transactions().txStart(PESSIMISTIC, READ_COMMITTED, 10_000, 0)) { + assertTrue(acquireLockForEntry(remote, entry, -1)); + tx.commit(); + } + + return null; + }); + + IgniteInternalFuture writer = null; + + try { + assertTrue(spi.waitForBlocked(1, 5_000)); + + IgniteCache competitorCache = ignite1.cache(DEFAULT_CACHE_NAME); + + writer = GridTestUtils.runAsync(() -> competitorCache.put(key, 1)); + + awaitBlockedLock(cache, key); + + assertFalse(writer.isDone()); + assertEquals(entry.version(), cache.getEntry(key).version()); + } + finally { + spi.stopBlock(); + } + + locker.get(10_000); + assertNotNull(writer); + writer.get(10_000); + + assertEquals(Integer.valueOf(1), cache.get(key)); + assertFalse(entry.version().equals(cache.getEntry(key).version())); + } + + /** Waits for a second transaction to enqueue on the primary, without relying on a sleep. */ + private void awaitBlockedLock(IgniteCache cache, int key) throws Exception { + GridCacheContext ctx = internalCache(cache).context(); + + if (ctx.isNear()) + ctx = ctx.near().dht().context(); + + GridCacheEntryEx primaryEntry = ctx.dht().entryEx(ctx.toCacheKeyObject(key)); + + assertTrue("The competing request must actually wait on the primary", + GridTestUtils.waitForCondition(() -> { + try { + return primaryEntry.localCandidates().size() >= 2; + } + catch (GridCacheEntryRemovedException e) { + return false; + } + }, 5_000)); + } + /** * Checks that the lock entry method returns {@code false} when can not wait for lock. * @@ -509,6 +809,6 @@ private static boolean acquireLockForEntries( List> entries, long timeout ) throws IgniteCheckedException { - return cache.unwrap(IgniteCacheProxy.class).internalProxy().lockTxEntries(entries, timeout); + return !cache.unwrap(IgniteCacheProxy.class).internalProxy().lockTxEntries(entries, timeout).containsValue(false); } } From 1266956375f49557dc57c3d2d485903bcb93cd75 Mon Sep 17 00:00:00 2001 From: Vladislav Pyatkov Date: Fri, 9 Oct 2026 23:55:14 +0300 Subject: [PATCH 2/2] The primary node enforces the lock wait timeout. The coordinator only enforces the transaction timeout. --- .../colocated/GridDhtColocatedLockFuture.java | 77 +++++++------------ .../distributed/near/GridNearLockFuture.java | 68 ++++++---------- ...heVersionedEntryTransactionalLockTest.java | 7 -- 3 files changed, 51 insertions(+), 101 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/colocated/GridDhtColocatedLockFuture.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/colocated/GridDhtColocatedLockFuture.java index 9f4f9efa93d8c..871a42f73d33b 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/colocated/GridDhtColocatedLockFuture.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/colocated/GridDhtColocatedLockFuture.java @@ -140,12 +140,12 @@ public final class GridDhtColocatedLockFuture extends GridCacheCompoundIdentityF /** Lock wait timeout. */ private final long waitTimeout; + /** Time when the first primary batch starts, or zero before mapping. */ + private long lockWaitStartTime; + /** Whether this operation conditionally locks a data version. */ private final boolean versionedLock; - /** Shared wait deadline for conditional primary batches. */ - private final long lockWaitEndTime; - /** Transaction. */ @GridToStringExclude private final GridNearTxLocal tx; @@ -233,9 +233,6 @@ public GridDhtColocatedLockFuture( this.retval = retval; this.timeout = timeout; this.waitTimeout = waitTimeout; - versionedLock = tx != null && keys.stream().anyMatch(key -> - tx.entry(cctx.txKey(key)) != null && tx.entry(cctx.txKey(key)).versionedLockPending()); - lockWaitEndTime = waitTimeout > 0 ? U.currentTimeMillis() + waitTimeout : 0; this.createTtl = createTtl; this.accessTtl = accessTtl; this.skipStore = skipStore; @@ -244,6 +241,9 @@ public GridDhtColocatedLockFuture( this.recovery = recovery; this.keepBinaryInInterceptor = keepBinaryInInterceptor; + versionedLock = tx != null && keys.stream().anyMatch(key -> + tx.entry(cctx.txKey(key)) != null && tx.entry(cctx.txKey(key)).versionedLockPending()); + ignoreInterrupts(); threadId = tx == null ? Thread.currentThread().getId() : tx.threadId(); @@ -786,7 +786,7 @@ void map() { if (isDone()) // Possible due to async rollback. return; - if (lockTimeout() > 0) { + if (timeout > 0) { timeoutObj = new LockTimeoutObject(); cctx.time().addTimeoutObject(timeoutObj); @@ -1229,12 +1229,24 @@ private void proceedMapping0() final Collection mappedKeys = map.distributedKeys(); final ClusterNode node = map.node(); - if (versionedLock) - req.waitTimeout(remainingWaitTimeout()); + long remainingWaitTimeout = waitTimeout; + + if (waitTimeout > 0) { + long now = U.currentTimeMillis(); + + if (lockWaitStartTime == 0) + lockWaitStartTime = now; + + long elapsed = now - lockWaitStartTime; + + remainingWaitTimeout = elapsed < waitTimeout ? waitTimeout - elapsed : -1; + } if (node.isLocal()) - lockLocally(mappedKeys, req.topologyVersion()); + lockLocally(mappedKeys, req.topologyVersion(), remainingWaitTimeout); else { + req.waitTimeout(remainingWaitTimeout); + final MiniFuture fut = new MiniFuture(node, mappedKeys, ++miniId); req.miniId(fut.futureId()); @@ -1262,10 +1274,12 @@ private void proceedMapping0() * Locks given keys directly through dht cache. * @param keys Collection of keys. * @param topVer Topology version to lock on. + * @param remainingWaitTimeout Remaining lock wait budget for this batch. */ private void lockLocally( final Collection keys, - AffinityTopologyVersion topVer + AffinityTopologyVersion topVer, + long remainingWaitTimeout ) { if (log.isDebugEnabled()) log.debug("Before locally locking keys : " + keys); @@ -1279,7 +1293,7 @@ private void lockLocally( read, retval, timeout, - remainingWaitTimeout(), + remainingWaitTimeout, createTtl, accessTtl, skipStore, @@ -1396,7 +1410,7 @@ private boolean mapAsPrimary(Collection keys, AffinityTopologyVe tx.addKeyMapping(cctx.txKey(key), cctx.localNode()); } - lockLocally(distributedKeys, topVer); + lockLocally(distributedKeys, topVer, waitTimeout); } GridDhtPartitionsExchangeFuture lastFinishedFut = cctx.shared().exchange().lastFinishedFuture(); @@ -1503,27 +1517,6 @@ private boolean errorOrTimeoutOnTopologyVersion(IgniteCheckedException e, boolea return false; } - /** - * @return Timeout value for this lock future. - */ - private long lockTimeout() { - // Wait for the primary's conditional rejection and cleanup; only the transaction deadline is local. - if (versionedLock) - return timeout; - - return CU.isWaitTimeoutExpiresFirst(waitTimeout, timeout) ? waitTimeout : timeout; - } - - /** @return Remaining shared wait budget for the next primary batch. */ - private long remainingWaitTimeout() { - if (!versionedLock || waitTimeout <= 0) - return waitTimeout; - - long remaining = lockWaitEndTime - U.currentTimeMillis(); - - return remaining > 0 ? remaining : -1; - } - /** * Lock request timeout object. */ @@ -1532,7 +1525,7 @@ private class LockTimeoutObject extends GridTimeoutObjectAdapter { * Default constructor. */ LockTimeoutObject() { - super(lockTimeout()); + super(timeout); } /** Requested keys. */ @@ -1543,20 +1536,6 @@ private class LockTimeoutObject extends GridTimeoutObjectAdapter { if (log.isDebugEnabled()) log.debug("Timed out waiting for lock response: " + this); - if (CU.isWaitTimeoutExpiresFirst(waitTimeout, timeout)) { - synchronized (GridDhtColocatedLockFuture.this) { - requestedKeys = requestedKeys0(); - - clear(); // Stop response processing. - } - - synchronized (this) { - onComplete(false, true, false); - } - - return; - } - if (inTx()) { if (cctx.tm().deadlockDetectionEnabled()) { synchronized (GridDhtColocatedLockFuture.this) { diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearLockFuture.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearLockFuture.java index edf00bd429d9e..5cbc9838cf284 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearLockFuture.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearLockFuture.java @@ -137,12 +137,12 @@ public final class GridNearLockFuture extends GridCacheCompoundIdentityFuture - tx.entry(cctx.txKey(key)) != null && tx.entry(cctx.txKey(key)).versionedLockPending()); - lockWaitEndTime = waitTimeout > 0 ? U.currentTimeMillis() + waitTimeout : 0; this.createTtl = createTtl; this.accessTtl = accessTtl; this.skipStore = skipStore; @@ -245,6 +242,9 @@ public GridNearLockFuture( ignoreInterrupts(); + versionedLock = tx != null && keys.stream().anyMatch(key -> + tx.entry(cctx.txKey(key)) != null && tx.entry(cctx.txKey(key)).versionedLockPending()); + threadId = tx == null ? Thread.currentThread().getId() : tx.threadId(); lockVer = tx != null ? tx.xidVersion() : cctx.cache().nextVersion(); @@ -365,7 +365,7 @@ private boolean locked(GridCacheEntryEx cached) throws GridCacheEntryRemovedExce threadId, lockVer, topVer, - lockTimeout(), + timeout, !inTx(), inTx(), implicitSingleTx(), @@ -380,16 +380,11 @@ private boolean locked(GridCacheEntryEx cached) throws GridCacheEntryRemovedExce entries.add(entry); - if (c == null && lockTimeout() < 0) { + if (c == null && timeout < 0) { if (log.isDebugEnabled()) log.debug("Failed to acquire lock with negative timeout: " + entry); - if (CU.isWaitTimeoutExpiresFirst(waitTimeout, timeout)) { - onComplete(false, true, false); - } - else { - onFailed(false); - } + onFailed(false); return null; } @@ -839,7 +834,7 @@ void map() { if (isDone()) // Possible due to async rollback. return; - if (lockTimeout() > 0) { + if (timeout > 0) { timeoutObj = new LockTimeoutObject(); cctx.time().addTimeoutObject(timeoutObj); @@ -1231,12 +1226,21 @@ private void proceedMapping0() final Collection mappedKeys = map.distributedKeys(); final ClusterNode node = map.node(); - if (versionedLock && waitTimeout > 0) { - long remaining = lockWaitEndTime - U.currentTimeMillis(); + long remainingWaitTimeout = waitTimeout; + + if (waitTimeout > 0) { + long now = U.currentTimeMillis(); - req.waitTimeout(remaining > 0 ? remaining : -1); + if (lockWaitStartTime == 0) + lockWaitStartTime = now; + + long elapsed = now - lockWaitStartTime; + + remainingWaitTimeout = elapsed < waitTimeout ? waitTimeout - elapsed : -1; } + req.waitTimeout(remainingWaitTimeout); + if (node.isLocal()) { req.miniId(-1); @@ -1459,18 +1463,6 @@ private ClusterTopologyCheckedException newTopologyException(@Nullable Throwable return topEx; } - /** - * @return Timeout value for this lock future. - */ - private long lockTimeout() { - // For conditional locks the primary owns the wait deadline and acknowledges rejection only after - // cleaning up its candidate. An earlier near-side rejection would race with a retry of the same key. - if (versionedLock) - return timeout; - - return CU.isWaitTimeoutExpiresFirst(waitTimeout, timeout) ? waitTimeout : timeout; - } - /** Applies one primary outcome, removing a rejected local candidate before proceeding to the next batch. */ private boolean processLockResult(KeyCacheObject key, boolean locked) { tx.entry(cctx.txKey(key)).versionedLockResult(locked); @@ -1498,7 +1490,7 @@ private class LockTimeoutObject extends GridTimeoutObjectAdapter { * Default constructor. */ LockTimeoutObject() { - super(lockTimeout()); + super(timeout); } /** Requested keys. */ @@ -1511,20 +1503,6 @@ private class LockTimeoutObject extends GridTimeoutObjectAdapter { timedOut = true; - if (CU.isWaitTimeoutExpiresFirst(waitTimeout, timeout)) { - synchronized (GridNearLockFuture.this) { - requestedKeys = requestedKeys0(); - - clear(); // Stop response processing. - } - - synchronized (this) { - onComplete(false, true, false); - } - - return; - } - txLockTimedOut = true; if (inTx()) { diff --git a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/CacheVersionedEntryTransactionalLockTest.java b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/CacheVersionedEntryTransactionalLockTest.java index d8dc7ae555e18..88b67bcaba990 100644 --- a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/CacheVersionedEntryTransactionalLockTest.java +++ b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/CacheVersionedEntryTransactionalLockTest.java @@ -414,13 +414,6 @@ public void testPerEntryResults() throws Exception { assertTrue(acquireLockForEntry(other, other.getEntry(keys.get(1)), -1)); } - CacheEntry current = cache.getEntry(keys.get(1)); - Map, Boolean> repeated = internalCache(cache) - .lockTxEntries(List.of(current, entries.get(1)), -1); - - assertEquals(Boolean.TRUE, repeated.get(current)); - assertEquals(Boolean.FALSE, repeated.get(entries.get(1))); - // Failed entries must not invalidate the transaction or release successful locks. cache.put(keys.get(0), 2); cache.put(keys.get(4), 2);