Skip to content

Commit 50f3be8

Browse files
fix: use the closed window in the eof response (#221)
Signed-off-by: Vaibhav Tiwari <vaibhav.tiwari33@gmail.com>
1 parent 7f129ed commit 50f3be8

3 files changed

Lines changed: 110 additions & 6 deletions

File tree

src/main/java/io/numaproj/numaflow/accumulator/AccumulatorActor.java

Lines changed: 18 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,7 @@ public Receive createReceive() {
3434
return ReceiveBuilder
3535
.create()
3636
.match(HandlerDatum.class, this::invokeHandler)
37+
.match(AccumulatorOuterClass.KeyedWindow.class, this::handleCloseWindow)
3738
.match(String.class, this::sendEOF)
3839
.build();
3940
}
@@ -42,17 +43,29 @@ private void invokeHandler(HandlerDatum handlerDatum) {
4243
this.accumulator.processMessage(handlerDatum, outputStream);
4344
}
4445

46+
// CLOSE: echo the exact close window (including slot)
47+
private void handleCloseWindow(AccumulatorOuterClass.KeyedWindow closeWindow) {
48+
sendEOFResponse(closeWindow);
49+
}
50+
51+
// Fallback: the input stream completed without a CLOSE (broadcast EOF). Keep prior
52+
// behavior — echo the OPEN window (start/end/keys).
4553
private void sendEOF(String EOF) {
54+
sendEOFResponse(AccumulatorOuterClass.KeyedWindow
55+
.newBuilder()
56+
.setStart(this.keyedWindow.getStart())
57+
.setEnd(this.keyedWindow.getEnd())
58+
.addAllKeys(this.keyedWindow.getKeysList())
59+
.build());
60+
}
61+
62+
private void sendEOFResponse(AccumulatorOuterClass.KeyedWindow eofWindow) {
4663
// invoke handleEndOfStream to materialize the messages received so far.
4764
this.accumulator.handleEndOfStream(outputStream);
4865

4966
AccumulatorOuterClass.AccumulatorResponse eofResponse = AccumulatorOuterClass.AccumulatorResponse
5067
.newBuilder()
51-
.setWindow(AccumulatorOuterClass.KeyedWindow
52-
.newBuilder()
53-
.setStart(this.keyedWindow.getStart())
54-
.setEnd(this.keyedWindow.getEnd())
55-
.addAllKeys(this.keyedWindow.getKeysList()))
68+
.setWindow(eofWindow)
5669
.setEOF(true)
5770
.build();
5871

src/main/java/io/numaproj/numaflow/accumulator/AccumulatorSupervisorActor.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -106,7 +106,8 @@ private void invokeActor(AccumulatorOuterClass.AccumulatorRequest request) {
106106
break;
107107
}
108108
case CLOSE: {
109-
actorsMap.get(uniqueId).tell(Constants.EOF, getSelf());
109+
// Send the CLOSE window to the child actor so it can echo it in the EOF response
110+
actorsMap.get(uniqueId).tell(request.getOperation().getKeyedWindow(), getSelf());
110111
actorsMap.remove(uniqueId);
111112
break;
112113
}

src/test/java/io/numaproj/numaflow/accumulator/ServerTest.java

Lines changed: 90 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
package io.numaproj.numaflow.accumulator;
22

33
import com.google.protobuf.ByteString;
4+
import com.google.protobuf.Timestamp;
45
import io.grpc.ManagedChannel;
56
import io.grpc.inprocess.InProcessChannelBuilder;
67
import io.grpc.inprocess.InProcessServerBuilder;
@@ -124,6 +125,95 @@ public void testAccumulatorSingleKey() {
124125
}
125126
}
126127

128+
@Test
129+
public void testAccumulatorEOFEchoesCloseWindow() {
130+
List<String> keys = List.of("test-accumulator");
131+
132+
AccumulatorOuterClass.KeyedWindow openWindow = AccumulatorOuterClass.KeyedWindow
133+
.newBuilder()
134+
.setStart(Timestamp.newBuilder().setSeconds(0).build())
135+
.setEnd(Timestamp.newBuilder().setSeconds(60).build())
136+
.setSlot("slot-0")
137+
.addAllKeys(keys)
138+
.build();
139+
140+
AccumulatorOuterClass.AccumulatorRequest openRequest = AccumulatorOuterClass.AccumulatorRequest
141+
.newBuilder()
142+
.setPayload(AccumulatorOuterClass.Payload
143+
.newBuilder()
144+
.setValue(ByteString.copyFromUtf8("test-payload"))
145+
.addAllKeys(keys)
146+
.build())
147+
.setOperation(AccumulatorOuterClass.AccumulatorRequest.WindowOperation
148+
.newBuilder()
149+
.setEvent(AccumulatorOuterClass.AccumulatorRequest.WindowOperation.Event.OPEN)
150+
.setKeyedWindow(openWindow)
151+
.build())
152+
.build();
153+
154+
AccumulatorOuterClass.AccumulatorRequest appendRequest = AccumulatorOuterClass.AccumulatorRequest
155+
.newBuilder()
156+
.setPayload(AccumulatorOuterClass.Payload
157+
.newBuilder()
158+
.setValue(ByteString.copyFromUtf8("test-payload"))
159+
.addAllKeys(keys)
160+
.build())
161+
.setOperation(AccumulatorOuterClass.AccumulatorRequest.WindowOperation
162+
.newBuilder()
163+
.setEvent(AccumulatorOuterClass.AccumulatorRequest.WindowOperation.Event.APPEND)
164+
.setKeyedWindow(openWindow)
165+
.build())
166+
.build();
167+
168+
// CLOSE carries a distinct window that must be echoed verbatim in the EOF response.
169+
AccumulatorOuterClass.KeyedWindow closeWindow = AccumulatorOuterClass.KeyedWindow
170+
.newBuilder()
171+
.setStart(Timestamp.newBuilder().setSeconds(1000).build())
172+
.setEnd(Timestamp.newBuilder().setSeconds(2000).build())
173+
.setSlot("slot-7")
174+
.addAllKeys(keys)
175+
.build();
176+
AccumulatorOuterClass.AccumulatorRequest closeRequest = AccumulatorOuterClass.AccumulatorRequest
177+
.newBuilder()
178+
.setOperation(AccumulatorOuterClass.AccumulatorRequest.WindowOperation
179+
.newBuilder()
180+
.setEvent(AccumulatorOuterClass.AccumulatorRequest.WindowOperation.Event.CLOSE)
181+
.setKeyedWindow(closeWindow)
182+
.build())
183+
.build();
184+
185+
// 2 data responses + 1 EOF response.
186+
AccumulatorStreamObserver responseObserver = new AccumulatorStreamObserver(3);
187+
188+
var stub = AccumulatorGrpc.newStub(inProcessChannel);
189+
var requestStreamObserver = stub.accumulateFn(responseObserver);
190+
191+
requestStreamObserver.onNext(openRequest);
192+
requestStreamObserver.onNext(appendRequest);
193+
requestStreamObserver.onNext(closeRequest);
194+
requestStreamObserver.onCompleted();
195+
196+
try {
197+
responseObserver.done.get();
198+
} catch (InterruptedException | ExecutionException e) {
199+
fail("Error while waiting for response" + e.getMessage());
200+
}
201+
202+
List<AccumulatorOuterClass.AccumulatorResponse> responses = responseObserver.getResponses();
203+
assertEquals(3, responses.size());
204+
205+
AccumulatorOuterClass.AccumulatorResponse eof = null;
206+
for (AccumulatorOuterClass.AccumulatorResponse response : responses) {
207+
if (response.getEOF()) {
208+
eof = response;
209+
}
210+
}
211+
assertEquals(1000, eof.getWindow().getStart().getSeconds());
212+
assertEquals(2000, eof.getWindow().getEnd().getSeconds());
213+
assertEquals("slot-7", eof.getWindow().getSlot());
214+
assertEquals(keys, eof.getWindow().getKeysList());
215+
}
216+
127217
private static class TestAccumFn extends Accumulator {
128218
@Override
129219
public void processMessage(Datum datum, OutputStreamObserver outputStream) {

0 commit comments

Comments
 (0)