Skip to content

Commit 71cf8d3

Browse files
authored
Merge pull request #10
feat-web-socket
2 parents 853f06d + 102c341 commit 71cf8d3

29 files changed

Lines changed: 588 additions & 98 deletions

.env.example

Lines changed: 4 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -8,12 +8,11 @@ APP_NAME="App"
88
SERVICE_NAME="api"
99
X_API_KEY=
1010

11-
# CORS Settings
12-
CORS_ALLOW_ORIGINS=["*"]
11+
CORS_ALLOW_ORIGINS='["*"]'
1312
CORS_ALLOW_CREDENTIALS=True
14-
CORS_ALLOW_METHODS=["*"]
15-
CORS_ALLOW_HEADERS=["*"]
16-
CORS_EXPOSE_HEADERS=["X-Total-Count", "X-Per-Page", "X-Current-Page", "X-Total-Pages", "X-User-Role"]
13+
CORS_ALLOW_METHODS='["*"]'
14+
CORS_ALLOW_HEADERS='["*"]'
15+
CORS_EXPOSE_HEADERS='["X-Total-Count", "X-Per-Page", "X-Current-Page", "X-Total-Pages", "X-User-Role"]'
1716

1817
SQLALCHEMY_DATABASE_URI=mysql+aiomysql://user:password@host/db_name
1918

pyproject.toml

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ dependencies = [
2323
"python-logging-loki>=0.3.1",
2424
"sqlalchemy>=2.0.41",
2525
"uvicorn-worker>=0.3.0",
26+
"uvicorn[standard]>=0.35.0",
2627
"uvloop>=0.21.0",
2728
]
2829

server.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -29,7 +29,7 @@ def load(self) -> BaseApplication:
2929
if __name__ == "__main__":
3030
options = {
3131
"bind": "0.0.0.0:8000",
32-
"workers": 1, # multiprocessing.cpu_count() * 2 + 1,
33-
"worker_class": "uvicorn.workers.UvicornWorker",
32+
"workers": 3, # multiprocessing.cpu_count() * 2 + 1,
33+
"worker_class": "uvicorn_worker.UvicornWorker",
3434
}
3535
StandaloneApplication(application=app, option=options).run()

src/api.py

Lines changed: 16 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@
1010
from src.app.auth.controller.auth_controller import AuthController
1111
from src.app.user.controller.user_controller import UserController
1212
from src.app.user_notification.controller.user_notification_controller import UserNotificationController
13+
from src.app.ws.route.ws_controller import WSController
1314
from src.core.di.container import Container
1415
from src.core.enum.env import Env
1516
from src.core.exception.error_no import ErrorNo
@@ -31,8 +32,22 @@ async def lifespan(api: FastAPI) -> AsyncGenerator[None]:
3132
AuthController(app=api, container=container)
3233
UserController(app=api, container=container)
3334
UserNotificationController(app=api, container=container)
35+
WSController(app=api, container=container)
36+
37+
# for route in api.routes:
38+
# if hasattr(route, "methods"): # HTTP
39+
# methods = ",".join(route.methods)
40+
# print(f"HTTP {methods:<10} {route.path}")
41+
# elif route.__class__.__name__ == "APIWebSocketRoute":
42+
# print(f"WS {'-':<10} {route.path}")
43+
# else:
44+
# print(f"UNKNOWN {'-':<10} {route.path}")
45+
3446
yield
3547
await container.db_config().close()
48+
await container.rmq_producer().close()
49+
await container.rmq_consumer().close()
50+
await container.ws_manager().close_all()
3651
container.unwire()
3752

3853

@@ -78,7 +93,7 @@ async def exception_handler(request: Request, e: Exception) -> JSONResponse:
7893
raise ValueError(e.args)
7994

8095
if isinstance(e, RequestValidationError):
81-
error = ApiResponseService.format_pydantic_error(errors=e.errors(), env=di.app_config().environment)
96+
error = ApiResponseService.format_pydantic_error(errors=e.errors())
8297
di.log().error(message=f"RequestValidationError: {e}", error=str(error), request=req)
8398
elif isinstance(e, DomainException):
8499
e_data = e.as_dict()

src/app/auth/controller/auth_controller.py

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -5,8 +5,7 @@
55
from src.app.auth.dto.login import LoginRequest
66
from src.app.auth.dto.re_send_confirm_email import ReSendConfirmEmailRequest
77
from src.app.auth.dto.sign_up import SignupRequest
8-
from src.cmd.worker.email.email_action import EmailAction
9-
from src.core.db.repository import Filter, FilterOperator
8+
from src.core.db.repository import Filter, Oper
109
from src.core.di.container import Container
1110
from src.core.dto.dto import Message
1211
from src.core.exception.error_no import ErrorNo
@@ -42,7 +41,7 @@ async def re_send_confirm_email(self, req: ReSendConfirmEmailRequest) -> JsonApi
4241
res = Message(message="Email successfully sent")
4342
user = await self.container.user_service().one(
4443
filters=[
45-
Filter("email", FilterOperator.EQ, req.email),
44+
Filter("email", Oper.EQ, req.email),
4645
]
4746
)
4847
if user is None:

src/app/auth/service/auth_service.py

Lines changed: 8 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@
22
from src.app.user.data.user_status import UserStatus
33
from src.app.user.model.user import User
44
from src.app.user.service.user_service import UserService
5-
from src.core.db.repository import Filter, FilterOperator
5+
from src.core.db.repository import Filter, Oper
66
from src.core.exception.error_no import ErrorNo
77
from src.core.exception.exceptions import UnauthorizedException, UnprocessableEntityException
88
from src.core.service.dto.token import TokenBearer, TokenType
@@ -24,7 +24,7 @@ def __init__(
2424
async def login(self, email: str, password: str) -> TokenBearer:
2525
user = await self.user_service.one(
2626
filters=[
27-
Filter("email", FilterOperator.EQ, email),
27+
Filter("email", Oper.EQ, email),
2828
]
2929
)
3030
if user is None:
@@ -34,7 +34,7 @@ async def login(self, email: str, password: str) -> TokenBearer:
3434

3535
if not self.hash_service.verify_password(password=password, hashed_password=user.hash_password):
3636
raise UnauthorizedException(error_no=ErrorNo.AUTHORIZATION_USER_PASSWORD_INVALID, message="Unauthorized!")
37-
user = await self.user_service.update(id_value=user.id, data={"session": self.hash_service.random_string()})
37+
user = await self.user_service.update(uid=user.id, data={"session": self.hash_service.random_string()})
3838

3939
return self.hash_service.create_token_bearer(user=user)
4040

@@ -47,7 +47,7 @@ async def signup(
4747
) -> User:
4848
user = await self.user_service.one(
4949
filters=[
50-
Filter("email", FilterOperator.EQ, email),
50+
Filter("email", Oper.EQ, email),
5151
]
5252
)
5353
if user is not None:
@@ -88,7 +88,7 @@ async def confirm_user(self, jwt: str) -> None:
8888
)
8989
user = await self.user_service.one(
9090
filters=[
91-
Filter("email", FilterOperator.EQ, token.email),
91+
Filter("email", Oper.EQ, token.email),
9292
]
9393
)
9494
if user is None:
@@ -99,7 +99,7 @@ async def confirm_user(self, jwt: str) -> None:
9999
raise UnprocessableEntityException(
100100
error_no=ErrorNo.CONFIRM_TOKEN_USER_SESSION_INVALID, message="Confirmation token is invalid"
101101
)
102-
await self.user_service.update(id_value=user.id, data={"status": UserStatus.ACTIVE})
102+
await self.user_service.update(uid=user.id, data={"status": UserStatus.ACTIVE})
103103

104104
async def refresh(self, jwt: str) -> TokenBearer:
105105
token = self.hash_service.verify_token(token=jwt)
@@ -108,7 +108,7 @@ async def refresh(self, jwt: str) -> TokenBearer:
108108
if token.token_type != TokenType.REFRESH:
109109
raise UnauthorizedException(error_no=ErrorNo.REFRESH_TOKEN_TYPE_INVALID, message="Unauthorized!")
110110

111-
user = await self.user_service.get_by_id(id_value=int(token.subject))
111+
user = await self.user_service.get_by_id(uid=int(token.subject))
112112
if user is None:
113113
raise UnauthorizedException(error_no=ErrorNo.REFRESH_TOKEN_USER_NOT_FOUND, message="Unauthorized!")
114114
if user.status != UserStatus.ACTIVE:
@@ -117,6 +117,6 @@ async def refresh(self, jwt: str) -> TokenBearer:
117117
if user.session != token.session:
118118
raise UnauthorizedException(error_no=ErrorNo.REFRESH_TOKEN_USER_SESSION_INVALID, message="Unauthorized!")
119119

120-
user = await self.user_service.update(id_value=user.id, data={"session": self.hash_service.random_string()})
120+
user = await self.user_service.update(uid=user.id, data={"session": self.hash_service.random_string()})
121121

122122
return self.hash_service.create_token_bearer(user=user)

src/app/user/controller/user_controller.py

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
from fastapi import APIRouter, Depends, FastAPI, Request
22

33
from src.app.user.dto.user import UserCreateRequest, UserListRequest
4-
from src.core.db.repository import Filter, FilterOperator, Pagination
4+
from src.core.db.repository import Filter, Oper, Pagination
55
from src.core.di.container import Container
66
from src.core.exception.error_no import ErrorNo
77
from src.core.exception.exceptions import UnprocessableEntityException
@@ -21,7 +21,7 @@ def __init__(self, app: FastAPI, container: Container) -> None:
2121
async def list(self, req: UserListRequest = Depends()) -> JsonApiResponse:
2222
users = await self.container.user_service().all(
2323
filters=[
24-
Filter("email", FilterOperator.EQ, req.email),
24+
Filter("email", Oper.EQ, req.email),
2525
],
2626
pagination=Pagination(
2727
per_page=req.per_page or 10,
@@ -41,7 +41,7 @@ async def create(
4141
) -> JsonApiResponse:
4242
user = await self.container.user_service().one(
4343
filters=[
44-
Filter("email", FilterOperator.EQ, req.email),
44+
Filter("email", Oper.EQ, req.email),
4545
]
4646
)
4747
if user is not None:

src/app/user/service/user_service.py

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -10,17 +10,17 @@ class UserService:
1010
def __init__(self, user_repository: UserRepository) -> None:
1111
self.user_repository = user_repository
1212

13-
async def get_by_id(self, id_value: int) -> User:
14-
return await self.user_repository.get_by_id(id_value=id_value)
13+
async def get_by_id(self, uid: int) -> User:
14+
return await self.user_repository.get_by_id(uid=uid)
1515

16-
async def find_by_id(self, id_value: int) -> User | None:
17-
return await self.user_repository.find_by_id(id_value=id_value)
16+
async def find_by_id(self, uid: int) -> User | None:
17+
return await self.user_repository.find_by_id(uid=uid)
1818

1919
async def create(self, data: dict[str, Any] | User) -> User:
2020
return await self.user_repository.create(data=data)
2121

22-
async def update(self, id_value: int, data: dict[str, Any]) -> User:
23-
return await self.user_repository.update(id_value=id_value, data=data)
22+
async def update(self, uid: int, data: dict[str, Any]) -> User:
23+
return await self.user_repository.update(uid=uid, data=data)
2424

2525
async def one(
2626
self,

src/app/user_notification/controller/user_notification_controller.py

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@
22

33
from src.app.user_notification.data.user_notification_status import UserNotificationStatus
44
from src.app.user_notification.dto.user_notification import UserNotificationCreateRequest, UserNotificationListRequest
5-
from src.core.db.repository import Filter, FilterOperator, Pagination
5+
from src.core.db.repository import Filter, Oper, Pagination
66
from src.core.di.container import Container
77
from src.core.http.controller import BaseController
88
from src.core.http.request.state import AuthState, get_auth_state
@@ -18,12 +18,14 @@ def __init__(self, app: FastAPI, container: Container) -> None:
1818
app.include_router(router=router)
1919

2020
async def user_list(
21-
self, state: AuthState = Depends(get_auth_state), req: UserNotificationListRequest = Depends()
21+
self,
22+
state: AuthState = Depends(get_auth_state),
23+
req: UserNotificationListRequest = Depends(),
2224
) -> JsonApiResponse:
2325
notifications = await self.container.user_notification_service().all(
2426
filters=[
25-
Filter("user_id", FilterOperator.EQ, state.user.id),
26-
Filter("status", FilterOperator.EQ, req.status),
27+
Filter("user_id", Oper.EQ, state.user.id),
28+
Filter("status", Oper.EQ, req.status),
2729
],
2830
pagination=Pagination(
2931
per_page=req.per_page or 10,

src/app/user_notification/model/user_notification.py

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,6 @@
1-
from sqlalchemy import JSON, DateTime, Integer, String, ForeignKey
1+
from datetime import datetime
2+
3+
from sqlalchemy import JSON, DateTime, ForeignKey, Integer
24
from sqlalchemy.orm import Mapped, mapped_column
35
from sqlalchemy.sql import func
46

@@ -18,11 +20,13 @@ class UserNotification(Entity):
1820
user_id: Mapped[int] = mapped_column(Integer, ForeignKey("users.id", ondelete="CASCADE"), nullable=False)
1921
data: Mapped[dict] = mapped_column(JSON, nullable=False, default=dict)
2022
status: Mapped[UserNotificationStatus] = mapped_column(
21-
IntEnum(UserNotificationStatus), nullable=False, default=UserStatus.PENDING, index=True
23+
IntEnum(UserNotificationStatus), nullable=False, default=UserNotificationStatus.NEW, index=True
2224
)
23-
updated_at: Mapped[DateTime] = mapped_column(
25+
updated_at: Mapped[datetime] = mapped_column(
2426
DateTime(timezone=True),
2527
server_default=func.now(),
2628
onupdate=func.now(),
2729
)
28-
created_at: Mapped[DateTime] = mapped_column(DateTime(timezone=True), server_default=func.now(), index=True)
30+
created_at: Mapped[datetime] = mapped_column(
31+
DateTime(timezone=True), server_default=func.now(), index=True
32+
)

0 commit comments

Comments
 (0)