A Jupyter notebook for experiments for the DEBS25 paper EPL

Fuente: Zenodo
Saved in:
Bibliographic Details
Main Authors: TOMMASINI, RICCARDO, Langhi, Samuele, Bonifati, Angela, Bernhardt, Thomas
Format: Recurso digital
Published: Zenodo 2025
Online Access:
Tags: Add Tag
No Tags, Be the first to tag this record!
_version_ 1866901897183494144
author TOMMASINI, RICCARDO
Langhi, Samuele
Bonifati, Angela
Bernhardt, Thomas
author_facet TOMMASINI, RICCARDO
Langhi, Samuele
Bonifati, Angela
Bernhardt, Thomas
contents <p># kEPLrbook</p> <p> </p> <p>A Jupyter notebook buildt for experimenting with EPL. <br>Experimental interactive environment, every notebook could be considered an experiment itself, so could be then exported and stored. </p> <p>## Docker</p> <p>To run the containerized version, just type `docker-compose up` inside the main repository folder. This will start 4 docker containers:<br>- runtime<br>- inputmanager<br>- outputmanager<br>- notebooks</p> <p>In order to access the notebooks, it is sufficient to connect to the browser at `localhost:8888`, no password required. </p> <p>## Experiment Design</p> <p>The experiment is based on 3 primary elements</p> <p>* Input: the events in input that are inserted into the runtime<br>* Query: the query that is used to process the input <br>* Expected Output (optional): the expected output that is eventually compared with the actual output</p> <p>Mirroring this classification the system is based on a 4-tier architecture in which all the backend servers are connected through web sockets.</p> <p>* Input Manager<br>* Output Manager<br>* Runtime Server<br>* Jupyter Server</p> <p><br>### Input Manager</p> <p>An input server that manages the input coming from the front-end of the application. The input is represented in YAML format and mapped using Jackson library ([Github](https://github.com/FasterXML/jackson)).</p> <p>#### Input Insertion</p> <p>  _The input is sent to the server in YAML format and then processed. In the input the events are represented as a time keyed map in which are listed all the events that should be sent in that timestamp ([Sample](#body)). The server reads it, mapping it into a special list object called **EventList** and separates the events into single, timestampend entities, using a wrapper called **TimestampedEvent**._</p> <p>* **URL**</p> <p>  /input</p> <p>* **Method:**<br>  <br>  `POST` </p> <p>* **<a name="body"> <br>Data Params<br></a>**</p> <p>```<br>  inputs:<br>     0 :<br>       - type: "NbaGame"<br>         id: 1<br>         awayTeam: "Sacramento Kings"<br>         homeTeam: "New York Knicks" <br>     1000 :<br>       - type: "NbaGame"<br>         id: 2<br>         awayTeam: "Sacramento Kings"<br>         homeTeam: "New York Knicks"</p> <p>       - type: "NbaGame"<br>         id: 3<br>         awayTeam: "Sacramento Kings"<br>         homeTeam: "New York Knicks<br>```</p> <p><br>* **Success Response:**<br>  <br>  * **Code:** 200<br>    **Content:** `Input file created.`<br> <br>* **Sample Call:**</p> <p>```<br>def sendMessage(self, content):<br>  r = self.session.post(self.host+"/input", content)<br>  print(r.text)<br>```<br>* **Notes:**</p> <p>  * _The server store the input received as a file in order to maintain the last input._ <br>  * _The serialization format of event sent to the runtime is JSON._<br>  * _After having sent all the events as TimestampedEvents the server send a "finish" message in order to indicate that there are no more input events._</p> <p>### Runtime Server</p> <p>The server on which is located the Esper Runtime. It receives through web socket the various input events and sends them into the Esper Runtime.</p> <p>#### Query Insertion</p> <p>_The server exposes a single RESTful endpoint that is used to post the EPL module used to process the various events._</p> <p>* **URL**</p> <p>  /query</p> <p>* **Method:**<br>  <br>  `POST` </p> <p>* **Data Params**</p> <p>```<br>create schema NbaGame as (type string, id int, awayTeam string, homeTeam string); <br>@name('Output')select sum(id) from NbaGame#time(500 milliseconds) output every 500 milliseconds; <br>```</p> <p>* **Success Response:**<br>  <br>  * **Code:** 200<br>    **Content:** `Query file created.`<br> <br>* **Sample Call:**</p> <p>```<br>def sendMessage(self, content):<br>  r = self.session.post(self.host+"/query", content)<br>  print(r.text)<br>```<br>* **Notes:**</p> <p>  * _The server store the module received as a file in order to always maintain and use at least a stored module in order to process the input._ <br>  * _The module **is not immediately deployed** but is simply stored._<br>  * _The module is deployed into the runtime by the web socket manager object when the first event arrives and the module file has been modified. In that case all the previous modules are undeployed and the new module is loaded._<br>  </p> <p>### Output Manager</p> <p>The Output Manager receives the resulting events from the runtime through the web socket. If an expected output is specified, this is received from the server and used to compare it with the actual output. The server receives the event **one by one** so basically it collects them and, once the "finished" message arrives, stores the event onto a file, using it after for comparison if the expected output is present. </p> <p>_In case the expected output is not specified the Output Manager will simply compose a representation of the actual output._</p> <p>#### Output Insertion</p> <p>  _The expected is sent to the server in YAML format and then processed. In the output the events are represented as a time keyed map in which are listed all the events that should be arrived in that timestamp ([Sample](#output)), just like the input. The server reads it, mapping it into the object **EventList**._</p> <p>* **URL**</p> <p>  /output</p> <p>* **Method:**<br>  <br>  `POST` </p> <p>* **<a name="output"> <br>Data Params<br></a>**</p> <p>```<br>  inputs:<br>     0 :<br>       - type: "NbaGame"<br>         id: 1<br>         awayTeam: "Sacramento Kings"<br>         homeTeam: "New York Knicks" <br>     1000 :<br>       - type: "NbaGame"<br>         id: 2<br>         awayTeam: "Sacramento Kings"<br>         homeTeam: "New York Knicks"</p> <p>       - type: "NbaGame"<br>         id: 3<br>         awayTeam: "Sacramento Kings"<br>         homeTeam: "New York Knicks<br>```</p> <p><br>* **Success Response:**<br>  <br>  * **Code:** 200<br>    **Content:** `Created expected output file.`<br> <br>* **Sample Call:**</p> <p>```<br>def sendMessage(self, content):<br>  r = self.session.post(self.host+"/output", content)<br>  print(r.text)<br>```<br>* **Notes:**</p> <p>  * _The server store the expected output received as a file in order to maintain the last expected output._ <br>  * _The serialization format of event received by the runtime is JSON._<br>  </p> <p>#### Output Link Retrieval</p> <p>  _At first the server checks that both the "actual output" and the "expected output" files exists, then if in the "expected output" file is present the server will proceed to writing the "comparison.txt" file. The comparison is done checking both the **EventList** objects (composed from the previously cited output files, using the jackson mapper), checking if both elements have at first the same timestamps, then comparing the event lists for each timestamp. At the end the method will return the **link to retrieve the comparison result**. If no expected output is provided (keyword "gotcha" in the "expected output" file) the comparison file will be written with the content of the "actual output" file_<br>  </p> <p>* **URL**</p> <p>  /output</p> <p>* **Method:**<br>  <br>  `GET` </p> <p>* **Success Response:**<br>  <br>  * **Code:** 200<br>    **Content:** `http://localhost:1234/final_output`<br> <br>* **Sample Call:**</p> <p>```<br>def getRequest(self):<br>  r = self.session.request('GET', self.host+'/output')<br>  print(r.text)<br>```<br>* **Notes:**</p> <p>  * _The sources of comparison **are the files**, this is due to the fact that they provide a possibility for monitoring, in fact if there is a problem with the parsing and mapping of the expected output (and also of the input in the Input Manager) i can check directly on the file the source of the problem, this is also valid because it gives an intermediate directly consultable representation of the actual output_<br>  * _The comparison file is a txt file since gives more freedom on the format of representation, since it must adapt to a comparison result (a simple text), and also to the actual output representation (text written in YAML format)_<br>  </p> <p>#### Output Retrieval</p> <p>_Endpoint for the retrieval of the final result, being it the actual output or the comparison result. The endpoint is reached when, through the browser, the result link is clicked, this way the server can set up the response of a get request from the browser with the final result._</p> <p><br>* **URL**</p> <p>  /final_output</p> <p>* **Method:**<br>  <br>  `GET` </p> <p>* **Success Response:**<br>  <br>  * **Code:** 200<br>    **Content:** <br>      <br>      * Actual Output<br>      ```<br>      ---<br>      inputs:<br>        0:<br>        - sum(id): 1<br>        1000:<br>        - sum(id): 3<br>        - sum(id): 6<br>      ```<br>      <br>      * Comparison Result<br>      <br>     ```<br>     Start Comparison:<br>      Time 0:<br>       - 0th events: <br>        - Actual event : <br>          {sum(id)=1}<br>        - Expected event: <br>          {type=NbaGamer, id=1, awayTeam=Sacramento Kings, homeTeam=New York Knicks}<br>        - Events are different<br>      ******************************<br>      Time 1000:<br>       - 0th events: <br>        - Actual event : <br>          {sum(id)=3}<br>        - Expected event: <br>          {type=NbaGamer, id=2, awayTeam=Sacramento Kings, homeTeam=New York Knicks}<br>        - Events are different<br>       - 1th events: <br>        - Actual event : <br>          {sum(id)=6}<br>        - Expected event: <br>          {type=NbaGamer, id=3, awayTeam=Sacramento Kings, homeTeam=New York Knicks}<br>        - Events are different<br>      ******************************<br>     ```</p> <p> <br>* **Sample Call:**</p> <p>```<br>def getRequest(self):<br>  r = self.session.request('GET', self.host+'/final_output')<br>  print(r.text)<br>```</p> <p>### Jupyter Server</p> <p>The Jupyter Server is where the jupyter kernel is running. The kernel is a python-wrapped kernel for jupyter. The implementation of the _doExecute_ method is essential, since it basically changes the destination of the requests according to presence of a magic in the cell.</p> <p>```<br>    def do_execute(self, code, silent,<br>                   store_history=True,<br>                   user_expressions=None,<br>                   allow_stdin=False):</p> <p>        functions = code.split('\n')</p> <p>        line = functions[0]<br>        n=0</p> <p>        # Processing the magics first, setting up the destination<br>        if line[0] == '%':<br>            self.connection.process_magics(line)<br>            n = 1<br>        else:<br>            self.connection.switchDest("/query")<br>            self.connection.switchHost("http://localhost:7890")<br>        res = self.connection.sendMessage('\n'.join(functions[n:]))</p> <p><br>        if not silent:<br>            # We send the standard output to the<br>            # client.<br>            self.send_response(<br>                self.iopub_socket,<br>                'stream', {<br>                    'name': 'stdout',<br>                    'data': ('Plotting {n} '<br>                             'function(s)'). \<br>                        format(n=len(functions))})</p> <p>            # We prepare the response, setting up the format.<br>            content = {<br>                'source': 'kernel',<br>                'data': {<br>                    'text/plain': res<br>                },</p> <p>            }</p> <p>            # We send the display_data message with<br>            # the contents.<br>            self.send_response(self.iopub_socket,<br>                               'display_data', content)</p> <p>        # We return the execution results.<br>        return {'status': 'ok',<br>                'execution_count':<br>                    self.execution_count,<br>                'payload': [],<br>                'user_expressions': {},<br>                }</p> <p>```</p> <p>Magic | Destination | Requests Succession <br>----- | ----------- | -------------------<br>%input | /input | POST request with the input as the body<br>%output | /output | POST request with the content of the cell as the body, then GET request<br>none | /query | POST request with the content of the cell as the body, containing the epl module</p> <p><br>The various response of the requests are then returned to response cell of the relative run cell, with a simple text representation. </p> <p>The notebook cells should be in this order <br>1. Query cell<br>2. Input cell<br>3. Output cell</p> <p>This is due to the fact that without an epl module the runtime cannot basically process events, so every events in input will be basically ignored. Point 2 and 3 are instead interchangeable.</p> <p>  <br>## Architecture </p> <p>The application is based on a 4-tier architecture. While the frontend communicates with the 3 servers through HTTP requests, these servers commmunicates using WebSockets, this way a constant, full-duplex communication can be established.</p> <p>![alt text](./kEPLrbook_servers/src/main/resources/images/architecture.png)</p> <p>In our case to initiate the socket connection the order of execution should be:</p> <p>1. Runtime server<br>2. Output Manager<br>3. Input Manager</p> <p>While point 2 and 3 are interchangeable, the precedence of the Runtime Server is essential, since in the setting up of the WebSocket communication we treat the _Output Manager and Input Manager as the WebSocketClients_ and the _Runtime Server as the WebSocketServer._</p> <p> </p> <p> </p> <p>### Events flow</p> <p>1. Events in input are embedded in the input YAML file, posted on the Input Manager by the Jupyter Server with a POST request<br>2. The file is read and mapped into **EventList** object<br>3. Separates the event into single entities and embed them into **TimestampedEvent** object<br>4. Serialize the created object into JSON format<br>5. Using a WebSocketClient instance (connection already established) we send the previously created serialization to the Runtime Server<br>6. The Runtime Server deserialize it and set the runtime time to the timestamp of the event before sending the event itself into the runtime<br>7. The **SendingListener** object attached to the statement with name "Prova" is also a WebSocketManager object, and sends through web socket using the _update_ method the events arrived<br>8. The Output Manager receives the events (as **TimestampedEvent** objects) and collect them, till the "finish" message is received<br>9. [OPTIONAL] The Output Manager receives the expected outputs and map the file into an **EventList** object<br>10. Once every event is collected into the **EventList** object the Server checks the "expected/output.yml" file<br>11. If the "expected/output.yml" file contains the keyword "gotcha" the Output manager return the Actual Output representation (in YAML)<br>12. Else, it computes the comparison between the 2 **EventList** objects and returns it</p> <p><br>## Guidelines for use (Local Version, not Docker)</p> <p>Requirements <br>* Jupyter Notebook application installed<br>* Maven to import the various libraries through the "pom.xml"</p> <p>To use it you have to install the jupyter kernel:</p> <p>```<br>jupyter kernelspec install --user <path-to-kEPLr_kernel-dir><br>```</p> <p>Type this line into the command line.<br>If it does not work try to modify appropiately the PYTHONPATH environment variable in the ".bash_profile" file.</p> <p># EPL Ambiguities</p> <p>The following are EPL query performance with two nested EVERY. This is to demonstrate the language risk in terms of performances. </p> <p>![alt text](./perfeveryevery.png)</p> <p><br># Future Works </p> <p>- [X] Implement a Retry system during the websocket communication start in order to dockerize everything, since we cannot control the sequentiality of the execution of the 3 servers.</p> <p> </p> <p> </p>
format Recurso digital
id zenodo_https___doi_org_10_5281_zenodo_15402318
institution Zenodo
language
publishDate 2025
publisher Zenodo
record_format zenodo
spellingShingle A Jupyter notebook for experiments for the DEBS25 paper EPL
TOMMASINI, RICCARDO
Langhi, Samuele
Bonifati, Angela
Bernhardt, Thomas
<p># kEPLrbook</p> <p> </p> <p>A Jupyter notebook buildt for experimenting with EPL. <br>Experimental interactive environment, every notebook could be considered an experiment itself, so could be then exported and stored. </p> <p>## Docker</p> <p>To run the containerized version, just type `docker-compose up` inside the main repository folder. This will start 4 docker containers:<br>- runtime<br>- inputmanager<br>- outputmanager<br>- notebooks</p> <p>In order to access the notebooks, it is sufficient to connect to the browser at `localhost:8888`, no password required. </p> <p>## Experiment Design</p> <p>The experiment is based on 3 primary elements</p> <p>* Input: the events in input that are inserted into the runtime<br>* Query: the query that is used to process the input <br>* Expected Output (optional): the expected output that is eventually compared with the actual output</p> <p>Mirroring this classification the system is based on a 4-tier architecture in which all the backend servers are connected through web sockets.</p> <p>* Input Manager<br>* Output Manager<br>* Runtime Server<br>* Jupyter Server</p> <p><br>### Input Manager</p> <p>An input server that manages the input coming from the front-end of the application. The input is represented in YAML format and mapped using Jackson library ([Github](https://github.com/FasterXML/jackson)).</p> <p>#### Input Insertion</p> <p>  _The input is sent to the server in YAML format and then processed. In the input the events are represented as a time keyed map in which are listed all the events that should be sent in that timestamp ([Sample](#body)). The server reads it, mapping it into a special list object called **EventList** and separates the events into single, timestampend entities, using a wrapper called **TimestampedEvent**._</p> <p>* **URL**</p> <p>  /input</p> <p>* **Method:**<br>  <br>  `POST` </p> <p>* **<a name="body"> <br>Data Params<br></a>**</p> <p>```<br>  inputs:<br>     0 :<br>       - type: "NbaGame"<br>         id: 1<br>         awayTeam: "Sacramento Kings"<br>         homeTeam: "New York Knicks" <br>     1000 :<br>       - type: "NbaGame"<br>         id: 2<br>         awayTeam: "Sacramento Kings"<br>         homeTeam: "New York Knicks"</p> <p>       - type: "NbaGame"<br>         id: 3<br>         awayTeam: "Sacramento Kings"<br>         homeTeam: "New York Knicks<br>```</p> <p><br>* **Success Response:**<br>  <br>  * **Code:** 200<br>    **Content:** `Input file created.`<br> <br>* **Sample Call:**</p> <p>```<br>def sendMessage(self, content):<br>  r = self.session.post(self.host+"/input", content)<br>  print(r.text)<br>```<br>* **Notes:**</p> <p>  * _The server store the input received as a file in order to maintain the last input._ <br>  * _The serialization format of event sent to the runtime is JSON._<br>  * _After having sent all the events as TimestampedEvents the server send a "finish" message in order to indicate that there are no more input events._</p> <p>### Runtime Server</p> <p>The server on which is located the Esper Runtime. It receives through web socket the various input events and sends them into the Esper Runtime.</p> <p>#### Query Insertion</p> <p>_The server exposes a single RESTful endpoint that is used to post the EPL module used to process the various events._</p> <p>* **URL**</p> <p>  /query</p> <p>* **Method:**<br>  <br>  `POST` </p> <p>* **Data Params**</p> <p>```<br>create schema NbaGame as (type string, id int, awayTeam string, homeTeam string); <br>@name('Output')select sum(id) from NbaGame#time(500 milliseconds) output every 500 milliseconds; <br>```</p> <p>* **Success Response:**<br>  <br>  * **Code:** 200<br>    **Content:** `Query file created.`<br> <br>* **Sample Call:**</p> <p>```<br>def sendMessage(self, content):<br>  r = self.session.post(self.host+"/query", content)<br>  print(r.text)<br>```<br>* **Notes:**</p> <p>  * _The server store the module received as a file in order to always maintain and use at least a stored module in order to process the input._ <br>  * _The module **is not immediately deployed** but is simply stored._<br>  * _The module is deployed into the runtime by the web socket manager object when the first event arrives and the module file has been modified. In that case all the previous modules are undeployed and the new module is loaded._<br>  </p> <p>### Output Manager</p> <p>The Output Manager receives the resulting events from the runtime through the web socket. If an expected output is specified, this is received from the server and used to compare it with the actual output. The server receives the event **one by one** so basically it collects them and, once the "finished" message arrives, stores the event onto a file, using it after for comparison if the expected output is present. </p> <p>_In case the expected output is not specified the Output Manager will simply compose a representation of the actual output._</p> <p>#### Output Insertion</p> <p>  _The expected is sent to the server in YAML format and then processed. In the output the events are represented as a time keyed map in which are listed all the events that should be arrived in that timestamp ([Sample](#output)), just like the input. The server reads it, mapping it into the object **EventList**._</p> <p>* **URL**</p> <p>  /output</p> <p>* **Method:**<br>  <br>  `POST` </p> <p>* **<a name="output"> <br>Data Params<br></a>**</p> <p>```<br>  inputs:<br>     0 :<br>       - type: "NbaGame"<br>         id: 1<br>         awayTeam: "Sacramento Kings"<br>         homeTeam: "New York Knicks" <br>     1000 :<br>       - type: "NbaGame"<br>         id: 2<br>         awayTeam: "Sacramento Kings"<br>         homeTeam: "New York Knicks"</p> <p>       - type: "NbaGame"<br>         id: 3<br>         awayTeam: "Sacramento Kings"<br>         homeTeam: "New York Knicks<br>```</p> <p><br>* **Success Response:**<br>  <br>  * **Code:** 200<br>    **Content:** `Created expected output file.`<br> <br>* **Sample Call:**</p> <p>```<br>def sendMessage(self, content):<br>  r = self.session.post(self.host+"/output", content)<br>  print(r.text)<br>```<br>* **Notes:**</p> <p>  * _The server store the expected output received as a file in order to maintain the last expected output._ <br>  * _The serialization format of event received by the runtime is JSON._<br>  </p> <p>#### Output Link Retrieval</p> <p>  _At first the server checks that both the "actual output" and the "expected output" files exists, then if in the "expected output" file is present the server will proceed to writing the "comparison.txt" file. The comparison is done checking both the **EventList** objects (composed from the previously cited output files, using the jackson mapper), checking if both elements have at first the same timestamps, then comparing the event lists for each timestamp. At the end the method will return the **link to retrieve the comparison result**. If no expected output is provided (keyword "gotcha" in the "expected output" file) the comparison file will be written with the content of the "actual output" file_<br>  </p> <p>* **URL**</p> <p>  /output</p> <p>* **Method:**<br>  <br>  `GET` </p> <p>* **Success Response:**<br>  <br>  * **Code:** 200<br>    **Content:** `http://localhost:1234/final_output`<br> <br>* **Sample Call:**</p> <p>```<br>def getRequest(self):<br>  r = self.session.request('GET', self.host+'/output')<br>  print(r.text)<br>```<br>* **Notes:**</p> <p>  * _The sources of comparison **are the files**, this is due to the fact that they provide a possibility for monitoring, in fact if there is a problem with the parsing and mapping of the expected output (and also of the input in the Input Manager) i can check directly on the file the source of the problem, this is also valid because it gives an intermediate directly consultable representation of the actual output_<br>  * _The comparison file is a txt file since gives more freedom on the format of representation, since it must adapt to a comparison result (a simple text), and also to the actual output representation (text written in YAML format)_<br>  </p> <p>#### Output Retrieval</p> <p>_Endpoint for the retrieval of the final result, being it the actual output or the comparison result. The endpoint is reached when, through the browser, the result link is clicked, this way the server can set up the response of a get request from the browser with the final result._</p> <p><br>* **URL**</p> <p>  /final_output</p> <p>* **Method:**<br>  <br>  `GET` </p> <p>* **Success Response:**<br>  <br>  * **Code:** 200<br>    **Content:** <br>      <br>      * Actual Output<br>      ```<br>      ---<br>      inputs:<br>        0:<br>        - sum(id): 1<br>        1000:<br>        - sum(id): 3<br>        - sum(id): 6<br>      ```<br>      <br>      * Comparison Result<br>      <br>     ```<br>     Start Comparison:<br>      Time 0:<br>       - 0th events: <br>        - Actual event : <br>          {sum(id)=1}<br>        - Expected event: <br>          {type=NbaGamer, id=1, awayTeam=Sacramento Kings, homeTeam=New York Knicks}<br>        - Events are different<br>      ******************************<br>      Time 1000:<br>       - 0th events: <br>        - Actual event : <br>          {sum(id)=3}<br>        - Expected event: <br>          {type=NbaGamer, id=2, awayTeam=Sacramento Kings, homeTeam=New York Knicks}<br>        - Events are different<br>       - 1th events: <br>        - Actual event : <br>          {sum(id)=6}<br>        - Expected event: <br>          {type=NbaGamer, id=3, awayTeam=Sacramento Kings, homeTeam=New York Knicks}<br>        - Events are different<br>      ******************************<br>     ```</p> <p> <br>* **Sample Call:**</p> <p>```<br>def getRequest(self):<br>  r = self.session.request('GET', self.host+'/final_output')<br>  print(r.text)<br>```</p> <p>### Jupyter Server</p> <p>The Jupyter Server is where the jupyter kernel is running. The kernel is a python-wrapped kernel for jupyter. The implementation of the _doExecute_ method is essential, since it basically changes the destination of the requests according to presence of a magic in the cell.</p> <p>```<br>    def do_execute(self, code, silent,<br>                   store_history=True,<br>                   user_expressions=None,<br>                   allow_stdin=False):</p> <p>        functions = code.split('\n')</p> <p>        line = functions[0]<br>        n=0</p> <p>        # Processing the magics first, setting up the destination<br>        if line[0] == '%':<br>            self.connection.process_magics(line)<br>            n = 1<br>        else:<br>            self.connection.switchDest("/query")<br>            self.connection.switchHost("http://localhost:7890")<br>        res = self.connection.sendMessage('\n'.join(functions[n:]))</p> <p><br>        if not silent:<br>            # We send the standard output to the<br>            # client.<br>            self.send_response(<br>                self.iopub_socket,<br>                'stream', {<br>                    'name': 'stdout',<br>                    'data': ('Plotting {n} '<br>                             'function(s)'). \<br>                        format(n=len(functions))})</p> <p>            # We prepare the response, setting up the format.<br>            content = {<br>                'source': 'kernel',<br>                'data': {<br>                    'text/plain': res<br>                },</p> <p>            }</p> <p>            # We send the display_data message with<br>            # the contents.<br>            self.send_response(self.iopub_socket,<br>                               'display_data', content)</p> <p>        # We return the execution results.<br>        return {'status': 'ok',<br>                'execution_count':<br>                    self.execution_count,<br>                'payload': [],<br>                'user_expressions': {},<br>                }</p> <p>```</p> <p>Magic | Destination | Requests Succession <br>----- | ----------- | -------------------<br>%input | /input | POST request with the input as the body<br>%output | /output | POST request with the content of the cell as the body, then GET request<br>none | /query | POST request with the content of the cell as the body, containing the epl module</p> <p><br>The various response of the requests are then returned to response cell of the relative run cell, with a simple text representation. </p> <p>The notebook cells should be in this order <br>1. Query cell<br>2. Input cell<br>3. Output cell</p> <p>This is due to the fact that without an epl module the runtime cannot basically process events, so every events in input will be basically ignored. Point 2 and 3 are instead interchangeable.</p> <p>  <br>## Architecture </p> <p>The application is based on a 4-tier architecture. While the frontend communicates with the 3 servers through HTTP requests, these servers commmunicates using WebSockets, this way a constant, full-duplex communication can be established.</p> <p>![alt text](./kEPLrbook_servers/src/main/resources/images/architecture.png)</p> <p>In our case to initiate the socket connection the order of execution should be:</p> <p>1. Runtime server<br>2. Output Manager<br>3. Input Manager</p> <p>While point 2 and 3 are interchangeable, the precedence of the Runtime Server is essential, since in the setting up of the WebSocket communication we treat the _Output Manager and Input Manager as the WebSocketClients_ and the _Runtime Server as the WebSocketServer._</p> <p> </p> <p> </p> <p>### Events flow</p> <p>1. Events in input are embedded in the input YAML file, posted on the Input Manager by the Jupyter Server with a POST request<br>2. The file is read and mapped into **EventList** object<br>3. Separates the event into single entities and embed them into **TimestampedEvent** object<br>4. Serialize the created object into JSON format<br>5. Using a WebSocketClient instance (connection already established) we send the previously created serialization to the Runtime Server<br>6. The Runtime Server deserialize it and set the runtime time to the timestamp of the event before sending the event itself into the runtime<br>7. The **SendingListener** object attached to the statement with name "Prova" is also a WebSocketManager object, and sends through web socket using the _update_ method the events arrived<br>8. The Output Manager receives the events (as **TimestampedEvent** objects) and collect them, till the "finish" message is received<br>9. [OPTIONAL] The Output Manager receives the expected outputs and map the file into an **EventList** object<br>10. Once every event is collected into the **EventList** object the Server checks the "expected/output.yml" file<br>11. If the "expected/output.yml" file contains the keyword "gotcha" the Output manager return the Actual Output representation (in YAML)<br>12. Else, it computes the comparison between the 2 **EventList** objects and returns it</p> <p><br>## Guidelines for use (Local Version, not Docker)</p> <p>Requirements <br>* Jupyter Notebook application installed<br>* Maven to import the various libraries through the "pom.xml"</p> <p>To use it you have to install the jupyter kernel:</p> <p>```<br>jupyter kernelspec install --user <path-to-kEPLr_kernel-dir><br>```</p> <p>Type this line into the command line.<br>If it does not work try to modify appropiately the PYTHONPATH environment variable in the ".bash_profile" file.</p> <p># EPL Ambiguities</p> <p>The following are EPL query performance with two nested EVERY. This is to demonstrate the language risk in terms of performances. </p> <p>![alt text](./perfeveryevery.png)</p> <p><br># Future Works </p> <p>- [X] Implement a Retry system during the websocket communication start in order to dockerize everything, since we cannot control the sequentiality of the execution of the 3 servers.</p> <p> </p> <p> </p>
title A Jupyter notebook for experiments for the DEBS25 paper EPL
url https://doi.org/10.5281/zenodo.15402318