Skip to content

Commit ade29b4

Browse files
committed
feat(celery): add async publication worker and beat
1 parent 686ca77 commit ade29b4

20 files changed

Lines changed: 539 additions & 278 deletions

Makefile

Lines changed: 14 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -13,18 +13,14 @@ endef
1313
grpc-compile \
1414
grpc-server-start \
1515
fastapi-server-start \
16+
celery-worker-start \
1617
run \
1718
payload-specs-fetch \
1819
payload-specs-build \
1920
payload-specs-compile \
2021
build-setup \
2122
migrate-up
2223

23-
24-
# ---------------------------------------------------------------------------
25-
# Build
26-
# ---------------------------------------------------------------------------
27-
2824
grpc-compile:
2925
$(call log,INFO,Compiling gRPC protos ...)
3026
@for v in v1 v2 v3; do \
@@ -61,21 +57,11 @@ payload-specs-compile: payload-specs-fetch payload-specs-build
6157

6258
build-setup: grpc-compile payload-specs-compile
6359

64-
65-
# ---------------------------------------------------------------------------
66-
# Database
67-
# ---------------------------------------------------------------------------
68-
6960
migrate-up:
7061
$(call log,INFO,Running database migrations ...)
7162
@$(python) -m alembic upgrade head
7263
$(call log,INFO,Migrations complete)
7364

74-
75-
# ---------------------------------------------------------------------------
76-
# Servers
77-
# ---------------------------------------------------------------------------
78-
7965
grpc-server-start:
8066
$(call log,INFO,Starting gRPC server ...)
8167
@$(python) -u grpc_server.py
@@ -84,13 +70,17 @@ fastapi-server-start:
8470
$(call log,INFO,Starting FastAPI server ...)
8571
@$(python) -m uvicorn app:app --workers 1 --host $(grpc_host) --port $(fastapi_port)
8672

73+
celery-worker-start:
74+
$(call log,INFO,Starting Celery worker ...)
75+
@$(python) -m celery -A tasks.celery_app:celery_app worker \
76+
--loglevel=info \
77+
--without-gossip \
78+
--without-mingle \
79+
--without-heartbeat
80+
81+
celery-beat-start:
82+
$(call log,INFO,Starting Celery beat scheduler ...)
83+
@$(python) -m celery -A tasks.celery_app:celery_app beat --loglevel=info
84+
8785
run:
88-
$(call log,INFO,Starting gRPC and FastAPI servers ...)
89-
@( \
90-
$(python) -u grpc_server.py & \
91-
GRPC_PID=$$!; \
92-
$(python) -m uvicorn app:app --workers 1 --host $(grpc_host) --port $(fastapi_port) & \
93-
FASTAPI_PID=$$!; \
94-
trap 'echo "Shutting down ..."; kill $$GRPC_PID $$FASTAPI_PID 2>/dev/null; wait' INT TERM; \
95-
wait \
96-
)
86+
@PYTHON=$(python) HOST=$(grpc_host) PORT=$(fastapi_port) ./scripts/run.sh

app.py

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -5,8 +5,8 @@
55
from fastapi import FastAPI, Request
66
from fastapi.responses import JSONResponse
77

8-
from db import dispose_engine
9-
from keys import initialize_server_identity_keys
8+
from db import dispose_engine, get_session
9+
from keys import KeyManager
1010
from platforms.adapter_manager import AdapterManager
1111
from rest_services.v1.routes import router as v1_router
1212
from utils import get_logger
@@ -17,7 +17,10 @@
1717
@asynccontextmanager
1818
async def lifespan(app: FastAPI):
1919
"""Handle application startup and shutdown."""
20-
initialize_server_identity_keys()
20+
with get_session() as db:
21+
key_manager = KeyManager(session=db)
22+
key_manager.initialize_server_identity_keys()
23+
2124
app.state.adapter_manager = AdapterManager()
2225
yield
2326
dispose_engine()

db.py

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -136,6 +136,10 @@ def _has_mysql_config() -> bool:
136136
return all(get_configs(key) for key in required)
137137

138138

139+
def _sql_echo_enabled() -> bool:
140+
return get_configs("LOG_LEVEL", default_value="info").lower() == "debug"
141+
142+
139143
def _create_engine() -> Engine:
140144
"""Create and configure database engine."""
141145
mode = get_configs("MODE", default_value="development")
@@ -169,7 +173,7 @@ def _create_engine() -> Engine:
169173
engine = create_engine(
170174
"sqlite://",
171175
creator=_make_sqlcipher3_creator(db_path, key),
172-
echo=False,
176+
echo=_sql_echo_enabled(),
173177
poolclass=QueuePool,
174178
pool_size=5,
175179
pool_pre_ping=True,
@@ -178,7 +182,7 @@ def _create_engine() -> Engine:
178182
url = _build_sqlite_url()
179183
engine = create_engine(
180184
url,
181-
echo=False,
185+
echo=_sql_echo_enabled(),
182186
connect_args={"check_same_thread": False, "timeout": 30},
183187
poolclass=QueuePool,
184188
pool_size=5,

grpc_server.py

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -10,9 +10,9 @@
1010
import grpc
1111
from grpc_interceptor import ServerInterceptor
1212

13-
from db import dispose_engine
13+
from db import dispose_engine, get_session
1414
from grpc_services.v3.service import PublisherServiceV3
15-
from keys import initialize_server_identity_keys
15+
from keys import KeyManager
1616
from logutils import get_logger
1717
from platforms.adapter_manager import AdapterManager
1818
from protos.v3 import publisher_pb2_grpc as v3_grpc
@@ -65,7 +65,10 @@ def _build_server(max_workers: int) -> grpc.Server:
6565
interceptors=[LoggingInterceptor()],
6666
)
6767

68-
initialize_server_identity_keys()
68+
with get_session() as db:
69+
key_manager = KeyManager(session=db)
70+
key_manager.initialize_server_identity_keys()
71+
6972
PublisherServiceV3.adapter_manager = AdapterManager()
7073
v3_grpc.add_PublisherServicer_to_server(PublisherServiceV3(), grpc_server)
7174

grpc_services/v3/revoke_oauth2_token.py

Lines changed: 21 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -7,10 +7,9 @@
77
import grpc
88

99
from db import get_session
10-
from keys import get_keys_for_decryption
10+
from keys import KeyManager, KeyManagerError
1111
from lib_relaysms_payload_specs.generated import relaysms_spec_payload as rrs
1212
from logutils import get_logger
13-
from models.server_identity_key import mark_key_used as mark_ss_key_used
1413
from platforms.adapter_ipc_handler import AdapterIPCHandler
1514
from protos.v3 import publisher_pb2
1615

@@ -26,7 +25,7 @@ def RevokeOAuth2Token(self, request, context):
2625
return auth_error
2726

2827
if not payload_bin:
29-
logger.error("missing token ciphertext in request payload")
28+
logger.error("Missing token ciphertext in revoke request")
3029
return self.handle_create_grpc_error_response(
3130
context,
3231
response,
@@ -42,9 +41,10 @@ def RevokeOAuth2Token(self, request, context):
4241

4342
try:
4443
with get_session() as s:
44+
key_manager = KeyManager(s)
4545
token, token_hash_obj, ss_kid, es_kid, _, ec_kid_pk = (
46-
get_keys_for_decryption(
47-
token_id=request.token_id, key_id=request.key_id, session=s
46+
key_manager.get_token_and_keys_for_decryption(
47+
token_id=request.token_id, key_id=request.key_id
4848
)
4949
)
5050

@@ -56,21 +56,23 @@ def RevokeOAuth2Token(self, request, context):
5656
key_id=request.key_id,
5757
ciphertext=payload_bin,
5858
)
59-
except rrs.V1CryptographicError.FailedToDecrypt as e:
59+
except rrs.V1CryptographicError.FailedToDecrypt:
60+
logger.exception("Token decryption failed: kid=%s", request.key_id)
6061
return self.handle_create_grpc_error_response(
6162
context,
6263
response,
63-
f"token decryption failed: {e}",
64+
"revocation failed",
6465
grpc.StatusCode.UNAUTHENTICATED,
6566
)
6667

6768
if not secrets.compare_digest(
6869
hashlib.sha256(decrypted_token).digest(), token_hash_obj.token_hash
6970
):
71+
logger.warning("Token hash mismatch: kid=%s", request.key_id)
7072
return self.handle_create_grpc_error_response(
7173
context,
7274
response,
73-
"token hash mismatch",
75+
"revocation failed",
7476
grpc.StatusCode.UNAUTHENTICATED,
7577
)
7678

@@ -87,21 +89,26 @@ def RevokeOAuth2Token(self, request, context):
8789
)
8890

8991
if pipe.get("error"):
90-
logger.error("adapter revocation failed: %s", pipe["error"])
92+
logger.error(
93+
"Adapter revocation failed for platform %r: %s",
94+
token.platform,
95+
pipe["error"],
96+
)
9197

9298
s.delete(token)
93-
mark_ss_key_used(request.key_id, s)
99+
key_manager.mark_identity_key_used(request.key_id)
100+
logger.info("Token revoked: platform=%r", token.platform)
94101

95102
return response(success=True, message="Successfully revoked and deleted token")
96103

97-
except ValueError as exc:
104+
except NotImplementedError as exc:
98105
return self.handle_create_grpc_error_response(
99-
context, response, exc, grpc.StatusCode.NOT_FOUND
106+
context, response, exc, grpc.StatusCode.UNIMPLEMENTED
100107
)
101108

102-
except NotImplementedError as exc:
109+
except KeyManagerError:
103110
return self.handle_create_grpc_error_response(
104-
context, response, exc, grpc.StatusCode.UNIMPLEMENTED
111+
context, response, "revocation failed", grpc.StatusCode.UNAUTHENTICATED
105112
)
106113

107114
except Exception as exc:

grpc_services/v3/revoke_pnba_token.py

Lines changed: 21 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -7,10 +7,9 @@
77
import grpc
88

99
from db import get_session
10-
from keys import get_keys_for_decryption
10+
from keys import KeyManager, KeyManagerError
1111
from lib_relaysms_payload_specs.generated import relaysms_spec_payload as rrs
1212
from logutils import get_logger
13-
from models.server_identity_key import mark_key_used as mark_ss_key_used
1413
from platforms.adapter_ipc_handler import AdapterIPCHandler
1514
from protos.v3 import publisher_pb2
1615

@@ -26,7 +25,7 @@ def RevokePNBAToken(self, request, context):
2625
return auth_error
2726

2827
if not payload_bin:
29-
logger.error("missing token ciphertext in request payload")
28+
logger.error("Missing token ciphertext in revoke request")
3029
return self.handle_create_grpc_error_response(
3130
context,
3231
response,
@@ -42,9 +41,10 @@ def RevokePNBAToken(self, request, context):
4241

4342
try:
4443
with get_session() as s:
44+
key_manager = KeyManager(s)
4545
token, token_hash_obj, ss_kid, es_kid, _, ec_kid_pk = (
46-
get_keys_for_decryption(
47-
token_id=request.token_id, key_id=request.key_id, session=s
46+
key_manager.get_token_and_keys_for_decryption(
47+
token_id=request.token_id, key_id=request.key_id
4848
)
4949
)
5050

@@ -56,21 +56,23 @@ def RevokePNBAToken(self, request, context):
5656
key_id=request.key_id,
5757
ciphertext=payload_bin,
5858
)
59-
except rrs.V1CryptographicError.FailedToDecrypt as e:
59+
except rrs.V1CryptographicError.FailedToDecrypt:
60+
logger.exception("Token decryption failed: kid=%s", request.key_id)
6061
return self.handle_create_grpc_error_response(
6162
context,
6263
response,
63-
f"token decryption failed: {e}",
64+
"revocation failed",
6465
grpc.StatusCode.UNAUTHENTICATED,
6566
)
6667

6768
if not secrets.compare_digest(
6869
hashlib.sha256(decrypted_token).digest(), token_hash_obj.token_hash
6970
):
71+
logger.error("Token hash mismatch: kid=%s", request.key_id)
7072
return self.handle_create_grpc_error_response(
7173
context,
7274
response,
73-
"token hash mismatch",
75+
"revocation failed",
7476
grpc.StatusCode.UNAUTHENTICATED,
7577
)
7678

@@ -87,21 +89,26 @@ def RevokePNBAToken(self, request, context):
8789
)
8890

8991
if pipe.get("error"):
90-
logger.error("adapter revocation failed: %s", pipe["error"])
92+
logger.error(
93+
"Adapter revocation failed for platform %r: %s",
94+
token.platform,
95+
pipe["error"],
96+
)
9197

9298
s.delete(token)
93-
mark_ss_key_used(request.key_id, s)
99+
key_manager.mark_identity_key_used(request.key_id)
100+
logger.info("Token revoked: platform=%r", token.platform)
94101

95102
return response(success=True, message="Successfully revoked and deleted token")
96103

97-
except ValueError as exc:
104+
except NotImplementedError as exc:
98105
return self.handle_create_grpc_error_response(
99-
context, response, exc, grpc.StatusCode.NOT_FOUND
106+
context, response, exc, grpc.StatusCode.UNIMPLEMENTED
100107
)
101108

102-
except NotImplementedError as exc:
109+
except KeyManagerError:
103110
return self.handle_create_grpc_error_response(
104-
context, response, exc, grpc.StatusCode.UNIMPLEMENTED
111+
context, response, "revocation failed", grpc.StatusCode.UNAUTHENTICATED
105112
)
106113

107114
except Exception as exc:

0 commit comments

Comments
 (0)