Skip to content

Commit dc69f56

Browse files
Fixing bug that dropped punctuations for Grouped Afa SingleEvent pipe when IsSyncTimeSimultaneityFree is true (#142)
1 parent c9d69c4 commit dc69f56

2 files changed

Lines changed: 16 additions & 8 deletions

File tree

Sources/Core/Microsoft.StreamProcessing/Operators/Afa/CompiledGroupedAfaPipe_SingleEvent.cs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -320,10 +320,10 @@ public override unsafe void OnNext(StreamMessage<TKey, TPayload> batch)
320320
if (this.IsDeterministic) break; // We are guaranteed to have only one start state
321321
}
322322
}
323-
else if (batch.vother.col[i] < 0 && !this.IsSyncTimeSimultaneityFree)
323+
else if (batch.vother.col[i] < 0)
324324
{
325325
long synctime = src_vsync[i];
326-
if (synctime > this.lastSyncTime) // move time forward
326+
if (!this.IsSyncTimeSimultaneityFree && synctime > this.lastSyncTime) // move time forward
327327
{
328328
this.seenEvent.Clear();
329329

Sources/Test/SimpleTesting/Streamables/AfaTests.cs

Lines changed: 14 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -844,16 +844,21 @@ public void GroupedAfa_IsSyncTimeSimultaneityFree()
844844
{
845845
var source = new StreamEvent<Tuple<string, int>>[]
846846
{
847-
StreamEvent.CreateStart(0, new Tuple<string, int>("A", 1)),
847+
StreamEvent.CreateStart(0, new Tuple<string, int>("A", 1)),
848848
StreamEvent.CreateStart(1, new Tuple<string, int>("A", 2)),
849849
StreamEvent.CreateStart(1, new Tuple<string, int>("B", 2)),
850-
StreamEvent.CreateStart(3, new Tuple<string, int>("A", 1)),
851-
StreamEvent.CreateStart(4, new Tuple<string, int>("B", 1)),
850+
StreamEvent.CreateStart(3, new Tuple<string, int>("A", 1)),
851+
852+
StreamEvent.CreatePunctuation<Tuple<string, int>>(4),
853+
854+
StreamEvent.CreateStart(4, new Tuple<string, int>("B", 1)),
852855
StreamEvent.CreateStart(4, new Tuple<string, int>("B", 2)),
853-
StreamEvent.CreateStart(5, new Tuple<string, int>("B", 1)),
854-
StreamEvent.CreateStart(5, new Tuple<string, int>("C", 1)),
856+
StreamEvent.CreateStart(5, new Tuple<string, int>("B", 1)),
857+
StreamEvent.CreateStart(5, new Tuple<string, int>("C", 1)),
855858
StreamEvent.CreateStart(6, new Tuple<string, int>("B", 2)),
856859
StreamEvent.CreateStart(7, new Tuple<string, int>("A", 2)),
860+
861+
StreamEvent.CreatePunctuation<Tuple<string, int>>(7),
857862
}.ToObservable()
858863
.ToStreamable()
859864
.AlterEventDuration(10);
@@ -881,18 +886,21 @@ public void GroupedAfa_IsSyncTimeSimultaneityFree()
881886

882887
var result = afa_compiled
883888
.ToStreamEventObservable()
884-
.Where(evt => evt.IsData)
885889
.ToEnumerable()
886890
.ToArray();
887891
var expected = new StreamEvent<Tuple<string, int>>[]
888892
{
889893
StreamEvent.CreateInterval(1, 11, new Tuple<string, int>("AB", 2)),
894+
StreamEvent.CreatePunctuation<Tuple<string, int>>(4),
890895
StreamEvent.CreateInterval(4, 13, new Tuple<string, int>("AB", 1)),
891896
StreamEvent.CreateInterval(4, 10, new Tuple<string, int>("AAB", 1)),
892897
StreamEvent.CreateInterval(4, 11, new Tuple<string, int>("ABB", 2)),
893898
StreamEvent.CreateInterval(5, 13, new Tuple<string, int>("ABB", 1)),
894899
StreamEvent.CreateInterval(5, 10, new Tuple<string, int>("AABB", 1)),
895900
StreamEvent.CreateInterval(6, 11, new Tuple<string, int>("ABBB", 2)),
901+
StreamEvent.CreatePunctuation<Tuple<string, int>>(7),
902+
903+
StreamEvent.CreatePunctuation<Tuple<string, int>>(StreamEvent.InfinitySyncTime),
896904
};
897905
Assert.IsTrue(result.SequenceEqual(expected));
898906
}

0 commit comments

Comments
 (0)