From 1d70dd85dcefc606be735a24aff09b0a48fd268c Mon Sep 17 00:00:00 2001 From: avillarreal Date: Mon, 31 Aug 2026 11:03:49 -0500 Subject: [PATCH 1/6] Add requireReserved() to FateMutator --- .../accumulo/core/fate/user/FateMutator.java | 5 +++++ .../accumulo/core/fate/user/FateMutatorImpl.java | 16 ++++++++++++---- 2 files changed, 17 insertions(+), 4 deletions(-) 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..1803eae6d0d 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(); + /** * 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..8d4c4b57992 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,16 @@ public FateMutator requireUnreserved() { return this; } + @Override + public FateMutator requireReserved() { + Preconditions.checkState(!requiredReserved); + Condition condition = new Condition(TxAdminColumnFamily.RESERVATION_COLUMN.getColumnFamily(), + TxAdminColumnFamily.RESERVATION_COLUMN.getColumnQualifier()); + mutation.addCondition(condition); + requiredReserved = true; + return this; + } + @Override public FateMutator requireAbsentKey() { Condition condition = new Condition(TxColumnFamily.TX_KEY_COLUMN.getColumnFamily(), @@ -125,10 +136,7 @@ public FateMutator putReservedTx(FateStore.FateReservation reservation) { @Override public FateMutator putUnreserveTx(FateStore.FateReservation reservation) { - Condition condition = new Condition(TxAdminColumnFamily.RESERVATION_COLUMN.getColumnFamily(), - TxAdminColumnFamily.RESERVATION_COLUMN.getColumnQualifier()) - .setValue(reservation.getSerialized()); - mutation.addCondition(condition); + requireReserved(); TxAdminColumnFamily.RESERVATION_COLUMN.putDelete(mutation); return this; } From 21414b06e8c76adeea22d520ba9698a09cd395bf Mon Sep 17 00:00:00 2001 From: avillarreal Date: Mon, 31 Aug 2026 11:40:42 -0500 Subject: [PATCH 2/6] Remove verifyReserved() from UserFateStore.java --- .../core/fate/user/FateMutatorImpl.java | 2 +- .../core/fate/user/UserFateStore.java | 20 ------------------- 2 files changed, 1 insertion(+), 21 deletions(-) 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 8d4c4b57992..bb03be1af41 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 @@ -113,7 +113,7 @@ public FateMutator requireUnreserved() { public FateMutator requireReserved() { Preconditions.checkState(!requiredReserved); Condition condition = new Condition(TxAdminColumnFamily.RESERVATION_COLUMN.getColumnFamily(), - TxAdminColumnFamily.RESERVATION_COLUMN.getColumnQualifier()); + TxAdminColumnFamily.RESERVATION_COLUMN.getColumnQualifier()); mutation.addCondition(condition); requiredReserved = true; return this; 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..9802e3afb35 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 @@ -546,8 +546,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 +560,6 @@ public Repo top() { @Override public List> getStack() { - verifyReservedAndNotDeleted(false); - return scanTx(scanner -> { scanner.setRange(getRow(fateId)); scanner.fetchColumnFamily(RepoColumnFamily.NAME); @@ -577,8 +573,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 +594,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,8 +604,6 @@ public long timeCreated() { @Override public void push(Repo repo) throws StackOverflowException { - verifyReservedAndNotDeleted(true); - Optional top = findTop(); if (top.filter(t -> t >= MAX_REPOS).isPresent()) { @@ -627,8 +617,6 @@ public void push(Repo repo) throws StackOverflowException { @Override public void pop() { - verifyReservedAndNotDeleted(true); - Optional top = findTop(); top.ifPresent(t -> newMutator(fateId).requireStatus(REQ_POP_STATUS.toArray(TStatus[]::new)) .deleteRepo(t).mutate()); @@ -636,16 +624,12 @@ public void pop() { @Override public void setStatus(TStatus status) { - verifyReservedAndNotDeleted(true); - newMutator(fateId).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(); @@ -653,8 +637,6 @@ public void setTransactionInfo(TxInfo txInfo, Serializable so) { @Override public void delete() { - verifyReservedAndNotDeleted(true); - var mutator = newMutator(fateId); mutator.requireStatus(REQ_DELETE_STATUS.toArray(TStatus[]::new)); mutator.delete().mutate(); @@ -663,8 +645,6 @@ public void delete() { @Override public void forceDelete() { - verifyReservedAndNotDeleted(true); - var mutator = newMutator(fateId); mutator.requireStatus(REQ_FORCE_DELETE_STATUS.toArray(TStatus[]::new)); mutator.delete().mutate(); From 5d446fcf677fc1605dad556989f2fd157aea83fc Mon Sep 17 00:00:00 2001 From: avillarreal Date: Mon, 31 Aug 2026 15:07:53 -0500 Subject: [PATCH 3/6] Add requireReserved() into some FateMutatorImpl methods used by UserFateStore --- .../org/apache/accumulo/core/fate/user/FateMutatorImpl.java | 5 +++++ 1 file changed, 5 insertions(+) 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 bb03be1af41..f0891e701ec 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 @@ -75,6 +75,7 @@ public FateMutatorImpl(ClientContext context, String tableName, FateId fateId, @Override public FateMutator putStatus(TStatus status) { + requireReserved(); TxAdminColumnFamily.STATUS_COLUMN.put(mutation, new Value(status.name())); return this; } @@ -173,6 +174,7 @@ public FateMutator putAgeOff(byte[] data) { @Override public FateMutator putTxInfo(TxInfo txInfo, byte[] data) { + requireReserved(); switch (txInfo) { case FATE_OP: putFateOp(data); @@ -197,6 +199,7 @@ public FateMutator putTxInfo(TxInfo txInfo, byte[] data) { @Override public FateMutator putRepo(int position, Repo repo) { + requireReserved(); final Text cq = invertRepo(position); // ensure this repo is not already set mutation.addCondition(new Condition(RepoColumnFamily.NAME, cq)); @@ -206,11 +209,13 @@ public FateMutator putRepo(int position, Repo repo) { @Override public FateMutator deleteRepo(int position) { + requireReserved(); mutation.putDelete(RepoColumnFamily.NAME, invertRepo(position)); return this; } public FateMutator delete() { + requireReserved(); try (Scanner scanner = context.createScanner(tableName, Authorizations.EMPTY)) { scanner.setRange(getRow(fateId)); scanner.forEach( From 0a8c9def7fe525569329cc71810b59c5d1c3bb98 Mon Sep 17 00:00:00 2001 From: avillarreal Date: Thu, 3 Sep 2026 11:58:26 -0500 Subject: [PATCH 4/6] Update requireReserved with FateReservation param, move usages to UserFateStore --- .../accumulo/core/fate/user/FateMutator.java | 2 +- .../core/fate/user/FateMutatorImpl.java | 15 ++++++------ .../core/fate/user/UserFateStore.java | 24 ++++++++++++------- 3 files changed, 24 insertions(+), 17 deletions(-) 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 1803eae6d0d..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 @@ -55,7 +55,7 @@ public interface FateMutator { /** * Require the transaction has a reservation. */ - FateMutator requireReserved(); + 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 f0891e701ec..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 @@ -75,7 +75,6 @@ public FateMutatorImpl(ClientContext context, String tableName, FateId fateId, @Override public FateMutator putStatus(TStatus status) { - requireReserved(); TxAdminColumnFamily.STATUS_COLUMN.put(mutation, new Value(status.name())); return this; } @@ -111,10 +110,11 @@ public FateMutator requireUnreserved() { } @Override - public FateMutator requireReserved() { + public FateMutator requireReserved(FateStore.FateReservation fateReservation) { Preconditions.checkState(!requiredReserved); Condition condition = new Condition(TxAdminColumnFamily.RESERVATION_COLUMN.getColumnFamily(), - TxAdminColumnFamily.RESERVATION_COLUMN.getColumnQualifier()); + TxAdminColumnFamily.RESERVATION_COLUMN.getColumnQualifier()) + .setValue(fateReservation.getSerialized()); mutation.addCondition(condition); requiredReserved = true; return this; @@ -137,7 +137,10 @@ public FateMutator putReservedTx(FateStore.FateReservation reservation) { @Override public FateMutator putUnreserveTx(FateStore.FateReservation reservation) { - requireReserved(); + Condition condition = new Condition(TxAdminColumnFamily.RESERVATION_COLUMN.getColumnFamily(), + TxAdminColumnFamily.RESERVATION_COLUMN.getColumnQualifier()) + .setValue(reservation.getSerialized()); + mutation.addCondition(condition); TxAdminColumnFamily.RESERVATION_COLUMN.putDelete(mutation); return this; } @@ -174,7 +177,6 @@ public FateMutator putAgeOff(byte[] data) { @Override public FateMutator putTxInfo(TxInfo txInfo, byte[] data) { - requireReserved(); switch (txInfo) { case FATE_OP: putFateOp(data); @@ -199,7 +201,6 @@ public FateMutator putTxInfo(TxInfo txInfo, byte[] data) { @Override public FateMutator putRepo(int position, Repo repo) { - requireReserved(); final Text cq = invertRepo(position); // ensure this repo is not already set mutation.addCondition(new Condition(RepoColumnFamily.NAME, cq)); @@ -209,13 +210,11 @@ public FateMutator putRepo(int position, Repo repo) { @Override public FateMutator deleteRepo(int position) { - requireReserved(); mutation.putDelete(RepoColumnFamily.NAME, invertRepo(position)); return this; } public FateMutator delete() { - requireReserved(); try (Scanner scanner = context.createScanner(tableName, Authorizations.EMPTY)) { scanner.setRange(getRow(fateId)); scanner.forEach( 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 9802e3afb35..b25b999e6aa 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 @@ -546,6 +546,8 @@ private FateTxStoreImpl(FateId fateId, FateReservation reservation) { @Override public Repo top() { + verifyReservedAndNotDeleted(false); + return scanTx(scanner -> { scanner.setRange(getRow(fateId)); scanner.setBatchSize(1); @@ -560,6 +562,8 @@ public Repo top() { @Override public List> getStack() { + verifyReservedAndNotDeleted(false); + return scanTx(scanner -> { scanner.setRange(getRow(fateId)); scanner.fetchColumnFamily(RepoColumnFamily.NAME); @@ -573,6 +577,8 @@ public List> getStack() { @Override public Serializable getTransactionInfo(TxInfo txInfo) { + verifyReservedAndNotDeleted(false); + try (Scanner scanner = context.createScanner(tableName, Authorizations.EMPTY)) { scanner.setRange(getRow(fateId)); @@ -594,6 +600,8 @@ public Serializable getTransactionInfo(TxInfo txInfo) { @Override public long timeCreated() { + verifyReservedAndNotDeleted(false); + return scanTx(scanner -> { scanner.setRange(getRow(fateId)); TxColumnFamily.CREATE_TIME_COLUMN.fetch(scanner); @@ -610,8 +618,8 @@ public void push(Repo repo) throws StackOverflowException { throw new StackOverflowException("Repo stack size too large"); } - FateMutator fateMutator = - newMutator(fateId).requireStatus(REQ_PUSH_STATUS.toArray(TStatus[]::new)); + FateMutator fateMutator = newMutator(fateId) + .requireStatus(REQ_PUSH_STATUS.toArray(TStatus[]::new)).requireReserved(reservation); fateMutator.putRepo(top.map(t -> t + 1).orElse(1), repo).mutate(); } @@ -619,26 +627,25 @@ public void push(Repo repo) throws StackOverflowException { public void pop() { Optional top = findTop(); top.ifPresent(t -> newMutator(fateId).requireStatus(REQ_POP_STATUS.toArray(TStatus[]::new)) - .deleteRepo(t).mutate()); + .requireReserved(reservation).deleteRepo(t).mutate()); } @Override public void setStatus(TStatus status) { - newMutator(fateId).putStatus(status).mutate(); + newMutator(fateId).requireReserved(reservation).putStatus(status).mutate(); observedStatus = status; } @Override public void setTransactionInfo(TxInfo txInfo, Serializable so) { final byte[] serialized = serializeTxInfo(so); - - newMutator(fateId).putTxInfo(txInfo, serialized).mutate(); + newMutator(fateId).requireReserved(reservation).putTxInfo(txInfo, serialized).mutate(); } @Override public void delete() { var mutator = newMutator(fateId); - mutator.requireStatus(REQ_DELETE_STATUS.toArray(TStatus[]::new)); + mutator.requireStatus(REQ_DELETE_STATUS.toArray(TStatus[]::new)).requireReserved(reservation); mutator.delete().mutate(); this.deleted = true; } @@ -646,7 +653,8 @@ public void delete() { @Override public void forceDelete() { var mutator = newMutator(fateId); - mutator.requireStatus(REQ_FORCE_DELETE_STATUS.toArray(TStatus[]::new)); + mutator.requireStatus(REQ_FORCE_DELETE_STATUS.toArray(TStatus[]::new)) + .requireReserved(reservation); mutator.delete().mutate(); this.deleted = true; } From 82b55037df3625b703d6809b535eef6ad29ece16 Mon Sep 17 00:00:00 2001 From: avillarreal Date: Thu, 3 Sep 2026 12:49:19 -0500 Subject: [PATCH 5/6] Replace verifyReserved() with requireReserved() to 4 more methods in UserFateStore --- .../org/apache/accumulo/core/fate/user/UserFateStore.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) 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 b25b999e6aa..f14bcf37cfe 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 @@ -546,7 +546,7 @@ private FateTxStoreImpl(FateId fateId, FateReservation reservation) { @Override public Repo top() { - verifyReservedAndNotDeleted(false); + newMutator(fateId).requireReserved(reservation); return scanTx(scanner -> { scanner.setRange(getRow(fateId)); @@ -562,7 +562,7 @@ public Repo top() { @Override public List> getStack() { - verifyReservedAndNotDeleted(false); + newMutator(fateId).requireReserved(reservation); return scanTx(scanner -> { scanner.setRange(getRow(fateId)); @@ -577,7 +577,7 @@ public List> getStack() { @Override public Serializable getTransactionInfo(TxInfo txInfo) { - verifyReservedAndNotDeleted(false); + newMutator(fateId).requireReserved(reservation); try (Scanner scanner = context.createScanner(tableName, Authorizations.EMPTY)) { scanner.setRange(getRow(fateId)); @@ -600,7 +600,7 @@ public Serializable getTransactionInfo(TxInfo txInfo) { @Override public long timeCreated() { - verifyReservedAndNotDeleted(false); + newMutator(fateId).requireReserved(reservation); return scanTx(scanner -> { scanner.setRange(getRow(fateId)); From 7b641deec3262e2b9b38e35be345022437d488c1 Mon Sep 17 00:00:00 2001 From: avillarreal Date: Fri, 4 Sep 2026 16:04:36 -0500 Subject: [PATCH 6/6] Create helper method newReservedMutator --- .../core/fate/user/UserFateStore.java | 35 +++++++++---------- 1 file changed, 16 insertions(+), 19 deletions(-) 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 f14bcf37cfe..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() { - newMutator(fateId).requireReserved(reservation); - return scanTx(scanner -> { scanner.setRange(getRow(fateId)); scanner.setBatchSize(1); @@ -562,8 +566,6 @@ public Repo top() { @Override public List> getStack() { - newMutator(fateId).requireReserved(reservation); - return scanTx(scanner -> { scanner.setRange(getRow(fateId)); scanner.fetchColumnFamily(RepoColumnFamily.NAME); @@ -577,8 +579,6 @@ public List> getStack() { @Override public Serializable getTransactionInfo(TxInfo txInfo) { - newMutator(fateId).requireReserved(reservation); - 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() { - newMutator(fateId).requireReserved(reservation); - return scanTx(scanner -> { scanner.setRange(getRow(fateId)); TxColumnFamily.CREATE_TIME_COLUMN.fetch(scanner); @@ -618,43 +616,42 @@ public void push(Repo repo) throws StackOverflowException { throw new StackOverflowException("Repo stack size too large"); } - FateMutator fateMutator = newMutator(fateId) - .requireStatus(REQ_PUSH_STATUS.toArray(TStatus[]::new)).requireReserved(reservation); + 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() { Optional top = findTop(); - top.ifPresent(t -> newMutator(fateId).requireStatus(REQ_POP_STATUS.toArray(TStatus[]::new)) - .requireReserved(reservation).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) { - newMutator(fateId).requireReserved(reservation).putStatus(status).mutate(); + newReservedMutator(fateId, reservation).putStatus(status).mutate(); observedStatus = status; } @Override public void setTransactionInfo(TxInfo txInfo, Serializable so) { final byte[] serialized = serializeTxInfo(so); - newMutator(fateId).requireReserved(reservation).putTxInfo(txInfo, serialized).mutate(); + newReservedMutator(fateId, reservation).putTxInfo(txInfo, serialized).mutate(); } @Override public void delete() { - var mutator = newMutator(fateId); - mutator.requireStatus(REQ_DELETE_STATUS.toArray(TStatus[]::new)).requireReserved(reservation); + var mutator = newReservedMutator(fateId, reservation); + mutator.requireStatus(REQ_DELETE_STATUS.toArray(TStatus[]::new)); mutator.delete().mutate(); this.deleted = true; } @Override public void forceDelete() { - var mutator = newMutator(fateId); - mutator.requireStatus(REQ_FORCE_DELETE_STATUS.toArray(TStatus[]::new)) - .requireReserved(reservation); + var mutator = newReservedMutator(fateId, reservation); + mutator.requireStatus(REQ_FORCE_DELETE_STATUS.toArray(TStatus[]::new)); mutator.delete().mutate(); this.deleted = true; }