Skip to content
This repository was archived by the owner on Mar 10, 2026. It is now read-only.

Commit 53d4db8

Browse files
authored
feat: StreamPosition type + fix: ditch Instant repr for timestamps (#18)
1 parent 9673d0e commit 53d4db8

12 files changed

Lines changed: 60 additions & 58 deletions

File tree

app/src/main/java/org/example/app/ManagedAppendSessionDemo.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -25,8 +25,8 @@ public class ManagedAppendSessionDemo {
2525
private static final Logger logger =
2626
LoggerFactory.getLogger(ManagedAppendSessionDemo.class.getName());
2727

28-
// 512KiB
29-
private static final Integer TARGET_BATCH_SIZE = 512 * 1024;
28+
// 128KiB
29+
private static final Integer TARGET_BATCH_SIZE = 128 * 1024;
3030

3131
public static void main(String[] args) throws Exception {
3232
final var authToken = System.getenv("S2_ACCESS_TOKEN");

app/src/main/java/org/example/app/ManagedReadSessionDemo.java

Lines changed: 3 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -61,12 +61,10 @@ public static void main(String[] args) throws Exception {
6161
if (elem instanceof Batch batch) {
6262
var size = batch.meteredBytes();
6363
logger.info(
64-
"batch of {} bytes, seqnums {}..={} / instants {}..={}",
64+
"batch of {} bytes, first={} ..= last={}",
6565
size,
66-
batch.firstSeqNum(),
67-
batch.lastSeqNum(),
68-
batch.firstTimestamp(),
69-
batch.lastTimestamp());
66+
batch.firstPosition(),
67+
batch.lastPosition());
7068
receivedBytes.addAndGet(size);
7169
} else {
7270
logger.info("non batch received: {}", elem);

gradle.properties

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1 +1 @@
1-
version=0.0.16-SNAPSHOT
1+
version=0.0.16

s2-internal/src/main/proto

Submodule proto updated from 4183c36 to b6add44

s2-sdk/src/main/java/s2/client/ManagedAppendSession.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -235,7 +235,7 @@ private synchronized Void cleanUp(Status fatal) throws InterruptedException {
235235
}
236236

237237
private void validate(InflightRecord record, AppendOutput output) {
238-
var numRecordsForAcknowledgement = output.endSeqNum - output.startSeqNum;
238+
var numRecordsForAcknowledgement = output.end.seqNum - output.start.seqNum;
239239
if (numRecordsForAcknowledgement != record.input.records.size()) {
240240
throw Status.INTERNAL
241241
.withDescription(

s2-sdk/src/main/java/s2/client/ReadSession.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -140,8 +140,8 @@ private ListenableFuture<Void> retrying() {
140140
resp -> {
141141
if (resp instanceof Batch) {
142142
final Batch batch = (Batch) resp;
143-
var lastRecordIdx = batch.lastSeqNum();
144-
lastRecordIdx.ifPresent(v -> nextStart.set(Start.seqNum(v + 1)));
143+
var lastPosition = batch.lastPosition();
144+
lastPosition.ifPresent(v -> nextStart.set(Start.seqNum(v.seqNum + 1)));
145145
consumedRecords.addAndGet(batch.sequencedRecordBatch.records.size());
146146
consumedBytes.addAndGet(batch.meteredBytes());
147147
}

s2-sdk/src/main/java/s2/client/StreamClient.java

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -21,12 +21,12 @@
2121
import s2.types.ReadOutput;
2222
import s2.types.ReadRequest;
2323
import s2.types.ReadSessionRequest;
24+
import s2.types.StreamPosition;
2425
import s2.v1alpha.AppendRequest;
2526
import s2.v1alpha.AppendResponse;
2627
import s2.v1alpha.AppendSessionRequest;
2728
import s2.v1alpha.AppendSessionResponse;
2829
import s2.v1alpha.CheckTailRequest;
29-
import s2.v1alpha.CheckTailResponse;
3030
import s2.v1alpha.StreamServiceGrpc;
3131
import s2.v1alpha.StreamServiceGrpc.StreamServiceFutureStub;
3232
import s2.v1alpha.StreamServiceGrpc.StreamServiceStub;
@@ -81,9 +81,9 @@ public static StreamClientBuilder newBuilder(Config config, String basinName, St
8181
/**
8282
* Check the sequence number that will be assigned to the next record on a stream.
8383
*
84-
* @return future of the next sequence number
84+
* @return future of the tail's position
8585
*/
86-
public ListenableFuture<Long> checkTail() {
86+
public ListenableFuture<StreamPosition> checkTail() {
8787
return withTimeout(
8888
() ->
8989
Futures.transform(
@@ -92,7 +92,7 @@ public ListenableFuture<Long> checkTail() {
9292
() ->
9393
this.futureStub.checkTail(
9494
CheckTailRequest.newBuilder().setStream(streamName).build())),
95-
CheckTailResponse::getNextSeqNum,
95+
(resp) -> new StreamPosition(resp.getNextSeqNum(), resp.getLastTimestamp()),
9696
executor));
9797
}
9898

Lines changed: 11 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -1,25 +1,25 @@
11
package s2.types;
22

33
public class AppendOutput {
4-
public final long startSeqNum;
5-
public final long endSeqNum;
6-
public final long nextSeqNum;
4+
public final StreamPosition start;
5+
public final StreamPosition end;
6+
public final StreamPosition tail;
77

8-
AppendOutput(long startSeqNum, long endSeqNum, long nextSeqNum) {
9-
this.startSeqNum = startSeqNum;
10-
this.endSeqNum = endSeqNum;
11-
this.nextSeqNum = nextSeqNum;
8+
AppendOutput(StreamPosition start, StreamPosition end, StreamPosition tail) {
9+
this.start = start;
10+
this.end = end;
11+
this.tail = tail;
1212
}
1313

1414
public static AppendOutput fromProto(s2.v1alpha.AppendOutput appendOutput) {
1515
return new AppendOutput(
16-
appendOutput.getStartSeqNum(), appendOutput.getEndSeqNum(), appendOutput.getNextSeqNum());
16+
new StreamPosition(appendOutput.getStartSeqNum(), appendOutput.getStartTimestamp()),
17+
new StreamPosition(appendOutput.getEndSeqNum(), appendOutput.getEndTimestamp()),
18+
new StreamPosition(appendOutput.getNextSeqNum(), appendOutput.getLastTimestamp()));
1719
}
1820

1921
@Override
2022
public String toString() {
21-
return String.format(
22-
"AppendOutput[startSeqNum=%s, endSeqNum=%s, nextSeqNum=%s]",
23-
startSeqNum, endSeqNum, nextSeqNum);
23+
return String.format("AppendOutput[start=%s, end=%s, tail=%s]", start, end, tail);
2424
}
2525
}

s2-sdk/src/main/java/s2/types/Batch.java

Lines changed: 8 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,5 @@
11
package s2.types;
22

3-
import java.time.Instant;
43
import java.util.Optional;
54

65
public final class Batch implements ReadOutput, MeteredBytes {
@@ -11,27 +10,17 @@ public final class Batch implements ReadOutput, MeteredBytes {
1110
this.sequencedRecordBatch = sequencedRecordBatch;
1211
}
1312

14-
public Optional<Long> firstSeqNum() {
15-
return this.sequencedRecordBatch.records.stream().findFirst().map(sr -> sr.seqNum);
13+
public Optional<StreamPosition> firstPosition() {
14+
return this.sequencedRecordBatch.records.stream()
15+
.findFirst()
16+
.map(sr -> new StreamPosition(sr.seqNum, sr.timestamp));
1617
}
1718

18-
public Optional<Long> lastSeqNum() {
19+
public Optional<StreamPosition> lastPosition() {
1920
if (!this.sequencedRecordBatch.records.isEmpty()) {
20-
return Optional.of(
21-
this.sequencedRecordBatch.records.get(sequencedRecordBatch.records.size() - 1).seqNum);
22-
} else {
23-
return Optional.empty();
24-
}
25-
}
26-
27-
public Optional<Instant> firstTimestamp() {
28-
return this.sequencedRecordBatch.records.stream().findFirst().map(sr -> sr.timestamp);
29-
}
30-
31-
public Optional<Instant> lastTimestamp() {
32-
if (!this.sequencedRecordBatch.records.isEmpty()) {
33-
return Optional.of(
34-
this.sequencedRecordBatch.records.get(sequencedRecordBatch.records.size() - 1).timestamp);
21+
var lastRecord =
22+
this.sequencedRecordBatch.records.get(this.sequencedRecordBatch.records.size() - 1);
23+
return Optional.of(new StreamPosition(lastRecord.seqNum, lastRecord.timestamp));
3524
} else {
3625
return Optional.empty();
3726
}

s2-sdk/src/main/java/s2/types/SequencedRecord.java

Lines changed: 3 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,17 +1,16 @@
11
package s2.types;
22

33
import com.google.protobuf.ByteString;
4-
import java.time.Instant;
54
import java.util.List;
65
import java.util.stream.Collectors;
76

87
public class SequencedRecord implements MeteredBytes {
98
public final long seqNum;
109
public final List<Header> headers;
1110
public final ByteString body;
12-
public final Instant timestamp;
11+
public final long timestamp;
1312

14-
SequencedRecord(long seqNum, List<Header> headers, ByteString body, Instant timestamp) {
13+
SequencedRecord(long seqNum, List<Header> headers, ByteString body, long timestamp) {
1514
this.seqNum = seqNum;
1615
this.headers = headers;
1716
this.body = body;
@@ -25,7 +24,7 @@ public static SequencedRecord fromProto(s2.v1alpha.SequencedRecord sequencedReco
2524
.map(Header::fromProto)
2625
.collect(Collectors.toList()),
2726
sequencedRecord.getBody(),
28-
Instant.ofEpochMilli(sequencedRecord.getTimestamp()));
27+
sequencedRecord.getTimestamp());
2928
}
3029

3130
@Override

0 commit comments

Comments
 (0)