Skip to main content
EVOKORE// BROWSE
>

./browse/prompts

14 NODES
📝text•3 hours ago

Narrative Control Prompt: Exhaustive System Architecture & Feature Reverse-Engineering

Narrative Control Prompt: Exhaustive System Architecture & Feature Reverse-Engineering User centric prompting for analyzing inspiration pages or instructing on internal analysis cycles

architecture
⭐1
# Narrative Control Prompt: Exhaustive System Architecture & Feature Reverse-Engineering You are an **Expert Enterprise Architect, Product Director, and Lead Engineer**. Your goal is to thoroughly analyze the provided documentation to architect a full-scale, competitive enterprise application. You must deconstruct the feature described in the content into **extensive, granular technical specifications** across multiple engineering disciplines. **CRITICAL INSTRUCTIONS**: - **DO NOT SUMMARIZE**. Be exhaustive. - For every category below, aim to list **10+ specific items** if possible. - Brainstorm every possible implication, edge case, and requirement derived from or inspired by the text. - If the text mentions a "search" feature, break it down into: Indexing, Query Parsing, UI Widgets, Highlighting, Filtering, Sort Logic, Caching, etc. Please output the analysis in the following Markdown format: # 1. Product Strategy & Scope * **Feature Name**: * **Core Value Proposition**: [Deep dive into why this exists] * **User Personas**: [List as many as applicable: e.g. Admin, Power User, Viewer, Auditor, API Consumer...] * **User Stories**: [Extensive list of 10+ granular user stories e.g. "As a User, I want to..."] * **Competitive Differentiators**: [What makes this specific implementation valuable?] # 2. Design & User Experience (UX/UI) * **Key Interface Components**: [List 10+ atoms/molecules: e.g. Data Grid, Filter Chips, Modals, Tooltips, Empty States, Toasts, Dropdowns...] * **Interaction Patterns**: [List 10+ patterns: e.g. Drag-and-drop, Double-click to edit, Hover states, Keyboard shortcuts, Infinite scroll...] * **Visual States**: [List all states: Loading, Success, Error, Warning, Partial Data, Offline...] * **Accessibility (a11y)**: [List 10+ checks: Contrast, ARIA labels, Focus management, Screen reader support, Resizing...] # 3. Frontend Engineering * **State Management**: [List 10+ state atoms: Upload progress %, Selected ID list, Sort order, Filter criteria, Current user permissions...] * **API Interactions**: [List 10+ potential endpoints: GET/POST/PUT/DELETE for main entities, Lookups, Search, Validation...] * **Component Architecture**: [List 10+ React/Vue components: Container, Presentation, Utility wrappers, HoC...] * **Client-Side Logic**: [Validation rules, Formatting (Dates/Currency), Debouncing, Caching...] # 4. Backend Engineering * **Data Models**: [List 10+ fields/entities: Table structure, Foreign keys, Indexes, JSONB fields, Audit columns...] * **API Specification**: [Detailed endpoint contract: Header requirements, Query params, Body schema, Error codes...] * **Business Logic**: [List 10+ rules: Permission checks, Data transformation, Workflows, Triggers, Notifications...] * **Security & Permissions**: [List 10+ checks: RBAC roles, Field-level security, API Rate limiting, CSRF protection...] # 5. Infrastructure & DevOps * **Storage Requirements**: [S3 buckets, Database types (SQL/NoSQL), Redis for cache, CDNs...] * **Compute Needs**: [Async workers, Scheduled cron jobs, Serverless functions, Container specs...] * **Background Jobs**: [List 10+ potential jobs: Email sending, File conversion, Indexing, Cleanup, Analytics aggregation...] * **Observability**: [Metrics to track: API latency, Error rates, Disk usage, Active users...] # 6. Quality Assurance (QA) * **Test Scenarios**: [List 10+ happy path scenarios] * **Edge Cases**: [List 10+ negative/edge cases: Network fail, Giant files, Concurrent edits, Invalid chars...] * **Performance Metrics**: [Specific SLAs: <200ms API, <1s Page load, 99.9% Uptime...] * **Security Testing**: [Pen-test vectors: XSS injection input, SQL injection, IDOR...] # 7. Documentation & Onboarding * **User Guides Needed**: [List 10+ articles to write based on this feature] * **Contextual Help**: [List 10+ places for Tooltips, Tours, Helper text...] * **API Documentation**: [Swagger/OpenAPI requirements] # 8. Implementation Roadmap * **Phase 1 (MVP)**: [List 10+ must-have tasks] * **Phase 2 (Enhanced)**: [List 10+ nice-to-have features] * **Phase 3 (Scale)**: [Optimization and enterprise hardening] --- **Context**: The content below is raw markdown from a help guide.
👍0
👁️0
🤖 Auto-discovered
📝text•3 hours ago

Narrative Control Prompt Exhaustive System Architecture & Feature Reverse-Engineering

Narrative Control Prompt: Exhaustive System Architecture & Feature Reverse-Engineering User concentric prompting for analyzing inspiration pages or instructing on internal analysis cycles

architecture
⭐1
# Narrative Control Prompt: Exhaustive System Architecture & Feature Reverse-Engineering You are an **Expert Enterprise Architect, Product Director, and Lead Engineer**. Your goal is to thoroughly analyze the provided documentation to architect a full-scale, competitive enterprise application. You must deconstruct the feature described in the content into **extensive, granular technical specifications** across multiple engineering disciplines. **CRITICAL INSTRUCTIONS**: - **DO NOT SUMMARIZE**. Be exhaustive. - For every category below, aim to list **10+ specific items** if possible. - Brainstorm every possible implication, edge case, and requirement derived from or inspired by the text. - If the text mentions a "search" feature, break it down into: Indexing, Query Parsing, UI Widgets, Highlighting, Filtering, Sort Logic, Caching, etc. Please output the analysis in the following Markdown format: # 1. Product Strategy & Scope * **Feature Name**: * **Core Value Proposition**: [Deep dive into why this exists] * **User Personas**: [List as many as applicable: e.g. Admin, Power User, Viewer, Auditor, API Consumer...] * **User Stories**: [Extensive list of 10+ granular user stories e.g. "As a User, I want to..."] * **Competitive Differentiators**: [What makes this specific implementation valuable?] # 2. Design & User Experience (UX/UI) * **Key Interface Components**: [List 10+ atoms/molecules: e.g. Data Grid, Filter Chips, Modals, Tooltips, Empty States, Toasts, Dropdowns...] * **Interaction Patterns**: [List 10+ patterns: e.g. Drag-and-drop, Double-click to edit, Hover states, Keyboard shortcuts, Infinite scroll...] * **Visual States**: [List all states: Loading, Success, Error, Warning, Partial Data, Offline...] * **Accessibility (a11y)**: [List 10+ checks: Contrast, ARIA labels, Focus management, Screen reader support, Resizing...] # 3. Frontend Engineering * **State Management**: [List 10+ state atoms: Upload progress %, Selected ID list, Sort order, Filter criteria, Current user permissions...] * **API Interactions**: [List 10+ potential endpoints: GET/POST/PUT/DELETE for main entities, Lookups, Search, Validation...] * **Component Architecture**: [List 10+ React/Vue components: Container, Presentation, Utility wrappers, HoC...] * **Client-Side Logic**: [Validation rules, Formatting (Dates/Currency), Debouncing, Caching...] # 4. Backend Engineering * **Data Models**: [List 10+ fields/entities: Table structure, Foreign keys, Indexes, JSONB fields, Audit columns...] * **API Specification**: [Detailed endpoint contract: Header requirements, Query params, Body schema, Error codes...] * **Business Logic**: [List 10+ rules: Permission checks, Data transformation, Workflows, Triggers, Notifications...] * **Security & Permissions**: [List 10+ checks: RBAC roles, Field-level security, API Rate limiting, CSRF protection...] # 5. Infrastructure & DevOps * **Storage Requirements**: [S3 buckets, Database types (SQL/NoSQL), Redis for cache, CDNs...] * **Compute Needs**: [Async workers, Scheduled cron jobs, Serverless functions, Container specs...] * **Background Jobs**: [List 10+ potential jobs: Email sending, File conversion, Indexing, Cleanup, Analytics aggregation...] * **Observability**: [Metrics to track: API latency, Error rates, Disk usage, Active users...] # 6. Quality Assurance (QA) * **Test Scenarios**: [List 10+ happy path scenarios] * **Edge Cases**: [List 10+ negative/edge cases: Network fail, Giant files, Concurrent edits, Invalid chars...] * **Performance Metrics**: [Specific SLAs: <200ms API, <1s Page load, 99.9% Uptime...] * **Security Testing**: [Pen-test vectors: XSS injection input, SQL injection, IDOR...] # 7. Documentation & Onboarding * **User Guides Needed**: [List 10+ articles to write based on this feature] * **Contextual Help**: [List 10+ places for Tooltips, Tours, Helper text...] * **API Documentation**: [Swagger/OpenAPI requirements] # 8. Implementation Roadmap * **Phase 1 (MVP)**: [List 10+ must-have tasks] * **Phase 2 (Enhanced)**: [List 10+ nice-to-have features] * **Phase 3 (Scale)**: [Optimization and enterprise hardening] --- **Context**: The content below is raw markdown from a help guide.
👍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

sql-optimization-patterns

Master SQL query optimization, indexing strategies, and EXPLAIN

coding
⭐1
# SQL Optimization Patterns Transform slow database queries into lightning-fast operations through systematic optimization, proper indexing, and query plan analysis. ## When to Use This Skill - Debugging slow-running queries - Designing performant database schemas - Optimizing application response times - Reducing database load and costs - Improving scalability for growing datasets - Analyzing EXPLAIN query plans - Implementing efficient indexes - Resolving N+1 query problems ## Core Concepts ### 1. Query Execution Plans (EXPLAIN) Understanding EXPLAIN output is fundamental to optimization. **PostgreSQL EXPLAIN:** ```sql -- Basic explain EXPLAIN SELECT * FROM users WHERE email = 'user@example.com'; -- With actual execution stats EXPLAIN ANALYZE SELECT * FROM users WHERE email = 'user@example.com'; -- Verbose output with more details EXPLAIN (ANALYZE, BUFFERS, VERBOSE) SELECT u.*, o.order_total FROM users u JOIN orders o ON u.id = o.user_id WHERE u.created_at > NOW() - INTERVAL '30 days'; ``` **Key Metrics to Watch:** - **Seq Scan**: Full table scan (usually slow for large tables) - **Index Scan**: Using index (good) - **Index Only Scan**: Using index without touching table (best) - **Nested Loop**: Join method (okay for small datasets) - **Hash Join**: Join method (good for larger datasets) - **Merge Join**: Join method (good for sorted data) - **Cost**: Estimated query cost (lower is better) - **Rows**: Estimated rows returned - **Actual Time**: Real execution time ### 2. Index Strategies Indexes are the most powerful optimization tool. **Index Types:** - **B-Tree**: Default, good for equality and range queries - **Hash**: Only for equality (=) comparisons - **GIN**: Full-text search, array queries, JSONB - **GiST**: Geometric data, full-text search - **BRIN**: Block Range INdex for very large tables with correlation ```sql -- Standard B-Tree index CREATE INDEX idx_users_email ON users(email); -- Composite index (order matters!) CREATE INDEX idx_orders_user_status ON orders(user_id, status); -- Partial index (index subset of rows) CREATE INDEX idx_active_users ON users(email) WHERE status = 'active'; -- Expression index CREATE INDEX idx_users_lower_email ON users(LOWER(email)); -- Covering index (include additional columns) CREATE INDEX idx_users_email_covering ON users(email) INCLUDE (name, created_at); -- Full-text search index CREATE INDEX idx_posts_search ON posts USING GIN(to_tsvector('english', title || ' ' || body)); -- JSONB index CREATE INDEX idx_metadata ON events USING GIN(metadata); ``` ### 3. Query Optimization Patterns **Avoid SELECT \*:** ```sql -- Bad: Fetches unnecessary columns SELECT * FROM users WHERE id = 123; -- Good: Fetch only what you need SELECT id, email, name FROM users WHERE id = 123; ``` **Use WHERE Clause Efficiently:** ```sql -- Bad: Function prevents index usage SELECT * FROM users WHERE LOWER(email) = 'user@example.com'; -- Good: Create functional index or use exact match CREATE INDEX idx_users_email_lower ON users(LOWER(email)); -- Then: SELECT * FROM users WHERE LOWER(email) = 'user@example.com'; -- Or store normalized data SELECT * FROM users WHERE email = 'user@example.com'; ``` **Optimize JOINs:** ```sql -- Bad: Cartesian product then filter SELECT u.name, o.total FROM users u, orders o WHERE u.id = o.user_id AND u.created_at > '2024-01-01'; -- Good: Filter before join SELECT u.name, o.total FROM users u JOIN orders o ON u.id = o.user_id WHERE u.created_at > '2024-01-01'; -- Better: Filter both tables SELECT u.name, o.total FROM (SELECT * FROM users WHERE created_at > '2024-01-01') u JOIN orders o ON u.id = o.user_id; ``` ## Optimization Patterns ### Pattern 1: Eliminate N+1 Queries **Problem: N+1 Query Anti-Pattern** ```python # Bad: Executes N+1 queries users = db.query("SELECT * FROM users LIMIT 10") for user in users: orders = db.query("SELECT * FROM orders WHERE user_id = ?", user.id) # Process orders ``` **Solution: Use JOINs or Batch Loading** ```sql -- Solution 1: JOIN SELECT u.id, u.name, o.id as order_id, o.total FROM users u LEFT JOIN orders o ON u.id = o.user_id WHERE u.id IN (1, 2, 3, 4, 5); -- Solution 2: Batch query SELECT * FROM orders WHERE user_id IN (1, 2, 3, 4, 5); ``` ```python # Good: Single query with JOIN or batch load # Using JOIN results = db.query(""" SELECT u.id, u.name, o.id as order_id, o.total FROM users u LEFT JOIN orders o ON u.id = o.user_id WHERE u.id IN (1, 2, 3, 4, 5) """) # Or batch load users = db.query("SELECT * FROM users LIMIT 10") user_ids = [u.id for u in users] orders = db.query( "SELECT * FROM orders WHERE user_id IN (?)", user_ids ) # Group orders by user_id orders_by_user = {} for order in orders: orders_by_user.setdefault(order.user_id, []).append(order) ``` ### Pattern 2: Optimize Pagination **Bad: OFFSET on Large Tables** ```sql -- Slow for large offsets SELECT * FROM users ORDER BY created_at DESC LIMIT 20 OFFSET 100000; -- Very slow! ``` **Good: Cursor-Based Pagination** ```sql -- Much faster: Use cursor (last seen ID) SELECT * FROM users WHERE created_at < '2024-01-15 10:30:00' -- Last cursor ORDER BY created_at DESC LIMIT 20; -- With composite sorting SELECT * FROM users WHERE (created_at, id) < ('2024-01-15 10:30:00', 12345) ORDER BY created_at DESC, id DESC LIMIT 20; -- Requires index CREATE INDEX idx_users_cursor ON users(created_at DESC, id DESC); ``` ### Pattern 3: Aggregate Efficiently **Optimize COUNT Queries:** ```sql -- Bad: Counts all rows SELECT COUNT(*) FROM orders; -- Slow on large tables -- Good: Use estimates for approximate counts SELECT reltuples::bigint AS estimate FROM pg_class WHERE relname = 'orders'; -- Good: Filter before counting SELECT COUNT(*) FROM orders WHERE created_at > NOW() - INTERVAL '7 days'; -- Better: Use index-only scan CREATE INDEX idx_orders_created ON orders(created_at); SELECT COUNT(*) FROM orders WHERE created_at > NOW() - INTERVAL '7 days'; ``` **Optimize GROUP BY:** ```sql -- Bad: Group by then filter SELECT user_id, COUNT(*) as order_count FROM orders GROUP BY user_id HAVING COUNT(*) > 10; -- Better: Filter first, then group (if possible) SELECT user_id, COUNT(*) as order_count FROM orders WHERE status = 'completed' GROUP BY user_id HAVING COUNT(*) > 10; -- Best: Use covering index CREATE INDEX idx_orders_user_status ON orders(user_id, status); ``` ### Pattern 4: Subquery Optimization **Transform Correlated Subqueries:** ```sql -- Bad: Correlated subquery (runs for each row) SELECT u.name, u.email, (SELECT COUNT(*) FROM orders o WHERE o.user_id = u.id) as order_count FROM users u; -- Good: JOIN with aggregation SELECT u.name, u.email, COUNT(o.id) as order_count FROM users u LEFT JOIN orders o ON o.user_id = u.id GROUP BY u.id, u.name, u.email; -- Better: Use window functions SELECT DISTINCT ON (u.id) u.name, u.email, COUNT(o.id) OVER (PARTITION BY u.id) as order_count FROM users u LEFT JOIN orders o ON o.user_id = u.id; ``` **Use CTEs for Clarity:** ```sql -- Using Common Table Expressions WITH recent_users AS ( SELECT id, name, email FROM users WHERE created_at > NOW() - INTERVAL '30 days' ), user_order_counts AS ( SELECT user_id, COUNT(*) as order_count FROM orders WHERE created_at > NOW() - INTERVAL '30 days' GROUP BY user_id ) SELECT ru.name, ru.email, COALESCE(uoc.order_count, 0) as orders FROM recent_users ru LEFT JOIN user_order_counts uoc ON ru.id = uoc.user_id; ``` ### Pattern 5: Batch Operations **Batch INSERT:** ```sql -- Bad: Multiple individual inserts INSERT INTO users (name, email) VALUES ('Alice', 'alice@example.com'); INSERT INTO users (name, email) VALUES ('Bob', 'bob@example.com'); INSERT INTO users (name, email) VALUES ('Carol', 'carol@example.com'); -- Good: Batch insert INSERT INTO users (name, email) VALUES ('Alice', 'alice@example.com'), ('Bob', 'bob@example.com'), ('Carol', 'carol@example.com'); -- Better: Use COPY for bulk inserts (PostgreSQL) COPY users (name, email) FROM '/tmp/users.csv' CSV HEADER; ``` **Batch UPDATE:** ```sql -- Bad: Update in loop UPDATE users SET status = 'active' WHERE id = 1; UPDATE users SET status = 'active' WHERE id = 2; -- ... repeat for many IDs -- Good: Single UPDATE with IN clause UPDATE users SET status = 'active' WHERE id IN (1, 2, 3, 4, 5, ...); -- Better: Use temporary table for large batches CREATE TEMP TABLE temp_user_updates (id INT, new_status VARCHAR); INSERT INTO temp_user_updates VALUES (1, 'active'), (2, 'active'), ...; UPDATE users u SET status = t.new_status FROM temp_user_updates t WHERE u.id = t.id; ``` ## Advanced Techniques ### Materialized Views Pre-compute expensive queries. ```sql -- Create materialized view CREATE MATERIALIZED VIEW user_order_summary AS SELECT u.id, u.name, COUNT(o.id) as total_orders, SUM(o.total) as total_spent, MAX(o.created_at) as last_order_date FROM users u LEFT JOIN orders o ON u.id = o.user_id GROUP BY u.id, u.name; -- Add index to materialized view CREATE INDEX idx_user_summary_spent ON user_order_summary(total_spent DESC); -- Refresh materialized view REFRESH MATERIALIZED VIEW user_order_summary; -- Concurrent refresh (PostgreSQL) REFRESH MATERIALIZED VIEW CONCURRENTLY user_order_summary; -- Query materialized view (very fast) SELECT * FROM user_order_summary WHERE total_spent > 1000 ORDER BY total_spent DESC; ``` ### Partitioning Split large tables for better performance. ```sql -- Range partitioning by date (PostgreSQL) CREATE TABLE orders ( id SERIAL, user_id INT, total DECIMAL, created_at TIMESTAMP ) PARTITION BY RANGE (created_at); -- Create partitions CREATE TABLE orders_2024_q1 PARTITION OF orders FOR VALUES FROM ('2024-01-01') TO ('2024-04-01'); CREATE TABLE orders_2024_q2 PARTITION OF orders FOR VALUES FROM ('2024-04-01') TO ('2024-07-01'); -- Queries automatically use appropriate partition SELECT * FROM orders WHERE created_at BETWEEN '2024-02-01' AND '2024-02-28'; -- Only scans orders_2024_q1 partition ``` ### Query Hints and Optimization ```sql -- Force index usage (MySQL) SELECT * FROM users USE INDEX (idx_users_email) WHERE email = 'user@example.com'; -- Parallel query (PostgreSQL) SET max_parallel_workers_per_gather = 4; SELECT * FROM large_table WHERE condition; -- Join hints (PostgreSQL) SET enable_nestloop = OFF; -- Force hash or merge join ``` ## Best Practices 1. **Index Selectively**: Too many indexes slow down writes 2. **Monitor Query Performance**: Use slow query logs 3. **Keep Statistics Updated**: Run ANALYZE regularly 4. **Use Appropriate Data Types**: Smaller types = better performance 5. **Normalize Thoughtfully**: Balance normalization vs performance 6. **Cache Frequently Accessed Data**: Use application-level caching 7. **Connection Pooling**: Reuse database connections 8. **Regular Maintenance**: VACUUM, ANALYZE, rebuild indexes ```sql -- Update statistics ANALYZE users; ANALYZE VERBOSE orders; -- Vacuum (PostgreSQL) VACUUM ANALYZE users; VACUUM FULL users; -- Reclaim space (locks table) -- Reindex REINDEX INDEX idx_users_email; REINDEX TABLE users; ``` ## Common Pitfalls - **Over-Indexing**: Each index slows down INSERT/UPDATE/DELETE - **Unused Indexes**: Waste space and slow writes - **Missing Indexes**: Slow queries, full table scans - **Implicit Type Conversion**: Prevents index usage - **OR Conditions**: Can't use indexes efficiently - **LIKE with Leading Wildcard**: `LIKE '%abc'` can't use index - **Function in WHERE**: Prevents index usage unless functional index exists ## Monitoring Queries ```sql -- Find slow queries (PostgreSQL) SELECT query, calls, total_time, mean_time FROM pg_stat_statements ORDER BY mean_time DESC LIMIT 10; -- Find missing indexes (PostgreSQL) SELECT schemaname, tablename, seq_scan, seq_tup_read, idx_scan, seq_tup_read / seq_scan AS avg_seq_tup_read FROM pg_stat_user_tables WHERE seq_scan > 0 ORDER BY seq_tup_read DESC LIMIT 10; -- Find unused indexes (PostgreSQL) SELECT schemaname, tablename, indexname, idx_scan, idx_tup_read, idx_tup_fetch FROM pg_stat_user_indexes WHERE idx_scan = 0 ORDER BY pg_relation_size(indexrelid) DESC; ``` ## Resources - **references/postgres-optimization-guide.md**: PostgreSQL-specific optimization - **references/mysql-optimization-guide.md**: MySQL/MariaDB optimization - **references/query-plan-analysis.md**: Deep dive into EXPLAIN plans - **assets/index-strategy-checklist.md**: When and how to create indexes - **assets/query-optimization-checklist.md**: Step-by-step optimization guide - **scripts/analyze-slow-queries.sql**: Identify slow queries in your database - **scripts/index-recommendations.sql**: Generate index recommendations
👍0
👁️0
🤖 Auto-discovered
🤖system prompt•7 months ago

distributed-tracing

Implement distributed tracing with Jaeger and Tempo to track

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

async-python-patterns

Master Python asyncio, concurrent programming, and async/await

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

temporal-python-testing

Test Temporal workflows with pytest, time-skipping, and mocking

coding
⭐1
# Temporal Python Testing Strategies Comprehensive testing approaches for Temporal workflows using pytest, progressive disclosure resources for specific testing scenarios. ## When to Use This Skill - **Unit testing workflows** - Fast tests with time-skipping - **Integration testing** - Workflows with mocked activities - **Replay testing** - Validate determinism against production histories - **Local development** - Set up Temporal server and pytest - **CI/CD integration** - Automated testing pipelines - **Coverage strategies** - Achieve ≥80% test coverage ## Testing Philosophy **Recommended Approach** (Source: docs.temporal.io/develop/python/testing-suite): - Write majority as integration tests - Use pytest with async fixtures - Time-skipping enables fast feedback (month-long workflows → seconds) - Mock activities to isolate workflow logic - Validate determinism with replay testing **Three Test Types**: 1. **Unit**: Workflows with time-skipping, activities with ActivityEnvironment 2. **Integration**: Workers with mocked activities 3. **End-to-end**: Full Temporal server with real activities (use sparingly) ## Available Resources This skill provides detailed guidance through progressive disclosure. Load specific resources based on your testing needs: ### Unit Testing Resources **File**: `resources/unit-testing.md` **When to load**: Testing individual workflows or activities in isolation **Contains**: - WorkflowEnvironment with time-skipping - ActivityEnvironment for activity testing - Fast execution of long-running workflows - Manual time advancement patterns - pytest fixtures and patterns ### Integration Testing Resources **File**: `resources/integration-testing.md` **When to load**: Testing workflows with mocked external dependencies **Contains**: - Activity mocking strategies - Error injection patterns - Multi-activity workflow testing - Signal and query testing - Coverage strategies ### Replay Testing Resources **File**: `resources/replay-testing.md` **When to load**: Validating determinism or deploying workflow changes **Contains**: - Determinism validation - Production history replay - CI/CD integration patterns - Version compatibility testing ### Local Development Resources **File**: `resources/local-setup.md` **When to load**: Setting up development environment **Contains**: - Docker Compose configuration - pytest setup and configuration - Coverage tool integration - Development workflow ## Quick Start Guide ### Basic Workflow Test ```python import pytest from temporalio.testing import WorkflowEnvironment from temporalio.worker import Worker @pytest.fixture async def workflow_env(): env = await WorkflowEnvironment.start_time_skipping() yield env await env.shutdown() @pytest.mark.asyncio async def test_workflow(workflow_env): async with Worker( workflow_env.client, task_queue="test-queue", workflows=[YourWorkflow], activities=[your_activity], ): result = await workflow_env.client.execute_workflow( YourWorkflow.run, args, id="test-wf-id", task_queue="test-queue", ) assert result == expected ``` ### Basic Activity Test ```python from temporalio.testing import ActivityEnvironment async def test_activity(): env = ActivityEnvironment() result = await env.run(your_activity, "test-input") assert result == expected_output ``` ## Coverage Targets **Recommended Coverage** (Source: docs.temporal.io best practices): - **Workflows**: ≥80% logic coverage - **Activities**: ≥80% logic coverage - **Integration**: Critical paths with mocked activities - **Replay**: All workflow versions before deployment ## Key Testing Principles 1. **Time-Skipping** - Month-long workflows test in seconds 2. **Mock Activities** - Isolate workflow logic from external dependencies 3. **Replay Testing** - Validate determinism before deployment 4. **High Coverage** - ≥80% target for production workflows 5. **Fast Feedback** - Unit tests run in milliseconds ## How to Use Resources **Load specific resource when needed**: - "Show me unit testing patterns" → Load `resources/unit-testing.md` - "How do I mock activities?" → Load `resources/integration-testing.md` - "Setup local Temporal server" → Load `resources/local-setup.md` - "Validate determinism" → Load `resources/replay-testing.md` ## Additional References - Python SDK Testing: docs.temporal.io/develop/python/testing-suite - Testing Patterns: github.com/temporalio/temporal/blob/main/docs/development/testing.md - Python Samples: github.com/temporalio/samples-python
👍0
👁️0
🤖 Auto-discovered
🤖system prompt•7 months ago

python-code-style

Python code style, linting, formatting, naming conventions, and

coding
⭐1
# Python Code Style & Documentation Consistent code style and clear documentation make codebases maintainable and collaborative. This skill covers modern Python tooling, naming conventions, and documentation standards. ## When to Use This Skill - Setting up linting and formatting for a new project - Writing or reviewing docstrings - Establishing team coding standards - Configuring ruff, mypy, or pyright - Reviewing code for style consistency - Creating project documentation ## Core Concepts ### 1. Automated Formatting Let tools handle formatting debates. Configure once, enforce automatically. ### 2. Consistent Naming Follow PEP 8 conventions with meaningful, descriptive names. ### 3. Documentation as Code Docstrings should be maintained alongside the code they describe. ### 4. Type Annotations Modern Python code should include type hints for all public APIs. ## Quick Start ```bash # Install modern tooling pip install ruff mypy # Configure in pyproject.toml [tool.ruff] line-length = 120 target-version = "py312" # Adjust based on your project's minimum Python version [tool.mypy] strict = true ``` ## Fundamental Patterns ### Pattern 1: Modern Python Tooling Use `ruff` as an all-in-one linter and formatter. It replaces flake8, isort, and black with a single fast tool. ```toml # pyproject.toml [tool.ruff] line-length = 120 target-version = "py312" # Adjust based on your project's minimum Python version [tool.ruff.lint] select = [ "E", # pycodestyle errors "W", # pycodestyle warnings "F", # pyflakes "I", # isort "B", # flake8-bugbear "C4", # flake8-comprehensions "UP", # pyupgrade "SIM", # flake8-simplify ] ignore = ["E501"] # Line length handled by formatter [tool.ruff.format] quote-style = "double" indent-style = "space" ``` Run with: ```bash ruff check --fix . # Lint and auto-fix ruff format . # Format code ``` ### Pattern 2: Type Checking Configuration Configure strict type checking for production code. ```toml # pyproject.toml [tool.mypy] python_version = "3.12" strict = true warn_return_any = true warn_unused_ignores = true disallow_untyped_defs = true disallow_incomplete_defs = true [[tool.mypy.overrides]] module = "tests.*" disallow_untyped_defs = false ``` Alternative: Use `pyright` for faster checking. ```toml [tool.pyright] pythonVersion = "3.12" typeCheckingMode = "strict" ``` ### Pattern 3: Naming Conventions Follow PEP 8 with emphasis on clarity over brevity. **Files and Modules:** ```python # Good: Descriptive snake_case user_repository.py order_processing.py http_client.py # Avoid: Abbreviations usr_repo.py ord_proc.py http_cli.py ``` **Classes and Functions:** ```python # Classes: PascalCase class UserRepository: pass class HTTPClientFactory: # Acronyms stay uppercase pass # Functions and variables: snake_case def get_user_by_email(email: str) -> User | None: retry_count = 3 max_connections = 100 ``` **Constants:** ```python # Module-level constants: SCREAMING_SNAKE_CASE MAX_RETRY_ATTEMPTS = 3 DEFAULT_TIMEOUT_SECONDS = 30 API_BASE_URL = "https://api.example.com" ``` ### Pattern 4: Import Organization Group imports in a consistent order: standard library, third-party, local. ```python # Standard library import os from collections.abc import Callable from typing import Any # Third-party packages import httpx from pydantic import BaseModel from sqlalchemy import Column # Local imports from myproject.models import User from myproject.services import UserService ``` Use absolute imports exclusively: ```python # Preferred from myproject.utils import retry_decorator # Avoid relative imports from ..utils import retry_decorator ``` ## Advanced Patterns ### Pattern 5: Google-Style Docstrings Write docstrings for all public classes, methods, and functions. **Simple Function:** ```python def get_user(user_id: str) -> User: """Retrieve a user by their unique identifier.""" ... ``` **Complex Function:** ```python def process_batch( items: list[Item], max_workers: int = 4, on_progress: Callable[[int, int], None] | None = None, ) -> BatchResult: """Process items concurrently using a worker pool. Processes each item in the batch using the configured number of workers. Progress can be monitored via the optional callback. Args: items: The items to process. Must not be empty. max_workers: Maximum concurrent workers. Defaults to 4. on_progress: Optional callback receiving (completed, total) counts. Returns: BatchResult containing succeeded items and any failures with their associated exceptions. Raises: ValueError: If items is empty. ProcessingError: If the batch cannot be processed. Example: >>> result = process_batch(items, max_workers=8) >>> print(f"Processed {len(result.succeeded)} items") """ ... ``` **Class Docstring:** ```python class UserService: """Service for managing user operations. Provides methods for creating, retrieving, updating, and deleting users with proper validation and error handling. Attributes: repository: The data access layer for user persistence. logger: Logger instance for operation tracking. Example: >>> service = UserService(repository, logger) >>> user = service.create_user(CreateUserInput(...)) """ def __init__(self, repository: UserRepository, logger: Logger) -> None: """Initialize the user service. Args: repository: Data access layer for users. logger: Logger for tracking operations. """ self.repository = repository self.logger = logger ``` ### Pattern 6: Line Length and Formatting Set line length to 120 characters for modern displays while maintaining readability. ```python # Good: Readable line breaks def create_user( email: str, name: str, role: UserRole = UserRole.MEMBER, notify: bool = True, ) -> User: ... # Good: Chain method calls clearly result = ( db.query(User) .filter(User.active == True) .order_by(User.created_at.desc()) .limit(10) .all() ) # Good: Format long strings error_message = ( f"Failed to process user {user_id}: " f"received status {response.status_code} " f"with body {response.text[:100]}" ) ``` ### Pattern 7: Project Documentation **README Structure:** ```markdown # Project Name Brief description of what the project does. ## Installation \`\`\`bash pip install myproject \`\`\` ## Quick Start \`\`\`python from myproject import Client client = Client(api_key="...") result = client.process(data) \`\`\` ## Configuration Document environment variables and configuration options. ## Development \`\`\`bash pip install -e ".[dev]" pytest \`\`\` ``` **CHANGELOG Format (Keep a Changelog):** ```markdown # Changelog ## [Unreleased] ### Added - New feature X ### Changed - Modified behavior of Y ### Fixed - Bug in Z ``` ## Best Practices Summary 1. **Use ruff** - Single tool for linting and formatting 2. **Enable strict mypy** - Catch type errors before runtime 3. **120 character lines** - Modern standard for readability 4. **Descriptive names** - Clarity over brevity 5. **Absolute imports** - More maintainable than relative 6. **Google-style docstrings** - Consistent, readable documentation 7. **Document public APIs** - Every public function needs a docstring 8. **Keep docs updated** - Treat documentation as code 9. **Automate in CI** - Run linters on every commit 10. **Target Python 3.10+** - For new projects, Python 3.12+ is recommended for modern language features
👍0
👁️0
🤖 Auto-discovered
🤖system prompt•7 months ago

python-type-safety

Python type safety with type hints, generics, protocols, and strict

coding
⭐1
# Python Type Safety Leverage Python's type system to catch errors at static analysis time. Type annotations serve as enforced documentation that tooling validates automatically. ## When to Use This Skill - Adding type hints to existing code - Creating generic, reusable classes - Defining structural interfaces with protocols - Configuring mypy or pyright for strict checking - Understanding type narrowing and guards - Building type-safe APIs and libraries ## Core Concepts ### 1. Type Annotations Declare expected types for function parameters, return values, and variables. ### 2. Generics Write reusable code that preserves type information across different types. ### 3. Protocols Define structural interfaces without inheritance (duck typing with type safety). ### 4. Type Narrowing Use guards and conditionals to narrow types within code blocks. ## Quick Start ```python def get_user(user_id: str) -> User | None: """Return type makes 'might not exist' explicit.""" ... # Type checker enforces handling None case user = get_user("123") if user is None: raise UserNotFoundError("123") print(user.name) # Type checker knows user is User here ``` ## Fundamental Patterns ### Pattern 1: Annotate All Public Signatures Every public function, method, and class should have type annotations. ```python def get_user(user_id: str) -> User: """Retrieve user by ID.""" ... def process_batch( items: list[Item], max_workers: int = 4, ) -> BatchResult[ProcessedItem]: """Process items concurrently.""" ... class UserRepository: def __init__(self, db: Database) -> None: self._db = db async def find_by_id(self, user_id: str) -> User | None: """Return User if found, None otherwise.""" ... async def find_by_email(self, email: str) -> User | None: ... async def save(self, user: User) -> User: """Save and return user with generated ID.""" ... ``` Use `mypy --strict` or `pyright` in CI to catch type errors early. For existing projects, enable strict mode incrementally using per-module overrides. ### Pattern 2: Use Modern Union Syntax Python 3.10+ provides cleaner union syntax. ```python # Preferred (3.10+) def find_user(user_id: str) -> User | None: ... def parse_value(v: str) -> int | float | str: ... # Older style (still valid, needed for 3.9) from typing import Optional, Union def find_user(user_id: str) -> Optional[User]: ... ``` ### Pattern 3: Type Narrowing with Guards Use conditionals to narrow types for the type checker. ```python def process_user(user_id: str) -> UserData: user = find_user(user_id) if user is None: raise UserNotFoundError(f"User {user_id} not found") # Type checker knows user is User here, not User | None return UserData( name=user.name, email=user.email, ) def process_items(items: list[Item | None]) -> list[ProcessedItem]: # Filter and narrow types valid_items = [item for item in items if item is not None] # valid_items is now list[Item] return [process(item) for item in valid_items] ``` ### Pattern 4: Generic Classes Create type-safe reusable containers. ```python from typing import TypeVar, Generic T = TypeVar("T") E = TypeVar("E", bound=Exception) class Result(Generic[T, E]): """Represents either a success value or an error.""" def __init__( self, value: T | None = None, error: E | None = None, ) -> None: if (value is None) == (error is None): raise ValueError("Exactly one of value or error must be set") self._value = value self._error = error @property def is_success(self) -> bool: return self._error is None @property def is_failure(self) -> bool: return self._error is not None def unwrap(self) -> T: """Get value or raise the error.""" if self._error is not None: raise self._error return self._value # type: ignore[return-value] def unwrap_or(self, default: T) -> T: """Get value or return default.""" if self._error is not None: return default return self._value # type: ignore[return-value] # Usage preserves types def parse_config(path: str) -> Result[Config, ConfigError]: try: return Result(value=Config.from_file(path)) except ConfigError as e: return Result(error=e) result = parse_config("config.yaml") if result.is_success: config = result.unwrap() # Type: Config ``` ## Advanced Patterns ### Pattern 5: Generic Repository Create type-safe data access patterns. ```python from typing import TypeVar, Generic from abc import ABC, abstractmethod T = TypeVar("T") ID = TypeVar("ID") class Repository(ABC, Generic[T, ID]): """Generic repository interface.""" @abstractmethod async def get(self, id: ID) -> T | None: """Get entity by ID.""" ... @abstractmethod async def save(self, entity: T) -> T: """Save and return entity.""" ... @abstractmethod async def delete(self, id: ID) -> bool: """Delete entity, return True if existed.""" ... class UserRepository(Repository[User, str]): """Concrete repository for Users with string IDs.""" async def get(self, id: str) -> User | None: row = await self._db.fetchrow( "SELECT * FROM users WHERE id = $1", id ) return User(**row) if row else None async def save(self, entity: User) -> User: ... async def delete(self, id: str) -> bool: ... ``` ### Pattern 6: TypeVar with Bounds Restrict generic parameters to specific types. ```python from typing import TypeVar from pydantic import BaseModel ModelT = TypeVar("ModelT", bound=BaseModel) def validate_and_create(model_cls: type[ModelT], data: dict) -> ModelT: """Create a validated Pydantic model from dict.""" return model_cls.model_validate(data) # Works with any BaseModel subclass class User(BaseModel): name: str email: str user = validate_and_create(User, {"name": "Alice", "email": "a@b.com"}) # user is typed as User # Type error: str is not a BaseModel subclass result = validate_and_create(str, {"name": "Alice"}) # Error! ``` ### Pattern 7: Protocols for Structural Typing Define interfaces without requiring inheritance. ```python from typing import Protocol, runtime_checkable @runtime_checkable class Serializable(Protocol): """Any class that can be serialized to/from dict.""" def to_dict(self) -> dict: ... @classmethod def from_dict(cls, data: dict) -> "Serializable": ... # User satisfies Serializable without inheriting from it class User: def __init__(self, id: str, name: str) -> None: self.id = id self.name = name def to_dict(self) -> dict: return {"id": self.id, "name": self.name} @classmethod def from_dict(cls, data: dict) -> "User": return cls(id=data["id"], name=data["name"]) def serialize(obj: Serializable) -> str: """Works with any Serializable object.""" return json.dumps(obj.to_dict()) # Works - User matches the protocol serialize(User("1", "Alice")) # Runtime checking with @runtime_checkable isinstance(User("1", "Alice"), Serializable) # True ``` ### Pattern 8: Common Protocol Patterns Define reusable structural interfaces. ```python from typing import Protocol class Closeable(Protocol): """Resource that can be closed.""" def close(self) -> None: ... class AsyncCloseable(Protocol): """Async resource that can be closed.""" async def close(self) -> None: ... class Readable(Protocol): """Object that can be read from.""" def read(self, n: int = -1) -> bytes: ... class HasId(Protocol): """Object with an ID property.""" @property def id(self) -> str: ... class Comparable(Protocol): """Object that supports comparison.""" def __lt__(self, other: "Comparable") -> bool: ... def __le__(self, other: "Comparable") -> bool: ... ``` ### Pattern 9: Type Aliases Create meaningful type names. **Note:** The `type` statement was introduced in Python 3.10 for simple aliases. Generic type statements require Python 3.12+. ```python # Python 3.10+ type statement for simple aliases type UserId = str type UserDict = dict[str, Any] # Python 3.12+ type statement with generics type Handler[T] = Callable[[Request], T] type AsyncHandler[T] = Callable[[Request], Awaitable[T]] # Python 3.9-3.11 style (needed for broader compatibility) from typing import TypeAlias from collections.abc import Callable, Awaitable UserId: TypeAlias = str Handler: TypeAlias = Callable[[Request], Response] # Usage def register_handler(path: str, handler: Handler[Response]) -> None: ... ``` ### Pattern 10: Callable Types Type function parameters and callbacks. ```python from collections.abc import Callable, Awaitable # Sync callback ProgressCallback = Callable[[int, int], None] # (current, total) # Async callback AsyncHandler = Callable[[Request], Awaitable[Response]] # With named parameters (using Protocol) class OnProgress(Protocol): def __call__( self, current: int, total: int, *, message: str = "", ) -> None: ... def process_items( items: list[Item], on_progress: ProgressCallback | None = None, ) -> list[Result]: for i, item in enumerate(items): if on_progress: on_progress(i, len(items)) ... ``` ## Configuration ### Strict Mode Checklist For `mypy --strict` compliance: ```toml # pyproject.toml [tool.mypy] python_version = "3.12" strict = true warn_return_any = true warn_unused_ignores = true disallow_untyped_defs = true disallow_incomplete_defs = true no_implicit_optional = true ``` Incremental adoption goals: - All function parameters annotated - All return types annotated - Class attributes annotated - Minimize `Any` usage (acceptable for truly dynamic data) - Generic collections use type parameters (`list[str]` not `list`) For existing codebases, enable strict mode per-module using `# mypy: strict` or configure per-module overrides in `pyproject.toml`. ## Best Practices Summary 1. **Annotate all public APIs** - Functions, methods, class attributes 2. **Use `T | None`** - Modern union syntax over `Optional[T]` 3. **Run strict type checking** - `mypy --strict` in CI 4. **Use generics** - Preserve type info in reusable code 5. **Define protocols** - Structural typing for interfaces 6. **Narrow types** - Use guards to help the type checker 7. **Bound type vars** - Restrict generics to meaningful types 8. **Create type aliases** - Meaningful names for complex types 9. **Minimize `Any`** - Use specific types or generics. `Any` is acceptable for truly dynamic data or when interfacing with untyped third-party code 10. **Document with types** - Types are enforceable documentation
👍0
👁️0
🤖 Auto-discovered
🤖system prompt•7 months ago

go-concurrency-patterns

Master Go concurrency with goroutines, channels, sync primitives,

coding
⭐1
# Go Concurrency Patterns Production patterns for Go concurrency including goroutines, channels, synchronization primitives, and context management. ## When to Use This Skill - Building concurrent Go applications - Implementing worker pools and pipelines - Managing goroutine lifecycles - Using channels for communication - Debugging race conditions - Implementing graceful shutdown ## Core Concepts ### 1. Go Concurrency Primitives | Primitive | Purpose | | ----------------- | -------------------------------- | | `goroutine` | Lightweight concurrent execution | | `channel` | Communication between goroutines | | `select` | Multiplex channel operations | | `sync.Mutex` | Mutual exclusion | | `sync.WaitGroup` | Wait for goroutines to complete | | `context.Context` | Cancellation and deadlines | ### 2. Go Concurrency Mantra ``` Don't communicate by sharing memory; share memory by communicating. ``` ## Quick Start ```go package main import ( "context" "fmt" "sync" "time" ) func main() { ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() results := make(chan string, 10) var wg sync.WaitGroup // Spawn workers for i := 0; i < 3; i++ { wg.Add(1) go worker(ctx, i, results, &wg) } // Close results when done go func() { wg.Wait() close(results) }() // Collect results for result := range results { fmt.Println(result) } } func worker(ctx context.Context, id int, results chan<- string, wg *sync.WaitGroup) { defer wg.Done() select { case <-ctx.Done(): return case results <- fmt.Sprintf("Worker %d done", id): } } ``` ## Patterns ### Pattern 1: Worker Pool ```go package main import ( "context" "fmt" "sync" ) type Job struct { ID int Data string } type Result struct { JobID int Output string Err error } func WorkerPool(ctx context.Context, numWorkers int, jobs <-chan Job) <-chan Result { results := make(chan Result, len(jobs)) var wg sync.WaitGroup for i := 0; i < numWorkers; i++ { wg.Add(1) go func(workerID int) { defer wg.Done() for job := range jobs { select { case <-ctx.Done(): return default: result := processJob(job) results <- result } } }(i) } go func() { wg.Wait() close(results) }() return results } func processJob(job Job) Result { // Simulate work return Result{ JobID: job.ID, Output: fmt.Sprintf("Processed: %s", job.Data), } } // Usage func main() { ctx, cancel := context.WithCancel(context.Background()) defer cancel() jobs := make(chan Job, 100) // Send jobs go func() { for i := 0; i < 50; i++ { jobs <- Job{ID: i, Data: fmt.Sprintf("job-%d", i)} } close(jobs) }() // Process with 5 workers results := WorkerPool(ctx, 5, jobs) for result := range results { fmt.Printf("Result: %+v\n", result) } } ``` ### Pattern 2: Fan-Out/Fan-In Pipeline ```go package main import ( "context" "sync" ) // Stage 1: Generate numbers func generate(ctx context.Context, nums ...int) <-chan int { out := make(chan int) go func() { defer close(out) for _, n := range nums { select { case <-ctx.Done(): return case out <- n: } } }() return out } // Stage 2: Square numbers (can run multiple instances) func square(ctx context.Context, in <-chan int) <-chan int { out := make(chan int) go func() { defer close(out) for n := range in { select { case <-ctx.Done(): return case out <- n * n: } } }() return out } // Fan-in: Merge multiple channels into one func merge(ctx context.Context, cs ...<-chan int) <-chan int { var wg sync.WaitGroup out := make(chan int) // Start output goroutine for each input channel output := func(c <-chan int) { defer wg.Done() for n := range c { select { case <-ctx.Done(): return case out <- n: } } } wg.Add(len(cs)) for _, c := range cs { go output(c) } // Close out after all inputs are done go func() { wg.Wait() close(out) }() return out } func main() { ctx, cancel := context.WithCancel(context.Background()) defer cancel() // Generate input in := generate(ctx, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10) // Fan out to multiple squarers c1 := square(ctx, in) c2 := square(ctx, in) c3 := square(ctx, in) // Fan in results for result := range merge(ctx, c1, c2, c3) { fmt.Println(result) } } ``` ### Pattern 3: Bounded Concurrency with Semaphore ```go package main import ( "context" "fmt" "golang.org/x/sync/semaphore" "sync" ) type RateLimitedWorker struct { sem *semaphore.Weighted } func NewRateLimitedWorker(maxConcurrent int64) *RateLimitedWorker { return &RateLimitedWorker{ sem: semaphore.NewWeighted(maxConcurrent), } } func (w *RateLimitedWorker) Do(ctx context.Context, tasks []func() error) []error { var ( wg sync.WaitGroup mu sync.Mutex errors []error ) for _, task := range tasks { // Acquire semaphore (blocks if at limit) if err := w.sem.Acquire(ctx, 1); err != nil { return []error{err} } wg.Add(1) go func(t func() error) { defer wg.Done() defer w.sem.Release(1) if err := t(); err != nil { mu.Lock() errors = append(errors, err) mu.Unlock() } }(task) } wg.Wait() return errors } // Alternative: Channel-based semaphore type Semaphore chan struct{} func NewSemaphore(n int) Semaphore { return make(chan struct{}, n) } func (s Semaphore) Acquire() { s <- struct{}{} } func (s Semaphore) Release() { <-s } ``` ### Pattern 4: Graceful Shutdown ```go package main import ( "context" "fmt" "os" "os/signal" "sync" "syscall" "time" ) type Server struct { shutdown chan struct{} wg sync.WaitGroup } func NewServer() *Server { return &Server{ shutdown: make(chan struct{}), } } func (s *Server) Start(ctx context.Context) { // Start workers for i := 0; i < 5; i++ { s.wg.Add(1) go s.worker(ctx, i) } } func (s *Server) worker(ctx context.Context, id int) { defer s.wg.Done() defer fmt.Printf("Worker %d stopped\n", id) ticker := time.NewTicker(time.Second) defer ticker.Stop() for { select { case <-ctx.Done(): // Cleanup fmt.Printf("Worker %d cleaning up...\n", id) time.Sleep(500 * time.Millisecond) // Simulated cleanup return case <-ticker.C: fmt.Printf("Worker %d working...\n", id) } } } func (s *Server) Shutdown(timeout time.Duration) { // Signal shutdown close(s.shutdown) // Wait with timeout done := make(chan struct{}) go func() { s.wg.Wait() close(done) }() select { case <-done: fmt.Println("Clean shutdown completed") case <-time.After(timeout): fmt.Println("Shutdown timed out, forcing exit") } } func main() { // Setup signal handling ctx, cancel := context.WithCancel(context.Background()) sigCh := make(chan os.Signal, 1) signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM) server := NewServer() server.Start(ctx) // Wait for signal sig := <-sigCh fmt.Printf("\nReceived signal: %v\n", sig) // Cancel context to stop workers cancel() // Wait for graceful shutdown server.Shutdown(5 * time.Second) } ``` ### Pattern 5: Error Group with Cancellation ```go package main import ( "context" "fmt" "golang.org/x/sync/errgroup" "net/http" ) func fetchAllURLs(ctx context.Context, urls []string) ([]string, error) { g, ctx := errgroup.WithContext(ctx) results := make([]string, len(urls)) for i, url := range urls { i, url := i, url // Capture loop variables g.Go(func() error { req, err := http.NewRequestWithContext(ctx, "GET", url, nil) if err != nil { return fmt.Errorf("creating request for %s: %w", url, err) } resp, err := http.DefaultClient.Do(req) if err != nil { return fmt.Errorf("fetching %s: %w", url, err) } defer resp.Body.Close() results[i] = fmt.Sprintf("%s: %d", url, resp.StatusCode) return nil }) } // Wait for all goroutines to complete or one to fail if err := g.Wait(); err != nil { return nil, err // First error cancels all others } return results, nil } // With concurrency limit func fetchWithLimit(ctx context.Context, urls []string, limit int) ([]string, error) { g, ctx := errgroup.WithContext(ctx) g.SetLimit(limit) // Max concurrent goroutines results := make([]string, len(urls)) var mu sync.Mutex for i, url := range urls { i, url := i, url g.Go(func() error { result, err := fetchURL(ctx, url) if err != nil { return err } mu.Lock() results[i] = result mu.Unlock() return nil }) } if err := g.Wait(); err != nil { return nil, err } return results, nil } ``` ### Pattern 6: Concurrent Map with sync.Map ```go package main import ( "sync" ) // For frequent reads, infrequent writes type Cache struct { m sync.Map } func (c *Cache) Get(key string) (interface{}, bool) { return c.m.Load(key) } func (c *Cache) Set(key string, value interface{}) { c.m.Store(key, value) } func (c *Cache) GetOrSet(key string, value interface{}) (interface{}, bool) { return c.m.LoadOrStore(key, value) } func (c *Cache) Delete(key string) { c.m.Delete(key) } // For write-heavy workloads, use sharded map type ShardedMap struct { shards []*shard numShards int } type shard struct { sync.RWMutex data map[string]interface{} } func NewShardedMap(numShards int) *ShardedMap { m := &ShardedMap{ shards: make([]*shard, numShards), numShards: numShards, } for i := range m.shards { m.shards[i] = &shard{data: make(map[string]interface{})} } return m } func (m *ShardedMap) getShard(key string) *shard { // Simple hash h := 0 for _, c := range key { h = 31*h + int(c) } return m.shards[h%m.numShards] } func (m *ShardedMap) Get(key string) (interface{}, bool) { shard := m.getShard(key) shard.RLock() defer shard.RUnlock() v, ok := shard.data[key] return v, ok } func (m *ShardedMap) Set(key string, value interface{}) { shard := m.getShard(key) shard.Lock() defer shard.Unlock() shard.data[key] = value } ``` ### Pattern 7: Select with Timeout and Default ```go func selectPatterns() { ch := make(chan int) // Timeout pattern select { case v := <-ch: fmt.Println("Received:", v) case <-time.After(time.Second): fmt.Println("Timeout!") } // Non-blocking send/receive select { case ch <- 42: fmt.Println("Sent") default: fmt.Println("Channel full, skipping") } // Priority select (check high priority first) highPriority := make(chan int) lowPriority := make(chan int) for { select { case msg := <-highPriority: fmt.Println("High priority:", msg) default: select { case msg := <-highPriority: fmt.Println("High priority:", msg) case msg := <-lowPriority: fmt.Println("Low priority:", msg) } } } } ``` ## Race Detection ```bash # Run tests with race detector go test -race ./... # Build with race detector go build -race . # Run with race detector go run -race main.go ``` ## Best Practices ### Do's - **Use context** - For cancellation and deadlines - **Close channels** - From sender side only - **Use errgroup** - For concurrent operations with errors - **Buffer channels** - When you know the count - **Prefer channels** - Over mutexes when possible ### Don'ts - **Don't leak goroutines** - Always have exit path - **Don't close from receiver** - Causes panic - **Don't use shared memory** - Unless necessary - **Don't ignore context cancellation** - Check ctx.Done() - **Don't use time.Sleep for sync** - Use proper primitives ## Resources - [Go Concurrency Patterns](https://go.dev/blog/pipelines) - [Effective Go - Concurrency](https://go.dev/doc/effective_go#concurrency) - [Go by Example - Goroutines](https://gobyexample.com/goroutines)
👍0
👁️0
🤖 Auto-discovered
🤖system prompt•7 months ago

rust-async-patterns

Master Rust async programming with Tokio, async traits, error

coding
⭐1
# Rust Async Patterns Production patterns for async Rust programming with Tokio runtime, including tasks, channels, streams, and error handling. ## When to Use This Skill - Building async Rust applications - Implementing concurrent network services - Using Tokio for async I/O - Handling async errors properly - Debugging async code issues - Optimizing async performance ## Core Concepts ### 1. Async Execution Model ``` Future (lazy) → poll() → Ready(value) | Pending ↑ ↓ Waker ← Runtime schedules ``` ### 2. Key Abstractions | Concept | Purpose | | ---------- | ---------------------------------------- | | `Future` | Lazy computation that may complete later | | `async fn` | Function returning impl Future | | `await` | Suspend until future completes | | `Task` | Spawned future running concurrently | | `Runtime` | Executor that polls futures | ## Quick Start ```toml # Cargo.toml [dependencies] tokio = { version = "1", features = ["full"] } futures = "0.3" async-trait = "0.1" anyhow = "1.0" tracing = "0.1" tracing-subscriber = "0.3" ``` ```rust use tokio::time::{sleep, Duration}; use anyhow::Result; #[tokio::main] async fn main() -> Result<()> { // Initialize tracing tracing_subscriber::fmt::init(); // Async operations let result = fetch_data("https://api.example.com").await?; println!("Got: {}", result); Ok(()) } async fn fetch_data(url: &str) -> Result<String> { // Simulated async operation sleep(Duration::from_millis(100)).await; Ok(format!("Data from {}", url)) } ``` ## Patterns ### Pattern 1: Concurrent Task Execution ```rust use tokio::task::JoinSet; use anyhow::Result; // Spawn multiple concurrent tasks async fn fetch_all_concurrent(urls: Vec<String>) -> Result<Vec<String>> { let mut set = JoinSet::new(); for url in urls { set.spawn(async move { fetch_data(&url).await }); } let mut results = Vec::new(); while let Some(res) = set.join_next().await { match res { Ok(Ok(data)) => results.push(data), Ok(Err(e)) => tracing::error!("Task failed: {}", e), Err(e) => tracing::error!("Join error: {}", e), } } Ok(results) } // With concurrency limit use futures::stream::{self, StreamExt}; async fn fetch_with_limit(urls: Vec<String>, limit: usize) -> Vec<Result<String>> { stream::iter(urls) .map(|url| async move { fetch_data(&url).await }) .buffer_unordered(limit) // Max concurrent tasks .collect() .await } // Select first to complete use tokio::select; async fn race_requests(url1: &str, url2: &str) -> Result<String> { select! { result = fetch_data(url1) => result, result = fetch_data(url2) => result, } } ``` ### Pattern 2: Channels for Communication ```rust use tokio::sync::{mpsc, broadcast, oneshot, watch}; // Multi-producer, single-consumer async fn mpsc_example() { let (tx, mut rx) = mpsc::channel::<String>(100); // Spawn producer let tx2 = tx.clone(); tokio::spawn(async move { tx2.send("Hello".to_string()).await.unwrap(); }); // Consume while let Some(msg) = rx.recv().await { println!("Got: {}", msg); } } // Broadcast: multi-producer, multi-consumer async fn broadcast_example() { let (tx, _) = broadcast::channel::<String>(100); let mut rx1 = tx.subscribe(); let mut rx2 = tx.subscribe(); tx.send("Event".to_string()).unwrap(); // Both receivers get the message let _ = rx1.recv().await; let _ = rx2.recv().await; } // Oneshot: single value, single use async fn oneshot_example() -> String { let (tx, rx) = oneshot::channel::<String>(); tokio::spawn(async move { tx.send("Result".to_string()).unwrap(); }); rx.await.unwrap() } // Watch: single producer, multi-consumer, latest value async fn watch_example() { let (tx, mut rx) = watch::channel("initial".to_string()); tokio::spawn(async move { loop { // Wait for changes rx.changed().await.unwrap(); println!("New value: {}", *rx.borrow()); } }); tx.send("updated".to_string()).unwrap(); } ``` ### Pattern 3: Async Error Handling ```rust use anyhow::{Context, Result, bail}; use thiserror::Error; #[derive(Error, Debug)] pub enum ServiceError { #[error("Network error: {0}")] Network(#[from] reqwest::Error), #[error("Database error: {0}")] Database(#[from] sqlx::Error), #[error("Not found: {0}")] NotFound(String), #[error("Timeout after {0:?}")] Timeout(std::time::Duration), } // Using anyhow for application errors async fn process_request(id: &str) -> Result<Response> { let data = fetch_data(id) .await .context("Failed to fetch data")?; let parsed = parse_response(&data) .context("Failed to parse response")?; Ok(parsed) } // Using custom errors for library code async fn get_user(id: &str) -> Result<User, ServiceError> { let result = db.query(id).await?; match result { Some(user) => Ok(user), None => Err(ServiceError::NotFound(id.to_string())), } } // Timeout wrapper use tokio::time::timeout; async fn with_timeout<T, F>(duration: Duration, future: F) -> Result<T, ServiceError> where F: std::future::Future<Output = Result<T, ServiceError>>, { timeout(duration, future) .await .map_err(|_| ServiceError::Timeout(duration))? } ``` ### Pattern 4: Graceful Shutdown ```rust use tokio::signal; use tokio::sync::broadcast; use tokio_util::sync::CancellationToken; async fn run_server() -> Result<()> { // Method 1: CancellationToken let token = CancellationToken::new(); let token_clone = token.clone(); // Spawn task that respects cancellation tokio::spawn(async move { loop { tokio::select! { _ = token_clone.cancelled() => { tracing::info!("Task shutting down"); break; } _ = do_work() => {} } } }); // Wait for shutdown signal signal::ctrl_c().await?; tracing::info!("Shutdown signal received"); // Cancel all tasks token.cancel(); // Give tasks time to cleanup tokio::time::sleep(Duration::from_secs(5)).await; Ok(()) } // Method 2: Broadcast channel for shutdown async fn run_with_broadcast() -> Result<()> { let (shutdown_tx, _) = broadcast::channel::<()>(1); let mut rx = shutdown_tx.subscribe(); tokio::spawn(async move { tokio::select! { _ = rx.recv() => { tracing::info!("Received shutdown"); } _ = async { loop { do_work().await } } => {} } }); signal::ctrl_c().await?; let _ = shutdown_tx.send(()); Ok(()) } ``` ### Pattern 5: Async Traits ```rust use async_trait::async_trait; #[async_trait] pub trait Repository { async fn get(&self, id: &str) -> Result<Entity>; async fn save(&self, entity: &Entity) -> Result<()>; async fn delete(&self, id: &str) -> Result<()>; } pub struct PostgresRepository { pool: sqlx::PgPool, } #[async_trait] impl Repository for PostgresRepository { async fn get(&self, id: &str) -> Result<Entity> { sqlx::query_as!(Entity, "SELECT * FROM entities WHERE id = $1", id) .fetch_one(&self.pool) .await .map_err(Into::into) } async fn save(&self, entity: &Entity) -> Result<()> { sqlx::query!( "INSERT INTO entities (id, data) VALUES ($1, $2) ON CONFLICT (id) DO UPDATE SET data = $2", entity.id, entity.data ) .execute(&self.pool) .await?; Ok(()) } async fn delete(&self, id: &str) -> Result<()> { sqlx::query!("DELETE FROM entities WHERE id = $1", id) .execute(&self.pool) .await?; Ok(()) } } // Trait object usage async fn process(repo: &dyn Repository, id: &str) -> Result<()> { let entity = repo.get(id).await?; // Process... repo.save(&entity).await } ``` ### Pattern 6: Streams and Async Iteration ```rust use futures::stream::{self, Stream, StreamExt}; use async_stream::stream; // Create stream from async iterator fn numbers_stream() -> impl Stream<Item = i32> { stream! { for i in 0..10 { tokio::time::sleep(Duration::from_millis(100)).await; yield i; } } } // Process stream async fn process_stream() { let stream = numbers_stream(); // Map and filter let processed: Vec<_> = stream .filter(|n| futures::future::ready(*n % 2 == 0)) .map(|n| n * 2) .collect() .await; println!("{:?}", processed); } // Chunked processing async fn process_in_chunks() { let stream = numbers_stream(); let mut chunks = stream.chunks(3); while let Some(chunk) = chunks.next().await { println!("Processing chunk: {:?}", chunk); } } // Merge multiple streams async fn merge_streams() { let stream1 = numbers_stream(); let stream2 = numbers_stream(); let merged = stream::select(stream1, stream2); merged .for_each(|n| async move { println!("Got: {}", n); }) .await; } ``` ### Pattern 7: Resource Management ```rust use std::sync::Arc; use tokio::sync::{Mutex, RwLock, Semaphore}; // Shared state with RwLock (prefer for read-heavy) struct Cache { data: RwLock<HashMap<String, String>>, } impl Cache { async fn get(&self, key: &str) -> Option<String> { self.data.read().await.get(key).cloned() } async fn set(&self, key: String, value: String) { self.data.write().await.insert(key, value); } } // Connection pool with semaphore struct Pool { semaphore: Semaphore, connections: Mutex<Vec<Connection>>, } impl Pool { fn new(size: usize) -> Self { Self { semaphore: Semaphore::new(size), connections: Mutex::new((0..size).map(|_| Connection::new()).collect()), } } async fn acquire(&self) -> PooledConnection<'_> { let permit = self.semaphore.acquire().await.unwrap(); let conn = self.connections.lock().await.pop().unwrap(); PooledConnection { pool: self, conn: Some(conn), _permit: permit } } } struct PooledConnection<'a> { pool: &'a Pool, conn: Option<Connection>, _permit: tokio::sync::SemaphorePermit<'a>, } impl Drop for PooledConnection<'_> { fn drop(&mut self) { if let Some(conn) = self.conn.take() { let pool = self.pool; tokio::spawn(async move { pool.connections.lock().await.push(conn); }); } } } ``` ## Debugging Tips ```rust // Enable tokio-console for runtime debugging // Cargo.toml: tokio = { features = ["tracing"] } // Run: RUSTFLAGS="--cfg tokio_unstable" cargo run // Then: tokio-console // Instrument async functions use tracing::instrument; #[instrument(skip(pool))] async fn fetch_user(pool: &PgPool, id: &str) -> Result<User> { tracing::debug!("Fetching user"); // ... } // Track task spawning let span = tracing::info_span!("worker", id = %worker_id); tokio::spawn(async move { // Enters span when polled }.instrument(span)); ``` ## Best Practices ### Do's - **Use `tokio::select!`** - For racing futures - **Prefer channels** - Over shared state when possible - **Use `JoinSet`** - For managing multiple tasks - **Instrument with tracing** - For debugging async code - **Handle cancellation** - Check `CancellationToken` ### Don'ts - **Don't block** - Never use `std::thread::sleep` in async - **Don't hold locks across awaits** - Causes deadlocks - **Don't spawn unboundedly** - Use semaphores for limits - **Don't ignore errors** - Propagate with `?` or log - **Don't forget Send bounds** - For spawned futures ## Resources - [Tokio Tutorial](https://tokio.rs/tokio/tutorial) - [Async Book](https://rust-lang.github.io/async-book/) - [Tokio Console](https://github.com/tokio-rs/console)
👍0
👁️0
🤖 Auto-discovered
🤖system prompt•7 months ago

python-background-jobs

Python background job patterns including task queues, workers, and

coding
⭐1
# Python Background Jobs & Task Queues Decouple long-running or unreliable work from request/response cycles. Return immediately to the user while background workers handle the heavy lifting asynchronously. ## When to Use This Skill - Processing tasks that take longer than a few seconds - Sending emails, notifications, or webhooks - Generating reports or exporting data - Processing uploads or media transformations - Integrating with unreliable external services - Building event-driven architectures ## Core Concepts ### 1. Task Queue Pattern API accepts request, enqueues a job, returns immediately with a job ID. Workers process jobs asynchronously. ### 2. Idempotency Tasks may be retried on failure. Design for safe re-execution. ### 3. Job State Machine Jobs transition through states: pending → running → succeeded/failed. ### 4. At-Least-Once Delivery Most queues guarantee at-least-once delivery. Your code must handle duplicates. ## Quick Start This skill uses Celery for examples, a widely adopted task queue. Alternatives like RQ, Dramatiq, and cloud-native solutions (AWS SQS, GCP Tasks) are equally valid choices. ```python from celery import Celery app = Celery("tasks", broker="redis://localhost:6379") @app.task def send_email(to: str, subject: str, body: str) -> None: # This runs in a background worker email_client.send(to, subject, body) # In your API handler send_email.delay("user@example.com", "Welcome!", "Thanks for signing up") ``` ## Fundamental Patterns ### Pattern 1: Return Job ID Immediately For operations exceeding a few seconds, return a job ID and process asynchronously. ```python from uuid import uuid4 from dataclasses import dataclass from enum import Enum from datetime import datetime class JobStatus(Enum): PENDING = "pending" RUNNING = "running" SUCCEEDED = "succeeded" FAILED = "failed" @dataclass class Job: id: str status: JobStatus created_at: datetime started_at: datetime | None = None completed_at: datetime | None = None result: dict | None = None error: str | None = None # API endpoint async def start_export(request: ExportRequest) -> JobResponse: """Start export job and return job ID.""" job_id = str(uuid4()) # Persist job record await jobs_repo.create(Job( id=job_id, status=JobStatus.PENDING, created_at=datetime.utcnow(), )) # Enqueue task for background processing await task_queue.enqueue( "export_data", job_id=job_id, params=request.model_dump(), ) # Return immediately with job ID return JobResponse( job_id=job_id, status="pending", poll_url=f"/jobs/{job_id}", ) ``` ### Pattern 2: Celery Task Configuration Configure Celery tasks with proper retry and timeout settings. ```python from celery import Celery app = Celery("tasks", broker="redis://localhost:6379") # Global configuration app.conf.update( task_time_limit=3600, # Hard limit: 1 hour task_soft_time_limit=3000, # Soft limit: 50 minutes task_acks_late=True, # Acknowledge after completion task_reject_on_worker_lost=True, worker_prefetch_multiplier=1, # Don't prefetch too many tasks ) @app.task( bind=True, max_retries=3, default_retry_delay=60, autoretry_for=(ConnectionError, TimeoutError), ) def process_payment(self, payment_id: str) -> dict: """Process payment with automatic retry on transient errors.""" try: result = payment_gateway.charge(payment_id) return {"status": "success", "transaction_id": result.id} except PaymentDeclinedError as e: # Don't retry permanent failures return {"status": "declined", "reason": str(e)} except TransientError as e: # Retry with exponential backoff raise self.retry(exc=e, countdown=2 ** self.request.retries * 60) ``` ### Pattern 3: Make Tasks Idempotent Workers may retry on crash or timeout. Design for safe re-execution. ```python @app.task(bind=True) def process_order(self, order_id: str) -> None: """Process order idempotently.""" order = orders_repo.get(order_id) # Already processed? Return early if order.status == OrderStatus.COMPLETED: logger.info("Order already processed", order_id=order_id) return # Already in progress? Check if we should continue if order.status == OrderStatus.PROCESSING: # Use idempotency key to avoid double-charging pass # Process with idempotency key result = payment_provider.charge( amount=order.total, idempotency_key=f"order-{order_id}", # Critical! ) orders_repo.update(order_id, status=OrderStatus.COMPLETED) ``` **Idempotency Strategies:** 1. **Check-before-write**: Verify state before action 2. **Idempotency keys**: Use unique tokens with external services 3. **Upsert patterns**: `INSERT ... ON CONFLICT UPDATE` 4. **Deduplication window**: Track processed IDs for N hours ### Pattern 4: Job State Management Persist job state transitions for visibility and debugging. ```python class JobRepository: """Repository for managing job state.""" async def create(self, job: Job) -> Job: """Create new job record.""" await self._db.execute( """INSERT INTO jobs (id, status, created_at) VALUES ($1, $2, $3)""", job.id, job.status.value, job.created_at, ) return job async def update_status( self, job_id: str, status: JobStatus, **fields, ) -> None: """Update job status with timestamp.""" updates = {"status": status.value, **fields} if status == JobStatus.RUNNING: updates["started_at"] = datetime.utcnow() elif status in (JobStatus.SUCCEEDED, JobStatus.FAILED): updates["completed_at"] = datetime.utcnow() await self._db.execute( "UPDATE jobs SET status = $1, ... WHERE id = $2", updates, job_id, ) logger.info( "Job status updated", job_id=job_id, status=status.value, ) ``` ## Advanced Patterns ### Pattern 5: Dead Letter Queue Handle permanently failed tasks for manual inspection. ```python @app.task(bind=True, max_retries=3) def process_webhook(self, webhook_id: str, payload: dict) -> None: """Process webhook with DLQ for failures.""" try: result = send_webhook(payload) if not result.success: raise WebhookFailedError(result.error) except Exception as e: if self.request.retries >= self.max_retries: # Move to dead letter queue for manual inspection dead_letter_queue.send({ "task": "process_webhook", "webhook_id": webhook_id, "payload": payload, "error": str(e), "attempts": self.request.retries + 1, "failed_at": datetime.utcnow().isoformat(), }) logger.error( "Webhook moved to DLQ after max retries", webhook_id=webhook_id, error=str(e), ) return # Exponential backoff retry raise self.retry(exc=e, countdown=2 ** self.request.retries * 60) ``` ### Pattern 6: Status Polling Endpoint Provide an endpoint for clients to check job status. ```python from fastapi import FastAPI, HTTPException app = FastAPI() @app.get("/jobs/{job_id}") async def get_job_status(job_id: str) -> JobStatusResponse: """Get current status of a background job.""" job = await jobs_repo.get(job_id) if job is None: raise HTTPException(404, f"Job {job_id} not found") return JobStatusResponse( job_id=job.id, status=job.status.value, created_at=job.created_at, started_at=job.started_at, completed_at=job.completed_at, result=job.result if job.status == JobStatus.SUCCEEDED else None, error=job.error if job.status == JobStatus.FAILED else None, # Helpful for clients is_terminal=job.status in (JobStatus.SUCCEEDED, JobStatus.FAILED), ) ``` ### Pattern 7: Task Chaining and Workflows Compose complex workflows from simple tasks. ```python from celery import chain, group, chord # Simple chain: A → B → C workflow = chain( extract_data.s(source_id), transform_data.s(), load_data.s(destination_id), ) # Parallel execution: A, B, C all at once parallel = group( send_email.s(user_email), send_sms.s(user_phone), update_analytics.s(event_data), ) # Chord: Run tasks in parallel, then a callback # Process all items, then send completion notification workflow = chord( [process_item.s(item_id) for item_id in item_ids], send_completion_notification.s(batch_id), ) workflow.apply_async() ``` ### Pattern 8: Alternative Task Queues Choose the right tool for your needs. **RQ (Redis Queue)**: Simple, Redis-based ```python from rq import Queue from redis import Redis queue = Queue(connection=Redis()) job = queue.enqueue(send_email, "user@example.com", "Subject", "Body") ``` **Dramatiq**: Modern Celery alternative ```python import dramatiq from dramatiq.brokers.redis import RedisBroker dramatiq.set_broker(RedisBroker()) @dramatiq.actor def send_email(to: str, subject: str, body: str) -> None: email_client.send(to, subject, body) ``` **Cloud-native options:** - AWS SQS + Lambda - Google Cloud Tasks - Azure Functions ## Best Practices Summary 1. **Return immediately** - Don't block requests for long operations 2. **Persist job state** - Enable status polling and debugging 3. **Make tasks idempotent** - Safe to retry on any failure 4. **Use idempotency keys** - For external service calls 5. **Set timeouts** - Both soft and hard limits 6. **Implement DLQ** - Capture permanently failed tasks 7. **Log transitions** - Track job state changes 8. **Retry appropriately** - Exponential backoff for transient errors 9. **Don't retry permanent failures** - Validation errors, invalid credentials 10. **Monitor queue depth** - Alert on backlog growth
👍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