From 13e7dd0226c3f9d6822976bf08bb10994cfa2a59 Mon Sep 17 00:00:00 2001 From: Jonathan Zernik Date: Fri, 3 Sep 2021 14:45:44 -0700 Subject: [PATCH] Only show followed squeaks in timeline (#1171) * Only show followed squeaks in timeline * Add rpc method to subscribe timeline squeaks * Use subscribe new timeline squeaks rpc method in frontend --- frontend/src/pages/timeline/Timeline.js | 4 ++-- frontend/src/squeakclient/requests.js | 17 +++++++++++++++++ proto/squeak_admin.proto | 8 ++++++++ squeaknode/admin/squeak_admin_server_handler.py | 17 +++++++++++++++++ .../admin/squeak_admin_server_servicer.py | 12 ++++++++++++ squeaknode/db/squeak_db.py | 5 +++++ squeaknode/node/squeak_controller.py | 7 +++++++ 7 files changed, 68 insertions(+), 2 deletions(-) diff --git a/frontend/src/pages/timeline/Timeline.js b/frontend/src/pages/timeline/Timeline.js index 23055689..a42851a7 100644 --- a/frontend/src/pages/timeline/Timeline.js +++ b/frontend/src/pages/timeline/Timeline.js @@ -27,7 +27,7 @@ import SqueakList from '../../components/SqueakList'; import { getTimelineSqueakDisplaysRequest, getNetworkRequest, - subscribeSqueakDisplaysRequest, + subscribeTimelineSqueakDisplaysRequest, } from '../../squeakclient/requests'; const SQUEAKS_PER_PAGE = 10; @@ -47,7 +47,7 @@ export default function TimelinePage() { setWaitingForTimeline(true); getTimelineSqueakDisplaysRequest(limit, blockHeight, squeakTime, squeakHash, handleLoadedTimeline, alertFailedRequest); }; - const subscribeNewSqueaks = () => subscribeSqueakDisplaysRequest(handleLoadedNewSqueak); + const subscribeNewSqueaks = () => subscribeTimelineSqueakDisplaysRequest(handleLoadedNewSqueak); const getNetwork = () => { getNetworkRequest(setNetwork); }; diff --git a/frontend/src/squeakclient/requests.js b/frontend/src/squeakclient/requests.js index 111f608d..d4615dec 100644 --- a/frontend/src/squeakclient/requests.js +++ b/frontend/src/squeakclient/requests.js @@ -70,6 +70,7 @@ import { SubscribeAddressSqueakDisplaysRequest, SubscribeAncestorSqueakDisplaysRequest, SubscribeSqueakDisplaysRequest, + SubscribeTimelineSqueakDisplaysRequest, } from '../proto/squeak_admin_pb'; import { SqueakAdminClient } from '../proto/squeak_admin_grpc_web_pb'; @@ -738,3 +739,19 @@ export function subscribeSqueakDisplaysRequest(handleResponse) { console.log(stream); return stream; } + +export function subscribeTimelineSqueakDisplaysRequest(handleResponse) { + const request = new SubscribeTimelineSqueakDisplaysRequest(); + const stream = client.subscribeSqueakDisplays(request); + stream.on('data', (response) => { + handleResponse(response.getSqueakDisplayEntry()); + }); + stream.on('end', (end) => { + // stream end signal + console.log(end); + alert(`Stream ended: ${end}`); + }); + console.log('Stream object:'); + console.log(stream); + return stream; +} diff --git a/proto/squeak_admin.proto b/proto/squeak_admin.proto index dd0b155c..4dd9f014 100644 --- a/proto/squeak_admin.proto +++ b/proto/squeak_admin.proto @@ -304,6 +304,10 @@ service SqueakAdmin { */ rpc SubscribeSqueakDisplays (SubscribeSqueakDisplaysRequest) returns (stream GetSqueakDisplayReply) {} + /** sqkadmin: `subscribetimelinesqueakdisplays` + */ + rpc SubscribeTimelineSqueakDisplays (SubscribeTimelineSqueakDisplaysRequest) returns (stream GetSqueakDisplayReply) {} + } message CreateSigningProfileRequest { @@ -1068,3 +1072,7 @@ message SubscribeAncestorSqueakDisplaysRequest { message SubscribeSqueakDisplaysRequest { } + +message SubscribeTimelineSqueakDisplaysRequest { +} + diff --git a/squeaknode/admin/squeak_admin_server_handler.py b/squeaknode/admin/squeak_admin_server_handler.py index e783fc7a..b6b647ac 100644 --- a/squeaknode/admin/squeak_admin_server_handler.py +++ b/squeaknode/admin/squeak_admin_server_handler.py @@ -924,3 +924,20 @@ class SqueakAdminServerHandler(object): yield squeak_admin_pb2.GetSqueakDisplayReply( squeak_display_entry=display_message ) + + def handle_subscribe_timeline_squeak_displays(self, request, stopped): + logger.info("Handle subscribe timeline squeak displays") + squeak_display_stream = self.squeak_controller.subscribe_timeline_squeak_entries( + stopped, + ) + for squeak_display in squeak_display_stream: + if squeak_display is None: + yield squeak_admin_pb2.GetSqueakDisplayReply( + squeak_display_entry=None + ) + else: + display_message = squeak_entry_to_message( + squeak_display) + yield squeak_admin_pb2.GetSqueakDisplayReply( + squeak_display_entry=display_message + ) diff --git a/squeaknode/admin/squeak_admin_server_servicer.py b/squeaknode/admin/squeak_admin_server_servicer.py index a7a2e0ec..e6a4cad6 100644 --- a/squeaknode/admin/squeak_admin_server_servicer.py +++ b/squeaknode/admin/squeak_admin_server_servicer.py @@ -363,3 +363,15 @@ class SqueakAdminServerServicer(squeak_admin_pb2_grpc.SqueakAdminServicer): request, stopped, ) + + def SubscribeTimelineSqueakDisplays(self, request, context): + stopped = threading.Event() + + def on_rpc_done(): + logger.info("Stopping SubscribeTimelineSqueakDisplays.") + stopped.set() + context.add_callback(on_rpc_done) + return self.handler.handle_subscribe_timeline_squeak_displays( + request, + stopped, + ) diff --git a/squeaknode/db/squeak_db.py b/squeaknode/db/squeak_db.py index 1ac1ffcd..1e265d8b 100644 --- a/squeaknode/db/squeak_db.py +++ b/squeaknode/db/squeak_db.py @@ -132,6 +132,10 @@ class SqueakDb: def profile_has_no_private_key(self): return self.profiles.c.private_key == None # noqa: E711 + @property + def profile_is_following(self): + return self.profiles.c.following == True # noqa: E711 + @property def received_offer_does_not_exist(self): return self.received_offers.c.squeak_hash == None # noqa: E711 @@ -264,6 +268,7 @@ class SqueakDb: self.profiles.c.address == self.squeaks.c.author_address, ) ) + .where(self.profile_is_following) .where( tuple_( self.squeaks.c.n_block_height, diff --git a/squeaknode/node/squeak_controller.py b/squeaknode/node/squeak_controller.py index afaffee0..19bf56af 100644 --- a/squeaknode/node/squeak_controller.py +++ b/squeaknode/node/squeak_controller.py @@ -720,3 +720,10 @@ class SqueakController: for item in self.new_squeak_listener.yield_items(stopped): squeak_hash = get_hash(item) yield self.get_squeak_entry(squeak_hash) + + def subscribe_timeline_squeak_entries(self, stopped: threading.Event): + for item in self.new_squeak_listener.yield_items(stopped): + followed_addresses = self.get_followed_addresses() + if str(item.GetAddress()) in set(followed_addresses): + squeak_hash = get_hash(item) + yield self.get_squeak_entry(squeak_hash)