From 4c4ce6d50991d4177eb1e9b98f78a6e22f251394 Mon Sep 17 00:00:00 2001 From: Lee moon soo Date: Tue, 13 Oct 2015 18:14:19 +0200 Subject: [PATCH 1/5] Don't update back to the browser where update the angular object --- .../zeppelin/socket/NotebookServer.java | 34 +++++++++++++++---- 1 file changed, 28 insertions(+), 6 deletions(-) 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 c6aa846d13e..0303cd0c3ce 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 @@ -261,6 +261,26 @@ private void broadcast(String noteId, Message m) { } } + private void broadcastExcept(String noteId, Message m, NotebookSocket exclude) { + synchronized (noteSocketMap) { + List socketLists = noteSocketMap.get(noteId); + if (socketLists == null || socketLists.size() == 0) { + return; + } + LOG.info("SEND >> " + m.op); + for (NotebookSocket conn : socketLists) { + if (exclude.equals(conn)) { + continue; + } + try { + conn.send(serializeMessage(m)); + } catch (IOException e) { + LOG.error("socket error", e); + } + } + } + } + private void broadcastAll(Message m) { synchronized (connectedSockets) { for (NotebookSocket conn : connectedSockets) { @@ -425,7 +445,7 @@ private void updateParagraph(NotebookSocket conn, Notebook notebook, note.persist(); broadcast(note.id(), new Message(OP.PARAGRAPH).put("paragraph", p)); } - + private void cloneNote(NotebookSocket conn, Notebook notebook, Message fromMessage) throws IOException, CloneNotSupportedException { String noteId = getOpenNoteId(conn); @@ -475,7 +495,7 @@ private void completion(NotebookSocket conn, Notebook notebook, * @param notebook the notebook. * @param fromMessage the message. */ - private void angularObjectUpdated(WebSocket conn, Notebook notebook, + private void angularObjectUpdated(NotebookSocket conn, Notebook notebook, Message fromMessage) { String noteId = (String) fromMessage.get("noteId"); String interpreterGroupId = (String) fromMessage.get("interpreterGroupId"); @@ -529,20 +549,22 @@ private void angularObjectUpdated(WebSocket conn, Notebook notebook, if (interpreterGroupId.equals(setting.getInterpreterGroup().getId())) { AngularObjectRegistry angularObjectRegistry = setting .getInterpreterGroup().getAngularObjectRegistry(); - this.broadcast( + this.broadcastExcept( n.id(), new Message(OP.ANGULAR_OBJECT_UPDATE).put("angularObject", ao) .put("interpreterGroupId", interpreterGroupId) - .put("noteId", n.id())); + .put("noteId", n.id()), + conn); } } } } else { // broadcast to all web session for the note - this.broadcast( + this.broadcastExcept( note.id(), new Message(OP.ANGULAR_OBJECT_UPDATE).put("angularObject", ao) .put("interpreterGroupId", interpreterGroupId) - .put("noteId", note.id())); + .put("noteId", note.id()), + conn); } } From 527c56f43ba532088a875afc49252e7045645779 Mon Sep 17 00:00:00 2001 From: Lee moon soo Date: Sun, 18 Oct 2015 03:18:31 +0900 Subject: [PATCH 2/5] Add test for broadcast for angularObjectUpdate --- .../zeppelin/socket/NotebookServerTest.java | 154 +++++++++++++++++- 1 file changed, 152 insertions(+), 2 deletions(-) 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 c17809ae91e..ea2c4fe405f 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 @@ -20,16 +20,100 @@ package org.apache.zeppelin.socket; import static org.junit.Assert.*; + +import java.io.File; import java.io.IOException; -import org.junit.Test; +import org.apache.zeppelin.conf.ZeppelinConfiguration; +import org.apache.zeppelin.conf.ZeppelinConfiguration.ConfVars; +import org.apache.zeppelin.interpreter.InterpreterFactory; +import org.apache.zeppelin.interpreter.InterpreterGroup; +import org.apache.zeppelin.interpreter.InterpreterOption; +import org.apache.zeppelin.interpreter.InterpreterSetting; +import org.apache.zeppelin.interpreter.mock.MockInterpreter1; +import org.apache.zeppelin.interpreter.mock.MockInterpreter2; +import org.apache.zeppelin.notebook.JobListenerFactory; +import org.apache.zeppelin.notebook.Note; +import org.apache.zeppelin.notebook.Notebook; +import org.apache.zeppelin.notebook.repo.NotebookRepo; +import org.apache.zeppelin.notebook.repo.VFSNotebookRepo; +import org.apache.zeppelin.scheduler.JobListener; +import org.apache.zeppelin.scheduler.SchedulerFactory; +import org.apache.zeppelin.server.ZeppelinServer; +import org.apache.zeppelin.socket.Message.OP; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; +import com.google.gson.Gson; import java.net.UnknownHostException; import java.net.InetAddress; +import java.util.List; +import javax.servlet.http.HttpServletRequest; +import static org.mockito.Mockito.*; + /** * BASIC Zeppelin rest api tests */ -public class NotebookServerTest { +public class NotebookServerTest implements JobListenerFactory { + + private File tmpDir; + private ZeppelinConfiguration conf; + private SchedulerFactory schedulerFactory; + private File notebookDir; + private Notebook notebook; + private NotebookRepo notebookRepo; + private InterpreterFactory factory; + private NotebookServer notebookServer; + private Gson gson; + + @Before + public void setUp() throws Exception { + gson = new Gson(); + tmpDir = new File(System.getProperty("java.io.tmpdir")+"/ZeppelinLTest_"+System.currentTimeMillis()); + tmpDir.mkdirs(); + new File(tmpDir, "conf").mkdirs(); + notebookDir = new File(System.getProperty("java.io.tmpdir")+"/ZeppelinLTest_"+System.currentTimeMillis()+"/notebook"); + notebookDir.mkdirs(); + + System.setProperty(ConfVars.ZEPPELIN_HOME.getVarName(), tmpDir.getAbsolutePath()); + System.setProperty(ConfVars.ZEPPELIN_NOTEBOOK_DIR.getVarName(), notebookDir.getAbsolutePath()); + System.setProperty(ConfVars.ZEPPELIN_INTERPRETERS.getVarName(), "org.apache.zeppelin.interpreter.mock.MockInterpreter1,org.apache.zeppelin.interpreter.mock.MockInterpreter2"); + + conf = ZeppelinConfiguration.create(); + + this.schedulerFactory = new SchedulerFactory(); + + MockInterpreter1.register("mock1", "org.apache.zeppelin.interpreter.mock.MockInterpreter1"); + MockInterpreter2.register("mock2", "org.apache.zeppelin.interpreter.mock.MockInterpreter2"); + + factory = new InterpreterFactory(conf, new InterpreterOption(false), null); + + notebookRepo = new VFSNotebookRepo(conf); + notebook = new Notebook(conf, notebookRepo, schedulerFactory, factory, this); + + notebookServer = new NotebookServer(); + ZeppelinServer.notebook = notebook; + ZeppelinServer.notebookServer = notebookServer; + } + + @After + public void tearDown() throws Exception { + delete(tmpDir); + } + + private void delete(File file){ + if(file.isFile()) file.delete(); + else if(file.isDirectory()){ + File [] files = file.listFiles(); + if(files!=null && files.length>0){ + for(File f : files){ + delete(f); + } + } + file.delete(); + } + } @Test public void checkOrigin() throws UnknownHostException { @@ -45,5 +129,71 @@ public void checkInvalidOrigin(){ NotebookServer server = new NotebookServer(); assertFalse(server.checkOrigin(new TestHttpServletRequest(), "http://evillocalhost:8080")); } + + @Test + public void testMakeSureNoAngularObjectBroadcastToWebsocketWhoFireTheEvent() throws IOException { + // create a notebook + Note note1 = notebook.createNote(); + + // get reference to interpreterGroup + InterpreterGroup interpreterGroup = null; + List settings = note1.getNoteReplLoader().getInterpreterSettings(); + for (InterpreterSetting setting : settings) { + if (setting.getInterpreterGroup() == null) { + continue; + } + + interpreterGroup = setting.getInterpreterGroup(); + break; + } + + // add angularObject + interpreterGroup.getAngularObjectRegistry().add("object1", "value1", note1.getId()); + + // create two sockets and open it + NotebookSocket sock1 = createWebSocket(); + NotebookSocket sock2 = createWebSocket(); + + assertEquals(sock1, sock1); + assertNotEquals(sock1, sock2); + + notebookServer.onOpen(sock1); + notebookServer.onOpen(sock2); + + // open the same notebook from sockets + notebookServer.onMessage(sock1, gson.toJson(new Message(OP.GET_NOTE).put("id", note1.getId()))); + notebookServer.onMessage(sock2, gson.toJson(new Message(OP.GET_NOTE).put("id", note1.getId()))); + + + // update object from sock1 + notebookServer.onMessage(sock1, gson.toJson( + new Message(OP.ANGULAR_OBJECT_UPDATED) + .put("noteId", note1.getId()) + .put("name", "object1") + .put("value", "value1") + .put("interpreterGroupId", interpreterGroup.getId()))); + + + // expect object is broadcasted except for where the update is created + verify(sock1, times(2)).send(anyString()); // getNote, getAngularObject + verify(sock2, times(3)).send(anyString()); // getNote, getAngularObject, updateAngularObject + + notebook.removeNote(note1.getId()); + } + + private NotebookSocket createWebSocket() { + NotebookSocket sock = mock(NotebookSocket.class); + when(sock.getRequest()).thenReturn(createHttpServletRequest()); + return sock; + } + + private HttpServletRequest createHttpServletRequest() { + return mock(HttpServletRequest.class); + } + + @Override + public JobListener getParagraphJobListener(Note note) { + return null; + } } From 8175c8d13b2f3637b9189451a4081ac7c2164a1d Mon Sep 17 00:00:00 2001 From: Lee moon soo Date: Sun, 18 Oct 2015 03:53:27 +0900 Subject: [PATCH 3/5] Add mock interpreter --- .../interpreter/mock/MockInterpreter1.java | 73 +++++++++++++++++++ .../zeppelin/socket/NotebookServerTest.java | 2 - 2 files changed, 73 insertions(+), 2 deletions(-) create mode 100644 zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/mock/MockInterpreter1.java diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/mock/MockInterpreter1.java b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/mock/MockInterpreter1.java new file mode 100644 index 00000000000..b76a8b2dc38 --- /dev/null +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/mock/MockInterpreter1.java @@ -0,0 +1,73 @@ +/* + * 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.interpreter.mock; + +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Properties; + +import org.apache.zeppelin.interpreter.Interpreter; +import org.apache.zeppelin.interpreter.InterpreterContext; +import org.apache.zeppelin.interpreter.InterpreterResult; +import org.apache.zeppelin.scheduler.Scheduler; +import org.apache.zeppelin.scheduler.SchedulerFactory; + +public class MockInterpreter1 extends Interpreter{ + Map vars = new HashMap(); + + public MockInterpreter1(Properties property) { + super(property); + } + + @Override + public void open() { + } + + @Override + public void close() { + } + + @Override + public InterpreterResult interpret(String st, InterpreterContext context) { + return new InterpreterResult(InterpreterResult.Code.SUCCESS, "repl1: "+st); + } + + @Override + public void cancel(InterpreterContext context) { + } + + @Override + public FormType getFormType() { + return FormType.SIMPLE; + } + + @Override + public int getProgress(InterpreterContext context) { + return 0; + } + + @Override + public Scheduler getScheduler() { + return SchedulerFactory.singleton().createOrGetFIFOScheduler("test_"+this.hashCode()); + } + + @Override + public List completion(String buf, int cursor) { + return null; + } +} 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 ea2c4fe405f..ecf0178e46a 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 @@ -31,7 +31,6 @@ import org.apache.zeppelin.interpreter.InterpreterOption; import org.apache.zeppelin.interpreter.InterpreterSetting; import org.apache.zeppelin.interpreter.mock.MockInterpreter1; -import org.apache.zeppelin.interpreter.mock.MockInterpreter2; import org.apache.zeppelin.notebook.JobListenerFactory; import org.apache.zeppelin.notebook.Note; import org.apache.zeppelin.notebook.Notebook; @@ -85,7 +84,6 @@ public void setUp() throws Exception { this.schedulerFactory = new SchedulerFactory(); MockInterpreter1.register("mock1", "org.apache.zeppelin.interpreter.mock.MockInterpreter1"); - MockInterpreter2.register("mock2", "org.apache.zeppelin.interpreter.mock.MockInterpreter2"); factory = new InterpreterFactory(conf, new InterpreterOption(false), null); From cd45a7fa8cdbc49310056db241d42285528ad956 Mon Sep 17 00:00:00 2001 From: Lee moon soo Date: Fri, 13 Nov 2015 15:45:27 +0900 Subject: [PATCH 4/5] Fix test --- .../zeppelin/socket/NotebookServerTest.java | 98 +++++-------------- 1 file changed, 25 insertions(+), 73 deletions(-) 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 ecf0178e46a..5275d81ac67 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 @@ -21,96 +21,51 @@ import static org.junit.Assert.*; -import java.io.File; import java.io.IOException; -import org.apache.zeppelin.conf.ZeppelinConfiguration; -import org.apache.zeppelin.conf.ZeppelinConfiguration.ConfVars; -import org.apache.zeppelin.interpreter.InterpreterFactory; import org.apache.zeppelin.interpreter.InterpreterGroup; -import org.apache.zeppelin.interpreter.InterpreterOption; import org.apache.zeppelin.interpreter.InterpreterSetting; -import org.apache.zeppelin.interpreter.mock.MockInterpreter1; -import org.apache.zeppelin.notebook.JobListenerFactory; import org.apache.zeppelin.notebook.Note; import org.apache.zeppelin.notebook.Notebook; -import org.apache.zeppelin.notebook.repo.NotebookRepo; -import org.apache.zeppelin.notebook.repo.VFSNotebookRepo; -import org.apache.zeppelin.scheduler.JobListener; -import org.apache.zeppelin.scheduler.SchedulerFactory; +import org.apache.zeppelin.rest.AbstractTestRestApi; import org.apache.zeppelin.server.ZeppelinServer; import org.apache.zeppelin.socket.Message.OP; -import org.junit.After; -import org.junit.Before; +import org.junit.AfterClass; +import org.junit.BeforeClass; import org.junit.Test; + import com.google.gson.Gson; + import java.net.UnknownHostException; import java.net.InetAddress; import java.util.List; + import javax.servlet.http.HttpServletRequest; + import static org.mockito.Mockito.*; /** * BASIC Zeppelin rest api tests */ -public class NotebookServerTest implements JobListenerFactory { - - private File tmpDir; - private ZeppelinConfiguration conf; - private SchedulerFactory schedulerFactory; - private File notebookDir; - private Notebook notebook; - private NotebookRepo notebookRepo; - private InterpreterFactory factory; - private NotebookServer notebookServer; - private Gson gson; - - @Before - public void setUp() throws Exception { - gson = new Gson(); - tmpDir = new File(System.getProperty("java.io.tmpdir")+"/ZeppelinLTest_"+System.currentTimeMillis()); - tmpDir.mkdirs(); - new File(tmpDir, "conf").mkdirs(); - notebookDir = new File(System.getProperty("java.io.tmpdir")+"/ZeppelinLTest_"+System.currentTimeMillis()+"/notebook"); - notebookDir.mkdirs(); - - System.setProperty(ConfVars.ZEPPELIN_HOME.getVarName(), tmpDir.getAbsolutePath()); - System.setProperty(ConfVars.ZEPPELIN_NOTEBOOK_DIR.getVarName(), notebookDir.getAbsolutePath()); - System.setProperty(ConfVars.ZEPPELIN_INTERPRETERS.getVarName(), "org.apache.zeppelin.interpreter.mock.MockInterpreter1,org.apache.zeppelin.interpreter.mock.MockInterpreter2"); - - conf = ZeppelinConfiguration.create(); - - this.schedulerFactory = new SchedulerFactory(); - - MockInterpreter1.register("mock1", "org.apache.zeppelin.interpreter.mock.MockInterpreter1"); +public class NotebookServerTest extends AbstractTestRestApi { - factory = new InterpreterFactory(conf, new InterpreterOption(false), null); - notebookRepo = new VFSNotebookRepo(conf); - notebook = new Notebook(conf, notebookRepo, schedulerFactory, factory, this); + private static Notebook notebook; + private static NotebookServer notebookServer; + private static Gson gson; - notebookServer = new NotebookServer(); - ZeppelinServer.notebook = notebook; - ZeppelinServer.notebookServer = notebookServer; - } - - @After - public void tearDown() throws Exception { - delete(tmpDir); + @BeforeClass + public static void init() throws Exception { + AbstractTestRestApi.startUp(); + gson = new Gson(); + notebook = ZeppelinServer.notebook; + notebookServer = ZeppelinServer.notebookServer; } - private void delete(File file){ - if(file.isFile()) file.delete(); - else if(file.isDirectory()){ - File [] files = file.listFiles(); - if(files!=null && files.length>0){ - for(File f : files){ - delete(f); - } - } - file.delete(); - } + @AfterClass + public static void destroy() throws Exception { + AbstractTestRestApi.shutDown(); } @Test @@ -157,11 +112,13 @@ public void testMakeSureNoAngularObjectBroadcastToWebsocketWhoFireTheEvent() thr notebookServer.onOpen(sock1); notebookServer.onOpen(sock2); - + verify(sock1, times(0)).send(anyString()); // getNote, getAngularObject // open the same notebook from sockets notebookServer.onMessage(sock1, gson.toJson(new Message(OP.GET_NOTE).put("id", note1.getId()))); notebookServer.onMessage(sock2, gson.toJson(new Message(OP.GET_NOTE).put("id", note1.getId()))); + reset(sock1); + reset(sock2); // update object from sock1 notebookServer.onMessage(sock1, gson.toJson( @@ -173,8 +130,8 @@ public void testMakeSureNoAngularObjectBroadcastToWebsocketWhoFireTheEvent() thr // expect object is broadcasted except for where the update is created - verify(sock1, times(2)).send(anyString()); // getNote, getAngularObject - verify(sock2, times(3)).send(anyString()); // getNote, getAngularObject, updateAngularObject + verify(sock1, times(0)).send(anyString()); + verify(sock2, times(1)).send(anyString()); notebook.removeNote(note1.getId()); } @@ -188,10 +145,5 @@ private NotebookSocket createWebSocket() { private HttpServletRequest createHttpServletRequest() { return mock(HttpServletRequest.class); } - - @Override - public JobListener getParagraphJobListener(Note note) { - return null; - } } From 40f8f848deae9d5113b2b7f485273d798741985e Mon Sep 17 00:00:00 2001 From: Lee moon soo Date: Sat, 14 Nov 2015 20:39:54 +0900 Subject: [PATCH 5/5] Change log level to debug for SEND message --- .../main/java/org/apache/zeppelin/socket/NotebookServer.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 0303cd0c3ce..8832c6fa912 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 @@ -267,7 +267,7 @@ private void broadcastExcept(String noteId, Message m, NotebookSocket exclude) { if (socketLists == null || socketLists.size() == 0) { return; } - LOG.info("SEND >> " + m.op); + LOG.debug("SEND >> " + m.op); for (NotebookSocket conn : socketLists) { if (exclude.equals(conn)) { continue;