-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathserver_persistent.py
More file actions
100 lines (78 loc) · 2.95 KB
/
Copy pathserver_persistent.py
File metadata and controls
100 lines (78 loc) · 2.95 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
import time
from flask import Blueprint, jsonify, request
from db.config import async_session, engine, Base
from db.AsyncDAL import DAL
server = Blueprint("server",__name__)
topic_lock = True
log_lock = True
@server.before_app_first_request
async def setup_db():
async with engine.begin() as conn:
# await conn.run_sync(Base.metadata.drop_all)
await conn.run_sync(Base.metadata.create_all, checkfirst=True)
global topic_lock
global log_lock
topic_lock, log_lock = True, True
@server.route("/")
def index():
return "<h1>Welcome to the Persistent Distributed server!</h1>"
@server.route("/status")
def status():
return jsonify({"status": "success"})
@server.route("/topics", methods=["POST"])
async def create_topic():
topic_name = request.json["topic_name"]
async with async_session() as session, session.begin():
db_dal = DAL(session)
return await db_dal.create_topic(topic_name)
@server.route("/topics", methods=["GET"])
async def list_topics():
async with async_session() as session, session.begin():
db_dal = DAL(session)
return await db_dal.list_topics()
@server.route("/consumer/register", methods=["POST"])
async def register_consumer():
topic_name = request.json["topic_name"]
async with async_session() as session, session.begin():
db_dal = DAL(session)
return await db_dal.register_consumer(topic_name)
@server.route("/producer/register", methods=["POST"])
async def register_producer():
topic_name = request.json["topic_name"]
global topic_lock
while not topic_lock:
time.sleep(1)
topic_lock = False
async with async_session() as session, session.begin():
db_dal = DAL(session)
respose = await db_dal.register_producer(topic_name)
topic_lock = True
return respose
@server.route("/producer/produce", methods=["POST"])
async def enqueue():
global log_lock
while not log_lock:
time.sleep(1)
log_lock = False
topic_name = request.json["topic_name"]
producer_id = request.json["producer_id"]
log_message = request.json["log_message"]
async with async_session() as session, session.begin():
db_dal = DAL(session)
responese = await db_dal.enqueue(topic_name, producer_id, log_message)
log_lock = True
return responese
@server.route("/consumer/consume", methods=["GET"])
async def dequeue():
topic_name = request.json["topic_name"]
consumer_id = request.json["consumer_id"]
async with async_session() as session, session.begin():
db_dal = DAL(session)
return await db_dal.dequeue(topic_name, consumer_id)
@server.route("/size", methods=["GET"])
async def size():
topic_name = request.json["topic_name"]
consumer_id = request.json["consumer_id"]
async with async_session() as session, session.begin():
db_dal = DAL(session)
return await db_dal.size(topic_name, consumer_id)