Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,11 @@ public interface FateMutator<T> {
*/
FateMutator<T> requireUnreserved();

/**
* Require the transaction is reserved with the given reservation.
*/
FateMutator<T> requireReserved(FateStore.FateReservation fateReservation);

/**
* Require the transaction has no fate key set.
*/
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,7 @@ public class FateMutatorImpl<T> implements FateMutator<T> {
private final ConditionalMutation mutation;
private final Supplier<ConditionalWriter> 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,
Expand Down Expand Up @@ -108,6 +109,17 @@ public FateMutator<T> requireUnreserved() {
return this;
}

@Override
public FateMutator<T> 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<T> requireAbsentKey() {
Condition condition = new Condition(TxColumnFamily.TX_KEY_COLUMN.getColumnFamily(),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -544,10 +544,17 @@ private FateTxStoreImpl(FateId fateId, FateReservation reservation) {
super(fateId, reservation);
}

private FateMutatorImpl<T> newReservedMutator() {
Preconditions.checkState(isReserved(),
"Attempted write on unreserved FATE transaction: " + fateId);
Preconditions.checkState(!deleted, "Attempted write on deleted FATE transaction: " + fateId);
FateMutatorImpl<T> mutator = newMutator(fateId);
mutator.requireReserved(reservation);
return mutator;
}
Comment thread
Amemeda marked this conversation as resolved.

@Override
public Repo<T> top() {
verifyReservedAndNotDeleted(false);
Comment thread
dlmarion marked this conversation as resolved.

return scanTx(scanner -> {
scanner.setRange(getRow(fateId));
scanner.setBatchSize(1);
Expand All @@ -562,8 +569,6 @@ public Repo<T> top() {

@Override
public List<ReadOnlyRepo<T>> getStack() {
verifyReservedAndNotDeleted(false);

return scanTx(scanner -> {
scanner.setRange(getRow(fateId));
scanner.fetchColumnFamily(RepoColumnFamily.NAME);
Expand All @@ -577,8 +582,6 @@ public List<ReadOnlyRepo<T>> getStack() {

@Override
public Serializable getTransactionInfo(TxInfo txInfo) {
verifyReservedAndNotDeleted(false);

try (Scanner scanner = context.createScanner(tableName, Authorizations.EMPTY)) {
scanner.setRange(getRow(fateId));

Expand All @@ -600,8 +603,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);
Expand All @@ -612,60 +613,47 @@ public long timeCreated() {

@Override
public void push(Repo<T> repo) throws StackOverflowException {
verifyReservedAndNotDeleted(true);

Optional<Integer> top = findTop();

if (top.filter(t -> t >= MAX_REPOS).isPresent()) {
throw new StackOverflowException("Repo stack size too large");
}

FateMutator<T> fateMutator =
newMutator(fateId).requireStatus(REQ_PUSH_STATUS.toArray(TStatus[]::new));
newReservedMutator().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<Integer> top = findTop();
top.ifPresent(t -> newMutator(fateId).requireStatus(REQ_POP_STATUS.toArray(TStatus[]::new))
top.ifPresent(t -> newReservedMutator().requireStatus(REQ_POP_STATUS.toArray(TStatus[]::new))
.deleteRepo(t).mutate());
}

@Override
public void setStatus(TStatus status) {
verifyReservedAndNotDeleted(true);

newMutator(fateId).putStatus(status).mutate();
newReservedMutator().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().putTxInfo(txInfo, serialized).mutate();
}

@Override
public void delete() {
verifyReservedAndNotDeleted(true);

var mutator = newMutator(fateId);
var mutator = newReservedMutator();
mutator.requireStatus(REQ_DELETE_STATUS.toArray(TStatus[]::new));
mutator.delete().mutate();
this.deleted = true;
}

@Override
public void forceDelete() {
verifyReservedAndNotDeleted(true);

var mutator = newMutator(fateId);
var mutator = newReservedMutator();
mutator.requireStatus(REQ_FORCE_DELETE_STATUS.toArray(TStatus[]::new));
mutator.delete().mutate();
this.deleted = true;
Expand Down
Loading