Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
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 has a 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 @@ -414,6 +414,12 @@ private FateMutatorImpl<T> newMutator(FateId fateId) {
return new FateMutatorImpl<>(context, tableName, fateId, writer);
}

public FateMutatorImpl<T> newReservedMutator(FateId fateId, FateReservation reservation) {
Preconditions.checkState(fateId != null,
"Attempted write on deleted FATE transaction: " + fateId);
return (FateMutatorImpl<T>) newMutator(fateId).requireReserved(reservation);
}

private <R> R scanTx(Function<Scanner,R> func) {
try (Scanner scanner = context.createScanner(tableName, Authorizations.EMPTY)) {
return func.apply(scanner);
Expand Down Expand Up @@ -546,8 +552,6 @@ private FateTxStoreImpl(FateId fateId, FateReservation reservation) {

@Override
public Repo<T> top() {
verifyReservedAndNotDeleted(false);

return scanTx(scanner -> {
scanner.setRange(getRow(fateId));
scanner.setBatchSize(1);
Expand All @@ -562,8 +566,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 +579,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 +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);
Expand All @@ -612,60 +610,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));
FateMutator<T> 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<Integer> 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;
}

@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;
Expand Down