Skip to content
Closed
Show file tree
Hide file tree
Changes from 1 commit
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
48 changes: 25 additions & 23 deletions examples/streamed-list-objects/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -57,50 +57,52 @@ export FGA_API_AUDIENCE=your_audience
```java
// Create a request
var request = new ClientListObjectsRequest()
.type("document")
.relation("owner")
.user("user:anne");

// Call the streaming API
var objectStream = fgaClient.streamedListObjects(request).get();
.type("document")
Comment thread
SoulPancake marked this conversation as resolved.
.relation("owner")
.user("user:anne");

// Call the streaming API and ensure proper resource cleanup
try (var objectStream = fgaClient.streamedListObjects(request).get()) {
// Collect all results
List<String> objects = objectStream
.map(StreamedListObjectsResponse::getObject)
.collect(Collectors.toList());
.map(StreamedListObjectsResponse::getObject)
.collect(Collectors.toList());
}
```

### Early Termination

```java
// Get only the first 10 results
var objectStream = fgaClient.streamedListObjects(request).get();
List<String> firstTen = objectStream
.map(StreamedListObjectsResponse::getObject)
.limit(10)
.collect(Collectors.toList());
// Get only the first 10 results, ensuring the stream is closed properly
try (var objectStream = fgaClient.streamedListObjects(request).get()) {
List<String> firstTen = objectStream
.map(StreamedListObjectsResponse::getObject)
.limit(10)
.collect(Collectors.toList());
}
```

### Process as You Go

```java
// Process each object immediately as it arrives
var objectStream = fgaClient.streamedListObjects(request).get();
objectStream
.map(StreamedListObjectsResponse::getObject)
.forEach(obj -> {
// Do something with each object
System.out.println("Processing: " + obj);
});
try (var objectStream = fgaClient.streamedListObjects(request).get()) {
objectStream
.map(StreamedListObjectsResponse::getObject)
.forEach(obj -> {
// Do something with each object
System.out.println("Processing: " + obj);
});
}
```

### With Options

```java
// Use options to specify consistency preference
var options = new ClientListObjectsOptions()
.consistency(ConsistencyPreference.HIGHER_CONSISTENCY)
.authorizationModelId("01GXSXXXXXXXXXXXXXXXX");
.consistency(ConsistencyPreference.HIGHER_CONSISTENCY)
.authorizationModelId("01GXSXXXXXXXXXXXXXXXX");

var objectStream = fgaClient.streamedListObjects(request, options).get();
```
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -138,8 +138,9 @@ private static int writeTuples(OpenFgaClient fgaClient, int quantity) throws Exc
*/
private static List<String> streamedListObjects(OpenFgaClient fgaClient, ClientListObjectsRequest request)
throws Exception {
var objectStream = fgaClient.streamedListObjects(request).get();
return objectStream.map(StreamedListObjectsResponse::getObject).collect(Collectors.toList());
try (var objectStream = fgaClient.streamedListObjects(request).get()) {
return objectStream.map(StreamedListObjectsResponse::getObject).collect(Collectors.toList());
}
}

/**
Expand Down
4 changes: 4 additions & 0 deletions src/main/java/dev/openfga/sdk/api/client/OpenFgaClient.java
Original file line number Diff line number Diff line change
Expand Up @@ -1175,6 +1175,10 @@ public CompletableFuture<Stream<StreamedListObjectsResponse>> streamedListObject
Stream<StreamedListObjectsResponse> stream = java.util.stream.StreamSupport.stream(
((Iterable<StreamedListObjectsResponse>) () -> iterator).spliterator(), false);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: do we need to cast to Iterable from a lambda to get a stream, or can we use any Spliterator functionality, like Spliterators.spliteratorUnknownSize or similar?

return stream.onClose(() -> {
try {
iterator.close();
} catch (java.io.IOException ignore) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

should we be ignoring exceptions closing the iterator or response body? Seems like it could cause issues to swallow these errors. If we can recover we should, if not can we throw back to caller with applicable information? Or if there is good reason to swallow exceptions, let's add clear code comment explaining why.

}
try {
srb.close();
} catch (java.io.IOException ignore) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@
* If an error is encountered in the stream (either from parsing or from an error
* response), it will be thrown as a StreamingException when hasNext() or next() is called.
*/
public class StreamedResponseIterator implements Iterator<StreamedListObjectsResponse> {
public class StreamedResponseIterator implements Iterator<StreamedListObjectsResponse>, AutoCloseable {
private final BufferedReader reader;
private final ObjectMapper objectMapper;
private StreamedListObjectsResponse nextItem;
Expand Down Expand Up @@ -92,4 +92,9 @@ public StreamedListObjectsResponse next() {

return current;
}

@Override
public void close() throws IOException {
reader.close();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -328,6 +328,33 @@ public void write_and_listObjects() throws Exception {
assertEquals(DEFAULT_DOC, response.getObjects().get(0));
}

@Test
public void write_and_streamedListObjects() throws Exception {
// Given
String storeName = thisTestName();
String storeId = createStore(storeName);
fga.setStoreId(storeId);
String authModelId = writeAuthModel(storeId);
fga.setAuthorizationModelId(authModelId);
ClientWriteRequest writeRequest = new ClientWriteRequest().writes(List.of(DEFAULT_TUPLE_KEY));
ClientListObjectsRequest listObjectsRequest = new ClientListObjectsRequest()
.user(DEFAULT_USER)
.relation("reader")
.type("document");

// When
fga.write(writeRequest).get();
List<String> objects;
try (var stream = fga.streamedListObjects(listObjectsRequest).get()) {
objects = stream.map(StreamedListObjectsResponse::getObject).collect(java.util.stream.Collectors.toList());
}

// Then
assertNotNull(objects);
assertEquals(1, objects.size());
assertEquals(DEFAULT_DOC, objects.get(0));
}

@Test
public void write_readAssertions() throws Exception {
// Given
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2610,10 +2610,11 @@ public void streamedListObjectsTest() throws Exception {
.user(DEFAULT_USER);

// When
Stream<StreamedListObjectsResponse> responseStream =
fga.streamedListObjects(request).get();
List<String> objects =
responseStream.map(StreamedListObjectsResponse::getObject).collect(Collectors.toList());
List<String> objects;
try (Stream<StreamedListObjectsResponse> responseStream =
fga.streamedListObjects(request).get()) {
objects = responseStream.map(StreamedListObjectsResponse::getObject).collect(Collectors.toList());
}

// Then
mockHttpClient.verify().post(postPath).withBody(is(expectedBody)).called(1);
Expand Down Expand Up @@ -2702,12 +2703,12 @@ public void streamedListObjects_errorInStream() throws Exception {
.user(DEFAULT_USER);

// When
Stream<StreamedListObjectsResponse> responseStream =
fga.streamedListObjects(request).get();

// Then - should throw when processing the stream
var exception = assertThrows(RuntimeException.class, () -> {
responseStream.map(StreamedListObjectsResponse::getObject).collect(Collectors.toList());
try (Stream<StreamedListObjectsResponse> responseStream =
fga.streamedListObjects(request).get()) {
responseStream.map(StreamedListObjectsResponse::getObject).collect(Collectors.toList());
}
});

assertTrue(exception.getMessage().contains("Error in streaming response"));
Expand Down
Loading