|
| 1 | +import mgp |
| 2 | +import json |
| 3 | + |
| 4 | + |
| 5 | +@mgp.transformation |
| 6 | +def satellite(messages: mgp.Messages |
| 7 | + ) -> mgp.Record(query=str, parameters=mgp.Nullable[mgp.Map]): |
| 8 | + result_queries = [] |
| 9 | + |
| 10 | + for i in range(messages.total_messages()): |
| 11 | + message = messages.message_at(i) |
| 12 | + json_message = json.loads(message.payload().decode('utf8')) |
| 13 | + # print(json_message) |
| 14 | + result_queries.append( |
| 15 | + mgp.Record( |
| 16 | + query=("MERGE (s:Satellite { id: toString($id) }) " |
| 17 | + "SET s.x = toFloat($x), " |
| 18 | + "s.y=toFloat($y), " |
| 19 | + "s.z=toFloat($z);"), |
| 20 | + parameters={ |
| 21 | + "id": json_message["id"], |
| 22 | + "x": json_message["x"], |
| 23 | + "y": json_message["y"], |
| 24 | + "z": json_message["z"]})) |
| 25 | + |
| 26 | + return result_queries |
| 27 | + |
| 28 | + |
| 29 | +@mgp.transformation |
| 30 | +def city(messages: mgp.Messages |
| 31 | + ) -> mgp.Record(query=str, parameters=mgp.Nullable[mgp.Map]): |
| 32 | + result_queries = [] |
| 33 | + |
| 34 | + for i in range(messages.total_messages()): |
| 35 | + message = messages.message_at(i) |
| 36 | + json_message = json.loads(message.payload().decode('utf8')) |
| 37 | + # print(json_message) |
| 38 | + result_queries.append( |
| 39 | + mgp.Record( |
| 40 | + query=("MERGE (s:City { id: toString($id) }) " |
| 41 | + "SET s.name = toString($name), " |
| 42 | + "s.x=toFloat($x), " |
| 43 | + "s.y=toFloat($y);"), |
| 44 | + parameters={ |
| 45 | + "id": json_message["id"], |
| 46 | + "name": json_message["name"], |
| 47 | + "x": json_message["x"], |
| 48 | + "y": json_message["y"]})) |
| 49 | + |
| 50 | + return result_queries |
| 51 | + |
| 52 | + |
| 53 | +@mgp.transformation |
| 54 | +def visible_from(messages: mgp.Messages |
| 55 | + ) -> mgp.Record(query=str, |
| 56 | + parameters=mgp.Nullable[mgp.Map]): |
| 57 | + result_queries = [] |
| 58 | + |
| 59 | + for i in range(messages.total_messages()): |
| 60 | + message = messages.message_at(i) |
| 61 | + json_message = json.loads(message.payload().decode('utf8')) |
| 62 | + # print(json_message) |
| 63 | + result_queries.append( |
| 64 | + mgp.Record( |
| 65 | + query=("MATCH (c:City { id: toString($city_id)}), (s:Satellite { id: toString($satellite_id) }) " |
| 66 | + "CREATE (s)-[r:VISIBLE_FROM { transmission_time: toFloat($transmission_time) }]->(c);"), |
| 67 | + parameters={ |
| 68 | + "city_id": json_message["city_id"], |
| 69 | + "satellite_id": json_message["satellite_id"], |
| 70 | + "transmission_time": json_message["transmission_time"]})) |
| 71 | + |
| 72 | + return result_queries |
| 73 | + |
| 74 | + |
| 75 | +@mgp.transformation |
| 76 | +def delete_visible_from(messages: mgp.Messages |
| 77 | + ) -> mgp.Record(query=str, |
| 78 | + parameters=mgp.Nullable[mgp.Map]): |
| 79 | + result_queries = [] |
| 80 | + |
| 81 | + for i in range(messages.total_messages()): |
| 82 | + message = messages.message_at(i) |
| 83 | + json_message = json.loads(message.payload().decode('utf8')) |
| 84 | + # print(json_message) |
| 85 | + result_queries.append( |
| 86 | + mgp.Record( |
| 87 | + query=( |
| 88 | + "MATCH (:Satellite)-[r]->(:City { id: toString($id) }) DELETE r;"), |
| 89 | + parameters={ |
| 90 | + "id": json_message["id"]})) |
| 91 | + |
| 92 | + return result_queries |
| 93 | + |
| 94 | + |
| 95 | +@mgp.transformation |
| 96 | +def laser_link(messages: mgp.Messages |
| 97 | + ) -> mgp.Record(query=str, |
| 98 | + parameters=mgp.Nullable[mgp.Map]): |
| 99 | + result_queries = [] |
| 100 | + |
| 101 | + for i in range(messages.total_messages()): |
| 102 | + message = messages.message_at(i) |
| 103 | + json_message = json.loads(message.payload().decode('utf8')) |
| 104 | + # print(json_message) |
| 105 | + result_queries.append( |
| 106 | + mgp.Record( |
| 107 | + query=("MATCH (a: Satellite { id: toString($laser_id) }), (b: Satellite { id: toString($moving_object_id) }) " |
| 108 | + "MERGE (a)-[c:CONNECTED_TO]->(b) " |
| 109 | + "SET c.transmission_time = toFloat($laser_transmission_time);"), |
| 110 | + parameters={ |
| 111 | + "laser_id": json_message["laser_id"], |
| 112 | + "moving_object_id": json_message["moving_object_id"], |
| 113 | + "laser_transmission_time": json_message["laser_transmission_time"]})) |
| 114 | + |
| 115 | + return result_queries |
0 commit comments