160 lines
5.0 KiB
Python
160 lines
5.0 KiB
Python
import datetime
|
|
import logging
|
|
import os
|
|
from pathlib import Path
|
|
|
|
import bcrypt
|
|
import psycopg2
|
|
from dotenv import load_dotenv
|
|
from fastapi import Depends, FastAPI, HTTPException, Response, UploadFile
|
|
from fastapi.security import HTTPAuthorizationCredentials, HTTPBearer
|
|
from psycopg2.extras import execute_values
|
|
|
|
load_dotenv()
|
|
|
|
LOG_DIR = Path("/logs")
|
|
LOG_DIR.mkdir(parents=True, exist_ok=True)
|
|
|
|
logger = logging.getLogger("file_uploader")
|
|
logger.setLevel(logging.INFO)
|
|
|
|
log_file = LOG_DIR / f"{datetime.datetime.now():%Y-%m}.log"
|
|
|
|
handler = logging.FileHandler(log_file)
|
|
handler.setFormatter(
|
|
logging.Formatter(
|
|
"%(asctime)s | %(levelname)s | %(message)s"
|
|
)
|
|
)
|
|
|
|
logger.addHandler(handler)
|
|
|
|
app = FastAPI()
|
|
|
|
UPLOAD_HASH = os.environ["UPLOAD_HASH"]
|
|
UPLOAD_PATH = "uploads"
|
|
|
|
security = HTTPBearer()
|
|
|
|
DB_HOST = os.getenv('DB_HOST')
|
|
DB_USER = os.getenv('DB_USER')
|
|
DB_PWD = os.getenv('DB_PWD')
|
|
DB_PORT = os.getenv('DB_PORT')
|
|
DB_NAME = os.getenv('DB_NAME')
|
|
|
|
|
|
def connect_db():
|
|
return psycopg2.connect(f"dbname=mf user={DB_USER} password={DB_PWD} host={DB_HOST} port={DB_PORT}")
|
|
|
|
|
|
def authenticate(
|
|
credentials: HTTPAuthorizationCredentials = Depends(security),
|
|
):
|
|
token = credentials.credentials.encode()
|
|
|
|
if not bcrypt.checkpw(token, UPLOAD_HASH.encode()):
|
|
raise HTTPException(
|
|
status_code=401,
|
|
detail="Invalid authentication credentials",
|
|
)
|
|
|
|
|
|
@app.get('/test', dependencies=[Depends(authenticate)])
|
|
async def file_test():
|
|
test_file = f"/{UPLOAD_PATH}/test.txt"
|
|
with open(test_file) as f:
|
|
content = f.read()
|
|
return {'message': content}
|
|
|
|
|
|
@app.post("/upload", dependencies=[Depends(authenticate)])
|
|
async def upload(file: UploadFile):
|
|
with open(Path(f"/{UPLOAD_PATH}/{file.filename}"), "wb") as f:
|
|
while chunk := await file.read(64 * 1024):
|
|
f.write(chunk)
|
|
|
|
await file.close()
|
|
|
|
return {"status": "success"}
|
|
|
|
|
|
@app.get('/test-db', dependencies=[Depends(authenticate)])
|
|
async def test_db():
|
|
conn = connect_db()
|
|
try:
|
|
with conn.cursor() as cur:
|
|
cur.execute("select 1")
|
|
cur.fetchall()
|
|
logger.info('test-db API was called and it was successful')
|
|
return "If you're reading this, you're ready to update NAVs"
|
|
except Exception as e:
|
|
logger.exception(f'test-db API was called and it failed with the following error: {e}')
|
|
return "No bueno" + e
|
|
|
|
|
|
@app.post("/update-nav", dependencies=[Depends(authenticate)])
|
|
async def update_nav(data: list[dict]):
|
|
"""Upsert the NAV data into the nav_history table"""
|
|
|
|
columns = ['amfi_code', 'date', 'nav']
|
|
|
|
query = f"""
|
|
INSERT INTO nav_history ({', '.join(columns)})
|
|
VALUES %s
|
|
on conflict(amfi_code, date) do update set nav=excluded.nav
|
|
"""
|
|
|
|
values = [
|
|
tuple(row[col] for col in columns)
|
|
for row in data
|
|
]
|
|
conn = connect_db()
|
|
try:
|
|
with conn.cursor() as cursor:
|
|
execute_values(cursor, query, values)
|
|
conn.commit()
|
|
logger.info('Update NAV was called and db update succeeded')
|
|
|
|
with conn.cursur() as cursor:
|
|
cursor.execute('refresh materialized view latest_nav')
|
|
conn.commit()
|
|
logger.info('Latest NAV materialized view refreshed')
|
|
|
|
return {'message': "Database updated succesfully"}
|
|
except Exception as e:
|
|
logger.exception(f"Update Db errored out: {e}")
|
|
return {'message': "database update failed. Check logs for details"}
|
|
|
|
|
|
@app.get('/update-navff', dependencies=[Depends(authenticate)])
|
|
async def update_nav_ff():
|
|
conn = connect_db()
|
|
update_query = """
|
|
insert into nav_history_ff
|
|
with cte1 as (
|
|
select date_value, amfi_code
|
|
from daily_date_series dds
|
|
cross join (select amfi_code from latest_nav)
|
|
where date_value between current_date - '7 days'::interval and current_date - '1 day'::interval
|
|
), cte2 as (
|
|
select cte1.*, nh.date, nh.nav, count(nav) over (partition by cte1.amfi_code order by date_value) as grouper
|
|
from cte1
|
|
left join nav_history nh on cte1.amfi_code = nh.amfi_code and cte1.date_value = nh.date
|
|
)
|
|
select date_value as date_ff, amfi_code, first_value(date) over (partition by amfi_code, grouper order by date_value) as actual_date,
|
|
first_value( nav) over (partition by amfi_code, grouper order by date_value) as nav_ff
|
|
from cte2
|
|
ON CONFLICT (amfi_code, date_ff)
|
|
DO UPDATE SET
|
|
nav_ff = EXCLUDED.nav_ff,
|
|
actual_date = EXCLUDED.actual_date;
|
|
"""
|
|
try:
|
|
with conn.cursor() as cur:
|
|
cur.execute(update_query)
|
|
conn.commit()
|
|
logger.info('Update NAV FF was called and db update succeeded')
|
|
return {'message': "Database updated succesfully"}
|
|
except Exception as e:
|
|
logger.exception(f"Update NAV FF errored out: {e}")
|
|
return {'message': "database update failed. Check logs for details"} |