Skip to content

Don't work remove_listener(channel, callback) #236

Description

@Belanchuk
  • asyncpg version: 0.13.0
  • PostgreSQL version: PostgreSQL 9.5.10 on x86_64-pc-linux-gnu
  • Python version: Python 3.6.3 :: Anaconda, Inc.
  • Platform: Ubuntu 5.4.0-6ubuntu1~16.04.4
  • Do you use pgbouncer?: no
  • Did you install asyncpg with pip?: yes (Anaconda)
  • Can the issue be reproduced under both asyncio and
    uvloop?
    : yes

Доброго дня!
Хочу управлять оповещениями Postgresql NOTIFY через websocket.
В качестве сервера использую aiohttp.
Успешно подключаюсь с вебстранички к серверу и отправляю по websocket'у {'subscribe':'channel_name'}
Подписка работает и уведомления соответствующего канала отправляются на веб-страничку.
Когда я хочу отписаться, я отправляю с веб-странички по websocket'у {'unsubscribe':'channel_name'}, но срабатывает только print('unsub') и подписка продолжает работать.

Буду очень признателен за помощь.

Код сервера:

import asyncio
# import uvloop
# asyncio.set_event_loop_policy(uvloop.EventLoopPolicy())
import asyncpg
import ujson
import aiohttp
from aiohttp import web
import routes


async def subscribe(data, ws, conn):

    def callback_listener(connection, pid, channel, payload):
        print(payload)
        ws.send_str(payload)

    if 'subscribe' in data.keys():
        print('sub')
        await conn.add_listener(data['subscribe'], callback_listener)                    
    if 'unsubscribe' in data.keys():
        print('unsub')
        await conn.remove_listener(data['unsubscribe'], callback_listener) 


async def websocket_handler(request):

    ws = web.WebSocketResponse()
    await ws.prepare(request)

    async for msg in ws:
        if msg.type == aiohttp.WSMsgType.TEXT:
            if msg.data == 'close':
                await ws.close()
            else:
                data = ujson.loads(msg.data)
                await subscribe(data, ws, request.app['db'])

        elif msg.type == aiohttp.WSMsgType.ERROR:
            print('ws connection closed with exception %s' %
                  ws.exception())

    print('websocket connection closed')
    return ws


async def init_pg(app):
    conn = await asyncpg.connect(host='localhost', database='db', user='user', password='pass')    
    app['db'] = conn
    print('init_pg')


async def close_pg(app):
    await app['db'].close()
    print('close_pg')


app = web.Application()
app.router.add_get('/ws', handler=websocket_handler, name='ws')
app.on_startup.append(init_pg)
app.on_cleanup.append(close_pg)
web.run_app(app, host='127.0.0.1', port=8080)

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions