Skip to content

Commit b8d8dd9

Browse files
committed
documenting algorithm
Signed-off-by: Attila Mészáros <a_meszaros@apple.com>
1 parent d674e30 commit b8d8dd9

3 files changed

Lines changed: 94 additions & 73 deletions

File tree

operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/informer/EventFilterSupport.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -91,7 +91,7 @@ public synchronized void addToOwnResourceVersions(ResourceID resourceId, String
9191
var window = eventFilterWindows.get(resourceId);
9292
if (window != null) {
9393
log.debug("Recording own resourceVersion. id={}, rv={}", resourceId, resourceVersion);
94-
window.addToOwnResourceVersions(resourceVersion);
94+
window.addToOwnUpdateVersions(resourceVersion);
9595
} else {
9696
log.debug(
9797
"addToOwnResourceVersions: no active window for id={}, rv={} (skipped)",

operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/source/informer/EventFilterWindow.java

Lines changed: 58 additions & 35 deletions
Original file line numberDiff line numberDiff line change
@@ -58,7 +58,7 @@ class EventFilterWindow {
5858
private static final Logger log = LoggerFactory.getLogger(EventFilterWindow.class);
5959

6060
private final SortedMap<Long, ExtendedResourceEvent> relatedEvents = new TreeMap<>();
61-
private final SortedSet<Long> ownResourceVersions = new TreeSet<>();
61+
private final SortedSet<Long> ownUpdateVersions = new TreeSet<>();
6262
private boolean reListOnGoing;
6363
private int activeUpdates = 0;
6464

@@ -82,25 +82,30 @@ public synchronized Optional<ExtendedResourceEvent> check() {
8282
private String snapshotState() {
8383
return String.format(
8484
"relatedEvents=%s, ownResourceVersions=%s, activeUpdates=%d, reListOnGoing=%s",
85-
relatedEvents.keySet(), ownResourceVersions, activeUpdates, reListOnGoing);
85+
relatedEvents.keySet(), ownUpdateVersions, activeUpdates, reListOnGoing);
8686
}
8787

8888
private Optional<ExtendedResourceEvent> doCheck() {
89+
// if we don't have related events we have nothing to mit
8990
if (relatedEvents.isEmpty()) {
9091
return Optional.empty();
9192
}
92-
if (activeUpdates == 0 && ownResourceVersions.isEmpty()) {
93-
return eventForRangeAndClear(relatedEvents, ownResourceVersions);
93+
// cleanup events which are not related to our updates
94+
if (activeUpdates == 0 && ownUpdateVersions.isEmpty()) {
95+
return eventForRangeAndClear(relatedEvents, ownUpdateVersions);
9496
}
95-
if (ownResourceVersions.isEmpty()
97+
// this is a special case that if we receive a delete event we
98+
// early clean it up, since we don't do filtering for deletes
99+
if (ownUpdateVersions.isEmpty()
96100
&& getFirstRelatedEvent().getAction().equals(ResourceAction.DELETED)) {
97-
return eventForRangeAndClear(relatedEvents, ownResourceVersions);
101+
return eventForRangeAndClear(relatedEvents, ownUpdateVersions);
98102
}
99-
100103
var lastEventVersion = relatedEvents.lastKey();
101104
var numberOwnUpdatesSelected = 0;
102105
long lastOwnVersion = -1;
103-
for (long ownVersion : ownResourceVersions) {
106+
// we find the last own update version for which we have event for
107+
// so those are the once we are going to clear our in this execution
108+
for (long ownVersion : ownUpdateVersions) {
104109
if (ownVersion <= lastEventVersion) {
105110
numberOwnUpdatesSelected++;
106111
lastOwnVersion = ownVersion;
@@ -109,61 +114,83 @@ && getFirstRelatedEvent().getAction().equals(ResourceAction.DELETED)) {
109114
}
110115
}
111116
if (numberOwnUpdatesSelected > 0) {
112-
if (numberOwnUpdatesSelected == ownResourceVersions.size() && activeUpdates == 0) {
113-
return eventForRangeAndClear(relatedEvents, ownResourceVersions);
117+
// If we selected all own update versions we process the whole range.
118+
// We check also if there is no active update, since if there still is
119+
// an event might have come which is newer than own version what it is for the ongoing update.
120+
// So If we have own version [1,3] and events [1,2,3,4] and active updates = 0
121+
// we select all events [1,2,3,4] because for the active
122+
// update we might add own version 4.
123+
if (numberOwnUpdatesSelected == ownUpdateVersions.size() && activeUpdates == 0) {
124+
return eventForRangeAndClear(relatedEvents, ownUpdateVersions);
114125
} else {
115-
if (numberOwnUpdatesSelected < ownResourceVersions.size()) {
126+
// if we select only a subset of own updates, we select related events
127+
// up to the next own version (what is not selected).
128+
// So If we have own updates version [1,3,5] and events [1,2,3,4]
129+
// we select all those events (also 4) that happened before own version 5
130+
// for which we don't have event yet.
131+
if (numberOwnUpdatesSelected < ownUpdateVersions.size()) {
116132
return eventForRangeAndClear(
117-
relatedEvents.headMap(ownResourceVersions.tailSet(lastOwnVersion + 1).first()),
118-
ownResourceVersions.headSet(lastOwnVersion + 1));
133+
relatedEvents.headMap(ownUpdateVersions.tailSet(lastOwnVersion + 1).first()),
134+
ownUpdateVersions.headSet(lastOwnVersion + 1));
119135
} else
136+
// this is essentially when we numberOwnUpdatesSelected == ownUpdateVersions.size() but
137+
// with active update > 0. In that case we:
138+
// So If we have own version [1,3] and events [1,2,3,4]
139+
// we select only events [1,2,3] (so no 4), because for the active
140+
// update we might add own version 4.
120141
return eventForRangeAndClear(
121142
relatedEvents.headMap(lastOwnVersion + 1),
122-
ownResourceVersions.headSet(lastOwnVersion + 1));
143+
ownUpdateVersions.headSet(lastOwnVersion + 1));
123144
}
124145
}
125146
return Optional.empty();
126147
}
127148

128-
// it has responsibility to clear those ranges and emit event if needed
149+
// calculates and clears events and own resources for a sorted range of events and own resources
129150
Optional<ExtendedResourceEvent> eventForRangeAndClear(
130151
SortedMap<Long, ExtendedResourceEvent> events, SortedSet<Long> ownResourceVersions) {
152+
131153
if (events.isEmpty()) {
132154
return Optional.empty();
133155
}
156+
157+
var lastEvent = getLastRelatedEvent(events);
158+
if (lastEvent.getAction() == ResourceAction.DELETED) {
159+
events.clear();
160+
ownResourceVersions.clear();
161+
return Optional.of(lastEvent);
162+
}
163+
164+
// if any of the events is part of re-list (including first delete) we detect it
134165
var isAnyEventFromReList =
135166
events.values().stream().anyMatch(ExtendedResourceEvent::isPartOfReList);
136167

137168
var first = getFirstRelatedEvent(events);
169+
// if delete event is first in the row and more events we can discard that
170+
// since won't play role in synthesized (synt) event.
138171
if (events.size() > 1 && first.getAction() == ResourceAction.DELETED) {
139172
events.remove(events.firstKey());
140173
first = getFirstRelatedEvent(events);
141174
}
142175

176+
// if all updates are related to own updates we don't return event.
177+
//
143178
if (events.keySet().equals(ownResourceVersions) && !isAnyEventFromReList) {
144-
ExtendedResourceEvent res = null;
145-
var lastEvent = getLastRelatedEvent(events);
146-
if (lastEvent.getAction() == ResourceAction.DELETED) {
147-
res = lastEvent;
148-
}
149179
events.clear();
150180
ownResourceVersions.clear();
151-
return Optional.ofNullable(res);
181+
return Optional.empty();
152182
}
153183

184+
// if only one event we return that
154185
if (events.size() == 1) {
155186
ownResourceVersions.clear();
156187
var res = Optional.of(events.values().iterator().next());
157188
events.clear();
158189
return res;
159190
}
160-
var lastEvent = getLastRelatedEvent(events);
161-
if (lastEvent.getAction() == ResourceAction.DELETED) {
162-
events.clear();
163-
ownResourceVersions.clear();
164-
return Optional.of(lastEvent);
165-
}
166191

192+
// if none above we create a synt event that contains from the oldest know resource
193+
// to the newest one. This is important to filters see the whole range
167194
var res =
168195
Optional.of(
169196
new ExtendedResourceEvent(
@@ -191,19 +218,15 @@ private ExtendedResourceEvent getLastRelatedEvent(SortedMap<Long, ExtendedResour
191218
return subMap.get(subMap.lastKey());
192219
}
193220

194-
private ExtendedResourceEvent getLastRelatedEvent() {
195-
return getLastRelatedEvent(relatedEvents);
196-
}
197-
198221
public synchronized boolean canBeRemoved() {
199-
if (activeUpdates == 0 && ownResourceVersions.isEmpty() && relatedEvents.isEmpty()) {
222+
if (activeUpdates == 0 && ownUpdateVersions.isEmpty() && relatedEvents.isEmpty()) {
200223
return true;
201224
}
202225
return false;
203226
}
204227

205-
public synchronized void addToOwnResourceVersions(String resourceVersion) {
206-
ownResourceVersions.add(Long.parseLong(resourceVersion));
228+
public synchronized void addToOwnUpdateVersions(String resourceVersion) {
229+
ownUpdateVersions.add(Long.parseLong(resourceVersion));
207230
}
208231

209232
public synchronized void addRelatedEvent(ExtendedResourceEvent event) {
@@ -237,6 +260,6 @@ synchronized SortedMap<Long, ExtendedResourceEvent> getRelatedEvents() {
237260
}
238261

239262
synchronized SortedSet<Long> getOwnResourceVersions() {
240-
return ownResourceVersions;
263+
return ownUpdateVersions;
241264
}
242265
}

0 commit comments

Comments
 (0)