Compare commits
14 Commits
1b629d8ea3
...
main
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1f5548ebb0 | ||
|
|
8d50163c72 | ||
|
|
58448f5faa | ||
|
|
661a08d0b7 | ||
|
|
d98ebdd152 | ||
|
|
29751d2d3d | ||
|
|
4395828ac7 | ||
|
|
0cd4d73f9b | ||
|
|
d195556f31 | ||
|
|
aecf42085f | ||
|
|
24dcc41720 | ||
|
|
bab21c2403 | ||
|
|
cde7d4b84b | ||
|
|
dd5dfe10f7 |
0
bootstrap_server.sh
Normal file → Executable file
0
bootstrap_server.sh
Normal file → Executable file
@@ -1,14 +0,0 @@
|
|||||||
apiVersion: kustomize.toolkit.fluxcd.io/v1
|
|
||||||
kind: Kustomization
|
|
||||||
metadata:
|
|
||||||
name: alaska
|
|
||||||
namespace: flux-system
|
|
||||||
spec:
|
|
||||||
interval: 10m
|
|
||||||
sourceRef:
|
|
||||||
kind: GitRepository
|
|
||||||
name: flux-system
|
|
||||||
path: ./manifests/alaska
|
|
||||||
prune: true
|
|
||||||
wait: true
|
|
||||||
timeout: 3m
|
|
||||||
@@ -1,14 +0,0 @@
|
|||||||
apiVersion: kustomize.toolkit.fluxcd.io/v1
|
|
||||||
kind: Kustomization
|
|
||||||
metadata:
|
|
||||||
name: core
|
|
||||||
namespace: flux-system
|
|
||||||
spec:
|
|
||||||
interval: 10m
|
|
||||||
sourceRef:
|
|
||||||
kind: GitRepository
|
|
||||||
name: flux-system
|
|
||||||
path: ./manifests/core
|
|
||||||
prune: true
|
|
||||||
wait: true
|
|
||||||
timeout: 3m
|
|
||||||
@@ -1,14 +0,0 @@
|
|||||||
apiVersion: kustomize.toolkit.fluxcd.io/v1
|
|
||||||
kind: Kustomization
|
|
||||||
metadata:
|
|
||||||
name: nearle
|
|
||||||
namespace: flux-system
|
|
||||||
spec:
|
|
||||||
interval: 10m
|
|
||||||
sourceRef:
|
|
||||||
kind: GitRepository
|
|
||||||
name: flux-system
|
|
||||||
path: ./manifests/nearle
|
|
||||||
prune: true
|
|
||||||
wait: true
|
|
||||||
timeout: 5m
|
|
||||||
File diff suppressed because it is too large
Load Diff
@@ -1,27 +0,0 @@
|
|||||||
# This manifest was generated by flux. DO NOT EDIT.
|
|
||||||
---
|
|
||||||
apiVersion: source.toolkit.fluxcd.io/v1
|
|
||||||
kind: GitRepository
|
|
||||||
metadata:
|
|
||||||
name: flux-system
|
|
||||||
namespace: flux-system
|
|
||||||
spec:
|
|
||||||
interval: 1m0s
|
|
||||||
ref:
|
|
||||||
branch: main
|
|
||||||
secretRef:
|
|
||||||
name: flux-system
|
|
||||||
url: https://gitapp.workolik.com/Nearle/kubernetes.git
|
|
||||||
---
|
|
||||||
apiVersion: kustomize.toolkit.fluxcd.io/v1
|
|
||||||
kind: Kustomization
|
|
||||||
metadata:
|
|
||||||
name: flux-system
|
|
||||||
namespace: flux-system
|
|
||||||
spec:
|
|
||||||
interval: 10m0s
|
|
||||||
path: ./clusters/production
|
|
||||||
prune: true
|
|
||||||
sourceRef:
|
|
||||||
kind: GitRepository
|
|
||||||
name: flux-system
|
|
||||||
@@ -1,5 +0,0 @@
|
|||||||
apiVersion: kustomize.config.k8s.io/v1beta1
|
|
||||||
kind: Kustomization
|
|
||||||
resources:
|
|
||||||
- gotk-components.yaml
|
|
||||||
- gotk-sync.yaml
|
|
||||||
@@ -1,20 +0,0 @@
|
|||||||
# The referenced Secret (gitea-webhook-token) is created directly on the
|
|
||||||
# cluster, not committed here - see the server-side command sequence.
|
|
||||||
apiVersion: notification.toolkit.fluxcd.io/v1
|
|
||||||
kind: Receiver
|
|
||||||
metadata:
|
|
||||||
name: gitea-receiver
|
|
||||||
namespace: flux-system
|
|
||||||
spec:
|
|
||||||
# This Flux version's Receiver CRD has no "gitea" type - valid values are
|
|
||||||
# generic, generic-hmac, generic-oidc, github, gitlab, bitbucket, harbor,
|
|
||||||
# dockerhub, quay, gcr, nexus, acr, cdevents. "generic" accepts any POST
|
|
||||||
# to the hook path without payload/signature parsing, which is fine here:
|
|
||||||
# worst case is an early sync trigger, not a permission escalation.
|
|
||||||
type: generic
|
|
||||||
secretRef:
|
|
||||||
name: gitea-webhook-token
|
|
||||||
resources:
|
|
||||||
- apiVersion: source.toolkit.fluxcd.io/v1
|
|
||||||
kind: GitRepository
|
|
||||||
name: flux-system
|
|
||||||
@@ -1,19 +0,0 @@
|
|||||||
# Exposes Flux's webhook-receiver directly via NodePort, bypassing
|
|
||||||
# Traefik/Ingress/DNS entirely - no A record needed. Selector matches the
|
|
||||||
# notification-controller pod labels from Flux's own bundled manifests; if
|
|
||||||
# `kubectl get endpoints webhook-receiver-external -n flux-system` comes up
|
|
||||||
# empty after this applies, check the pod's real labels with
|
|
||||||
# `kubectl get pods -n flux-system --show-labels` and fix the selector below.
|
|
||||||
apiVersion: v1
|
|
||||||
kind: Service
|
|
||||||
metadata:
|
|
||||||
name: webhook-receiver-external
|
|
||||||
namespace: flux-system
|
|
||||||
spec:
|
|
||||||
type: NodePort
|
|
||||||
selector:
|
|
||||||
app: notification-controller
|
|
||||||
ports:
|
|
||||||
- port: 80
|
|
||||||
targetPort: 9292
|
|
||||||
protocol: TCP
|
|
||||||
3
deploy-alaska.sh
Normal file → Executable file
3
deploy-alaska.sh
Normal file → Executable file
@@ -7,8 +7,7 @@ set -euo pipefail
|
|||||||
NAMESPACE="alaska"
|
NAMESPACE="alaska"
|
||||||
|
|
||||||
echo "🚀 Deploying Alaska Stack..."
|
echo "🚀 Deploying Alaska Stack..."
|
||||||
kubectl apply -f manifests/alaska/alaska.yaml
|
kubectl apply -k manifests/alaska
|
||||||
kubectl apply -f manifests/alaska/k8s-dashboard.yaml
|
|
||||||
|
|
||||||
echo ""
|
echo ""
|
||||||
echo "✅ Alaska stack deployment applied."
|
echo "✅ Alaska stack deployment applied."
|
||||||
|
|||||||
0
deploy-core-stack.sh
Normal file → Executable file
0
deploy-core-stack.sh
Normal file → Executable file
28
deploy-nearle-stack.sh
Normal file → Executable file
28
deploy-nearle-stack.sh
Normal file → Executable file
@@ -6,32 +6,8 @@ set -euo pipefail
|
|||||||
|
|
||||||
NAMESPACE="nearle"
|
NAMESPACE="nearle"
|
||||||
|
|
||||||
echo "🔎 Ensuring namespace '${NAMESPACE}' exists..."
|
echo "🚀 Deploying Nearle Stack..."
|
||||||
kubectl apply -f manifests/nearle/nearle-namespace.yaml
|
kubectl apply -k manifests/nearle
|
||||||
|
|
||||||
echo "🔐 Deploying Configs & Secrets..."
|
|
||||||
kubectl apply -f manifests/nearle/nearle-secrets.yaml
|
|
||||||
kubectl apply -f manifests/nearle/nearle-app-secrets.yaml
|
|
||||||
kubectl apply -f manifests/nearle/nearle-config.yaml
|
|
||||||
kubectl apply -f manifests/nearle/fiesta-gateway.yaml
|
|
||||||
|
|
||||||
echo "🚀 Deploying Services..."
|
|
||||||
kubectl apply -f manifests/nearle/jupiter-sts.yaml
|
|
||||||
kubectl apply -f manifests/nearle/jupiter-svc.yaml
|
|
||||||
|
|
||||||
kubectl apply -f manifests/nearle/nearle-titan.yaml
|
|
||||||
|
|
||||||
kubectl apply -f manifests/nearle/fiesta-sts.yaml
|
|
||||||
kubectl apply -f manifests/nearle/fiesta-svc.yaml
|
|
||||||
|
|
||||||
kubectl apply -f manifests/nearle/nearle-ariane.yaml
|
|
||||||
|
|
||||||
kubectl apply -f manifests/nearle/atlantis-sts.yaml
|
|
||||||
kubectl apply -f manifests/nearle/atlantis-svc.yaml
|
|
||||||
|
|
||||||
echo "🌐 Deploying Gateway Routes..."
|
|
||||||
kubectl apply -f manifests/nearle/nearle-gateway.yaml
|
|
||||||
kubectl apply -f manifests/nearle/nearle-reference-grant.yaml
|
|
||||||
|
|
||||||
echo ""
|
echo ""
|
||||||
echo "✅ Nearle stack deployment applied."
|
echo "✅ Nearle stack deployment applied."
|
||||||
|
|||||||
@@ -21,6 +21,16 @@ services:
|
|||||||
- "traefik.http.routers.queue-api.tls.certresolver=letsencrypt"
|
- "traefik.http.routers.queue-api.tls.certresolver=letsencrypt"
|
||||||
- "traefik.http.routers.queue-api.priority=100"
|
- "traefik.http.routers.queue-api.priority=100"
|
||||||
- "traefik.http.services.queue-api.loadbalancer.server.port=8201"
|
- "traefik.http.services.queue-api.loadbalancer.server.port=8201"
|
||||||
|
# Edge rate limit (per source IP, Traefik default). Predates this repo -
|
||||||
|
# the running container (created 2026-05-30) had average=150/burst=50
|
||||||
|
# baked in at creation time, undocumented here since labels aren't
|
||||||
|
# re-read on an existing container. Raised to fit real bulk-order
|
||||||
|
# traffic: NATS + the worker pool absorb bursts fine once a request
|
||||||
|
# reaches the gateway, so this only needs to block actual flood-scale
|
||||||
|
# abuse, not legitimate customers batching orders (2026-07-30).
|
||||||
|
- "traefik.http.middlewares.queue-rl.ratelimit.average=300"
|
||||||
|
- "traefik.http.middlewares.queue-rl.ratelimit.burst=400"
|
||||||
|
- "traefik.http.routers.queue-api.middlewares=queue-rl"
|
||||||
|
|
||||||
# Kubernetes dashboard proxy (kube.workolik.com ? K8s dashboard)
|
# Kubernetes dashboard proxy (kube.workolik.com ? K8s dashboard)
|
||||||
k8s-dashboard-proxy:
|
k8s-dashboard-proxy:
|
||||||
|
|||||||
@@ -49,14 +49,16 @@ kubectl get svc -n kubernetes-dashboard
|
|||||||
|
|
||||||
### Login to Dashboard
|
### Login to Dashboard
|
||||||
|
|
||||||
The dashboard is configured with `--enable-skip-login`, so you can skip the login screen. However, if you need to authenticate:
|
The dashboard is configured with `--enable-skip-login`, so it opens straight to the UI - no token needed to look around.
|
||||||
|
|
||||||
1. Get the token:
|
Skip-login runs as the dashboard's own `kubernetes-dashboard` ServiceAccount, which is **view-only** (get/list/watch). If you need to edit, delete, or exec into something:
|
||||||
|
|
||||||
|
1. Get an admin token:
|
||||||
```bash
|
```bash
|
||||||
kubectl -n kubernetes-dashboard create token admin-user
|
kubectl -n kubernetes-dashboard create token admin-user
|
||||||
```
|
```
|
||||||
|
|
||||||
2. Copy the token and paste it in the dashboard login screen.
|
2. Click "Sign In" on the dashboard and paste the token.
|
||||||
|
|
||||||
### What You Can See
|
### What You Can See
|
||||||
|
|
||||||
|
|||||||
@@ -62,6 +62,13 @@ spec:
|
|||||||
args:
|
args:
|
||||||
- --auto-generate-certificates
|
- --auto-generate-certificates
|
||||||
- --namespace=kubernetes-dashboard
|
- --namespace=kubernetes-dashboard
|
||||||
|
- --enable-skip-login
|
||||||
|
# Skip-login uses the "kubernetes-dashboard" ServiceAccount below,
|
||||||
|
# which only has get/list/watch (view-only) - so opening the
|
||||||
|
# dashboard needs no token, but it can't edit/delete/exec.
|
||||||
|
# For write access, still log in with the admin-user token
|
||||||
|
# (kubectl -n kubernetes-dashboard create token admin-user).
|
||||||
|
#
|
||||||
# --token-ttl=0 was tried here to disable the 15-min idle
|
# --token-ttl=0 was tried here to disable the 15-min idle
|
||||||
# timeout, but login broke immediately after that pod came up -
|
# timeout, but login broke immediately after that pod came up -
|
||||||
# in this dashboard version, 0 appears to mean "expire
|
# in this dashboard version, 0 appears to mean "expire
|
||||||
|
|||||||
@@ -43,20 +43,13 @@ spec:
|
|||||||
- host: fiesta.nearle.app
|
- host: fiesta.nearle.app
|
||||||
http:
|
http:
|
||||||
paths:
|
paths:
|
||||||
- path: /live/api/v1/mob/orders/createorder
|
# Real routing for this host is the standalone Docker/nginx proxy
|
||||||
pathType: Prefix
|
# (conf/nginx-fiesta.conf), which forwards everything to NodePort
|
||||||
backend:
|
# 30823 (this backend, port 80) - no path splitting, no gateway
|
||||||
service:
|
# sidecar. This Ingress has no controller on the cluster (confirmed:
|
||||||
name: fiesta
|
# no Traefik/Envoy pod, no IngressClass) so it's inert either way,
|
||||||
port:
|
# but kept as a single correct catch-all rather than a stale
|
||||||
number: 8000
|
# path-split that pointed part of it at a dead gateway sidecar.
|
||||||
- path: /live/api/v1/web/products/create
|
|
||||||
pathType: Prefix
|
|
||||||
backend:
|
|
||||||
service:
|
|
||||||
name: fiesta
|
|
||||||
port:
|
|
||||||
number: 8000
|
|
||||||
- path: /
|
- path: /
|
||||||
pathType: Prefix
|
pathType: Prefix
|
||||||
backend:
|
backend:
|
||||||
|
|||||||
@@ -105,7 +105,7 @@ spec:
|
|||||||
key: api_key
|
key: api_key
|
||||||
optional: true
|
optional: true
|
||||||
- name: EXTERNAL_BASE_URL
|
- name: EXTERNAL_BASE_URL
|
||||||
value: "http://10.43.224.63"
|
value: "http://jupiter.nearle"
|
||||||
resources:
|
resources:
|
||||||
requests:
|
requests:
|
||||||
memory: "128Mi"
|
memory: "128Mi"
|
||||||
|
|||||||
@@ -48,6 +48,20 @@ data:
|
|||||||
# waiting on a slow-but-legitimate response, causing a duplicate forward.
|
# waiting on a slow-but-legitimate response, causing a duplicate forward.
|
||||||
ACK_WAIT_SECONDS = int(os.getenv("ACK_WAIT_SECONDS", "60"))
|
ACK_WAIT_SECONDS = int(os.getenv("ACK_WAIT_SECONDS", "60"))
|
||||||
|
|
||||||
|
# Fiesta is a *separate application* (separate repo, separate binary,
|
||||||
|
# separate business line - jupiter and Fiesta only happen to share one
|
||||||
|
# Postgres instance). Its tenants' orders need Fiesta's own CreateOrderv3
|
||||||
|
# (item validation, atomic order numbers, stock-insufficient checks) -
|
||||||
|
# jupiter has no equivalent logic and was never meant to process this
|
||||||
|
# tenant population. jupiter and Fiesta share one `tenants` table with
|
||||||
|
# no single clean column to tell them apart, so this is an explicit
|
||||||
|
# allowlist rather than a heuristic. Expand FIESTA_TENANT_IDS as more
|
||||||
|
# Fiesta tenants are identified (2026-07-29).
|
||||||
|
FIESTA_BASE_URL = os.getenv("FIESTA_BASE_URL", "http://fiesta.nearle")
|
||||||
|
FIESTA_TENANT_IDS = set(
|
||||||
|
int(t) for t in os.getenv("FIESTA_TENANT_IDS", "").split(",") if t.strip()
|
||||||
|
)
|
||||||
|
|
||||||
# Endpoint mapping is still useful for constructing the target URL
|
# Endpoint mapping is still useful for constructing the target URL
|
||||||
ENDPOINT_MAPPING = {
|
ENDPOINT_MAPPING = {
|
||||||
"/live/api/v1/deliveries/createdeliveries": f"{BASE_URL}/live/api/v1/deliveries/createdeliveries",
|
"/live/api/v1/deliveries/createdeliveries": f"{BASE_URL}/live/api/v1/deliveries/createdeliveries",
|
||||||
@@ -56,10 +70,17 @@ data:
|
|||||||
"/live/api/v2/deliveries/createdeliverylog": f"{BASE_URL}/live/api/v2/deliveries/createdeliverylog",
|
"/live/api/v2/deliveries/createdeliverylog": f"{BASE_URL}/live/api/v2/deliveries/createdeliverylog",
|
||||||
"/live/api/v2/partners/createbreaklog": f"{BASE_URL}/live/api/v2/partners/createbreaklog",
|
"/live/api/v2/partners/createbreaklog": f"{BASE_URL}/live/api/v2/partners/createbreaklog",
|
||||||
"/live/api/v2/partners/updatebreaklog": f"{BASE_URL}/live/api/v2/partners/updatebreaklog",
|
"/live/api/v2/partners/updatebreaklog": f"{BASE_URL}/live/api/v2/partners/updatebreaklog",
|
||||||
"/live/api/v1/mob/orders/createorder": f"{BASE_URL}/live/api/v1/mob/orders/createorder",
|
# Default target for non-Fiesta tenants (jupiter). v1 CreateOrder
|
||||||
"/live/api/v1/web/products/create": f"{BASE_URL}/live/api/v1/web/products/create",
|
# only ever writes the order header - it never loops over "items".
|
||||||
"/live/api/v1/mob/customers/login": f"{BASE_URL}/live/api/v1/mob/customers/login",
|
# CreateOrderv3 does, and the loop is a no-op when Items is empty,
|
||||||
"/live/api/v1/mob/customers/create": f"{BASE_URL}/live/api/v1/mob/customers/create",
|
# so jupiter-native tenants who never send items (e.g. 916/908)
|
||||||
|
# behave identically either way. Fiesta tenants are redirected to
|
||||||
|
# their own backend below, in forward_to_external - this mapping
|
||||||
|
# is only the fallback for everyone else.
|
||||||
|
"/live/api/v1/mob/orders/createorder": f"{BASE_URL}/live/api/v3/orders/createorder",
|
||||||
|
"/live/api/v1/web/products/create": f"{BASE_URL}/live/api/v1/products/create",
|
||||||
|
"/live/api/v1/mob/customers/login": f"{BASE_URL}/live/api/v1/customers/login",
|
||||||
|
"/live/api/v1/mob/customers/create": f"{BASE_URL}/live/api/v1/customers/create",
|
||||||
}
|
}
|
||||||
|
|
||||||
# Global State
|
# Global State
|
||||||
@@ -98,6 +119,27 @@ data:
|
|||||||
return "DROP", None
|
return "DROP", None
|
||||||
data_to_forward = payload["data"]
|
data_to_forward = payload["data"]
|
||||||
|
|
||||||
|
# Fiesta-tenant orders go to Fiesta's own backend instead of jupiter
|
||||||
|
# - see FIESTA_TENANT_IDS above. tenantid may arrive nested under
|
||||||
|
# "orders" (mobile app shape) or flat (direct/API-tool shape); this
|
||||||
|
# is only used to pick the target, Fiesta's own handler does its own
|
||||||
|
# (more thorough) parsing of the actual body.
|
||||||
|
if endpoint == "/live/api/v1/mob/orders/createorder" and FIESTA_TENANT_IDS:
|
||||||
|
tenantid = None
|
||||||
|
if isinstance(data_to_forward, dict):
|
||||||
|
orders_obj = data_to_forward.get("orders")
|
||||||
|
if isinstance(orders_obj, dict) and "tenantid" in orders_obj:
|
||||||
|
tenantid = orders_obj.get("tenantid")
|
||||||
|
elif "tenantid" in data_to_forward:
|
||||||
|
tenantid = data_to_forward.get("tenantid")
|
||||||
|
try:
|
||||||
|
tenantid = int(tenantid) if tenantid is not None else None
|
||||||
|
except (TypeError, ValueError):
|
||||||
|
tenantid = None
|
||||||
|
if tenantid in FIESTA_TENANT_IDS:
|
||||||
|
external_url = f"{FIESTA_BASE_URL}{endpoint}"
|
||||||
|
print(f"?? Routing tenant {tenantid} createorder to Fiesta backend")
|
||||||
|
|
||||||
# Payload Normalization for logs
|
# Payload Normalization for logs
|
||||||
if endpoint == "/live/api/v2/deliveries/createdeliverylog":
|
if endpoint == "/live/api/v2/deliveries/createdeliverylog":
|
||||||
if isinstance(data_to_forward, dict):
|
if isinstance(data_to_forward, dict):
|
||||||
|
|||||||
@@ -97,7 +97,15 @@ spec:
|
|||||||
key: api_key
|
key: api_key
|
||||||
optional: true
|
optional: true
|
||||||
- name: EXTERNAL_BASE_URL
|
- name: EXTERNAL_BASE_URL
|
||||||
value: "http://10.43.229.168"
|
value: "http://jupiter.nearle"
|
||||||
|
- name: FIESTA_BASE_URL
|
||||||
|
value: "http://fiesta.nearle"
|
||||||
|
# Known Fiesta tenants (R mart=1147, Suriya Store=1135, confirmed
|
||||||
|
# 2026-07-29). Expand as more are identified - see worker.py's
|
||||||
|
# FIESTA_TENANT_IDS comment for why this is an explicit list rather
|
||||||
|
# than a DB heuristic.
|
||||||
|
- name: FIESTA_TENANT_IDS
|
||||||
|
value: "1147,1135"
|
||||||
resources:
|
resources:
|
||||||
requests:
|
requests:
|
||||||
memory: "128Mi"
|
memory: "128Mi"
|
||||||
@@ -207,7 +215,7 @@ spec:
|
|||||||
name: nats-credentials
|
name: nats-credentials
|
||||||
key: password
|
key: password
|
||||||
- name: EXTERNAL_BASE_URL
|
- name: EXTERNAL_BASE_URL
|
||||||
value: "http://10.43.224.63"
|
value: "http://jupiter.nearle"
|
||||||
resources:
|
resources:
|
||||||
requests:
|
requests:
|
||||||
memory: "128Mi"
|
memory: "128Mi"
|
||||||
@@ -323,7 +331,7 @@ spec:
|
|||||||
key: api_key
|
key: api_key
|
||||||
optional: true
|
optional: true
|
||||||
- name: EXTERNAL_BASE_URL
|
- name: EXTERNAL_BASE_URL
|
||||||
value: "http://10.43.229.168"
|
value: "http://jupiter.nearle"
|
||||||
resources:
|
resources:
|
||||||
requests:
|
requests:
|
||||||
memory: "128Mi"
|
memory: "128Mi"
|
||||||
@@ -439,7 +447,7 @@ spec:
|
|||||||
key: api_key
|
key: api_key
|
||||||
optional: true
|
optional: true
|
||||||
- name: EXTERNAL_BASE_URL
|
- name: EXTERNAL_BASE_URL
|
||||||
value: "http://10.43.224.63"
|
value: "http://jupiter.nearle"
|
||||||
resources:
|
resources:
|
||||||
requests:
|
requests:
|
||||||
memory: "128Mi"
|
memory: "128Mi"
|
||||||
@@ -555,7 +563,7 @@ spec:
|
|||||||
key: api_key
|
key: api_key
|
||||||
optional: true
|
optional: true
|
||||||
- name: EXTERNAL_BASE_URL
|
- name: EXTERNAL_BASE_URL
|
||||||
value: "http://10.43.229.168"
|
value: "http://jupiter.nearle"
|
||||||
resources:
|
resources:
|
||||||
requests:
|
requests:
|
||||||
memory: "128Mi"
|
memory: "128Mi"
|
||||||
|
|||||||
@@ -1,514 +0,0 @@
|
|||||||
apiVersion: v1
|
|
||||||
kind: ConfigMap
|
|
||||||
metadata:
|
|
||||||
name: fiesta-gateway-script
|
|
||||||
namespace: nearle
|
|
||||||
data:
|
|
||||||
app.py: |
|
|
||||||
#!/usr/bin/env python3
|
|
||||||
"""
|
|
||||||
FastAPI application with multiple endpoints that publish to NATS JetStream
|
|
||||||
Each endpoint corresponds to an external API that workers will forward to
|
|
||||||
"""
|
|
||||||
import os
|
|
||||||
import json
|
|
||||||
import asyncio
|
|
||||||
from fastapi import FastAPI, HTTPException, Request, Body
|
|
||||||
from fastapi.responses import JSONResponse
|
|
||||||
from fastapi.middleware.cors import CORSMiddleware
|
|
||||||
from starlette.middleware.base import BaseHTTPMiddleware
|
|
||||||
from starlette.requests import Request as StarletteRequest
|
|
||||||
from starlette.responses import Response as StarletteResponse
|
|
||||||
from pydantic import BaseModel
|
|
||||||
import nats
|
|
||||||
import nats.errors
|
|
||||||
|
|
||||||
from prometheus_client import Counter, Histogram, generate_latest, REGISTRY
|
|
||||||
from starlette.responses import Response
|
|
||||||
import uvicorn
|
|
||||||
from typing import Optional, Dict, Any, List, Union
|
|
||||||
|
|
||||||
# Prometheus metrics
|
|
||||||
request_count = Counter('http_requests_total', 'Total HTTP requests', ['method', 'endpoint', 'status'])
|
|
||||||
request_duration = Histogram('http_request_duration_seconds', 'HTTP request duration', ['method', 'endpoint'])
|
|
||||||
|
|
||||||
app = FastAPI(title="NATS Backend API - Multi-Endpoint", version="1.0.0")
|
|
||||||
|
|
||||||
# Custom CORS middleware to ensure headers are always added
|
|
||||||
class CORSHeaderMiddleware(BaseHTTPMiddleware):
|
|
||||||
async def dispatch(self, request: StarletteRequest, call_next):
|
|
||||||
origin = request.headers.get("origin", "*")
|
|
||||||
response = await call_next(request)
|
|
||||||
response.headers["Access-Control-Allow-Origin"] = origin
|
|
||||||
response.headers["Access-Control-Allow-Methods"] = "GET, POST, PUT, PATCH, DELETE, OPTIONS"
|
|
||||||
response.headers["Access-Control-Allow-Headers"] = "*"
|
|
||||||
response.headers["Access-Control-Allow-Credentials"] = "false"
|
|
||||||
response.headers["Access-Control-Max-Age"] = "600"
|
|
||||||
return response
|
|
||||||
|
|
||||||
# Add custom CORS middleware first
|
|
||||||
app.add_middleware(CORSHeaderMiddleware)
|
|
||||||
|
|
||||||
# Also add FastAPI's CORS middleware as backup
|
|
||||||
app.add_middleware(
|
|
||||||
CORSMiddleware,
|
|
||||||
allow_origin_regex=r".*", # Match all origins
|
|
||||||
allow_credentials=False,
|
|
||||||
allow_methods=["*"], # Allow all HTTP methods
|
|
||||||
allow_headers=["*"], # Allow all headers
|
|
||||||
expose_headers=["*"], # Expose all headers
|
|
||||||
max_age=600,
|
|
||||||
)
|
|
||||||
|
|
||||||
# NATS connection (will be initialized on startup)
|
|
||||||
nc = None
|
|
||||||
js = None
|
|
||||||
|
|
||||||
# Configuration
|
|
||||||
NATS_URL = os.getenv("NATS_URL", "nats://nats-server:4222")
|
|
||||||
NATS_USER = os.getenv("NATS_USER", "admin")
|
|
||||||
NATS_PASSWORD = os.getenv("NATS_PASSWORD", "")
|
|
||||||
|
|
||||||
# Endpoint to NATS subject mapping
|
|
||||||
ENDPOINT_ROUTES = {
|
|
||||||
"/live/api/v1/deliveries/createdeliveries": "api.v1.deliveries.createdeliveries",
|
|
||||||
"/live/api/v1/deliveries/updatedelivery": "api.v1.deliveries.updatedelivery",
|
|
||||||
"/live/api/v2/partners/createriderlog": "api.v2.partners.createriderlog",
|
|
||||||
"/live/api/v2/deliveries/createdeliverylog": "api.v2.deliveries.createdeliverylog",
|
|
||||||
"/live/api/v2/partners/createbreaklog": "api.v2.partners.createbreaklog",
|
|
||||||
"/live/api/v2/partners/updatebreaklog": "api.v2.partners.updatebreaklog",
|
|
||||||
"/live/api/v1/mob/orders/createorder": "api.v1.mob.orders.createorder",
|
|
||||||
"/live/api/v1/web/products/create": "api.v1.web.products.create",
|
|
||||||
# Customer Endpoints (Sync)
|
|
||||||
"/live/api/v1/mob/customers/login": "api.v1.mob.customers.login",
|
|
||||||
"/live/api/v1/mob/customers/create": "api.v1.mob.customers.create",
|
|
||||||
}
|
|
||||||
|
|
||||||
@app.on_event("startup")
|
|
||||||
async def startup():
|
|
||||||
"""Initialize NATS connection on startup"""
|
|
||||||
global nc, js
|
|
||||||
try:
|
|
||||||
print(f"Connecting to NATS at {NATS_URL}...")
|
|
||||||
nc = await nats.connect(
|
|
||||||
servers=[NATS_URL],
|
|
||||||
user=NATS_USER,
|
|
||||||
password=NATS_PASSWORD,
|
|
||||||
reconnect_time_wait=2,
|
|
||||||
max_reconnect_attempts=10
|
|
||||||
)
|
|
||||||
js = nc.jetstream()
|
|
||||||
print("✅ Connected to NATS JetStream")
|
|
||||||
print(f"✅ Configured {len(ENDPOINT_ROUTES)} endpoint routes")
|
|
||||||
except Exception as e:
|
|
||||||
print(f"❌ Failed to connect to NATS: {e}")
|
|
||||||
raise
|
|
||||||
|
|
||||||
@app.on_event("shutdown")
|
|
||||||
async def shutdown():
|
|
||||||
"""Close NATS connection on shutdown"""
|
|
||||||
global nc
|
|
||||||
if nc:
|
|
||||||
await nc.close()
|
|
||||||
print("NATS connection closed")
|
|
||||||
|
|
||||||
@app.options("/{full_path:path}")
|
|
||||||
async def options_handler(full_path: str, request: Request):
|
|
||||||
"""Handle OPTIONS requests for CORS preflight"""
|
|
||||||
origin = request.headers.get("origin")
|
|
||||||
return JSONResponse(
|
|
||||||
status_code=200,
|
|
||||||
content={},
|
|
||||||
headers={
|
|
||||||
"Access-Control-Allow-Origin": origin if origin else "*",
|
|
||||||
"Access-Control-Allow-Methods": "GET, POST, PUT, PATCH, DELETE, OPTIONS",
|
|
||||||
"Access-Control-Allow-Headers": "*",
|
|
||||||
"Access-Control-Allow-Credentials": "false",
|
|
||||||
"Access-Control-Max-Age": "600",
|
|
||||||
}
|
|
||||||
)
|
|
||||||
|
|
||||||
@app.get("/health")
|
|
||||||
async def health():
|
|
||||||
"""Health check endpoint"""
|
|
||||||
return {"status": "healthy", "nats_connected": nc.is_connected if nc else False}
|
|
||||||
|
|
||||||
@app.get("/ready")
|
|
||||||
async def ready():
|
|
||||||
"""Readiness check endpoint"""
|
|
||||||
if nc and nc.is_connected:
|
|
||||||
return {"status": "ready"}
|
|
||||||
raise HTTPException(status_code=503, detail="Not ready")
|
|
||||||
|
|
||||||
async def publish_to_nats(endpoint: str, data: Union[Dict[str, Any], List[Dict[str, Any]]], request_method: str = "POST"):
|
|
||||||
"""Publish message to NATS with endpoint metadata"""
|
|
||||||
payload = {
|
|
||||||
"endpoint": endpoint,
|
|
||||||
"method": request_method,
|
|
||||||
"data": data,
|
|
||||||
"received_at": int(asyncio.get_event_loop().time() * 1000),
|
|
||||||
"original_path": endpoint
|
|
||||||
}
|
|
||||||
|
|
||||||
# Get NATS subject for this endpoint
|
|
||||||
subject = ENDPOINT_ROUTES.get(endpoint, "api.unknown")
|
|
||||||
|
|
||||||
if js:
|
|
||||||
try:
|
|
||||||
ack = await js.publish(subject, json.dumps(payload).encode())
|
|
||||||
return ack.seq
|
|
||||||
except Exception as e:
|
|
||||||
print(f"❌ Failed to publish to NATS: {e}")
|
|
||||||
raise HTTPException(status_code=500, detail=f"Failed to publish message: {str(e)}")
|
|
||||||
else:
|
|
||||||
raise HTTPException(status_code=503, detail="NATS not connected")
|
|
||||||
|
|
||||||
async def publish_request_to_nats(endpoint: str, data: Dict[str, Any], request_method: str = "POST", timeout: int = 10):
|
|
||||||
"""
|
|
||||||
Publish to NATS and WAIT for a reply (Request-Reply pattern).
|
|
||||||
Used for synchronous endpoints like Login.
|
|
||||||
"""
|
|
||||||
payload = {
|
|
||||||
"endpoint": endpoint,
|
|
||||||
"method": request_method,
|
|
||||||
"data": data,
|
|
||||||
"received_at": int(asyncio.get_event_loop().time() * 1000),
|
|
||||||
"original_path": endpoint
|
|
||||||
}
|
|
||||||
|
|
||||||
subject = ENDPOINT_ROUTES.get(endpoint, "api.unknown")
|
|
||||||
|
|
||||||
if not js:
|
|
||||||
raise HTTPException(status_code=503, detail="NATS not connected")
|
|
||||||
|
|
||||||
try:
|
|
||||||
# Create a unique inbox for the reply
|
|
||||||
inbox = nc.new_inbox()
|
|
||||||
|
|
||||||
# Subscribe to the inbox first
|
|
||||||
sub = await nc.subscribe(inbox, max_msgs=1)
|
|
||||||
|
|
||||||
# Publish request with reply inbox
|
|
||||||
# Note: We use js.publish to ensure it goes to the Stream (Queue), but attach a reply subject
|
|
||||||
await js.publish(subject, json.dumps(payload).encode(), reply=inbox)
|
|
||||||
|
|
||||||
# Wait for valid response
|
|
||||||
try:
|
|
||||||
msg = await sub.next_msg(timeout=timeout)
|
|
||||||
response_data = json.loads(msg.data.decode())
|
|
||||||
return response_data
|
|
||||||
except nats.errors.TimeoutError:
|
|
||||||
raise HTTPException(status_code=504, detail="Gateway Timeout: Upstream service did not respond in time")
|
|
||||||
finally:
|
|
||||||
await sub.unsubscribe()
|
|
||||||
|
|
||||||
except HTTPException:
|
|
||||||
raise
|
|
||||||
except Exception as e:
|
|
||||||
print(f"❌ Failed to request from NATS: {e}")
|
|
||||||
raise HTTPException(status_code=500, detail=f"RPC Error: {str(e)}")
|
|
||||||
|
|
||||||
# Endpoint 1: Update Delivery (v1) - External API requires PUT
|
|
||||||
@app.put("/live/api/v1/deliveries/updatedelivery")
|
|
||||||
async def update_delivery_v1(data: Dict[str, Any], request: Request):
|
|
||||||
"""Update Delivery endpoint - forwards to NATS (PUT only, as external API requires PUT)"""
|
|
||||||
start_time = asyncio.get_event_loop().time()
|
|
||||||
endpoint = "/live/api/v1/deliveries/updatedelivery"
|
|
||||||
method = "PUT" # Always use PUT for this endpoint
|
|
||||||
|
|
||||||
try:
|
|
||||||
message_id = await publish_to_nats(endpoint, data, method)
|
|
||||||
duration = asyncio.get_event_loop().time() - start_time
|
|
||||||
request_duration.labels(method="PUT", endpoint=endpoint).observe(duration)
|
|
||||||
request_count.labels(method="PUT", endpoint=endpoint, status="200").inc()
|
|
||||||
|
|
||||||
return JSONResponse(
|
|
||||||
status_code=200,
|
|
||||||
content={"status": "accepted", "message_id": message_id, "endpoint": endpoint}
|
|
||||||
)
|
|
||||||
except HTTPException:
|
|
||||||
raise
|
|
||||||
except Exception as e:
|
|
||||||
request_count.labels(method="PUT", endpoint=endpoint, status="500").inc()
|
|
||||||
print(f"❌ Error processing {endpoint}: {e}")
|
|
||||||
raise HTTPException(status_code=500, detail=str(e))
|
|
||||||
|
|
||||||
# Endpoint: Create Deliveries (v1)
|
|
||||||
@app.post("/live/api/v1/deliveries/createdeliveries")
|
|
||||||
async def create_deliveries_v1(data: Union[Dict[str, Any], List[Dict[str, Any]]], request: Request):
|
|
||||||
"""Create Deliveries endpoint - forwards payload to NATS"""
|
|
||||||
start_time = asyncio.get_event_loop().time()
|
|
||||||
endpoint = "/live/api/v1/deliveries/createdeliveries"
|
|
||||||
method = request.method
|
|
||||||
|
|
||||||
try:
|
|
||||||
message_id = await publish_to_nats(endpoint, data, method)
|
|
||||||
duration = asyncio.get_event_loop().time() - start_time
|
|
||||||
request_duration.labels(method=method, endpoint=endpoint).observe(duration)
|
|
||||||
request_count.labels(method=method, endpoint=endpoint, status="200").inc()
|
|
||||||
|
|
||||||
return JSONResponse(
|
|
||||||
status_code=200,
|
|
||||||
content={"status": "accepted", "message_id": message_id, "endpoint": endpoint}
|
|
||||||
)
|
|
||||||
except HTTPException:
|
|
||||||
raise
|
|
||||||
except Exception as e:
|
|
||||||
request_count.labels(method=method, endpoint=endpoint, status="500").inc()
|
|
||||||
print(f"❌ Error processing {endpoint}: {e}")
|
|
||||||
raise HTTPException(status_code=500, detail=str(e))
|
|
||||||
|
|
||||||
# Endpoint 2: Create Rider Log (v2)
|
|
||||||
@app.post("/live/api/v2/partners/createriderlog")
|
|
||||||
async def create_rider_log_v2(data: Dict[str, Any], request: Request):
|
|
||||||
"""Create Rider Log endpoint - forwards to NATS"""
|
|
||||||
start_time = asyncio.get_event_loop().time()
|
|
||||||
endpoint = "/live/api/v2/partners/createriderlog"
|
|
||||||
method = request.method
|
|
||||||
|
|
||||||
try:
|
|
||||||
message_id = await publish_to_nats(endpoint, data, method)
|
|
||||||
duration = asyncio.get_event_loop().time() - start_time
|
|
||||||
request_duration.labels(method=method, endpoint=endpoint).observe(duration)
|
|
||||||
request_count.labels(method=method, endpoint=endpoint, status="200").inc()
|
|
||||||
|
|
||||||
return JSONResponse(
|
|
||||||
status_code=200,
|
|
||||||
content={"status": "accepted", "message_id": message_id, "endpoint": endpoint}
|
|
||||||
)
|
|
||||||
except HTTPException:
|
|
||||||
raise
|
|
||||||
except Exception as e:
|
|
||||||
request_count.labels(method=method, endpoint=endpoint, status="500").inc()
|
|
||||||
print(f"❌ Error processing {endpoint}: {e}")
|
|
||||||
raise HTTPException(status_code=500, detail=str(e))
|
|
||||||
|
|
||||||
# Endpoint 3: Create Delivery Log (v2)
|
|
||||||
@app.post("/live/api/v2/deliveries/createdeliverylog")
|
|
||||||
async def create_delivery_log_v2(
|
|
||||||
request: Request,
|
|
||||||
data: List[Dict[str, Any]] = Body(
|
|
||||||
...,
|
|
||||||
example=[{
|
|
||||||
"logid": 0,
|
|
||||||
"tenantid": 1,
|
|
||||||
"partnerid": 44,
|
|
||||||
"locationid": 1,
|
|
||||||
"orderheaderid": 123456,
|
|
||||||
"deliveryid": 654321,
|
|
||||||
"userid": 1111,
|
|
||||||
"orderid": "1-20231624",
|
|
||||||
"orderstatus": "active",
|
|
||||||
"starttime": "2025-12-10 17:51:03",
|
|
||||||
"logdate": "2025-12-10 18:14:04",
|
|
||||||
"latitude": "11.0050664",
|
|
||||||
"longitude": "76.9508776"
|
|
||||||
}]
|
|
||||||
)
|
|
||||||
):
|
|
||||||
"""
|
|
||||||
Create Delivery Log endpoint - forwards to NATS.
|
|
||||||
Accepts either a single dict or a list of dicts to align with external API expectations.
|
|
||||||
"""
|
|
||||||
start_time = asyncio.get_event_loop().time()
|
|
||||||
endpoint = "/live/api/v2/deliveries/createdeliverylog"
|
|
||||||
method = request.method
|
|
||||||
|
|
||||||
try:
|
|
||||||
# Expect only a list of dicts
|
|
||||||
if not isinstance(data, list) or not all(isinstance(item, dict) for item in data):
|
|
||||||
raise HTTPException(status_code=422, detail="Body must be a list of objects")
|
|
||||||
|
|
||||||
normalized_data = data
|
|
||||||
|
|
||||||
message_id = await publish_to_nats(endpoint, normalized_data, method)
|
|
||||||
duration = asyncio.get_event_loop().time() - start_time
|
|
||||||
request_duration.labels(method=method, endpoint=endpoint).observe(duration)
|
|
||||||
request_count.labels(method=method, endpoint=endpoint, status="200").inc()
|
|
||||||
|
|
||||||
return JSONResponse(
|
|
||||||
status_code=200,
|
|
||||||
content={"status": "accepted", "message_id": message_id, "endpoint": endpoint}
|
|
||||||
)
|
|
||||||
except HTTPException:
|
|
||||||
raise
|
|
||||||
except Exception as e:
|
|
||||||
request_count.labels(method=method, endpoint=endpoint, status="500").inc()
|
|
||||||
print(f"❌ Error processing {endpoint}: {e}")
|
|
||||||
raise HTTPException(status_code=500, detail=str(e))
|
|
||||||
|
|
||||||
# Endpoint 4: Create Break Rider Log (v2)
|
|
||||||
@app.post("/live/api/v2/partners/createbreaklog")
|
|
||||||
async def create_break_log_v2(data: Dict[str, Any], request: Request):
|
|
||||||
"""Create Break Rider Log endpoint - forwards to NATS"""
|
|
||||||
start_time = asyncio.get_event_loop().time()
|
|
||||||
endpoint = "/live/api/v2/partners/createbreaklog"
|
|
||||||
method = request.method
|
|
||||||
|
|
||||||
try:
|
|
||||||
message_id = await publish_to_nats(endpoint, data, method)
|
|
||||||
duration = asyncio.get_event_loop().time() - start_time
|
|
||||||
request_duration.labels(method=method, endpoint=endpoint).observe(duration)
|
|
||||||
request_count.labels(method=method, endpoint=endpoint, status="200").inc()
|
|
||||||
|
|
||||||
return JSONResponse(
|
|
||||||
status_code=200,
|
|
||||||
content={"status": "accepted", "message_id": message_id, "endpoint": endpoint}
|
|
||||||
)
|
|
||||||
except HTTPException:
|
|
||||||
raise
|
|
||||||
except Exception as e:
|
|
||||||
request_count.labels(method=method, endpoint=endpoint, status="500").inc()
|
|
||||||
print(f"❌ Error processing {endpoint}: {e}")
|
|
||||||
raise HTTPException(status_code=500, detail=str(e))
|
|
||||||
|
|
||||||
# Endpoint 5: Update Break Rider Log (v2) - Supports both POST and PUT
|
|
||||||
@app.post("/live/api/v2/partners/updatebreaklog")
|
|
||||||
@app.put("/live/api/v2/partners/updatebreaklog")
|
|
||||||
async def update_break_log_v2(data: Dict[str, Any], request: Request):
|
|
||||||
"""Update Break Rider Log endpoint - forwards to NATS"""
|
|
||||||
start_time = asyncio.get_event_loop().time()
|
|
||||||
endpoint = "/live/api/v2/partners/updatebreaklog"
|
|
||||||
method = request.method # Will be POST or PUT
|
|
||||||
|
|
||||||
try:
|
|
||||||
message_id = await publish_to_nats(endpoint, data, method)
|
|
||||||
duration = asyncio.get_event_loop().time() - start_time
|
|
||||||
request_duration.labels(method=method, endpoint=endpoint).observe(duration)
|
|
||||||
request_count.labels(method=method, endpoint=endpoint, status="200").inc()
|
|
||||||
|
|
||||||
return JSONResponse(
|
|
||||||
status_code=200,
|
|
||||||
content={"status": "accepted", "message_id": message_id, "endpoint": endpoint}
|
|
||||||
)
|
|
||||||
except HTTPException:
|
|
||||||
raise
|
|
||||||
except Exception as e:
|
|
||||||
request_count.labels(method=method, endpoint=endpoint, status="500").inc()
|
|
||||||
print(f"❌ Error processing {endpoint}: {e}")
|
|
||||||
raise HTTPException(status_code=500, detail=str(e))
|
|
||||||
|
|
||||||
@app.get("/metrics")
|
|
||||||
async def metrics():
|
|
||||||
"""Prometheus metrics endpoint"""
|
|
||||||
return Response(content=generate_latest(REGISTRY), media_type="text/plain")
|
|
||||||
|
|
||||||
@app.get("/routes")
|
|
||||||
async def list_routes():
|
|
||||||
"""List all configured routes"""
|
|
||||||
return {
|
|
||||||
"routes": ENDPOINT_ROUTES,
|
|
||||||
"total": len(ENDPOINT_ROUTES)
|
|
||||||
}
|
|
||||||
|
|
||||||
|
|
||||||
# Endpoint 6: Create Order (Mob V1)
|
|
||||||
@app.post("/live/api/v1/mob/orders/createorder")
|
|
||||||
async def create_order_mob_v1(data: Dict[str, Any], request: Request):
|
|
||||||
"""Create Order endpoint (Mobile) - forwards to NATS"""
|
|
||||||
start_time = asyncio.get_event_loop().time()
|
|
||||||
endpoint = "/live/api/v1/mob/orders/createorder"
|
|
||||||
method = request.method
|
|
||||||
|
|
||||||
try:
|
|
||||||
message_id = await publish_to_nats(endpoint, data, method)
|
|
||||||
duration = asyncio.get_event_loop().time() - start_time
|
|
||||||
request_duration.labels(method=method, endpoint=endpoint).observe(duration)
|
|
||||||
request_count.labels(method=method, endpoint=endpoint, status="200").inc()
|
|
||||||
|
|
||||||
return JSONResponse(
|
|
||||||
status_code=200,
|
|
||||||
content={"status": "accepted", "message_id": message_id, "endpoint": endpoint}
|
|
||||||
)
|
|
||||||
except HTTPException:
|
|
||||||
raise
|
|
||||||
except Exception as e:
|
|
||||||
request_count.labels(method=method, endpoint=endpoint, status="500").inc()
|
|
||||||
print(f"❌ Error processing {endpoint}: {e}")
|
|
||||||
raise HTTPException(status_code=500, detail=str(e))
|
|
||||||
|
|
||||||
# Endpoint 7: Create Product (Web V1)
|
|
||||||
@app.post("/live/api/v1/web/products/create")
|
|
||||||
async def create_product_web_v1(data: Dict[str, Any], request: Request):
|
|
||||||
"""Create Product endpoint (Web) - forwards to NATS"""
|
|
||||||
start_time = asyncio.get_event_loop().time()
|
|
||||||
endpoint = "/live/api/v1/web/products/create"
|
|
||||||
method = request.method
|
|
||||||
|
|
||||||
try:
|
|
||||||
message_id = await publish_to_nats(endpoint, data, method)
|
|
||||||
duration = asyncio.get_event_loop().time() - start_time
|
|
||||||
request_duration.labels(method=method, endpoint=endpoint).observe(duration)
|
|
||||||
request_count.labels(method=method, endpoint=endpoint, status="200").inc()
|
|
||||||
|
|
||||||
return JSONResponse(
|
|
||||||
status_code=200,
|
|
||||||
content={"status": "accepted", "message_id": message_id, "endpoint": endpoint}
|
|
||||||
)
|
|
||||||
except HTTPException:
|
|
||||||
raise
|
|
||||||
except Exception as e:
|
|
||||||
request_count.labels(method=method, endpoint=endpoint, status="500").inc()
|
|
||||||
print(f"❌ Error processing {endpoint}: {e}")
|
|
||||||
raise HTTPException(status_code=500, detail=str(e))
|
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
# Endpoint 8: Customer Login (Sync Request-Reply)
|
|
||||||
@app.post("/live/api/v1/mob/customers/login")
|
|
||||||
async def customer_login(data: Dict[str, Any], request: Request):
|
|
||||||
"""Customer Login - Waits for response from worker"""
|
|
||||||
start_time = asyncio.get_event_loop().time()
|
|
||||||
endpoint = "/live/api/v1/mob/customers/login"
|
|
||||||
method = request.method
|
|
||||||
|
|
||||||
try:
|
|
||||||
# Wait for reply!
|
|
||||||
response_data = await publish_request_to_nats(endpoint, data, method)
|
|
||||||
|
|
||||||
duration = asyncio.get_event_loop().time() - start_time
|
|
||||||
request_duration.labels(method=method, endpoint=endpoint).observe(duration)
|
|
||||||
request_count.labels(method=method, endpoint=endpoint, status="200").inc()
|
|
||||||
|
|
||||||
# Return the actual backend response
|
|
||||||
return JSONResponse(
|
|
||||||
status_code=200,
|
|
||||||
content=response_data
|
|
||||||
)
|
|
||||||
except HTTPException:
|
|
||||||
raise
|
|
||||||
except Exception as e:
|
|
||||||
request_count.labels(method=method, endpoint=endpoint, status="500").inc()
|
|
||||||
print(f"❌ Error processing {endpoint}: {e}")
|
|
||||||
raise HTTPException(status_code=500, detail=str(e))
|
|
||||||
|
|
||||||
# Endpoint 9: Customer Create (Sync Request-Reply)
|
|
||||||
@app.post("/live/api/v1/mob/customers/create")
|
|
||||||
async def customer_create(data: Dict[str, Any], request: Request):
|
|
||||||
"""Customer Create - Waits for response from worker"""
|
|
||||||
start_time = asyncio.get_event_loop().time()
|
|
||||||
endpoint = "/live/api/v1/mob/customers/create"
|
|
||||||
method = request.method
|
|
||||||
|
|
||||||
try:
|
|
||||||
# Wait for reply!
|
|
||||||
response_data = await publish_request_to_nats(endpoint, data, method)
|
|
||||||
|
|
||||||
duration = asyncio.get_event_loop().time() - start_time
|
|
||||||
request_duration.labels(method=method, endpoint=endpoint).observe(duration)
|
|
||||||
request_count.labels(method=method, endpoint=endpoint, status="200").inc()
|
|
||||||
|
|
||||||
return JSONResponse(
|
|
||||||
status_code=200,
|
|
||||||
content=response_data
|
|
||||||
)
|
|
||||||
except HTTPException:
|
|
||||||
raise
|
|
||||||
except Exception as e:
|
|
||||||
request_count.labels(method=method, endpoint=endpoint, status="500").inc()
|
|
||||||
print(f"❌ Error processing {endpoint}: {e}")
|
|
||||||
raise HTTPException(status_code=500, detail=str(e))
|
|
||||||
|
|
||||||
if __name__ == "__main__":
|
|
||||||
uvicorn.run(app, host="0.0.0.0", port=8000)
|
|
||||||
|
|
||||||
@@ -6,7 +6,6 @@ resources:
|
|||||||
- nearle-config.yaml
|
- nearle-config.yaml
|
||||||
- nearle-secrets.yaml
|
- nearle-secrets.yaml
|
||||||
- nearle-app-secrets.yaml
|
- nearle-app-secrets.yaml
|
||||||
- fiesta-gateway.yaml
|
|
||||||
- nearle-fiesta.yaml
|
- nearle-fiesta.yaml
|
||||||
- nearle-jupiter.yaml
|
- nearle-jupiter.yaml
|
||||||
- jupiter-cors-proxy.yaml
|
- jupiter-cors-proxy.yaml
|
||||||
|
|||||||
@@ -7,17 +7,14 @@ metadata:
|
|||||||
app: fiesta
|
app: fiesta
|
||||||
type: Opaque
|
type: Opaque
|
||||||
stringData:
|
stringData:
|
||||||
# The IP of your BigRock Server
|
|
||||||
DATABASE_HOST: "66.116.207.225"
|
DATABASE_HOST: "66.116.207.225"
|
||||||
DB_HOST: "66.116.207.225"
|
DB_HOST: "66.116.207.225"
|
||||||
|
|
||||||
# The user we confirmed works
|
|
||||||
DATABASE_USERNAME: "admin"
|
DATABASE_USERNAME: "admin"
|
||||||
DB_USER: "admin"
|
DB_USER: "admin"
|
||||||
|
|
||||||
# The password we confirmed works
|
|
||||||
DATABASE_PASSWORD: "Package@123#"
|
DATABASE_PASSWORD: "Package@123#"
|
||||||
DB_PASSWORD: "Package@123#"
|
DB_PASSWORD: "Package@123#"
|
||||||
|
|
||||||
# The rest...
|
|
||||||
JWT_SECRET_KEY: "nearle"
|
JWT_SECRET_KEY: "nearle"
|
||||||
|
CATALOGUE_DB_USER: "admin"
|
||||||
|
CATALOGUE_DB_PASSWORD: "'Package@321#'"
|
||||||
|
S3_ACCESS_KEY: "DO801G8Q8JAZKF49U3WJ"
|
||||||
|
S3_SECRET_KEY: "lBQExYfkVqH+ybmGVmQH5MkThBbrIohA/VQLgcPUvug"
|
||||||
|
|||||||
@@ -8,16 +8,21 @@ metadata:
|
|||||||
data:
|
data:
|
||||||
NATS_URL: "nats://66.116.226.161:4222"
|
NATS_URL: "nats://66.116.226.161:4222"
|
||||||
LOG_LEVEL: "info"
|
LOG_LEVEL: "info"
|
||||||
# The stream is set to ORDERS as per your request
|
|
||||||
NATS_STREAM: "ORDERS"
|
NATS_STREAM: "ORDERS"
|
||||||
# Base subject pattern - the app likely appends the path to this or uses it as a listener filter
|
NATS_SUBJECT: "api.>"
|
||||||
NATS_SUBJECT: "api.>"
|
|
||||||
ALLOWED_ORIGINS: "*"
|
ALLOWED_ORIGINS: "*"
|
||||||
ENV: "production"
|
ENV: "production"
|
||||||
DATABASE_NAME: "nearledb"
|
DATABASE_NAME: "nearledb"
|
||||||
DB_NAME: "nearledb"
|
DB_NAME: "nearledb"
|
||||||
DATABASE_PORT: "5432"
|
DATABASE_PORT: "5433"
|
||||||
DB_PORT: "5432"
|
DB_PORT: "5433"
|
||||||
DATABASE_SERVER_HOST: "66.116.207.225"
|
DATABASE_SERVER_HOST: "66.116.207.225"
|
||||||
DB_HOST: "66.116.207.225"
|
DB_HOST: "66.116.207.225"
|
||||||
USER_CONTEXT_KEY: "nearle"
|
USER_CONTEXT_KEY: "nearle"
|
||||||
|
CATALOGUE_DB_HOST: "31.97.228.132"
|
||||||
|
CATALOGUE_DB_PORT: "6054"
|
||||||
|
CATALOGUE_DB_NAME: "pgvector"
|
||||||
|
USE_S3: "true"
|
||||||
|
S3_ENDPOINT: "https://nearle.sgp1.digitaloceanspaces.com"
|
||||||
|
S3_BUCKET: "nearle"
|
||||||
|
S3_REGION: "sgp1"
|
||||||
|
|||||||
@@ -47,7 +47,7 @@ spec:
|
|||||||
app: fiesta
|
app: fiesta
|
||||||
containers:
|
containers:
|
||||||
- name: backend
|
- name: backend
|
||||||
image: nearlecommerce/fiesta:v1.3.78
|
image: nearlecommerce/fiesta:v1.3.93
|
||||||
imagePullPolicy: Always
|
imagePullPolicy: Always
|
||||||
securityContext:
|
securityContext:
|
||||||
allowPrivilegeEscalation: false
|
allowPrivilegeEscalation: false
|
||||||
@@ -74,39 +74,22 @@ spec:
|
|||||||
secretKeyRef:
|
secretKeyRef:
|
||||||
name: nats-credentials
|
name: nats-credentials
|
||||||
key: password
|
key: password
|
||||||
- name: gateway
|
- name: MQTT_URL
|
||||||
image: workolik360/alaska:v1.2.0
|
value: "tcp://66.116.225.226:1883"
|
||||||
imagePullPolicy: Always
|
- name: MQTT_USER
|
||||||
securityContext:
|
value: "pos_ingest"
|
||||||
allowPrivilegeEscalation: false
|
- name: MQTT_PASSWORD
|
||||||
capabilities:
|
value: "AXbEPrNDWnMLdp7T1tFETwyU"
|
||||||
drop:
|
- name: REDIS_HOST
|
||||||
- ALL
|
value: "66.116.226.255"
|
||||||
ports:
|
- name: REDIS_PORT
|
||||||
- containerPort: 8000
|
value: "6379"
|
||||||
name: http
|
- name: REDIS_USER
|
||||||
volumeMounts:
|
value: "default"
|
||||||
- name: gateway-script
|
- name: REDIS_PASSWORD
|
||||||
mountPath: /app/app.py
|
value: "Package@324969#"
|
||||||
subPath: app.py
|
- name: REDIS_DB
|
||||||
envFrom:
|
value: "0"
|
||||||
- configMapRef:
|
|
||||||
name: nearle-config
|
|
||||||
env:
|
|
||||||
- name: NATS_USER
|
|
||||||
valueFrom:
|
|
||||||
secretKeyRef:
|
|
||||||
name: nats-credentials
|
|
||||||
key: username
|
|
||||||
- name: NATS_PASSWORD
|
|
||||||
valueFrom:
|
|
||||||
secretKeyRef:
|
|
||||||
name: nats-credentials
|
|
||||||
key: password
|
|
||||||
volumes:
|
|
||||||
- name: gateway-script
|
|
||||||
configMap:
|
|
||||||
name: fiesta-gateway-script
|
|
||||||
---
|
---
|
||||||
apiVersion: v1
|
apiVersion: v1
|
||||||
kind: Service
|
kind: Service
|
||||||
@@ -123,45 +106,5 @@ spec:
|
|||||||
nodePort: 30823
|
nodePort: 30823
|
||||||
protocol: TCP
|
protocol: TCP
|
||||||
name: main
|
name: main
|
||||||
- port: 8000
|
|
||||||
targetPort: 8000
|
|
||||||
name: gateway
|
|
||||||
protocol: TCP
|
|
||||||
selector:
|
selector:
|
||||||
app: fiesta
|
app: fiesta
|
||||||
---
|
|
||||||
apiVersion: gateway.networking.k8s.io/v1
|
|
||||||
kind: HTTPRoute
|
|
||||||
metadata:
|
|
||||||
name: fiesta-route
|
|
||||||
namespace: nearle
|
|
||||||
labels:
|
|
||||||
app: fiesta
|
|
||||||
spec:
|
|
||||||
parentRefs:
|
|
||||||
- name: gateway
|
|
||||||
namespace: alaska
|
|
||||||
hostnames:
|
|
||||||
- "fiesta.nearle.app"
|
|
||||||
rules:
|
|
||||||
- matches:
|
|
||||||
- path:
|
|
||||||
type: PathPrefix
|
|
||||||
value: /live/api/v1/mob/orders/createorder
|
|
||||||
backendRefs:
|
|
||||||
- name: fiesta
|
|
||||||
port: 8000
|
|
||||||
- matches:
|
|
||||||
- path:
|
|
||||||
type: PathPrefix
|
|
||||||
value: /live/api/v1/web/products/create
|
|
||||||
backendRefs:
|
|
||||||
- name: fiesta
|
|
||||||
port: 8000
|
|
||||||
- matches:
|
|
||||||
- path:
|
|
||||||
type: PathPrefix
|
|
||||||
value: /
|
|
||||||
backendRefs:
|
|
||||||
- name: fiesta
|
|
||||||
port: 80
|
|
||||||
|
|||||||
@@ -47,7 +47,7 @@ spec:
|
|||||||
app: jupiter
|
app: jupiter
|
||||||
containers:
|
containers:
|
||||||
- name: jupiter
|
- name: jupiter
|
||||||
image: nearlecommerce/jupiter:v2.7.55
|
image: nearlecommerce/jupiter:v2.7.59
|
||||||
imagePullPolicy: Always
|
imagePullPolicy: Always
|
||||||
securityContext:
|
securityContext:
|
||||||
allowPrivilegeEscalation: false
|
allowPrivilegeEscalation: false
|
||||||
|
|||||||
144
manifests/nearle/riderlogs-retention-cronjob.yaml
Normal file
144
manifests/nearle/riderlogs-retention-cronjob.yaml
Normal file
@@ -0,0 +1,144 @@
|
|||||||
|
apiVersion: v1
|
||||||
|
kind: Secret
|
||||||
|
metadata:
|
||||||
|
name: nearle-redis-secrets
|
||||||
|
namespace: nearle
|
||||||
|
type: Opaque
|
||||||
|
stringData:
|
||||||
|
REDIS_HOST: "66.116.226.255"
|
||||||
|
REDIS_PORT: "6379"
|
||||||
|
REDIS_USER: "default"
|
||||||
|
REDIS_PASSWORD: "Package@324969#"
|
||||||
|
---
|
||||||
|
apiVersion: v1
|
||||||
|
kind: ConfigMap
|
||||||
|
metadata:
|
||||||
|
name: riderlogs-retention-script
|
||||||
|
namespace: nearle
|
||||||
|
data:
|
||||||
|
trim.py: |
|
||||||
|
#!/usr/bin/env python3
|
||||||
|
"""
|
||||||
|
Trims the riderlogs Redis list down to a rolling retention window.
|
||||||
|
riderlogs is append-only (RPUSH), so it's roughly chronologically
|
||||||
|
ordered; we binary-search for the first entry within the retention
|
||||||
|
window and LTRIM everything before it, rather than scanning the
|
||||||
|
whole (900K+ entry) list.
|
||||||
|
"""
|
||||||
|
import os
|
||||||
|
import json
|
||||||
|
import sys
|
||||||
|
from datetime import datetime, timedelta
|
||||||
|
|
||||||
|
import redis
|
||||||
|
|
||||||
|
RETENTION_DAYS = int(os.getenv("RETENTION_DAYS", "90"))
|
||||||
|
KEY = os.getenv("REDIS_KEY", "riderlogs")
|
||||||
|
|
||||||
|
r = redis.Redis(
|
||||||
|
host=os.environ["REDIS_HOST"],
|
||||||
|
port=int(os.environ.get("REDIS_PORT", "6379")),
|
||||||
|
username=os.environ.get("REDIS_USER", "default"),
|
||||||
|
password=os.environ["REDIS_PASSWORD"],
|
||||||
|
socket_timeout=30,
|
||||||
|
)
|
||||||
|
|
||||||
|
cutoff = datetime.utcnow() - timedelta(days=RETENTION_DAYS)
|
||||||
|
|
||||||
|
def get_logdate(idx):
|
||||||
|
v = r.lindex(KEY, idx)
|
||||||
|
if v is None:
|
||||||
|
return None
|
||||||
|
try:
|
||||||
|
d = json.loads(v)
|
||||||
|
ld = d.get("logdate")
|
||||||
|
if not ld:
|
||||||
|
return None
|
||||||
|
return datetime.strptime(ld, "%Y-%m-%d %H:%M:%S")
|
||||||
|
except Exception:
|
||||||
|
return None
|
||||||
|
|
||||||
|
n = r.llen(KEY)
|
||||||
|
if n == 0:
|
||||||
|
print(f"{KEY}: empty, nothing to do")
|
||||||
|
sys.exit(0)
|
||||||
|
|
||||||
|
oldest = get_logdate(0)
|
||||||
|
if oldest is None or oldest >= cutoff:
|
||||||
|
print(f"{KEY}: oldest entry ({oldest}) already within the {RETENTION_DAYS}-day window, nothing to trim")
|
||||||
|
sys.exit(0)
|
||||||
|
|
||||||
|
lo, hi = 0, n - 1
|
||||||
|
while lo < hi:
|
||||||
|
mid = (lo + hi) // 2
|
||||||
|
d = get_logdate(mid)
|
||||||
|
if d is None or d < cutoff:
|
||||||
|
lo = mid + 1
|
||||||
|
else:
|
||||||
|
hi = mid
|
||||||
|
|
||||||
|
before = n
|
||||||
|
r.ltrim(KEY, lo, -1)
|
||||||
|
after = r.llen(KEY)
|
||||||
|
print(f"{KEY}: cutoff={cutoff.isoformat()} trimmed {before - after} entries ({before} -> {after})")
|
||||||
|
---
|
||||||
|
apiVersion: batch/v1
|
||||||
|
kind: CronJob
|
||||||
|
metadata:
|
||||||
|
name: riderlogs-retention
|
||||||
|
namespace: nearle
|
||||||
|
spec:
|
||||||
|
schedule: "0 3 * * *"
|
||||||
|
timeZone: "Asia/Kolkata"
|
||||||
|
concurrencyPolicy: Forbid
|
||||||
|
successfulJobsHistoryLimit: 3
|
||||||
|
failedJobsHistoryLimit: 3
|
||||||
|
jobTemplate:
|
||||||
|
spec:
|
||||||
|
activeDeadlineSeconds: 300
|
||||||
|
backoffLimit: 1
|
||||||
|
template:
|
||||||
|
spec:
|
||||||
|
restartPolicy: Never
|
||||||
|
containers:
|
||||||
|
- name: trim
|
||||||
|
image: python:3.11-slim
|
||||||
|
command: ["sh", "-c", "pip install -q redis && python3 /scripts/trim.py"]
|
||||||
|
env:
|
||||||
|
- name: RETENTION_DAYS
|
||||||
|
value: "90"
|
||||||
|
- name: REDIS_KEY
|
||||||
|
value: "riderlogs"
|
||||||
|
- name: REDIS_HOST
|
||||||
|
valueFrom:
|
||||||
|
secretKeyRef:
|
||||||
|
name: nearle-redis-secrets
|
||||||
|
key: REDIS_HOST
|
||||||
|
- name: REDIS_PORT
|
||||||
|
valueFrom:
|
||||||
|
secretKeyRef:
|
||||||
|
name: nearle-redis-secrets
|
||||||
|
key: REDIS_PORT
|
||||||
|
- name: REDIS_USER
|
||||||
|
valueFrom:
|
||||||
|
secretKeyRef:
|
||||||
|
name: nearle-redis-secrets
|
||||||
|
key: REDIS_USER
|
||||||
|
- name: REDIS_PASSWORD
|
||||||
|
valueFrom:
|
||||||
|
secretKeyRef:
|
||||||
|
name: nearle-redis-secrets
|
||||||
|
key: REDIS_PASSWORD
|
||||||
|
resources:
|
||||||
|
requests:
|
||||||
|
memory: "64Mi"
|
||||||
|
cpu: "50m"
|
||||||
|
limits:
|
||||||
|
memory: "128Mi"
|
||||||
|
volumeMounts:
|
||||||
|
- name: script
|
||||||
|
mountPath: /scripts
|
||||||
|
volumes:
|
||||||
|
- name: script
|
||||||
|
configMap:
|
||||||
|
name: riderlogs-retention-script
|
||||||
Reference in New Issue
Block a user