Make start/stop twitter stream thread-safe (#1859)

* Make start/stop twitter stream thread-safe

* Add rpc for get twitter stream status

* Add get twitter stream status rpc to flask

* Display twitter stream connection status in twitter page

* Fix twitter forwarder task stop

* Update frontend build
This commit is contained in:
Jonathan Zernik 2021-12-15 14:29:36 -08:00 committed by GitHub
parent 6f08e8c0f2
commit 8e975ca811
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
19 changed files with 138 additions and 59 deletions

View file

@ -37,7 +37,10 @@ export default function SetBearerTokenDialog({
const handleResponse = (response) => {
// goToProfilePage(history, response.getProfileId());
// TODO
reloadBearerTokenFn();
if (reloadBearerTokenFn) {
reloadBearerTokenFn();
}
};
const handleErr = (err) => {

View file

@ -8,7 +8,6 @@ import {
Box,
CircularProgress,
Typography,
TextField,
} from '@material-ui/core';
// styles
@ -22,17 +21,19 @@ import SetBearerTokenDialog from '../../components/SetBearerTokenDialog';
import AddTwitterAccountDialog from '../../components/AddTwitterAccountDialog';
import TwitterAccountListItem from '../../components/TwitterAccountListItem';
import CloudIcon from '@mui/icons-material/Cloud';
import CloudOffIcon from '@mui/icons-material/CloudOff';
import {
getTwitterBearerTokenRequest,
getTwitterAccountsRequest,
getTwitterStreamStatusRequest,
} from '../../squeakclient/requests';
export default function Twitter() {
const classes = useStyles();
const [bearerToken, setBearerToken] = useState('');
const [accounts, setAccounts] = useState([]);
const [waitingForBearerToken, setWaitingForBearerToken] = useState(false);
const [streamStatus, setStreamStatus] = useState(null);
const [waitingForAccounts, setWaitingForAccounts] = useState(false);
const [setBearerTokenDialogOpen, setSetBearerTokenDialogOpen] = useState(false);
const [addAccountDialogOpen, setAddAccountDialogOpen] = useState(false);
@ -45,14 +46,6 @@ export default function Twitter() {
};
}
const getBearerToken = () => {
setWaitingForBearerToken(true);
getTwitterBearerTokenRequest((resp) => {
setWaitingForBearerToken(false);
setBearerToken(resp);
});
};
const getAccounts = () => {
setWaitingForAccounts(true);
getTwitterAccountsRequest((resp) => {
@ -61,6 +54,12 @@ export default function Twitter() {
});
};
const getStreamStatus = () => {
getTwitterStreamStatusRequest((resp) => {
setStreamStatus(resp);
});
};
const handleClickOpenSetBearerTokenDialog = () => {
setSetBearerTokenDialogOpen(true);
};
@ -78,10 +77,10 @@ export default function Twitter() {
};
useEffect(() => {
getBearerToken();
getAccounts();
}, []);
useEffect(() => {
getAccounts();
getStreamStatus();
}, []);
function TabPanel(props) {
@ -104,29 +103,44 @@ export default function Twitter() {
);
}
function BearerTokenSummary() {
const bearerTokenText = (bearerToken ? bearerToken : 'not configured')
function StreamStatusSummary() {
const isStreamActive = (streamStatus ? streamStatus.getIsStreamActive() : false);
return (
<Grid item xs={12}>
<Box
p={1}
>
<Typography variant="h5" component="h5">
{`Bearer Token:`}
</Typography>
<TextField
id="standard-textarea"
value={bearerTokenText}
fullWidth
inputProps={{
readOnly: true,
}}
/>
{isStreamActive
? StreamActiveDisplay()
: StreamNotActiveDisplay()
}
</Box>
</Grid>
);
}
function StreamActiveDisplay() {
return (
<>
<Typography variant="h5" component="h5">
Twitter account connected
</Typography>
<CloudIcon fontSize="large" style={{ fill: 'green' }} />
</>
);
}
function StreamNotActiveDisplay() {
return (
<>
<Typography variant="h5" component="h5">
Twitter account not connected
</Typography>
<CloudOffIcon fontSize="large" style={{ fill: 'red' }} />
</>
);
}
function AccountsGridItem(accounts) {
return (
<Grid item xs={12}>
@ -214,7 +228,7 @@ export default function Twitter() {
<Widget disableWidgetMenu>
{SetBearerTokenButton()}
{AddAccountButton()}
{BearerTokenSummary()}
{StreamStatusSummary()}
{AccountsContent()}
</Widget>
</Grid>
@ -232,7 +246,7 @@ export default function Twitter() {
</Tabs>
</AppBar>
<TabPanel value={0} index={0}>
{(waitingForBearerToken || waitingForAccounts)
{( waitingForAccounts)
? WaitingIndicator()
: TwitterAccountsContent()}
</TabPanel>
@ -246,7 +260,6 @@ export default function Twitter() {
<SetBearerTokenDialog
open={setBearerTokenDialogOpen}
handleClose={handleCloseSetBearerTokenDialog}
reloadBearerTokenFn={getBearerToken}
/>
</>
);

View file

@ -143,6 +143,8 @@ import {
GetSellPriceReply,
ClearSellPriceRequest,
ClearSellPriceReply,
GetTwitterStreamStatusRequest,
GetTwitterStreamStatusReply,
} from '../proto/squeak_admin_pb';
console.log('The value of REACT_APP_DEV_MODE_ENABLED is:', Boolean(process.env.REACT_APP_DEV_MODE_ENABLED));
@ -1216,6 +1218,16 @@ export function deleteTwitterAccountRequest(twitterAccountId, handleResponse) {
);
}
export function getTwitterStreamStatusRequest(handleResponse) {
const request = new GetTwitterStreamStatusRequest();
makeRequest(
'gettwitterstreamstatus',
request,
GetTwitterStreamStatusReply.deserializeBinary,
handleResponse,
);
}
export function setSellPriceRequest(priceMsat, handleResponse) {
const request = new SetSellPriceRequest();
request.setPriceMsat(priceMsat);

View file

@ -356,6 +356,10 @@ service SqueakAdmin {
*/
rpc DeleteTwitterAccount (DeleteTwitterAccountRequest) returns (DeleteTwitterAccountReply) {}
/** sqkadmin: `gettwitterstreamstatus`
*/
rpc GetTwitterStreamStatus (GetTwitterStreamStatusRequest) returns (GetTwitterStreamStatusReply) {}
}
message CreateSigningProfileRequest {
@ -1314,3 +1318,11 @@ message DeleteTwitterAccountRequest {
message DeleteTwitterAccountReply {
}
message GetTwitterStreamStatusRequest {
}
message GetTwitterStreamStatusReply {
/// Is twitter stream active
bool is_stream_active = 1;
}

View file

@ -1105,3 +1105,10 @@ class SqueakAdminServerHandler(object):
twitter_account_id,
)
return squeak_admin_pb2.DeleteTwitterAccountReply()
def handle_get_twitter_stream_status(self, request):
logger.info("Handle get twitter stream status")
twitter_stream_status = self.squeak_controller.get_twitter_stream_status()
return squeak_admin_pb2.GetTwitterStreamStatusReply(
is_stream_active=twitter_stream_status,
)

View file

@ -400,3 +400,6 @@ class SqueakAdminServerServicer(squeak_admin_pb2_grpc.SqueakAdminServicer):
def DeleteTwitterAccount(self, request, context):
return self.handler.handle_delete_twitter_account(request)
def GetTwitterStreamStatus(self, request, context):
return self.handler.handle_get_twitter_stream_status(request)

View file

@ -559,6 +559,12 @@ def create_app(handler, username, password):
def deletetwitteraccount(msg):
return handler.handle_delete_twitter_account(msg)
@app.route("/gettwitterstreamstatus", methods=["POST"])
@login_required
@protobuf_serialized(squeak_admin_pb2.GetTwitterStreamStatusRequest())
def gettwitterstreamstatus(msg):
return handler.handle_get_twitter_stream_status(msg)
return app

View file

@ -1,15 +1,15 @@
{
"files": {
"main.js": "/static/js/main.8b406438.chunk.js",
"main.js.map": "/static/js/main.8b406438.chunk.js.map",
"main.js": "/static/js/main.fcaf0671.chunk.js",
"main.js.map": "/static/js/main.fcaf0671.chunk.js.map",
"runtime-main.js": "/static/js/runtime-main.48a1d1b9.js",
"runtime-main.js.map": "/static/js/runtime-main.48a1d1b9.js.map",
"static/css/2.f68b60d4.chunk.css": "/static/css/2.f68b60d4.chunk.css",
"static/js/2.ba0fb77d.chunk.js": "/static/js/2.ba0fb77d.chunk.js",
"static/js/2.ba0fb77d.chunk.js.map": "/static/js/2.ba0fb77d.chunk.js.map",
"static/js/2.695b0a70.chunk.js": "/static/js/2.695b0a70.chunk.js",
"static/js/2.695b0a70.chunk.js.map": "/static/js/2.695b0a70.chunk.js.map",
"index.html": "/index.html",
"static/css/2.f68b60d4.chunk.css.map": "/static/css/2.f68b60d4.chunk.css.map",
"static/js/2.ba0fb77d.chunk.js.LICENSE.txt": "/static/js/2.ba0fb77d.chunk.js.LICENSE.txt",
"static/js/2.695b0a70.chunk.js.LICENSE.txt": "/static/js/2.695b0a70.chunk.js.LICENSE.txt",
"static/media/font-awesome.min.css": "/static/media/fontawesome-webfont.f691f37e.woff",
"static/media/google.a4fa9b6f.svg": "/static/media/google.a4fa9b6f.svg",
"static/media/logo.31df0ce8.svg": "/static/media/logo.31df0ce8.svg"
@ -17,7 +17,7 @@
"entrypoints": [
"static/js/runtime-main.48a1d1b9.js",
"static/css/2.f68b60d4.chunk.css",
"static/js/2.ba0fb77d.chunk.js",
"static/js/main.8b406438.chunk.js"
"static/js/2.695b0a70.chunk.js",
"static/js/main.fcaf0671.chunk.js"
]
}

View file

@ -1 +1 @@
<!doctype html><html lang="en"><head><meta charset="utf-8"/><link rel="shortcut icon" href="/favicon.ico"/><meta name="viewport" content="width=device-width,initial-scale=1,shrink-to-fit=no"/><meta name="theme-color" content="#000000"/><link rel="manifest" href="/manifest.json"/><title>Squeaknode</title><meta name="description" content="Squeaknode is a frontend for accessing a squeak node"><meta name="keywords" content="squeak, bitcoin, lightning"><meta name="author" content="Flatlogic LLC."><link href="/static/css/2.f68b60d4.chunk.css" rel="stylesheet"></head><body style="font-family:Roboto,sans-serif"><noscript>You need to enable JavaScript to run this app.</noscript><div id="root"></div><script>!function(e){function r(r){for(var n,f,l=r[0],a=r[1],i=r[2],c=0,s=[];c<l.length;c++)f=l[c],Object.prototype.hasOwnProperty.call(o,f)&&o[f]&&s.push(o[f][0]),o[f]=0;for(n in a)Object.prototype.hasOwnProperty.call(a,n)&&(e[n]=a[n]);for(p&&p(r);s.length;)s.shift()();return u.push.apply(u,i||[]),t()}function t(){for(var e,r=0;r<u.length;r++){for(var t=u[r],n=!0,l=1;l<t.length;l++){var a=t[l];0!==o[a]&&(n=!1)}n&&(u.splice(r--,1),e=f(f.s=t[0]))}return e}var n={},o={1:0},u=[];function f(r){if(n[r])return n[r].exports;var t=n[r]={i:r,l:!1,exports:{}};return e[r].call(t.exports,t,t.exports,f),t.l=!0,t.exports}f.m=e,f.c=n,f.d=function(e,r,t){f.o(e,r)||Object.defineProperty(e,r,{enumerable:!0,get:t})},f.r=function(e){"undefined"!=typeof Symbol&&Symbol.toStringTag&&Object.defineProperty(e,Symbol.toStringTag,{value:"Module"}),Object.defineProperty(e,"__esModule",{value:!0})},f.t=function(e,r){if(1&r&&(e=f(e)),8&r)return e;if(4&r&&"object"==typeof e&&e&&e.__esModule)return e;var t=Object.create(null);if(f.r(t),Object.defineProperty(t,"default",{enumerable:!0,value:e}),2&r&&"string"!=typeof e)for(var n in e)f.d(t,n,function(r){return e[r]}.bind(null,n));return t},f.n=function(e){var r=e&&e.__esModule?function(){return e.default}:function(){return e};return f.d(r,"a",r),r},f.o=function(e,r){return Object.prototype.hasOwnProperty.call(e,r)},f.p="/";var l=this["webpackJsonpsqueak-node-frontend"]=this["webpackJsonpsqueak-node-frontend"]||[],a=l.push.bind(l);l.push=r,l=l.slice();for(var i=0;i<l.length;i++)r(l[i]);var p=a;t()}([])</script><script src="/static/js/2.ba0fb77d.chunk.js"></script><script src="/static/js/main.8b406438.chunk.js"></script></body></html>
<!doctype html><html lang="en"><head><meta charset="utf-8"/><link rel="shortcut icon" href="/favicon.ico"/><meta name="viewport" content="width=device-width,initial-scale=1,shrink-to-fit=no"/><meta name="theme-color" content="#000000"/><link rel="manifest" href="/manifest.json"/><title>Squeaknode</title><meta name="description" content="Squeaknode is a frontend for accessing a squeak node"><meta name="keywords" content="squeak, bitcoin, lightning"><meta name="author" content="Flatlogic LLC."><link href="/static/css/2.f68b60d4.chunk.css" rel="stylesheet"></head><body style="font-family:Roboto,sans-serif"><noscript>You need to enable JavaScript to run this app.</noscript><div id="root"></div><script>!function(e){function r(r){for(var n,f,l=r[0],a=r[1],i=r[2],c=0,s=[];c<l.length;c++)f=l[c],Object.prototype.hasOwnProperty.call(o,f)&&o[f]&&s.push(o[f][0]),o[f]=0;for(n in a)Object.prototype.hasOwnProperty.call(a,n)&&(e[n]=a[n]);for(p&&p(r);s.length;)s.shift()();return u.push.apply(u,i||[]),t()}function t(){for(var e,r=0;r<u.length;r++){for(var t=u[r],n=!0,l=1;l<t.length;l++){var a=t[l];0!==o[a]&&(n=!1)}n&&(u.splice(r--,1),e=f(f.s=t[0]))}return e}var n={},o={1:0},u=[];function f(r){if(n[r])return n[r].exports;var t=n[r]={i:r,l:!1,exports:{}};return e[r].call(t.exports,t,t.exports,f),t.l=!0,t.exports}f.m=e,f.c=n,f.d=function(e,r,t){f.o(e,r)||Object.defineProperty(e,r,{enumerable:!0,get:t})},f.r=function(e){"undefined"!=typeof Symbol&&Symbol.toStringTag&&Object.defineProperty(e,Symbol.toStringTag,{value:"Module"}),Object.defineProperty(e,"__esModule",{value:!0})},f.t=function(e,r){if(1&r&&(e=f(e)),8&r)return e;if(4&r&&"object"==typeof e&&e&&e.__esModule)return e;var t=Object.create(null);if(f.r(t),Object.defineProperty(t,"default",{enumerable:!0,value:e}),2&r&&"string"!=typeof e)for(var n in e)f.d(t,n,function(r){return e[r]}.bind(null,n));return t},f.n=function(e){var r=e&&e.__esModule?function(){return e.default}:function(){return e};return f.d(r,"a",r),r},f.o=function(e,r){return Object.prototype.hasOwnProperty.call(e,r)},f.p="/";var l=this["webpackJsonpsqueak-node-frontend"]=this["webpackJsonpsqueak-node-frontend"]||[],a=l.push.bind(l);l.push=r,l=l.slice();for(var i=0;i<l.length;i++)r(l[i]);var p=a;t()}([])</script><script src="/static/js/2.695b0a70.chunk.js"></script><script src="/static/js/main.fcaf0671.chunk.js"></script></body></html>

File diff suppressed because one or more lines are too long

File diff suppressed because one or more lines are too long

File diff suppressed because one or more lines are too long

File diff suppressed because one or more lines are too long

File diff suppressed because one or more lines are too long

File diff suppressed because one or more lines are too long

View file

@ -907,6 +907,9 @@ class SqueakController:
return None
return user_config.twitter_bearer_token
def get_twitter_stream_status(self) -> bool:
return self.tweet_forwarder.is_processing()
def add_twitter_account(self, handle: str, profile_id: int) -> Optional[int]:
twitter_account = TwitterAccount(
twitter_account_id=None,

View file

@ -57,6 +57,12 @@ class TwitterForwarder:
if self.current_task is not None:
self.current_task.stop_processing()
def is_processing(self) -> bool:
with self.lock:
if self.current_task is not None:
return self.current_task.is_processing()
return False
class TwitterForwarderTask:
@ -69,6 +75,7 @@ class TwitterForwarderTask:
self.retry_s = retry_s
self.stopped = threading.Event()
self.tweet_stream = None
self.lock = threading.Lock()
def start_processing(self):
logger.info("Starting twitter forwarder task.")
@ -77,13 +84,31 @@ class TwitterForwarderTask:
daemon=True,
).start()
def setup_stream(self, bearer_token, handles):
logger.info("Starting Twitter stream with bearer token: {} and twitter handles: {}".format(
bearer_token,
handles,
))
with self.lock:
if self.stopped.is_set():
return
twitter_stream = TwitterStream(bearer_token, handles)
self.tweet_stream = twitter_stream.get_tweets()
def stop_processing(self):
logger.info("Stopping twitter forwarder task.")
self.stopped.set()
if self.tweet_stream is not None:
self.tweet_stream.cancel_fn()
with self.lock:
self.stopped.set()
if self.tweet_stream is not None:
self.tweet_stream.cancel_fn()
def is_processing(self) -> bool:
return self.tweet_stream is not None
def process_forward_tweets(self):
# Do exponential backoff
wait_s = self.retry_s
while not self.stopped.is_set():
try:
bearer_token = self.get_bearer_token()
@ -92,24 +117,19 @@ class TwitterForwarderTask:
return
if not handles:
return
logger.info("Starting forward tweets with bearer token: {} and twitter handles: {}".format(
bearer_token,
handles,
))
twitter_stream = TwitterStream(bearer_token, handles)
self.tweet_stream = twitter_stream.get_tweets()
if self.stopped.is_set():
return
self.setup_stream(bearer_token, handles)
for tweet in self.tweet_stream.result_stream:
self.handle_tweet(tweet)
# TODO: use more specific error.
except Exception:
self.tweet_stream = None
logger.exception(
"Unable to subscribe tweet stream. Retrying in {} seconds...".format(
self.retry_s,
wait_s,
),
)
self.stopped.wait(self.retry_s)
self.stopped.wait(wait_s)
wait_s *= 2
def get_bearer_token(self) -> str:
return self.squeak_controller.get_twitter_bearer_token() or ''