Update to websockets
This commit is contained in:
@@ -27,28 +27,29 @@ html = """
|
|||||||
<ul id='messages'>
|
<ul id='messages'>
|
||||||
</ul>
|
</ul>
|
||||||
<script>
|
<script>
|
||||||
var client_id = Date.now()
|
var client_id = Date.now();
|
||||||
document.querySelector("#ws-id").textContent = client_id;
|
document.querySelector("#ws-id").textContent = client_id;
|
||||||
// var ws = new WebSocket(`ws://localhost:5005/ws/${client_id}`);
|
var ws = new WebSocket(`ws://localhost:5005/ws/${client_id}`);
|
||||||
var ws = new WebSocket(`ws://localhost:5005/ws`);
|
// var ws = new WebSocket(`ws://localhost:5005/ws`);
|
||||||
//var ws = new WebSocket("ws://localhost:8000/ws");
|
// var ws = new WebSocket("ws://localhost:8000/ws");
|
||||||
//var ws = new WebSocket("ws://fastapi.localhost/ws");
|
// var ws = new WebSocket("ws://fastapi.localhost/ws");
|
||||||
ws.onmessage = function(event) {
|
ws.onmessage = function(event) {
|
||||||
var messages = document.getElementById('messages')
|
var messages = document.getElementById('messages');
|
||||||
var message = document.createElement('li')
|
var message = document.createElement('li');
|
||||||
var content = document.createTextNode(event.data)
|
var content = document.createTextNode(event.data);
|
||||||
message.appendChild(content)
|
message.appendChild(content);
|
||||||
messages.appendChild(message)
|
messages.appendChild(message);
|
||||||
};
|
};
|
||||||
function sendMessage(event) {
|
function sendMessage(event) {
|
||||||
var input = document.getElementById("messageText")
|
console.log('*** sendMessage() ***');
|
||||||
|
var input = document.getElementById("messageText");
|
||||||
var data = { 'client_id': client_id, 'message': input.value };
|
var data = { 'client_id': client_id, 'message': input.value };
|
||||||
var data_json_str = JSON.stringify(data);
|
var data_json_str = JSON.stringify(data);
|
||||||
ws.send(data_json_str);
|
ws.send(data_json_str);
|
||||||
|
|
||||||
//ws.send(input.value)
|
//ws.send(input.value);
|
||||||
input.value = ''
|
input.value = '';
|
||||||
event.preventDefault()
|
event.preventDefault();
|
||||||
}
|
}
|
||||||
</script>
|
</script>
|
||||||
</body>
|
</body>
|
||||||
@@ -63,8 +64,12 @@ async def get(response: Response = Response):
|
|||||||
return HTMLResponse(html)
|
return HTMLResponse(html)
|
||||||
|
|
||||||
|
|
||||||
@router.websocket("/ws")
|
@router.websocket("/ws/{client_id}")
|
||||||
async def websocket_endpoint(websocket: WebSocket, response: Response = Response):
|
async def websocket_endpoint(
|
||||||
|
websocket: WebSocket,
|
||||||
|
client_id: int,
|
||||||
|
response: Response = Response,
|
||||||
|
):
|
||||||
log.setLevel(logging.DEBUG)
|
log.setLevel(logging.DEBUG)
|
||||||
log.debug(locals())
|
log.debug(locals())
|
||||||
|
|
||||||
@@ -75,8 +80,9 @@ async def websocket_endpoint(websocket: WebSocket, response: Response = Response
|
|||||||
|
|
||||||
|
|
||||||
async def redis_connector(
|
async def redis_connector(
|
||||||
websocket: WebSocket, redis_uri: str = "redis://localhost:6379"
|
websocket: WebSocket,
|
||||||
):
|
redis_url: str = "redis://localhost:6379",
|
||||||
|
):
|
||||||
log.setLevel(logging.DEBUG)
|
log.setLevel(logging.DEBUG)
|
||||||
log.debug(locals())
|
log.debug(locals())
|
||||||
|
|
||||||
@@ -108,7 +114,11 @@ async def redis_connector(
|
|||||||
# TODO this needs handling better
|
# TODO this needs handling better
|
||||||
log.error(exc)
|
log.error(exc)
|
||||||
|
|
||||||
redis = await aioredis.create_pool(redis_uri)
|
# redis = await aioredis.create_pool(redis_url)
|
||||||
|
# Redis client bound to pool of connections (auto-reconnecting).
|
||||||
|
redis = aioredis.from_url(
|
||||||
|
redis_url, encoding="utf-8", decode_responses=True
|
||||||
|
)
|
||||||
|
|
||||||
consumer_task = consumer_handler(websocket, redis)
|
consumer_task = consumer_handler(websocket, redis)
|
||||||
producer_task = producer_handler(redis, websocket)
|
producer_task = producer_handler(redis, websocket)
|
||||||
@@ -119,5 +129,5 @@ async def redis_connector(
|
|||||||
for task in pending:
|
for task in pending:
|
||||||
log.debug(f"Canceling task: {task}")
|
log.debug(f"Canceling task: {task}")
|
||||||
task.cancel()
|
task.cancel()
|
||||||
redis.close()
|
await redis.close()
|
||||||
await redis.wait_closed()
|
# await redis.wait_closed()
|
||||||
|
|||||||
Reference in New Issue
Block a user