diff --git a/core/src/main/java/org/apache/accumulo/core/fate/user/FateMutator.java b/core/src/main/java/org/apache/accumulo/core/fate/user/FateMutator.java index f8a65011473..22cb97ca649 100644 --- a/core/src/main/java/org/apache/accumulo/core/fate/user/FateMutator.java +++ b/core/src/main/java/org/apache/accumulo/core/fate/user/FateMutator.java @@ -52,6 +52,11 @@ public interface FateMutator { */ FateMutator requireUnreserved(); + /** + * Require the transaction has a reservation. + */ + FateMutator requireReserved(FateStore.FateReservation fateReservation); + /** * Require the transaction has no fate key set. */ diff --git a/core/src/main/java/org/apache/accumulo/core/fate/user/FateMutatorImpl.java b/core/src/main/java/org/apache/accumulo/core/fate/user/FateMutatorImpl.java index 26971e2c479..34dd4c1b670 100644 --- a/core/src/main/java/org/apache/accumulo/core/fate/user/FateMutatorImpl.java +++ b/core/src/main/java/org/apache/accumulo/core/fate/user/FateMutatorImpl.java @@ -61,6 +61,7 @@ public class FateMutatorImpl implements FateMutator { private final ConditionalMutation mutation; private final Supplier writer; private boolean requiredUnreserved = false; + private boolean requiredReserved = false; public static final int INITIAL_ITERATOR_PRIO = 1000000; public FateMutatorImpl(ClientContext context, String tableName, FateId fateId, @@ -108,6 +109,17 @@ public FateMutator requireUnreserved() { return this; } + @Override + public FateMutator requireReserved(FateStore.FateReservation fateReservation) { + Preconditions.checkState(!requiredReserved); + Condition condition = new Condition(TxAdminColumnFamily.RESERVATION_COLUMN.getColumnFamily(), + TxAdminColumnFamily.RESERVATION_COLUMN.getColumnQualifier()) + .setValue(fateReservation.getSerialized()); + mutation.addCondition(condition); + requiredReserved = true; + return this; + } + @Override public FateMutator requireAbsentKey() { Condition condition = new Condition(TxColumnFamily.TX_KEY_COLUMN.getColumnFamily(), diff --git a/core/src/main/java/org/apache/accumulo/core/fate/user/UserFateStore.java b/core/src/main/java/org/apache/accumulo/core/fate/user/UserFateStore.java index 158c2ce1aa2..aea79faf625 100644 --- a/core/src/main/java/org/apache/accumulo/core/fate/user/UserFateStore.java +++ b/core/src/main/java/org/apache/accumulo/core/fate/user/UserFateStore.java @@ -414,6 +414,12 @@ private FateMutatorImpl newMutator(FateId fateId) { return new FateMutatorImpl<>(context, tableName, fateId, writer); } + public FateMutatorImpl newReservedMutator(FateId fateId, FateReservation reservation) { + Preconditions.checkState(fateId != null, + "Attempted write on deleted FATE transaction: " + fateId); + return (FateMutatorImpl) newMutator(fateId).requireReserved(reservation); + } + private R scanTx(Function func) { try (Scanner scanner = context.createScanner(tableName, Authorizations.EMPTY)) { return func.apply(scanner); @@ -546,8 +552,6 @@ private FateTxStoreImpl(FateId fateId, FateReservation reservation) { @Override public Repo top() { - verifyReservedAndNotDeleted(false); - return scanTx(scanner -> { scanner.setRange(getRow(fateId)); scanner.setBatchSize(1); @@ -562,8 +566,6 @@ public Repo top() { @Override public List> getStack() { - verifyReservedAndNotDeleted(false); - return scanTx(scanner -> { scanner.setRange(getRow(fateId)); scanner.fetchColumnFamily(RepoColumnFamily.NAME); @@ -577,8 +579,6 @@ public List> getStack() { @Override public Serializable getTransactionInfo(TxInfo txInfo) { - verifyReservedAndNotDeleted(false); - try (Scanner scanner = context.createScanner(tableName, Authorizations.EMPTY)) { scanner.setRange(getRow(fateId)); @@ -600,8 +600,6 @@ public Serializable getTransactionInfo(TxInfo txInfo) { @Override public long timeCreated() { - verifyReservedAndNotDeleted(false); - return scanTx(scanner -> { scanner.setRange(getRow(fateId)); TxColumnFamily.CREATE_TIME_COLUMN.fetch(scanner); @@ -612,50 +610,39 @@ public long timeCreated() { @Override public void push(Repo repo) throws StackOverflowException { - verifyReservedAndNotDeleted(true); - Optional top = findTop(); if (top.filter(t -> t >= MAX_REPOS).isPresent()) { throw new StackOverflowException("Repo stack size too large"); } - FateMutator fateMutator = - newMutator(fateId).requireStatus(REQ_PUSH_STATUS.toArray(TStatus[]::new)); + FateMutator fateMutator = newReservedMutator(fateId, reservation) + .requireStatus(REQ_PUSH_STATUS.toArray(TStatus[]::new)); fateMutator.putRepo(top.map(t -> t + 1).orElse(1), repo).mutate(); } @Override public void pop() { - verifyReservedAndNotDeleted(true); - Optional top = findTop(); - top.ifPresent(t -> newMutator(fateId).requireStatus(REQ_POP_STATUS.toArray(TStatus[]::new)) - .deleteRepo(t).mutate()); + top.ifPresent(t -> newReservedMutator(fateId, reservation) + .requireStatus(REQ_POP_STATUS.toArray(TStatus[]::new)).deleteRepo(t).mutate()); } @Override public void setStatus(TStatus status) { - verifyReservedAndNotDeleted(true); - - newMutator(fateId).putStatus(status).mutate(); + newReservedMutator(fateId, reservation).putStatus(status).mutate(); observedStatus = status; } @Override public void setTransactionInfo(TxInfo txInfo, Serializable so) { - verifyReservedAndNotDeleted(true); - final byte[] serialized = serializeTxInfo(so); - - newMutator(fateId).putTxInfo(txInfo, serialized).mutate(); + newReservedMutator(fateId, reservation).putTxInfo(txInfo, serialized).mutate(); } @Override public void delete() { - verifyReservedAndNotDeleted(true); - - var mutator = newMutator(fateId); + var mutator = newReservedMutator(fateId, reservation); mutator.requireStatus(REQ_DELETE_STATUS.toArray(TStatus[]::new)); mutator.delete().mutate(); this.deleted = true; @@ -663,9 +650,7 @@ public void delete() { @Override public void forceDelete() { - verifyReservedAndNotDeleted(true); - - var mutator = newMutator(fateId); + var mutator = newReservedMutator(fateId, reservation); mutator.requireStatus(REQ_FORCE_DELETE_STATUS.toArray(TStatus[]::new)); mutator.delete().mutate(); this.deleted = true;