|
16 | 16 | */ |
17 | 17 | package org.apache.zeppelin.interpreter; |
18 | 18 |
|
| 19 | +import static org.awaitility.Awaitility.await; |
19 | 20 | import static org.junit.jupiter.api.Assertions.assertFalse; |
20 | 21 | import static org.junit.jupiter.api.Assertions.assertThrows; |
21 | 22 | import static org.junit.jupiter.api.Assertions.assertTrue; |
|
45 | 46 | import java.util.concurrent.Executors; |
46 | 47 | import java.util.concurrent.Future; |
47 | 48 | import java.util.concurrent.TimeUnit; |
48 | | -import java.util.concurrent.TimeoutException; |
49 | 49 | import java.util.concurrent.atomic.AtomicBoolean; |
50 | 50 | import java.util.concurrent.atomic.AtomicReference; |
51 | 51 | import org.apache.zeppelin.conf.ZeppelinConfiguration; |
@@ -174,26 +174,67 @@ void sameNoteBoundaryWaitsForInFlightAppend(String operation) throws Exception { |
174 | 174 | RemoteInterpreterEventServer server = serverWithListener(listener); |
175 | 175 | CountDownLatch entered = new CountDownLatch(1); |
176 | 176 | CountDownLatch release = new CountDownLatch(1); |
177 | | - CountDownLatch started = new CountDownLatch(1); |
| 177 | + AtomicReference<Thread> callerThread = new AtomicReference<>(); |
| 178 | + AtomicBoolean appendCallbackActive = new AtomicBoolean(); |
| 179 | + AtomicBoolean appendCompleted = new AtomicBoolean(); |
| 180 | + AtomicBoolean boundaryOverlappedAppend = new AtomicBoolean(); |
| 181 | + AtomicBoolean rpcReturnedEarly = new AtomicBoolean(); |
| 182 | + |
178 | 183 | doAnswer(invocation -> { |
| 184 | + appendCallbackActive.set(true); |
179 | 185 | entered.countDown(); |
180 | | - assertTrue(release.await(5, TimeUnit.SECONDS)); |
| 186 | + try { |
| 187 | + assertTrue(release.await(5, TimeUnit.SECONDS)); |
| 188 | + } finally { |
| 189 | + appendCallbackActive.set(false); |
| 190 | + appendCompleted.set(true); |
| 191 | + } |
181 | 192 | return null; |
182 | 193 | }).when(listener).onParagraphOutputAppend("note", "first", 0, null, "old"); |
| 194 | + doAnswer(call -> { |
| 195 | + if (appendCallbackActive.get() || !appendCompleted.get()) { |
| 196 | + boundaryOverlappedAppend.set(true); |
| 197 | + } |
| 198 | + return null; |
| 199 | + }).when(listener).onParagraphOutputUpdated( |
| 200 | + "note", "second", 0, null, InterpreterResult.Type.TEXT, "replacement"); |
| 201 | + doAnswer(call -> { |
| 202 | + if (appendCallbackActive.get() || !appendCompleted.get()) { |
| 203 | + boundaryOverlappedAppend.set(true); |
| 204 | + } |
| 205 | + return null; |
| 206 | + }).when(listener).onParagraphOutputClear("note", "second", null); |
| 207 | + doAnswer(call -> { |
| 208 | + if (appendCallbackActive.get() || !appendCompleted.get()) { |
| 209 | + boundaryOverlappedAppend.set(true); |
| 210 | + } |
| 211 | + return null; |
| 212 | + }).when(listener).checkpointOutput("note", "second"); |
| 213 | + |
183 | 214 | ExecutorService callers = Executors.newSingleThreadExecutor(); |
184 | 215 | try { |
185 | 216 | server.appendOutput(new OutputAppendEvent("note", "first", 0, "old", null, null)); |
186 | 217 | dispatcherOf(server).flush(); |
187 | 218 | assertTrue(entered.await(5, TimeUnit.SECONDS)); |
188 | 219 | Future<?> boundary = callers.submit(() -> { |
189 | | - started.countDown(); |
| 220 | + callerThread.set(Thread.currentThread()); |
190 | 221 | callBoundary(server, "note", "second", operation); |
| 222 | + rpcReturnedEarly.set(!appendCompleted.get()); |
191 | 223 | return null; |
192 | 224 | }); |
193 | | - assertTrue(started.await(5, TimeUnit.SECONDS)); |
194 | | - assertThrows(TimeoutException.class, () -> boundary.get(100, TimeUnit.MILLISECONDS)); |
| 225 | + await().atMost(5, TimeUnit.SECONDS).until(() -> boundary.isDone() |
| 226 | + || (callerThread.get() != null |
| 227 | + && callerThread.get().getState() == Thread.State.WAITING)); |
| 228 | + |
| 229 | + assertFalse(appendCompleted.get()); |
| 230 | + assertFalse(boundary.isDone()); |
| 231 | + |
195 | 232 | release.countDown(); |
196 | 233 | boundary.get(5, TimeUnit.SECONDS); |
| 234 | + |
| 235 | + assertFalse(rpcReturnedEarly.get(), "RPC must wait for the earlier append"); |
| 236 | + assertFalse(boundaryOverlappedAppend.get(), "Boundary must wait for append completion"); |
| 237 | + |
197 | 238 | InOrder order = inOrder(listener); |
198 | 239 | order.verify(listener).onParagraphOutputAppend("note", "first", 0, null, "old"); |
199 | 240 | if ("UPDATE".equals(operation)) { |
@@ -259,47 +300,73 @@ void laterAppendCannotOvertakeBoundaryCallback(String operation) throws Exceptio |
259 | 300 | RemoteInterpreterEventServer server = serverWithListener(listener); |
260 | 301 | CountDownLatch entered = new CountDownLatch(1); |
261 | 302 | CountDownLatch release = new CountDownLatch(1); |
| 303 | + AtomicBoolean boundaryCallbackActive = new AtomicBoolean(); |
262 | 304 | AtomicBoolean boundaryCompleted = new AtomicBoolean(); |
| 305 | + AtomicBoolean rpcReturnedEarly = new AtomicBoolean(); |
263 | 306 | AtomicBoolean appendOverlappedBoundary = new AtomicBoolean(); |
| 307 | + |
264 | 308 | doAnswer(call -> { |
| 309 | + boundaryCallbackActive.set(true); |
265 | 310 | entered.countDown(); |
266 | | - assertTrue(release.await(5, TimeUnit.SECONDS)); |
| 311 | + try { |
| 312 | + assertTrue(release.await(5, TimeUnit.SECONDS)); |
| 313 | + } finally { |
| 314 | + boundaryCallbackActive.set(false); |
| 315 | + } |
267 | 316 | return null; |
268 | 317 | }).when(listener).onParagraphOutputClear("note", "para", null); |
269 | 318 | doAnswer(call -> { |
| 319 | + boundaryCallbackActive.set(true); |
270 | 320 | entered.countDown(); |
271 | | - assertTrue(release.await(5, TimeUnit.SECONDS)); |
272 | | - boundaryCompleted.set(true); |
| 321 | + try { |
| 322 | + assertTrue(release.await(5, TimeUnit.SECONDS)); |
| 323 | + } finally { |
| 324 | + boundaryCallbackActive.set(false); |
| 325 | + boundaryCompleted.set(true); |
| 326 | + } |
273 | 327 | return null; |
274 | 328 | }).when(listener).checkpointOutput("note", "para"); |
275 | 329 | doAnswer(call -> { |
| 330 | + boundaryCallbackActive.set(true); |
276 | 331 | entered.countDown(); |
277 | | - assertTrue(release.await(5, TimeUnit.SECONDS)); |
278 | | - boundaryCompleted.set(true); |
| 332 | + try { |
| 333 | + assertTrue(release.await(5, TimeUnit.SECONDS)); |
| 334 | + } finally { |
| 335 | + boundaryCallbackActive.set(false); |
| 336 | + boundaryCompleted.set(true); |
| 337 | + } |
279 | 338 | return null; |
280 | 339 | }).when(listener).onParagraphOutputUpdated("note", "para", 0, null, InterpreterResult.Type.TEXT, "replacement"); |
281 | 340 | doAnswer(call -> { |
282 | | - if (!boundaryCompleted.get()) { |
| 341 | + if (boundaryCallbackActive.get() || !boundaryCompleted.get()) { |
283 | 342 | appendOverlappedBoundary.set(true); |
284 | 343 | } |
285 | 344 | return null; |
286 | 345 | }).when(listener).onParagraphOutputAppend("note", "para", 0, null, "later"); |
| 346 | + |
287 | 347 | ExecutorService callers = Executors.newSingleThreadExecutor(); |
288 | 348 | try { |
289 | 349 | Future<?> boundary = callers.submit(() -> { |
290 | 350 | callBoundary(server, "note", "para", operation); |
| 351 | + rpcReturnedEarly.set(!boundaryCompleted.get()); |
291 | 352 | return null; |
292 | 353 | }); |
293 | 354 | assertTrue(entered.await(5, TimeUnit.SECONDS)); |
294 | 355 | server.appendOutput(new OutputAppendEvent("note", "para", 0, "later", null, null)); |
295 | 356 | dispatcherOf(server).flush(); |
| 357 | + |
296 | 358 | verify(listener, never()).onParagraphOutputAppend("note", "para", 0, null, "later"); |
297 | | - assertThrows(TimeoutException.class, () -> boundary.get(100, TimeUnit.MILLISECONDS)); |
| 359 | + assertFalse(boundaryCompleted.get()); |
| 360 | + assertFalse(boundary.isDone()); |
| 361 | + |
298 | 362 | release.countDown(); |
299 | 363 | boundary.get(5, TimeUnit.SECONDS); |
300 | 364 | server.checkpointOutput("note", "drained"); |
| 365 | + |
301 | 366 | assertTrue(boundaryCompleted.get()); |
| 367 | + assertFalse(rpcReturnedEarly.get(), "RPC must wait for boundary completion"); |
302 | 368 | assertFalse(appendOverlappedBoundary.get(), "Append must wait for the entire boundary"); |
| 369 | + |
303 | 370 | InOrder order = inOrder(listener); |
304 | 371 | if ("UPDATE".equals(operation)) { |
305 | 372 | order.verify(listener).onParagraphOutputUpdated("note", "para", 0, null, InterpreterResult.Type.TEXT, "replacement"); |
|
0 commit comments