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 @@ -3,6 +3,7 @@
import com.google.cloud.firestore.DocumentReference;
import com.google.cloud.firestore.DocumentSnapshot;
import com.google.cloud.firestore.Firestore;
import com.google.cloud.firestore.Precondition;
import com.google.cloud.firestore.Query;
import com.google.cloud.firestore.QueryDocumentSnapshot;
import com.google.cloud.firestore.QuerySnapshot;
Expand All @@ -18,8 +19,10 @@
import java.util.List;
import java.util.Map;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.atomic.AtomicInteger;

import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;

@TestMethodOrder(MethodOrderer.OrderAnnotation.class)
class FirestoreTest {
Expand Down Expand Up @@ -171,4 +174,80 @@ void deleteDocument() throws ExecutionException, InterruptedException {
DocumentSnapshot snapshot = docRef.get().get();
assertThat(snapshot.exists()).isFalse();
}

@Test
@Order(7)
void createFailsWhenDocumentAlreadyExists() throws ExecutionException, InterruptedException {
DocumentReference docRef = firestore.collection(COLLECTION)
.document(TestFixtures.uniqueName("precond"));
docRef.create(Map.of("name", "Alice")).get();
try {
assertThatThrownBy(() -> docRef.create(Map.of("name", "Bob")).get())
.isInstanceOf(ExecutionException.class)
.cause()
.hasMessageContaining("already exists");

DocumentSnapshot snapshot = docRef.get().get();
assertThat(snapshot.getString("name")).isEqualTo("Alice");
} finally {
docRef.delete().get();
}
}

@Test
@Order(8)
void updateFailsWhenDocumentMissing() {
DocumentReference docRef = firestore.collection(COLLECTION)
.document(TestFixtures.uniqueName("missing"));

assertThatThrownBy(() -> docRef.update("name", "Alice").get())
.isInstanceOf(ExecutionException.class)
.cause()
.hasMessageContaining("No document to update");
}

@Test
@Order(9)
void staleUpdateTimePreconditionFails() throws ExecutionException, InterruptedException {
DocumentReference docRef = firestore.collection(COLLECTION)
.document(TestFixtures.uniqueName("precond"));
WriteResult first = docRef.set(Map.of("version", 1L)).get();
WriteResult second = docRef.set(Map.of("version", 2L)).get();
try {
assertThatThrownBy(() -> docRef.update(
Map.of("version", 3L), Precondition.updatedAt(first.getUpdateTime())).get())
.isInstanceOf(ExecutionException.class);

// matching precondition succeeds
docRef.update(Map.of("version", 3L), Precondition.updatedAt(second.getUpdateTime())).get();
assertThat(docRef.get().get().getLong("version")).isEqualTo(3L);
} finally {
docRef.delete().get();
}
}

@Test
@Order(10)
void transactionRetriesOnConcurrentModification() throws ExecutionException, InterruptedException {
DocumentReference docRef = firestore.collection(COLLECTION)
.document(TestFixtures.uniqueName("counter"));
docRef.set(Map.of("count", 0L)).get();
AtomicInteger attempts = new AtomicInteger();
try {
firestore.runTransaction(transaction -> {
long current = transaction.get(docRef).get().getLong("count");
if (attempts.incrementAndGet() == 1) {
// out-of-band write invalidates the transaction's read set
docRef.update("count", 100L).get();
}
transaction.update(docRef, "count", current + 1);
return null;
}).get();

assertThat(attempts.get()).isEqualTo(2);
assertThat(docRef.get().get().getLong("count")).isEqualTo(101L);
} finally {
docRef.delete().get();
}
}
}
33 changes: 33 additions & 0 deletions docs/services/firestore.md
Original file line number Diff line number Diff line change
Expand Up @@ -117,6 +117,25 @@ db.runTransaction(transaction -> {
}).get();
```

Transactions use optimistic concurrency, matching real Firestore: every document
read inside a transaction is tracked, and the commit fails with `ABORTED` if any
of those documents changed after being read. SDK clients retry aborted
transactions automatically, so concurrent increments like the example above never
lose updates.

## Preconditions

Write preconditions (`currentDocument`) are enforced on all write paths:

- `create()` fails with `ALREADY_EXISTS` if the document exists
- `update()` fails with `NOT_FOUND` if the document does not exist
- `Precondition.updatedAt(...)` fails with `FAILED_PRECONDITION` if the
document's update time no longer matches

`Commit` is atomic: if any write's precondition fails, no writes in the request
are applied. `BatchWrite` is non-atomic and reports a per-write status, matching
real Firestore.

## Batch Writes

```java
Expand Down Expand Up @@ -164,3 +183,17 @@ registration.remove();
- `Write` (streaming)
- `Listen` (real-time change streams)
- `ListCollectionIds`

## Deviations from real Firestore

- Transaction conflict detection covers documents actually read. Queries inside
a transaction track only the documents they return, so phantom reads (a
concurrent write creating a document that would have matched the query) do not
abort the transaction.
- `RunAggregationQuery` ignores `transaction` and `newTransaction`; aggregations
read outside the transaction and are not part of its read set.
- Transactions expire after 15 minutes instead of Firestore's shorter
server-side deadlines.
- Transaction state is held in memory; transactions do not survive an emulator
restart. A commit against an unknown transaction id is applied without conflict
validation.
4 changes: 4 additions & 0 deletions src/main/java/io/floci/gcp/core/common/GcpException.java
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,10 @@ public static GcpException outOfRange(String message) {
return new GcpException(416, "OUT_OF_RANGE", Status.Code.OUT_OF_RANGE, message);
}

public static GcpException aborted(String message) {
return new GcpException(409, "ABORTED", Status.Code.ABORTED, message);
}

public static GcpException failedPrecondition(String message) {
return new GcpException(400, "FAILED_PRECONDITION", Status.Code.FAILED_PRECONDITION, message);
}
Expand Down
115 changes: 88 additions & 27 deletions src/main/java/io/floci/gcp/services/firestore/FirestoreController.java
Original file line number Diff line number Diff line change
Expand Up @@ -33,8 +33,9 @@ public void commit(CommitRequest request, StreamObserver<CommitResponse> respons
CommitResponse.Builder response = CommitResponse.newBuilder()
.setCommitTime(toTimestamp(commitTime.toString()));

for (Write write : request.getWritesList()) {
FirestoreService.WriteCommitResult result = service.applyWrite(write, commitTime);
List<FirestoreService.WriteCommitResult> results = service.commit(
request.getWritesList(), request.getTransaction().toByteArray(), commitTime);
for (FirestoreService.WriteCommitResult result : results) {
WriteResult.Builder wr = WriteResult.newBuilder();
if (result.updateTime() != null) {
wr.setUpdateTime(toTimestamp(result.updateTime()));
Expand All @@ -54,10 +55,14 @@ public void commit(CommitRequest request, StreamObserver<CommitResponse> respons
public void getDocument(GetDocumentRequest request, StreamObserver<Document> responseObserver) {
LOG.debugf("getDocument name=%s", request.getName());
try {
StoredDocument stored = service.getDocument(request.getName())
.orElseThrow(() -> io.floci.gcp.core.common.GcpException.notFound(
"Document not found: " + request.getName()));
responseObserver.onNext(toProto(stored));
Optional<StoredDocument> stored = service.getDocument(request.getName());
if (!request.getTransaction().isEmpty()) {
service.recordTransactionRead(request.getTransaction().toByteArray(), request.getName(),
stored.map(StoredDocument::getUpdateTime).orElse(null));
}
responseObserver.onNext(toProto(stored.orElseThrow(
() -> io.floci.gcp.core.common.GcpException.notFound(
"Document not found: " + request.getName()))));
responseObserver.onCompleted();
} catch (Exception e) {
LOG.warnf("getDocument failed: %s", e.getMessage());
Expand All @@ -71,17 +76,37 @@ public void batchGetDocuments(BatchGetDocumentsRequest request,
LOG.debugf("batchGetDocuments database=%s docs=%d", request.getDatabase(), request.getDocumentsCount());
try {
Instant readTime = Instant.now();
ByteString txId = request.getTransaction();
boolean newTransaction = request.hasNewTransaction();
if (newTransaction) {
txId = ByteString.copyFrom(service.beginTransaction());
}
boolean first = true;
for (String docName : request.getDocumentsList()) {
Optional<StoredDocument> stored = service.getDocument(docName);
if (!txId.isEmpty()) {
service.recordTransactionRead(txId.toByteArray(), docName,
stored.map(StoredDocument::getUpdateTime).orElse(null));
}
BatchGetDocumentsResponse.Builder resp = BatchGetDocumentsResponse.newBuilder()
.setReadTime(toTimestamp(readTime.toString()));
if (newTransaction && first) {
resp.setTransaction(txId);
first = false;
}
if (stored.isPresent()) {
resp.setFound(toProto(stored.get()));
} else {
resp.setMissing(docName);
}
responseObserver.onNext(resp.build());
}
if (newTransaction && first) {
responseObserver.onNext(BatchGetDocumentsResponse.newBuilder()
.setReadTime(toTimestamp(readTime.toString()))
.setTransaction(txId)
.build());
}
responseObserver.onCompleted();
} catch (Exception e) {
LOG.warnf("batchGetDocuments failed: %s", e.getMessage());
Expand All @@ -94,20 +119,36 @@ public void runQuery(RunQueryRequest request, StreamObserver<RunQueryResponse> r
LOG.debugf("runQuery parent=%s", request.getParent());
try {
Instant readTime = Instant.now();
ByteString txId = request.getTransaction();
boolean newTransaction = request.hasNewTransaction();
if (newTransaction) {
txId = ByteString.copyFrom(service.beginTransaction());
}
List<StoredDocument> results = service.runQuery(request.getParent(), request.getStructuredQuery());

boolean first = true;
for (StoredDocument doc : results) {
responseObserver.onNext(RunQueryResponse.newBuilder()
if (!txId.isEmpty()) {
service.recordTransactionRead(txId.toByteArray(), doc.getName(), doc.getUpdateTime());
}
RunQueryResponse.Builder resp = RunQueryResponse.newBuilder()
.setDocument(toProto(doc))
.setReadTime(toTimestamp(readTime.toString()))
.build());
.setReadTime(toTimestamp(readTime.toString()));
if (newTransaction && first) {
resp.setTransaction(txId);
first = false;
}
responseObserver.onNext(resp.build());
}

// terminal message
responseObserver.onNext(RunQueryResponse.newBuilder()
RunQueryResponse.Builder done = RunQueryResponse.newBuilder()
.setReadTime(toTimestamp(readTime.toString()))
.setDone(true)
.build());
.setDone(true);
if (newTransaction && first) {
done.setTransaction(txId);
}
responseObserver.onNext(done.build());
responseObserver.onCompleted();
} catch (Exception e) {
LOG.warnf("runQuery failed: %s", e.getMessage());
Expand All @@ -134,6 +175,7 @@ public void beginTransaction(BeginTransactionRequest request,
@Override
public void rollback(RollbackRequest request, StreamObserver<Empty> responseObserver) {
LOG.debugf("rollback database=%s", request.getDatabase());
service.rollback(request.getTransaction().toByteArray());
responseObserver.onNext(Empty.getDefaultInstance());
responseObserver.onCompleted();
}
Expand Down Expand Up @@ -163,11 +205,13 @@ public void listDocuments(ListDocumentsRequest request, StreamObserver<ListDocum
public void updateDocument(UpdateDocumentRequest request, StreamObserver<Document> responseObserver) {
LOG.debugf("updateDocument name=%s", request.getDocument().getName());
try {
Write write = Write.newBuilder()
Write.Builder write = Write.newBuilder()
.setUpdate(request.getDocument())
.setUpdateMask(request.getUpdateMask())
.build();
service.applyWrite(write, Instant.now());
.setUpdateMask(request.getUpdateMask());
if (request.hasCurrentDocument()) {
write.setCurrentDocument(request.getCurrentDocument());
}
service.applyWrite(write.build(), Instant.now());
StoredDocument stored = service.getDocument(request.getDocument().getName())
.orElseThrow();
responseObserver.onNext(toProto(stored));
Expand All @@ -182,8 +226,11 @@ public void updateDocument(UpdateDocumentRequest request, StreamObserver<Documen
public void deleteDocument(DeleteDocumentRequest request, StreamObserver<Empty> responseObserver) {
LOG.debugf("deleteDocument name=%s", request.getName());
try {
Write write = Write.newBuilder().setDelete(request.getName()).build();
service.applyWrite(write, Instant.now());
Write.Builder write = Write.newBuilder().setDelete(request.getName());
if (request.hasCurrentDocument()) {
write.setCurrentDocument(request.getCurrentDocument());
}
service.applyWrite(write.build(), Instant.now());
responseObserver.onNext(Empty.getDefaultInstance());
responseObserver.onCompleted();
} catch (Exception e) {
Expand All @@ -201,7 +248,10 @@ public void createDocument(CreateDocumentRequest request, StreamObserver<Documen
: request.getDocumentId();
String name = request.getParent() + "/" + request.getCollectionId() + "/" + docId;
Document doc = request.getDocument().toBuilder().setName(name).build();
Write write = Write.newBuilder().setUpdate(doc).build();
Write write = Write.newBuilder()
.setUpdate(doc)
.setCurrentDocument(Precondition.newBuilder().setExists(false).build())
.build();
service.applyWrite(write, Instant.now());
responseObserver.onNext(toProto(service.getDocument(name).orElseThrow()));
responseObserver.onCompleted();
Expand Down Expand Up @@ -232,14 +282,25 @@ public void batchWrite(BatchWriteRequest request, StreamObserver<BatchWriteRespo
try {
Instant commitTime = Instant.now();
BatchWriteResponse.Builder resp = BatchWriteResponse.newBuilder();
// BatchWrite is non-atomic: each write succeeds or fails independently
for (Write write : request.getWritesList()) {
FirestoreService.WriteCommitResult result = service.applyWrite(write, commitTime);
WriteResult.Builder wr = WriteResult.newBuilder();
if (result.updateTime() != null) {
wr.setUpdateTime(toTimestamp(result.updateTime()));
try {
FirestoreService.WriteCommitResult result = service.applyWrite(write, commitTime);
WriteResult.Builder wr = WriteResult.newBuilder();
if (result.updateTime() != null) {
wr.setUpdateTime(toTimestamp(result.updateTime()));
}
resp.addWriteResults(wr.build());
resp.addStatus(Status.newBuilder().setCode(0).build());
} catch (RuntimeException e) {
io.grpc.Status mapped = GcpGrpcController.grpcException(e).getStatus();
String message = mapped.getDescription();
resp.addWriteResults(WriteResult.getDefaultInstance());
resp.addStatus(Status.newBuilder()
.setCode(mapped.getCode().value())
.setMessage(message == null ? "" : message)
.build());
}
resp.addWriteResults(wr.build());
resp.addStatus(com.google.rpc.Status.newBuilder().setCode(0).build());
}
responseObserver.onNext(resp.build());
responseObserver.onCompleted();
Expand Down Expand Up @@ -294,8 +355,8 @@ public void onNext(WriteRequest req) {
WriteResponse.Builder resp = WriteResponse.newBuilder()
.setStreamId(req.getStreamId())
.setCommitTime(toTimestamp(now.toString()));
for (Write w : req.getWritesList()) {
FirestoreService.WriteCommitResult r = service.applyWrite(w, now);
for (FirestoreService.WriteCommitResult r
: service.commit(req.getWritesList(), null, now)) {
WriteResult.Builder wr = WriteResult.newBuilder();
if (r.updateTime() != null) {
wr.setUpdateTime(toTimestamp(r.updateTime()));
Expand Down
Loading