From d98ebdd1522b6a1b27630d9c4e7ee568fc917f08 Mon Sep 17 00:00:00 2001 From: Suriya Date: Wed, 29 Jul 2026 17:34:03 +0530 Subject: [PATCH] clean up dead Fiesta NATS gateway sidecar and its routing scraps The "gateway" sidecar in the fiesta StatefulSet (workolik360/alaska image, running fiesta-gateway-script) and the k8s routing objects meant to reach it (fiesta-route HTTPRoute, the port-8000 rules in nearle-ingress) never actually served real traffic. fiesta.nearle.app is routed by a standalone Docker/nginx proxy (conf/nginx-fiesta.conf, unchanged since the initial commit) straight to NodePort 30823 - the real backend, no path splitting. The Kubernetes-native routing objects have no working controller on this cluster (no Traefik/Envoy pod; the shared Gateway resource's own status is "Waiting for controller"), so they were inert either way. The sidecar's only real effect, when briefly reachable, was publishing Fiesta orders onto the same NATS subject worker-orders drains - which forwards to jupiter, not this backend - producing the header-only phantom orders fixed in the app on 2026-07-29. It's not a working parallel path, it was the source of that bug. Removed: the gateway container + gateway-script volume from the fiesta StatefulSet, port 8000 from the fiesta Service, the fiesta-route HTTPRoute, the fiesta-gateway-script ConfigMap (fiesta-gateway.yaml deleted entirely), and the two dead port-8000 path rules in nearle-ingress's fiesta host block (kept the correct catch-all). queue.workolik.com and the core-namespace NATS worker pipeline are untouched - separate system entirely. Verified live: fiesta pods rolled to single-container, service has only port 80, and a real request through fiesta.nearle.app still behaves correctly post-cleanup. --- manifests/core/ingress-unified.yaml | 21 +- manifests/nearle/fiesta-gateway.yaml | 514 --------------------------- manifests/nearle/kustomization.yaml | 1 - manifests/nearle/nearle-fiesta.yaml | 70 ---- 4 files changed, 7 insertions(+), 599 deletions(-) delete mode 100644 manifests/nearle/fiesta-gateway.yaml diff --git a/manifests/core/ingress-unified.yaml b/manifests/core/ingress-unified.yaml index 27acb2b..1b9aa5f 100644 --- a/manifests/core/ingress-unified.yaml +++ b/manifests/core/ingress-unified.yaml @@ -43,20 +43,13 @@ spec: - host: fiesta.nearle.app http: paths: - - path: /live/api/v1/mob/orders/createorder - pathType: Prefix - backend: - service: - name: fiesta - port: - number: 8000 - - path: /live/api/v1/web/products/create - pathType: Prefix - backend: - service: - name: fiesta - port: - number: 8000 + # Real routing for this host is the standalone Docker/nginx proxy + # (conf/nginx-fiesta.conf), which forwards everything to NodePort + # 30823 (this backend, port 80) - no path splitting, no gateway + # sidecar. This Ingress has no controller on the cluster (confirmed: + # no Traefik/Envoy pod, no IngressClass) so it's inert either way, + # but kept as a single correct catch-all rather than a stale + # path-split that pointed part of it at a dead gateway sidecar. - path: / pathType: Prefix backend: diff --git a/manifests/nearle/fiesta-gateway.yaml b/manifests/nearle/fiesta-gateway.yaml deleted file mode 100644 index 9befdc6..0000000 --- a/manifests/nearle/fiesta-gateway.yaml +++ /dev/null @@ -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) - diff --git a/manifests/nearle/kustomization.yaml b/manifests/nearle/kustomization.yaml index fc97cd4..a683661 100644 --- a/manifests/nearle/kustomization.yaml +++ b/manifests/nearle/kustomization.yaml @@ -6,7 +6,6 @@ resources: - nearle-config.yaml - nearle-secrets.yaml - nearle-app-secrets.yaml - - fiesta-gateway.yaml - nearle-fiesta.yaml - nearle-jupiter.yaml - jupiter-cors-proxy.yaml diff --git a/manifests/nearle/nearle-fiesta.yaml b/manifests/nearle/nearle-fiesta.yaml index aee1d8e..d3e9c22 100644 --- a/manifests/nearle/nearle-fiesta.yaml +++ b/manifests/nearle/nearle-fiesta.yaml @@ -74,39 +74,6 @@ spec: secretKeyRef: name: nats-credentials key: password - - name: gateway - image: workolik360/alaska:v1.2.0 - imagePullPolicy: Always - securityContext: - allowPrivilegeEscalation: false - capabilities: - drop: - - ALL - ports: - - containerPort: 8000 - name: http - volumeMounts: - - name: gateway-script - mountPath: /app/app.py - subPath: app.py - envFrom: - - 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 kind: Service @@ -123,42 +90,5 @@ spec: nodePort: 30823 protocol: TCP name: main - - port: 8000 - targetPort: 8000 - name: gateway - protocol: TCP selector: 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: - # mob/orders/createorder and web/products/create used to be routed to the - # NATS gateway sidecar (port 8000, ConfigMap fiesta-gateway-script). That - # sidecar publishes onto the shared "ORDERS"/"PRODUCTS" NATS subjects, - # whose only consumer (worker-orders/worker-products in the core namespace) - # has EXTERNAL_BASE_URL hardcoded to jupiter, not this backend. Fiesta has - # no NATS consumer of its own, so every order/product submitted through - # that path was silently created on jupiter (header-only, wrong schema) - # instead of ever reaching CreateOrderv3 here. Both paths are already - # registered directly on this backend (routes/orderroutes.go, - # routes/productroutes.go), so route them straight through synchronously - # via the catch-all rule below instead of overriding them to the sidecar. - - matches: - - path: - type: PathPrefix - value: / - backendRefs: - - name: fiesta - port: 80