Compare commits

...

2 Commits

Author SHA1 Message Date
Ivan Vazhenin 87cb86a7f9 Add put_object operation 2023-09-19 20:23:13 +03:00
Ivan Vazhenin e0fdd25ca4 Add active checking 2023-09-18 19:32:13 +03:00
4 changed files with 98 additions and 5 deletions
+9 -1
View File
@@ -1,6 +1,6 @@
from datetime import datetime
from typing import List, Optional
from sqlalchemy import create_engine, String, select, ForeignKey
from sqlalchemy import create_engine, String, select, ForeignKey, Enum
from sqlalchemy.orm import Session, DeclarativeBase, Mapped, mapped_column, relationship
@@ -124,5 +124,13 @@ class Queue(Base):
}
class Schemas(Base):
__tablename__ = 'schemas'
id: Mapped[int] = mapped_column(primary_key=True)
schema: Mapped[str]
schema_type: Mapped[str]
def connect_db():
return create_engine("postgresql+psycopg://postgres:Root12345678@10.10.8.83:32101/db")
+31
View File
@@ -0,0 +1,31 @@
from db import Base, User, connect_db
from sqlalchemy.orm import Session
def init():
sync_engine = connect_db()
Base.metadata.drop_all(sync_engine)
Base.metadata.create_all(sync_engine)
with Session(sync_engine) as session:
bnd127 = User(
username="bnd127",
passwd="gost_2012$a742ec53198ec2a5027086fba8814a89982a57112d1a72d02260161108f39b50",
bndname="bnd127",
newbnd=True,
active=True,
upstream=False
)
bnd128 = User(
username="bnd128",
passwd="gost_2012$a742ec53198ec2a5027086fba8814a89982a57112d1a72d02260161108f39b50",
bndname="bnd128",
newbnd=True,
active=True,
upstream=False
)
session.add_all([bnd127, bnd128])
session.commit()
if __name__ == '__main__':
init()
+55 -4
View File
@@ -25,6 +25,9 @@ from zip import Zip
import boto3
import db
from sqlalchemy.orm import Session
from fastapi import FastAPI, Response, Form, UploadFile, File
import uvicorn
from typing import Annotated
NEW_REPLICATION_REQUEST = 99
@@ -341,6 +344,10 @@ def get_metadata(params, files, url):
proxy.send(response_params, response_files, Config.ret_path)
def put_object(params, files, url):
pass
def apply_commits(params, files, url):
logger.warning(params, files, url)
assert len(files) == 1
@@ -376,6 +383,8 @@ def run_task(query_type, params, files, url):
tasks.put(lambda: get_objects(params, files, url))
if query_type == 24:
tasks.put(lambda: get_metadata(params, files, url))
if query_type == 6:
tasks.put(lambda: put_object(params, files, url))
if query_type == NEW_REPLICATION_REQUEST:
tasks.put(lambda: apply_commits(params, files, url))
@@ -407,9 +416,7 @@ def bnd_disconnected(bnd_name: str):
connected.remove(bnd_name)
def main():
logger.setLevel(logging.INFO)
logger.warning('Use Control-C to exit')
def xmlrpc_task():
server = SimpleXMLRPCServer(('0.0.0.0', 9000), logRequests=False, allow_none=True)
server.register_function(list_contents)
server.register_function(aud_add)
@@ -420,6 +427,50 @@ def main():
server.register_function(onDelivered)
server.register_function(bnd_connected)
server.register_function(bnd_disconnected)
server.serve_forever()
app = FastAPI()
@app.get("/")
async def fa_root():
return {"message": "Hello World"}
@app.post("/login")
async def fa_login(response: Response):
response.set_cookie(key='sessionid', value='87654321')
response.set_cookie(key='csrftoken', value='12345678')
return {"login": "Ok"}
@app.post("/api/easo/PutObject")
async def fa_put_object(response: Response,
object_attrs: Annotated[str, Form()],
object_file: Annotated[UploadFile, File()]
):
date = datetime.datetime.now()
files = [{'name': object_file.filename, 'url': object_file.file.name, 'size': os.path.getsize(object_file.size)}]
request_params = {
'from': '',
'to': '',
'ts_added': date.timestamp(),
'user_id': '0',
'user_id_to': '0',
'query_type': 6,
'query_data': object_attrs,
}
put_object(request_params, files, None)
return {"Upload": "Ok"}
def main():
logger.setLevel(logging.INFO)
logger.warning('Use Control-C to exit')
xmlrpc_thread = threading.Thread(target=xmlrpc_task)
xmlrpc_thread.start()
thread = threading.Thread(target=run_tasks)
thread.start()
@@ -432,7 +483,7 @@ def main():
try:
logger.warning('Start server')
server.serve_forever()
uvicorn.run(app, host="0.0.0.0", port=9001)
except KeyboardInterrupt:
logger.warning('Exiting')
+3
View File
@@ -23,3 +23,6 @@ typing_extensions==4.7.1
urllib3==2.0.4
yarl==1.9.2
SQLAlchemy~=2.0.20
fastapi=~0.103.1
uvicorn=~0.23.2
python-multipart=~0.0.6