Skip to main content
EVOKORE// BROWSE
>

./browse/prompts

11 NODES
šŸ¤–system prompt•7 months ago

task-coordination-strategies

Decompose complex tasks, design dependency graphs, and coordinate

coding
⭐1
# Task Coordination Strategies Strategies for decomposing complex tasks into parallelizable units, designing dependency graphs, writing effective task descriptions, and monitoring workload across agent teams. ## When to Use This Skill - Breaking down a complex task for parallel execution - Designing task dependency relationships (blockedBy/blocks) - Writing task descriptions with clear acceptance criteria - Monitoring and rebalancing workload across teammates - Identifying the critical path in a multi-task workflow ## Task Decomposition Strategies ### By Layer Split work by architectural layer: - Frontend components - Backend API endpoints - Database migrations/models - Test suites **Best for**: Full-stack features, vertical slices ### By Component Split work by functional component: - Authentication module - User profile module - Notification module **Best for**: Microservices, modular architectures ### By Concern Split work by cross-cutting concern: - Security review - Performance review - Architecture review **Best for**: Code reviews, audits ### By File Ownership Split work by file/directory boundaries: - `src/components/` — Implementer 1 - `src/api/` — Implementer 2 - `src/utils/` — Implementer 3 **Best for**: Parallel implementation, conflict avoidance ## Dependency Graph Design ### Principles 1. **Minimize chain depth** — Prefer wide, shallow graphs over deep chains 2. **Identify the critical path** — The longest chain determines minimum completion time 3. **Use blockedBy sparingly** — Only add dependencies that are truly required 4. **Avoid circular dependencies** — Task A blocks B blocks A is a deadlock ### Patterns **Independent (Best parallelism)**: ``` Task A ─┐ Task B ─┼─→ Integration Task C ā”€ā”˜ ``` **Sequential (Necessary dependencies)**: ``` Task A → Task B → Task C ``` **Diamond (Mixed)**: ``` ā”Œā†’ Task B ─┐ Task A ─┤ ā”œā†’ Task D └→ Task C ā”€ā”˜ ``` ### Using blockedBy/blocks ``` TaskCreate: { subject: "Build API endpoints" } → Task #1 TaskCreate: { subject: "Build frontend components" } → Task #2 TaskCreate: { subject: "Integration testing" } → Task #3 TaskUpdate: { taskId: "3", addBlockedBy: ["1", "2"] } → #3 waits for #1 and #2 ``` ## Task Description Best Practices Every task should include: 1. **Objective** — What needs to be accomplished (1-2 sentences) 2. **Owned Files** — Explicit list of files/directories this teammate may modify 3. **Requirements** — Specific deliverables or behaviors expected 4. **Interface Contracts** — How this work connects to other teammates' work 5. **Acceptance Criteria** — How to verify the task is done correctly 6. **Scope Boundaries** — What is explicitly out of scope ### Template ``` ## Objective Build the user authentication API endpoints. ## Owned Files - src/api/auth.ts - src/api/middleware/auth-middleware.ts - src/types/auth.ts (shared — read only, do not modify) ## Requirements - POST /api/login — accepts email/password, returns JWT - POST /api/register — creates new user, returns JWT - GET /api/me — returns current user profile (requires auth) ## Interface Contract - Import User type from src/types/auth.ts (owned by implementer-1) - Export AuthResponse type for frontend consumption ## Acceptance Criteria - All endpoints return proper HTTP status codes - JWT tokens expire after 24 hours - Passwords are hashed with bcrypt ## Out of Scope - OAuth/social login - Password reset flow - Rate limiting ``` ## Workload Monitoring ### Indicators of Imbalance | Signal | Meaning | Action | | -------------------------- | ------------------- | --------------------------- | | Teammate idle, others busy | Uneven distribution | Reassign pending tasks | | Teammate stuck on one task | Possible blocker | Check in, offer help | | All tasks blocked | Dependency issue | Resolve critical path first | | One teammate has 3x others | Overloaded | Split tasks or reassign | ### Rebalancing Steps 1. Call `TaskList` to assess current state 2. Identify idle or overloaded teammates 3. Use `TaskUpdate` to reassign tasks 4. Use `SendMessage` to notify affected teammates 5. Monitor for improved throughput
šŸ‘0
šŸ‘ļø0
šŸ¤– Auto-discovered
šŸ¤–system prompt•7 months ago

fastapi-templates

Create production-ready FastAPI projects with async patterns,

coding
⭐1
# FastAPI Project Templates Production-ready FastAPI project structures with async patterns, dependency injection, middleware, and best practices for building high-performance APIs. ## When to Use This Skill - Starting new FastAPI projects from scratch - Implementing async REST APIs with Python - Building high-performance web services and microservices - Creating async applications with PostgreSQL, MongoDB - Setting up API projects with proper structure and testing ## Core Concepts ### 1. Project Structure **Recommended Layout:** ``` app/ ā”œā”€ā”€ api/ # API routes │ ā”œā”€ā”€ v1/ │ │ ā”œā”€ā”€ endpoints/ │ │ │ ā”œā”€ā”€ users.py │ │ │ ā”œā”€ā”€ auth.py │ │ │ └── items.py │ │ └── router.py │ └── dependencies.py # Shared dependencies ā”œā”€ā”€ core/ # Core configuration │ ā”œā”€ā”€ config.py │ ā”œā”€ā”€ security.py │ └── database.py ā”œā”€ā”€ models/ # Database models │ ā”œā”€ā”€ user.py │ └── item.py ā”œā”€ā”€ schemas/ # Pydantic schemas │ ā”œā”€ā”€ user.py │ └── item.py ā”œā”€ā”€ services/ # Business logic │ ā”œā”€ā”€ user_service.py │ └── auth_service.py ā”œā”€ā”€ repositories/ # Data access │ ā”œā”€ā”€ user_repository.py │ └── item_repository.py └── main.py # Application entry ``` ### 2. Dependency Injection FastAPI's built-in DI system using `Depends`: - Database session management - Authentication/authorization - Shared business logic - Configuration injection ### 3. Async Patterns Proper async/await usage: - Async route handlers - Async database operations - Async background tasks - Async middleware ## Implementation Patterns ### Pattern 1: Complete FastAPI Application ```python # main.py from fastapi import FastAPI, Depends from fastapi.middleware.cors import CORSMiddleware from contextlib import asynccontextmanager @asynccontextmanager async def lifespan(app: FastAPI): """Application lifespan events.""" # Startup await database.connect() yield # Shutdown await database.disconnect() app = FastAPI( title="API Template", version="1.0.0", lifespan=lifespan ) # CORS middleware app.add_middleware( CORSMiddleware, allow_origins=["*"], allow_credentials=True, allow_methods=["*"], allow_headers=["*"], ) # Include routers from app.api.v1.router import api_router app.include_router(api_router, prefix="/api/v1") # core/config.py from pydantic_settings import BaseSettings from functools import lru_cache class Settings(BaseSettings): """Application settings.""" DATABASE_URL: str SECRET_KEY: str ACCESS_TOKEN_EXPIRE_MINUTES: int = 30 API_V1_STR: str = "/api/v1" class Config: env_file = ".env" @lru_cache() def get_settings() -> Settings: return Settings() # core/database.py from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession from sqlalchemy.ext.declarative import declarative_base from sqlalchemy.orm import sessionmaker from app.core.config import get_settings settings = get_settings() engine = create_async_engine( settings.DATABASE_URL, echo=True, future=True ) AsyncSessionLocal = sessionmaker( engine, class_=AsyncSession, expire_on_commit=False ) Base = declarative_base() async def get_db() -> AsyncSession: """Dependency for database session.""" async with AsyncSessionLocal() as session: try: yield session await session.commit() except Exception: await session.rollback() raise finally: await session.close() ``` ### Pattern 2: CRUD Repository Pattern ```python # repositories/base_repository.py from typing import Generic, TypeVar, Type, Optional, List from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy import select from pydantic import BaseModel ModelType = TypeVar("ModelType") CreateSchemaType = TypeVar("CreateSchemaType", bound=BaseModel) UpdateSchemaType = TypeVar("UpdateSchemaType", bound=BaseModel) class BaseRepository(Generic[ModelType, CreateSchemaType, UpdateSchemaType]): """Base repository for CRUD operations.""" def __init__(self, model: Type[ModelType]): self.model = model async def get(self, db: AsyncSession, id: int) -> Optional[ModelType]: """Get by ID.""" result = await db.execute( select(self.model).where(self.model.id == id) ) return result.scalars().first() async def get_multi( self, db: AsyncSession, skip: int = 0, limit: int = 100 ) -> List[ModelType]: """Get multiple records.""" result = await db.execute( select(self.model).offset(skip).limit(limit) ) return result.scalars().all() async def create( self, db: AsyncSession, obj_in: CreateSchemaType ) -> ModelType: """Create new record.""" db_obj = self.model(**obj_in.dict()) db.add(db_obj) await db.flush() await db.refresh(db_obj) return db_obj async def update( self, db: AsyncSession, db_obj: ModelType, obj_in: UpdateSchemaType ) -> ModelType: """Update record.""" update_data = obj_in.dict(exclude_unset=True) for field, value in update_data.items(): setattr(db_obj, field, value) await db.flush() await db.refresh(db_obj) return db_obj async def delete(self, db: AsyncSession, id: int) -> bool: """Delete record.""" obj = await self.get(db, id) if obj: await db.delete(obj) return True return False # repositories/user_repository.py from app.repositories.base_repository import BaseRepository from app.models.user import User from app.schemas.user import UserCreate, UserUpdate class UserRepository(BaseRepository[User, UserCreate, UserUpdate]): """User-specific repository.""" async def get_by_email(self, db: AsyncSession, email: str) -> Optional[User]: """Get user by email.""" result = await db.execute( select(User).where(User.email == email) ) return result.scalars().first() async def is_active(self, db: AsyncSession, user_id: int) -> bool: """Check if user is active.""" user = await self.get(db, user_id) return user.is_active if user else False user_repository = UserRepository(User) ``` ### Pattern 3: Service Layer ```python # services/user_service.py from typing import Optional from sqlalchemy.ext.asyncio import AsyncSession from app.repositories.user_repository import user_repository from app.schemas.user import UserCreate, UserUpdate, User from app.core.security import get_password_hash, verify_password class UserService: """Business logic for users.""" def __init__(self): self.repository = user_repository async def create_user( self, db: AsyncSession, user_in: UserCreate ) -> User: """Create new user with hashed password.""" # Check if email exists existing = await self.repository.get_by_email(db, user_in.email) if existing: raise ValueError("Email already registered") # Hash password user_in_dict = user_in.dict() user_in_dict["hashed_password"] = get_password_hash(user_in_dict.pop("password")) # Create user user = await self.repository.create(db, UserCreate(**user_in_dict)) return user async def authenticate( self, db: AsyncSession, email: str, password: str ) -> Optional[User]: """Authenticate user.""" user = await self.repository.get_by_email(db, email) if not user: return None if not verify_password(password, user.hashed_password): return None return user async def update_user( self, db: AsyncSession, user_id: int, user_in: UserUpdate ) -> Optional[User]: """Update user.""" user = await self.repository.get(db, user_id) if not user: return None if user_in.password: user_in_dict = user_in.dict(exclude_unset=True) user_in_dict["hashed_password"] = get_password_hash( user_in_dict.pop("password") ) user_in = UserUpdate(**user_in_dict) return await self.repository.update(db, user, user_in) user_service = UserService() ``` ### Pattern 4: API Endpoints with Dependencies ```python # api/v1/endpoints/users.py from fastapi import APIRouter, Depends, HTTPException, status from sqlalchemy.ext.asyncio import AsyncSession from typing import List from app.core.database import get_db from app.schemas.user import User, UserCreate, UserUpdate from app.services.user_service import user_service from app.api.dependencies import get_current_user router = APIRouter() @router.post("/", response_model=User, status_code=status.HTTP_201_CREATED) async def create_user( user_in: UserCreate, db: AsyncSession = Depends(get_db) ): """Create new user.""" try: user = await user_service.create_user(db, user_in) return user except ValueError as e: raise HTTPException(status_code=400, detail=str(e)) @router.get("/me", response_model=User) async def read_current_user( current_user: User = Depends(get_current_user) ): """Get current user.""" return current_user @router.get("/{user_id}", response_model=User) async def read_user( user_id: int, db: AsyncSession = Depends(get_db), current_user: User = Depends(get_current_user) ): """Get user by ID.""" user = await user_service.repository.get(db, user_id) if not user: raise HTTPException(status_code=404, detail="User not found") return user @router.patch("/{user_id}", response_model=User) async def update_user( user_id: int, user_in: UserUpdate, db: AsyncSession = Depends(get_db), current_user: User = Depends(get_current_user) ): """Update user.""" if current_user.id != user_id: raise HTTPException(status_code=403, detail="Not authorized") user = await user_service.update_user(db, user_id, user_in) if not user: raise HTTPException(status_code=404, detail="User not found") return user @router.delete("/{user_id}", status_code=status.HTTP_204_NO_CONTENT) async def delete_user( user_id: int, db: AsyncSession = Depends(get_db), current_user: User = Depends(get_current_user) ): """Delete user.""" if current_user.id != user_id: raise HTTPException(status_code=403, detail="Not authorized") deleted = await user_service.repository.delete(db, user_id) if not deleted: raise HTTPException(status_code=404, detail="User not found") ``` ### Pattern 5: Authentication & Authorization ```python # core/security.py from datetime import datetime, timedelta from typing import Optional from jose import JWTError, jwt from passlib.context import CryptContext from app.core.config import get_settings settings = get_settings() pwd_context = CryptContext(schemes=["bcrypt"], deprecated="auto") ALGORITHM = "HS256" def create_access_token(data: dict, expires_delta: Optional[timedelta] = None): """Create JWT access token.""" to_encode = data.copy() if expires_delta: expire = datetime.utcnow() + expires_delta else: expire = datetime.utcnow() + timedelta(minutes=15) to_encode.update({"exp": expire}) encoded_jwt = jwt.encode(to_encode, settings.SECRET_KEY, algorithm=ALGORITHM) return encoded_jwt def verify_password(plain_password: str, hashed_password: str) -> bool: """Verify password against hash.""" return pwd_context.verify(plain_password, hashed_password) def get_password_hash(password: str) -> str: """Hash password.""" return pwd_context.hash(password) # api/dependencies.py from fastapi import Depends, HTTPException, status from fastapi.security import OAuth2PasswordBearer from jose import JWTError, jwt from sqlalchemy.ext.asyncio import AsyncSession from app.core.database import get_db from app.core.security import ALGORITHM from app.core.config import get_settings from app.repositories.user_repository import user_repository oauth2_scheme = OAuth2PasswordBearer(tokenUrl=f"{settings.API_V1_STR}/auth/login") async def get_current_user( db: AsyncSession = Depends(get_db), token: str = Depends(oauth2_scheme) ): """Get current authenticated user.""" credentials_exception = HTTPException( status_code=status.HTTP_401_UNAUTHORIZED, detail="Could not validate credentials", headers={"WWW-Authenticate": "Bearer"}, ) try: payload = jwt.decode(token, settings.SECRET_KEY, algorithms=[ALGORITHM]) user_id: int = payload.get("sub") if user_id is None: raise credentials_exception except JWTError: raise credentials_exception user = await user_repository.get(db, user_id) if user is None: raise credentials_exception return user ``` ## Testing ```python # tests/conftest.py import pytest import asyncio from httpx import AsyncClient from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession from sqlalchemy.orm import sessionmaker from app.main import app from app.core.database import get_db, Base TEST_DATABASE_URL = "sqlite+aiosqlite:///:memory:" @pytest.fixture(scope="session") def event_loop(): loop = asyncio.get_event_loop_policy().new_event_loop() yield loop loop.close() @pytest.fixture async def db_session(): engine = create_async_engine(TEST_DATABASE_URL, echo=True) async with engine.begin() as conn: await conn.run_sync(Base.metadata.create_all) AsyncSessionLocal = sessionmaker( engine, class_=AsyncSession, expire_on_commit=False ) async with AsyncSessionLocal() as session: yield session @pytest.fixture async def client(db_session): async def override_get_db(): yield db_session app.dependency_overrides[get_db] = override_get_db async with AsyncClient(app=app, base_url="http://test") as client: yield client # tests/test_users.py import pytest @pytest.mark.asyncio async def test_create_user(client): response = await client.post( "/api/v1/users/", json={ "email": "test@example.com", "password": "testpass123", "name": "Test User" } ) assert response.status_code == 201 data = response.json() assert data["email"] == "test@example.com" assert "id" in data ``` ## Resources - **references/fastapi-architecture.md**: Detailed architecture guide - **references/async-best-practices.md**: Async/await patterns - **references/testing-strategies.md**: Comprehensive testing guide - **assets/project-template/**: Complete FastAPI project - **assets/docker-compose.yml**: Development environment setup ## Best Practices 1. **Async All The Way**: Use async for database, external APIs 2. **Dependency Injection**: Leverage FastAPI's DI system 3. **Repository Pattern**: Separate data access from business logic 4. **Service Layer**: Keep business logic out of routes 5. **Pydantic Schemas**: Strong typing for request/response 6. **Error Handling**: Consistent error responses 7. **Testing**: Test all layers independently ## Common Pitfalls - **Blocking Code in Async**: Using synchronous database drivers - **No Service Layer**: Business logic in route handlers - **Missing Type Hints**: Loses FastAPI's benefits - **Ignoring Sessions**: Not properly managing database sessions - **No Testing**: Skipping integration tests - **Tight Coupling**: Direct database access in routes
šŸ‘0
šŸ‘ļø0
šŸ¤– Auto-discovered
šŸ¤–system prompt•7 months ago

architecture-patterns

Implement proven backend architecture patterns including Clean

coding
⭐1
# Architecture Patterns Master proven backend architecture patterns including Clean Architecture, Hexagonal Architecture, and Domain-Driven Design to build maintainable, testable, and scalable systems. ## When to Use This Skill - Designing new backend systems from scratch - Refactoring monolithic applications for better maintainability - Establishing architecture standards for your team - Migrating from tightly coupled to loosely coupled architectures - Implementing domain-driven design principles - Creating testable and mockable codebases - Planning microservices decomposition ## Core Concepts ### 1. Clean Architecture (Uncle Bob) **Layers (dependency flows inward):** - **Entities**: Core business models - **Use Cases**: Application business rules - **Interface Adapters**: Controllers, presenters, gateways - **Frameworks & Drivers**: UI, database, external services **Key Principles:** - Dependencies point inward - Inner layers know nothing about outer layers - Business logic independent of frameworks - Testable without UI, database, or external services ### 2. Hexagonal Architecture (Ports and Adapters) **Components:** - **Domain Core**: Business logic - **Ports**: Interfaces defining interactions - **Adapters**: Implementations of ports (database, REST, message queue) **Benefits:** - Swap implementations easily (mock for testing) - Technology-agnostic core - Clear separation of concerns ### 3. Domain-Driven Design (DDD) **Strategic Patterns:** - **Bounded Contexts**: Separate models for different domains - **Context Mapping**: How contexts relate - **Ubiquitous Language**: Shared terminology **Tactical Patterns:** - **Entities**: Objects with identity - **Value Objects**: Immutable objects defined by attributes - **Aggregates**: Consistency boundaries - **Repositories**: Data access abstraction - **Domain Events**: Things that happened ## Clean Architecture Pattern ### Directory Structure ``` app/ ā”œā”€ā”€ domain/ # Entities & business rules │ ā”œā”€ā”€ entities/ │ │ ā”œā”€ā”€ user.py │ │ └── order.py │ ā”œā”€ā”€ value_objects/ │ │ ā”œā”€ā”€ email.py │ │ └── money.py │ └── interfaces/ # Abstract interfaces │ ā”œā”€ā”€ user_repository.py │ └── payment_gateway.py ā”œā”€ā”€ use_cases/ # Application business rules │ ā”œā”€ā”€ create_user.py │ ā”œā”€ā”€ process_order.py │ └── send_notification.py ā”œā”€ā”€ adapters/ # Interface implementations │ ā”œā”€ā”€ repositories/ │ │ ā”œā”€ā”€ postgres_user_repository.py │ │ └── redis_cache_repository.py │ ā”œā”€ā”€ controllers/ │ │ └── user_controller.py │ └── gateways/ │ ā”œā”€ā”€ stripe_payment_gateway.py │ └── sendgrid_email_gateway.py └── infrastructure/ # Framework & external concerns ā”œā”€ā”€ database.py ā”œā”€ā”€ config.py └── logging.py ``` ### Implementation Example ```python # domain/entities/user.py from dataclasses import dataclass from datetime import datetime from typing import Optional @dataclass class User: """Core user entity - no framework dependencies.""" id: str email: str name: str created_at: datetime is_active: bool = True def deactivate(self): """Business rule: deactivating user.""" self.is_active = False def can_place_order(self) -> bool: """Business rule: active users can order.""" return self.is_active # domain/interfaces/user_repository.py from abc import ABC, abstractmethod from typing import Optional, List from domain.entities.user import User class IUserRepository(ABC): """Port: defines contract, no implementation.""" @abstractmethod async def find_by_id(self, user_id: str) -> Optional[User]: pass @abstractmethod async def find_by_email(self, email: str) -> Optional[User]: pass @abstractmethod async def save(self, user: User) -> User: pass @abstractmethod async def delete(self, user_id: str) -> bool: pass # use_cases/create_user.py from domain.entities.user import User from domain.interfaces.user_repository import IUserRepository from dataclasses import dataclass from datetime import datetime import uuid @dataclass class CreateUserRequest: email: str name: str @dataclass class CreateUserResponse: user: User success: bool error: Optional[str] = None class CreateUserUseCase: """Use case: orchestrates business logic.""" def __init__(self, user_repository: IUserRepository): self.user_repository = user_repository async def execute(self, request: CreateUserRequest) -> CreateUserResponse: # Business validation existing = await self.user_repository.find_by_email(request.email) if existing: return CreateUserResponse( user=None, success=False, error="Email already exists" ) # Create entity user = User( id=str(uuid.uuid4()), email=request.email, name=request.name, created_at=datetime.now(), is_active=True ) # Persist saved_user = await self.user_repository.save(user) return CreateUserResponse( user=saved_user, success=True ) # adapters/repositories/postgres_user_repository.py from domain.interfaces.user_repository import IUserRepository from domain.entities.user import User from typing import Optional import asyncpg class PostgresUserRepository(IUserRepository): """Adapter: PostgreSQL implementation.""" def __init__(self, pool: asyncpg.Pool): self.pool = pool async def find_by_id(self, user_id: str) -> Optional[User]: async with self.pool.acquire() as conn: row = await conn.fetchrow( "SELECT * FROM users WHERE id = $1", user_id ) return self._to_entity(row) if row else None async def find_by_email(self, email: str) -> Optional[User]: async with self.pool.acquire() as conn: row = await conn.fetchrow( "SELECT * FROM users WHERE email = $1", email ) return self._to_entity(row) if row else None async def save(self, user: User) -> User: async with self.pool.acquire() as conn: await conn.execute( """ INSERT INTO users (id, email, name, created_at, is_active) VALUES ($1, $2, $3, $4, $5) ON CONFLICT (id) DO UPDATE SET email = $2, name = $3, is_active = $5 """, user.id, user.email, user.name, user.created_at, user.is_active ) return user async def delete(self, user_id: str) -> bool: async with self.pool.acquire() as conn: result = await conn.execute( "DELETE FROM users WHERE id = $1", user_id ) return result == "DELETE 1" def _to_entity(self, row) -> User: """Map database row to entity.""" return User( id=row["id"], email=row["email"], name=row["name"], created_at=row["created_at"], is_active=row["is_active"] ) # adapters/controllers/user_controller.py from fastapi import APIRouter, Depends, HTTPException from use_cases.create_user import CreateUserUseCase, CreateUserRequest from pydantic import BaseModel router = APIRouter() class CreateUserDTO(BaseModel): email: str name: str @router.post("/users") async def create_user( dto: CreateUserDTO, use_case: CreateUserUseCase = Depends(get_create_user_use_case) ): """Controller: handles HTTP concerns only.""" request = CreateUserRequest(email=dto.email, name=dto.name) response = await use_case.execute(request) if not response.success: raise HTTPException(status_code=400, detail=response.error) return {"user": response.user} ``` ## Hexagonal Architecture Pattern ```python # Core domain (hexagon center) class OrderService: """Domain service - no infrastructure dependencies.""" def __init__( self, order_repository: OrderRepositoryPort, payment_gateway: PaymentGatewayPort, notification_service: NotificationPort ): self.orders = order_repository self.payments = payment_gateway self.notifications = notification_service async def place_order(self, order: Order) -> OrderResult: # Business logic if not order.is_valid(): return OrderResult(success=False, error="Invalid order") # Use ports (interfaces) payment = await self.payments.charge( amount=order.total, customer=order.customer_id ) if not payment.success: return OrderResult(success=False, error="Payment failed") order.mark_as_paid() saved_order = await self.orders.save(order) await self.notifications.send( to=order.customer_email, subject="Order confirmed", body=f"Order {order.id} confirmed" ) return OrderResult(success=True, order=saved_order) # Ports (interfaces) class OrderRepositoryPort(ABC): @abstractmethod async def save(self, order: Order) -> Order: pass class PaymentGatewayPort(ABC): @abstractmethod async def charge(self, amount: Money, customer: str) -> PaymentResult: pass class NotificationPort(ABC): @abstractmethod async def send(self, to: str, subject: str, body: str): pass # Adapters (implementations) class StripePaymentAdapter(PaymentGatewayPort): """Primary adapter: connects to Stripe API.""" def __init__(self, api_key: str): self.stripe = stripe self.stripe.api_key = api_key async def charge(self, amount: Money, customer: str) -> PaymentResult: try: charge = self.stripe.Charge.create( amount=amount.cents, currency=amount.currency, customer=customer ) return PaymentResult(success=True, transaction_id=charge.id) except stripe.error.CardError as e: return PaymentResult(success=False, error=str(e)) class MockPaymentAdapter(PaymentGatewayPort): """Test adapter: no external dependencies.""" async def charge(self, amount: Money, customer: str) -> PaymentResult: return PaymentResult(success=True, transaction_id="mock-123") ``` ## Domain-Driven Design Pattern ```python # Value Objects (immutable) from dataclasses import dataclass from typing import Optional @dataclass(frozen=True) class Email: """Value object: validated email.""" value: str def __post_init__(self): if "@" not in self.value: raise ValueError("Invalid email") @dataclass(frozen=True) class Money: """Value object: amount with currency.""" amount: int # cents currency: str def add(self, other: "Money") -> "Money": if self.currency != other.currency: raise ValueError("Currency mismatch") return Money(self.amount + other.amount, self.currency) # Entities (with identity) class Order: """Entity: has identity, mutable state.""" def __init__(self, id: str, customer: Customer): self.id = id self.customer = customer self.items: List[OrderItem] = [] self.status = OrderStatus.PENDING self._events: List[DomainEvent] = [] def add_item(self, product: Product, quantity: int): """Business logic in entity.""" item = OrderItem(product, quantity) self.items.append(item) self._events.append(ItemAddedEvent(self.id, item)) def total(self) -> Money: """Calculated property.""" return sum(item.subtotal() for item in self.items) def submit(self): """State transition with business rules.""" if not self.items: raise ValueError("Cannot submit empty order") if self.status != OrderStatus.PENDING: raise ValueError("Order already submitted") self.status = OrderStatus.SUBMITTED self._events.append(OrderSubmittedEvent(self.id)) # Aggregates (consistency boundary) class Customer: """Aggregate root: controls access to entities.""" def __init__(self, id: str, email: Email): self.id = id self.email = email self._addresses: List[Address] = [] self._orders: List[str] = [] # Order IDs, not full objects def add_address(self, address: Address): """Aggregate enforces invariants.""" if len(self._addresses) >= 5: raise ValueError("Maximum 5 addresses allowed") self._addresses.append(address) @property def primary_address(self) -> Optional[Address]: return next((a for a in self._addresses if a.is_primary), None) # Domain Events @dataclass class OrderSubmittedEvent: order_id: str occurred_at: datetime = field(default_factory=datetime.now) # Repository (aggregate persistence) class OrderRepository: """Repository: persist/retrieve aggregates.""" async def find_by_id(self, order_id: str) -> Optional[Order]: """Reconstitute aggregate from storage.""" pass async def save(self, order: Order): """Persist aggregate and publish events.""" await self._persist(order) await self._publish_events(order._events) order._events.clear() ``` ## Resources - **references/clean-architecture-guide.md**: Detailed layer breakdown - **references/hexagonal-architecture-guide.md**: Ports and adapters patterns - **references/ddd-tactical-patterns.md**: Entities, value objects, aggregates - **assets/clean-architecture-template/**: Complete project structure - **assets/ddd-examples/**: Domain modeling examples ## Best Practices 1. **Dependency Rule**: Dependencies always point inward 2. **Interface Segregation**: Small, focused interfaces 3. **Business Logic in Domain**: Keep frameworks out of core 4. **Test Independence**: Core testable without infrastructure 5. **Bounded Contexts**: Clear domain boundaries 6. **Ubiquitous Language**: Consistent terminology 7. **Thin Controllers**: Delegate to use cases 8. **Rich Domain Models**: Behavior with data ## Common Pitfalls - **Anemic Domain**: Entities with only data, no behavior - **Framework Coupling**: Business logic depends on frameworks - **Fat Controllers**: Business logic in controllers - **Repository Leakage**: Exposing ORM objects - **Missing Abstractions**: Concrete dependencies in core - **Over-Engineering**: Clean architecture for simple CRUD
šŸ‘0
šŸ‘ļø0
šŸ¤– Auto-discovered
šŸ¤–system prompt•7 months ago

microservices-patterns

Design microservices architectures with service boundaries,

coding
⭐1
# Microservices Patterns Master microservices architecture patterns including service boundaries, inter-service communication, data management, and resilience patterns for building distributed systems. ## When to Use This Skill - Decomposing monoliths into microservices - Designing service boundaries and contracts - Implementing inter-service communication - Managing distributed data and transactions - Building resilient distributed systems - Implementing service discovery and load balancing - Designing event-driven architectures ## Core Concepts ### 1. Service Decomposition Strategies **By Business Capability** - Organize services around business functions - Each service owns its domain - Example: OrderService, PaymentService, InventoryService **By Subdomain (DDD)** - Core domain, supporting subdomains - Bounded contexts map to services - Clear ownership and responsibility **Strangler Fig Pattern** - Gradually extract from monolith - New functionality as microservices - Proxy routes to old/new systems ### 2. Communication Patterns **Synchronous (Request/Response)** - REST APIs - gRPC - GraphQL **Asynchronous (Events/Messages)** - Event streaming (Kafka) - Message queues (RabbitMQ, SQS) - Pub/Sub patterns ### 3. Data Management **Database Per Service** - Each service owns its data - No shared databases - Loose coupling **Saga Pattern** - Distributed transactions - Compensating actions - Eventual consistency ### 4. Resilience Patterns **Circuit Breaker** - Fail fast on repeated errors - Prevent cascade failures **Retry with Backoff** - Transient fault handling - Exponential backoff **Bulkhead** - Isolate resources - Limit impact of failures ## Service Decomposition Patterns ### Pattern 1: By Business Capability ```python # E-commerce example # Order Service class OrderService: """Handles order lifecycle.""" async def create_order(self, order_data: dict) -> Order: order = Order.create(order_data) # Publish event for other services await self.event_bus.publish( OrderCreatedEvent( order_id=order.id, customer_id=order.customer_id, items=order.items, total=order.total ) ) return order # Payment Service (separate service) class PaymentService: """Handles payment processing.""" async def process_payment(self, payment_request: PaymentRequest) -> PaymentResult: # Process payment result = await self.payment_gateway.charge( amount=payment_request.amount, customer=payment_request.customer_id ) if result.success: await self.event_bus.publish( PaymentCompletedEvent( order_id=payment_request.order_id, transaction_id=result.transaction_id ) ) return result # Inventory Service (separate service) class InventoryService: """Handles inventory management.""" async def reserve_items(self, order_id: str, items: List[OrderItem]) -> ReservationResult: # Check availability for item in items: available = await self.inventory_repo.get_available(item.product_id) if available < item.quantity: return ReservationResult( success=False, error=f"Insufficient inventory for {item.product_id}" ) # Reserve items reservation = await self.create_reservation(order_id, items) await self.event_bus.publish( InventoryReservedEvent( order_id=order_id, reservation_id=reservation.id ) ) return ReservationResult(success=True, reservation=reservation) ``` ### Pattern 2: API Gateway ```python from fastapi import FastAPI, HTTPException, Depends import httpx from circuitbreaker import circuit app = FastAPI() class APIGateway: """Central entry point for all client requests.""" def __init__(self): self.order_service_url = "http://order-service:8000" self.payment_service_url = "http://payment-service:8001" self.inventory_service_url = "http://inventory-service:8002" self.http_client = httpx.AsyncClient(timeout=5.0) @circuit(failure_threshold=5, recovery_timeout=30) async def call_order_service(self, path: str, method: str = "GET", **kwargs): """Call order service with circuit breaker.""" response = await self.http_client.request( method, f"{self.order_service_url}{path}", **kwargs ) response.raise_for_status() return response.json() async def create_order_aggregate(self, order_id: str) -> dict: """Aggregate data from multiple services.""" # Parallel requests order, payment, inventory = await asyncio.gather( self.call_order_service(f"/orders/{order_id}"), self.call_payment_service(f"/payments/order/{order_id}"), self.call_inventory_service(f"/reservations/order/{order_id}"), return_exceptions=True ) # Handle partial failures result = {"order": order} if not isinstance(payment, Exception): result["payment"] = payment if not isinstance(inventory, Exception): result["inventory"] = inventory return result @app.post("/api/orders") async def create_order( order_data: dict, gateway: APIGateway = Depends() ): """API Gateway endpoint.""" try: # Route to order service order = await gateway.call_order_service( "/orders", method="POST", json=order_data ) return {"order": order} except httpx.HTTPError as e: raise HTTPException(status_code=503, detail="Order service unavailable") ``` ## Communication Patterns ### Pattern 1: Synchronous REST Communication ```python # Service A calls Service B import httpx from tenacity import retry, stop_after_attempt, wait_exponential class ServiceClient: """HTTP client with retries and timeout.""" def __init__(self, base_url: str): self.base_url = base_url self.client = httpx.AsyncClient( timeout=httpx.Timeout(5.0, connect=2.0), limits=httpx.Limits(max_keepalive_connections=20) ) @retry( stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, min=2, max=10) ) async def get(self, path: str, **kwargs): """GET with automatic retries.""" response = await self.client.get(f"{self.base_url}{path}", **kwargs) response.raise_for_status() return response.json() async def post(self, path: str, **kwargs): """POST request.""" response = await self.client.post(f"{self.base_url}{path}", **kwargs) response.raise_for_status() return response.json() # Usage payment_client = ServiceClient("http://payment-service:8001") result = await payment_client.post("/payments", json=payment_data) ``` ### Pattern 2: Asynchronous Event-Driven ```python # Event-driven communication with Kafka from aiokafka import AIOKafkaProducer, AIOKafkaConsumer import json from dataclasses import dataclass, asdict from datetime import datetime @dataclass class DomainEvent: event_id: str event_type: str aggregate_id: str occurred_at: datetime data: dict class EventBus: """Event publishing and subscription.""" def __init__(self, bootstrap_servers: List[str]): self.bootstrap_servers = bootstrap_servers self.producer = None async def start(self): self.producer = AIOKafkaProducer( bootstrap_servers=self.bootstrap_servers, value_serializer=lambda v: json.dumps(v).encode() ) await self.producer.start() async def publish(self, event: DomainEvent): """Publish event to Kafka topic.""" topic = event.event_type await self.producer.send_and_wait( topic, value=asdict(event), key=event.aggregate_id.encode() ) async def subscribe(self, topic: str, handler: callable): """Subscribe to events.""" consumer = AIOKafkaConsumer( topic, bootstrap_servers=self.bootstrap_servers, value_deserializer=lambda v: json.loads(v.decode()), group_id="my-service" ) await consumer.start() try: async for message in consumer: event_data = message.value await handler(event_data) finally: await consumer.stop() # Order Service publishes event async def create_order(order_data: dict): order = await save_order(order_data) event = DomainEvent( event_id=str(uuid.uuid4()), event_type="OrderCreated", aggregate_id=order.id, occurred_at=datetime.now(), data={ "order_id": order.id, "customer_id": order.customer_id, "total": order.total } ) await event_bus.publish(event) # Inventory Service listens for OrderCreated async def handle_order_created(event_data: dict): """React to order creation.""" order_id = event_data["data"]["order_id"] items = event_data["data"]["items"] # Reserve inventory await reserve_inventory(order_id, items) ``` ### Pattern 3: Saga Pattern (Distributed Transactions) ```python # Saga orchestration for order fulfillment from enum import Enum from typing import List, Callable class SagaStep: """Single step in saga.""" def __init__( self, name: str, action: Callable, compensation: Callable ): self.name = name self.action = action self.compensation = compensation class SagaStatus(Enum): PENDING = "pending" COMPLETED = "completed" COMPENSATING = "compensating" FAILED = "failed" class OrderFulfillmentSaga: """Orchestrated saga for order fulfillment.""" def __init__(self): self.steps: List[SagaStep] = [ SagaStep( "create_order", action=self.create_order, compensation=self.cancel_order ), SagaStep( "reserve_inventory", action=self.reserve_inventory, compensation=self.release_inventory ), SagaStep( "process_payment", action=self.process_payment, compensation=self.refund_payment ), SagaStep( "confirm_order", action=self.confirm_order, compensation=self.cancel_order_confirmation ) ] async def execute(self, order_data: dict) -> SagaResult: """Execute saga steps.""" completed_steps = [] context = {"order_data": order_data} try: for step in self.steps: # Execute step result = await step.action(context) if not result.success: # Compensate await self.compensate(completed_steps, context) return SagaResult( status=SagaStatus.FAILED, error=result.error ) completed_steps.append(step) context.update(result.data) return SagaResult(status=SagaStatus.COMPLETED, data=context) except Exception as e: # Compensate on error await self.compensate(completed_steps, context) return SagaResult(status=SagaStatus.FAILED, error=str(e)) async def compensate(self, completed_steps: List[SagaStep], context: dict): """Execute compensating actions in reverse order.""" for step in reversed(completed_steps): try: await step.compensation(context) except Exception as e: # Log compensation failure print(f"Compensation failed for {step.name}: {e}") # Step implementations async def create_order(self, context: dict) -> StepResult: order = await order_service.create(context["order_data"]) return StepResult(success=True, data={"order_id": order.id}) async def cancel_order(self, context: dict): await order_service.cancel(context["order_id"]) async def reserve_inventory(self, context: dict) -> StepResult: result = await inventory_service.reserve( context["order_id"], context["order_data"]["items"] ) return StepResult( success=result.success, data={"reservation_id": result.reservation_id} ) async def release_inventory(self, context: dict): await inventory_service.release(context["reservation_id"]) async def process_payment(self, context: dict) -> StepResult: result = await payment_service.charge( context["order_id"], context["order_data"]["total"] ) return StepResult( success=result.success, data={"transaction_id": result.transaction_id}, error=result.error ) async def refund_payment(self, context: dict): await payment_service.refund(context["transaction_id"]) ``` ## Resilience Patterns ### Circuit Breaker Pattern ```python from enum import Enum from datetime import datetime, timedelta from typing import Callable, Any class CircuitState(Enum): CLOSED = "closed" # Normal operation OPEN = "open" # Failing, reject requests HALF_OPEN = "half_open" # Testing if recovered class CircuitBreaker: """Circuit breaker for service calls.""" def __init__( self, failure_threshold: int = 5, recovery_timeout: int = 30, success_threshold: int = 2 ): self.failure_threshold = failure_threshold self.recovery_timeout = recovery_timeout self.success_threshold = success_threshold self.failure_count = 0 self.success_count = 0 self.state = CircuitState.CLOSED self.opened_at = None async def call(self, func: Callable, *args, **kwargs) -> Any: """Execute function with circuit breaker.""" if self.state == CircuitState.OPEN: if self._should_attempt_reset(): self.state = CircuitState.HALF_OPEN else: raise CircuitBreakerOpenError("Circuit breaker is open") try: result = await func(*args, **kwargs) self._on_success() return result except Exception as e: self._on_failure() raise def _on_success(self): """Handle successful call.""" self.failure_count = 0 if self.state == CircuitState.HALF_OPEN: self.success_count += 1 if self.success_count >= self.success_threshold: self.state = CircuitState.CLOSED self.success_count = 0 def _on_failure(self): """Handle failed call.""" self.failure_count += 1 if self.failure_count >= self.failure_threshold: self.state = CircuitState.OPEN self.opened_at = datetime.now() if self.state == CircuitState.HALF_OPEN: self.state = CircuitState.OPEN self.opened_at = datetime.now() def _should_attempt_reset(self) -> bool: """Check if enough time passed to try again.""" return ( datetime.now() - self.opened_at > timedelta(seconds=self.recovery_timeout) ) # Usage breaker = CircuitBreaker(failure_threshold=5, recovery_timeout=30) async def call_payment_service(payment_data: dict): return await breaker.call( payment_client.process_payment, payment_data ) ``` ## Resources - **references/service-decomposition-guide.md**: Breaking down monoliths - **references/communication-patterns.md**: Sync vs async patterns - **references/saga-implementation.md**: Distributed transactions - **assets/circuit-breaker.py**: Production circuit breaker - **assets/event-bus-template.py**: Kafka event bus implementation - **assets/api-gateway-template.py**: Complete API gateway ## Best Practices 1. **Service Boundaries**: Align with business capabilities 2. **Database Per Service**: No shared databases 3. **API Contracts**: Versioned, backward compatible 4. **Async When Possible**: Events over direct calls 5. **Circuit Breakers**: Fail fast on service failures 6. **Distributed Tracing**: Track requests across services 7. **Service Registry**: Dynamic service discovery 8. **Health Checks**: Liveness and readiness probes ## Common Pitfalls - **Distributed Monolith**: Tightly coupled services - **Chatty Services**: Too many inter-service calls - **Shared Databases**: Tight coupling through data - **No Circuit Breakers**: Cascade failures - **Synchronous Everything**: Tight coupling, poor resilience - **Premature Microservices**: Starting with microservices - **Ignoring Network Failures**: Assuming reliable network - **No Compensation Logic**: Can't undo failed transactions
šŸ‘0
šŸ‘ļø0
šŸ¤– Auto-discovered
šŸ¤–system prompt•7 months ago

saga-orchestration

Implement saga patterns for distributed transactions and

coding
⭐1
# Saga Orchestration Patterns for managing distributed transactions and long-running business processes. ## When to Use This Skill - Coordinating multi-service transactions - Implementing compensating transactions - Managing long-running business workflows - Handling failures in distributed systems - Building order fulfillment processes - Implementing approval workflows ## Core Concepts ### 1. Saga Types ``` Choreography Orchestration ā”Œā”€ā”€ā”€ā”€ā”€ā” ā”Œā”€ā”€ā”€ā”€ā”€ā” ā”Œā”€ā”€ā”€ā”€ā”€ā” ā”Œā”€ā”€ā”€ā”€ā”€ā”€ā”€ā”€ā”€ā”€ā”€ā”€ā”€ā” │Svc A│─►│Svc B│─►│Svc C│ │ Orchestrator│ ā””ā”€ā”€ā”€ā”€ā”€ā”˜ ā””ā”€ā”€ā”€ā”€ā”€ā”˜ ā””ā”€ā”€ā”€ā”€ā”€ā”˜ ā””ā”€ā”€ā”€ā”€ā”€ā”€ā”¬ā”€ā”€ā”€ā”€ā”€ā”€ā”˜ │ │ │ │ ā–¼ ā–¼ ā–¼ ā”Œā”€ā”€ā”€ā”€ā”€ā”¼ā”€ā”€ā”€ā”€ā”€ā” Event Event Event ā–¼ ā–¼ ā–¼ ā”Œā”€ā”€ā”€ā”€ā”ā”Œā”€ā”€ā”€ā”€ā”ā”Œā”€ā”€ā”€ā”€ā” │Svc1││Svc2││Svc3│ ā””ā”€ā”€ā”€ā”€ā”˜ā””ā”€ā”€ā”€ā”€ā”˜ā””ā”€ā”€ā”€ā”€ā”˜ ``` ### 2. Saga Execution States | State | Description | | ---------------- | ------------------------------ | | **Started** | Saga initiated | | **Pending** | Waiting for step completion | | **Compensating** | Rolling back due to failure | | **Completed** | All steps succeeded | | **Failed** | Saga failed after compensation | ## Templates ### Template 1: Saga Orchestrator Base ```python from abc import ABC, abstractmethod from dataclasses import dataclass, field from enum import Enum from typing import List, Dict, Any, Optional from datetime import datetime import uuid class SagaState(Enum): STARTED = "started" PENDING = "pending" COMPENSATING = "compensating" COMPLETED = "completed" FAILED = "failed" @dataclass class SagaStep: name: str action: str compensation: str status: str = "pending" result: Optional[Dict] = None error: Optional[str] = None executed_at: Optional[datetime] = None compensated_at: Optional[datetime] = None @dataclass class Saga: saga_id: str saga_type: str state: SagaState data: Dict[str, Any] steps: List[SagaStep] current_step: int = 0 created_at: datetime = field(default_factory=datetime.utcnow) updated_at: datetime = field(default_factory=datetime.utcnow) class SagaOrchestrator(ABC): """Base class for saga orchestrators.""" def __init__(self, saga_store, event_publisher): self.saga_store = saga_store self.event_publisher = event_publisher @abstractmethod def define_steps(self, data: Dict) -> List[SagaStep]: """Define the saga steps.""" pass @property @abstractmethod def saga_type(self) -> str: """Unique saga type identifier.""" pass async def start(self, data: Dict) -> Saga: """Start a new saga.""" saga = Saga( saga_id=str(uuid.uuid4()), saga_type=self.saga_type, state=SagaState.STARTED, data=data, steps=self.define_steps(data) ) await self.saga_store.save(saga) await self._execute_next_step(saga) return saga async def handle_step_completed(self, saga_id: str, step_name: str, result: Dict): """Handle successful step completion.""" saga = await self.saga_store.get(saga_id) # Update step for step in saga.steps: if step.name == step_name: step.status = "completed" step.result = result step.executed_at = datetime.utcnow() break saga.current_step += 1 saga.updated_at = datetime.utcnow() # Check if saga is complete if saga.current_step >= len(saga.steps): saga.state = SagaState.COMPLETED await self.saga_store.save(saga) await self._on_saga_completed(saga) else: saga.state = SagaState.PENDING await self.saga_store.save(saga) await self._execute_next_step(saga) async def handle_step_failed(self, saga_id: str, step_name: str, error: str): """Handle step failure - start compensation.""" saga = await self.saga_store.get(saga_id) # Mark step as failed for step in saga.steps: if step.name == step_name: step.status = "failed" step.error = error break saga.state = SagaState.COMPENSATING saga.updated_at = datetime.utcnow() await self.saga_store.save(saga) # Start compensation from current step backwards await self._compensate(saga) async def _execute_next_step(self, saga: Saga): """Execute the next step in the saga.""" if saga.current_step >= len(saga.steps): return step = saga.steps[saga.current_step] step.status = "executing" await self.saga_store.save(saga) # Publish command to execute step await self.event_publisher.publish( step.action, { "saga_id": saga.saga_id, "step_name": step.name, **saga.data } ) async def _compensate(self, saga: Saga): """Execute compensation for completed steps.""" # Compensate in reverse order for i in range(saga.current_step - 1, -1, -1): step = saga.steps[i] if step.status == "completed": step.status = "compensating" await self.saga_store.save(saga) await self.event_publisher.publish( step.compensation, { "saga_id": saga.saga_id, "step_name": step.name, "original_result": step.result, **saga.data } ) async def handle_compensation_completed(self, saga_id: str, step_name: str): """Handle compensation completion.""" saga = await self.saga_store.get(saga_id) for step in saga.steps: if step.name == step_name: step.status = "compensated" step.compensated_at = datetime.utcnow() break # Check if all compensations complete all_compensated = all( s.status in ("compensated", "pending", "failed") for s in saga.steps ) if all_compensated: saga.state = SagaState.FAILED await self._on_saga_failed(saga) await self.saga_store.save(saga) async def _on_saga_completed(self, saga: Saga): """Called when saga completes successfully.""" await self.event_publisher.publish( f"{self.saga_type}Completed", {"saga_id": saga.saga_id, **saga.data} ) async def _on_saga_failed(self, saga: Saga): """Called when saga fails after compensation.""" await self.event_publisher.publish( f"{self.saga_type}Failed", {"saga_id": saga.saga_id, "error": "Saga failed", **saga.data} ) ``` ### Template 2: Order Fulfillment Saga ```python class OrderFulfillmentSaga(SagaOrchestrator): """Orchestrates order fulfillment across services.""" @property def saga_type(self) -> str: return "OrderFulfillment" def define_steps(self, data: Dict) -> List[SagaStep]: return [ SagaStep( name="reserve_inventory", action="InventoryService.ReserveItems", compensation="InventoryService.ReleaseReservation" ), SagaStep( name="process_payment", action="PaymentService.ProcessPayment", compensation="PaymentService.RefundPayment" ), SagaStep( name="create_shipment", action="ShippingService.CreateShipment", compensation="ShippingService.CancelShipment" ), SagaStep( name="send_confirmation", action="NotificationService.SendOrderConfirmation", compensation="NotificationService.SendCancellationNotice" ) ] # Usage async def create_order(order_data: Dict): saga = OrderFulfillmentSaga(saga_store, event_publisher) return await saga.start({ "order_id": order_data["order_id"], "customer_id": order_data["customer_id"], "items": order_data["items"], "payment_method": order_data["payment_method"], "shipping_address": order_data["shipping_address"] }) # Event handlers in each service class InventoryService: async def handle_reserve_items(self, command: Dict): try: # Reserve inventory reservation = await self.reserve( command["items"], command["order_id"] ) # Report success await self.event_publisher.publish( "SagaStepCompleted", { "saga_id": command["saga_id"], "step_name": "reserve_inventory", "result": {"reservation_id": reservation.id} } ) except InsufficientInventoryError as e: await self.event_publisher.publish( "SagaStepFailed", { "saga_id": command["saga_id"], "step_name": "reserve_inventory", "error": str(e) } ) async def handle_release_reservation(self, command: Dict): # Compensating action await self.release_reservation( command["original_result"]["reservation_id"] ) await self.event_publisher.publish( "SagaCompensationCompleted", { "saga_id": command["saga_id"], "step_name": "reserve_inventory" } ) ``` ### Template 3: Choreography-Based Saga ```python from dataclasses import dataclass from typing import Dict, Any import asyncio @dataclass class SagaContext: """Passed through choreographed saga events.""" saga_id: str step: int data: Dict[str, Any] completed_steps: list class OrderChoreographySaga: """Choreography-based saga using events.""" def __init__(self, event_bus): self.event_bus = event_bus self._register_handlers() def _register_handlers(self): self.event_bus.subscribe("OrderCreated", self._on_order_created) self.event_bus.subscribe("InventoryReserved", self._on_inventory_reserved) self.event_bus.subscribe("PaymentProcessed", self._on_payment_processed) self.event_bus.subscribe("ShipmentCreated", self._on_shipment_created) # Compensation handlers self.event_bus.subscribe("PaymentFailed", self._on_payment_failed) self.event_bus.subscribe("ShipmentFailed", self._on_shipment_failed) async def _on_order_created(self, event: Dict): """Step 1: Order created, reserve inventory.""" await self.event_bus.publish("ReserveInventory", { "saga_id": event["order_id"], "order_id": event["order_id"], "items": event["items"] }) async def _on_inventory_reserved(self, event: Dict): """Step 2: Inventory reserved, process payment.""" await self.event_bus.publish("ProcessPayment", { "saga_id": event["saga_id"], "order_id": event["order_id"], "amount": event["total_amount"], "reservation_id": event["reservation_id"] }) async def _on_payment_processed(self, event: Dict): """Step 3: Payment done, create shipment.""" await self.event_bus.publish("CreateShipment", { "saga_id": event["saga_id"], "order_id": event["order_id"], "payment_id": event["payment_id"] }) async def _on_shipment_created(self, event: Dict): """Step 4: Complete - send confirmation.""" await self.event_bus.publish("OrderFulfilled", { "saga_id": event["saga_id"], "order_id": event["order_id"], "tracking_number": event["tracking_number"] }) # Compensation handlers async def _on_payment_failed(self, event: Dict): """Payment failed - release inventory.""" await self.event_bus.publish("ReleaseInventory", { "saga_id": event["saga_id"], "reservation_id": event["reservation_id"] }) await self.event_bus.publish("OrderFailed", { "order_id": event["order_id"], "reason": "Payment failed" }) async def _on_shipment_failed(self, event: Dict): """Shipment failed - refund payment and release inventory.""" await self.event_bus.publish("RefundPayment", { "saga_id": event["saga_id"], "payment_id": event["payment_id"] }) await self.event_bus.publish("ReleaseInventory", { "saga_id": event["saga_id"], "reservation_id": event["reservation_id"] }) ``` ### Template 4: Saga with Timeouts ```python class TimeoutSagaOrchestrator(SagaOrchestrator): """Saga orchestrator with step timeouts.""" def __init__(self, saga_store, event_publisher, scheduler): super().__init__(saga_store, event_publisher) self.scheduler = scheduler async def _execute_next_step(self, saga: Saga): if saga.current_step >= len(saga.steps): return step = saga.steps[saga.current_step] step.status = "executing" step.timeout_at = datetime.utcnow() + timedelta(minutes=5) await self.saga_store.save(saga) # Schedule timeout check await self.scheduler.schedule( f"saga_timeout_{saga.saga_id}_{step.name}", self._check_timeout, {"saga_id": saga.saga_id, "step_name": step.name}, run_at=step.timeout_at ) await self.event_publisher.publish( step.action, {"saga_id": saga.saga_id, "step_name": step.name, **saga.data} ) async def _check_timeout(self, data: Dict): """Check if step has timed out.""" saga = await self.saga_store.get(data["saga_id"]) step = next(s for s in saga.steps if s.name == data["step_name"]) if step.status == "executing": # Step timed out - fail it await self.handle_step_failed( data["saga_id"], data["step_name"], "Step timed out" ) ``` ## Best Practices ### Do's - **Make steps idempotent** - Safe to retry - **Design compensations carefully** - They must work - **Use correlation IDs** - For tracing across services - **Implement timeouts** - Don't wait forever - **Log everything** - For debugging failures ### Don'ts - **Don't assume instant completion** - Sagas take time - **Don't skip compensation testing** - Most critical part - **Don't couple services** - Use async messaging - **Don't ignore partial failures** - Handle gracefully ## Resources - [Saga Pattern](https://microservices.io/patterns/data/saga.html) - [Designing Data-Intensive Applications](https://dataintensive.net/)
šŸ‘0
šŸ‘ļø0
šŸ¤– Auto-discovered
šŸ¤–system prompt•7 months ago

architecture-decision-records

Write and maintain Architecture Decision Records (ADRs) following

coding
⭐1
# Architecture Decision Records Comprehensive patterns for creating, maintaining, and managing Architecture Decision Records (ADRs) that capture the context and rationale behind significant technical decisions. ## When to Use This Skill - Making significant architectural decisions - Documenting technology choices - Recording design trade-offs - Onboarding new team members - Reviewing historical decisions - Establishing decision-making processes ## Core Concepts ### 1. What is an ADR? An Architecture Decision Record captures: - **Context**: Why we needed to make a decision - **Decision**: What we decided - **Consequences**: What happens as a result ### 2. When to Write an ADR | Write ADR | Skip ADR | | -------------------------- | ---------------------- | | New framework adoption | Minor version upgrades | | Database technology choice | Bug fixes | | API design patterns | Implementation details | | Security architecture | Routine maintenance | | Integration patterns | Configuration changes | ### 3. ADR Lifecycle ``` Proposed → Accepted → Deprecated → Superseded ↓ Rejected ``` ## Templates ### Template 1: Standard ADR (MADR Format) ```markdown # ADR-0001: Use PostgreSQL as Primary Database ## Status Accepted ## Context We need to select a primary database for our new e-commerce platform. The system will handle: - ~10,000 concurrent users - Complex product catalog with hierarchical categories - Transaction processing for orders and payments - Full-text search for products - Geospatial queries for store locator The team has experience with MySQL, PostgreSQL, and MongoDB. We need ACID compliance for financial transactions. ## Decision Drivers - **Must have ACID compliance** for payment processing - **Must support complex queries** for reporting - **Should support full-text search** to reduce infrastructure complexity - **Should have good JSON support** for flexible product attributes - **Team familiarity** reduces onboarding time ## Considered Options ### Option 1: PostgreSQL - **Pros**: ACID compliant, excellent JSON support (JSONB), built-in full-text search, PostGIS for geospatial, team has experience - **Cons**: Slightly more complex replication setup than MySQL ### Option 2: MySQL - **Pros**: Very familiar to team, simple replication, large community - **Cons**: Weaker JSON support, no built-in full-text search (need Elasticsearch), no geospatial without extensions ### Option 3: MongoDB - **Pros**: Flexible schema, native JSON, horizontal scaling - **Cons**: No ACID for multi-document transactions (at decision time), team has limited experience, requires schema design discipline ## Decision We will use **PostgreSQL 15** as our primary database. ## Rationale PostgreSQL provides the best balance of: 1. **ACID compliance** essential for e-commerce transactions 2. **Built-in capabilities** (full-text search, JSONB, PostGIS) reduce infrastructure complexity 3. **Team familiarity** with SQL databases reduces learning curve 4. **Mature ecosystem** with excellent tooling and community support The slight complexity in replication is outweighed by the reduction in additional services (no separate Elasticsearch needed). ## Consequences ### Positive - Single database handles transactions, search, and geospatial queries - Reduced operational complexity (fewer services to manage) - Strong consistency guarantees for financial data - Team can leverage existing SQL expertise ### Negative - Need to learn PostgreSQL-specific features (JSONB, full-text search syntax) - Vertical scaling limits may require read replicas sooner - Some team members need PostgreSQL-specific training ### Risks - Full-text search may not scale as well as dedicated search engines - Mitigation: Design for potential Elasticsearch addition if needed ## Implementation Notes - Use JSONB for flexible product attributes - Implement connection pooling with PgBouncer - Set up streaming replication for read replicas - Use pg_trgm extension for fuzzy search ## Related Decisions - ADR-0002: Caching Strategy (Redis) - complements database choice - ADR-0005: Search Architecture - may supersede if Elasticsearch needed ## References - [PostgreSQL JSON Documentation](https://www.postgresql.org/docs/current/datatype-json.html) - [PostgreSQL Full Text Search](https://www.postgresql.org/docs/current/textsearch.html) - Internal: Performance benchmarks in `/docs/benchmarks/database-comparison.md` ``` ### Template 2: Lightweight ADR ```markdown # ADR-0012: Adopt TypeScript for Frontend Development **Status**: Accepted **Date**: 2024-01-15 **Deciders**: @alice, @bob, @charlie ## Context Our React codebase has grown to 50+ components with increasing bug reports related to prop type mismatches and undefined errors. PropTypes provide runtime-only checking. ## Decision Adopt TypeScript for all new frontend code. Migrate existing code incrementally. ## Consequences **Good**: Catch type errors at compile time, better IDE support, self-documenting code. **Bad**: Learning curve for team, initial slowdown, build complexity increase. **Mitigations**: TypeScript training sessions, allow gradual adoption with `allowJs: true`. ``` ### Template 3: Y-Statement Format ```markdown # ADR-0015: API Gateway Selection In the context of **building a microservices architecture**, facing **the need for centralized API management, authentication, and rate limiting**, we decided for **Kong Gateway** and against **AWS API Gateway and custom Nginx solution**, to achieve **vendor independence, plugin extensibility, and team familiarity with Lua**, accepting that **we need to manage Kong infrastructure ourselves**. ``` ### Template 4: ADR for Deprecation ```markdown # ADR-0020: Deprecate MongoDB in Favor of PostgreSQL ## Status Accepted (Supersedes ADR-0003) ## Context ADR-0003 (2021) chose MongoDB for user profile storage due to schema flexibility needs. Since then: - MongoDB's multi-document transactions remain problematic for our use case - Our schema has stabilized and rarely changes - We now have PostgreSQL expertise from other services - Maintaining two databases increases operational burden ## Decision Deprecate MongoDB and migrate user profiles to PostgreSQL. ## Migration Plan 1. **Phase 1** (Week 1-2): Create PostgreSQL schema, dual-write enabled 2. **Phase 2** (Week 3-4): Backfill historical data, validate consistency 3. **Phase 3** (Week 5): Switch reads to PostgreSQL, monitor 4. **Phase 4** (Week 6): Remove MongoDB writes, decommission ## Consequences ### Positive - Single database technology reduces operational complexity - ACID transactions for user data - Team can focus PostgreSQL expertise ### Negative - Migration effort (~4 weeks) - Risk of data issues during migration - Lose some schema flexibility ## Lessons Learned Document from ADR-0003 experience: - Schema flexibility benefits were overestimated - Operational cost of multiple databases was underestimated - Consider long-term maintenance in technology decisions ``` ### Template 5: Request for Comments (RFC) Style ```markdown # RFC-0025: Adopt Event Sourcing for Order Management ## Summary Propose adopting event sourcing pattern for the order management domain to improve auditability, enable temporal queries, and support business analytics. ## Motivation Current challenges: 1. Audit requirements need complete order history 2. "What was the order state at time X?" queries are impossible 3. Analytics team needs event stream for real-time dashboards 4. Order state reconstruction for customer support is manual ## Detailed Design ### Event Store ``` OrderCreated { orderId, customerId, items[], timestamp } OrderItemAdded { orderId, item, timestamp } OrderItemRemoved { orderId, itemId, timestamp } PaymentReceived { orderId, amount, paymentId, timestamp } OrderShipped { orderId, trackingNumber, timestamp } ``` ### Projections - **CurrentOrderState**: Materialized view for queries - **OrderHistory**: Complete timeline for audit - **DailyOrderMetrics**: Analytics aggregation ### Technology - Event Store: EventStoreDB (purpose-built, handles projections) - Alternative considered: Kafka + custom projection service ## Drawbacks - Learning curve for team - Increased complexity vs. CRUD - Need to design events carefully (immutable once stored) - Storage growth (events never deleted) ## Alternatives 1. **Audit tables**: Simpler but doesn't enable temporal queries 2. **CDC from existing DB**: Complex, doesn't change data model 3. **Hybrid**: Event source only for order state changes ## Unresolved Questions - [ ] Event schema versioning strategy - [ ] Retention policy for events - [ ] Snapshot frequency for performance ## Implementation Plan 1. Prototype with single order type (2 weeks) 2. Team training on event sourcing (1 week) 3. Full implementation and migration (4 weeks) 4. Monitoring and optimization (ongoing) ## References - [Event Sourcing by Martin Fowler](https://martinfowler.com/eaaDev/EventSourcing.html) - [EventStoreDB Documentation](https://www.eventstore.com/docs) ``` ## ADR Management ### Directory Structure ``` docs/ ā”œā”€ā”€ adr/ │ ā”œā”€ā”€ README.md # Index and guidelines │ ā”œā”€ā”€ template.md # Team's ADR template │ ā”œā”€ā”€ 0001-use-postgresql.md │ ā”œā”€ā”€ 0002-caching-strategy.md │ ā”œā”€ā”€ 0003-mongodb-user-profiles.md # [DEPRECATED] │ └── 0020-deprecate-mongodb.md # Supersedes 0003 ``` ### ADR Index (README.md) ```markdown # Architecture Decision Records This directory contains Architecture Decision Records (ADRs) for [Project Name]. ## Index | ADR | Title | Status | Date | | ------------------------------------- | ---------------------------------- | ---------- | ---------- | | [0001](0001-use-postgresql.md) | Use PostgreSQL as Primary Database | Accepted | 2024-01-10 | | [0002](0002-caching-strategy.md) | Caching Strategy with Redis | Accepted | 2024-01-12 | | [0003](0003-mongodb-user-profiles.md) | MongoDB for User Profiles | Deprecated | 2023-06-15 | | [0020](0020-deprecate-mongodb.md) | Deprecate MongoDB | Accepted | 2024-01-15 | ## Creating a New ADR 1. Copy `template.md` to `NNNN-title-with-dashes.md` 2. Fill in the template 3. Submit PR for review 4. Update this index after approval ## ADR Status - **Proposed**: Under discussion - **Accepted**: Decision made, implementing - **Deprecated**: No longer relevant - **Superseded**: Replaced by another ADR - **Rejected**: Considered but not adopted ``` ### Automation (adr-tools) ```bash # Install adr-tools brew install adr-tools # Initialize ADR directory adr init docs/adr # Create new ADR adr new "Use PostgreSQL as Primary Database" # Supersede an ADR adr new -s 3 "Deprecate MongoDB in Favor of PostgreSQL" # Generate table of contents adr generate toc > docs/adr/README.md # Link related ADRs adr link 2 "Complements" 1 "Is complemented by" ``` ## Review Process ```markdown ## ADR Review Checklist ### Before Submission - [ ] Context clearly explains the problem - [ ] All viable options considered - [ ] Pros/cons balanced and honest - [ ] Consequences (positive and negative) documented - [ ] Related ADRs linked ### During Review - [ ] At least 2 senior engineers reviewed - [ ] Affected teams consulted - [ ] Security implications considered - [ ] Cost implications documented - [ ] Reversibility assessed ### After Acceptance - [ ] ADR index updated - [ ] Team notified - [ ] Implementation tickets created - [ ] Related documentation updated ``` ## Best Practices ### Do's - **Write ADRs early** - Before implementation starts - **Keep them short** - 1-2 pages maximum - **Be honest about trade-offs** - Include real cons - **Link related decisions** - Build decision graph - **Update status** - Deprecate when superseded ### Don'ts - **Don't change accepted ADRs** - Write new ones to supersede - **Don't skip context** - Future readers need background - **Don't hide failures** - Rejected decisions are valuable - **Don't be vague** - Specific decisions, specific consequences - **Don't forget implementation** - ADR without action is waste ## Resources - [Documenting Architecture Decisions (Michael Nygard)](https://cognitect.com/blog/2011/11/15/documenting-architecture-decisions) - [MADR Template](https://adr.github.io/madr/) - [ADR GitHub Organization](https://adr.github.io/) - [adr-tools](https://github.com/npryce/adr-tools)
šŸ‘0
šŸ‘ļø0
šŸ¤– Auto-discovered
šŸ¤–system prompt•7 months ago

nodejs-backend-patterns

Build production-ready Node.js backend services with

coding
⭐1
# Node.js Backend Patterns Comprehensive guidance for building scalable, maintainable, and production-ready Node.js backend applications with modern frameworks, architectural patterns, and best practices. ## When to Use This Skill - Building REST APIs or GraphQL servers - Creating microservices with Node.js - Implementing authentication and authorization - Designing scalable backend architectures - Setting up middleware and error handling - Integrating databases (SQL and NoSQL) - Building real-time applications with WebSockets - Implementing background job processing ## Core Frameworks ### Express.js - Minimalist Framework **Basic Setup:** ```typescript import express, { Request, Response, NextFunction } from "express"; import helmet from "helmet"; import cors from "cors"; import compression from "compression"; const app = express(); // Security middleware app.use(helmet()); app.use(cors({ origin: process.env.ALLOWED_ORIGINS?.split(",") })); app.use(compression()); // Body parsing app.use(express.json({ limit: "10mb" })); app.use(express.urlencoded({ extended: true, limit: "10mb" })); // Request logging app.use((req: Request, res: Response, next: NextFunction) => { console.log(`${req.method} ${req.path}`); next(); }); const PORT = process.env.PORT || 3000; app.listen(PORT, () => { console.log(`Server running on port ${PORT}`); }); ``` ### Fastify - High Performance Framework **Basic Setup:** ```typescript import Fastify from "fastify"; import helmet from "@fastify/helmet"; import cors from "@fastify/cors"; import compress from "@fastify/compress"; const fastify = Fastify({ logger: { level: process.env.LOG_LEVEL || "info", transport: { target: "pino-pretty", options: { colorize: true }, }, }, }); // Plugins await fastify.register(helmet); await fastify.register(cors, { origin: true }); await fastify.register(compress); // Type-safe routes with schema validation fastify.post<{ Body: { name: string; email: string }; Reply: { id: string; name: string }; }>( "/users", { schema: { body: { type: "object", required: ["name", "email"], properties: { name: { type: "string", minLength: 1 }, email: { type: "string", format: "email" }, }, }, }, }, async (request, reply) => { const { name, email } = request.body; return { id: "123", name }; }, ); await fastify.listen({ port: 3000, host: "0.0.0.0" }); ``` ## Architectural Patterns ### Pattern 1: Layered Architecture **Structure:** ``` src/ ā”œā”€ā”€ controllers/ # Handle HTTP requests/responses ā”œā”€ā”€ services/ # Business logic ā”œā”€ā”€ repositories/ # Data access layer ā”œā”€ā”€ models/ # Data models ā”œā”€ā”€ middleware/ # Express/Fastify middleware ā”œā”€ā”€ routes/ # Route definitions ā”œā”€ā”€ utils/ # Helper functions ā”œā”€ā”€ config/ # Configuration └── types/ # TypeScript types ``` **Controller Layer:** ```typescript // controllers/user.controller.ts import { Request, Response, NextFunction } from "express"; import { UserService } from "../services/user.service"; import { CreateUserDTO, UpdateUserDTO } from "../types/user.types"; export class UserController { constructor(private userService: UserService) {} async createUser(req: Request, res: Response, next: NextFunction) { try { const userData: CreateUserDTO = req.body; const user = await this.userService.createUser(userData); res.status(201).json(user); } catch (error) { next(error); } } async getUser(req: Request, res: Response, next: NextFunction) { try { const { id } = req.params; const user = await this.userService.getUserById(id); res.json(user); } catch (error) { next(error); } } async updateUser(req: Request, res: Response, next: NextFunction) { try { const { id } = req.params; const updates: UpdateUserDTO = req.body; const user = await this.userService.updateUser(id, updates); res.json(user); } catch (error) { next(error); } } async deleteUser(req: Request, res: Response, next: NextFunction) { try { const { id } = req.params; await this.userService.deleteUser(id); res.status(204).send(); } catch (error) { next(error); } } } ``` **Service Layer:** ```typescript // services/user.service.ts import { UserRepository } from "../repositories/user.repository"; import { CreateUserDTO, UpdateUserDTO, User } from "../types/user.types"; import { NotFoundError, ValidationError } from "../utils/errors"; import bcrypt from "bcrypt"; export class UserService { constructor(private userRepository: UserRepository) {} async createUser(userData: CreateUserDTO): Promise<User> { // Validation const existingUser = await this.userRepository.findByEmail(userData.email); if (existingUser) { throw new ValidationError("Email already exists"); } // Hash password const hashedPassword = await bcrypt.hash(userData.password, 10); // Create user const user = await this.userRepository.create({ ...userData, password: hashedPassword, }); // Remove password from response const { password, ...userWithoutPassword } = user; return userWithoutPassword as User; } async getUserById(id: string): Promise<User> { const user = await this.userRepository.findById(id); if (!user) { throw new NotFoundError("User not found"); } const { password, ...userWithoutPassword } = user; return userWithoutPassword as User; } async updateUser(id: string, updates: UpdateUserDTO): Promise<User> { const user = await this.userRepository.update(id, updates); if (!user) { throw new NotFoundError("User not found"); } const { password, ...userWithoutPassword } = user; return userWithoutPassword as User; } async deleteUser(id: string): Promise<void> { const deleted = await this.userRepository.delete(id); if (!deleted) { throw new NotFoundError("User not found"); } } } ``` **Repository Layer:** ```typescript // repositories/user.repository.ts import { Pool } from "pg"; import { CreateUserDTO, UpdateUserDTO, UserEntity } from "../types/user.types"; export class UserRepository { constructor(private db: Pool) {} async create( userData: CreateUserDTO & { password: string }, ): Promise<UserEntity> { const query = ` INSERT INTO users (name, email, password) VALUES ($1, $2, $3) RETURNING id, name, email, password, created_at, updated_at `; const { rows } = await this.db.query(query, [ userData.name, userData.email, userData.password, ]); return rows[0]; } async findById(id: string): Promise<UserEntity | null> { const query = "SELECT * FROM users WHERE id = $1"; const { rows } = await this.db.query(query, [id]); return rows[0] || null; } async findByEmail(email: string): Promise<UserEntity | null> { const query = "SELECT * FROM users WHERE email = $1"; const { rows } = await this.db.query(query, [email]); return rows[0] || null; } async update(id: string, updates: UpdateUserDTO): Promise<UserEntity | null> { const fields = Object.keys(updates); const values = Object.values(updates); const setClause = fields .map((field, idx) => `${field} = $${idx + 2}`) .join(", "); const query = ` UPDATE users SET ${setClause}, updated_at = CURRENT_TIMESTAMP WHERE id = $1 RETURNING * `; const { rows } = await this.db.query(query, [id, ...values]); return rows[0] || null; } async delete(id: string): Promise<boolean> { const query = "DELETE FROM users WHERE id = $1"; const { rowCount } = await this.db.query(query, [id]); return rowCount > 0; } } ``` ### Pattern 2: Dependency Injection **DI Container:** ```typescript // di-container.ts import { Pool } from "pg"; import { UserRepository } from "./repositories/user.repository"; import { UserService } from "./services/user.service"; import { UserController } from "./controllers/user.controller"; import { AuthService } from "./services/auth.service"; class Container { private instances = new Map<string, any>(); register<T>(key: string, factory: () => T): void { this.instances.set(key, factory); } resolve<T>(key: string): T { const factory = this.instances.get(key); if (!factory) { throw new Error(`No factory registered for ${key}`); } return factory(); } singleton<T>(key: string, factory: () => T): void { let instance: T; this.instances.set(key, () => { if (!instance) { instance = factory(); } return instance; }); } } export const container = new Container(); // Register dependencies container.singleton( "db", () => new Pool({ host: process.env.DB_HOST, port: parseInt(process.env.DB_PORT || "5432"), database: process.env.DB_NAME, user: process.env.DB_USER, password: process.env.DB_PASSWORD, max: 20, idleTimeoutMillis: 30000, connectionTimeoutMillis: 2000, }), ); container.singleton( "userRepository", () => new UserRepository(container.resolve("db")), ); container.singleton( "userService", () => new UserService(container.resolve("userRepository")), ); container.register( "userController", () => new UserController(container.resolve("userService")), ); container.singleton( "authService", () => new AuthService(container.resolve("userRepository")), ); ``` ## Middleware Patterns ### Authentication Middleware ```typescript // middleware/auth.middleware.ts import { Request, Response, NextFunction } from "express"; import jwt from "jsonwebtoken"; import { UnauthorizedError } from "../utils/errors"; interface JWTPayload { userId: string; email: string; } declare global { namespace Express { interface Request { user?: JWTPayload; } } } export const authenticate = async ( req: Request, res: Response, next: NextFunction, ) => { try { const token = req.headers.authorization?.replace("Bearer ", ""); if (!token) { throw new UnauthorizedError("No token provided"); } const payload = jwt.verify(token, process.env.JWT_SECRET!) as JWTPayload; req.user = payload; next(); } catch (error) { next(new UnauthorizedError("Invalid token")); } }; export const authorize = (...roles: string[]) => { return async (req: Request, res: Response, next: NextFunction) => { if (!req.user) { return next(new UnauthorizedError("Not authenticated")); } // Check if user has required role const hasRole = roles.some((role) => req.user?.roles?.includes(role)); if (!hasRole) { return next(new UnauthorizedError("Insufficient permissions")); } next(); }; }; ``` ### Validation Middleware ```typescript // middleware/validation.middleware.ts import { Request, Response, NextFunction } from "express"; import { AnyZodObject, ZodError } from "zod"; import { ValidationError } from "../utils/errors"; export const validate = (schema: AnyZodObject) => { return async (req: Request, res: Response, next: NextFunction) => { try { await schema.parseAsync({ body: req.body, query: req.query, params: req.params, }); next(); } catch (error) { if (error instanceof ZodError) { const errors = error.errors.map((err) => ({ field: err.path.join("."), message: err.message, })); next(new ValidationError("Validation failed", errors)); } else { next(error); } } }; }; // Usage with Zod import { z } from "zod"; const createUserSchema = z.object({ body: z.object({ name: z.string().min(1), email: z.string().email(), password: z.string().min(8), }), }); router.post("/users", validate(createUserSchema), userController.createUser); ``` ### Rate Limiting Middleware ```typescript // middleware/rate-limit.middleware.ts import rateLimit from "express-rate-limit"; import RedisStore from "rate-limit-redis"; import Redis from "ioredis"; const redis = new Redis({ host: process.env.REDIS_HOST, port: parseInt(process.env.REDIS_PORT || "6379"), }); export const apiLimiter = rateLimit({ store: new RedisStore({ client: redis, prefix: "rl:", }), windowMs: 15 * 60 * 1000, // 15 minutes max: 100, // Limit each IP to 100 requests per windowMs message: "Too many requests from this IP, please try again later", standardHeaders: true, legacyHeaders: false, }); export const authLimiter = rateLimit({ store: new RedisStore({ client: redis, prefix: "rl:auth:", }), windowMs: 15 * 60 * 1000, max: 5, // Stricter limit for auth endpoints skipSuccessfulRequests: true, }); ``` ### Request Logging Middleware ```typescript // middleware/logger.middleware.ts import { Request, Response, NextFunction } from "express"; import pino from "pino"; const logger = pino({ level: process.env.LOG_LEVEL || "info", transport: { target: "pino-pretty", options: { colorize: true }, }, }); export const requestLogger = ( req: Request, res: Response, next: NextFunction, ) => { const start = Date.now(); // Log response when finished res.on("finish", () => { const duration = Date.now() - start; logger.info({ method: req.method, url: req.url, status: res.statusCode, duration: `${duration}ms`, userAgent: req.headers["user-agent"], ip: req.ip, }); }); next(); }; export { logger }; ``` ## Error Handling ### Custom Error Classes ```typescript // utils/errors.ts export class AppError extends Error { constructor( public message: string, public statusCode: number = 500, public isOperational: boolean = true, ) { super(message); Object.setPrototypeOf(this, AppError.prototype); Error.captureStackTrace(this, this.constructor); } } export class ValidationError extends AppError { constructor( message: string, public errors?: any[], ) { super(message, 400); } } export class NotFoundError extends AppError { constructor(message: string = "Resource not found") { super(message, 404); } } export class UnauthorizedError extends AppError { constructor(message: string = "Unauthorized") { super(message, 401); } } export class ForbiddenError extends AppError { constructor(message: string = "Forbidden") { super(message, 403); } } export class ConflictError extends AppError { constructor(message: string) { super(message, 409); } } ``` ### Global Error Handler ```typescript // middleware/error-handler.ts import { Request, Response, NextFunction } from "express"; import { AppError } from "../utils/errors"; import { logger } from "./logger.middleware"; export const errorHandler = ( err: Error, req: Request, res: Response, next: NextFunction, ) => { if (err instanceof AppError) { return res.status(err.statusCode).json({ status: "error", message: err.message, ...(err instanceof ValidationError && { errors: err.errors }), }); } // Log unexpected errors logger.error({ error: err.message, stack: err.stack, url: req.url, method: req.method, }); // Don't leak error details in production const message = process.env.NODE_ENV === "production" ? "Internal server error" : err.message; res.status(500).json({ status: "error", message, }); }; // Async error wrapper export const asyncHandler = ( fn: (req: Request, res: Response, next: NextFunction) => Promise<any>, ) => { return (req: Request, res: Response, next: NextFunction) => { Promise.resolve(fn(req, res, next)).catch(next); }; }; ``` ## Database Patterns ### PostgreSQL with Connection Pool ```typescript // config/database.ts import { Pool, PoolConfig } from "pg"; const poolConfig: PoolConfig = { host: process.env.DB_HOST, port: parseInt(process.env.DB_PORT || "5432"), database: process.env.DB_NAME, user: process.env.DB_USER, password: process.env.DB_PASSWORD, max: 20, idleTimeoutMillis: 30000, connectionTimeoutMillis: 2000, }; export const pool = new Pool(poolConfig); // Test connection pool.on("connect", () => { console.log("Database connected"); }); pool.on("error", (err) => { console.error("Unexpected database error", err); process.exit(-1); }); // Graceful shutdown export const closeDatabase = async () => { await pool.end(); console.log("Database connection closed"); }; ``` ### MongoDB with Mongoose ```typescript // config/mongoose.ts import mongoose from "mongoose"; const connectDB = async () => { try { await mongoose.connect(process.env.MONGODB_URI!, { maxPoolSize: 10, serverSelectionTimeoutMS: 5000, socketTimeoutMS: 45000, }); console.log("MongoDB connected"); } catch (error) { console.error("MongoDB connection error:", error); process.exit(1); } }; mongoose.connection.on("disconnected", () => { console.log("MongoDB disconnected"); }); mongoose.connection.on("error", (err) => { console.error("MongoDB error:", err); }); export { connectDB }; // Model example import { Schema, model, Document } from "mongoose"; interface IUser extends Document { name: string; email: string; password: string; createdAt: Date; updatedAt: Date; } const userSchema = new Schema<IUser>( { name: { type: String, required: true }, email: { type: String, required: true, unique: true }, password: { type: String, required: true }, }, { timestamps: true, }, ); // Indexes userSchema.index({ email: 1 }); export const User = model<IUser>("User", userSchema); ``` ### Transaction Pattern ```typescript // services/order.service.ts import { Pool } from "pg"; export class OrderService { constructor(private db: Pool) {} async createOrder(userId: string, items: any[]) { const client = await this.db.connect(); try { await client.query("BEGIN"); // Create order const orderResult = await client.query( "INSERT INTO orders (user_id, total) VALUES ($1, $2) RETURNING id", [userId, calculateTotal(items)], ); const orderId = orderResult.rows[0].id; // Create order items for (const item of items) { await client.query( "INSERT INTO order_items (order_id, product_id, quantity, price) VALUES ($1, $2, $3, $4)", [orderId, item.productId, item.quantity, item.price], ); // Update inventory await client.query( "UPDATE products SET stock = stock - $1 WHERE id = $2", [item.quantity, item.productId], ); } await client.query("COMMIT"); return orderId; } catch (error) { await client.query("ROLLBACK"); throw error; } finally { client.release(); } } } ``` ## Authentication & Authorization ### JWT Authentication ```typescript // services/auth.service.ts import jwt from
šŸ‘0
šŸ‘ļø0
šŸ¤– Auto-discovered
šŸ¤–system prompt•7 months ago

k8s-manifest-generator

Create production-ready Kubernetes manifests for Deployments,

architecture
⭐1
# Kubernetes Manifest Generator Step-by-step guidance for creating production-ready Kubernetes manifests including Deployments, Services, ConfigMaps, Secrets, and PersistentVolumeClaims. ## Purpose This skill provides comprehensive guidance for generating well-structured, secure, and production-ready Kubernetes manifests following cloud-native best practices and Kubernetes conventions. ## When to Use This Skill Use this skill when you need to: - Create new Kubernetes Deployment manifests - Define Service resources for network connectivity - Generate ConfigMap and Secret resources for configuration management - Create PersistentVolumeClaim manifests for stateful workloads - Follow Kubernetes best practices and naming conventions - Implement resource limits, health checks, and security contexts - Design manifests for multi-environment deployments ## Step-by-Step Workflow ### 1. Gather Requirements **Understand the workload:** - Application type (stateless/stateful) - Container image and version - Environment variables and configuration needs - Storage requirements - Network exposure requirements (internal/external) - Resource requirements (CPU, memory) - Scaling requirements - Health check endpoints **Questions to ask:** - What is the application name and purpose? - What container image and tag will be used? - Does the application need persistent storage? - What ports does the application expose? - Are there any secrets or configuration files needed? - What are the CPU and memory requirements? - Does the application need to be exposed externally? ### 2. Create Deployment Manifest **Follow this structure:** ```yaml apiVersion: apps/v1 kind: Deployment metadata: name: <app-name> namespace: <namespace> labels: app: <app-name> version: <version> spec: replicas: 3 selector: matchLabels: app: <app-name> template: metadata: labels: app: <app-name> version: <version> spec: containers: - name: <container-name> image: <image>:<tag> ports: - containerPort: <port> name: http resources: requests: memory: "256Mi" cpu: "250m" limits: memory: "512Mi" cpu: "500m" livenessProbe: httpGet: path: /health port: http initialDelaySeconds: 30 periodSeconds: 10 readinessProbe: httpGet: path: /ready port: http initialDelaySeconds: 5 periodSeconds: 5 env: - name: ENV_VAR value: "value" envFrom: - configMapRef: name: <app-name>-config - secretRef: name: <app-name>-secret ``` **Best practices to apply:** - Always set resource requests and limits - Implement both liveness and readiness probes - Use specific image tags (never `:latest`) - Apply security context for non-root users - Use labels for organization and selection - Set appropriate replica count based on availability needs **Reference:** See `references/deployment-spec.md` for detailed deployment options ### 3. Create Service Manifest **Choose the appropriate Service type:** **ClusterIP (internal only):** ```yaml apiVersion: v1 kind: Service metadata: name: <app-name> namespace: <namespace> labels: app: <app-name> spec: type: ClusterIP selector: app: <app-name> ports: - name: http port: 80 targetPort: 8080 protocol: TCP ``` **LoadBalancer (external access):** ```yaml apiVersion: v1 kind: Service metadata: name: <app-name> namespace: <namespace> labels: app: <app-name> annotations: service.beta.kubernetes.io/aws-load-balancer-type: nlb spec: type: LoadBalancer selector: app: <app-name> ports: - name: http port: 80 targetPort: 8080 protocol: TCP ``` **Reference:** See `references/service-spec.md` for service types and networking ### 4. Create ConfigMap **For application configuration:** ```yaml apiVersion: v1 kind: ConfigMap metadata: name: <app-name>-config namespace: <namespace> data: APP_MODE: production LOG_LEVEL: info DATABASE_HOST: db.example.com # For config files app.properties: | server.port=8080 server.host=0.0.0.0 logging.level=INFO ``` **Best practices:** - Use ConfigMaps for non-sensitive data only - Organize related configuration together - Use meaningful names for keys - Consider using one ConfigMap per component - Version ConfigMaps when making changes **Reference:** See `assets/configmap-template.yaml` for examples ### 5. Create Secret **For sensitive data:** ```yaml apiVersion: v1 kind: Secret metadata: name: <app-name>-secret namespace: <namespace> type: Opaque stringData: DATABASE_PASSWORD: "changeme" API_KEY: "secret-api-key" # For certificate files tls.crt: | -----BEGIN CERTIFICATE----- ... -----END CERTIFICATE----- tls.key: | -----BEGIN PRIVATE KEY----- ... -----END PRIVATE KEY----- ``` **Security considerations:** - Never commit secrets to Git in plain text - Use Sealed Secrets, External Secrets Operator, or Vault - Rotate secrets regularly - Use RBAC to limit secret access - Consider using Secret type: `kubernetes.io/tls` for TLS secrets ### 6. Create PersistentVolumeClaim (if needed) **For stateful applications:** ```yaml apiVersion: v1 kind: PersistentVolumeClaim metadata: name: <app-name>-data namespace: <namespace> spec: accessModes: - ReadWriteOnce storageClassName: gp3 resources: requests: storage: 10Gi ``` **Mount in Deployment:** ```yaml spec: template: spec: containers: - name: app volumeMounts: - name: data mountPath: /var/lib/app volumes: - name: data persistentVolumeClaim: claimName: <app-name>-data ``` **Storage considerations:** - Choose appropriate StorageClass for performance needs - Use ReadWriteOnce for single-pod access - Use ReadWriteMany for multi-pod shared storage - Consider backup strategies - Set appropriate retention policies ### 7. Apply Security Best Practices **Add security context to Deployment:** ```yaml spec: template: spec: securityContext: runAsNonRoot: true runAsUser: 1000 fsGroup: 1000 seccompProfile: type: RuntimeDefault containers: - name: app securityContext: allowPrivilegeEscalation: false readOnlyRootFilesystem: true capabilities: drop: - ALL ``` **Security checklist:** - [ ] Run as non-root user - [ ] Drop all capabilities - [ ] Use read-only root filesystem - [ ] Disable privilege escalation - [ ] Set seccomp profile - [ ] Use Pod Security Standards ### 8. Add Labels and Annotations **Standard labels (recommended):** ```yaml metadata: labels: app.kubernetes.io/name: <app-name> app.kubernetes.io/instance: <instance-name> app.kubernetes.io/version: "1.0.0" app.kubernetes.io/component: backend app.kubernetes.io/part-of: <system-name> app.kubernetes.io/managed-by: kubectl ``` **Useful annotations:** ```yaml metadata: annotations: description: "Application description" contact: "team@example.com" prometheus.io/scrape: "true" prometheus.io/port: "9090" prometheus.io/path: "/metrics" ``` ### 9. Organize Multi-Resource Manifests **File organization options:** **Option 1: Single file with `---` separator** ```yaml # app-name.yaml --- apiVersion: v1 kind: ConfigMap ... --- apiVersion: v1 kind: Secret ... --- apiVersion: apps/v1 kind: Deployment ... --- apiVersion: v1 kind: Service ... ``` **Option 2: Separate files** ``` manifests/ ā”œā”€ā”€ configmap.yaml ā”œā”€ā”€ secret.yaml ā”œā”€ā”€ deployment.yaml ā”œā”€ā”€ service.yaml └── pvc.yaml ``` **Option 3: Kustomize structure** ``` base/ ā”œā”€ā”€ kustomization.yaml ā”œā”€ā”€ deployment.yaml ā”œā”€ā”€ service.yaml └── configmap.yaml overlays/ ā”œā”€ā”€ dev/ │ └── kustomization.yaml └── prod/ └── kustomization.yaml ``` ### 10. Validate and Test **Validation steps:** ```bash # Dry-run validation kubectl apply -f manifest.yaml --dry-run=client # Server-side validation kubectl apply -f manifest.yaml --dry-run=server # Validate with kubeval kubeval manifest.yaml # Validate with kube-score kube-score score manifest.yaml # Check with kube-linter kube-linter lint manifest.yaml ``` **Testing checklist:** - [ ] Manifest passes dry-run validation - [ ] All required fields are present - [ ] Resource limits are reasonable - [ ] Health checks are configured - [ ] Security context is set - [ ] Labels follow conventions - [ ] Namespace exists or is created ## Common Patterns ### Pattern 1: Simple Stateless Web Application **Use case:** Standard web API or microservice **Components needed:** - Deployment (3 replicas for HA) - ClusterIP Service - ConfigMap for configuration - Secret for API keys - HorizontalPodAutoscaler (optional) **Reference:** See `assets/deployment-template.yaml` ### Pattern 2: Stateful Database Application **Use case:** Database or persistent storage application **Components needed:** - StatefulSet (not Deployment) - Headless Service - PersistentVolumeClaim template - ConfigMap for DB configuration - Secret for credentials ### Pattern 3: Background Job or Cron **Use case:** Scheduled tasks or batch processing **Components needed:** - CronJob or Job - ConfigMap for job parameters - Secret for credentials - ServiceAccount with RBAC ### Pattern 4: Multi-Container Pod **Use case:** Application with sidecar containers **Components needed:** - Deployment with multiple containers - Shared volumes between containers - Init containers for setup - Service (if needed) ## Templates The following templates are available in the `assets/` directory: - `deployment-template.yaml` - Standard deployment with best practices - `service-template.yaml` - Service configurations (ClusterIP, LoadBalancer, NodePort) - `configmap-template.yaml` - ConfigMap examples with different data types - `secret-template.yaml` - Secret examples (to be generated, not committed) - `pvc-template.yaml` - PersistentVolumeClaim templates ## Reference Documentation - `references/deployment-spec.md` - Detailed Deployment specification - `references/service-spec.md` - Service types and networking details ## Best Practices Summary 1. **Always set resource requests and limits** - Prevents resource starvation 2. **Implement health checks** - Ensures Kubernetes can manage your application 3. **Use specific image tags** - Avoid unpredictable deployments 4. **Apply security contexts** - Run as non-root, drop capabilities 5. **Use ConfigMaps and Secrets** - Separate config from code 6. **Label everything** - Enables filtering and organization 7. **Follow naming conventions** - Use standard Kubernetes labels 8. **Validate before applying** - Use dry-run and validation tools 9. **Version your manifests** - Keep in Git with version control 10. **Document with annotations** - Add context for other developers ## Troubleshooting **Pods not starting:** - Check image pull errors: `kubectl describe pod <pod-name>` - Verify resource availability: `kubectl get nodes` - Check events: `kubectl get events --sort-by='.lastTimestamp'` **Service not accessible:** - Verify selector matches pod labels: `kubectl get endpoints <service-name>` - Check service type and port configuration - Test from within cluster: `kubectl run debug --rm -it --image=busybox -- sh` **ConfigMap/Secret not loading:** - Verify names match in Deployment - Check namespace - Ensure resources exist: `kubectl get configmap,secret` ## Next Steps After creating manifests: 1. Store in Git repository 2. Set up CI/CD pipeline for deployment 3. Consider using Helm or Kustomize for templating 4. Implement GitOps with ArgoCD or Flux 5. Add monitoring and observability ## Related Skills - `helm-chart-scaffolding` - For templating and packaging - `gitops-workflow` - For automated deployments - `k8s-security-policies` - For advanced security configurations
šŸ‘0
šŸ‘ļø0
šŸ¤– Auto-discovered
šŸ¤–system prompt•7 months ago

distributed-tracing

Implement distributed tracing with Jaeger and Tempo to track

coding
⭐1
# Distributed Tracing Implement distributed tracing with Jaeger and Tempo for request flow visibility across microservices. ## Purpose Track requests across distributed systems to understand latency, dependencies, and failure points. ## When to Use - Debug latency issues - Understand service dependencies - Identify bottlenecks - Trace error propagation - Analyze request paths ## Distributed Tracing Concepts ### Trace Structure ``` Trace (Request ID: abc123) ↓ Span (frontend) [100ms] ↓ Span (api-gateway) [80ms] ā”œā†’ Span (auth-service) [10ms] └→ Span (user-service) [60ms] └→ Span (database) [40ms] ``` ### Key Components - **Trace** - End-to-end request journey - **Span** - Single operation within a trace - **Context** - Metadata propagated between services - **Tags** - Key-value pairs for filtering - **Logs** - Timestamped events within a span ## Jaeger Setup ### Kubernetes Deployment ```bash # Deploy Jaeger Operator kubectl create namespace observability kubectl create -f https://github.com/jaegertracing/jaeger-operator/releases/download/v1.51.0/jaeger-operator.yaml -n observability # Deploy Jaeger instance kubectl apply -f - <<EOF apiVersion: jaegertracing.io/v1 kind: Jaeger metadata: name: jaeger namespace: observability spec: strategy: production storage: type: elasticsearch options: es: server-urls: http://elasticsearch:9200 ingress: enabled: true EOF ``` ### Docker Compose ```yaml version: "3.8" services: jaeger: image: jaegertracing/all-in-one:latest ports: - "5775:5775/udp" - "6831:6831/udp" - "6832:6832/udp" - "5778:5778" - "16686:16686" # UI - "14268:14268" # Collector - "14250:14250" # gRPC - "9411:9411" # Zipkin environment: - COLLECTOR_ZIPKIN_HOST_PORT=:9411 ``` **Reference:** See `references/jaeger-setup.md` ## Application Instrumentation ### OpenTelemetry (Recommended) #### Python (Flask) ```python from opentelemetry import trace from opentelemetry.exporter.jaeger.thrift import JaegerExporter from opentelemetry.sdk.resources import SERVICE_NAME, Resource from opentelemetry.sdk.trace import TracerProvider from opentelemetry.sdk.trace.export import BatchSpanProcessor from opentelemetry.instrumentation.flask import FlaskInstrumentor from flask import Flask # Initialize tracer resource = Resource(attributes={SERVICE_NAME: "my-service"}) provider = TracerProvider(resource=resource) processor = BatchSpanProcessor(JaegerExporter( agent_host_name="jaeger", agent_port=6831, )) provider.add_span_processor(processor) trace.set_tracer_provider(provider) # Instrument Flask app = Flask(__name__) FlaskInstrumentor().instrument_app(app) @app.route('/api/users') def get_users(): tracer = trace.get_tracer(__name__) with tracer.start_as_current_span("get_users") as span: span.set_attribute("user.count", 100) # Business logic users = fetch_users_from_db() return {"users": users} def fetch_users_from_db(): tracer = trace.get_tracer(__name__) with tracer.start_as_current_span("database_query") as span: span.set_attribute("db.system", "postgresql") span.set_attribute("db.statement", "SELECT * FROM users") # Database query return query_database() ``` #### Node.js (Express) ```javascript const { NodeTracerProvider } = require("@opentelemetry/sdk-trace-node"); const { JaegerExporter } = require("@opentelemetry/exporter-jaeger"); const { BatchSpanProcessor } = require("@opentelemetry/sdk-trace-base"); const { registerInstrumentations } = require("@opentelemetry/instrumentation"); const { HttpInstrumentation } = require("@opentelemetry/instrumentation-http"); const { ExpressInstrumentation, } = require("@opentelemetry/instrumentation-express"); // Initialize tracer const provider = new NodeTracerProvider({ resource: { attributes: { "service.name": "my-service" } }, }); const exporter = new JaegerExporter({ endpoint: "http://jaeger:14268/api/traces", }); provider.addSpanProcessor(new BatchSpanProcessor(exporter)); provider.register(); // Instrument libraries registerInstrumentations({ instrumentations: [new HttpInstrumentation(), new ExpressInstrumentation()], }); const express = require("express"); const app = express(); app.get("/api/users", async (req, res) => { const tracer = trace.getTracer("my-service"); const span = tracer.startSpan("get_users"); try { const users = await fetchUsers(); span.setAttributes({ "user.count": users.length }); res.json({ users }); } finally { span.end(); } }); ``` #### Go ```go package main import ( "context" "go.opentelemetry.io/otel" "go.opentelemetry.io/otel/exporters/jaeger" "go.opentelemetry.io/otel/sdk/resource" sdktrace "go.opentelemetry.io/otel/sdk/trace" semconv "go.opentelemetry.io/otel/semconv/v1.4.0" ) func initTracer() (*sdktrace.TracerProvider, error) { exporter, err := jaeger.New(jaeger.WithCollectorEndpoint( jaeger.WithEndpoint("http://jaeger:14268/api/traces"), )) if err != nil { return nil, err } tp := sdktrace.NewTracerProvider( sdktrace.WithBatcher(exporter), sdktrace.WithResource(resource.NewWithAttributes( semconv.SchemaURL, semconv.ServiceNameKey.String("my-service"), )), ) otel.SetTracerProvider(tp) return tp, nil } func getUsers(ctx context.Context) ([]User, error) { tracer := otel.Tracer("my-service") ctx, span := tracer.Start(ctx, "get_users") defer span.End() span.SetAttributes(attribute.String("user.filter", "active")) users, err := fetchUsersFromDB(ctx) if err != nil { span.RecordError(err) return nil, err } span.SetAttributes(attribute.Int("user.count", len(users))) return users, nil } ``` **Reference:** See `references/instrumentation.md` ## Context Propagation ### HTTP Headers ``` traceparent: 00-0af7651916cd43dd8448eb211c80319c-b7ad6b7169203331-01 tracestate: congo=t61rcWkgMzE ``` ### Propagation in HTTP Requests #### Python ```python from opentelemetry.propagate import inject headers = {} inject(headers) # Injects trace context response = requests.get('http://downstream-service/api', headers=headers) ``` #### Node.js ```javascript const { propagation } = require("@opentelemetry/api"); const headers = {}; propagation.inject(context.active(), headers); axios.get("http://downstream-service/api", { headers }); ``` ## Tempo Setup (Grafana) ### Kubernetes Deployment ```yaml apiVersion: v1 kind: ConfigMap metadata: name: tempo-config data: tempo.yaml: | server: http_listen_port: 3200 distributor: receivers: jaeger: protocols: thrift_http: grpc: otlp: protocols: http: grpc: storage: trace: backend: s3 s3: bucket: tempo-traces endpoint: s3.amazonaws.com querier: frontend_worker: frontend_address: tempo-query-frontend:9095 --- apiVersion: apps/v1 kind: Deployment metadata: name: tempo spec: replicas: 1 template: spec: containers: - name: tempo image: grafana/tempo:latest args: - -config.file=/etc/tempo/tempo.yaml volumeMounts: - name: config mountPath: /etc/tempo volumes: - name: config configMap: name: tempo-config ``` **Reference:** See `assets/jaeger-config.yaml.template` ## Sampling Strategies ### Probabilistic Sampling ```yaml # Sample 1% of traces sampler: type: probabilistic param: 0.01 ``` ### Rate Limiting Sampling ```yaml # Sample max 100 traces per second sampler: type: ratelimiting param: 100 ``` ### Adaptive Sampling ```python from opentelemetry.sdk.trace.sampling import ParentBased, TraceIdRatioBased # Sample based on trace ID (deterministic) sampler = ParentBased(root=TraceIdRatioBased(0.01)) ``` ## Trace Analysis ### Finding Slow Requests **Jaeger Query:** ``` service=my-service duration > 1s ``` ### Finding Errors **Jaeger Query:** ``` service=my-service error=true tags.http.status_code >= 500 ``` ### Service Dependency Graph Jaeger automatically generates service dependency graphs showing: - Service relationships - Request rates - Error rates - Average latencies ## Best Practices 1. **Sample appropriately** (1-10% in production) 2. **Add meaningful tags** (user_id, request_id) 3. **Propagate context** across all service boundaries 4. **Log exceptions** in spans 5. **Use consistent naming** for operations 6. **Monitor tracing overhead** (<1% CPU impact) 7. **Set up alerts** for trace errors 8. **Implement distributed context** (baggage) 9. **Use span events** for important milestones 10. **Document instrumentation** standards ## Integration with Logging ### Correlated Logs ```python import logging from opentelemetry import trace logger = logging.getLogger(__name__) def process_request(): span = trace.get_current_span() trace_id = span.get_span_context().trace_id logger.info( "Processing request", extra={"trace_id": format(trace_id, '032x')} ) ``` ## Troubleshooting **No traces appearing:** - Check collector endpoint - Verify network connectivity - Check sampling configuration - Review application logs **High latency overhead:** - Reduce sampling rate - Use batch span processor - Check exporter configuration ## Reference Files - `references/jaeger-setup.md` - Jaeger installation - `references/instrumentation.md` - Instrumentation patterns - `assets/jaeger-config.yaml.template` - Jaeger configuration ## Related Skills - `prometheus-configuration` - For metrics - `grafana-dashboards` - For visualization - `slo-implementation` - For latency SLOs
šŸ‘0
šŸ‘ļø0
šŸ¤– Auto-discovered
šŸ¤–system prompt•7 months ago

async-python-patterns

Master Python asyncio, concurrent programming, and async/await

coding
⭐1
# Async Python Patterns Comprehensive guidance for implementing asynchronous Python applications using asyncio, concurrent programming patterns, and async/await for building high-performance, non-blocking systems. ## When to Use This Skill - Building async web APIs (FastAPI, aiohttp, Sanic) - Implementing concurrent I/O operations (database, file, network) - Creating web scrapers with concurrent requests - Developing real-time applications (WebSocket servers, chat systems) - Processing multiple independent tasks simultaneously - Building microservices with async communication - Optimizing I/O-bound workloads - Implementing async background tasks and queues ## Sync vs Async Decision Guide Before adopting async, consider whether it's the right choice for your use case. | Use Case | Recommended Approach | |----------|---------------------| | Many concurrent network/DB calls | `asyncio` | | CPU-bound computation | `multiprocessing` or thread pool | | Mixed I/O + CPU | Offload CPU work with `asyncio.to_thread()` | | Simple scripts, few connections | Sync (simpler, easier to debug) | | Web APIs with high concurrency | Async frameworks (FastAPI, aiohttp) | **Key Rule:** Stay fully sync or fully async within a call path. Mixing creates hidden blocking and complexity. ## Core Concepts ### 1. Event Loop The event loop is the heart of asyncio, managing and scheduling asynchronous tasks. **Key characteristics:** - Single-threaded cooperative multitasking - Schedules coroutines for execution - Handles I/O operations without blocking - Manages callbacks and futures ### 2. Coroutines Functions defined with `async def` that can be paused and resumed. **Syntax:** ```python async def my_coroutine(): result = await some_async_operation() return result ``` ### 3. Tasks Scheduled coroutines that run concurrently on the event loop. ### 4. Futures Low-level objects representing eventual results of async operations. ### 5. Async Context Managers Resources that support `async with` for proper cleanup. ### 6. Async Iterators Objects that support `async for` for iterating over async data sources. ## Quick Start ```python import asyncio async def main(): print("Hello") await asyncio.sleep(1) print("World") # Python 3.7+ asyncio.run(main()) ``` ## Fundamental Patterns ### Pattern 1: Basic Async/Await ```python import asyncio async def fetch_data(url: str) -> dict: """Fetch data from URL asynchronously.""" await asyncio.sleep(1) # Simulate I/O return {"url": url, "data": "result"} async def main(): result = await fetch_data("https://api.example.com") print(result) asyncio.run(main()) ``` ### Pattern 2: Concurrent Execution with gather() ```python import asyncio from typing import List async def fetch_user(user_id: int) -> dict: """Fetch user data.""" await asyncio.sleep(0.5) return {"id": user_id, "name": f"User {user_id}"} async def fetch_all_users(user_ids: List[int]) -> List[dict]: """Fetch multiple users concurrently.""" tasks = [fetch_user(uid) for uid in user_ids] results = await asyncio.gather(*tasks) return results async def main(): user_ids = [1, 2, 3, 4, 5] users = await fetch_all_users(user_ids) print(f"Fetched {len(users)} users") asyncio.run(main()) ``` ### Pattern 3: Task Creation and Management ```python import asyncio async def background_task(name: str, delay: int): """Long-running background task.""" print(f"{name} started") await asyncio.sleep(delay) print(f"{name} completed") return f"Result from {name}" async def main(): # Create tasks task1 = asyncio.create_task(background_task("Task 1", 2)) task2 = asyncio.create_task(background_task("Task 2", 1)) # Do other work print("Main: doing other work") await asyncio.sleep(0.5) # Wait for tasks result1 = await task1 result2 = await task2 print(f"Results: {result1}, {result2}") asyncio.run(main()) ``` ### Pattern 4: Error Handling in Async Code ```python import asyncio from typing import List, Optional async def risky_operation(item_id: int) -> dict: """Operation that might fail.""" await asyncio.sleep(0.1) if item_id % 3 == 0: raise ValueError(f"Item {item_id} failed") return {"id": item_id, "status": "success"} async def safe_operation(item_id: int) -> Optional[dict]: """Wrapper with error handling.""" try: return await risky_operation(item_id) except ValueError as e: print(f"Error: {e}") return None async def process_items(item_ids: List[int]): """Process multiple items with error handling.""" tasks = [safe_operation(iid) for iid in item_ids] results = await asyncio.gather(*tasks, return_exceptions=True) # Filter out failures successful = [r for r in results if r is not None and not isinstance(r, Exception)] failed = [r for r in results if isinstance(r, Exception)] print(f"Success: {len(successful)}, Failed: {len(failed)}") return successful asyncio.run(process_items([1, 2, 3, 4, 5, 6])) ``` ### Pattern 5: Timeout Handling ```python import asyncio async def slow_operation(delay: int) -> str: """Operation that takes time.""" await asyncio.sleep(delay) return f"Completed after {delay}s" async def with_timeout(): """Execute operation with timeout.""" try: result = await asyncio.wait_for(slow_operation(5), timeout=2.0) print(result) except asyncio.TimeoutError: print("Operation timed out") asyncio.run(with_timeout()) ``` ## Advanced Patterns ### Pattern 6: Async Context Managers ```python import asyncio from typing import Optional class AsyncDatabaseConnection: """Async database connection context manager.""" def __init__(self, dsn: str): self.dsn = dsn self.connection: Optional[object] = None async def __aenter__(self): print("Opening connection") await asyncio.sleep(0.1) # Simulate connection self.connection = {"dsn": self.dsn, "connected": True} return self.connection async def __aexit__(self, exc_type, exc_val, exc_tb): print("Closing connection") await asyncio.sleep(0.1) # Simulate cleanup self.connection = None async def query_database(): """Use async context manager.""" async with AsyncDatabaseConnection("postgresql://localhost") as conn: print(f"Using connection: {conn}") await asyncio.sleep(0.2) # Simulate query return {"rows": 10} asyncio.run(query_database()) ``` ### Pattern 7: Async Iterators and Generators ```python import asyncio from typing import AsyncIterator async def async_range(start: int, end: int, delay: float = 0.1) -> AsyncIterator[int]: """Async generator that yields numbers with delay.""" for i in range(start, end): await asyncio.sleep(delay) yield i async def fetch_pages(url: str, max_pages: int) -> AsyncIterator[dict]: """Fetch paginated data asynchronously.""" for page in range(1, max_pages + 1): await asyncio.sleep(0.2) # Simulate API call yield { "page": page, "url": f"{url}?page={page}", "data": [f"item_{page}_{i}" for i in range(5)] } async def consume_async_iterator(): """Consume async iterator.""" async for number in async_range(1, 5): print(f"Number: {number}") print("\nFetching pages:") async for page_data in fetch_pages("https://api.example.com/items", 3): print(f"Page {page_data['page']}: {len(page_data['data'])} items") asyncio.run(consume_async_iterator()) ``` ### Pattern 8: Producer-Consumer Pattern ```python import asyncio from asyncio import Queue from typing import Optional async def producer(queue: Queue, producer_id: int, num_items: int): """Produce items and put them in queue.""" for i in range(num_items): item = f"Item-{producer_id}-{i}" await queue.put(item) print(f"Producer {producer_id} produced: {item}") await asyncio.sleep(0.1) await queue.put(None) # Signal completion async def consumer(queue: Queue, consumer_id: int): """Consume items from queue.""" while True: item = await queue.get() if item is None: queue.task_done() break print(f"Consumer {consumer_id} processing: {item}") await asyncio.sleep(0.2) # Simulate work queue.task_done() async def producer_consumer_example(): """Run producer-consumer pattern.""" queue = Queue(maxsize=10) # Create tasks producers = [ asyncio.create_task(producer(queue, i, 5)) for i in range(2) ] consumers = [ asyncio.create_task(consumer(queue, i)) for i in range(3) ] # Wait for producers await asyncio.gather(*producers) # Wait for queue to be empty await queue.join() # Cancel consumers for c in consumers: c.cancel() asyncio.run(producer_consumer_example()) ``` ### Pattern 9: Semaphore for Rate Limiting ```python import asyncio from typing import List async def api_call(url: str, semaphore: asyncio.Semaphore) -> dict: """Make API call with rate limiting.""" async with semaphore: print(f"Calling {url}") await asyncio.sleep(0.5) # Simulate API call return {"url": url, "status": 200} async def rate_limited_requests(urls: List[str], max_concurrent: int = 5): """Make multiple requests with rate limiting.""" semaphore = asyncio.Semaphore(max_concurrent) tasks = [api_call(url, semaphore) for url in urls] results = await asyncio.gather(*tasks) return results async def main(): urls = [f"https://api.example.com/item/{i}" for i in range(20)] results = await rate_limited_requests(urls, max_concurrent=3) print(f"Completed {len(results)} requests") asyncio.run(main()) ``` ### Pattern 10: Async Locks and Synchronization ```python import asyncio class AsyncCounter: """Thread-safe async counter.""" def __init__(self): self.value = 0 self.lock = asyncio.Lock() async def increment(self): """Safely increment counter.""" async with self.lock: current = self.value await asyncio.sleep(0.01) # Simulate work self.value = current + 1 async def get_value(self) -> int: """Get current value.""" async with self.lock: return self.value async def worker(counter: AsyncCounter, worker_id: int): """Worker that increments counter.""" for _ in range(10): await counter.increment() print(f"Worker {worker_id} incremented") async def test_counter(): """Test concurrent counter.""" counter = AsyncCounter() workers = [asyncio.create_task(worker(counter, i)) for i in range(5)] await asyncio.gather(*workers) final_value = await counter.get_value() print(f"Final counter value: {final_value}") asyncio.run(test_counter()) ``` ## Real-World Applications ### Web Scraping with aiohttp ```python import asyncio import aiohttp from typing import List, Dict async def fetch_url(session: aiohttp.ClientSession, url: str) -> Dict: """Fetch single URL.""" try: async with session.get(url, timeout=aiohttp.ClientTimeout(total=10)) as response: text = await response.text() return { "url": url, "status": response.status, "length": len(text) } except Exception as e: return {"url": url, "error": str(e)} async def scrape_urls(urls: List[str]) -> List[Dict]: """Scrape multiple URLs concurrently.""" async with aiohttp.ClientSession() as session: tasks = [fetch_url(session, url) for url in urls] results = await asyncio.gather(*tasks) return results async def main(): urls = [ "https://httpbin.org/delay/1", "https://httpbin.org/delay/2", "https://httpbin.org/status/404", ] results = await scrape_urls(urls) for result in results: print(result) asyncio.run(main()) ``` ### Async Database Operations ```python import asyncio from typing import List, Optional # Simulated async database client class AsyncDB: """Simulated async database.""" async def execute(self, query: str) -> List[dict]: """Execute query.""" await asyncio.sleep(0.1) return [{"id": 1, "name": "Example"}] async def fetch_one(self, query: str) -> Optional[dict]: """Fetch single row.""" await asyncio.sleep(0.1) return {"id": 1, "name": "Example"} async def get_user_data(db: AsyncDB, user_id: int) -> dict: """Fetch user and related data concurrently.""" user_task = db.fetch_one(f"SELECT * FROM users WHERE id = {user_id}") orders_task = db.execute(f"SELECT * FROM orders WHERE user_id = {user_id}") profile_task = db.fetch_one(f"SELECT * FROM profiles WHERE user_id = {user_id}") user, orders, profile = await asyncio.gather(user_task, orders_task, profile_task) return { "user": user, "orders": orders, "profile": profile } async def main(): db = AsyncDB() user_data = await get_user_data(db, 1) print(user_data) asyncio.run(main()) ``` ### WebSocket Server ```python import asyncio from typing import Set # Simulated WebSocket connection class WebSocket: """Simulated WebSocket.""" def __init__(self, client_id: str): self.client_id = client_id async def send(self, message: str): """Send message.""" print(f"Sending to {self.client_id}: {message}") await asyncio.sleep(0.01) async def recv(self) -> str: """Receive message.""" await asyncio.sleep(1) return f"Message from {self.client_id}" class WebSocketServer: """Simple WebSocket server.""" def __init__(self): self.clients: Set[WebSocket] = set() async def register(self, websocket: WebSocket): """Register new client.""" self.clients.add(websocket) print(f"Client {websocket.client_id} connected") async def unregister(self, websocket: WebSocket): """Unregister client.""" self.clients.remove(websocket) print(f"Client {websocket.client_id} disconnected") async def broadcast(self, message: str): """Broadcast message to all clients.""" if self.clients: tasks = [client.send(message) for client in self.clients] await asyncio.gather(*tasks) async def handle_client(self, websocket: WebSocket): """Handle individual client connection.""" await self.register(websocket) try: async for message in self.message_iterator(websocket): await self.broadcast(f"{websocket.client_id}: {message}") finally: await self.unregister(websocket) async def message_iterator(self, websocket: WebSocket): """Iterate over messages from client.""" for _ in range(3): # Simulate 3 messages yield await websocket.recv() ``` ## Performance Best Practices ### 1. Use Connection Pools ```python import asyncio import aiohttp async def with_connection_pool(): """Use connection pool for efficiency.""" connector = aiohttp.TCPConnector(limit=100, limit_per_host=10) async with aiohttp.ClientSession(connector=connector) as session: tasks = [session.get(f"https://api.example.com/item/{i}") for i in range(50)] responses = await asyncio.gather(*tasks) return responses ``` ### 2. Batch Operations ```python async def batch_process(items: List[str], batch_size: int = 10): """Process items in batches.""" for i in range(0, len(items), batch_size): batch = items[i:i + batch_size] tasks = [process_item(item) for item in batch] await asyncio.gather(*tasks) print(f"Processed batch {i // batch_size + 1}") async def process_item(item: str): """Process single item.""" await asyncio.sleep(0.1) return f"Processed: {item}" ``` ### 3. Avoid Blocking Operations Never block the event loop with synchronous operations. A single blocking call stalls all concurrent tasks. ```python # BAD - blocks the entire event loop async def fetch_data_bad(): import time import requests time.sleep(1) # Blocks! response = requests.get(url) # Also blocks! # GOOD - use async-native libraries (e.g., httpx for async HTTP) import httpx async def fetch_data_good(url: str): await asyncio.sleep(1) async with httpx.AsyncClient() as client: response = await client.get(url) ``` **Wrapping Blocking Code with `asyncio.to_thread()` (Python 3.9+):** When you must use synchronous libraries, offload to a thread pool: ```python import asyncio from pathlib import Path async def read_file_async(path: str) -> str: """Read file without blocking event loop.""" # asyncio.to_thread() runs sync code in a thread pool return await asyncio.to_thread(Path(path).read_text) async def call_sync_library(data: dict) -> dict: """Wrap a synchronous library call.""" # Useful for sync database drivers, file I/O, CPU work return await asyncio.to_thread(sync_library.process, data) ``` **Lower-level approach with `run_in_executor()`:** ```python import asyncio import concurrent.futures from typing import Any def blocking_operation(data: Any) -> Any: """CPU-intensive blocking operation.""" import time time.sleep(1) return data * 2 async def run_in_executor(data: Any) -> Any: """Run blocking operation in thread pool.""" loop = asyncio.get_running_loop() with concurrent.futures.ThreadPoolExecutor() as pool: result = await loop.run_in_executor(pool, blocking_operation, data) return result async def main(): results = await asyncio.gather(*[run_in_executor(i) for i in range(5)]) print(results) asyncio.run(main()) ``` ## Common Pitfalls ### 1. Forgetting await ```python # Wrong - returns coroutine object, doesn't execute result = async_function() # Correct result = await async_function() ``` ### 2. Blocking the Event Loop ```python # Wrong - blocks event loop import time async def bad(): time.sleep(1) # Blocks! # Correct async def good(): await asyncio.sleep(1) # Non-blocking ``` ### 3. Not Handling Cancellation ```python async def cancelable_task(): """Task that handles cancellation.""" try: while True: await asyncio.sleep(1) print("Working...") except asyncio.CancelledError: print("Task cancelled, cleaning up...") # Perform cleanup raise # Re-raise to propagate cancellation ``` ### 4. Mixing Sync and Async Code ```python # Wrong - can't call async from sync directly def sync_function(): result = await async_function() # SyntaxError! # Correct def sync_function(): result = asyncio.run(async_function()) ``` ## Testing Async Code ```python import asyncio import pytest # Using pytest-asyncio @pytest.mark.asyncio async def test_async_function(): """Test async function.""" result = await fetch_data("https
šŸ‘0
šŸ‘ļø0
šŸ¤– Auto-discovered
šŸ¤–system prompt•7 months ago

python-resilience

Python resilience patterns including automatic retries, exponential

coding
⭐1
# Python Resilience Patterns Build fault-tolerant Python applications that gracefully handle transient failures, network issues, and service outages. Resilience patterns keep systems running when dependencies are unreliable. ## When to Use This Skill - Adding retry logic to external service calls - Implementing timeouts for network operations - Building fault-tolerant microservices - Handling rate limiting and backpressure - Creating infrastructure decorators - Designing circuit breakers ## Core Concepts ### 1. Transient vs Permanent Failures Retry transient errors (network timeouts, temporary service issues). Don't retry permanent errors (invalid credentials, bad requests). ### 2. Exponential Backoff Increase wait time between retries to avoid overwhelming recovering services. ### 3. Jitter Add randomness to backoff to prevent thundering herd when many clients retry simultaneously. ### 4. Bounded Retries Cap both attempt count and total duration to prevent infinite retry loops. ## Quick Start ```python from tenacity import retry, stop_after_attempt, wait_exponential_jitter @retry( stop=stop_after_attempt(3), wait=wait_exponential_jitter(initial=1, max=10), ) def call_external_service(request: dict) -> dict: return httpx.post("https://api.example.com", json=request).json() ``` ## Fundamental Patterns ### Pattern 1: Basic Retry with Tenacity Use the `tenacity` library for production-grade retry logic. For simpler cases, consider built-in retry functionality or a lightweight custom implementation. ```python from tenacity import ( retry, stop_after_attempt, stop_after_delay, wait_exponential_jitter, retry_if_exception_type, ) TRANSIENT_ERRORS = (ConnectionError, TimeoutError, OSError) @retry( retry=retry_if_exception_type(TRANSIENT_ERRORS), stop=stop_after_attempt(5) | stop_after_delay(60), wait=wait_exponential_jitter(initial=1, max=30), ) def fetch_data(url: str) -> dict: """Fetch data with automatic retry on transient failures.""" response = httpx.get(url, timeout=30) response.raise_for_status() return response.json() ``` ### Pattern 2: Retry Only Appropriate Errors Whitelist specific transient exceptions. Never retry: - `ValueError`, `TypeError` - These are bugs, not transient issues - `AuthenticationError` - Invalid credentials won't become valid - HTTP 4xx errors (except 429) - Client errors are permanent ```python from tenacity import retry, retry_if_exception_type import httpx # Define what's retryable RETRYABLE_EXCEPTIONS = ( ConnectionError, TimeoutError, httpx.ConnectTimeout, httpx.ReadTimeout, ) @retry( retry=retry_if_exception_type(RETRYABLE_EXCEPTIONS), stop=stop_after_attempt(3), wait=wait_exponential_jitter(initial=1, max=10), ) def resilient_api_call(endpoint: str) -> dict: """Make API call with retry on network issues.""" return httpx.get(endpoint, timeout=10).json() ``` ### Pattern 3: HTTP Status Code Retries Retry specific HTTP status codes that indicate transient issues. ```python from tenacity import retry, retry_if_result, stop_after_attempt import httpx RETRY_STATUS_CODES = {429, 502, 503, 504} def should_retry_response(response: httpx.Response) -> bool: """Check if response indicates a retryable error.""" return response.status_code in RETRY_STATUS_CODES @retry( retry=retry_if_result(should_retry_response), stop=stop_after_attempt(3), wait=wait_exponential_jitter(initial=1, max=10), ) def http_request(method: str, url: str, **kwargs) -> httpx.Response: """Make HTTP request with retry on transient status codes.""" return httpx.request(method, url, timeout=30, **kwargs) ``` ### Pattern 4: Combined Exception and Status Retry Handle both network exceptions and HTTP status codes. ```python from tenacity import ( retry, retry_if_exception_type, retry_if_result, stop_after_attempt, wait_exponential_jitter, before_sleep_log, ) import logging import httpx logger = logging.getLogger(__name__) TRANSIENT_EXCEPTIONS = ( ConnectionError, TimeoutError, httpx.ConnectError, httpx.ReadTimeout, ) RETRY_STATUS_CODES = {429, 500, 502, 503, 504} def is_retryable_response(response: httpx.Response) -> bool: return response.status_code in RETRY_STATUS_CODES @retry( retry=( retry_if_exception_type(TRANSIENT_EXCEPTIONS) | retry_if_result(is_retryable_response) ), stop=stop_after_attempt(5), wait=wait_exponential_jitter(initial=1, max=30), before_sleep=before_sleep_log(logger, logging.WARNING), ) def robust_http_call( method: str, url: str, **kwargs, ) -> httpx.Response: """HTTP call with comprehensive retry handling.""" return httpx.request(method, url, timeout=30, **kwargs) ``` ## Advanced Patterns ### Pattern 5: Logging Retry Attempts Track retry behavior for debugging and alerting. ```python from tenacity import retry, stop_after_attempt, wait_exponential import structlog logger = structlog.get_logger() def log_retry_attempt(retry_state): """Log detailed retry information.""" exception = retry_state.outcome.exception() logger.warning( "Retrying operation", attempt=retry_state.attempt_number, exception_type=type(exception).__name__, exception_message=str(exception), next_wait_seconds=retry_state.next_action.sleep if retry_state.next_action else None, ) @retry( stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, max=10), before_sleep=log_retry_attempt, ) def call_with_logging(request: dict) -> dict: """External call with retry logging.""" ... ``` ### Pattern 6: Timeout Decorator Create reusable timeout decorators for consistent timeout handling. ```python import asyncio from functools import wraps from typing import TypeVar, Callable T = TypeVar("T") def with_timeout(seconds: float): """Decorator to add timeout to async functions.""" def decorator(func: Callable[..., T]) -> Callable[..., T]: @wraps(func) async def wrapper(*args, **kwargs) -> T: return await asyncio.wait_for( func(*args, **kwargs), timeout=seconds, ) return wrapper return decorator @with_timeout(30) async def fetch_with_timeout(url: str) -> dict: """Fetch URL with 30 second timeout.""" async with httpx.AsyncClient() as client: response = await client.get(url) return response.json() ``` ### Pattern 7: Cross-Cutting Concerns via Decorators Stack decorators to separate infrastructure from business logic. ```python from functools import wraps from typing import TypeVar, Callable import structlog logger = structlog.get_logger() T = TypeVar("T") def traced(name: str | None = None): """Add tracing to function calls.""" def decorator(func: Callable[..., T]) -> Callable[..., T]: span_name = name or func.__name__ @wraps(func) async def wrapper(*args, **kwargs) -> T: logger.info("Operation started", operation=span_name) try: result = await func(*args, **kwargs) logger.info("Operation completed", operation=span_name) return result except Exception as e: logger.error("Operation failed", operation=span_name, error=str(e)) raise return wrapper return decorator # Stack multiple concerns @traced("fetch_user_data") @with_timeout(30) @retry(stop=stop_after_attempt(3), wait=wait_exponential_jitter()) async def fetch_user_data(user_id: str) -> dict: """Fetch user with tracing, timeout, and retry.""" ... ``` ### Pattern 8: Dependency Injection for Testability Pass infrastructure components through constructors for easy testing. ```python from dataclasses import dataclass from typing import Protocol class Logger(Protocol): def info(self, msg: str, **kwargs) -> None: ... def error(self, msg: str, **kwargs) -> None: ... class MetricsClient(Protocol): def increment(self, metric: str, tags: dict | None = None) -> None: ... def timing(self, metric: str, value: float) -> None: ... @dataclass class UserService: """Service with injected infrastructure.""" repository: UserRepository logger: Logger metrics: MetricsClient async def get_user(self, user_id: str) -> User: self.logger.info("Fetching user", user_id=user_id) start = time.perf_counter() try: user = await self.repository.get(user_id) self.metrics.increment("user.fetch.success") return user except Exception as e: self.metrics.increment("user.fetch.error") self.logger.error("Failed to fetch user", user_id=user_id, error=str(e)) raise finally: elapsed = time.perf_counter() - start self.metrics.timing("user.fetch.duration", elapsed) # Easy to test with fakes service = UserService( repository=FakeRepository(), logger=FakeLogger(), metrics=FakeMetrics(), ) ``` ### Pattern 9: Fail-Safe Defaults Degrade gracefully when non-critical operations fail. ```python from typing import TypeVar from collections.abc import Callable T = TypeVar("T") def fail_safe(default: T, log_failure: bool = True): """Return default value on failure instead of raising.""" def decorator(func: Callable[..., T]) -> Callable[..., T]: @wraps(func) async def wrapper(*args, **kwargs) -> T: try: return await func(*args, **kwargs) except Exception as e: if log_failure: logger.warning( "Operation failed, using default", function=func.__name__, error=str(e), ) return default return wrapper return decorator @fail_safe(default=[]) async def get_recommendations(user_id: str) -> list[str]: """Get recommendations, return empty list on failure.""" ... ``` ## Best Practices Summary 1. **Retry only transient errors** - Don't retry bugs or authentication failures 2. **Use exponential backoff** - Give services time to recover 3. **Add jitter** - Prevent thundering herd from synchronized retries 4. **Cap total duration** - `stop_after_attempt(5) | stop_after_delay(60)` 5. **Log every retry** - Silent retries hide systemic problems 6. **Use decorators** - Keep retry logic separate from business logic 7. **Inject dependencies** - Make infrastructure testable 8. **Set timeouts everywhere** - Every network call needs a timeout 9. **Fail gracefully** - Return cached/default values for non-critical paths 10. **Monitor retry rates** - High retry rates indicate underlying issues
šŸ‘0
šŸ‘ļø0
šŸ¤– Auto-discovered