core/resolvers/topics.py

77 lines
2.4 KiB
Python
Raw Normal View History

2021-12-11 12:17:59 +00:00
from orm import Topic, TopicSubscription, Shout, User
2021-10-28 10:42:34 +00:00
from orm.base import local_session
from resolvers.base import mutation, query, subscription
2021-11-10 08:42:29 +00:00
from resolvers.zine import ShoutSubscriptions
2021-10-28 10:42:34 +00:00
from auth.authenticate import login_required
import asyncio
2021-11-27 06:15:02 +00:00
@query.field("topicsAll")
2021-11-27 07:06:20 +00:00
async def topics_all(_, info):
2021-11-27 06:15:02 +00:00
topics = []
with local_session() as session:
topics = session.query(Topic)
return topics
2021-10-28 10:42:34 +00:00
@query.field("topicsBySlugs")
async def topics_by_slugs(_, info, slugs):
topics = []
with local_session() as session:
2021-12-11 12:17:59 +00:00
topics = session.query(Topic).filter(Topic.slug.in_(slugs))
2021-10-28 10:42:34 +00:00
return topics
@query.field("topicsByCommunity")
async def topics_by_community(_, info, community):
topics = []
with local_session() as session:
topics = session.query(Topic).filter(Topic.community == community)
return topics
@query.field("topicsByAuthor")
async def topics_by_author(_, info, author):
2021-12-11 12:17:59 +00:00
topics = {}
2021-10-28 10:42:34 +00:00
with local_session() as session:
2021-12-11 12:17:59 +00:00
shouts = session.query(Shout).\
filter(Shout.authors.any(User.slug == author))
for shout in shouts:
topics.update(dict([(topic.slug, topic) for topic in shout.topics]))
return topics.values()
2021-10-28 10:42:34 +00:00
@mutation.field("topicSubscribe")
@login_required
async def topic_subscribe(_, info, slug):
auth = info.context["request"].auth
user_id = auth.user_id
sub = TopicSubscription.create({ user: user_id, topic: slug })
return {} # type Result
@mutation.field("topicUnsubscribe")
@login_required
async def topic_unsubscribe(_, info, slug):
auth = info.context["request"].auth
user_id = auth.user_id
sub = session.query(TopicSubscription).filter(TopicSubscription.user == user_id and TopicSubscription.topic == slug).first()
with local_session() as session:
session.delete(sub)
return {} # type Result
return { "error": "session error" }
2021-11-04 16:37:41 +00:00
@subscription.source("topicUpdated")
async def new_shout_generator(obj, info, user_id):
2021-11-10 08:42:29 +00:00
try:
with local_session() as session:
topics = session.query(TopicSubscription.topic).filter(TopicSubscription.user == user_id).all()
topics = set([item.topic for item in topics])
shouts_queue = asyncio.Queue()
await ShoutSubscriptions.register_subscription(shouts_queue)
while True:
shout = await shouts_queue.get()
if topics.intersection(set(shout.topic_ids)):
yield shout
finally:
await ShoutSubscriptions.del_subscription(shouts_queue)
2021-11-04 16:37:41 +00:00
@subscription.field("topicUpdated")
def shout_resolver(shout, info, user_id):
return shout