diff --git a/.github/workflows/frontend.yml b/.github/workflows/frontend.yml index 08af6bf6998f..dedfadbb7ba0 100644 --- a/.github/workflows/frontend.yml +++ b/.github/workflows/frontend.yml @@ -49,9 +49,12 @@ jobs: env: # Use VFS storage instead of Git to avoid Git-related issues in CI ZEPPELIN_NOTEBOOK_STORAGE: org.apache.zeppelin.notebook.repo.VFSNotebookRepo + ZEPPELIN_JOBMANAGER_ENABLE: 'true' + ZEPPELIN_EVENTBUS_ENABLED: ${{ matrix.eventbus == 'eventbus' }} strategy: matrix: mode: [anonymous, auth] + eventbus: [legacy, eventbus] python: [ 3.9 ] steps: - name: Checkout @@ -120,7 +123,7 @@ jobs: echo "Created test notebook directory: $ZEPPELIN_E2E_TEST_NOTEBOOK_DIR" - name: Run headless E2E test with Maven # Classic UI e2e runs only on the anonymous leg, like the legacy Protractor suite - run: xvfb-run --auto-servernum --server-args="-screen 0 1024x768x24" ./mvnw verify -pl zeppelin-web-angular -Pweb-e2e -Dweb.e2e.classic.disabled=${{ matrix.mode != 'anonymous' }} ${MAVEN_ARGS} + run: xvfb-run --auto-servernum --server-args="-screen 0 1024x768x24" ./mvnw verify -pl zeppelin-web-angular -Pweb-e2e -Dweb.e2e.classic.disabled=${{ matrix.mode != 'anonymous' || matrix.eventbus != 'legacy' }} ${MAVEN_ARGS} - name: Run revision isolation E2E test with Git storage env: CI: 'true' @@ -145,7 +148,7 @@ jobs: uses: actions/upload-artifact@v6 if: always() with: - name: playwright-report-${{ matrix.mode }} + name: playwright-report-${{ matrix.mode }}-${{ matrix.eventbus }} path: | zeppelin-web-angular/playwright-report/ zeppelin-web-angular/playwright-report-classic/ diff --git a/conf/zeppelin-site.xml.template b/conf/zeppelin-site.xml.template index 4107a6799cf3..72960ca0ffef 100755 --- a/conf/zeppelin-site.xml.template +++ b/conf/zeppelin-site.xml.template @@ -840,4 +840,10 @@ fields to be excluded from being saved in note files, with Paragraph prefix mean the fields in Paragraph, e.g. Paragraph.results + + zeppelin.eventbus.enabled + false + Enables the new event-driven architecture using an in-process EventBus + + diff --git a/zeppelin-server/pom.xml b/zeppelin-server/pom.xml index 505c5ed75bbd..bd74a408007e 100644 --- a/zeppelin-server/pom.xml +++ b/zeppelin-server/pom.xml @@ -97,6 +97,12 @@ jakarta.inject-api + + io.reactivex.rxjava3 + rxjava + 3.1.12 + + commons-io commons-io diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/conf/ZeppelinConfiguration.java b/zeppelin-server/src/main/java/org/apache/zeppelin/conf/ZeppelinConfiguration.java index 01ad388ebac4..ca15098a0f54 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/conf/ZeppelinConfiguration.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/conf/ZeppelinConfiguration.java @@ -941,6 +941,8 @@ public boolean isPrometheusMetricEnabled() { return getBoolean(ConfVars.ZEPPELIN_METRIC_ENABLE_PROMETHEUS); } + public boolean isEventBusEnabled() { return getBoolean(ConfVars.ZEPPELIN_EVENTBUS_ENABLED); } + public DEFAULT_UI getDefaultUi() { return DEFAULT_UI.valueOf(getString(ConfVars.ZEPPELIN_DEFAULT_UI).toUpperCase()); } @@ -1187,7 +1189,8 @@ public enum ConfVars { ZEPPELIN_SPARK_ONLY_YARN_CLUSTER("zeppelin.spark.only_yarn_cluster", false), ZEPPELIN_SESSION_CHECK_INTERVAL("zeppelin.session.check_interval", 60 * 10 * 1000), ZEPPELIN_NOTE_CACHE_THRESHOLD("zeppelin.note.cache.threshold", 50), - ZEPPELIN_NOTE_FILE_EXCLUDE_FIELDS("zeppelin.note.file.exclude.fields", ""); + ZEPPELIN_NOTE_FILE_EXCLUDE_FIELDS("zeppelin.note.file.exclude.fields", ""), + ZEPPELIN_EVENTBUS_ENABLED("zeppelin.eventbus.enabled", false); private String varName; private Class varClass; diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/eventbus/EventBus.java b/zeppelin-server/src/main/java/org/apache/zeppelin/eventbus/EventBus.java new file mode 100644 index 000000000000..cff47114adfb --- /dev/null +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/eventbus/EventBus.java @@ -0,0 +1,27 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.zeppelin.eventbus; + +import io.reactivex.rxjava3.core.Observable; + +public interface EventBus { + + void post(ZeppelinEvent event); + + Observable observe(Class eventType); +} diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/eventbus/NoOpEventBus.java b/zeppelin-server/src/main/java/org/apache/zeppelin/eventbus/NoOpEventBus.java new file mode 100644 index 000000000000..cf9eb180a087 --- /dev/null +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/eventbus/NoOpEventBus.java @@ -0,0 +1,43 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.zeppelin.eventbus; + +import io.reactivex.rxjava3.core.Observable; +import jakarta.inject.Inject; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +public class NoOpEventBus implements EventBus { + + private static final Logger LOGGER = LoggerFactory.getLogger(NoOpEventBus.class); + + @Inject + public NoOpEventBus() { + LOGGER.info("Starting NoOpEventBus"); + } + + @Override + public void post(ZeppelinEvent event) { + LOGGER.debug("Posting event: {}", event.getClass().getName()); + } + + @Override + public Observable observe(Class eventType) { + return Observable.empty(); + } +} diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/eventbus/NoteEvent.java b/zeppelin-server/src/main/java/org/apache/zeppelin/eventbus/NoteEvent.java new file mode 100644 index 000000000000..2b6b806a1499 --- /dev/null +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/eventbus/NoteEvent.java @@ -0,0 +1,28 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.zeppelin.eventbus; + +import org.apache.zeppelin.notebook.Note; +import org.apache.zeppelin.user.AuthenticationInfo; + +public interface NoteEvent extends ZeppelinEvent { + + Note getNote(); + + AuthenticationInfo getSubject(); +} diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/eventbus/NoteRemovedEvent.java b/zeppelin-server/src/main/java/org/apache/zeppelin/eventbus/NoteRemovedEvent.java new file mode 100644 index 000000000000..de7ca8874953 --- /dev/null +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/eventbus/NoteRemovedEvent.java @@ -0,0 +1,43 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.zeppelin.eventbus; + +import org.apache.zeppelin.notebook.Note; +import org.apache.zeppelin.user.AuthenticationInfo; + +public class NoteRemovedEvent implements NoteEvent { + + private final Note note; + + private final AuthenticationInfo subject; + + public NoteRemovedEvent(Note note, AuthenticationInfo subject) { + this.note = note; + this.subject = subject; + } + + @Override + public Note getNote() { + return this.note; + } + + @Override + public AuthenticationInfo getSubject() { + return subject; + } +} diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/eventbus/ZeppelinEvent.java b/zeppelin-server/src/main/java/org/apache/zeppelin/eventbus/ZeppelinEvent.java new file mode 100644 index 000000000000..2f181ef7e8d6 --- /dev/null +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/eventbus/ZeppelinEvent.java @@ -0,0 +1,21 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.zeppelin.eventbus; + +public interface ZeppelinEvent { +} diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/eventbus/ZeppelinEventBus.java b/zeppelin-server/src/main/java/org/apache/zeppelin/eventbus/ZeppelinEventBus.java new file mode 100644 index 000000000000..652609b01451 --- /dev/null +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/eventbus/ZeppelinEventBus.java @@ -0,0 +1,53 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.zeppelin.eventbus; + +import io.reactivex.rxjava3.core.Observable; +import io.reactivex.rxjava3.subjects.PublishSubject; +import io.reactivex.rxjava3.subjects.Subject; +import jakarta.inject.Inject; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +public class ZeppelinEventBus implements EventBus { + + private static final Logger LOGGER = LoggerFactory.getLogger(ZeppelinEventBus.class); + + private final Subject eventBus; + + @Inject + public ZeppelinEventBus() { + LOGGER.info("Starting ZeppelinEventBus"); + + eventBus = PublishSubject.create().toSerialized(); + } + + @Override + public void post(ZeppelinEvent event) { + LOGGER.debug("Posting event: {}", event.getClass().getName()); + + eventBus.onNext(event); + } + + @Override + public Observable observe(Class eventType) { + LOGGER.debug("Observing event: {}", eventType.getName()); + + return eventBus.ofType(eventType); + } +} diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/notebook/Notebook.java b/zeppelin-server/src/main/java/org/apache/zeppelin/notebook/Notebook.java index 10c2abca7792..9805300ab29b 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/notebook/Notebook.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/notebook/Notebook.java @@ -44,6 +44,8 @@ import org.apache.zeppelin.conf.ZeppelinConfiguration.ConfVars; import org.apache.zeppelin.display.AngularObject; import org.apache.zeppelin.display.AngularObjectRegistry; +import org.apache.zeppelin.eventbus.EventBus; +import org.apache.zeppelin.eventbus.NoteRemovedEvent; import org.apache.zeppelin.interpreter.Interpreter; import org.apache.zeppelin.interpreter.InterpreterFactory; import org.apache.zeppelin.interpreter.InterpreterGroup; @@ -86,11 +88,13 @@ public class Notebook { private Credentials credentials; private final List> initConsumers; private ExecutorService initExecutor; + private EventBus eventBus; /** * Main constructor \w manual Dependency Injection * * @throws IOException + * * @throws SchedulerException */ public Notebook( @@ -100,8 +104,8 @@ public Notebook( NoteManager noteManager, InterpreterFactory replFactory, InterpreterSettingManager interpreterSettingManager, - Credentials credentials) - { + Credentials credentials, + EventBus eventBus) { this.zConf = zConf; this.authorizationService = authorizationService; this.noteManager = noteManager; @@ -111,6 +115,7 @@ public Notebook( // TODO(zjffdu) cycle refer, not a good solution this.interpreterSettingManager.setNotebook(this); this.credentials = credentials; + this.eventBus = eventBus; addNotebookEventListener(this.interpreterSettingManager); initConsumers = new LinkedList<>(); } @@ -221,7 +226,8 @@ public Notebook( InterpreterFactory replFactory, InterpreterSettingManager interpreterSettingManager, Credentials credentials, - NoteEventListener noteEventListener) + NoteEventListener noteEventListener, + EventBus eventBus) throws IOException { this( zConf, @@ -230,7 +236,8 @@ public Notebook( noteManager, replFactory, interpreterSettingManager, - credentials); + credentials, + eventBus); if (null != noteEventListener) { addNotebookEventListener(noteEventListener); } @@ -432,7 +439,11 @@ private void removeNote(Note note, AuthenticationInfo subject) throws IOExceptio note.setRemoved(true); noteManager.removeNote(note.getId(), subject); authorizationService.removeNoteAuth(note.getId()); + fireNoteRemoveEvent(note, subject); + if (zConf.isEventBusEnabled()) { + eventBus.post(new NoteRemovedEvent(note, subject)); + } } public void removeCorruptedNote(String noteId, AuthenticationInfo subject) throws IOException { @@ -564,6 +575,9 @@ public void removeFolder(String folderPath, AuthenticationInfo subject) throws I note.setRemoved(true); authorizationService.removeNoteAuth(note.getId()); fireNoteRemoveEvent(note, subject); + if (zConf.isEventBusEnabled()) { + eventBus.post(new NoteRemovedEvent(note, subject)); + } } return null; }); @@ -843,6 +857,7 @@ private void fireNoteUpdateEvent(Note note, AuthenticationInfo subject) { } } + private void fireNoteRemoveEvent(Note note, AuthenticationInfo subject) { for (NoteEventListener listener : noteEventListeners) { listener.onNoteRemove(note, subject); diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/server/ZeppelinServer.java b/zeppelin-server/src/main/java/org/apache/zeppelin/server/ZeppelinServer.java index 6dece8a13d33..07d46b089525 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/server/ZeppelinServer.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/server/ZeppelinServer.java @@ -64,6 +64,9 @@ import org.apache.zeppelin.conf.ZeppelinConfiguration.ConfVars; import org.apache.zeppelin.conf.ZeppelinConfiguration.DEFAULT_UI; import org.apache.zeppelin.display.AngularObjectRegistryListener; +import org.apache.zeppelin.eventbus.EventBus; +import org.apache.zeppelin.eventbus.NoOpEventBus; +import org.apache.zeppelin.eventbus.ZeppelinEventBus; import org.apache.zeppelin.healthcheck.HealthChecks; import org.apache.zeppelin.helium.ApplicationEventListener; import org.apache.zeppelin.helium.Helium; @@ -179,6 +182,11 @@ protected void configure() { bind(storage).to(ConfigStorage.class); bindAsContract(PluginManager.class).in(Singleton.class); bind(GsonNoteParser.class).to(NoteParser.class).in(Singleton.class); + if (zConf.isEventBusEnabled()) { + bind(ZeppelinEventBus.class).to(EventBus.class).in(Singleton.class); + } else { + bind(NoOpEventBus.class).to(EventBus.class).in(Singleton.class); + } bindAsContract(InterpreterFactory.class).in(Singleton.class); bindAsContract(NotebookRepoSync.class).to(NotebookRepo.class).in(Singleton.class); bindAsContract(Helium.class).in(Singleton.class); @@ -362,6 +370,7 @@ public void shutdown(int exitCode) { if (!zConf.isRecoveryEnabled()) { sharedServiceLocator.getService(InterpreterSettingManager.class).close(); } + sharedServiceLocator.getService(NotebookServer.class).close(); sharedServiceLocator.getService(Notebook.class).close(); } } catch (Exception e) { diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java b/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java index 4d31a06558cc..eae23f90408a 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java @@ -42,6 +42,8 @@ import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicReference; + +import io.reactivex.rxjava3.disposables.CompositeDisposable; import jakarta.inject.Inject; import jakarta.inject.Provider; import jakarta.websocket.CloseReason; @@ -64,6 +66,9 @@ import org.apache.zeppelin.display.AngularObjectRegistryListener; import org.apache.zeppelin.display.GUI; import org.apache.zeppelin.display.Input; +import org.apache.zeppelin.eventbus.EventBus; +import org.apache.zeppelin.eventbus.NoteEvent; +import org.apache.zeppelin.eventbus.NoteRemovedEvent; import org.apache.zeppelin.helium.ApplicationEventListener; import org.apache.zeppelin.helium.HeliumPackage; import org.apache.zeppelin.interpreter.InterpreterGroup; @@ -119,7 +124,8 @@ public class NotebookServer implements AngularObjectRegistryListener, RemoteInterpreterProcessListener, ApplicationEventListener, ParagraphJobListener, - NoteEventListener { + NoteEventListener, + AutoCloseable { /** * Job manager service type. @@ -162,6 +168,7 @@ String getKey() { private Provider notebookServiceProvider; private AuthorizationService authorizationService; private Provider jobManagerServiceProvider; + private CompositeDisposable subscriptions; public NotebookServer() { NotebookServer.self.set(this); @@ -173,6 +180,31 @@ public void setZeppelinConfiguration(ZeppelinConfiguration zConf) { this.zConf = zConf; } + @Inject + public void registerEventBus(EventBus eventBus, ZeppelinConfiguration zConf) { + if (!zConf.isEventBusEnabled()) { + LOGGER.debug("ZeppelinEventBus is disabled"); + return; + } + subscriptions = new CompositeDisposable(); + + subscriptions.add(eventBus.observe(NoteEvent.class) + .subscribe(event -> { + try { + handleNoteEvent(event); + } catch (Exception e) { + LOGGER.error("Failed to handle note event: {}", event, e); + } + })); + } + + @Override + public void close() { + if (subscriptions != null && !subscriptions.isDisposed()) { + subscriptions.dispose(); + } + } + @Inject public void setNoteParser(Provider noteParser) { this.noteParser = noteParser; @@ -1967,19 +1999,12 @@ public void onParagraphRemove(Paragraph p) { @Override public void onNoteRemove(Note note, AuthenticationInfo subject) { - try { - broadcastUpdateNoteJobInfo(note, System.currentTimeMillis() - 5000); - } catch (IOException e) { - LOGGER.warn("can not broadcast for job manager: {}", e.getMessage(), e); - } - - try { - getJobManagerService().removeNoteJobInfo(note.getId(), null, - new JobManagerServiceCallback()); - } catch (IOException e) { - LOGGER.warn("can not broadcast for job manager: {}", e.getMessage(), e); + if (zConf.isEventBusEnabled()) { + LOGGER.debug("ZeppelinEventBus is enabled"); + return; } + handleNoteRemove(note); } @Override @@ -2421,4 +2446,28 @@ public void onFailure(Exception ex, ServiceContext context) throws IOException { } } } + + private void handleNoteEvent(NoteEvent event) { + if (event instanceof NoteRemovedEvent) { + Note note = event.getNote(); + handleNoteRemove(note); + } else { + LOGGER.warn("Unknown event type: {}", event.getClass().getName()); + } + } + + private void handleNoteRemove(Note note) { + try { + broadcastUpdateNoteJobInfo(note, System.currentTimeMillis() - 5000); + } catch (IOException e) { + LOGGER.warn("can not broadcast for job manager: {}", e.getMessage(), e); + } + + try { + getJobManagerService().removeNoteJobInfo(note.getId(), null, + new JobManagerServiceCallback()); + } catch (IOException e) { + LOGGER.warn("can not broadcast for job manager: {}", e.getMessage(), e); + } + } } diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/eventbus/RxJavaTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/eventbus/RxJavaTest.java new file mode 100644 index 000000000000..f3d861bb3e63 --- /dev/null +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/eventbus/RxJavaTest.java @@ -0,0 +1,34 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.zeppelin.eventbus; + +import io.reactivex.rxjava3.core.Observable; +import org.junit.jupiter.api.Test; + +class RxJavaTest { + + @Test + void testObservable() { + Observable observable = Observable.just("Hello", "RxJava", "Test"); + + observable.test() + .assertValues("Hello", "RxJava", "Test") + .assertComplete() + .assertNoErrors(); + } +} diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/eventbus/ZeppelinEventBusTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/eventbus/ZeppelinEventBusTest.java new file mode 100644 index 000000000000..509195bfdf03 --- /dev/null +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/eventbus/ZeppelinEventBusTest.java @@ -0,0 +1,96 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.zeppelin.eventbus; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import io.reactivex.rxjava3.disposables.Disposable; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import org.junit.jupiter.api.Test; + +class ZeppelinEventBusTest { + + @Test + void eventBusFlow() throws InterruptedException { + ZeppelinEventBus bus = new ZeppelinEventBus(); + Publisher publisher = new Publisher(bus); + Subscriber subscriber = new Subscriber(bus); + + String payload = "data"; + publisher.createNote(payload); + + assertTrue(subscriber.awaitEvent()); + + List received = subscriber.collection; + assertEquals(1, received.size()); + assertEquals(payload, received.get(0)); + assertTrue(received.contains(payload)); + + subscriber.stopListening(); + } + + private static class MockEvent implements ZeppelinEvent { + private final String payload; + + MockEvent(String payload) { + this.payload = payload; + } + } + + private static class Publisher { + private final ZeppelinEventBus eventBus; + + Publisher(ZeppelinEventBus eventBus) { + this.eventBus = eventBus; + } + + void createNote(String noteId) { + eventBus.post(new MockEvent(noteId)); + } + } + + private static class Subscriber { + private final List collection = new ArrayList<>(); + + private final CountDownLatch eventReceived = new CountDownLatch(1); + + private final Disposable disposable; + + Subscriber(ZeppelinEventBus eventBus) { + this.disposable = eventBus.observe(MockEvent.class) + .subscribe(event -> { + collection.add(event.payload); + eventReceived.countDown(); + }); + } + + boolean awaitEvent() throws InterruptedException { + return eventReceived.await(1, TimeUnit.SECONDS); + } + + void stopListening() { + if (!disposable.isDisposed()) { + disposable.dispose(); + } + } + } +} diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/helium/HeliumApplicationFactoryTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/helium/HeliumApplicationFactoryTest.java index 2f356281ab72..d66662ebfad8 100644 --- a/zeppelin-server/src/test/java/org/apache/zeppelin/helium/HeliumApplicationFactoryTest.java +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/helium/HeliumApplicationFactoryTest.java @@ -70,7 +70,8 @@ public void setUp() throws Exception { new NoteManager(notebookRepo, ZeppelinConfiguration.load()), interpreterFactory, interpreterSettingManager, - new Credentials()); + new Credentials(), + eventBus); heliumAppFactory = new HeliumApplicationFactory(notebook, null); diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/AbstractInterpreterTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/AbstractInterpreterTest.java index 0bfd19014b7c..2f7b6a65cba0 100644 --- a/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/AbstractInterpreterTest.java +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/AbstractInterpreterTest.java @@ -21,6 +21,7 @@ import org.apache.commons.io.FileUtils; import org.apache.zeppelin.conf.ZeppelinConfiguration; import org.apache.zeppelin.display.AngularObjectRegistryListener; +import org.apache.zeppelin.eventbus.ZeppelinEventBus; import org.apache.zeppelin.helium.ApplicationEventListener; import org.apache.zeppelin.interpreter.remote.RemoteInterpreterProcessListener; import org.apache.zeppelin.notebook.AuthorizationService; @@ -65,6 +66,7 @@ public abstract class AbstractInterpreterTest { protected ZeppelinConfiguration zConf; protected ConfigStorage storage; protected PluginManager pluginManager; + protected ZeppelinEventBus eventBus; @BeforeEach public void setUp() throws Exception { @@ -100,12 +102,12 @@ public void setUp() throws Exception { zConf.setProperty(ZeppelinConfiguration.ConfVars.ZEPPELIN_INTERPRETER_GROUP_DEFAULT.getVarName(), "test"); - NotebookRepo notebookRepo = new InMemoryNotebookRepo(); NoteManager noteManager = new NoteManager(notebookRepo, zConf); noteParser = new GsonNoteParser(zConf); storage = ConfigStorage.createConfigStorage(zConf); pluginManager = new PluginManager(zConf); + eventBus = new ZeppelinEventBus(); AuthorizationService authorizationService = new AuthorizationService(noteManager, zConf, storage); @@ -114,7 +116,7 @@ public void setUp() throws Exception { mock(ApplicationEventListener.class), storage, pluginManager); interpreterFactory = new InterpreterFactory(interpreterSettingManager); Credentials credentials = new Credentials(zConf, storage); - notebook = new Notebook(zConf, authorizationService, notebookRepo, noteManager, interpreterFactory, interpreterSettingManager, credentials); + notebook = new Notebook(zConf, authorizationService, notebookRepo, noteManager, interpreterFactory, interpreterSettingManager, credentials, eventBus); interpreterSettingManager.setNotebook(notebook); } diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/notebook/NotebookTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/notebook/NotebookTest.java index f7113944f895..b9031d8474a3 100644 --- a/zeppelin-server/src/test/java/org/apache/zeppelin/notebook/NotebookTest.java +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/notebook/NotebookTest.java @@ -21,6 +21,7 @@ import org.apache.zeppelin.conf.ZeppelinConfiguration; import org.apache.zeppelin.conf.ZeppelinConfiguration.ConfVars; import org.apache.zeppelin.display.AngularObjectRegistry; +import org.apache.zeppelin.eventbus.ZeppelinEventBus; import org.apache.zeppelin.interpreter.AbstractInterpreterTest; import org.apache.zeppelin.interpreter.ExecutionContext; import org.apache.zeppelin.interpreter.InterpreterException; @@ -107,7 +108,8 @@ public void setUp() throws Exception { authorizationService = new AuthorizationService(noteManager, zConf, storage); credentials = new Credentials(zConf, storage); - notebook = new Notebook(zConf, authorizationService, notebookRepo, noteManager, interpreterFactory, interpreterSettingManager, credentials, null); + ZeppelinEventBus eventBus = new ZeppelinEventBus(); + notebook = new Notebook(zConf, authorizationService, notebookRepo, noteManager, interpreterFactory, interpreterSettingManager, credentials, eventBus); notebook.setParagraphJobListener(this); schedulerService = new QuartzSchedulerService(zConf, notebook); notebook.initNotebook(); @@ -1710,7 +1712,7 @@ public void onParagraphStatusChange(Paragraph p, Status status) { } @Test - void testRemoveFolderFiresNoteRemoveEventForEachNote() throws IOException { + void testRemoveFolderFiresNoteRemovedEventForEachNote() throws IOException { final AtomicInteger onNoteRemove = new AtomicInteger(0); notebook.addNotebookEventListener(new NoteEventListener() { @Override diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/notebook/repo/NotebookRepoSyncTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/notebook/repo/NotebookRepoSyncTest.java index 3a90e4f60348..b7cb75729eac 100644 --- a/zeppelin-server/src/test/java/org/apache/zeppelin/notebook/repo/NotebookRepoSyncTest.java +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/notebook/repo/NotebookRepoSyncTest.java @@ -39,6 +39,8 @@ import org.apache.zeppelin.conf.ZeppelinConfiguration; import org.apache.zeppelin.conf.ZeppelinConfiguration.ConfVars; import org.apache.zeppelin.display.AngularObjectRegistryListener; +import org.apache.zeppelin.eventbus.EventBus; +import org.apache.zeppelin.eventbus.ZeppelinEventBus; import org.apache.zeppelin.helium.ApplicationEventListener; import org.apache.zeppelin.interpreter.InterpreterFactory; import org.apache.zeppelin.interpreter.InterpreterSettingManager; @@ -76,6 +78,7 @@ class NotebookRepoSyncTest { private InterpreterFactory factory; private InterpreterSettingManager interpreterSettingManager; private Credentials credentials; + private EventBus eventBus; private AuthenticationInfo anonymous; private NoteManager noteManager; private AuthorizationService authorizationService; @@ -116,7 +119,8 @@ public void setUp() throws Exception { noteManager = new NoteManager(notebookRepoSync, zConf); authorizationService = new AuthorizationService(noteManager, zConf, storage); credentials = new Credentials(zConf, storage); - notebook = new Notebook(zConf, authorizationService, notebookRepoSync, noteManager, factory, interpreterSettingManager, credentials, null); + eventBus = new ZeppelinEventBus(); + notebook = new Notebook(zConf, authorizationService, notebookRepoSync, noteManager, factory, interpreterSettingManager, credentials, eventBus); anonymous = new AuthenticationInfo("anonymous"); } diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/service/NotebookServiceTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/service/NotebookServiceTest.java index 1fb67977bcca..d9c23718e9e3 100644 --- a/zeppelin-server/src/test/java/org/apache/zeppelin/service/NotebookServiceTest.java +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/service/NotebookServiceTest.java @@ -49,6 +49,8 @@ import org.apache.commons.lang3.StringUtils; import org.apache.zeppelin.conf.ZeppelinConfiguration; +import org.apache.zeppelin.eventbus.EventBus; +import org.apache.zeppelin.eventbus.ZeppelinEventBus; import org.apache.zeppelin.interpreter.Interpreter; import org.apache.zeppelin.interpreter.Interpreter.FormType; import org.apache.zeppelin.interpreter.InterpreterFactory; @@ -84,9 +86,6 @@ import org.mockito.ArgumentCaptor; import org.junit.jupiter.api.TestInfo; - -import com.google.gson.Gson; - class NotebookServiceTest { private static NotebookService notebookService; @@ -149,6 +148,7 @@ void setUp(TestInfo testInfo) throws Exception { Credentials credentials = new Credentials(); NoteManager noteManager = new NoteManager(notebookRepo, zConf); authorizationService = new AuthorizationService(noteManager, zConf, storage); + EventBus eventBus = new ZeppelinEventBus(); notebook = new Notebook( zConf, @@ -158,7 +158,7 @@ void setUp(TestInfo testInfo) throws Exception { mockInterpreterFactory, mockInterpreterSettingManager, credentials, - null); + eventBus); searchService = new LuceneSearch(zConf, notebook); QuartzSchedulerService schedulerService = new QuartzSchedulerService(zConf, notebook); notebook.initNotebook(); diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/socket/NotebookServerTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/socket/NotebookServerTest.java index d288851fbff4..5aeae7b597b4 100644 --- a/zeppelin-server/src/test/java/org/apache/zeppelin/socket/NotebookServerTest.java +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/socket/NotebookServerTest.java @@ -54,6 +54,9 @@ import org.apache.zeppelin.conf.ZeppelinConfiguration; import org.apache.zeppelin.display.AngularObject; import org.apache.zeppelin.display.AngularObjectBuilder; +import org.apache.zeppelin.eventbus.EventBus; +import org.apache.zeppelin.eventbus.NoteRemovedEvent; +import org.apache.zeppelin.eventbus.ZeppelinEventBus; import org.apache.zeppelin.interpreter.InterpreterGroup; import org.apache.zeppelin.interpreter.InterpreterSetting; import org.apache.zeppelin.interpreter.remote.RemoteAngularObjectRegistry; @@ -72,6 +75,7 @@ import org.apache.zeppelin.rest.AbstractTestRestApi; import org.apache.zeppelin.scheduler.Job; import org.apache.zeppelin.scheduler.Job.Status; +import org.apache.zeppelin.service.JobManagerService; import org.apache.zeppelin.service.NotebookService; import org.apache.zeppelin.service.ServiceContext; import org.apache.zeppelin.user.AuthenticationInfo; @@ -195,6 +199,32 @@ void testOnNoteRemove_whenJobManagerDisabled() { } } + @Test + void testEventBusHandlesNoteRemovedEvent() throws IOException { + ZeppelinConfiguration conf = ZeppelinConfiguration.load(); + conf.setProperty(ZeppelinConfiguration.ConfVars.ZEPPELIN_EVENTBUS_ENABLED.getVarName(), "true"); + + EventBus eventBus = new ZeppelinEventBus(); + NotebookServer server = new NotebookServer(); + AuthorizationService mockAuthorizationService = mock(AuthorizationService.class); + JobManagerService mockJobManagerService = mock(JobManagerService.class); + Note note = new Note(); + + when(mockAuthorizationService.getOwners(note.getId())).thenReturn(new HashSet<>()); + + server.setZeppelinConfiguration(conf); + server.setAuthorizationService(mockAuthorizationService); + server.setJobManagerService(() -> mockJobManagerService); + server.registerEventBus(eventBus, conf); + + eventBus.post(new NoteRemovedEvent(note, AuthenticationInfo.ANONYMOUS)); + + verify(mockJobManagerService).getNoteJobInfoByUnixTime(Mockito.anyLong(), Mockito.any(), Mockito.any()); + verify(mockJobManagerService).removeNoteJobInfo(eq(note.getId()), Mockito.isNull(), Mockito.any()); + + server.close(); + } + @Test void testOnParagraphCreate_whenJobManagerDisabled() { boolean originalFlag = disableJobManagerAndBackupFlag(); diff --git a/zeppelin-web-angular/e2e/tests/workspace/job-manager/eventbus-removal-parity.spec.ts b/zeppelin-web-angular/e2e/tests/workspace/job-manager/eventbus-removal-parity.spec.ts new file mode 100644 index 000000000000..a2c99ac80b72 --- /dev/null +++ b/zeppelin-web-angular/e2e/tests/workspace/job-manager/eventbus-removal-parity.spec.ts @@ -0,0 +1,203 @@ +/* + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * http://www.apache.org/licenses/LICENSE-2.0 + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +import { Browser, expect, Page, test } from '@playwright/test'; +import { JobManagerPage } from 'e2e/models/job-manager-page'; +import { LoginPage } from 'e2e/models/login-page'; +import { LoginTestUtil, TestCredentials } from 'e2e/models/login-page.util'; +import { addPageAnnotationBeforeEach, PAGES } from '../../../utils'; + +interface NoteJob { + noteId?: string; + isRunningJob?: boolean; + isRemoved?: boolean; + unixTimeLastRun?: number; +} + +interface JobManagerMessage { + op?: string; + data?: { + noteJobs?: { jobs?: NoteJob[] }; + noteRunningJobs?: { jobs?: NoteJob[] }; + }; +} + +class JobManagerMessageRecorder { + private readonly messages: JobManagerMessage[] = []; + + constructor(page: Page) { + page.on('websocket', socket => { + if (new URL(socket.url()).pathname !== '/ws') { + return; + } + socket.on('framereceived', frame => { + const payload = frame.payload.toString(); + if (payload.startsWith('{')) { + this.messages.push(JSON.parse(payload) as JobManagerMessage); + } + }); + }); + } + + hasInitialList(): boolean { + return this.messages.some(message => message.op === 'LIST_NOTE_JOBS'); + } + + initialNoteIds(): string[] { + const initial = [...this.messages].reverse().find(message => message.op === 'LIST_NOTE_JOBS'); + return initial?.data?.noteJobs?.jobs?.flatMap(job => (job.noteId ? [job.noteId] : [])) ?? []; + } + + removals(): NoteJob[] { + return this.messages + .filter(message => message.op === 'LIST_UPDATE_NOTE_JOBS') + .flatMap(message => message.data?.noteRunningJobs?.jobs ?? []) + .filter(job => job.isRemoved === true); + } +} + +const createNote = async (page: Page, label: string): Promise<{ noteId: string; noteName: string }> => { + const noteName = `EventBusParity_${label}_${Date.now()}_${Math.random().toString(36).slice(2, 8)}`; + const response = await page.request.post('/api/notebook', { + data: { + notePath: `E2E_TEST_FOLDER/${noteName}`, + defaultInterpreterGroup: 'python', + addingEmptyParagraph: true + }, + failOnStatusCode: false + }); + await expect(response).toBeOK(); + + const body = (await response.json()) as { body?: string }; + expect(body.body).toEqual(expect.any(String)); + return { noteId: body.body as string, noteName }; +}; + +const deleteNote = async (page: Page, noteId: string): Promise => { + const response = await page.request.delete(`/api/notebook/${noteId}`, { failOnStatusCode: false }); + await expect(response).toBeOK(); +}; + +const expectNoteMissing = async (page: Page, noteId: string): Promise => { + await expect + .poll(async () => (await page.request.get(`/api/notebook/${noteId}`, { failOnStatusCode: false })).status()) + .toBe(404); +}; + +const login = async (page: Page, credentials: TestCredentials): Promise => { + const loginPage = new LoginPage(page); + await loginPage.navigate(); + await expect(loginPage.formContainer).toBeVisible(); + await loginPage.login(credentials.username, credentials.password); + await expect(loginPage.formContainer).toBeHidden(); +}; + +const expectServerEventBusMode = async (page: Page): Promise => { + const response = await page.request.get('/api/configurations/prefix/zeppelin.eventbus.enabled', { + failOnStatusCode: false + }); + await expect(response).toBeOK(); + + const body = (await response.json()) as { body?: Record }; + expect(body.body?.['zeppelin.eventbus.enabled']).toBe( + process.env.ZEPPELIN_EVENTBUS_ENABLED === 'true' ? 'true' : 'false' + ); +}; + +const verifyRemovalParity = async ( + browser: Browser, + ownerPage: Page, + observerCredentials: TestCredentials | undefined, + observerOwnsNotes: boolean +): Promise<{ ownerRemovalCount: number; observerRemovalCount: number }> => { + await expectServerEventBusMode(ownerPage); + const targetNote = await createNote(ownerPage, 'target'); + const barrierNote = await createNote(ownerPage, 'barrier'); + const observerContext = await browser.newContext({ storageState: { cookies: [], origins: [] } }); + const observerPage = await observerContext.newPage(); + + try { + const ownerRecorder = new JobManagerMessageRecorder(ownerPage); + const observerRecorder = new JobManagerMessageRecorder(observerPage); + + if (observerCredentials) { + await login(observerPage, observerCredentials); + } + + const ownerJobManager = new JobManagerPage(ownerPage); + const observerJobManager = new JobManagerPage(observerPage); + + await Promise.all([ownerJobManager.navigate(), observerJobManager.navigate()]); + await expect.poll(() => ownerRecorder.hasInitialList()).toBe(true); + await expect.poll(() => observerRecorder.hasInitialList()).toBe(true); + + expect(ownerRecorder.initialNoteIds()).toEqual(expect.arrayContaining([targetNote.noteId, barrierNote.noteId])); + expect(observerRecorder.initialNoteIds().includes(targetNote.noteId)).toBe(observerOwnsNotes); + expect(observerRecorder.initialNoteIds().includes(barrierNote.noteId)).toBe(observerOwnsNotes); + await expect(ownerJobManager.jobItemByName(targetNote.noteName)).toBeVisible(); + await expect(observerJobManager.jobItemByName(targetNote.noteName)).toHaveCount(observerOwnsNotes ? 1 : 0); + + await deleteNote(ownerPage, targetNote.noteId); + await expectNoteMissing(ownerPage, targetNote.noteId); + await deleteNote(ownerPage, barrierNote.noteId); + await expectNoteMissing(ownerPage, barrierNote.noteId); + + const expectedRemovalOrder = [targetNote.noteId, barrierNote.noteId]; + await expect.poll(() => ownerRecorder.removals().map(job => job.noteId)).toEqual(expectedRemovalOrder); + await expect.poll(() => observerRecorder.removals().map(job => job.noteId)).toEqual(expectedRemovalOrder); + await expect(ownerJobManager.jobItemByName(targetNote.noteName)).toHaveCount(0); + + for (const recorder of [ownerRecorder, observerRecorder]) { + expect(recorder.removals()).toEqual([ + { noteId: targetNote.noteId, isRunningJob: false, isRemoved: true, unixTimeLastRun: 0 }, + { noteId: barrierNote.noteId, isRunningJob: false, isRemoved: true, unixTimeLastRun: 0 } + ]); + } + + return { + ownerRemovalCount: ownerRecorder.removals().length, + observerRemovalCount: observerRecorder.removals().length + }; + } finally { + await observerContext.close(); + await ownerPage.request.delete(`/api/notebook/${targetNote.noteId}`, { failOnStatusCode: false }); + await ownerPage.request.delete(`/api/notebook/${barrierNote.noteId}`, { failOnStatusCode: false }); + } +}; + +test.describe('Job Manager note-removal EventBus parity', () => { + addPageAnnotationBeforeEach(PAGES.WORKSPACE.JOB_MANAGER); + + test.beforeEach(async ({}, testInfo) => { + testInfo.annotations.push({ + type: 'eventbus-mode', + description: process.env.ZEPPELIN_EVENTBUS_ENABLED === 'true' ? 'eventbus' : 'legacy' + }); + }); + + test('anonymous viewers receive one ordered removal payload per deleted note', async ({ browser, page }) => { + test.skip(await LoginTestUtil.isShiroEnabled(), 'ZEPPELIN-6698 requires the anonymous server matrix'); + const removalCounts = await verifyRemovalParity(browser, page, undefined, true); + expect(removalCounts).toEqual({ ownerRemovalCount: 2, observerRemovalCount: 2 }); + }); + + test('an authenticated non-owner receives one ordered removal payload per deleted note', async ({ + browser, + page + }) => { + test.skip(!(await LoginTestUtil.isShiroEnabled()), 'ZEPPELIN-6698 requires the authenticated server matrix'); + const credentials = await LoginTestUtil.getTestCredentials(); + expect(credentials.user2).toBeDefined(); + const removalCounts = await verifyRemovalParity(browser, page, credentials.user2, false); + expect(removalCounts).toEqual({ ownerRemovalCount: 2, observerRemovalCount: 2 }); + }); +});