Skip to main content
EVOKORE// BROWSE
>

./browse/prompts

73 NODES
πŸ“textβ€’6 months ago

PR Merge Boundary Validation Runbook

Validate chained pull requests at each merge boundary with explicit checks, stale-approval handling, and rollback awareness.

devops
⭐1
# PR Merge Boundary Validation Runbook Imported from curated first-party documentation sources. ## What this covers Use this runbook when a merge sequence spans dependent pull requests and each boundary needs fresh validation. ## Use this when - Managing dependency-chained pull requests - Refreshing stale approvals near merge time - Reducing merge surprises through checkpoint validation ## Expected outcomes - Each merge boundary has a concrete verification step - Review freshness stays visible as branches evolve - Rollback planning is captured before the merge happens ## Source synthesis - EVOKORE-MCP/docs/PR_MERGE_RUNBOOK.md (https://github.com/mattmre/EVOKORE-MCP/blob/main/docs/PR_MERGE_RUNBOOK.md) ## Dedupe notes Keeps the merge-boundary validation concept separate from more general workflow or orchestration docs. ## Source excerpts ### EVOKORE-MCP/docs/PR_MERGE_RUNBOOK.md Operator runbook for reliable merges and context-rot prevention. ## Pre-merge Checklist - [ ] PR scope matches approved plan - [ ] PR description is filled using `.github/PULL_REQUEST_TEMPLATE.md` - [ ] PR metadata automation check (`scripts/validate-pr-metadata.js`) is passing for pull_request CI runs - [ ] Required tests pass locally/CI - [ ] Docs updated for user-facing behavior changes - [ ] Release-impacting changes called out - [ ] Follow-up issues captured (if any) ## Required Checks by Change Type Use this as the minimum check set before approval and merge: | Change type | Required checks | | --- | --- | | Docs-only changes | `node test-docs-canonical-links.js` | | Ops/docs process changes (`docs/PR_MERGE_RUNBOOK.md`, `next-session.md`, orchestration docs) | `node test-ops-docs-validation.js` and `node test-docs-canonical-links.js` | | Source/tooling/config changes (`src/`, `scripts/`, workflow/config files) | Relevant targeted tests for touched area plus CI-required suite | | Release-flow changes | `node test-npm-release-flow-validation.js` plus docs/link checks | If a PR spans multiple change types, run the union of required checks. ## Reviewer Responsibilities - Confirm PR scope and dependency assumptions are explicit in description. - Confirm PR metadata fields from `.github/PULL_REQUEST_TEMPLATE.md` are complete. - Verify required checks for each change type are attached in PR evidence. - Block approval if dependency base PR is not merged or branch is stale. - Approve only the current chain head; do not pre-approve non-head dependent PRs. - Confirm all review conversations are resolved before final approval. - Ensure merge strategy and rollback notes are documented for risky changes. ## Merge Steps 1. Rebase or update branch with latest target branch. 2. Re-run required validations. 3. Confirm reviewer approvals and resolved conversations. 4. Merge via approved strategy. 5. Record merge commit/PR number in tracker. ## Merge-boundary Checkpoints At every dependency merge boundary (`base -> dependent`): 1. Rebase dependent PR branch on latest `main` immediately after parent merge. 2. Re-run required checks and attach updated evidence in PR metadata. 3. Revalidate approvals for the new head state (stale approvals must be refreshed). 4. Confirm merge-boundary revalidation notes are updated before merge. ## Merge-order Controls (Dependency Chain) 1. Define merge order explicitly in PR descriptions (`base -> dependent`). 2. Merge only the current chain head; hold dependents until parent merge is complete. 3. After each merge, rebase ...
πŸ‘0
πŸ‘οΈ0
docs
πŸ€–system promptβ€’6 months ago

Cross-CLI MCP Config Sync

Keep Claude, Cursor, Gemini, and related CLI integrations aligned with a repeatable dry-run and apply workflow.

productivity
⭐1
# Cross-CLI MCP Config Sync Imported from curated first-party documentation sources. ## What this covers Use this skill when multiple AI clients need the same MCP configuration without drifting out of sync. ## Use this when - Rolling out a shared MCP config to multiple clients - Previewing config changes before applying them - Standardizing developer setup across tools ## Expected outcomes - Cross-client MCP setup becomes easier to repeat - Dry-run and apply modes reduce accidental changes - Environment-specific details stay documented near the workflow ## Source synthesis - EVOKORE-MCP/docs/CLI_INTEGRATION.md (https://github.com/mattmre/EVOKORE-MCP/blob/main/docs/CLI_INTEGRATION.md) ## Dedupe notes Uses the dedicated CLI integration guide as the canonical source for config sync instead of duplicating setup notes elsewhere. ## Source excerpts ### EVOKORE-MCP/docs/CLI_INTEGRATION.md EVOKORE-MCP isn't just an MCP Server-Ò€—it also ships with natively integrated UI hooks designed to make your AI CLI experience (like Gemini CLI or Claude Code) significantly more powerful and transparent. ## 🍨 The Interactive Status Line When you connect EVOKORE-MCP to your AI Assistant, you can optionally enable the **EVOKORE Status Line**. Every time the AI finishes a thought or a tool execution, this hook intercepts the internal JSON payload and renders a beautiful, color-coded ASCII status bar at the bottom of your terminal showing: - **Location**: Your current working directory. - **Model Identity**: The exact LLM model currently loaded. - **Skill Count**: A live count of the MCP Agent Skills currently indexed in your library. - **Context Window Health**: A dynamic, color-coded progress bar showing exactly how many tokens you have consumed. --- ### 💜 Enabling in Gemini CLI Gemini CLI features a robust native hook engine. You can configure it to execute the EVOKORE Status Line immediately after every model response (`AfterModel`). **Step 1:** Locate your global settings file (`~/.gemini/settings.json`). **Step 2:** Ensure hooks are enabled, and add the `AfterModel` event array to the root of the JSON object: ```json { "enableHooks": true, "hooks": { "AfterModel": [ { "type": "command", "command": "node /absolute/path/to/EVOKORE-MCP/scripts/status.js" } ] } } ``` **Step 3:** Restart your Gemini CLI! --- ### 💜 Enabling in Claude Code Claude Code features an undocumented internal hook architecture that natively supports this status line. *(Note: Because this feature is currently undocumented by Anthropic, Claude Code's `doctor` command will display "Found 1 settings issue". This is perfectly normal and the status line will still execute successfully)*. **Step 1:** Locate your Claude settings file (`~/.claude/settings.json`). **Step 2:** Add the `statusLine` block to the root of the JSON object: ```json { "statusLine": { "type": "command", "command": "node /absolute/path/to/EVOKORE-MCP/scripts/status.js" } } ``` **Step 3:** Restart Claude Code. --- ### Òő ï¸ A Note on GitHub Copilot and Codex Microsoft's GitHub Copilot CLI and OpenAI's Codex CLI **do not natively support** these JSON hook configurations. If you want the EVOKORE Status Line to appear after commands in these tools, you must configure a native PowerShell/Bash alias wrapper around the CLI execution. **Example (PowerShell Profile):** ```powershell function copilot-evokore { gh copilot $args node "/absolute/path/to/EV ...
πŸ‘0
πŸ‘οΈ0
docs
πŸ€–system promptβ€’6 months ago

Hook Observability and Session Replay

Instrument hooks with replayable logs, task-state visibility, and non-blocking observability around agent sessions.

devops
⭐1
# Hook Observability and Session Replay Imported from curated first-party documentation sources. ## What this covers Use this skill when you need to understand what hooks fired, what they emitted, and how a session can be replayed after the fact. ## Use this when - Debugging hook-driven automation - Replaying session events after failures - Adding observability without blocking interactive work ## Expected outcomes - Hook activity becomes inspectable and replayable - Operational state survives beyond a single terminal session - Observability stays useful without overwhelming operators ## Source synthesis - EVOKORE-MCP/docs/VOICE_AND_HOOKS.md (https://github.com/mattmre/EVOKORE-MCP/blob/main/docs/VOICE_AND_HOOKS.md) - EVOKORE-MCP/docs/USE_CASES_AND_WALKTHROUGHS.md (https://github.com/mattmre/EVOKORE-MCP/blob/main/docs/USE_CASES_AND_WALKTHROUGHS.md) ## Dedupe notes Focuses on replay and observability instead of importing the broader voice-sidecar guide verbatim. ## Source excerpts ### EVOKORE-MCP/docs/VOICE_AND_HOOKS.md EVOKORE currently has three separate voice-related systems plus a set of hook and observability utilities. They overlap in operator workflows, but they are not the same runtime. ## The three voice-related systems ### 1. ElevenLabs MCP proxy This is the optional `elevenlabs` child server configured in `mcp.config.json`. What it is: - proxied through the EVOKORE router - exposed as prefixed MCP tools - available to any EVOKORE-connected MCP client when configured successfully What it is for: - text-to-speech and other ElevenLabs MCP operations as tools - routing voice-related actions through the standard EVOKORE proxy/security stack Requirements: - `uvx` available on PATH - `ELEVENLABS_API_KEY` set ### 2. VoiceMode VoiceMode is a separate voice-conversation system for Claude Code. What it is: - registered separately from EVOKORE - not routed through EVOKOREÒ€ℒs stdio server - used for bidirectional voice conversation in Claude Code What it is for: - speaking to Claude and hearing spoken responses - using `OPENAI_API_KEY` and VoiceModeÒ€ℒs own runtime Windows note: - VoiceMode relies on `uvx` being directly available - set `OPENAI_API_KEY` in the shell that launches Claude Code ### 3. VoiceSidecar VoiceSidecar is a standalone WebSocket server implemented in `src/Voice ... ### EVOKORE-MCP/docs/USE_CASES_AND_WALKTHROUGHS.md This guide turns the runtime contracts into practical operator flows. ## Walkthrough 1: Adopt a workflow from the skill library Use this when you want EVOKORE to retrieve process guidance before the model starts acting. ### Goal Find and adopt an existing workflow such as `session-wrap`. ### Steps 1. Ask the client to search skills: ```text Search the MCP for a workflow about session wrap-up and continuity. ``` 2. EVOKORE uses `search_skills` and returns matching skills. 3. Ask for a specific skill: ```text Show me help for the session-wrap skill. ``` 4. EVOKORE uses `get_skill_help` and returns the skillÒ€ℒs internal instructions. 5. For broader task matching, ask: ```text Resolve a workflow for wrapping this session, documenting open risks, and preparing the next handoff. ``` 6. EVOKORE uses `resolve_workflow` and injects the top 1-3 relevant workflows directly into the tool response. ### Why this matters - keeps the model grounded in repo-specific process - reduces prompt drift - makes handoff and governance behavior repeatable ## Walkthrough 2: Use a proxied tool that requires HITL approval Use this when the tool is configured as `require_approval` in `permissions.yml`. ### Goal Allow a protected proxied tool call such as `fs_write_f ...
πŸ‘0
πŸ‘οΈ0
docs
πŸ€–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

cqrs-implementation

Implement Command Query Responsibility Segregation for scalable

coding
⭐1
# CQRS Implementation Comprehensive guide to implementing CQRS (Command Query Responsibility Segregation) patterns. ## When to Use This Skill - Separating read and write concerns - Scaling reads independently from writes - Building event-sourced systems - Optimizing complex query scenarios - Different read/write data models needed - High-performance reporting requirements ## Core Concepts ### 1. CQRS Architecture ``` β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ Client β”‚ β””β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”˜ β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ β”‚ β–Ό β–Ό β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ Commands β”‚ β”‚ Queries β”‚ β”‚ API β”‚ β”‚ API β”‚ β””β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”˜ β”‚ β”‚ β–Ό β–Ό β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ Command β”‚ β”‚ Query β”‚ β”‚ Handlers β”‚ β”‚ Handlers β”‚ β””β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”˜ β”‚ β”‚ β–Ό β–Ό β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ Write │─────────►│ Read β”‚ β”‚ Model β”‚ Events β”‚ Model β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ ``` ### 2. Key Components | Component | Responsibility | | ------------------- | ------------------------------- | | **Command** | Intent to change state | | **Command Handler** | Validates and executes commands | | **Event** | Record of state change | | **Query** | Request for data | | **Query Handler** | Retrieves data from read model | | **Projector** | Updates read model from events | ## Templates ### Template 1: Command Infrastructure ```python from abc import ABC, abstractmethod from dataclasses import dataclass from typing import TypeVar, Generic, Dict, Any, Type from datetime import datetime import uuid # Command base @dataclass class Command: command_id: str = None timestamp: datetime = None def __post_init__(self): self.command_id = self.command_id or str(uuid.uuid4()) self.timestamp = self.timestamp or datetime.utcnow() # Concrete commands @dataclass class CreateOrder(Command): customer_id: str items: list shipping_address: dict @dataclass class AddOrderItem(Command): order_id: str product_id: str quantity: int price: float @dataclass class CancelOrder(Command): order_id: str reason: str # Command handler base T = TypeVar('T', bound=Command) class CommandHandler(ABC, Generic[T]): @abstractmethod async def handle(self, command: T) -> Any: pass # Command bus class CommandBus: def __init__(self): self._handlers: Dict[Type[Command], CommandHandler] = {} def register(self, command_type: Type[Command], handler: CommandHandler): self._handlers[command_type] = handler async def dispatch(self, command: Command) -> Any: handler = self._handlers.get(type(command)) if not handler: raise ValueError(f"No handler for {type(command).__name__}") return await handler.handle(command) # Command handler implementation class CreateOrderHandler(CommandHandler[CreateOrder]): def __init__(self, order_repository, event_store): self.order_repository = order_repository self.event_store = event_store async def handle(self, command: CreateOrder) -> str: # Validate if not command.items: raise ValueError("Order must have at least one item") # Create aggregate order = Order.create( customer_id=command.customer_id, items=command.items, shipping_address=command.shipping_address ) # Persist events await self.event_store.append_events( stream_id=f"Order-{order.id}", stream_type="Order", events=order.uncommitted_events ) return order.id ``` ### Template 2: Query Infrastructure ```python from abc import ABC, abstractmethod from dataclasses import dataclass from typing import TypeVar, Generic, List, Optional # Query base @dataclass class Query: pass # Concrete queries @dataclass class GetOrderById(Query): order_id: str @dataclass class GetCustomerOrders(Query): customer_id: str status: Optional[str] = None page: int = 1 page_size: int = 20 @dataclass class SearchOrders(Query): query: str filters: dict = None sort_by: str = "created_at" sort_order: str = "desc" # Query result types @dataclass class OrderView: order_id: str customer_id: str status: str total_amount: float item_count: int created_at: datetime shipped_at: Optional[datetime] = None @dataclass class PaginatedResult(Generic[T]): items: List[T] total: int page: int page_size: int @property def total_pages(self) -> int: return (self.total + self.page_size - 1) // self.page_size # Query handler base T = TypeVar('T', bound=Query) R = TypeVar('R') class QueryHandler(ABC, Generic[T, R]): @abstractmethod async def handle(self, query: T) -> R: pass # Query bus class QueryBus: def __init__(self): self._handlers: Dict[Type[Query], QueryHandler] = {} def register(self, query_type: Type[Query], handler: QueryHandler): self._handlers[query_type] = handler async def dispatch(self, query: Query) -> Any: handler = self._handlers.get(type(query)) if not handler: raise ValueError(f"No handler for {type(query).__name__}") return await handler.handle(query) # Query handler implementation class GetOrderByIdHandler(QueryHandler[GetOrderById, Optional[OrderView]]): def __init__(self, read_db): self.read_db = read_db async def handle(self, query: GetOrderById) -> Optional[OrderView]: async with self.read_db.acquire() as conn: row = await conn.fetchrow( """ SELECT order_id, customer_id, status, total_amount, item_count, created_at, shipped_at FROM order_views WHERE order_id = $1 """, query.order_id ) if row: return OrderView(**dict(row)) return None class GetCustomerOrdersHandler(QueryHandler[GetCustomerOrders, PaginatedResult[OrderView]]): def __init__(self, read_db): self.read_db = read_db async def handle(self, query: GetCustomerOrders) -> PaginatedResult[OrderView]: async with self.read_db.acquire() as conn: # Build query with optional status filter where_clause = "customer_id = $1" params = [query.customer_id] if query.status: where_clause += " AND status = $2" params.append(query.status) # Get total count total = await conn.fetchval( f"SELECT COUNT(*) FROM order_views WHERE {where_clause}", *params ) # Get paginated results offset = (query.page - 1) * query.page_size rows = await conn.fetch( f""" SELECT order_id, customer_id, status, total_amount, item_count, created_at, shipped_at FROM order_views WHERE {where_clause} ORDER BY created_at DESC LIMIT ${len(params) + 1} OFFSET ${len(params) + 2} """, *params, query.page_size, offset ) return PaginatedResult( items=[OrderView(**dict(row)) for row in rows], total=total, page=query.page, page_size=query.page_size ) ``` ### Template 3: FastAPI CQRS Application ```python from fastapi import FastAPI, HTTPException, Depends from pydantic import BaseModel from typing import List, Optional app = FastAPI() # Request/Response models class CreateOrderRequest(BaseModel): customer_id: str items: List[dict] shipping_address: dict class OrderResponse(BaseModel): order_id: str customer_id: str status: str total_amount: float item_count: int created_at: datetime # Dependency injection def get_command_bus() -> CommandBus: return app.state.command_bus def get_query_bus() -> QueryBus: return app.state.query_bus # Command endpoints (POST, PUT, DELETE) @app.post("/orders", response_model=dict) async def create_order( request: CreateOrderRequest, command_bus: CommandBus = Depends(get_command_bus) ): command = CreateOrder( customer_id=request.customer_id, items=request.items, shipping_address=request.shipping_address ) order_id = await command_bus.dispatch(command) return {"order_id": order_id} @app.post("/orders/{order_id}/items") async def add_item( order_id: str, product_id: str, quantity: int, price: float, command_bus: CommandBus = Depends(get_command_bus) ): command = AddOrderItem( order_id=order_id, product_id=product_id, quantity=quantity, price=price ) await command_bus.dispatch(command) return {"status": "item_added"} @app.delete("/orders/{order_id}") async def cancel_order( order_id: str, reason: str, command_bus: CommandBus = Depends(get_command_bus) ): command = CancelOrder(order_id=order_id, reason=reason) await command_bus.dispatch(command) return {"status": "cancelled"} # Query endpoints (GET) @app.get("/orders/{order_id}", response_model=OrderResponse) async def get_order( order_id: str, query_bus: QueryBus = Depends(get_query_bus) ): query = GetOrderById(order_id=order_id) result = await query_bus.dispatch(query) if not result: raise HTTPException(status_code=404, detail="Order not found") return result @app.get("/customers/{customer_id}/orders") async def get_customer_orders( customer_id: str, status: Optional[str] = None, page: int = 1, page_size: int = 20, query_bus: QueryBus = Depends(get_query_bus) ): query = GetCustomerOrders( customer_id=customer_id, status=status, page=page, page_size=page_size ) return await query_bus.dispatch(query) @app.get("/orders/search") async def search_orders( q: str, sort_by: str = "created_at", query_bus: QueryBus = Depends(get_query_bus) ): query = SearchOrders(query=q, sort_by=sort_by) return await query_bus.dispatch(query) ``` ### Template 4: Read Model Synchronization ```python class ReadModelSynchronizer: """Keeps read models in sync with events.""" def __init__(self, event_store, read_db, projections: List[Projection]): self.event_store = event_store self.read_db = read_db self.projections = {p.name: p for p in projections} async def run(self): """Continuously sync read models.""" while True: for name, projection in self.projections.items(): await self._sync_projection(projection) await asyncio.sleep(0.1) async def _sync_projection(self, projection: Projection): checkpoint = await self._get_checkpoint(projection.name) events = await self.event_store.read_all( from_position=checkpoint, limit=100 ) for event in events: if event.event_type in projection.handles(): try: await projection.apply(event) except Exception as e: # Log error, possibly retry or skip logger.error(f"Projection error: {e}") continue await self._save_checkpoint(projection.name, event.global_position) async def rebuild_projection(self, projection_name: str): """Rebuild a projection from scratch.""" projection = self.projections[projection_name] # Clear existing data await projection.clear() # Reset checkpoint await self._save_checkpoint(projection_name, 0) # Rebuild while True: checkpoint = await self._get_checkpoint(projection_name) events = await self.event_store.read_all(checkpoint, 1000) if not events: break for event in events: if event.event_type in projection.handles(): await projection.apply(event) await self._save_checkpoint( projection_name, events[-1].global_position ) ``` ### Template 5: Eventual Consistency Handling ```python class ConsistentQueryHandler: """Query handler that can wait for consistency.""" def __init__(self, read_db, event_store): self.read_db = read_db self.event_store = event_store async def query_after_command( self, query: Query, expected_version: int, stream_id: str, timeout: float = 5.0 ): """ Execute query, ensuring read model is at expected version. Used for read-your-writes consistency. """ start_time = time.time() while time.time() - start_time < timeout: # Check if read model is caught up projection_version = await self._get_projection_version(stream_id) if projection_version >= expected_version: return await self.execute_query(query) # Wait a bit and retry await asyncio.sleep(0.1) # Timeout - return stale data with warning return { "data": await self.execute_query(query), "_warning": "Data may be stale" } async def _get_projection_version(self, stream_id: str) -> int: """Get the last processed event version for a stream.""" async with self.read_db.acquire() as conn: return await conn.fetchval( "SELECT last_event_version FROM projection_state WHERE stream_id = $1", stream_id ) or 0 ``` ## Best Practices ### Do's - **Separate command and query models** - Different needs - **Use eventual consistency** - Accept propagation delay - **Validate in command handlers** - Before state change - **Denormalize read models** - Optimize for queries - **Version your events** - For schema evolution ### Don'ts - **Don't query in commands** - Use only for writes - **Don't couple read/write schemas** - Independent evolution - **Don't over-engineer** - Start simple - **Don't ignore consistency SLAs** - Define acceptable lag ## Resources - [CQRS Pattern](https://martinfowler.com/bliki/CQRS.html) - [Microsoft CQRS Guidance](https://docs.microsoft.com/en-us/azure/architecture/patterns/cqrs)
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

event-store-design

Design and implement event stores for event-sourced systems. Use

coding
⭐1
# Event Store Design Comprehensive guide to designing event stores for event-sourced applications. ## When to Use This Skill - Designing event sourcing infrastructure - Choosing between event store technologies - Implementing custom event stores - Optimizing event storage and retrieval - Setting up event store schemas - Planning for event store scaling ## Core Concepts ### 1. Event Store Architecture ``` β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ Event Store β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ β”‚ β”‚ Stream 1 β”‚ β”‚ Stream 2 β”‚ β”‚ Stream 3 β”‚ β”‚ β”‚ β”‚ (Aggregate) β”‚ β”‚ (Aggregate) β”‚ β”‚ (Aggregate) β”‚ β”‚ β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ β”‚ β”‚ Event 1 β”‚ β”‚ Event 1 β”‚ β”‚ Event 1 β”‚ β”‚ β”‚ β”‚ Event 2 β”‚ β”‚ Event 2 β”‚ β”‚ Event 2 β”‚ β”‚ β”‚ β”‚ Event 3 β”‚ β”‚ ... β”‚ β”‚ Event 3 β”‚ β”‚ β”‚ β”‚ ... β”‚ β”‚ β”‚ β”‚ Event 4 β”‚ β”‚ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ Global Position: 1 β†’ 2 β†’ 3 β†’ 4 β†’ 5 β†’ 6 β†’ ... β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ ``` ### 2. Event Store Requirements | Requirement | Description | | ----------------- | ---------------------------------- | | **Append-only** | Events are immutable, only appends | | **Ordered** | Per-stream and global ordering | | **Versioned** | Optimistic concurrency control | | **Subscriptions** | Real-time event notifications | | **Idempotent** | Handle duplicate writes safely | ## Technology Comparison | Technology | Best For | Limitations | | ---------------- | ------------------------- | -------------------------------- | | **EventStoreDB** | Pure event sourcing | Single-purpose | | **PostgreSQL** | Existing Postgres stack | Manual implementation | | **Kafka** | High-throughput streaming | Not ideal for per-stream queries | | **DynamoDB** | Serverless, AWS-native | Query limitations | | **Marten** | .NET ecosystems | .NET specific | ## Templates ### Template 1: PostgreSQL Event Store Schema ```sql -- Events table CREATE TABLE events ( id UUID PRIMARY KEY DEFAULT gen_random_uuid(), stream_id VARCHAR(255) NOT NULL, stream_type VARCHAR(255) NOT NULL, event_type VARCHAR(255) NOT NULL, event_data JSONB NOT NULL, metadata JSONB DEFAULT '{}', version BIGINT NOT NULL, global_position BIGSERIAL, created_at TIMESTAMPTZ DEFAULT NOW(), CONSTRAINT unique_stream_version UNIQUE (stream_id, version) ); -- Index for stream queries CREATE INDEX idx_events_stream_id ON events(stream_id, version); -- Index for global subscription CREATE INDEX idx_events_global_position ON events(global_position); -- Index for event type queries CREATE INDEX idx_events_event_type ON events(event_type); -- Index for time-based queries CREATE INDEX idx_events_created_at ON events(created_at); -- Snapshots table CREATE TABLE snapshots ( stream_id VARCHAR(255) PRIMARY KEY, stream_type VARCHAR(255) NOT NULL, snapshot_data JSONB NOT NULL, version BIGINT NOT NULL, created_at TIMESTAMPTZ DEFAULT NOW() ); -- Subscriptions checkpoint table CREATE TABLE subscription_checkpoints ( subscription_id VARCHAR(255) PRIMARY KEY, last_position BIGINT NOT NULL DEFAULT 0, updated_at TIMESTAMPTZ DEFAULT NOW() ); ``` ### Template 2: Python Event Store Implementation ```python from dataclasses import dataclass, field from datetime import datetime from typing import Any, Optional, List from uuid import UUID, uuid4 import json import asyncpg @dataclass class Event: stream_id: str event_type: str data: dict metadata: dict = field(default_factory=dict) event_id: UUID = field(default_factory=uuid4) version: Optional[int] = None global_position: Optional[int] = None created_at: datetime = field(default_factory=datetime.utcnow) class EventStore: def __init__(self, pool: asyncpg.Pool): self.pool = pool async def append_events( self, stream_id: str, stream_type: str, events: List[Event], expected_version: Optional[int] = None ) -> List[Event]: """Append events to a stream with optimistic concurrency.""" async with self.pool.acquire() as conn: async with conn.transaction(): # Check expected version if expected_version is not None: current = await conn.fetchval( "SELECT MAX(version) FROM events WHERE stream_id = $1", stream_id ) current = current or 0 if current != expected_version: raise ConcurrencyError( f"Expected version {expected_version}, got {current}" ) # Get starting version start_version = await conn.fetchval( "SELECT COALESCE(MAX(version), 0) + 1 FROM events WHERE stream_id = $1", stream_id ) # Insert events saved_events = [] for i, event in enumerate(events): event.version = start_version + i row = await conn.fetchrow( """ INSERT INTO events (id, stream_id, stream_type, event_type, event_data, metadata, version, created_at) VALUES ($1, $2, $3, $4, $5, $6, $7, $8) RETURNING global_position """, event.event_id, stream_id, stream_type, event.event_type, json.dumps(event.data), json.dumps(event.metadata), event.version, event.created_at ) event.global_position = row['global_position'] saved_events.append(event) return saved_events async def read_stream( self, stream_id: str, from_version: int = 0, limit: int = 1000 ) -> List[Event]: """Read events from a stream.""" async with self.pool.acquire() as conn: rows = await conn.fetch( """ SELECT id, stream_id, event_type, event_data, metadata, version, global_position, created_at FROM events WHERE stream_id = $1 AND version >= $2 ORDER BY version LIMIT $3 """, stream_id, from_version, limit ) return [self._row_to_event(row) for row in rows] async def read_all( self, from_position: int = 0, limit: int = 1000 ) -> List[Event]: """Read all events globally.""" async with self.pool.acquire() as conn: rows = await conn.fetch( """ SELECT id, stream_id, event_type, event_data, metadata, version, global_position, created_at FROM events WHERE global_position > $1 ORDER BY global_position LIMIT $2 """, from_position, limit ) return [self._row_to_event(row) for row in rows] async def subscribe( self, subscription_id: str, handler, from_position: int = 0, batch_size: int = 100 ): """Subscribe to all events from a position.""" # Get checkpoint async with self.pool.acquire() as conn: checkpoint = await conn.fetchval( """ SELECT last_position FROM subscription_checkpoints WHERE subscription_id = $1 """, subscription_id ) position = checkpoint or from_position while True: events = await self.read_all(position, batch_size) if not events: await asyncio.sleep(1) # Poll interval continue for event in events: await handler(event) position = event.global_position # Save checkpoint async with self.pool.acquire() as conn: await conn.execute( """ INSERT INTO subscription_checkpoints (subscription_id, last_position) VALUES ($1, $2) ON CONFLICT (subscription_id) DO UPDATE SET last_position = $2, updated_at = NOW() """, subscription_id, position ) def _row_to_event(self, row) -> Event: return Event( event_id=row['id'], stream_id=row['stream_id'], event_type=row['event_type'], data=json.loads(row['event_data']), metadata=json.loads(row['metadata']), version=row['version'], global_position=row['global_position'], created_at=row['created_at'] ) class ConcurrencyError(Exception): """Raised when optimistic concurrency check fails.""" pass ``` ### Template 3: EventStoreDB Usage ```python from esdbclient import EventStoreDBClient, NewEvent, StreamState import json # Connect client = EventStoreDBClient(uri="esdb://localhost:2113?tls=false") # Append events def append_events(stream_name: str, events: list, expected_revision=None): new_events = [ NewEvent( type=event['type'], data=json.dumps(event['data']).encode(), metadata=json.dumps(event.get('metadata', {})).encode() ) for event in events ] if expected_revision is None: state = StreamState.ANY elif expected_revision == -1: state = StreamState.NO_STREAM else: state = expected_revision return client.append_to_stream( stream_name=stream_name, events=new_events, current_version=state ) # Read stream def read_stream(stream_name: str, from_revision: int = 0): events = client.get_stream( stream_name=stream_name, stream_position=from_revision ) return [ { 'type': event.type, 'data': json.loads(event.data), 'metadata': json.loads(event.metadata) if event.metadata else {}, 'stream_position': event.stream_position, 'commit_position': event.commit_position } for event in events ] # Subscribe to all async def subscribe_to_all(handler, from_position: int = 0): subscription = client.subscribe_to_all(commit_position=from_position) async for event in subscription: await handler({ 'type': event.type, 'data': json.loads(event.data), 'stream_id': event.stream_name, 'position': event.commit_position }) # Category projection ($ce-Category) def read_category(category: str): """Read all events for a category using system projection.""" return read_stream(f"$ce-{category}") ``` ### Template 4: DynamoDB Event Store ```python import boto3 from boto3.dynamodb.conditions import Key from datetime import datetime import json import uuid class DynamoEventStore: def __init__(self, table_name: str): self.dynamodb = boto3.resource('dynamodb') self.table = self.dynamodb.Table(table_name) def append_events(self, stream_id: str, events: list, expected_version: int = None): """Append events with conditional write for concurrency.""" with self.table.batch_writer() as batch: for i, event in enumerate(events): version = (expected_version or 0) + i + 1 item = { 'PK': f"STREAM#{stream_id}", 'SK': f"VERSION#{version:020d}", 'GSI1PK': 'EVENTS', 'GSI1SK': datetime.utcnow().isoformat(), 'event_id': str(uuid.uuid4()), 'stream_id': stream_id, 'event_type': event['type'], 'event_data': json.dumps(event['data']), 'version': version, 'created_at': datetime.utcnow().isoformat() } batch.put_item(Item=item) return events def read_stream(self, stream_id: str, from_version: int = 0): """Read events from a stream.""" response = self.table.query( KeyConditionExpression=Key('PK').eq(f"STREAM#{stream_id}") & Key('SK').gte(f"VERSION#{from_version:020d}") ) return [ { 'event_type': item['event_type'], 'data': json.loads(item['event_data']), 'version': item['version'] } for item in response['Items'] ] # Table definition (CloudFormation/Terraform) """ DynamoDB Table: - PK (Partition Key): String - SK (Sort Key): String - GSI1PK, GSI1SK for global ordering Capacity: On-demand or provisioned based on throughput needs """ ``` ## Best Practices ### Do's - **Use stream IDs that include aggregate type** - `Order-{uuid}` - **Include correlation/causation IDs** - For tracing - **Version events from day one** - Plan for schema evolution - **Implement idempotency** - Use event IDs for deduplication - **Index appropriately** - For your query patterns ### Don'ts - **Don't update or delete events** - They're immutable facts - **Don't store large payloads** - Keep events small - **Don't skip optimistic concurrency** - Prevents data corruption - **Don't ignore backpressure** - Handle slow consumers ## Resources - [EventStoreDB](https://www.eventstore.com/) - [Marten Events](https://martendb.io/events/) - [Event Sourcing Pattern](https://docs.microsoft.com/en-us/azure/architecture/patterns/event-sourcing)
πŸ‘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

workflow-orchestration-patterns

Design durable workflows with Temporal for distributed systems.

coding
⭐1
# Workflow Orchestration Patterns Master workflow orchestration architecture with Temporal, covering fundamental design decisions, resilience patterns, and best practices for building reliable distributed systems. ## When to Use Workflow Orchestration ### Ideal Use Cases (Source: docs.temporal.io) - **Multi-step processes** spanning machines/services/databases - **Distributed transactions** requiring all-or-nothing semantics - **Long-running workflows** (hours to years) with automatic state persistence - **Failure recovery** that must resume from last successful step - **Business processes**: bookings, orders, campaigns, approvals - **Entity lifecycle management**: inventory tracking, account management, cart workflows - **Infrastructure automation**: CI/CD pipelines, provisioning, deployments - **Human-in-the-loop** systems requiring timeouts and escalations ### When NOT to Use - Simple CRUD operations (use direct API calls) - Pure data processing pipelines (use Airflow, batch processing) - Stateless request/response (use standard APIs) - Real-time streaming (use Kafka, event processors) ## Critical Design Decision: Workflows vs Activities **The Fundamental Rule** (Source: temporal.io/blog/workflow-engine-principles): - **Workflows** = Orchestration logic and decision-making - **Activities** = External interactions (APIs, databases, network calls) ### Workflows (Orchestration) **Characteristics:** - Contain business logic and coordination - **MUST be deterministic** (same inputs β†’ same outputs) - **Cannot** perform direct external calls - State automatically preserved across failures - Can run for years despite infrastructure failures **Example workflow tasks:** - Decide which steps to execute - Handle compensation logic - Manage timeouts and retries - Coordinate child workflows ### Activities (External Interactions) **Characteristics:** - Handle all external system interactions - Can be non-deterministic (API calls, DB writes) - Include built-in timeouts and retry logic - **Must be idempotent** (calling N times = calling once) - Short-lived (seconds to minutes typically) **Example activity tasks:** - Call payment gateway API - Write to database - Send emails or notifications - Query external services ### Design Decision Framework ``` Does it touch external systems? β†’ Activity Is it orchestration/decision logic? β†’ Workflow ``` ## Core Workflow Patterns ### 1. Saga Pattern with Compensation **Purpose**: Implement distributed transactions with rollback capability **Pattern** (Source: temporal.io/blog/compensating-actions-part-of-a-complete-breakfast-with-sagas): ``` For each step: 1. Register compensation BEFORE executing 2. Execute the step (via activity) 3. On failure, run all compensations in reverse order (LIFO) ``` **Example: Payment Workflow** 1. Reserve inventory (compensation: release inventory) 2. Charge payment (compensation: refund payment) 3. Fulfill order (compensation: cancel fulfillment) **Critical Requirements:** - Compensations must be idempotent - Register compensation BEFORE executing step - Run compensations in reverse order - Handle partial failures gracefully ### 2. Entity Workflows (Actor Model) **Purpose**: Long-lived workflow representing single entity instance **Pattern** (Source: docs.temporal.io/evaluate/use-cases-design-patterns): - One workflow execution = one entity (cart, account, inventory item) - Workflow persists for entity lifetime - Receives signals for state changes - Supports queries for current state **Example Use Cases:** - Shopping cart (add items, checkout, expiration) - Bank account (deposits, withdrawals, balance checks) - Product inventory (stock updates, reservations) **Benefits:** - Encapsulates entity behavior - Guarantees consistency per entity - Natural event sourcing ### 3. Fan-Out/Fan-In (Parallel Execution) **Purpose**: Execute multiple tasks in parallel, aggregate results **Pattern:** - Spawn child workflows or parallel activities - Wait for all to complete - Aggregate results - Handle partial failures **Scaling Rule** (Source: temporal.io/blog/workflow-engine-principles): - Don't scale individual workflows - For 1M tasks: spawn 1K child workflows Γ— 1K tasks each - Keep each workflow bounded ### 4. Async Callback Pattern **Purpose**: Wait for external event or human approval **Pattern:** - Workflow sends request and waits for signal - External system processes asynchronously - Sends signal to resume workflow - Workflow continues with response **Use Cases:** - Human approval workflows - Webhook callbacks - Long-running external processes ## State Management and Determinism ### Automatic State Preservation **How Temporal Works** (Source: docs.temporal.io/workflows): - Complete program state preserved automatically - Event History records every command and event - Seamless recovery from crashes - Applications restore pre-failure state ### Determinism Constraints **Workflows Execute as State Machines**: - Replay behavior must be consistent - Same inputs β†’ identical outputs every time **Prohibited in Workflows** (Source: docs.temporal.io/workflows): - ❌ Threading, locks, synchronization primitives - ❌ Random number generation (`random()`) - ❌ Global state or static variables - ❌ System time (`datetime.now()`) - ❌ Direct file I/O or network calls - ❌ Non-deterministic libraries **Allowed in Workflows**: - βœ… `workflow.now()` (deterministic time) - βœ… `workflow.random()` (deterministic random) - βœ… Pure functions and calculations - βœ… Calling activities (non-deterministic operations) ### Versioning Strategies **Challenge**: Changing workflow code while old executions still running **Solutions**: 1. **Versioning API**: Use `workflow.get_version()` for safe changes 2. **New Workflow Type**: Create new workflow, route new executions to it 3. **Backward Compatibility**: Ensure old events replay correctly ## Resilience and Error Handling ### Retry Policies **Default Behavior**: Temporal retries activities forever **Configure Retry**: - Initial retry interval - Backoff coefficient (exponential backoff) - Maximum interval (cap retry delay) - Maximum attempts (eventually fail) **Non-Retryable Errors**: - Invalid input (validation failures) - Business rule violations - Permanent failures (resource not found) ### Idempotency Requirements **Why Critical** (Source: docs.temporal.io/activities): - Activities may execute multiple times - Network failures trigger retries - Duplicate execution must be safe **Implementation Strategies**: - Idempotency keys (deduplication) - Check-then-act with unique constraints - Upsert operations instead of insert - Track processed request IDs ### Activity Heartbeats **Purpose**: Detect stalled long-running activities **Pattern**: - Activity sends periodic heartbeat - Includes progress information - Timeout if no heartbeat received - Enables progress-based retry ## Best Practices ### Workflow Design 1. **Keep workflows focused** - Single responsibility per workflow 2. **Small workflows** - Use child workflows for scalability 3. **Clear boundaries** - Workflow orchestrates, activities execute 4. **Test locally** - Use time-skipping test environment ### Activity Design 1. **Idempotent operations** - Safe to retry 2. **Short-lived** - Seconds to minutes, not hours 3. **Timeout configuration** - Always set timeouts 4. **Heartbeat for long tasks** - Report progress 5. **Error handling** - Distinguish retryable vs non-retryable ### Common Pitfalls **Workflow Violations**: - Using `datetime.now()` instead of `workflow.now()` - Threading or async operations in workflow code - Calling external APIs directly from workflow - Non-deterministic logic in workflows **Activity Mistakes**: - Non-idempotent operations (can't handle retries) - Missing timeouts (activities run forever) - No error classification (retry validation errors) - Ignoring payload limits (2MB per argument) ### Operational Considerations **Monitoring**: - Workflow execution duration - Activity failure rates - Retry attempts and backoff - Pending workflow counts **Scalability**: - Horizontal scaling with workers - Task queue partitioning - Child workflow decomposition - Activity batching when appropriate ## Additional Resources **Official Documentation**: - Temporal Core Concepts: docs.temporal.io/workflows - Workflow Patterns: docs.temporal.io/evaluate/use-cases-design-patterns - Best Practices: docs.temporal.io/develop/best-practices - Saga Pattern: temporal.io/blog/saga-pattern-made-easy **Key Principles**: 1. Workflows = orchestration, Activities = external calls 2. Determinism is non-negotiable for workflows 3. Idempotency is critical for activities 4. State preservation is automatic 5. Design for failure and recovery
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

kpi-dashboard-design

Design effective KPI dashboards with metrics selection,

coding
⭐1
# KPI Dashboard Design Comprehensive patterns for designing effective Key Performance Indicator (KPI) dashboards that drive business decisions. ## When to Use This Skill - Designing executive dashboards - Selecting meaningful KPIs - Building real-time monitoring displays - Creating department-specific metrics views - Improving existing dashboard layouts - Establishing metric governance ## Core Concepts ### 1. KPI Framework | Level | Focus | Update Frequency | Audience | | --------------- | ---------------- | ----------------- | ---------- | | **Strategic** | Long-term goals | Monthly/Quarterly | Executives | | **Tactical** | Department goals | Weekly/Monthly | Managers | | **Operational** | Day-to-day | Real-time/Daily | Teams | ### 2. SMART KPIs ``` Specific: Clear definition Measurable: Quantifiable Achievable: Realistic targets Relevant: Aligned to goals Time-bound: Defined period ``` ### 3. Dashboard Hierarchy ``` β”œβ”€β”€ Executive Summary (1 page) β”‚ β”œβ”€β”€ 4-6 headline KPIs β”‚ β”œβ”€β”€ Trend indicators β”‚ └── Key alerts β”œβ”€β”€ Department Views β”‚ β”œβ”€β”€ Sales Dashboard β”‚ β”œβ”€β”€ Marketing Dashboard β”‚ β”œβ”€β”€ Operations Dashboard β”‚ └── Finance Dashboard └── Detailed Drilldowns β”œβ”€β”€ Individual metrics └── Root cause analysis ``` ## Common KPIs by Department ### Sales KPIs ```yaml Revenue Metrics: - Monthly Recurring Revenue (MRR) - Annual Recurring Revenue (ARR) - Average Revenue Per User (ARPU) - Revenue Growth Rate Pipeline Metrics: - Sales Pipeline Value - Win Rate - Average Deal Size - Sales Cycle Length Activity Metrics: - Calls/Emails per Rep - Demos Scheduled - Proposals Sent - Close Rate ``` ### Marketing KPIs ```yaml Acquisition: - Cost Per Acquisition (CPA) - Customer Acquisition Cost (CAC) - Lead Volume - Marketing Qualified Leads (MQL) Engagement: - Website Traffic - Conversion Rate - Email Open/Click Rate - Social Engagement ROI: - Marketing ROI - Campaign Performance - Channel Attribution - CAC Payback Period ``` ### Product KPIs ```yaml Usage: - Daily/Monthly Active Users (DAU/MAU) - Session Duration - Feature Adoption Rate - Stickiness (DAU/MAU) Quality: - Net Promoter Score (NPS) - Customer Satisfaction (CSAT) - Bug/Issue Count - Time to Resolution Growth: - User Growth Rate - Activation Rate - Retention Rate - Churn Rate ``` ### Finance KPIs ```yaml Profitability: - Gross Margin - Net Profit Margin - EBITDA - Operating Margin Liquidity: - Current Ratio - Quick Ratio - Cash Flow - Working Capital Efficiency: - Revenue per Employee - Operating Expense Ratio - Days Sales Outstanding - Inventory Turnover ``` ## Dashboard Layout Patterns ### Pattern 1: Executive Summary ``` β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ EXECUTIVE DASHBOARD [Date Range β–Ό] β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ REVENUE β”‚ PROFIT β”‚ CUSTOMERS β”‚ NPS SCORE β”‚ β”‚ $2.4M β”‚ $450K β”‚ 12,450 β”‚ 72 β”‚ β”‚ β–² 12% β”‚ β–² 8% β”‚ β–² 15% β”‚ β–² 5pts β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ β”‚ β”‚ Revenue Trend β”‚ Revenue by Product β”‚ β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ β”‚ β”‚ /\ /\ β”‚ β”‚ β”‚ β–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆ 45% β”‚ β”‚ β”‚ β”‚ / \ / \ /\ β”‚ β”‚ β”‚ β–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆ 32% β”‚ β”‚ β”‚ β”‚ / \/ \ / \ β”‚ β”‚ β”‚ β–ˆβ–ˆβ–ˆβ–ˆ 18% β”‚ β”‚ β”‚ β”‚ / \/ \ β”‚ β”‚ β”‚ β–ˆβ–ˆ 5% β”‚ β”‚ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β”‚ β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ πŸ”΄ Alert: Churn rate exceeded threshold (>5%) β”‚ β”‚ 🟑 Warning: Support ticket volume 20% above average β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ ``` ### Pattern 2: SaaS Metrics Dashboard ``` β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ SAAS METRICS Jan 2024 [Monthly β–Ό] β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ MRR GROWTH β”‚ β”‚ β”‚ MRR β”‚ β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ β”‚ β”‚ $125,000 β”‚ β”‚ β”‚ /── β”‚ β”‚ β”‚ β”‚ β–² 8% β”‚ β”‚ β”‚ /────/ β”‚ β”‚ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β”‚ /────/ β”‚ β”‚ β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ β”‚ /────/ β”‚ β”‚ β”‚ β”‚ ARR β”‚ β”‚ β”‚ /────/ β”‚ β”‚ β”‚ β”‚ $1,500,000 β”‚ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β”‚ β”‚ β–² 15% β”‚ β”‚ J F M A M J J A S O N D β”‚ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”Όβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ UNIT ECONOMICS β”‚ COHORT RETENTION β”‚ β”‚ β”‚ β”‚ β”‚ CAC: $450 β”‚ Month 1: β–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆ 100% β”‚ β”‚ LTV: $2,700 β”‚ Month 3: β–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆ 85% β”‚ β”‚ LTV/CAC: 6.0x β”‚ Month 6: β–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆ 80% β”‚ β”‚ β”‚ Month 12: β–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆ 72% β”‚ β”‚ Payback: 4 months β”‚ β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ CHURN ANALYSIS β”‚ β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ β”‚ β”‚ Gross β”‚ Net β”‚ Logo β”‚ Expansion β”‚ β”‚ β”‚ β”‚ 4.2% β”‚ 1.8% β”‚ 3.1% β”‚ 2.4% β”‚ β”‚ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ ``` ### Pattern 3: Real-time Operations ``` β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ OPERATIONS CENTER Live ● Last: 10:42:15 β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ SYSTEM HEALTH β”‚ SERVICE STATUS β”‚ β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ β”‚ β”‚ β”‚ CPU MEM DISK β”‚ β”‚ ● API Gateway Healthy β”‚ β”‚ β”‚ 45% 72% 58% β”‚ β”‚ ● User Service Healthy β”‚ β”‚ β”‚ β–ˆβ–ˆβ–ˆ β–ˆβ–ˆβ–ˆβ–ˆ β–ˆβ–ˆβ–ˆ β”‚ β”‚ ● Payment Service Degraded β”‚ β”‚ β”‚ β–ˆβ–ˆβ–ˆ β–ˆβ–ˆβ–ˆβ–ˆ β–ˆβ–ˆβ–ˆ β”‚ β”‚ ● Database Healthy β”‚ β”‚ β”‚ β–ˆβ–ˆβ–ˆ β–ˆβ–ˆβ–ˆβ–ˆ β–ˆβ–ˆβ–ˆ β”‚ β”‚ ● Cache Healthy β”‚ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”Όβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ REQUEST THROUGHPUT β”‚ ERROR RATE β”‚ β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ β”‚ β”‚ β–β–‚β–ƒβ–„β–…β–†β–‡β–ˆβ–‡β–†β–…β–„β–ƒβ–‚β–β–‚β–ƒβ–„β–… β”‚ β”‚ β”‚ ▁▁▁▁▁▂▁▁▁▁▁▁▁▁▁▁▁▁▁▁ β”‚ β”‚ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β”‚ Current: 12,450 req/s β”‚ Current: 0.02% β”‚ β”‚ Peak: 18,200 req/s β”‚ Threshold: 1.0% β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ RECENT ALERTS β”‚ β”‚ 10:40 🟑 High latency on payment-service (p99 > 500ms) β”‚ β”‚ 10:35 🟒 Resolved: Database connection pool recovered β”‚ β”‚ 10:22 πŸ”΄ Payment service circuit breaker tripped β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ ``` ## Implementation Patterns ### SQL for KPI Calculations ```sql -- Monthly Recurring Revenue (MRR) WITH mrr_calculation AS ( SELECT DATE_TRUNC('month', billing_date) AS month, SUM( CASE subscription_interval WHEN 'monthly' THEN amount WHEN 'yearly' THEN amount / 12 WHEN 'quarterly' THEN amount / 3 END ) AS mrr FROM subscriptions WHERE status = 'active' GROUP BY DATE_TRUNC('month', billing_date) ) SELECT month, mrr, LAG(mrr) OVER (ORDER BY month) AS prev_mrr, (mrr - LAG(mrr) OVER (ORDER BY month)) / LAG(mrr) OVER (ORDER BY month) * 100 AS growth_pct FROM mrr_calculation; -- Cohort Retention WITH cohorts AS ( SELECT user_id, DATE_TRUNC('month', created_at) AS cohort_month FROM users ), activity AS ( SELECT user_id, DATE_TRUNC('month', event_date) AS activity_month FROM user_events WHERE event_type = 'active_session' ) SELECT c.cohort_month, EXTRACT(MONTH FROM age(a.activity_month, c.cohort_month)) AS months_since_signup, COUNT(DISTINCT a.user_id) AS active_users, COUNT(DISTINCT a.user_id)::FLOAT / COUNT(DISTINCT c.user_id) * 100 AS retention_rate FROM cohorts c LEFT JOIN activity a ON c.user_id = a.user_id AND a.activity_month >= c.cohort_month GROUP BY c.cohort_month, EXTRACT(MONTH FROM age(a.activity_month, c.cohort_month)) ORDER BY c.cohort_month, months_since_signup; -- Customer Acquisition Cost (CAC) SELECT DATE_TRUNC('month', acquired_date) AS month, SUM(marketing_spend) / NULLIF(COUNT(new_customers), 0) AS cac, SUM(marketing_spend) AS total_spend, COUNT(new_customers) AS customers_acquired FROM ( SELECT DATE_TRUNC('month', u.created_at) AS acquired_date, u.id AS new_customers, m.spend AS marketing_spend FROM users u JOIN marketing_spend m ON DATE_TRUNC('month', u.created_at) = m.month WHERE u.source = 'marketing' ) acquisition GROUP BY DATE_TRUNC('month', acquired_date); ``` ### Python Dashboard Code (Streamlit) ```python import streamlit as st import pandas as pd import plotly.express as px import plotly.graph_objects as go st.set_page_config(page_title="KPI Dashboard", layout="wide") # Header with date filter col1, col2 = st.columns([3, 1]) with col1: st.title("Executive Dashboard") with col2: date_range = st.selectbox( "Period", ["Last 7 Days", "Last 30 Days", "Last Quarter", "YTD"] ) # KPI Cards def metric_card(label, value, delta, prefix="", suffix=""): delta_color = "green" if delta >= 0 else "red" delta_arrow = "β–²" if delta >= 0 else "β–Ό" st.metric( label=label, value=f"{prefix}{value:,.0f}{suffix}", delta=f"{delta_arrow} {abs(delta):.1f}%" ) col1, col2, col3, col4 = st.columns(4) with col1: metric_card("Revenue", 2400000, 12.5, prefix="$") with col2: metric_card("Customers", 12450, 15.2) with col3: metric_card("NPS Score", 72, 5.0) with col4: metric_card("Churn Rate", 4.2, -0.8, suffix="%") # Charts col1, col2 = st.columns(2) with col1: st.subheader("Revenue Trend") revenue_data = pd.DataFrame({ 'Month': pd.date_range('2024-01-01', periods=12, freq='M'), 'Revenue': [180000, 195000, 210000, 225000, 240000, 255000, 270000, 285000, 300000, 315000, 330000, 345000] }) fig = px.line(revenue_data, x='Month', y='Revenue', line_shape='spline', markers=True) fig.update_layout(height=300) st.plotly_chart(fig, use_container_width=True) with col2: st.subheader("Revenue by Product") product_data = pd.DataFrame({ 'Product': ['Enterprise', 'Professional', 'Starter', 'Other'], 'Revenue': [45, 32, 18, 5] }) fig = px.pie(product_data, values='Revenue', names='Product', hole=0.4) fig.update_layout(height=300) st.plotly_chart(fig, use_container_width=True) # Cohort Heatmap st.subheader("Cohort Retention") cohort_data = pd.DataFrame({ 'Cohort': ['Jan', 'Feb', 'Mar', 'Apr', 'May'], 'M0': [100, 100, 100, 100, 100], 'M1': [85, 87, 84, 86, 88], 'M2': [78, 80, 76, 79, None], 'M3': [72, 74, 70, None, None], 'M4': [68, 70, None, None, None], }) fig = go.Figure(data=go.Heatmap( z=cohort_data.iloc[:, 1:].values, x=['M0', 'M1', 'M2', 'M3', 'M4'], y=cohort_data['Cohort'], colorscale='Blues', text=cohort_data.iloc[:, 1:].values, texttemplate='%{text}%', textfont={"size": 12}, )) fig.update_layout(height=250) st.plotly_chart(fig, use_container_width=True) # Alerts Section st.subheader("Alerts") alerts = [ {"level": "error", "message": "Churn rate exceeded threshold (>5%)"}, {"level": "warning", "message": "Support ticket volume 20% above average"}, ] for alert in alerts: if alert["level"] == "error": st.error(f"πŸ”΄ {alert['message']}") elif alert["level"] == "warning": st.warning(f"🟑 {alert['message']}") ``` ## Best Practices ### Do's - **Limit to 5-7 KPIs** - Focus on what matters - **Show context** - Comparisons, trends, targets - **Use consistent colors** - Red=bad, green=good - **Enable drilldown** - From summary to detail - **Update appropriately** - Match metric frequency ### Don'ts - **Don't show vanity metrics** - Focus on actionable data - **Don't overcrowd** - White space aids comprehension - **Don't use 3D charts** - They distort perception - **Don't hide methodology** - Document calculations - **Don't ignore mobile** - Ensure responsive design ## Resources - [Stephen Few's Dashboard Design](https://www.perceptualedge.com/articles/visual_business_intelligence/rules_for_using_color.pdf) - [Edward Tufte's Principles](https://www.edwardtufte.com/tufte/) - [Google Data Studio Gallery](https://datastudio.google.com/gallery)
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

secrets-management

Implement secure secrets management for CI/CD pipelines using

coding
⭐1
# Secrets Management Secure secrets management practices for CI/CD pipelines using Vault, AWS Secrets Manager, and other tools. ## Purpose Implement secure secrets management in CI/CD pipelines without hardcoding sensitive information. ## When to Use - Store API keys and credentials - Manage database passwords - Handle TLS certificates - Rotate secrets automatically - Implement least-privilege access ## Secrets Management Tools ### HashiCorp Vault - Centralized secrets management - Dynamic secrets generation - Secret rotation - Audit logging - Fine-grained access control ### AWS Secrets Manager - AWS-native solution - Automatic rotation - Integration with RDS - CloudFormation support ### Azure Key Vault - Azure-native solution - HSM-backed keys - Certificate management - RBAC integration ### Google Secret Manager - GCP-native solution - Versioning - IAM integration ## HashiCorp Vault Integration ### Setup Vault ```bash # Start Vault dev server vault server -dev # Set environment export VAULT_ADDR='http://127.0.0.1:8200' export VAULT_TOKEN='root' # Enable secrets engine vault secrets enable -path=secret kv-v2 # Store secret vault kv put secret/database/config username=admin password=secret ``` ### GitHub Actions with Vault ```yaml name: Deploy with Vault Secrets on: [push] jobs: deploy: runs-on: ubuntu-latest steps: - uses: actions/checkout@v4 - name: Import Secrets from Vault uses: hashicorp/vault-action@v2 with: url: https://vault.example.com:8200 token: ${{ secrets.VAULT_TOKEN }} secrets: | secret/data/database username | DB_USERNAME ; secret/data/database password | DB_PASSWORD ; secret/data/api key | API_KEY - name: Use secrets run: | echo "Connecting to database as $DB_USERNAME" # Use $DB_PASSWORD, $API_KEY ``` ### GitLab CI with Vault ```yaml deploy: image: vault:latest before_script: - export VAULT_ADDR=https://vault.example.com:8200 - export VAULT_TOKEN=$VAULT_TOKEN - apk add curl jq script: - | DB_PASSWORD=$(vault kv get -field=password secret/database/config) API_KEY=$(vault kv get -field=key secret/api/credentials) echo "Deploying with secrets..." # Use $DB_PASSWORD, $API_KEY ``` **Reference:** See `references/vault-setup.md` ## AWS Secrets Manager ### Store Secret ```bash aws secretsmanager create-secret \ --name production/database/password \ --secret-string "super-secret-password" ``` ### Retrieve in GitHub Actions ```yaml - name: Configure AWS credentials uses: aws-actions/configure-aws-credentials@v4 with: aws-access-key-id: ${{ secrets.AWS_ACCESS_KEY_ID }} aws-secret-access-key: ${{ secrets.AWS_SECRET_ACCESS_KEY }} aws-region: us-west-2 - name: Get secret from AWS run: | SECRET=$(aws secretsmanager get-secret-value \ --secret-id production/database/password \ --query SecretString \ --output text) echo "::add-mask::$SECRET" echo "DB_PASSWORD=$SECRET" >> $GITHUB_ENV - name: Use secret run: | # Use $DB_PASSWORD ./deploy.sh ``` ### Terraform with AWS Secrets Manager ```hcl data "aws_secretsmanager_secret_version" "db_password" { secret_id = "production/database/password" } resource "aws_db_instance" "main" { allocated_storage = 100 engine = "postgres" instance_class = "db.t3.large" username = "admin" password = jsondecode(data.aws_secretsmanager_secret_version.db_password.secret_string)["password"] } ``` ## GitHub Secrets ### Organization/Repository Secrets ```yaml - name: Use GitHub secret run: | echo "API Key: ${{ secrets.API_KEY }}" echo "Database URL: ${{ secrets.DATABASE_URL }}" ``` ### Environment Secrets ```yaml deploy: runs-on: ubuntu-latest environment: production steps: - name: Deploy run: | echo "Deploying with ${{ secrets.PROD_API_KEY }}" ``` **Reference:** See `references/github-secrets.md` ## GitLab CI/CD Variables ### Project Variables ```yaml deploy: script: - echo "Deploying with $API_KEY" - echo "Database: $DATABASE_URL" ``` ### Protected and Masked Variables - Protected: Only available in protected branches - Masked: Hidden in job logs - File type: Stored as file ## Best Practices 1. **Never commit secrets** to Git 2. **Use different secrets** per environment 3. **Rotate secrets regularly** 4. **Implement least-privilege access** 5. **Enable audit logging** 6. **Use secret scanning** (GitGuardian, TruffleHog) 7. **Mask secrets in logs** 8. **Encrypt secrets at rest** 9. **Use short-lived tokens** when possible 10. **Document secret requirements** ## Secret Rotation ### Automated Rotation with AWS ```python import boto3 import json def lambda_handler(event, context): client = boto3.client('secretsmanager') # Get current secret response = client.get_secret_value(SecretId='my-secret') current_secret = json.loads(response['SecretString']) # Generate new password new_password = generate_strong_password() # Update database password update_database_password(new_password) # Update secret client.put_secret_value( SecretId='my-secret', SecretString=json.dumps({ 'username': current_secret['username'], 'password': new_password }) ) return {'statusCode': 200} ``` ### Manual Rotation Process 1. Generate new secret 2. Update secret in secret store 3. Update applications to use new secret 4. Verify functionality 5. Revoke old secret ## External Secrets Operator ### Kubernetes Integration ```yaml apiVersion: external-secrets.io/v1beta1 kind: SecretStore metadata: name: vault-backend namespace: production spec: provider: vault: server: "https://vault.example.com:8200" path: "secret" version: "v2" auth: kubernetes: mountPath: "kubernetes" role: "production" --- apiVersion: external-secrets.io/v1beta1 kind: ExternalSecret metadata: name: database-credentials namespace: production spec: refreshInterval: 1h secretStoreRef: name: vault-backend kind: SecretStore target: name: database-credentials creationPolicy: Owner data: - secretKey: username remoteRef: key: database/config property: username - secretKey: password remoteRef: key: database/config property: password ``` ## Secret Scanning ### Pre-commit Hook ```bash #!/bin/bash # .git/hooks/pre-commit # Check for secrets with TruffleHog docker run --rm -v "$(pwd):/repo" \ trufflesecurity/trufflehog:latest \ filesystem --directory=/repo if [ $? -ne 0 ]; then echo "❌ Secret detected! Commit blocked." exit 1 fi ``` ### CI/CD Secret Scanning ```yaml secret-scan: stage: security image: trufflesecurity/trufflehog:latest script: - trufflehog filesystem . allow_failure: false ``` ## Reference Files - `references/vault-setup.md` - HashiCorp Vault configuration - `references/github-secrets.md` - GitHub Secrets best practices ## Related Skills - `github-actions-templates` - For GitHub Actions integration - `gitlab-ci-patterns` - For GitLab CI integration - `deployment-pipeline-design` - For pipeline architecture
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

linkerd-patterns

Implement Linkerd service mesh patterns for lightweight,

architecture
⭐1
# Linkerd Patterns Production patterns for Linkerd service mesh - the lightweight, security-first service mesh for Kubernetes. ## When to Use This Skill - Setting up a lightweight service mesh - Implementing automatic mTLS - Configuring traffic splits for canary deployments - Setting up service profiles for per-route metrics - Implementing retries and timeouts - Multi-cluster service mesh ## Core Concepts ### 1. Linkerd Architecture ``` β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ Control Plane β”‚ β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ β”‚ β”‚ destiny β”‚ β”‚ identity β”‚ β”‚ proxy-inject β”‚ β”‚ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ Data Plane β”‚ β”‚ β”Œβ”€β”€β”€β”€β”€β” β”Œβ”€β”€β”€β”€β”€β” β”Œβ”€β”€β”€β”€β”€β” β”‚ β”‚ β”‚proxy│────│proxy│────│proxyβ”‚ β”‚ β”‚ β””β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”˜ β”‚ β”‚ β”‚ β”‚ β”‚ β”‚ β”‚ β”Œβ”€β”€β”΄β”€β”€β” β”Œβ”€β”€β”΄β”€β”€β” β”Œβ”€β”€β”΄β”€β”€β” β”‚ β”‚ β”‚ app β”‚ β”‚ app β”‚ β”‚ app β”‚ β”‚ β”‚ β””β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”˜ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ ``` ### 2. Key Resources | Resource | Purpose | | ----------------------- | ------------------------------------ | | **ServiceProfile** | Per-route metrics, retries, timeouts | | **TrafficSplit** | Canary deployments, A/B testing | | **Server** | Define server-side policies | | **ServerAuthorization** | Access control policies | ## Templates ### Template 1: Mesh Installation ```bash # Install CLI curl --proto '=https' --tlsv1.2 -sSfL https://run.linkerd.io/install | sh # Validate cluster linkerd check --pre # Install CRDs linkerd install --crds | kubectl apply -f - # Install control plane linkerd install | kubectl apply -f - # Verify installation linkerd check # Install viz extension (optional) linkerd viz install | kubectl apply -f - ``` ### Template 2: Inject Namespace ```yaml # Automatic injection for namespace apiVersion: v1 kind: Namespace metadata: name: my-app annotations: linkerd.io/inject: enabled --- # Or inject specific deployment apiVersion: apps/v1 kind: Deployment metadata: name: my-app annotations: linkerd.io/inject: enabled spec: template: metadata: annotations: linkerd.io/inject: enabled ``` ### Template 3: Service Profile with Retries ```yaml apiVersion: linkerd.io/v1alpha2 kind: ServiceProfile metadata: name: my-service.my-namespace.svc.cluster.local namespace: my-namespace spec: routes: - name: GET /api/users condition: method: GET pathRegex: /api/users responseClasses: - condition: status: min: 500 max: 599 isFailure: true isRetryable: true - name: POST /api/users condition: method: POST pathRegex: /api/users # POST not retryable by default isRetryable: false - name: GET /api/users/{id} condition: method: GET pathRegex: /api/users/[^/]+ timeout: 5s isRetryable: true retryBudget: retryRatio: 0.2 minRetriesPerSecond: 10 ttl: 10s ``` ### Template 4: Traffic Split (Canary) ```yaml apiVersion: split.smi-spec.io/v1alpha1 kind: TrafficSplit metadata: name: my-service-canary namespace: my-namespace spec: service: my-service backends: - service: my-service-stable weight: 900m # 90% - service: my-service-canary weight: 100m # 10% ``` ### Template 5: Server Authorization Policy ```yaml # Define the server apiVersion: policy.linkerd.io/v1beta1 kind: Server metadata: name: my-service-http namespace: my-namespace spec: podSelector: matchLabels: app: my-service port: http proxyProtocol: HTTP/1 --- # Allow traffic from specific clients apiVersion: policy.linkerd.io/v1beta1 kind: ServerAuthorization metadata: name: allow-frontend namespace: my-namespace spec: server: name: my-service-http client: meshTLS: serviceAccounts: - name: frontend namespace: my-namespace --- # Allow unauthenticated traffic (e.g., from ingress) apiVersion: policy.linkerd.io/v1beta1 kind: ServerAuthorization metadata: name: allow-ingress namespace: my-namespace spec: server: name: my-service-http client: unauthenticated: true networks: - cidr: 10.0.0.0/8 ``` ### Template 6: HTTPRoute for Advanced Routing ```yaml apiVersion: policy.linkerd.io/v1beta2 kind: HTTPRoute metadata: name: my-route namespace: my-namespace spec: parentRefs: - name: my-service kind: Service group: core port: 8080 rules: - matches: - path: type: PathPrefix value: /api/v2 - headers: - name: x-api-version value: v2 backendRefs: - name: my-service-v2 port: 8080 - matches: - path: type: PathPrefix value: /api backendRefs: - name: my-service-v1 port: 8080 ``` ### Template 7: Multi-cluster Setup ```bash # On each cluster, install with cluster credentials linkerd multicluster install | kubectl apply -f - # Link clusters linkerd multicluster link --cluster-name west \ --api-server-address https://west.example.com:6443 \ | kubectl apply -f - # Export a service to other clusters kubectl label svc/my-service mirror.linkerd.io/exported=true # Verify cross-cluster connectivity linkerd multicluster check linkerd multicluster gateways ``` ## Monitoring Commands ```bash # Live traffic view linkerd viz top deploy/my-app # Per-route metrics linkerd viz routes deploy/my-app # Check proxy status linkerd viz stat deploy -n my-namespace # View service dependencies linkerd viz edges deploy -n my-namespace # Dashboard linkerd viz dashboard ``` ## Debugging ```bash # Check injection status linkerd check --proxy -n my-namespace # View proxy logs kubectl logs deploy/my-app -c linkerd-proxy # Debug identity/TLS linkerd identity -n my-namespace # Tap traffic (live) linkerd viz tap deploy/my-app --to deploy/my-backend ``` ## Best Practices ### Do's - **Enable mTLS everywhere** - It's automatic with Linkerd - **Use ServiceProfiles** - Get per-route metrics and retries - **Set retry budgets** - Prevent retry storms - **Monitor golden metrics** - Success rate, latency, throughput ### Don'ts - **Don't skip check** - Always run `linkerd check` after changes - **Don't over-configure** - Linkerd defaults are sensible - **Don't ignore ServiceProfiles** - They unlock advanced features - **Don't forget timeouts** - Set appropriate values per route ## Resources - [Linkerd Documentation](https://linkerd.io/2.14/overview/) - [Service Profiles](https://linkerd.io/2.14/features/service-profiles/) - [Authorization Policy](https://linkerd.io/2.14/features/server-policy/)
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

service-mesh-observability

Implement comprehensive observability for service meshes including

architecture
⭐1
# Service Mesh Observability Complete guide to observability patterns for Istio, Linkerd, and service mesh deployments. ## When to Use This Skill - Setting up distributed tracing across services - Implementing service mesh metrics and dashboards - Debugging latency and error issues - Defining SLOs for service communication - Visualizing service dependencies - Troubleshooting mesh connectivity ## Core Concepts ### 1. Three Pillars of Observability ``` β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ Observability β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ Metrics β”‚ Traces β”‚ Logs β”‚ β”‚ β”‚ β”‚ β”‚ β”‚ β€’ Request rate β”‚ β€’ Span context β”‚ β€’ Access logs β”‚ β”‚ β€’ Error rate β”‚ β€’ Latency β”‚ β€’ Error details β”‚ β”‚ β€’ Latency P50 β”‚ β€’ Dependencies β”‚ β€’ Debug info β”‚ β”‚ β€’ Saturation β”‚ β€’ Bottlenecks β”‚ β€’ Audit trail β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ ``` ### 2. Golden Signals for Mesh | Signal | Description | Alert Threshold | | -------------- | ------------------------- | ----------------- | | **Latency** | Request duration P50, P99 | P99 > 500ms | | **Traffic** | Requests per second | Anomaly detection | | **Errors** | 5xx error rate | > 1% | | **Saturation** | Resource utilization | > 80% | ## Templates ### Template 1: Istio with Prometheus & Grafana ```yaml # Install Prometheus apiVersion: v1 kind: ConfigMap metadata: name: prometheus namespace: istio-system data: prometheus.yml: | global: scrape_interval: 15s scrape_configs: - job_name: 'istio-mesh' kubernetes_sd_configs: - role: endpoints namespaces: names: - istio-system relabel_configs: - source_labels: [__meta_kubernetes_service_name] action: keep regex: istio-telemetry --- # ServiceMonitor for Prometheus Operator apiVersion: monitoring.coreos.com/v1 kind: ServiceMonitor metadata: name: istio-mesh namespace: istio-system spec: selector: matchLabels: app: istiod endpoints: - port: http-monitoring interval: 15s ``` ### Template 2: Key Istio Metrics Queries ```promql # Request rate by service sum(rate(istio_requests_total{reporter="destination"}[5m])) by (destination_service_name) # Error rate (5xx) sum(rate(istio_requests_total{reporter="destination", response_code=~"5.."}[5m])) / sum(rate(istio_requests_total{reporter="destination"}[5m])) * 100 # P99 latency histogram_quantile(0.99, sum(rate(istio_request_duration_milliseconds_bucket{reporter="destination"}[5m])) by (le, destination_service_name)) # TCP connections sum(istio_tcp_connections_opened_total{reporter="destination"}) by (destination_service_name) # Request size histogram_quantile(0.99, sum(rate(istio_request_bytes_bucket{reporter="destination"}[5m])) by (le, destination_service_name)) ``` ### Template 3: Jaeger Distributed Tracing ```yaml # Jaeger installation for Istio apiVersion: install.istio.io/v1alpha1 kind: IstioOperator spec: meshConfig: enableTracing: true defaultConfig: tracing: sampling: 100.0 # 100% in dev, lower in prod zipkin: address: jaeger-collector.istio-system:9411 --- # Jaeger deployment apiVersion: apps/v1 kind: Deployment metadata: name: jaeger namespace: istio-system spec: selector: matchLabels: app: jaeger template: metadata: labels: app: jaeger spec: containers: - name: jaeger image: jaegertracing/all-in-one:1.50 ports: - containerPort: 5775 # UDP - containerPort: 6831 # Thrift - containerPort: 6832 # Thrift - containerPort: 5778 # Config - containerPort: 16686 # UI - containerPort: 14268 # HTTP - containerPort: 14250 # gRPC - containerPort: 9411 # Zipkin env: - name: COLLECTOR_ZIPKIN_HOST_PORT value: ":9411" ``` ### Template 4: Linkerd Viz Dashboard ```bash # Install Linkerd viz extension linkerd viz install | kubectl apply -f - # Access dashboard linkerd viz dashboard # CLI commands for observability # Top requests linkerd viz top deploy/my-app # Per-route metrics linkerd viz routes deploy/my-app --to deploy/backend # Live traffic inspection linkerd viz tap deploy/my-app --to deploy/backend # Service edges (dependencies) linkerd viz edges deployment -n my-namespace ``` ### Template 5: Grafana Dashboard JSON ```json { "dashboard": { "title": "Service Mesh Overview", "panels": [ { "title": "Request Rate", "type": "graph", "targets": [ { "expr": "sum(rate(istio_requests_total{reporter=\"destination\"}[5m])) by (destination_service_name)", "legendFormat": "{{destination_service_name}}" } ] }, { "title": "Error Rate", "type": "gauge", "targets": [ { "expr": "sum(rate(istio_requests_total{response_code=~\"5..\"}[5m])) / sum(rate(istio_requests_total[5m])) * 100" } ], "fieldConfig": { "defaults": { "thresholds": { "steps": [ { "value": 0, "color": "green" }, { "value": 1, "color": "yellow" }, { "value": 5, "color": "red" } ] } } } }, { "title": "P99 Latency", "type": "graph", "targets": [ { "expr": "histogram_quantile(0.99, sum(rate(istio_request_duration_milliseconds_bucket{reporter=\"destination\"}[5m])) by (le, destination_service_name))", "legendFormat": "{{destination_service_name}}" } ] }, { "title": "Service Topology", "type": "nodeGraph", "targets": [ { "expr": "sum(rate(istio_requests_total{reporter=\"destination\"}[5m])) by (source_workload, destination_service_name)" } ] } ] } } ``` ### Template 6: Kiali Service Mesh Visualization ```yaml # Kiali installation apiVersion: kiali.io/v1alpha1 kind: Kiali metadata: name: kiali namespace: istio-system spec: auth: strategy: anonymous # or openid, token deployment: accessible_namespaces: - "**" external_services: prometheus: url: http://prometheus.istio-system:9090 tracing: url: http://jaeger-query.istio-system:16686 grafana: url: http://grafana.istio-system:3000 ``` ### Template 7: OpenTelemetry Integration ```yaml # OpenTelemetry Collector for mesh apiVersion: v1 kind: ConfigMap metadata: name: otel-collector-config data: config.yaml: | receivers: otlp: protocols: grpc: endpoint: 0.0.0.0:4317 http: endpoint: 0.0.0.0:4318 zipkin: endpoint: 0.0.0.0:9411 processors: batch: timeout: 10s exporters: jaeger: endpoint: jaeger-collector:14250 tls: insecure: true prometheus: endpoint: 0.0.0.0:8889 service: pipelines: traces: receivers: [otlp, zipkin] processors: [batch] exporters: [jaeger] metrics: receivers: [otlp] processors: [batch] exporters: [prometheus] --- # Istio Telemetry v2 with OTel apiVersion: telemetry.istio.io/v1alpha1 kind: Telemetry metadata: name: mesh-default namespace: istio-system spec: tracing: - providers: - name: otel randomSamplingPercentage: 10 ``` ## Alerting Rules ```yaml apiVersion: monitoring.coreos.com/v1 kind: PrometheusRule metadata: name: mesh-alerts namespace: istio-system spec: groups: - name: mesh.rules rules: - alert: HighErrorRate expr: | sum(rate(istio_requests_total{response_code=~"5.."}[5m])) by (destination_service_name) / sum(rate(istio_requests_total[5m])) by (destination_service_name) > 0.05 for: 5m labels: severity: critical annotations: summary: "High error rate for {{ $labels.destination_service_name }}" - alert: HighLatency expr: | histogram_quantile(0.99, sum(rate(istio_request_duration_milliseconds_bucket[5m])) by (le, destination_service_name)) > 1000 for: 5m labels: severity: warning annotations: summary: "High P99 latency for {{ $labels.destination_service_name }}" - alert: MeshCertExpiring expr: | (certmanager_certificate_expiration_timestamp_seconds - time()) / 86400 < 7 labels: severity: warning annotations: summary: "Mesh certificate expiring in less than 7 days" ``` ## Best Practices ### Do's - **Sample appropriately** - 100% in dev, 1-10% in prod - **Use trace context** - Propagate headers consistently - **Set up alerts** - For golden signals - **Correlate metrics/traces** - Use exemplars - **Retain strategically** - Hot/cold storage tiers ### Don'ts - **Don't over-sample** - Storage costs add up - **Don't ignore cardinality** - Limit label values - **Don't skip dashboards** - Visualize dependencies - **Don't forget costs** - Monitor observability costs ## Resources - [Istio Observability](https://istio.io/latest/docs/tasks/observability/) - [Linkerd Observability](https://linkerd.io/2.14/features/dashboard/) - [OpenTelemetry](https://opentelemetry.io/) - [Kiali](https://kiali.io/)
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

track-management

Use this skill when creating, managing, or working with Conductor

coding
⭐1
# Track Management Guide for creating, managing, and completing Conductor tracks - the logical work units that organize features, bugs, and refactors through specification, planning, and implementation phases. ## When to Use This Skill - Creating new feature, bug, or refactor tracks - Writing or reviewing spec.md files - Creating or updating plan.md files - Managing track lifecycle from creation to completion - Understanding track status markers and conventions - Working with the tracks.md registry - Interpreting or updating track metadata ## Track Concept A track is a logical work unit that encapsulates a complete piece of work. Each track has: - A unique identifier - A specification defining requirements - A phased plan breaking work into tasks - Metadata tracking status and progress Tracks provide semantic organization for work, enabling: - Clear scope boundaries - Progress tracking - Git-aware operations (revert by track) - Team coordination ## Track Types ### feature New functionality or capabilities. Use for: - New user-facing features - New API endpoints - New integrations - Significant enhancements ### bug Defect fixes. Use for: - Incorrect behavior - Error conditions - Performance regressions - Security vulnerabilities ### chore Maintenance and housekeeping. Use for: - Dependency updates - Configuration changes - Documentation updates - Cleanup tasks ### refactor Code improvement without behavior change. Use for: - Code restructuring - Pattern adoption - Technical debt reduction - Performance optimization (same behavior, better performance) ## Track ID Format Track IDs follow the pattern: `{shortname}_{YYYYMMDD}` - **shortname**: 2-4 word kebab-case description (e.g., `user-auth`, `api-rate-limit`) - **YYYYMMDD**: Creation date in ISO format Examples: - `user-auth_20250115` - `fix-login-error_20250115` - `upgrade-deps_20250115` - `refactor-api-client_20250115` ## Track Lifecycle ### 1. Creation (newTrack) **Define Requirements** 1. Gather requirements through interactive Q&A 2. Identify acceptance criteria 3. Determine scope boundaries 4. Identify dependencies **Generate Specification** 1. Create `spec.md` with structured requirements 2. Document functional and non-functional requirements 3. Define acceptance criteria 4. List dependencies and constraints **Generate Plan** 1. Create `plan.md` with phased task breakdown 2. Organize tasks into logical phases 3. Add verification tasks after phases 4. Estimate effort and complexity **Register Track** 1. Add entry to `tracks.md` registry 2. Create track directory structure 3. Generate `metadata.json` 4. Create track `index.md` ### 2. Implementation **Execute Tasks** 1. Select next pending task from plan 2. Mark task as in-progress 3. Implement following workflow (TDD) 4. Mark task complete with commit SHA **Update Status** 1. Update task markers in plan.md 2. Record commit SHAs for traceability 3. Update phase progress 4. Update track status in tracks.md **Verify Progress** 1. Complete verification tasks 2. Wait for checkpoint approval 3. Record checkpoint commits ### 3. Completion **Sync Documentation** 1. Update product.md if features added 2. Update tech-stack.md if dependencies changed 3. Verify all acceptance criteria met **Archive or Delete** 1. Mark track as completed in tracks.md 2. Record completion date 3. Archive or retain track directory ## Specification (spec.md) Structure ```markdown # {Track Title} ## Overview Brief description of what this track accomplishes and why. ## Functional Requirements ### FR-1: {Requirement Name} Description of the functional requirement. - Acceptance: How to verify this requirement is met ### FR-2: {Requirement Name} ... ## Non-Functional Requirements ### NFR-1: {Requirement Name} Description of the non-functional requirement (performance, security, etc.) - Target: Specific measurable target - Verification: How to test ## Acceptance Criteria - [ ] Criterion 1: Specific, testable condition - [ ] Criterion 2: Specific, testable condition - [ ] Criterion 3: Specific, testable condition ## Scope ### In Scope - Explicitly included items - Features to implement - Components to modify ### Out of Scope - Explicitly excluded items - Future considerations - Related but separate work ## Dependencies ### Internal - Other tracks or components this depends on - Required context artifacts ### External - Third-party services or APIs - External dependencies ## Risks and Mitigations | Risk | Impact | Mitigation | | ---------------- | --------------- | ------------------- | | Risk description | High/Medium/Low | Mitigation strategy | ## Open Questions - [ ] Question that needs resolution - [x] Resolved question - Answer ``` ## Plan (plan.md) Structure ```markdown # Implementation Plan: {Track Title} Track ID: `{track-id}` Created: YYYY-MM-DD Status: pending | in-progress | completed ## Overview Brief description of implementation approach. ## Phase 1: {Phase Name} ### Tasks - [ ] **Task 1.1**: Task description - Sub-task or detail - Sub-task or detail - [ ] **Task 1.2**: Task description - [ ] **Task 1.3**: Task description ### Verification - [ ] **Verify 1.1**: Verification step for phase ## Phase 2: {Phase Name} ### Tasks - [ ] **Task 2.1**: Task description - [ ] **Task 2.2**: Task description ### Verification - [ ] **Verify 2.1**: Verification step for phase ## Phase 3: Finalization ### Tasks - [ ] **Task 3.1**: Update documentation - [ ] **Task 3.2**: Final integration test ### Verification - [ ] **Verify 3.1**: All acceptance criteria met ## Checkpoints | Phase | Checkpoint SHA | Date | Status | | ------- | -------------- | ---- | ------- | | Phase 1 | | | pending | | Phase 2 | | | pending | | Phase 3 | | | pending | ``` ## Status Marker Conventions Use consistent markers in plan.md: | Marker | Meaning | Usage | | ------ | ----------- | --------------------------- | | `[ ]` | Pending | Task not started | | `[~]` | In Progress | Currently being worked | | `[x]` | Complete | Task finished (include SHA) | | `[-]` | Skipped | Intentionally not done | | `[!]` | Blocked | Waiting on dependency | Example: ```markdown - [x] **Task 1.1**: Set up database schema `abc1234` - [~] **Task 1.2**: Implement user model - [ ] **Task 1.3**: Add validation logic - [!] **Task 1.4**: Integrate auth service (blocked: waiting for API key) - [-] **Task 1.5**: Legacy migration (skipped: not needed) ``` ## Track Registry (tracks.md) Format ```markdown # Track Registry ## Active Tracks | Track ID | Type | Status | Phase | Started | Assignee | | ------------------------------------------------ | ------- | ----------- | ----- | ---------- | ---------- | | [user-auth_20250115](tracks/user-auth_20250115/) | feature | in-progress | 2/3 | 2025-01-15 | @developer | | [fix-login_20250114](tracks/fix-login_20250114/) | bug | pending | 0/2 | 2025-01-14 | - | ## Completed Tracks | Track ID | Type | Completed | Duration | | ---------------------------------------------- | ----- | ---------- | -------- | | [setup-ci_20250110](tracks/setup-ci_20250110/) | chore | 2025-01-12 | 2 days | ## Archived Tracks | Track ID | Reason | Archived | | ---------------------------------------------------- | ---------- | ---------- | | [old-feature_20241201](tracks/old-feature_20241201/) | Superseded | 2025-01-05 | ``` ## Metadata (metadata.json) Fields ```json { "id": "user-auth_20250115", "title": "User Authentication System", "type": "feature", "status": "in-progress", "priority": "high", "created": "2025-01-15T10:30:00Z", "updated": "2025-01-15T14:45:00Z", "started": "2025-01-15T11:00:00Z", "completed": null, "assignee": "@developer", "phases": { "total": 3, "current": 2, "completed": 1 }, "tasks": { "total": 12, "completed": 5, "in_progress": 1, "pending": 6 }, "checkpoints": [ { "phase": 1, "sha": "abc1234", "date": "2025-01-15T13:00:00Z" } ], "dependencies": [], "tags": ["auth", "security"] } ``` ## Track Operations ### Creating a Track 1. Run `/conductor:new-track` 2. Answer interactive questions 3. Review generated spec.md 4. Review generated plan.md 5. Confirm track creation ### Starting Implementation 1. Read spec.md and plan.md 2. Verify context artifacts are current 3. Mark first task as `[~]` 4. Begin TDD workflow ### Completing a Phase 1. Ensure all phase tasks are `[x]` 2. Complete verification tasks 3. Wait for checkpoint approval 4. Record checkpoint SHA 5. Proceed to next phase ### Completing a Track 1. Verify all phases complete 2. Verify all acceptance criteria met 3. Update product.md if needed 4. Mark track completed in tracks.md 5. Update metadata.json ### Reverting a Track 1. Run `/conductor:revert` 2. Select track to revert 3. Choose granularity (track/phase/task) 4. Confirm revert operation 5. Update status markers ## Handling Track Dependencies ### Identifying Dependencies During track creation, identify: - **Hard dependencies**: Must complete before this track can start - **Soft dependencies**: Can proceed in parallel but may affect integration - **External dependencies**: Third-party services, APIs, or team decisions ### Documenting Dependencies In spec.md, list dependencies with: - Dependency type (hard/soft/external) - Current status (available/pending/blocked) - Resolution path (what needs to happen) ### Managing Blocked Tracks When a track is blocked: 1. Mark blocked tasks with `[!]` and reason 2. Update tracks.md status 3. Document blocker in metadata.json 4. Consider creating dependency track if needed ## Track Sizing Guidelines ### Right-Sized Tracks Aim for tracks that: - Complete in 1-5 days of work - Have 2-4 phases - Contain 8-20 tasks total - Deliver a coherent, testable unit ### Too Large Signs a track is too large: - More than 5 phases - More than 25 tasks - Multiple unrelated features - Estimated duration > 1 week Solution: Split into multiple tracks with clear boundaries. ### Too Small Signs a track is too small: - Single phase with 1-2 tasks - No meaningful verification needed - Could be a sub-task of another track - Less than a few hours of work Solution: Combine with related work or handle as part of existing track. ## Specification Quality Checklist Before finalizing spec.md, verify: ### Requirements Quality - [ ] Each requirement has clear acceptance criteria - [ ] Requirements are testable - [ ] Requirements are independent (can verify separately) - [ ] No ambiguous language ("should be fast" β†’ "response < 200ms") ### Scope Clarity - [ ] In-scope items are specific - [ ] Out-of-scope items prevent scope creep - [ ] Boundaries are clear to implementer ### Dependencies Identified - [ ] All internal dependencies listed - [ ] External dependencies have owners/contacts - [ ] Dependency status is current ### Risks Addressed - [ ] Major risks identified - [ ] Impact assessment realistic - [ ] Mitigations are actionable ## Plan Quality Checklist Before starting implementation, verify plan.md: ### Task Quality - [ ] Tasks are atomic (one logical action) - [ ] Tasks are independently verifiable - [ ] Task descriptions are clear - [ ] Sub-tasks provide helpful detail ### Phase Organization - [ ] Phases group related tasks - [ ] Each phase delivers something testable - [ ] Verification tasks after each phase - [ ] Phases build on each other logically ### Completeness - [ ] All spec requirements have corresponding tasks - [ ] Documentation tasks included - [ ] Testing tasks included - [ ] Integration tasks included ## Common Track Patterns ### Feature Track Pattern ``` Phase 1: Foundation - Data models - Database migrations - Basic API structure Phase 2: Core Logic - Business logic implementation - Input validation - Error handling Phase 3: Integration - UI integration - API documentation - End-to-end tests ``` ### Bug Fix Track Pattern ``` Phase 1: Reproduction - Write failing test capturing bug - Document reproduction steps Phase 2: Fix - Implement fix - Verify test passes - Check for regressions Phase 3: Verification - Manual verification - Update documentation if needed ``` ### Refactor Track Pattern ``` Phase 1: Preparation - Add characterization tests - Document current behavior Phase 2: Refactoring - Apply changes incrementally - Maintain green tests throughout Phase 3: Cleanup - Remove dead code - Update documentation ``` ## Best Practices 1. **One track, one concern**: Keep tracks focused on a single logical change 2. **Small phases**: Break work into phases of 3-5 tasks maximum 3. **Verification after phases**: Always include verification tasks 4. **Update markers immediately**: Mark task status as you work 5. **Record SHAs**: Always note commit SHAs for completed tasks 6. **Review specs before planning**: Ensure spec is complete before creating plan 7. **Link dependencies**: Explicitly note track dependencies 8. **Archive, don't delete**: Preserve completed tracks for reference 9. **Size appropriately**: Keep tracks between 1-5 days of work 10. **Clear acceptance criteria**: Every requirement must be testable
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

workflow-patterns

Use this skill when implementing tasks according to Conductor's TDD

coding
⭐1
# Workflow Patterns Guide for implementing tasks using Conductor's TDD workflow, managing phase checkpoints, handling git commits, and executing the verification protocol that ensures quality throughout implementation. ## When to Use This Skill - Implementing tasks from a track's plan.md - Following TDD red-green-refactor cycle - Completing phase checkpoints - Managing git commits and notes - Understanding quality assurance gates - Handling verification protocols - Recording progress in plan files ## TDD Task Lifecycle Follow these 11 steps for each task: ### Step 1: Select Next Task Read plan.md and identify the next pending `[ ]` task. Select tasks in order within the current phase. Do not skip ahead to later phases. ### Step 2: Mark as In Progress Update plan.md to mark the task as `[~]`: ```markdown - [~] **Task 2.1**: Implement user validation ``` Commit this status change separately from implementation. ### Step 3: RED - Write Failing Tests Write tests that define the expected behavior before writing implementation: - Create test file if needed - Write test cases covering happy path - Write test cases covering edge cases - Write test cases covering error conditions - Run tests - they should FAIL Example: ```python def test_validate_user_email_valid(): user = User(email="test@example.com") assert user.validate_email() is True def test_validate_user_email_invalid(): user = User(email="invalid") assert user.validate_email() is False ``` ### Step 4: GREEN - Implement Minimum Code Write the minimum code necessary to make tests pass: - Focus on making tests green, not perfection - Avoid premature optimization - Keep implementation simple - Run tests - they should PASS ### Step 5: REFACTOR - Improve Clarity With green tests, improve the code: - Extract common patterns - Improve naming - Remove duplication - Simplify logic - Run tests after each change - they should remain GREEN ### Step 6: Verify Coverage Check test coverage meets the 80% target: ```bash pytest --cov=module --cov-report=term-missing ``` If coverage is below 80%: - Identify uncovered lines - Add tests for missing paths - Re-run coverage check ### Step 7: Document Deviations If implementation deviated from plan or introduced new dependencies: - Update tech-stack.md with new dependencies - Note deviations in plan.md task comments - Update spec.md if requirements changed ### Step 8: Commit Implementation Create a focused commit for the task: ```bash git add -A git commit -m "feat(user): implement email validation - Add validate_email method to User class - Handle empty and malformed emails - Add comprehensive test coverage Task: 2.1 Track: user-auth_20250115" ``` Commit message format: - Type: feat, fix, refactor, test, docs, chore - Scope: affected module or component - Summary: imperative, present tense - Body: bullet points of changes - Footer: task and track references ### Step 9: Attach Git Notes Add rich task summary as git note: ```bash git notes add -m "Task 2.1: Implement user validation Summary: - Added email validation using regex pattern - Handles edge cases: empty, no @, no domain - Coverage: 94% on validation module Files changed: - src/models/user.py (modified) - tests/test_user.py (modified) Decisions: - Used simple regex over email-validator library - Reason: No external dependency for basic validation" ``` ### Step 10: Update Plan with SHA Update plan.md to mark task complete with commit SHA: ```markdown - [x] **Task 2.1**: Implement user validation `abc1234` ``` ### Step 11: Commit Plan Update Commit the plan status update: ```bash git add conductor/tracks/*/plan.md git commit -m "docs: update plan - task 2.1 complete Track: user-auth_20250115" ``` ## Phase Completion Protocol When all tasks in a phase are complete, execute the verification protocol: ### Identify Changed Files List all files modified since the last checkpoint: ```bash git diff --name-only <last-checkpoint-sha>..HEAD ``` ### Ensure Test Coverage For each modified file: 1. Identify corresponding test file 2. Verify tests exist for new/changed code 3. Run coverage for modified modules 4. Add tests if coverage < 80% ### Run Full Test Suite Execute complete test suite: ```bash pytest -v --tb=short ``` All tests must pass before proceeding. ### Generate Manual Verification Steps Create checklist of manual verifications: ```markdown ## Phase 1 Verification Checklist - [ ] User can register with valid email - [ ] Invalid email shows appropriate error - [ ] Database stores user correctly - [ ] API returns expected response codes ``` ### WAIT for User Approval Present verification checklist to user: ``` Phase 1 complete. Please verify: 1. [ ] Test suite passes (automated) 2. [ ] Coverage meets target (automated) 3. [ ] Manual verification items (requires human) Respond with 'approved' to continue, or note issues. ``` Do NOT proceed without explicit approval. ### Create Checkpoint Commit After approval, create checkpoint commit: ```bash git add -A git commit -m "checkpoint: phase 1 complete - user-auth_20250115 Verified: - All tests passing - Coverage: 87% - Manual verification approved Phase 1 tasks: - [x] Task 1.1: Setup database schema - [x] Task 1.2: Implement user model - [x] Task 1.3: Add validation logic" ``` ### Record Checkpoint SHA Update plan.md checkpoints table: ```markdown ## Checkpoints | Phase | Checkpoint SHA | Date | Status | | ------- | -------------- | ---------- | -------- | | Phase 1 | def5678 | 2025-01-15 | verified | | Phase 2 | | | pending | ``` ## Quality Assurance Gates Before marking any task complete, verify these gates: ### Passing Tests - All existing tests pass - New tests pass - No test regressions ### Coverage >= 80% - New code has 80%+ coverage - Overall project coverage maintained - Critical paths fully covered ### Style Compliance - Code follows style guides - Linting passes - Formatting correct ### Documentation - Public APIs documented - Complex logic explained - README updated if needed ### Type Safety - Type hints present (if applicable) - Type checker passes - No type: ignore without reason ### No Linting Errors - Zero linter errors - Warnings addressed or justified - Static analysis clean ### Mobile Compatibility If applicable: - Responsive design verified - Touch interactions work - Performance acceptable ### Security Audit - No secrets in code - Input validation present - Authentication/authorization correct - Dependencies vulnerability-free ## Git Integration ### Commit Message Format ``` <type>(<scope>): <subject> <body> <footer> ``` Types: - `feat`: New feature - `fix`: Bug fix - `refactor`: Code change without feature/fix - `test`: Adding tests - `docs`: Documentation - `chore`: Maintenance ### Git Notes for Rich Summaries Attach detailed notes to commits: ```bash git notes add -m "<detailed summary>" ``` View notes: ```bash git log --show-notes ``` Benefits: - Preserves context without cluttering commit message - Enables semantic queries across commits - Supports track-based operations ### SHA Recording in plan.md Always record the commit SHA when completing tasks: ```markdown - [x] **Task 1.1**: Setup schema `abc1234` - [x] **Task 1.2**: Add model `def5678` ``` This enables: - Traceability from plan to code - Semantic revert operations - Progress auditing ## Verification Checkpoints ### Why Checkpoints Matter Checkpoints create restore points for semantic reversion: - Revert to end of any phase - Maintain logical code state - Enable safe experimentation ### When to Create Checkpoints Create checkpoint after: - All phase tasks complete - All phase verifications pass - User approval received ### Checkpoint Commit Content Include in checkpoint commit: - All uncommitted changes - Updated plan.md - Updated metadata.json - Any documentation updates ### How to Use Checkpoints For reverting: ```bash # Revert to end of Phase 1 git revert --no-commit <phase-2-commits>... git commit -m "revert: rollback to phase 1 checkpoint" ``` For review: ```bash # See what changed in Phase 2 git diff <phase-1-sha>..<phase-2-sha> ``` ## Handling Deviations During implementation, deviations from the plan may occur. Handle them systematically: ### Types of Deviations **Scope Addition** Discovered requirement not in original spec. - Document in spec.md as new requirement - Add tasks to plan.md - Note addition in task comments **Scope Reduction** Feature deemed unnecessary during implementation. - Mark tasks as `[-]` (skipped) with reason - Update spec.md scope section - Document decision rationale **Technical Deviation** Different implementation approach than planned. - Note deviation in task completion comment - Update tech-stack.md if dependencies changed - Document why original approach was unsuitable **Requirement Change** Understanding of requirement changes during work. - Update spec.md with corrected requirement - Adjust plan.md tasks if needed - Re-verify acceptance criteria ### Deviation Documentation Format When completing a task with deviation: ```markdown - [x] **Task 2.1**: Implement validation `abc1234` - DEVIATION: Used library instead of custom code - Reason: Better edge case handling - Impact: Added email-validator to dependencies ``` ## Error Recovery ### Failed Tests After GREEN If tests fail after reaching GREEN: 1. Do NOT proceed to REFACTOR 2. Identify which test started failing 3. Check if refactoring broke something 4. Revert to last known GREEN state 5. Re-approach the implementation ### Checkpoint Rejection If user rejects a checkpoint: 1. Note rejection reason in plan.md 2. Create tasks to address issues 3. Complete remediation tasks 4. Request checkpoint approval again ### Blocked by Dependency If task cannot proceed: 1. Mark task as `[!]` with blocker description 2. Check if other tasks can proceed 3. Document expected resolution timeline 4. Consider creating dependency resolution track ## TDD Variations by Task Type ### Data Model Tasks ``` RED: Write test for model creation and validation GREEN: Implement model class with fields REFACTOR: Add computed properties, improve types ``` ### API Endpoint Tasks ``` RED: Write test for request/response contract GREEN: Implement endpoint handler REFACTOR: Extract validation, improve error handling ``` ### Integration Tasks ``` RED: Write test for component interaction GREEN: Wire components together REFACTOR: Improve error propagation, add logging ``` ### Refactoring Tasks ``` RED: Add characterization tests for current behavior GREEN: Apply refactoring (tests should stay green) REFACTOR: Clean up any introduced complexity ``` ## Working with Existing Tests When modifying code with existing tests: ### Extend, Don't Replace - Keep existing tests passing - Add new tests for new behavior - Update tests only when requirements change ### Test Migration When refactoring changes test structure: 1. Run existing tests (should pass) 2. Add new tests for refactored code 3. Migrate test cases to new structure 4. Remove old tests only after new tests pass ### Regression Prevention After any change: 1. Run full test suite 2. Check for unexpected failures 3. Investigate any new failures 4. Fix regressions before proceeding ## Checkpoint Verification Details ### Automated Verification Run before requesting approval: ```bash # Test suite pytest -v --tb=short # Coverage pytest --cov=src --cov-report=term-missing # Linting ruff check src/ tests/ # Type checking (if applicable) mypy src/ ``` ### Manual Verification Guidance For manual items, provide specific instructions: ```markdown ## Manual Verification Steps ### User Registration 1. Navigate to /register 2. Enter valid email: test@example.com 3. Enter password meeting requirements 4. Click Submit 5. Verify success message appears 6. Verify user appears in database ### Error Handling 1. Enter invalid email: "notanemail" 2. Verify error message shows 3. Verify form retains other entered data ``` ## Performance Considerations ### Test Suite Performance Keep test suite fast: - Use fixtures to avoid redundant setup - Mock slow external calls - Run subset during development, full suite at checkpoints ### Commit Performance Keep commits atomic: - One logical change per commit - Complete thought, not work-in-progress - Tests should pass after every commit ## Best Practices 1. **Never skip RED**: Always write failing tests first 2. **Small commits**: One logical change per commit 3. **Immediate updates**: Update plan.md right after task completion 4. **Wait for approval**: Never skip checkpoint verification 5. **Rich git notes**: Include context that helps future understanding 6. **Coverage discipline**: Don't accept coverage below target 7. **Quality gates**: Check all gates before marking complete 8. **Sequential phases**: Complete phases in order 9. **Document deviations**: Note any changes from original plan 10. **Clean state**: Each commit should leave code in working state 11. **Fast feedback**: Run relevant tests frequently during development 12. **Clear blockers**: Address blockers promptly, don't work around them
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

auth-implementation-patterns

Master authentication and authorization patterns including JWT,

coding
⭐1
# Authentication & Authorization Implementation Patterns Build secure, scalable authentication and authorization systems using industry-standard patterns and modern best practices. ## When to Use This Skill - Implementing user authentication systems - Securing REST or GraphQL APIs - Adding OAuth2/social login - Implementing role-based access control (RBAC) - Designing session management - Migrating authentication systems - Debugging auth issues - Implementing SSO or multi-tenancy ## Core Concepts ### 1. Authentication vs Authorization **Authentication (AuthN)**: Who are you? - Verifying identity (username/password, OAuth, biometrics) - Issuing credentials (sessions, tokens) - Managing login/logout **Authorization (AuthZ)**: What can you do? - Permission checking - Role-based access control (RBAC) - Resource ownership validation - Policy enforcement ### 2. Authentication Strategies **Session-Based:** - Server stores session state - Session ID in cookie - Traditional, simple, stateful **Token-Based (JWT):** - Stateless, self-contained - Scales horizontally - Can store claims **OAuth2/OpenID Connect:** - Delegate authentication - Social login (Google, GitHub) - Enterprise SSO ## JWT Authentication ### Pattern 1: JWT Implementation ```typescript // JWT structure: header.payload.signature import jwt from "jsonwebtoken"; import { Request, Response, NextFunction } from "express"; interface JWTPayload { userId: string; email: string; role: string; iat: number; exp: number; } // Generate JWT function generateTokens(userId: string, email: string, role: string) { const accessToken = jwt.sign( { userId, email, role }, process.env.JWT_SECRET!, { expiresIn: "15m" }, // Short-lived ); const refreshToken = jwt.sign( { userId }, process.env.JWT_REFRESH_SECRET!, { expiresIn: "7d" }, // Long-lived ); return { accessToken, refreshToken }; } // Verify JWT function verifyToken(token: string): JWTPayload { try { return jwt.verify(token, process.env.JWT_SECRET!) as JWTPayload; } catch (error) { if (error instanceof jwt.TokenExpiredError) { throw new Error("Token expired"); } if (error instanceof jwt.JsonWebTokenError) { throw new Error("Invalid token"); } throw error; } } // Middleware function authenticate(req: Request, res: Response, next: NextFunction) { const authHeader = req.headers.authorization; if (!authHeader?.startsWith("Bearer ")) { return res.status(401).json({ error: "No token provided" }); } const token = authHeader.substring(7); try { const payload = verifyToken(token); req.user = payload; // Attach user to request next(); } catch (error) { return res.status(401).json({ error: "Invalid token" }); } } // Usage app.get("/api/profile", authenticate, (req, res) => { res.json({ user: req.user }); }); ``` ### Pattern 2: Refresh Token Flow ```typescript interface StoredRefreshToken { token: string; userId: string; expiresAt: Date; createdAt: Date; } class RefreshTokenService { // Store refresh token in database async storeRefreshToken(userId: string, refreshToken: string) { const expiresAt = new Date(Date.now() + 7 * 24 * 60 * 60 * 1000); await db.refreshTokens.create({ token: await hash(refreshToken), // Hash before storing userId, expiresAt, }); } // Refresh access token async refreshAccessToken(refreshToken: string) { // Verify refresh token let payload; try { payload = jwt.verify(refreshToken, process.env.JWT_REFRESH_SECRET!) as { userId: string; }; } catch { throw new Error("Invalid refresh token"); } // Check if token exists in database const storedToken = await db.refreshTokens.findOne({ where: { token: await hash(refreshToken), userId: payload.userId, expiresAt: { $gt: new Date() }, }, }); if (!storedToken) { throw new Error("Refresh token not found or expired"); } // Get user const user = await db.users.findById(payload.userId); if (!user) { throw new Error("User not found"); } // Generate new access token const accessToken = jwt.sign( { userId: user.id, email: user.email, role: user.role }, process.env.JWT_SECRET!, { expiresIn: "15m" }, ); return { accessToken }; } // Revoke refresh token (logout) async revokeRefreshToken(refreshToken: string) { await db.refreshTokens.deleteOne({ token: await hash(refreshToken), }); } // Revoke all user tokens (logout all devices) async revokeAllUserTokens(userId: string) { await db.refreshTokens.deleteMany({ userId }); } } // API endpoints app.post("/api/auth/refresh", async (req, res) => { const { refreshToken } = req.body; try { const { accessToken } = await refreshTokenService.refreshAccessToken(refreshToken); res.json({ accessToken }); } catch (error) { res.status(401).json({ error: "Invalid refresh token" }); } }); app.post("/api/auth/logout", authenticate, async (req, res) => { const { refreshToken } = req.body; await refreshTokenService.revokeRefreshToken(refreshToken); res.json({ message: "Logged out successfully" }); }); ``` ## Session-Based Authentication ### Pattern 1: Express Session ```typescript import session from "express-session"; import RedisStore from "connect-redis"; import { createClient } from "redis"; // Setup Redis for session storage const redisClient = createClient({ url: process.env.REDIS_URL, }); await redisClient.connect(); app.use( session({ store: new RedisStore({ client: redisClient }), secret: process.env.SESSION_SECRET!, resave: false, saveUninitialized: false, cookie: { secure: process.env.NODE_ENV === "production", // HTTPS only httpOnly: true, // No JavaScript access maxAge: 24 * 60 * 60 * 1000, // 24 hours sameSite: "strict", // CSRF protection }, }), ); // Login app.post("/api/auth/login", async (req, res) => { const { email, password } = req.body; const user = await db.users.findOne({ email }); if (!user || !(await verifyPassword(password, user.passwordHash))) { return res.status(401).json({ error: "Invalid credentials" }); } // Store user in session req.session.userId = user.id; req.session.role = user.role; res.json({ user: { id: user.id, email: user.email, role: user.role } }); }); // Session middleware function requireAuth(req: Request, res: Response, next: NextFunction) { if (!req.session.userId) { return res.status(401).json({ error: "Not authenticated" }); } next(); } // Protected route app.get("/api/profile", requireAuth, async (req, res) => { const user = await db.users.findById(req.session.userId); res.json({ user }); }); // Logout app.post("/api/auth/logout", (req, res) => { req.session.destroy((err) => { if (err) { return res.status(500).json({ error: "Logout failed" }); } res.clearCookie("connect.sid"); res.json({ message: "Logged out successfully" }); }); }); ``` ## OAuth2 / Social Login ### Pattern 1: OAuth2 with Passport.js ```typescript import passport from "passport"; import { Strategy as GoogleStrategy } from "passport-google-oauth20"; import { Strategy as GitHubStrategy } from "passport-github2"; // Google OAuth passport.use( new GoogleStrategy( { clientID: process.env.GOOGLE_CLIENT_ID!, clientSecret: process.env.GOOGLE_CLIENT_SECRET!, callbackURL: "/api/auth/google/callback", }, async (accessToken, refreshToken, profile, done) => { try { // Find or create user let user = await db.users.findOne({ googleId: profile.id, }); if (!user) { user = await db.users.create({ googleId: profile.id, email: profile.emails?.[0]?.value, name: profile.displayName, avatar: profile.photos?.[0]?.value, }); } return done(null, user); } catch (error) { return done(error, undefined); } }, ), ); // Routes app.get( "/api/auth/google", passport.authenticate("google", { scope: ["profile", "email"], }), ); app.get( "/api/auth/google/callback", passport.authenticate("google", { session: false }), (req, res) => { // Generate JWT const tokens = generateTokens(req.user.id, req.user.email, req.user.role); // Redirect to frontend with token res.redirect( `${process.env.FRONTEND_URL}/auth/callback?token=${tokens.accessToken}`, ); }, ); ``` ## Authorization Patterns ### Pattern 1: Role-Based Access Control (RBAC) ```typescript enum Role { USER = "user", MODERATOR = "moderator", ADMIN = "admin", } const roleHierarchy: Record<Role, Role[]> = { [Role.ADMIN]: [Role.ADMIN, Role.MODERATOR, Role.USER], [Role.MODERATOR]: [Role.MODERATOR, Role.USER], [Role.USER]: [Role.USER], }; function hasRole(userRole: Role, requiredRole: Role): boolean { return roleHierarchy[userRole].includes(requiredRole); } // Middleware function requireRole(...roles: Role[]) { return (req: Request, res: Response, next: NextFunction) => { if (!req.user) { return res.status(401).json({ error: "Not authenticated" }); } if (!roles.some((role) => hasRole(req.user.role, role))) { return res.status(403).json({ error: "Insufficient permissions" }); } next(); }; } // Usage app.delete( "/api/users/:id", authenticate, requireRole(Role.ADMIN), async (req, res) => { // Only admins can delete users await db.users.delete(req.params.id); res.json({ message: "User deleted" }); }, ); ``` ### Pattern 2: Permission-Based Access Control ```typescript enum Permission { READ_USERS = "read:users", WRITE_USERS = "write:users", DELETE_USERS = "delete:users", READ_POSTS = "read:posts", WRITE_POSTS = "write:posts", } const rolePermissions: Record<Role, Permission[]> = { [Role.USER]: [Permission.READ_POSTS, Permission.WRITE_POSTS], [Role.MODERATOR]: [ Permission.READ_POSTS, Permission.WRITE_POSTS, Permission.READ_USERS, ], [Role.ADMIN]: Object.values(Permission), }; function hasPermission(userRole: Role, permission: Permission): boolean { return rolePermissions[userRole]?.includes(permission) ?? false; } function requirePermission(...permissions: Permission[]) { return (req: Request, res: Response, next: NextFunction) => { if (!req.user) { return res.status(401).json({ error: "Not authenticated" }); } const hasAllPermissions = permissions.every((permission) => hasPermission(req.user.role, permission), ); if (!hasAllPermissions) { return res.status(403).json({ error: "Insufficient permissions" }); } next(); }; } // Usage app.get( "/api/users", authenticate, requirePermission(Permission.READ_USERS), async (req, res) => { const users = await db.users.findAll(); res.json({ users }); }, ); ``` ### Pattern 3: Resource Ownership ```typescript // Check if user owns resource async function requireOwnership( resourceType: "post" | "comment", resourceIdParam: string = "id", ) { return async (req: Request, res: Response, next: NextFunction) => { if (!req.user) { return res.status(401).json({ error: "Not authenticated" }); } const resourceId = req.params[resourceIdParam]; // Admins can access anything if (req.user.role === Role.ADMIN) { return next(); } // Check ownership let resource; if (resourceType === "post") { resource = await db.posts.findById(resourceId); } else if (resourceType === "comment") { resource = await db.comments.findById(resourceId); } if (!resource) { return res.status(404).json({ error: "Resource not found" }); } if (resource.userId !== req.user.userId) { return res.status(403).json({ error: "Not authorized" }); } next(); }; } // Usage app.put( "/api/posts/:id", authenticate, requireOwnership("post"), async (req, res) => { // User can only update their own posts const post = await db.posts.update(req.params.id, req.body); res.json({ post }); }, ); ``` ## Security Best Practices ### Pattern 1: Password Security ```typescript import bcrypt from "bcrypt"; import { z } from "zod"; // Password validation schema const passwordSchema = z .string() .min(12, "Password must be at least 12 characters") .regex(/[A-Z]/, "Password must contain uppercase letter") .regex(/[a-z]/, "Password must contain lowercase letter") .regex(/[0-9]/, "Password must contain number") .regex(/[^A-Za-z0-9]/, "Password must contain special character"); // Hash password async function hashPassword(password: string): Promise<string> { const saltRounds = 12; // 2^12 iterations return bcrypt.hash(password, saltRounds); } // Verify password async function verifyPassword( password: string, hash: string, ): Promise<boolean> { return bcrypt.compare(password, hash); } // Registration with password validation app.post("/api/auth/register", async (req, res) => { try { const { email, password } = req.body; // Validate password passwordSchema.parse(password); // Check if user exists const existingUser = await db.users.findOne({ email }); if (existingUser) { return res.status(400).json({ error: "Email already registered" }); } // Hash password const passwordHash = await hashPassword(password); // Create user const user = await db.users.create({ email, passwordHash, }); // Generate tokens const tokens = generateTokens(user.id, user.email, user.role); res.status(201).json({ user: { id: user.id, email: user.email }, ...tokens, }); } catch (error) { if (error instanceof z.ZodError) { return res.status(400).json({ error: error.errors[0].message }); } res.status(500).json({ error: "Registration failed" }); } }); ``` ### Pattern 2: Rate Limiting ```typescript import rateLimit from "express-rate-limit"; import RedisStore from "rate-limit-redis"; // Login rate limiter const loginLimiter = rateLimit({ store: new RedisStore({ client: redisClient }), windowMs: 15 * 60 * 1000, // 15 minutes max: 5, // 5 attempts message: "Too many login attempts, please try again later", standardHeaders: true, legacyHeaders: false, }); // API rate limiter const apiLimiter = rateLimit({ windowMs: 60 * 1000, // 1 minute max: 100, // 100 requests per minute standardHeaders: true, }); // Apply to routes app.post("/api/auth/login", loginLimiter, async (req, res) => { // Login logic }); app.use("/api/", apiLimiter); ``` ## Best Practices 1. **Never Store Plain Passwords**: Always hash with bcrypt/argon2 2. **Use HTTPS**: Encrypt data in transit 3. **Short-Lived Access Tokens**: 15-30 minutes max 4. **Secure Cookies**: httpOnly, secure, sameSite flags 5. **Validate All Input**: Email format, password strength 6. **Rate Limit Auth Endpoints**: Prevent brute force attacks 7. **Implement CSRF Protection**: For session-based auth 8. **Rotate Secrets Regularly**: JWT secrets, session secrets 9. **Log Security Events**: Login attempts, failed auth 10. **Use MFA When Possible**: Extra security layer ## Common Pitfalls - **Weak Passwords**: Enforce strong password policies - **JWT in localStorage**: Vulnerable to XSS, use httpOnly cookies - **No Token Expiration**: Tokens should expire - **Client-Side Auth Checks Only**: Always validate server-side - **Insecure Password Reset**: Use secure tokens with expiration - **No Rate Limiting**: Vulnerable to brute force - **Trusting Client Data**: Always validate on server ## Resources - **references/jwt-best-practices.md**: JWT implementation guide - **references/oauth2-flows.md**: OAuth2 flow diagrams and examples - **references/session-security.md**: Secure session management - **assets/auth-security-checklist.md**: Security review checklist - **assets/password-policy-template.md**: Password requirements template - **scripts/token-validator.ts**: JWT validation utility
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

airflow-dag-patterns

Build production Apache Airflow DAGs with best practices for

data
⭐1
# Apache Airflow DAG Patterns Production-ready patterns for Apache Airflow including DAG design, operators, sensors, testing, and deployment strategies. ## When to Use This Skill - Creating data pipeline orchestration with Airflow - Designing DAG structures and dependencies - Implementing custom operators and sensors - Testing Airflow DAGs locally - Setting up Airflow in production - Debugging failed DAG runs ## Core Concepts ### 1. DAG Design Principles | Principle | Description | | --------------- | ----------------------------------- | | **Idempotent** | Running twice produces same result | | **Atomic** | Tasks succeed or fail completely | | **Incremental** | Process only new/changed data | | **Observable** | Logs, metrics, alerts at every step | ### 2. Task Dependencies ```python # Linear task1 >> task2 >> task3 # Fan-out task1 >> [task2, task3, task4] # Fan-in [task1, task2, task3] >> task4 # Complex task1 >> task2 >> task4 task1 >> task3 >> task4 ``` ## Quick Start ```python # dags/example_dag.py from datetime import datetime, timedelta from airflow import DAG from airflow.operators.python import PythonOperator from airflow.operators.empty import EmptyOperator default_args = { 'owner': 'data-team', 'depends_on_past': False, 'email_on_failure': True, 'email_on_retry': False, 'retries': 3, 'retry_delay': timedelta(minutes=5), 'retry_exponential_backoff': True, 'max_retry_delay': timedelta(hours=1), } with DAG( dag_id='example_etl', default_args=default_args, description='Example ETL pipeline', schedule='0 6 * * *', # Daily at 6 AM start_date=datetime(2024, 1, 1), catchup=False, tags=['etl', 'example'], max_active_runs=1, ) as dag: start = EmptyOperator(task_id='start') def extract_data(**context): execution_date = context['ds'] # Extract logic here return {'records': 1000} extract = PythonOperator( task_id='extract', python_callable=extract_data, ) end = EmptyOperator(task_id='end') start >> extract >> end ``` ## Patterns ### Pattern 1: TaskFlow API (Airflow 2.0+) ```python # dags/taskflow_example.py from datetime import datetime from airflow.decorators import dag, task from airflow.models import Variable @dag( dag_id='taskflow_etl', schedule='@daily', start_date=datetime(2024, 1, 1), catchup=False, tags=['etl', 'taskflow'], ) def taskflow_etl(): """ETL pipeline using TaskFlow API""" @task() def extract(source: str) -> dict: """Extract data from source""" import pandas as pd df = pd.read_csv(f's3://bucket/{source}/{{ ds }}.csv') return {'data': df.to_dict(), 'rows': len(df)} @task() def transform(extracted: dict) -> dict: """Transform extracted data""" import pandas as pd df = pd.DataFrame(extracted['data']) df['processed_at'] = datetime.now() df = df.dropna() return {'data': df.to_dict(), 'rows': len(df)} @task() def load(transformed: dict, target: str): """Load data to target""" import pandas as pd df = pd.DataFrame(transformed['data']) df.to_parquet(f's3://bucket/{target}/{{ ds }}.parquet') return transformed['rows'] @task() def notify(rows_loaded: int): """Send notification""" print(f'Loaded {rows_loaded} rows') # Define dependencies with XCom passing extracted = extract(source='raw_data') transformed = transform(extracted) loaded = load(transformed, target='processed_data') notify(loaded) # Instantiate the DAG taskflow_etl() ``` ### Pattern 2: Dynamic DAG Generation ```python # dags/dynamic_dag_factory.py from datetime import datetime, timedelta from airflow import DAG from airflow.operators.python import PythonOperator from airflow.models import Variable import json # Configuration for multiple similar pipelines PIPELINE_CONFIGS = [ {'name': 'customers', 'schedule': '@daily', 'source': 's3://raw/customers'}, {'name': 'orders', 'schedule': '@hourly', 'source': 's3://raw/orders'}, {'name': 'products', 'schedule': '@weekly', 'source': 's3://raw/products'}, ] def create_dag(config: dict) -> DAG: """Factory function to create DAGs from config""" dag_id = f"etl_{config['name']}" default_args = { 'owner': 'data-team', 'retries': 3, 'retry_delay': timedelta(minutes=5), } dag = DAG( dag_id=dag_id, default_args=default_args, schedule=config['schedule'], start_date=datetime(2024, 1, 1), catchup=False, tags=['etl', 'dynamic', config['name']], ) with dag: def extract_fn(source, **context): print(f"Extracting from {source} for {context['ds']}") def transform_fn(**context): print(f"Transforming data for {context['ds']}") def load_fn(table_name, **context): print(f"Loading to {table_name} for {context['ds']}") extract = PythonOperator( task_id='extract', python_callable=extract_fn, op_kwargs={'source': config['source']}, ) transform = PythonOperator( task_id='transform', python_callable=transform_fn, ) load = PythonOperator( task_id='load', python_callable=load_fn, op_kwargs={'table_name': config['name']}, ) extract >> transform >> load return dag # Generate DAGs for config in PIPELINE_CONFIGS: globals()[f"dag_{config['name']}"] = create_dag(config) ``` ### Pattern 3: Branching and Conditional Logic ```python # dags/branching_example.py from airflow.decorators import dag, task from airflow.operators.python import BranchPythonOperator from airflow.operators.empty import EmptyOperator from airflow.utils.trigger_rule import TriggerRule @dag( dag_id='branching_pipeline', schedule='@daily', start_date=datetime(2024, 1, 1), catchup=False, ) def branching_pipeline(): @task() def check_data_quality() -> dict: """Check data quality and return metrics""" quality_score = 0.95 # Simulated return {'score': quality_score, 'rows': 10000} def choose_branch(**context) -> str: """Determine which branch to execute""" ti = context['ti'] metrics = ti.xcom_pull(task_ids='check_data_quality') if metrics['score'] >= 0.9: return 'high_quality_path' elif metrics['score'] >= 0.7: return 'medium_quality_path' else: return 'low_quality_path' quality_check = check_data_quality() branch = BranchPythonOperator( task_id='branch', python_callable=choose_branch, ) high_quality = EmptyOperator(task_id='high_quality_path') medium_quality = EmptyOperator(task_id='medium_quality_path') low_quality = EmptyOperator(task_id='low_quality_path') # Join point - runs after any branch completes join = EmptyOperator( task_id='join', trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS, ) quality_check >> branch >> [high_quality, medium_quality, low_quality] >> join branching_pipeline() ``` ### Pattern 4: Sensors and External Dependencies ```python # dags/sensor_patterns.py from datetime import datetime, timedelta from airflow import DAG from airflow.sensors.filesystem import FileSensor from airflow.providers.amazon.aws.sensors.s3 import S3KeySensor from airflow.sensors.external_task import ExternalTaskSensor from airflow.operators.python import PythonOperator with DAG( dag_id='sensor_example', schedule='@daily', start_date=datetime(2024, 1, 1), catchup=False, ) as dag: # Wait for file on S3 wait_for_file = S3KeySensor( task_id='wait_for_s3_file', bucket_name='data-lake', bucket_key='raw/{{ ds }}/data.parquet', aws_conn_id='aws_default', timeout=60 * 60 * 2, # 2 hours poke_interval=60 * 5, # Check every 5 minutes mode='reschedule', # Free up worker slot while waiting ) # Wait for another DAG to complete wait_for_upstream = ExternalTaskSensor( task_id='wait_for_upstream_dag', external_dag_id='upstream_etl', external_task_id='final_task', execution_date_fn=lambda dt: dt, # Same execution date timeout=60 * 60 * 3, mode='reschedule', ) # Custom sensor using @task.sensor decorator @task.sensor(poke_interval=60, timeout=3600, mode='reschedule') def wait_for_api() -> PokeReturnValue: """Custom sensor for API availability""" import requests response = requests.get('https://api.example.com/health') is_done = response.status_code == 200 return PokeReturnValue(is_done=is_done, xcom_value=response.json()) api_ready = wait_for_api() def process_data(**context): api_result = context['ti'].xcom_pull(task_ids='wait_for_api') print(f"API returned: {api_result}") process = PythonOperator( task_id='process', python_callable=process_data, ) [wait_for_file, wait_for_upstream, api_ready] >> process ``` ### Pattern 5: Error Handling and Alerts ```python # dags/error_handling.py from datetime import datetime, timedelta from airflow import DAG from airflow.operators.python import PythonOperator from airflow.utils.trigger_rule import TriggerRule from airflow.models import Variable def task_failure_callback(context): """Callback on task failure""" task_instance = context['task_instance'] exception = context.get('exception') # Send to Slack/PagerDuty/etc message = f""" Task Failed! DAG: {task_instance.dag_id} Task: {task_instance.task_id} Execution Date: {context['ds']} Error: {exception} Log URL: {task_instance.log_url} """ # send_slack_alert(message) print(message) def dag_failure_callback(context): """Callback on DAG failure""" # Aggregate failures, send summary pass with DAG( dag_id='error_handling_example', schedule='@daily', start_date=datetime(2024, 1, 1), catchup=False, on_failure_callback=dag_failure_callback, default_args={ 'on_failure_callback': task_failure_callback, 'retries': 3, 'retry_delay': timedelta(minutes=5), }, ) as dag: def might_fail(**context): import random if random.random() < 0.3: raise ValueError("Random failure!") return "Success" risky_task = PythonOperator( task_id='risky_task', python_callable=might_fail, ) def cleanup(**context): """Cleanup runs regardless of upstream failures""" print("Cleaning up...") cleanup_task = PythonOperator( task_id='cleanup', python_callable=cleanup, trigger_rule=TriggerRule.ALL_DONE, # Run even if upstream fails ) def notify_success(**context): """Only runs if all upstream succeeded""" print("All tasks succeeded!") success_notification = PythonOperator( task_id='notify_success', python_callable=notify_success, trigger_rule=TriggerRule.ALL_SUCCESS, ) risky_task >> [cleanup_task, success_notification] ``` ### Pattern 6: Testing DAGs ```python # tests/test_dags.py import pytest from datetime import datetime from airflow.models import DagBag @pytest.fixture def dagbag(): return DagBag(dag_folder='dags/', include_examples=False) def test_dag_loaded(dagbag): """Test that all DAGs load without errors""" assert len(dagbag.import_errors) == 0, f"DAG import errors: {dagbag.import_errors}" def test_dag_structure(dagbag): """Test specific DAG structure""" dag = dagbag.get_dag('example_etl') assert dag is not None assert len(dag.tasks) == 3 assert dag.schedule_interval == '0 6 * * *' def test_task_dependencies(dagbag): """Test task dependencies are correct""" dag = dagbag.get_dag('example_etl') extract_task = dag.get_task('extract') assert 'start' in [t.task_id for t in extract_task.upstream_list] assert 'end' in [t.task_id for t in extract_task.downstream_list] def test_dag_integrity(dagbag): """Test DAG has no cycles and is valid""" for dag_id, dag in dagbag.dags.items(): assert dag.test_cycle() is None, f"Cycle detected in {dag_id}" # Test individual task logic def test_extract_function(): """Unit test for extract function""" from dags.example_dag import extract_data result = extract_data(ds='2024-01-01') assert 'records' in result assert isinstance(result['records'], int) ``` ## Project Structure ``` airflow/ β”œβ”€β”€ dags/ β”‚ β”œβ”€β”€ __init__.py β”‚ β”œβ”€β”€ common/ β”‚ β”‚ β”œβ”€β”€ __init__.py β”‚ β”‚ β”œβ”€β”€ operators.py # Custom operators β”‚ β”‚ β”œβ”€β”€ sensors.py # Custom sensors β”‚ β”‚ └── callbacks.py # Alert callbacks β”‚ β”œβ”€β”€ etl/ β”‚ β”‚ β”œβ”€β”€ customers.py β”‚ β”‚ └── orders.py β”‚ └── ml/ β”‚ └── training.py β”œβ”€β”€ plugins/ β”‚ └── custom_plugin.py β”œβ”€β”€ tests/ β”‚ β”œβ”€β”€ __init__.py β”‚ β”œβ”€β”€ test_dags.py β”‚ └── test_operators.py β”œβ”€β”€ docker-compose.yml └── requirements.txt ``` ## Best Practices ### Do's - **Use TaskFlow API** - Cleaner code, automatic XCom - **Set timeouts** - Prevent zombie tasks - **Use `mode='reschedule'`** - For sensors, free up workers - **Test DAGs** - Unit tests and integration tests - **Idempotent tasks** - Safe to retry ### Don'ts - **Don't use `depends_on_past=True`** - Creates bottlenecks - **Don't hardcode dates** - Use `{{ ds }}` macros - **Don't use global state** - Tasks should be stateless - **Don't skip catchup blindly** - Understand implications - **Don't put heavy logic in DAG file** - Import from modules ## Resources - [Airflow Documentation](https://airflow.apache.org/docs/) - [Astronomer Guides](https://docs.astronomer.io/learn) - [TaskFlow API](https://airflow.apache.org/docs/apache-airflow/stable/tutorial/taskflow.html)
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

e2e-testing-patterns

Master end-to-end testing with Playwright and Cypress to build

coding
⭐1
# E2E Testing Patterns Build reliable, fast, and maintainable end-to-end test suites that provide confidence to ship code quickly and catch regressions before users do. ## When to Use This Skill - Implementing end-to-end test automation - Debugging flaky or unreliable tests - Testing critical user workflows - Setting up CI/CD test pipelines - Testing across multiple browsers - Validating accessibility requirements - Testing responsive designs - Establishing E2E testing standards ## Core Concepts ### 1. E2E Testing Fundamentals **What to Test with E2E:** - Critical user journeys (login, checkout, signup) - Complex interactions (drag-and-drop, multi-step forms) - Cross-browser compatibility - Real API integration - Authentication flows **What NOT to Test with E2E:** - Unit-level logic (use unit tests) - API contracts (use integration tests) - Edge cases (too slow) - Internal implementation details ### 2. Test Philosophy **The Testing Pyramid:** ``` /\ /E2E\ ← Few, focused on critical paths /─────\ /Integr\ ← More, test component interactions /────────\ /Unit Tests\ ← Many, fast, isolated /────────────\ ``` **Best Practices:** - Test user behavior, not implementation - Keep tests independent - Make tests deterministic - Optimize for speed - Use data-testid, not CSS selectors ## Playwright Patterns ### Setup and Configuration ```typescript // playwright.config.ts import { defineConfig, devices } from "@playwright/test"; export default defineConfig({ testDir: "./e2e", timeout: 30000, expect: { timeout: 5000, }, fullyParallel: true, forbidOnly: !!process.env.CI, retries: process.env.CI ? 2 : 0, workers: process.env.CI ? 1 : undefined, reporter: [["html"], ["junit", { outputFile: "results.xml" }]], use: { baseURL: "http://localhost:3000", trace: "on-first-retry", screenshot: "only-on-failure", video: "retain-on-failure", }, projects: [ { name: "chromium", use: { ...devices["Desktop Chrome"] } }, { name: "firefox", use: { ...devices["Desktop Firefox"] } }, { name: "webkit", use: { ...devices["Desktop Safari"] } }, { name: "mobile", use: { ...devices["iPhone 13"] } }, ], }); ``` ### Pattern 1: Page Object Model ```typescript // pages/LoginPage.ts import { Page, Locator } from "@playwright/test"; export class LoginPage { readonly page: Page; readonly emailInput: Locator; readonly passwordInput: Locator; readonly loginButton: Locator; readonly errorMessage: Locator; constructor(page: Page) { this.page = page; this.emailInput = page.getByLabel("Email"); this.passwordInput = page.getByLabel("Password"); this.loginButton = page.getByRole("button", { name: "Login" }); this.errorMessage = page.getByRole("alert"); } async goto() { await this.page.goto("/login"); } async login(email: string, password: string) { await this.emailInput.fill(email); await this.passwordInput.fill(password); await this.loginButton.click(); } async getErrorMessage(): Promise<string> { return (await this.errorMessage.textContent()) ?? ""; } } // Test using Page Object import { test, expect } from "@playwright/test"; import { LoginPage } from "./pages/LoginPage"; test("successful login", async ({ page }) => { const loginPage = new LoginPage(page); await loginPage.goto(); await loginPage.login("user@example.com", "password123"); await expect(page).toHaveURL("/dashboard"); await expect(page.getByRole("heading", { name: "Dashboard" })).toBeVisible(); }); test("failed login shows error", async ({ page }) => { const loginPage = new LoginPage(page); await loginPage.goto(); await loginPage.login("invalid@example.com", "wrong"); const error = await loginPage.getErrorMessage(); expect(error).toContain("Invalid credentials"); }); ``` ### Pattern 2: Fixtures for Test Data ```typescript // fixtures/test-data.ts import { test as base } from "@playwright/test"; type TestData = { testUser: { email: string; password: string; name: string; }; adminUser: { email: string; password: string; }; }; export const test = base.extend<TestData>({ testUser: async ({}, use) => { const user = { email: `test-${Date.now()}@example.com`, password: "Test123!@#", name: "Test User", }; // Setup: Create user in database await createTestUser(user); await use(user); // Teardown: Clean up user await deleteTestUser(user.email); }, adminUser: async ({}, use) => { await use({ email: "admin@example.com", password: process.env.ADMIN_PASSWORD!, }); }, }); // Usage in tests import { test } from "./fixtures/test-data"; test("user can update profile", async ({ page, testUser }) => { await page.goto("/login"); await page.getByLabel("Email").fill(testUser.email); await page.getByLabel("Password").fill(testUser.password); await page.getByRole("button", { name: "Login" }).click(); await page.goto("/profile"); await page.getByLabel("Name").fill("Updated Name"); await page.getByRole("button", { name: "Save" }).click(); await expect(page.getByText("Profile updated")).toBeVisible(); }); ``` ### Pattern 3: Waiting Strategies ```typescript // ❌ Bad: Fixed timeouts await page.waitForTimeout(3000); // Flaky! // βœ… Good: Wait for specific conditions await page.waitForLoadState("networkidle"); await page.waitForURL("/dashboard"); await page.waitForSelector('[data-testid="user-profile"]'); // βœ… Better: Auto-waiting with assertions await expect(page.getByText("Welcome")).toBeVisible(); await expect(page.getByRole("button", { name: "Submit" })).toBeEnabled(); // Wait for API response const responsePromise = page.waitForResponse( (response) => response.url().includes("/api/users") && response.status() === 200, ); await page.getByRole("button", { name: "Load Users" }).click(); const response = await responsePromise; const data = await response.json(); expect(data.users).toHaveLength(10); // Wait for multiple conditions await Promise.all([ page.waitForURL("/success"), page.waitForLoadState("networkidle"), expect(page.getByText("Payment successful")).toBeVisible(), ]); ``` ### Pattern 4: Network Mocking and Interception ```typescript // Mock API responses test("displays error when API fails", async ({ page }) => { await page.route("**/api/users", (route) => { route.fulfill({ status: 500, contentType: "application/json", body: JSON.stringify({ error: "Internal Server Error" }), }); }); await page.goto("/users"); await expect(page.getByText("Failed to load users")).toBeVisible(); }); // Intercept and modify requests test("can modify API request", async ({ page }) => { await page.route("**/api/users", async (route) => { const request = route.request(); const postData = JSON.parse(request.postData() || "{}"); // Modify request postData.role = "admin"; await route.continue({ postData: JSON.stringify(postData), }); }); // Test continues... }); // Mock third-party services test("payment flow with mocked Stripe", async ({ page }) => { await page.route("**/api/stripe/**", (route) => { route.fulfill({ status: 200, body: JSON.stringify({ id: "mock_payment_id", status: "succeeded", }), }); }); // Test payment flow with mocked response }); ``` ## Cypress Patterns ### Setup and Configuration ```typescript // cypress.config.ts import { defineConfig } from "cypress"; export default defineConfig({ e2e: { baseUrl: "http://localhost:3000", viewportWidth: 1280, viewportHeight: 720, video: false, screenshotOnRunFailure: true, defaultCommandTimeout: 10000, requestTimeout: 10000, setupNodeEvents(on, config) { // Implement node event listeners }, }, }); ``` ### Pattern 1: Custom Commands ```typescript // cypress/support/commands.ts declare global { namespace Cypress { interface Chainable { login(email: string, password: string): Chainable<void>; createUser(userData: UserData): Chainable<User>; dataCy(value: string): Chainable<JQuery<HTMLElement>>; } } } Cypress.Commands.add("login", (email: string, password: string) => { cy.visit("/login"); cy.get('[data-testid="email"]').type(email); cy.get('[data-testid="password"]').type(password); cy.get('[data-testid="login-button"]').click(); cy.url().should("include", "/dashboard"); }); Cypress.Commands.add("createUser", (userData: UserData) => { return cy.request("POST", "/api/users", userData).its("body"); }); Cypress.Commands.add("dataCy", (value: string) => { return cy.get(`[data-cy="${value}"]`); }); // Usage cy.login("user@example.com", "password"); cy.dataCy("submit-button").click(); ``` ### Pattern 2: Cypress Intercept ```typescript // Mock API calls cy.intercept("GET", "/api/users", { statusCode: 200, body: [ { id: 1, name: "John" }, { id: 2, name: "Jane" }, ], }).as("getUsers"); cy.visit("/users"); cy.wait("@getUsers"); cy.get('[data-testid="user-list"]').children().should("have.length", 2); // Modify responses cy.intercept("GET", "/api/users", (req) => { req.reply((res) => { // Modify response res.body.users = res.body.users.slice(0, 5); res.send(); }); }); // Simulate slow network cy.intercept("GET", "/api/data", (req) => { req.reply((res) => { res.delay(3000); // 3 second delay res.send(); }); }); ``` ## Advanced Patterns ### Pattern 1: Visual Regression Testing ```typescript // With Playwright import { test, expect } from "@playwright/test"; test("homepage looks correct", async ({ page }) => { await page.goto("/"); await expect(page).toHaveScreenshot("homepage.png", { fullPage: true, maxDiffPixels: 100, }); }); test("button in all states", async ({ page }) => { await page.goto("/components"); const button = page.getByRole("button", { name: "Submit" }); // Default state await expect(button).toHaveScreenshot("button-default.png"); // Hover state await button.hover(); await expect(button).toHaveScreenshot("button-hover.png"); // Disabled state await button.evaluate((el) => el.setAttribute("disabled", "true")); await expect(button).toHaveScreenshot("button-disabled.png"); }); ``` ### Pattern 2: Parallel Testing with Sharding ```typescript // playwright.config.ts export default defineConfig({ projects: [ { name: "shard-1", use: { ...devices["Desktop Chrome"] }, grepInvert: /@slow/, shard: { current: 1, total: 4 }, }, { name: "shard-2", use: { ...devices["Desktop Chrome"] }, shard: { current: 2, total: 4 }, }, // ... more shards ], }); // Run in CI // npx playwright test --shard=1/4 // npx playwright test --shard=2/4 ``` ### Pattern 3: Accessibility Testing ```typescript // Install: npm install @axe-core/playwright import { test, expect } from "@playwright/test"; import AxeBuilder from "@axe-core/playwright"; test("page should not have accessibility violations", async ({ page }) => { await page.goto("/"); const accessibilityScanResults = await new AxeBuilder({ page }) .exclude("#third-party-widget") .analyze(); expect(accessibilityScanResults.violations).toEqual([]); }); test("form is accessible", async ({ page }) => { await page.goto("/signup"); const results = await new AxeBuilder({ page }).include("form").analyze(); expect(results.violations).toEqual([]); }); ``` ## Best Practices 1. **Use Data Attributes**: `data-testid` or `data-cy` for stable selectors 2. **Avoid Brittle Selectors**: Don't rely on CSS classes or DOM structure 3. **Test User Behavior**: Click, type, see - not implementation details 4. **Keep Tests Independent**: Each test should run in isolation 5. **Clean Up Test Data**: Create and destroy test data in each test 6. **Use Page Objects**: Encapsulate page logic 7. **Meaningful Assertions**: Check actual user-visible behavior 8. **Optimize for Speed**: Mock when possible, parallel execution ```typescript // ❌ Bad selectors cy.get(".btn.btn-primary.submit-button").click(); cy.get("div > form > div:nth-child(2) > input").type("text"); // βœ… Good selectors cy.getByRole("button", { name: "Submit" }).click(); cy.getByLabel("Email address").type("user@example.com"); cy.get('[data-testid="email-input"]').type("user@example.com"); ``` ## Common Pitfalls - **Flaky Tests**: Use proper waits, not fixed timeouts - **Slow Tests**: Mock external APIs, use parallel execution - **Over-Testing**: Don't test every edge case with E2E - **Coupled Tests**: Tests should not depend on each other - **Poor Selectors**: Avoid CSS classes and nth-child - **No Cleanup**: Clean up test data after each test - **Testing Implementation**: Test user behavior, not internals ## Debugging Failing Tests ```typescript // Playwright debugging // 1. Run in headed mode npx playwright test --headed // 2. Run in debug mode npx playwright test --debug // 3. Use trace viewer await page.screenshot({ path: 'screenshot.png' }); await page.video()?.saveAs('video.webm'); // 4. Add test.step for better reporting test('checkout flow', async ({ page }) => { await test.step('Add item to cart', async () => { await page.goto('/products'); await page.getByRole('button', { name: 'Add to Cart' }).click(); }); await test.step('Proceed to checkout', async () => { await page.goto('/cart'); await page.getByRole('button', { name: 'Checkout' }).click(); }); }); // 5. Inspect page state await page.pause(); // Pauses execution, opens inspector ``` ## Resources - **references/playwright-best-practices.md**: Playwright-specific patterns - **references/cypress-best-practices.md**: Cypress-specific patterns - **references/flaky-test-debugging.md**: Debugging unreliable tests - **assets/e2e-testing-checklist.md**: What to test with E2E - **assets/selector-strategies.md**: Finding reliable selectors - **scripts/test-analyzer.ts**: Analyze test flakiness and duration
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

error-handling-patterns

Master error handling patterns across languages including

coding
⭐1
# Error Handling Patterns Build resilient applications with robust error handling strategies that gracefully handle failures and provide excellent debugging experiences. ## When to Use This Skill - Implementing error handling in new features - Designing error-resilient APIs - Debugging production issues - Improving application reliability - Creating better error messages for users and developers - Implementing retry and circuit breaker patterns - Handling async/concurrent errors - Building fault-tolerant distributed systems ## Core Concepts ### 1. Error Handling Philosophies **Exceptions vs Result Types:** - **Exceptions**: Traditional try-catch, disrupts control flow - **Result Types**: Explicit success/failure, functional approach - **Error Codes**: C-style, requires discipline - **Option/Maybe Types**: For nullable values **When to Use Each:** - Exceptions: Unexpected errors, exceptional conditions - Result Types: Expected errors, validation failures - Panics/Crashes: Unrecoverable errors, programming bugs ### 2. Error Categories **Recoverable Errors:** - Network timeouts - Missing files - Invalid user input - API rate limits **Unrecoverable Errors:** - Out of memory - Stack overflow - Programming bugs (null pointer, etc.) ## Language-Specific Patterns ### Python Error Handling **Custom Exception Hierarchy:** ```python class ApplicationError(Exception): """Base exception for all application errors.""" def __init__(self, message: str, code: str = None, details: dict = None): super().__init__(message) self.code = code self.details = details or {} self.timestamp = datetime.utcnow() class ValidationError(ApplicationError): """Raised when validation fails.""" pass class NotFoundError(ApplicationError): """Raised when resource not found.""" pass class ExternalServiceError(ApplicationError): """Raised when external service fails.""" def __init__(self, message: str, service: str, **kwargs): super().__init__(message, **kwargs) self.service = service # Usage def get_user(user_id: str) -> User: user = db.query(User).filter_by(id=user_id).first() if not user: raise NotFoundError( f"User not found", code="USER_NOT_FOUND", details={"user_id": user_id} ) return user ``` **Context Managers for Cleanup:** ```python from contextlib import contextmanager @contextmanager def database_transaction(session): """Ensure transaction is committed or rolled back.""" try: yield session session.commit() except Exception as e: session.rollback() raise finally: session.close() # Usage with database_transaction(db.session) as session: user = User(name="Alice") session.add(user) # Automatic commit or rollback ``` **Retry with Exponential Backoff:** ```python import time from functools import wraps from typing import TypeVar, Callable T = TypeVar('T') def retry( max_attempts: int = 3, backoff_factor: float = 2.0, exceptions: tuple = (Exception,) ): """Retry decorator with exponential backoff.""" def decorator(func: Callable[..., T]) -> Callable[..., T]: @wraps(func) def wrapper(*args, **kwargs) -> T: last_exception = None for attempt in range(max_attempts): try: return func(*args, **kwargs) except exceptions as e: last_exception = e if attempt < max_attempts - 1: sleep_time = backoff_factor ** attempt time.sleep(sleep_time) continue raise raise last_exception return wrapper return decorator # Usage @retry(max_attempts=3, exceptions=(NetworkError,)) def fetch_data(url: str) -> dict: response = requests.get(url, timeout=5) response.raise_for_status() return response.json() ``` ### TypeScript/JavaScript Error Handling **Custom Error Classes:** ```typescript // Custom error classes class ApplicationError extends Error { constructor( message: string, public code: string, public statusCode: number = 500, public details?: Record<string, any>, ) { super(message); this.name = this.constructor.name; Error.captureStackTrace(this, this.constructor); } } class ValidationError extends ApplicationError { constructor(message: string, details?: Record<string, any>) { super(message, "VALIDATION_ERROR", 400, details); } } class NotFoundError extends ApplicationError { constructor(resource: string, id: string) { super(`${resource} not found`, "NOT_FOUND", 404, { resource, id }); } } // Usage function getUser(id: string): User { const user = users.find((u) => u.id === id); if (!user) { throw new NotFoundError("User", id); } return user; } ``` **Result Type Pattern:** ```typescript // Result type for explicit error handling type Result<T, E = Error> = { ok: true; value: T } | { ok: false; error: E }; // Helper functions function Ok<T>(value: T): Result<T, never> { return { ok: true, value }; } function Err<E>(error: E): Result<never, E> { return { ok: false, error }; } // Usage function parseJSON<T>(json: string): Result<T, SyntaxError> { try { const value = JSON.parse(json) as T; return Ok(value); } catch (error) { return Err(error as SyntaxError); } } // Consuming Result const result = parseJSON<User>(userJson); if (result.ok) { console.log(result.value.name); } else { console.error("Parse failed:", result.error.message); } // Chaining Results function chain<T, U, E>( result: Result<T, E>, fn: (value: T) => Result<U, E>, ): Result<U, E> { return result.ok ? fn(result.value) : result; } ``` **Async Error Handling:** ```typescript // Async/await with proper error handling async function fetchUserOrders(userId: string): Promise<Order[]> { try { const user = await getUser(userId); const orders = await getOrders(user.id); return orders; } catch (error) { if (error instanceof NotFoundError) { return []; // Return empty array for not found } if (error instanceof NetworkError) { // Retry logic return retryFetchOrders(userId); } // Re-throw unexpected errors throw error; } } // Promise error handling function fetchData(url: string): Promise<Data> { return fetch(url) .then((response) => { if (!response.ok) { throw new NetworkError(`HTTP ${response.status}`); } return response.json(); }) .catch((error) => { console.error("Fetch failed:", error); throw error; }); } ``` ### Rust Error Handling **Result and Option Types:** ```rust use std::fs::File; use std::io::{self, Read}; // Result type for operations that can fail fn read_file(path: &str) -> Result<String, io::Error> { let mut file = File::open(path)?; // ? operator propagates errors let mut contents = String::new(); file.read_to_string(&mut contents)?; Ok(contents) } // Custom error types #[derive(Debug)] enum AppError { Io(io::Error), Parse(std::num::ParseIntError), NotFound(String), Validation(String), } impl From<io::Error> for AppError { fn from(error: io::Error) -> Self { AppError::Io(error) } } // Using custom error type fn read_number_from_file(path: &str) -> Result<i32, AppError> { let contents = read_file(path)?; // Auto-converts io::Error let number = contents.trim().parse() .map_err(AppError::Parse)?; // Explicitly convert ParseIntError Ok(number) } // Option for nullable values fn find_user(id: &str) -> Option<User> { users.iter().find(|u| u.id == id).cloned() } // Combining Option and Result fn get_user_age(id: &str) -> Result<u32, AppError> { find_user(id) .ok_or_else(|| AppError::NotFound(id.to_string())) .map(|user| user.age) } ``` ### Go Error Handling **Explicit Error Returns:** ```go // Basic error handling func getUser(id string) (*User, error) { user, err := db.QueryUser(id) if err != nil { return nil, fmt.Errorf("failed to query user: %w", err) } if user == nil { return nil, errors.New("user not found") } return user, nil } // Custom error types type ValidationError struct { Field string Message string } func (e *ValidationError) Error() string { return fmt.Sprintf("validation failed for %s: %s", e.Field, e.Message) } // Sentinel errors for comparison var ( ErrNotFound = errors.New("not found") ErrUnauthorized = errors.New("unauthorized") ErrInvalidInput = errors.New("invalid input") ) // Error checking user, err := getUser("123") if err != nil { if errors.Is(err, ErrNotFound) { // Handle not found } else { // Handle other errors } } // Error wrapping and unwrapping func processUser(id string) error { user, err := getUser(id) if err != nil { return fmt.Errorf("process user failed: %w", err) } // Process user return nil } // Unwrap errors err := processUser("123") if err != nil { var valErr *ValidationError if errors.As(err, &valErr) { fmt.Printf("Validation error: %s\n", valErr.Field) } } ``` ## Universal Patterns ### Pattern 1: Circuit Breaker Prevent cascading failures in distributed systems. ```python from enum import Enum from datetime import datetime, timedelta from typing import Callable, TypeVar T = TypeVar('T') class CircuitState(Enum): CLOSED = "closed" # Normal operation OPEN = "open" # Failing, reject requests HALF_OPEN = "half_open" # Testing if recovered class CircuitBreaker: def __init__( self, failure_threshold: int = 5, timeout: timedelta = timedelta(seconds=60), success_threshold: int = 2 ): self.failure_threshold = failure_threshold self.timeout = timeout self.success_threshold = success_threshold self.failure_count = 0 self.success_count = 0 self.state = CircuitState.CLOSED self.last_failure_time = None def call(self, func: Callable[[], T]) -> T: if self.state == CircuitState.OPEN: if datetime.now() - self.last_failure_time > self.timeout: self.state = CircuitState.HALF_OPEN self.success_count = 0 else: raise Exception("Circuit breaker is OPEN") try: result = func() self.on_success() return result except Exception as e: self.on_failure() raise def on_success(self): 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): self.failure_count += 1 self.last_failure_time = datetime.now() if self.failure_count >= self.failure_threshold: self.state = CircuitState.OPEN # Usage circuit_breaker = CircuitBreaker() def fetch_data(): return circuit_breaker.call(lambda: external_api.get_data()) ``` ### Pattern 2: Error Aggregation Collect multiple errors instead of failing on first error. ```typescript class ErrorCollector { private errors: Error[] = []; add(error: Error): void { this.errors.push(error); } hasErrors(): boolean { return this.errors.length > 0; } getErrors(): Error[] { return [...this.errors]; } throw(): never { if (this.errors.length === 1) { throw this.errors[0]; } throw new AggregateError( this.errors, `${this.errors.length} errors occurred`, ); } } // Usage: Validate multiple fields function validateUser(data: any): User { const errors = new ErrorCollector(); if (!data.email) { errors.add(new ValidationError("Email is required")); } else if (!isValidEmail(data.email)) { errors.add(new ValidationError("Email is invalid")); } if (!data.name || data.name.length < 2) { errors.add(new ValidationError("Name must be at least 2 characters")); } if (!data.age || data.age < 18) { errors.add(new ValidationError("Age must be 18 or older")); } if (errors.hasErrors()) { errors.throw(); } return data as User; } ``` ### Pattern 3: Graceful Degradation Provide fallback functionality when errors occur. ```python from typing import Optional, Callable, TypeVar T = TypeVar('T') def with_fallback( primary: Callable[[], T], fallback: Callable[[], T], log_error: bool = True ) -> T: """Try primary function, fall back to fallback on error.""" try: return primary() except Exception as e: if log_error: logger.error(f"Primary function failed: {e}") return fallback() # Usage def get_user_profile(user_id: str) -> UserProfile: return with_fallback( primary=lambda: fetch_from_cache(user_id), fallback=lambda: fetch_from_database(user_id) ) # Multiple fallbacks def get_exchange_rate(currency: str) -> float: return ( try_function(lambda: api_provider_1.get_rate(currency)) or try_function(lambda: api_provider_2.get_rate(currency)) or try_function(lambda: cache.get_rate(currency)) or DEFAULT_RATE ) def try_function(func: Callable[[], Optional[T]]) -> Optional[T]: try: return func() except Exception: return None ``` ## Best Practices 1. **Fail Fast**: Validate input early, fail quickly 2. **Preserve Context**: Include stack traces, metadata, timestamps 3. **Meaningful Messages**: Explain what happened and how to fix it 4. **Log Appropriately**: Error = log, expected failure = don't spam logs 5. **Handle at Right Level**: Catch where you can meaningfully handle 6. **Clean Up Resources**: Use try-finally, context managers, defer 7. **Don't Swallow Errors**: Log or re-throw, don't silently ignore 8. **Type-Safe Errors**: Use typed errors when possible ```python # Good error handling example def process_order(order_id: str) -> Order: """Process order with comprehensive error handling.""" try: # Validate input if not order_id: raise ValidationError("Order ID is required") # Fetch order order = db.get_order(order_id) if not order: raise NotFoundError("Order", order_id) # Process payment try: payment_result = payment_service.charge(order.total) except PaymentServiceError as e: # Log and wrap external service error logger.error(f"Payment failed for order {order_id}: {e}") raise ExternalServiceError( f"Payment processing failed", service="payment_service", details={"order_id": order_id, "amount": order.total} ) from e # Update order order.status = "completed" order.payment_id = payment_result.id db.save(order) return order except ApplicationError: # Re-raise known application errors raise except Exception as e: # Log unexpected errors logger.exception(f"Unexpected error processing order {order_id}") raise ApplicationError( "Order processing failed", code="INTERNAL_ERROR" ) from e ``` ## Common Pitfalls - **Catching Too Broadly**: `except Exception` hides bugs - **Empty Catch Blocks**: Silently swallowing errors - **Logging and Re-throwing**: Creates duplicate log entries - **Not Cleaning Up**: Forgetting to close files, connections - **Poor Error Messages**: "Error occurred" is not helpful - **Returning Error Codes**: Use exceptions or Result types - **Ignoring Async Errors**: Unhandled promise rejections ## Resources - **references/exception-hierarchy-design.md**: Designing error class hierarchies - **references/error-recovery-strategies.md**: Recovery patterns for different scenarios - **references/async-error-handling.md**: Handling errors in concurrent code - **assets/error-handling-checklist.md**: Review checklist for error handling - **assets/error-message-guide.md**: Writing helpful error messages - **scripts/error-analyzer.py**: Analyze error patterns in logs
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

nx-workspace-patterns

Configure and optimize Nx monorepo workspaces. Use when setting up

coding
⭐1
# Nx Workspace Patterns Production patterns for Nx monorepo management. ## When to Use This Skill - Setting up new Nx workspaces - Configuring project boundaries - Optimizing CI with affected commands - Implementing remote caching - Managing dependencies between projects - Migrating to Nx ## Core Concepts ### 1. Nx Architecture ``` workspace/ β”œβ”€β”€ apps/ # Deployable applications β”‚ β”œβ”€β”€ web/ β”‚ └── api/ β”œβ”€β”€ libs/ # Shared libraries β”‚ β”œβ”€β”€ shared/ β”‚ β”‚ β”œβ”€β”€ ui/ β”‚ β”‚ └── utils/ β”‚ └── feature/ β”‚ β”œβ”€β”€ auth/ β”‚ └── dashboard/ β”œβ”€β”€ tools/ # Custom executors/generators β”œβ”€β”€ nx.json # Nx configuration └── workspace.json # Project configuration ``` ### 2. Library Types | Type | Purpose | Example | | --------------- | -------------------------------- | ------------------- | | **feature** | Smart components, business logic | `feature-auth` | | **ui** | Presentational components | `ui-buttons` | | **data-access** | API calls, state management | `data-access-users` | | **util** | Pure functions, helpers | `util-formatting` | | **shell** | App bootstrapping | `shell-web` | ## Templates ### Template 1: nx.json Configuration ```json { "$schema": "./node_modules/nx/schemas/nx-schema.json", "npmScope": "myorg", "affected": { "defaultBase": "main" }, "tasksRunnerOptions": { "default": { "runner": "nx/tasks-runners/default", "options": { "cacheableOperations": [ "build", "lint", "test", "e2e", "build-storybook" ], "parallel": 3 } } }, "targetDefaults": { "build": { "dependsOn": ["^build"], "inputs": ["production", "^production"], "cache": true }, "test": { "inputs": ["default", "^production", "{workspaceRoot}/jest.preset.js"], "cache": true }, "lint": { "inputs": ["default", "{workspaceRoot}/.eslintrc.json"], "cache": true }, "e2e": { "inputs": ["default", "^production"], "cache": true } }, "namedInputs": { "default": ["{projectRoot}/**/*", "sharedGlobals"], "production": [ "default", "!{projectRoot}/**/?(*.)+(spec|test).[jt]s?(x)?(.snap)", "!{projectRoot}/tsconfig.spec.json", "!{projectRoot}/jest.config.[jt]s", "!{projectRoot}/.eslintrc.json" ], "sharedGlobals": [ "{workspaceRoot}/babel.config.json", "{workspaceRoot}/tsconfig.base.json" ] }, "generators": { "@nx/react": { "application": { "style": "css", "linter": "eslint", "bundler": "webpack" }, "library": { "style": "css", "linter": "eslint" }, "component": { "style": "css" } } } } ``` ### Template 2: Project Configuration ```json // apps/web/project.json { "name": "web", "$schema": "../../node_modules/nx/schemas/project-schema.json", "sourceRoot": "apps/web/src", "projectType": "application", "tags": ["type:app", "scope:web"], "targets": { "build": { "executor": "@nx/webpack:webpack", "outputs": ["{options.outputPath}"], "defaultConfiguration": "production", "options": { "compiler": "babel", "outputPath": "dist/apps/web", "index": "apps/web/src/index.html", "main": "apps/web/src/main.tsx", "tsConfig": "apps/web/tsconfig.app.json", "assets": ["apps/web/src/assets"], "styles": ["apps/web/src/styles.css"] }, "configurations": { "development": { "extractLicenses": false, "optimization": false, "sourceMap": true }, "production": { "optimization": true, "outputHashing": "all", "sourceMap": false, "extractLicenses": true } } }, "serve": { "executor": "@nx/webpack:dev-server", "defaultConfiguration": "development", "options": { "buildTarget": "web:build" }, "configurations": { "development": { "buildTarget": "web:build:development" }, "production": { "buildTarget": "web:build:production" } } }, "test": { "executor": "@nx/jest:jest", "outputs": ["{workspaceRoot}/coverage/{projectRoot}"], "options": { "jestConfig": "apps/web/jest.config.ts", "passWithNoTests": true } }, "lint": { "executor": "@nx/eslint:lint", "outputs": ["{options.outputFile}"], "options": { "lintFilePatterns": ["apps/web/**/*.{ts,tsx,js,jsx}"] } } } } ``` ### Template 3: Module Boundary Rules ```json // .eslintrc.json { "root": true, "ignorePatterns": ["**/*"], "plugins": ["@nx"], "overrides": [ { "files": ["*.ts", "*.tsx", "*.js", "*.jsx"], "rules": { "@nx/enforce-module-boundaries": [ "error", { "enforceBuildableLibDependency": true, "allow": [], "depConstraints": [ { "sourceTag": "type:app", "onlyDependOnLibsWithTags": [ "type:feature", "type:ui", "type:data-access", "type:util" ] }, { "sourceTag": "type:feature", "onlyDependOnLibsWithTags": [ "type:ui", "type:data-access", "type:util" ] }, { "sourceTag": "type:ui", "onlyDependOnLibsWithTags": ["type:ui", "type:util"] }, { "sourceTag": "type:data-access", "onlyDependOnLibsWithTags": ["type:data-access", "type:util"] }, { "sourceTag": "type:util", "onlyDependOnLibsWithTags": ["type:util"] }, { "sourceTag": "scope:web", "onlyDependOnLibsWithTags": ["scope:web", "scope:shared"] }, { "sourceTag": "scope:api", "onlyDependOnLibsWithTags": ["scope:api", "scope:shared"] }, { "sourceTag": "scope:shared", "onlyDependOnLibsWithTags": ["scope:shared"] } ] } ] } } ] } ``` ### Template 4: Custom Generator ```typescript // tools/generators/feature-lib/index.ts import { Tree, formatFiles, generateFiles, joinPathFragments, names, readProjectConfiguration, } from "@nx/devkit"; import { libraryGenerator } from "@nx/react"; interface FeatureLibraryGeneratorSchema { name: string; scope: string; directory?: string; } export default async function featureLibraryGenerator( tree: Tree, options: FeatureLibraryGeneratorSchema, ) { const { name, scope, directory } = options; const projectDirectory = directory ? `${directory}/${name}` : `libs/${scope}/feature-${name}`; // Generate base library await libraryGenerator(tree, { name: `feature-${name}`, directory: projectDirectory, tags: `type:feature,scope:${scope}`, style: "css", skipTsConfig: false, skipFormat: true, unitTestRunner: "jest", linter: "eslint", }); // Add custom files const projectConfig = readProjectConfiguration( tree, `${scope}-feature-${name}`, ); const projectNames = names(name); generateFiles( tree, joinPathFragments(__dirname, "files"), projectConfig.sourceRoot, { ...projectNames, scope, tmpl: "", }, ); await formatFiles(tree); } ``` ### Template 5: CI Configuration with Affected ```yaml # .github/workflows/ci.yml name: CI on: push: branches: [main] pull_request: branches: [main] env: NX_CLOUD_ACCESS_TOKEN: ${{ secrets.NX_CLOUD_ACCESS_TOKEN }} jobs: main: runs-on: ubuntu-latest steps: - uses: actions/checkout@v4 with: fetch-depth: 0 - uses: actions/setup-node@v4 with: node-version: 20 cache: "npm" - name: Install dependencies run: npm ci - name: Derive SHAs for affected commands uses: nrwl/nx-set-shas@v4 - name: Run affected lint run: npx nx affected -t lint --parallel=3 - name: Run affected test run: npx nx affected -t test --parallel=3 --configuration=ci - name: Run affected build run: npx nx affected -t build --parallel=3 - name: Run affected e2e run: npx nx affected -t e2e --parallel=1 ``` ### Template 6: Remote Caching Setup ```typescript // nx.json with Nx Cloud { "tasksRunnerOptions": { "default": { "runner": "nx-cloud", "options": { "cacheableOperations": ["build", "lint", "test", "e2e"], "accessToken": "your-nx-cloud-token", "parallel": 3, "cacheDirectory": ".nx/cache" } } }, "nxCloudAccessToken": "your-nx-cloud-token" } // Self-hosted cache with S3 { "tasksRunnerOptions": { "default": { "runner": "@nx-aws-cache/nx-aws-cache", "options": { "cacheableOperations": ["build", "lint", "test"], "awsRegion": "us-east-1", "awsBucket": "my-nx-cache-bucket", "awsProfile": "default" } } } } ``` ## Common Commands ```bash # Generate new library nx g @nx/react:lib feature-auth --directory=libs/web --tags=type:feature,scope:web # Run affected tests nx affected -t test --base=main # View dependency graph nx graph # Run specific project nx build web --configuration=production # Reset cache nx reset # Run migrations nx migrate latest nx migrate --run-migrations ``` ## Best Practices ### Do's - **Use tags consistently** - Enforce with module boundaries - **Enable caching early** - Significant CI savings - **Keep libs focused** - Single responsibility - **Use generators** - Ensure consistency - **Document boundaries** - Help new developers ### Don'ts - **Don't create circular deps** - Graph should be acyclic - **Don't skip affected** - Test only what changed - **Don't ignore boundaries** - Tech debt accumulates - **Don't over-granularize** - Balance lib count ## Resources - [Nx Documentation](https://nx.dev/getting-started/intro) - [Module Boundaries](https://nx.dev/core-features/enforce-module-boundaries) - [Nx Cloud](https://nx.app/)
πŸ‘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