Skip to main content
EVOKORE// BROWSE
>

./browse/prompts

57 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β€’6 months ago

Semantic Session Checkpointing Pattern

Record lightweight semantic checkpoints at planning, risk, failure, and handoff boundaries while keeping heartbeat outside the model loop.

productivity
⭐1
# Semantic Session Checkpointing Pattern Imported from curated first-party documentation sources. ## What this covers Use this pattern when session continuity matters and you need checkpoints that survive model drift or hard crashes. ## Use this when - Capturing handoff context during long sessions - Checkpointing before risky refactors - Separating liveness tracking from agent behavior ## Expected outcomes - Checkpoint timing becomes explicit and repeatable - Session heartbeat stays reliable even if the model stalls - Handoffs capture task, next action, files touched, and risk ## Source synthesis - REVOKORE/docs/MCP-Integration.md ## Dedupe notes Uses the compact REVOKORE integration guide as a distinct session-state pattern that is not already covered by EVOKORE-MCP or AGENT33 docs. ## Source excerpts ### REVOKORE/docs/MCP-Integration.md ## Recommended MCP Pattern Expose a tool that writes a short checkpoint into the active session: - current task - next action - files touched - blocker or risk The simplest implementation is to call: ```powershell REVOKORE/scripts/Write-AiCliCheckpoint.ps1 -Kind note -Message "working on auth bug" ``` From an MCP server, you can call the same script directly or append JSON lines to `REVOKORE_CHECKPOINT_PATH`.
πŸ‘0
πŸ‘οΈ14
docs
πŸ“textβ€’6 months ago

ARCH-AEP Tiered Remediation Cycle

Run a scope-locked remediation cycle that normalizes findings, revalidates them in parallel, and clears tiers with verification evidence.

architecture
⭐1
# ARCH-AEP Tiered Remediation Cycle Imported from curated first-party documentation sources. ## What this covers Use this workflow when you need an orchestrated engineering review cycle with strict scope lock, tiered execution, and audit-grade evidence. ## Use this when - Large remediation efforts spanning multiple PRs - Parallel specialist review with shared backlog authority - Tracking verification evidence through closure ## Expected outcomes - Scope changes are controlled instead of creeping mid-cycle - Severity tiers drive remediation order predictably - Verification evidence is treated as part of done ## Source synthesis - CHELATEDAI/docs/ARCH AGENTIC ENGINEERING AND PLANNING/workflow.md (https://github.com/mattmre/CHELATEDAI/blob/main/docs/ARCH%20AGENTIC%20ENGINEERING%20AND%20PLANNING/workflow.md) - CHELATEDAI/docs/ARCH AGENTIC ENGINEERING AND PLANNING/orchestrator-briefing.md (https://github.com/mattmre/CHELATEDAI/blob/main/docs/ARCH%20AGENTIC%20ENGINEERING%20AND%20PLANNING/orchestrator-briefing.md) ## Dedupe notes Condenses the core ARCH-AEP workflow and briefing docs into a single remediation-cycle entry instead of importing both separately. ## Source excerpts ### CHELATEDAI/docs/ARCH AGENTIC ENGINEERING AND PLANNING/workflow.md ## Workflow (ARCH phase) 1. Scope lock - Define PR range and time window. - Freeze inputs (refinement report + framework + PR list). 2. Discovery + normalization - Orchestrator ingests `agentic-review-framework.md` and the latest refinement report. - Normalize all findings into a single backlog with unique IDs. - De-duplicate, merge overlaps, and assign provisional severity. - Record backlog in `docs/ARCH AGENTIC ENGINEERING AND PLANNING/backlog-YYYY-MM-DD.md`. 3. Parallel re-validation - Spawn specialist agents to re-validate findings against current main. - Each agent must attach exact file paths and line ranges. - Propose the smallest safe fix, acceptance criteria, and effort sizing (S/M/L). 4. Architecture and planning synthesis - Orchestrator merges validated findings into a master backlog. - Re-score and re-rank using the sorting model. - Identify dependency chains and blockers. ### CHELATEDAI/docs/ARCH AGENTIC ENGINEERING AND PLANNING/orchestrator-briefing.md Purpose: Quick reference for the ARCH-AEP documentation set, plus the narrative to start a new session. Full file index: `docs/INDEX.md` ## What Exists In This Folder - `README.md`: narrative overview and intent for ARCH-AEP. - `workflow.md`: end-to-end workflow specification with phases, guardrails, and enhancements. - `templates.md`: ID and branch conventions, tracker table format, and status log. - `backlog-template.md`: template for the master backlog file. - `backlog-index.md`: index of backlog files. - `next-session.md`: short handoff checklist for resuming work. - `phase-planning.md`: long-running planning record for the current cycle. - `schedule-and-tracking.md`: cadence, gates, and milestone tracking. - `agent-learning.md`: cross-session learnings and reusable patterns. - `change-log.md`: scope/defer decisions with audit context. - `risk-memo-template.md`: Critical/High risk memo template. - `risk-memos/README.md`: storage location for risk memos. - `phase-summary-template.md`: template for per-PR or per-phase summaries. - `phase-summaries/README.md`: storage location for phase summaries. - `verification-log.md`: test/build evidence mirror for PRs. - `tracker-pointer.md`: pointer to the active tracker file. - `tracker-index.md`: index of tracker files. - `scope-lock-t ...
πŸ‘0
πŸ‘οΈ0
docs
πŸ€–system promptβ€’6 months ago

Cross-CLI MCP Config Sync

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

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

Hook Observability and Session Replay

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

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

Multi-Server MCP Aggregation Pattern

Aggregate tools across multiple MCP child servers with prefixing, collision avoidance, and routing rules that stay deterministic.

architecture
⭐1
# Multi-Server MCP Aggregation Pattern Imported from curated first-party documentation sources. ## What this covers Use this pattern when an MCP host must broker tools from multiple child servers without sacrificing clarity or control. ## Use this when - Combining tools from several MCP backends - Avoiding tool-name collisions across providers - Keeping origin and routing visible during execution ## Expected outcomes - Server-prefixed names make tool origins obvious - Collisions are avoided without brittle manual renaming - Operators can extend the tool surface without losing determinism ## Source synthesis - EVOKORE-MCP/docs/AGENT33_IMPROVEMENT_INSTRUCTIONS.md (https://github.com/mattmre/EVOKORE-MCP/blob/main/docs/AGENT33_IMPROVEMENT_INSTRUCTIONS.md) - EVOKORE-MCP/docs/TOOLS_AND_DISCOVERY.md (https://github.com/mattmre/EVOKORE-MCP/blob/main/docs/TOOLS_AND_DISCOVERY.md) ## Dedupe notes Uses the improvement-transfer doc as the primary source, with the discovery doc covering prefixing and compatibility details. ## Source excerpts ### EVOKORE-MCP/docs/AGENT33_IMPROVEMENT_INSTRUCTIONS.md > **Purpose**: Feed this file into Claude Code CLI when working on the Agent33 repo. It contains patterns, architectures, and capabilities proven in EVOKORE-MCP that Agent33 should adopt. --- ## 1. Multi-Server MCP Aggregation Pattern **What Agent33 lacks**: Agent33's MCP server (Phase 43) is a single-endpoint bridge. It doesn't aggregate multiple child MCP servers behind a unified namespace. **What to build**: A proxy layer that spawns and manages multiple child MCP servers from a single config file, presenting them as one unified tool surface. ### Implementation spec: ``` mcp.config.json { "servers": { "github": { "command": "npx", "args": ["-y", "@modelcontextprotocol/server-github"], "env": { "GITHUB_TOKEN": "${GITHUB_TOKEN}" } }, "fs": { "command": "npx", "args": ["-y", "@modelcontextprotocol/server-filesystem", "./"] }, "elevenlabs": { "command": "uvx", "args": ["elevenlabs-mcp"], "env": { "ELEVENLABS_API_KEY": "${ELEVENLABS_API_KEY}" } } } } ``` **Key patterns from EVOKORE**: - **Tool name prefixing**: Every proxied tool gets renamed `{serverId}_{originalName}` to prevent namespace collisions (e.g., `github_create_issue`, `fs_read_file`). First-registration-wins for duplicates. - **Environment interpolation**: `${VAR}` syntax in `env` blocks resolved ... ### EVOKORE-MCP/docs/TOOLS_AND_DISCOVERY.md This page explains how EVOKORE presents tools, how proxy names are built, and how `discover_tools` changes the visible tool surface. ## Two tool populations ### Native EVOKORE tools These tools are defined by EVOKORE itself: - `docs_architect` - `skill_creator` - `resolve_workflow` - `search_skills` - `get_skill_help` - `discover_tools` Properties: - always available - always visible - not subject to proxy prefixing ### Proxied child-server tools These come from child servers in `mcp.config.json`. Current configured sources: - `github` - `fs` - optional `elevenlabs` Properties: - fetched from child servers at startup - renamed with server prefixes - governed by `permissions.yml` - routed through `ProxyManager` ## Prefixing and compatibility EVOKORE rewrites proxied tool names to: ```text ${serverId}_${tool.name} ``` Why this exists: - prevents tool-name collisions across child servers - makes origin obvious during execution and review - keeps exact-name routing deterministic Examples: | Upstream tool | EVOKORE-exposed tool | |---|---| | `read_file` from `fs` | `fs_read_file` | | `create_issue` from `github` | `github_create_issue` | ### Duplicate-prefixed name policy If two child registrations would create the same final prefixed name: - the first registrati ...
πŸ‘0
πŸ‘οΈ0
docs
πŸ€–system promptβ€’6 months ago

Human-in-the-Loop Approval Token Workflow

Gate risky tool execution behind approval tokens so agents can retry safely after a human reviewer signs off.

security
⭐1
# Human-in-the-Loop Approval Token Workflow Imported from curated first-party documentation sources. ## What this covers Use this workflow when autonomous execution needs a durable handoff from human review back into the agent loop. ## Use this when - Retrying a blocked tool call after approval - Adding a human checkpoint to sensitive actions - Reducing insecure workarounds around approval flows ## Expected outcomes - Approval becomes a reusable tokenized workflow - Agents resume work without losing execution context - Security controls stay explicit instead of being implied ## Source synthesis - EVOKORE-MCP/docs/V2_MULTI_AGENT_WORKFLOWS.md (https://github.com/mattmre/EVOKORE-MCP/blob/main/docs/V2_MULTI_AGENT_WORKFLOWS.md) - EVOKORE-MCP/docs/AGENT33_IMPROVEMENT_INSTRUCTIONS.md (https://github.com/mattmre/EVOKORE-MCP/blob/main/docs/AGENT33_IMPROVEMENT_INSTRUCTIONS.md) ## Dedupe notes Synthesizes the approval-token concept from the main workflow doc and the Agent33 improvement transfer notes. ## Source excerpts ### EVOKORE-MCP/docs/V2_MULTI_AGENT_WORKFLOWS.md With EVOKORE-MCP v2.0 fully operational, we can now leverage 40+ proxied GitHub and Filesystem tools seamlessly within complex multi-agent workflows. The core architectural advancements include: ## 1. Dynamic Tool Prefixing & Indexing Instead of loading 40+ tools into an LLM's context window statically (which causes massive bloat), EVOKORE's `ProxyManager` boots child servers (like `@modelcontextprotocol/server-github` and `@modelcontextprotocol/server-filesystem`) and dynamically prefixes their tools (`github_create_issue`, `fs_write_file`). This prevents namespace collisions while keeping the tools accessible to native skills. ## 2. Human-in-the-Loop (HITL) Security Interceptor Automated multi-agent workflows involving sensitive endpoints (like GitHub write access or file deletion) are governed by EVOKORE's stateless `_evokore_approval_token` architecture. - When an agent attempts a restricted action, the tool call is intercepted and blocked. - The server returns an error explicitly commanding the agent to prompt the human for approval. - Upon approval, the agent retries the exact tool call with the injected token, securely fulfilling the workflow without severing the conversational context. ## 3. Active Skill Orchestration (Native Harnessing) Unlike v1.0 where skills merel ... ### EVOKORE-MCP/docs/AGENT33_IMPROVEMENT_INSTRUCTIONS.md > **Purpose**: Feed this file into Claude Code CLI when working on the Agent33 repo. It contains patterns, architectures, and capabilities proven in EVOKORE-MCP that Agent33 should adopt. --- ## 1. Multi-Server MCP Aggregation Pattern **What Agent33 lacks**: Agent33's MCP server (Phase 43) is a single-endpoint bridge. It doesn't aggregate multiple child MCP servers behind a unified namespace. **What to build**: A proxy layer that spawns and manages multiple child MCP servers from a single config file, presenting them as one unified tool surface. ### Implementation spec: ``` mcp.config.json { "servers": { "github": { "command": "npx", "args": ["-y", "@modelcontextprotocol/server-github"], "env": { "GITHUB_TOKEN": "${GITHUB_TOKEN}" } }, "fs": { "command": "npx", "args": ["-y", "@modelcontextprotocol/server-filesystem", "./"] }, "elevenlabs": { "command": "uvx", "args": ["elevenlabs-mcp"], "env": { "ELEVENLABS_API_KEY": "${ELEVENLABS_API_KEY}" } } } } ``` **Key patterns from EVOKORE**: - **Tool name prefixing**: Every proxied tool gets renamed `{serverId}_{originalName}` to prevent namespace collisions (e.g., `github_create_issue`, `fs_read_file`). First-registration-wins for duplicates. - **Environment interpolation**: `${VAR}` syntax in `env` blocks resolved ...
πŸ‘0
πŸ‘οΈ0
docs
πŸ“textβ€’6 months ago

Release Command Center Workflow

Coordinate release freezes, candidate promotion, validation, publication, and rollback handling through a defined state model.

devops
⭐1
# Release Command Center Workflow Imported from curated first-party documentation sources. ## What this covers Use this workflow to move a release from planning to publication with explicit validation checkpoints and rollback branches. ## Use this when - Managing release freezes and candidate promotion - Tracking validation state before publication - Documenting rollback-friendly release operations ## Expected outcomes - Release state changes are visible and repeatable - Rollback paths are treated as first-class outcomes - Validation discipline is preserved between phases ## Source synthesis - AGENT33/docs/functionality-and-workflows.md (https://github.com/mattmre/AGENT33/blob/main/docs/functionality-and-workflows.md) - AGENT33/docs/use-cases.md (https://github.com/mattmre/AGENT33/blob/main/docs/use-cases.md) ## Dedupe notes Merges AGENT33 release lifecycle details with release-oriented use-case guidance. ## Source excerpts ### AGENT33/docs/functionality-and-workflows.md ### 4.2 Release Lifecycle States: - `planned -> frozen -> rc -> validating -> released` - Failure/rollback branches: `failed`, `rolled_back` Main APIs: - `/v1/releases/{id}/freeze` - `/v1/releases/{id}/rc` - `/v1/releases/{id}/validate` - `/v1/releases/{id}/publish` - `/v1/releases/{id}/rollback` ### AGENT33/docs/use-cases.md ## 2. Release Command Center Goal: - Move releases through frozen -> RC -> validate -> publish with checklist controls. Use these modules: - `api/routes/releases.py` - `release/service.py` - `release/checklist.py` - `release/sync.py` - `release/rollback.py` Typical flow: 1. Create release (`POST /v1/releases`). 2. Freeze, cut RC, validate. 3. Update/evaluate checklist items. 4. Publish when checks pass. 5. Run sync dry-runs or real syncs. 6. Initiate rollback if needed. Best fit: - Teams with repeatable release compliance requirements.
πŸ‘0
πŸ‘οΈ0
docs
πŸ€–system promptβ€’7 months ago

workflow-orchestration-patterns

Design durable workflows with Temporal for distributed systems.

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

dbt-transformation-patterns

Master dbt (data build tool) for analytics engineering with model

data
⭐1
# dbt Transformation Patterns Production-ready patterns for dbt (data build tool) including model organization, testing strategies, documentation, and incremental processing. ## When to Use This Skill - Building data transformation pipelines with dbt - Organizing models into staging, intermediate, and marts layers - Implementing data quality tests - Creating incremental models for large datasets - Documenting data models and lineage - Setting up dbt project structure ## Core Concepts ### 1. Model Layers (Medallion Architecture) ``` sources/ Raw data definitions ↓ staging/ 1:1 with source, light cleaning ↓ intermediate/ Business logic, joins, aggregations ↓ marts/ Final analytics tables ``` ### 2. Naming Conventions | Layer | Prefix | Example | | ------------ | -------------- | ----------------------------- | | Staging | `stg_` | `stg_stripe__payments` | | Intermediate | `int_` | `int_payments_pivoted` | | Marts | `dim_`, `fct_` | `dim_customers`, `fct_orders` | ## Quick Start ```yaml # dbt_project.yml name: "analytics" version: "1.0.0" profile: "analytics" model-paths: ["models"] analysis-paths: ["analyses"] test-paths: ["tests"] seed-paths: ["seeds"] macro-paths: ["macros"] vars: start_date: "2020-01-01" models: analytics: staging: +materialized: view +schema: staging intermediate: +materialized: ephemeral marts: +materialized: table +schema: analytics ``` ``` # Project structure models/ β”œβ”€β”€ staging/ β”‚ β”œβ”€β”€ stripe/ β”‚ β”‚ β”œβ”€β”€ _stripe__sources.yml β”‚ β”‚ β”œβ”€β”€ _stripe__models.yml β”‚ β”‚ β”œβ”€β”€ stg_stripe__customers.sql β”‚ β”‚ └── stg_stripe__payments.sql β”‚ └── shopify/ β”‚ β”œβ”€β”€ _shopify__sources.yml β”‚ └── stg_shopify__orders.sql β”œβ”€β”€ intermediate/ β”‚ └── finance/ β”‚ └── int_payments_pivoted.sql └── marts/ β”œβ”€β”€ core/ β”‚ β”œβ”€β”€ _core__models.yml β”‚ β”œβ”€β”€ dim_customers.sql β”‚ └── fct_orders.sql └── finance/ └── fct_revenue.sql ``` ## Patterns ### Pattern 1: Source Definitions ```yaml # models/staging/stripe/_stripe__sources.yml version: 2 sources: - name: stripe description: Raw Stripe data loaded via Fivetran database: raw schema: stripe loader: fivetran loaded_at_field: _fivetran_synced freshness: warn_after: { count: 12, period: hour } error_after: { count: 24, period: hour } tables: - name: customers description: Stripe customer records columns: - name: id description: Primary key tests: - unique - not_null - name: email description: Customer email - name: created description: Account creation timestamp - name: payments description: Stripe payment transactions columns: - name: id tests: - unique - not_null - name: customer_id tests: - not_null - relationships: to: source('stripe', 'customers') field: id ``` ### Pattern 2: Staging Models ```sql -- models/staging/stripe/stg_stripe__customers.sql with source as ( select * from {{ source('stripe', 'customers') }} ), renamed as ( select -- ids id as customer_id, -- strings lower(email) as email, name as customer_name, -- timestamps created as created_at, -- metadata _fivetran_synced as _loaded_at from source ) select * from renamed ``` ```sql -- models/staging/stripe/stg_stripe__payments.sql {{ config( materialized='incremental', unique_key='payment_id', on_schema_change='append_new_columns' ) }} with source as ( select * from {{ source('stripe', 'payments') }} {% if is_incremental() %} where _fivetran_synced > (select max(_loaded_at) from {{ this }}) {% endif %} ), renamed as ( select -- ids id as payment_id, customer_id, invoice_id, -- amounts (convert cents to dollars) amount / 100.0 as amount, amount_refunded / 100.0 as amount_refunded, -- status status as payment_status, -- timestamps created as created_at, -- metadata _fivetran_synced as _loaded_at from source ) select * from renamed ``` ### Pattern 3: Intermediate Models ```sql -- models/intermediate/finance/int_payments_pivoted_to_customer.sql with payments as ( select * from {{ ref('stg_stripe__payments') }} ), customers as ( select * from {{ ref('stg_stripe__customers') }} ), payment_summary as ( select customer_id, count(*) as total_payments, count(case when payment_status = 'succeeded' then 1 end) as successful_payments, sum(case when payment_status = 'succeeded' then amount else 0 end) as total_amount_paid, min(created_at) as first_payment_at, max(created_at) as last_payment_at from payments group by customer_id ) select customers.customer_id, customers.email, customers.created_at as customer_created_at, coalesce(payment_summary.total_payments, 0) as total_payments, coalesce(payment_summary.successful_payments, 0) as successful_payments, coalesce(payment_summary.total_amount_paid, 0) as lifetime_value, payment_summary.first_payment_at, payment_summary.last_payment_at from customers left join payment_summary using (customer_id) ``` ### Pattern 4: Mart Models (Dimensions and Facts) ```sql -- models/marts/core/dim_customers.sql {{ config( materialized='table', unique_key='customer_id' ) }} with customers as ( select * from {{ ref('int_payments_pivoted_to_customer') }} ), orders as ( select * from {{ ref('stg_shopify__orders') }} ), order_summary as ( select customer_id, count(*) as total_orders, sum(total_price) as total_order_value, min(created_at) as first_order_at, max(created_at) as last_order_at from orders group by customer_id ), final as ( select -- surrogate key {{ dbt_utils.generate_surrogate_key(['customers.customer_id']) }} as customer_key, -- natural key customers.customer_id, -- attributes customers.email, customers.customer_created_at, -- payment metrics customers.total_payments, customers.successful_payments, customers.lifetime_value, customers.first_payment_at, customers.last_payment_at, -- order metrics coalesce(order_summary.total_orders, 0) as total_orders, coalesce(order_summary.total_order_value, 0) as total_order_value, order_summary.first_order_at, order_summary.last_order_at, -- calculated fields case when customers.lifetime_value >= 1000 then 'high' when customers.lifetime_value >= 100 then 'medium' else 'low' end as customer_tier, -- timestamps current_timestamp as _loaded_at from customers left join order_summary using (customer_id) ) select * from final ``` ```sql -- models/marts/core/fct_orders.sql {{ config( materialized='incremental', unique_key='order_id', incremental_strategy='merge' ) }} with orders as ( select * from {{ ref('stg_shopify__orders') }} {% if is_incremental() %} where updated_at > (select max(updated_at) from {{ this }}) {% endif %} ), customers as ( select * from {{ ref('dim_customers') }} ), final as ( select -- keys orders.order_id, customers.customer_key, orders.customer_id, -- dimensions orders.order_status, orders.fulfillment_status, orders.payment_status, -- measures orders.subtotal, orders.tax, orders.shipping, orders.total_price, orders.total_discount, orders.item_count, -- timestamps orders.created_at, orders.updated_at, orders.fulfilled_at, -- metadata current_timestamp as _loaded_at from orders left join customers on orders.customer_id = customers.customer_id ) select * from final ``` ### Pattern 5: Testing and Documentation ```yaml # models/marts/core/_core__models.yml version: 2 models: - name: dim_customers description: Customer dimension with payment and order metrics columns: - name: customer_key description: Surrogate key for the customer dimension tests: - unique - not_null - name: customer_id description: Natural key from source system tests: - unique - not_null - name: email description: Customer email address tests: - not_null - name: customer_tier description: Customer value tier based on lifetime value tests: - accepted_values: values: ["high", "medium", "low"] - name: lifetime_value description: Total amount paid by customer tests: - dbt_utils.expression_is_true: expression: ">= 0" - name: fct_orders description: Order fact table with all order transactions tests: - dbt_utils.recency: datepart: day field: created_at interval: 1 columns: - name: order_id tests: - unique - not_null - name: customer_key tests: - not_null - relationships: to: ref('dim_customers') field: customer_key ``` ### Pattern 6: Macros and DRY Code ```sql -- macros/cents_to_dollars.sql {% macro cents_to_dollars(column_name, precision=2) %} round({{ column_name }} / 100.0, {{ precision }}) {% endmacro %} -- macros/generate_schema_name.sql {% macro generate_schema_name(custom_schema_name, node) %} {%- set default_schema = target.schema -%} {%- if custom_schema_name is none -%} {{ default_schema }} {%- else -%} {{ default_schema }}_{{ custom_schema_name }} {%- endif -%} {% endmacro %} -- macros/limit_data_in_dev.sql {% macro limit_data_in_dev(column_name, days=3) %} {% if target.name == 'dev' %} where {{ column_name }} >= dateadd(day, -{{ days }}, current_date) {% endif %} {% endmacro %} -- Usage in model select * from {{ ref('stg_orders') }} {{ limit_data_in_dev('created_at') }} ``` ### Pattern 7: Incremental Strategies ```sql -- Delete+Insert (default for most warehouses) {{ config( materialized='incremental', unique_key='id', incremental_strategy='delete+insert' ) }} -- Merge (best for late-arriving data) {{ config( materialized='incremental', unique_key='id', incremental_strategy='merge', merge_update_columns=['status', 'amount', 'updated_at'] ) }} -- Insert Overwrite (partition-based) {{ config( materialized='incremental', incremental_strategy='insert_overwrite', partition_by={ "field": "created_date", "data_type": "date", "granularity": "day" } ) }} select *, date(created_at) as created_date from {{ ref('stg_events') }} {% if is_incremental() %} where created_date >= dateadd(day, -3, current_date) {% endif %} ``` ## dbt Commands ```bash # Development dbt run # Run all models dbt run --select staging # Run staging models only dbt run --select +fct_orders # Run fct_orders and its upstream dbt run --select fct_orders+ # Run fct_orders and its downstream dbt run --full-refresh # Rebuild incremental models # Testing dbt test # Run all tests dbt test --select stg_stripe # Test specific models dbt build # Run + test in DAG order # Documentation dbt docs generate # Generate docs dbt docs serve # Serve docs locally # Debugging dbt compile # Compile SQL without running dbt debug # Test connection dbt ls --select tag:critical # List models by tag ``` ## Best Practices ### Do's - **Use staging layer** - Clean data once, use everywhere - **Test aggressively** - Not null, unique, relationships - **Document everything** - Column descriptions, model descriptions - **Use incremental** - For tables > 1M rows - **Version control** - dbt project in Git ### Don'ts - **Don't skip staging** - Raw β†’ mart is tech debt - **Don't hardcode dates** - Use `{{ var('start_date') }}` - **Don't repeat logic** - Extract to macros - **Don't test in prod** - Use dev target - **Don't ignore freshness** - Monitor source data ## Resources - [dbt Documentation](https://docs.getdbt.com/) - [dbt Best Practices](https://docs.getdbt.com/guides/best-practices) - [dbt-utils Package](https://hub.getdbt.com/dbt-labs/dbt_utils/latest/) - [dbt Discourse](https://discourse.getdbt.com/)
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

data-quality-frameworks

Implement data quality validation with Great Expectations, dbt

data
⭐1
# Data Quality Frameworks Production patterns for implementing data quality with Great Expectations, dbt tests, and data contracts to ensure reliable data pipelines. ## When to Use This Skill - Implementing data quality checks in pipelines - Setting up Great Expectations validation - Building comprehensive dbt test suites - Establishing data contracts between teams - Monitoring data quality metrics - Automating data validation in CI/CD ## Core Concepts ### 1. Data Quality Dimensions | Dimension | Description | Example Check | | ---------------- | ------------------------ | -------------------------------------------------- | | **Completeness** | No missing values | `expect_column_values_to_not_be_null` | | **Uniqueness** | No duplicates | `expect_column_values_to_be_unique` | | **Validity** | Values in expected range | `expect_column_values_to_be_in_set` | | **Accuracy** | Data matches reality | Cross-reference validation | | **Consistency** | No contradictions | `expect_column_pair_values_A_to_be_greater_than_B` | | **Timeliness** | Data is recent | `expect_column_max_to_be_between` | ### 2. Testing Pyramid for Data ``` /\ / \ Integration Tests (cross-table) /────\ / \ Unit Tests (single column) /────────\ / \ Schema Tests (structure) /────────────\ ``` ## Quick Start ### Great Expectations Setup ```bash # Install pip install great_expectations # Initialize project great_expectations init # Create datasource great_expectations datasource new ``` ```python # great_expectations/checkpoints/daily_validation.yml import great_expectations as gx # Create context context = gx.get_context() # Create expectation suite suite = context.add_expectation_suite("orders_suite") # Add expectations suite.add_expectation( gx.expectations.ExpectColumnValuesToNotBeNull(column="order_id") ) suite.add_expectation( gx.expectations.ExpectColumnValuesToBeUnique(column="order_id") ) # Validate results = context.run_checkpoint(checkpoint_name="daily_orders") ``` ## Patterns ### Pattern 1: Great Expectations Suite ```python # expectations/orders_suite.py import great_expectations as gx from great_expectations.core import ExpectationSuite from great_expectations.core.expectation_configuration import ExpectationConfiguration def build_orders_suite() -> ExpectationSuite: """Build comprehensive orders expectation suite""" suite = ExpectationSuite(expectation_suite_name="orders_suite") # Schema expectations suite.add_expectation(ExpectationConfiguration( expectation_type="expect_table_columns_to_match_set", kwargs={ "column_set": ["order_id", "customer_id", "amount", "status", "created_at"], "exact_match": False # Allow additional columns } )) # Primary key suite.add_expectation(ExpectationConfiguration( expectation_type="expect_column_values_to_not_be_null", kwargs={"column": "order_id"} )) suite.add_expectation(ExpectationConfiguration( expectation_type="expect_column_values_to_be_unique", kwargs={"column": "order_id"} )) # Foreign key suite.add_expectation(ExpectationConfiguration( expectation_type="expect_column_values_to_not_be_null", kwargs={"column": "customer_id"} )) # Categorical values suite.add_expectation(ExpectationConfiguration( expectation_type="expect_column_values_to_be_in_set", kwargs={ "column": "status", "value_set": ["pending", "processing", "shipped", "delivered", "cancelled"] } )) # Numeric ranges suite.add_expectation(ExpectationConfiguration( expectation_type="expect_column_values_to_be_between", kwargs={ "column": "amount", "min_value": 0, "max_value": 100000, "strict_min": True # amount > 0 } )) # Date validity suite.add_expectation(ExpectationConfiguration( expectation_type="expect_column_values_to_be_dateutil_parseable", kwargs={"column": "created_at"} )) # Freshness - data should be recent suite.add_expectation(ExpectationConfiguration( expectation_type="expect_column_max_to_be_between", kwargs={ "column": "created_at", "min_value": {"$PARAMETER": "now - timedelta(days=1)"}, "max_value": {"$PARAMETER": "now"} } )) # Row count sanity suite.add_expectation(ExpectationConfiguration( expectation_type="expect_table_row_count_to_be_between", kwargs={ "min_value": 1000, # Expect at least 1000 rows "max_value": 10000000 } )) # Statistical expectations suite.add_expectation(ExpectationConfiguration( expectation_type="expect_column_mean_to_be_between", kwargs={ "column": "amount", "min_value": 50, "max_value": 500 } )) return suite ``` ### Pattern 2: Great Expectations Checkpoint ```yaml # great_expectations/checkpoints/orders_checkpoint.yml name: orders_checkpoint config_version: 1.0 class_name: Checkpoint run_name_template: "%Y%m%d-%H%M%S-orders-validation" validations: - batch_request: datasource_name: warehouse data_connector_name: default_inferred_data_connector_name data_asset_name: orders data_connector_query: index: -1 # Latest batch expectation_suite_name: orders_suite action_list: - name: store_validation_result action: class_name: StoreValidationResultAction - name: store_evaluation_parameters action: class_name: StoreEvaluationParametersAction - name: update_data_docs action: class_name: UpdateDataDocsAction # Slack notification on failure - name: send_slack_notification action: class_name: SlackNotificationAction slack_webhook: ${SLACK_WEBHOOK} notify_on: failure renderer: module_name: great_expectations.render.renderer.slack_renderer class_name: SlackRenderer ``` ```python # Run checkpoint import great_expectations as gx context = gx.get_context() result = context.run_checkpoint(checkpoint_name="orders_checkpoint") if not result.success: failed_expectations = [ r for r in result.run_results.values() if not r.success ] raise ValueError(f"Data quality check failed: {failed_expectations}") ``` ### Pattern 3: dbt Data Tests ```yaml # models/marts/core/_core__models.yml version: 2 models: - name: fct_orders description: Order fact table tests: # Table-level tests - dbt_utils.recency: datepart: day field: created_at interval: 1 - dbt_utils.at_least_one - dbt_utils.expression_is_true: expression: "total_amount >= 0" columns: - name: order_id description: Primary key tests: - unique - not_null - name: customer_id description: Foreign key to dim_customers tests: - not_null - relationships: to: ref('dim_customers') field: customer_id - name: order_status tests: - accepted_values: values: ["pending", "processing", "shipped", "delivered", "cancelled"] - name: total_amount tests: - not_null - dbt_utils.expression_is_true: expression: ">= 0" - name: created_at tests: - not_null - dbt_utils.expression_is_true: expression: "<= current_timestamp" - name: dim_customers columns: - name: customer_id tests: - unique - not_null - name: email tests: - unique - not_null # Custom regex test - dbt_utils.expression_is_true: expression: "email ~ '^[A-Za-z0-9._%+-]+@[A-Za-z0-9.-]+\\.[A-Za-z]{2,}$'" ``` ### Pattern 4: Custom dbt Tests ```sql -- tests/generic/test_row_count_in_range.sql {% test row_count_in_range(model, min_count, max_count) %} with row_count as ( select count(*) as cnt from {{ model }} ) select cnt from row_count where cnt < {{ min_count }} or cnt > {{ max_count }} {% endtest %} -- Usage in schema.yml: -- tests: -- - row_count_in_range: -- min_count: 1000 -- max_count: 10000000 ``` ```sql -- tests/generic/test_sequential_values.sql {% test sequential_values(model, column_name, interval=1) %} with lagged as ( select {{ column_name }}, lag({{ column_name }}) over (order by {{ column_name }}) as prev_value from {{ model }} ) select * from lagged where {{ column_name }} - prev_value != {{ interval }} and prev_value is not null {% endtest %} ``` ```sql -- tests/singular/assert_orders_customers_match.sql -- Singular test: specific business rule with orders_customers as ( select distinct customer_id from {{ ref('fct_orders') }} ), dim_customers as ( select customer_id from {{ ref('dim_customers') }} ), orphaned_orders as ( select o.customer_id from orders_customers o left join dim_customers c using (customer_id) where c.customer_id is null ) select * from orphaned_orders -- Test passes if this returns 0 rows ``` ### Pattern 5: Data Contracts ```yaml # contracts/orders_contract.yaml apiVersion: datacontract.com/v1.0.0 kind: DataContract metadata: name: orders version: 1.0.0 owner: data-platform-team contact: data-team@company.com info: title: Orders Data Contract description: Contract for order event data from the ecommerce platform purpose: Analytics, reporting, and ML features servers: production: type: snowflake account: company.us-east-1 database: ANALYTICS schema: CORE terms: usage: Internal analytics only limitations: PII must not be exposed in downstream marts billing: Charged per query TB scanned schema: type: object properties: order_id: type: string format: uuid description: Unique order identifier required: true unique: true pii: false customer_id: type: string format: uuid description: Customer identifier required: true pii: true piiClassification: indirect total_amount: type: number minimum: 0 maximum: 100000 description: Order total in USD created_at: type: string format: date-time description: Order creation timestamp required: true status: type: string enum: [pending, processing, shipped, delivered, cancelled] description: Current order status quality: type: SodaCL specification: checks for orders: - row_count > 0 - missing_count(order_id) = 0 - duplicate_count(order_id) = 0 - invalid_count(status) = 0: valid values: [pending, processing, shipped, delivered, cancelled] - freshness(created_at) < 24h sla: availability: 99.9% freshness: 1 hour latency: 5 minutes ``` ### Pattern 6: Automated Quality Pipeline ```python # quality_pipeline.py from dataclasses import dataclass from typing import List, Dict, Any import great_expectations as gx from datetime import datetime @dataclass class QualityResult: table: str passed: bool total_expectations: int failed_expectations: int details: List[Dict[str, Any]] timestamp: datetime class DataQualityPipeline: """Orchestrate data quality checks across tables""" def __init__(self, context: gx.DataContext): self.context = context self.results: List[QualityResult] = [] def validate_table(self, table: str, suite: str) -> QualityResult: """Validate a single table against expectation suite""" checkpoint_config = { "name": f"{table}_validation", "config_version": 1.0, "class_name": "Checkpoint", "validations": [{ "batch_request": { "datasource_name": "warehouse", "data_asset_name": table, }, "expectation_suite_name": suite, }], } result = self.context.run_checkpoint(**checkpoint_config) # Parse results validation_result = list(result.run_results.values())[0] results = validation_result.results failed = [r for r in results if not r.success] return QualityResult( table=table, passed=result.success, total_expectations=len(results), failed_expectations=len(failed), details=[{ "expectation": r.expectation_config.expectation_type, "success": r.success, "observed_value": r.result.get("observed_value"), } for r in results], timestamp=datetime.now() ) def run_all(self, tables: Dict[str, str]) -> Dict[str, QualityResult]: """Run validation for all tables""" results = {} for table, suite in tables.items(): print(f"Validating {table}...") results[table] = self.validate_table(table, suite) return results def generate_report(self, results: Dict[str, QualityResult]) -> str: """Generate quality report""" report = ["# Data Quality Report", f"Generated: {datetime.now()}", ""] total_passed = sum(1 for r in results.values() if r.passed) total_tables = len(results) report.append(f"## Summary: {total_passed}/{total_tables} tables passed") report.append("") for table, result in results.items(): status = "βœ…" if result.passed else "❌" report.append(f"### {status} {table}") report.append(f"- Expectations: {result.total_expectations}") report.append(f"- Failed: {result.failed_expectations}") if not result.passed: report.append("- Failed checks:") for detail in result.details: if not detail["success"]: report.append(f" - {detail['expectation']}: {detail['observed_value']}") report.append("") return "\n".join(report) # Usage context = gx.get_context() pipeline = DataQualityPipeline(context) tables_to_validate = { "orders": "orders_suite", "customers": "customers_suite", "products": "products_suite", } results = pipeline.run_all(tables_to_validate) report = pipeline.generate_report(results) # Fail pipeline if any table failed if not all(r.passed for r in results.values()): print(report) raise ValueError("Data quality checks failed!") ``` ## Best Practices ### Do's - **Test early** - Validate source data before transformations - **Test incrementally** - Add tests as you find issues - **Document expectations** - Clear descriptions for each test - **Alert on failures** - Integrate with monitoring - **Version contracts** - Track schema changes ### Don'ts - **Don't test everything** - Focus on critical columns - **Don't ignore warnings** - They often precede failures - **Don't skip freshness** - Stale data is bad data - **Don't hardcode thresholds** - Use dynamic baselines - **Don't test in isolation** - Test relationships too ## Resources - [Great Expectations Documentation](https://docs.greatexpectations.io/) - [dbt Testing Documentation](https://docs.getdbt.com/docs/build/tests) - [Data Contract Specification](https://datacontract.com/) - [Soda Core](https://docs.soda.io/soda-core/overview.html)
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

spark-optimization

Optimize Apache Spark jobs with partitioning, caching, shuffle

data
⭐1
# Apache Spark Optimization Production patterns for optimizing Apache Spark jobs including partitioning strategies, memory management, shuffle optimization, and performance tuning. ## When to Use This Skill - Optimizing slow Spark jobs - Tuning memory and executor configuration - Implementing efficient partitioning strategies - Debugging Spark performance issues - Scaling Spark pipelines for large datasets - Reducing shuffle and data skew ## Core Concepts ### 1. Spark Execution Model ``` Driver Program ↓ Job (triggered by action) ↓ Stages (separated by shuffles) ↓ Tasks (one per partition) ``` ### 2. Key Performance Factors | Factor | Impact | Solution | | ----------------- | --------------------- | ----------------------------- | | **Shuffle** | Network I/O, disk I/O | Minimize wide transformations | | **Data Skew** | Uneven task duration | Salting, broadcast joins | | **Serialization** | CPU overhead | Use Kryo, columnar formats | | **Memory** | GC pressure, spills | Tune executor memory | | **Partitions** | Parallelism | Right-size partitions | ## Quick Start ```python from pyspark.sql import SparkSession from pyspark.sql import functions as F # Create optimized Spark session spark = (SparkSession.builder .appName("OptimizedJob") .config("spark.sql.adaptive.enabled", "true") .config("spark.sql.adaptive.coalescePartitions.enabled", "true") .config("spark.sql.adaptive.skewJoin.enabled", "true") .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer") .config("spark.sql.shuffle.partitions", "200") .getOrCreate()) # Read with optimized settings df = (spark.read .format("parquet") .option("mergeSchema", "false") .load("s3://bucket/data/")) # Efficient transformations result = (df .filter(F.col("date") >= "2024-01-01") .select("id", "amount", "category") .groupBy("category") .agg(F.sum("amount").alias("total"))) result.write.mode("overwrite").parquet("s3://bucket/output/") ``` ## Patterns ### Pattern 1: Optimal Partitioning ```python # Calculate optimal partition count def calculate_partitions(data_size_gb: float, partition_size_mb: int = 128) -> int: """ Optimal partition size: 128MB - 256MB Too few: Under-utilization, memory pressure Too many: Task scheduling overhead """ return max(int(data_size_gb * 1024 / partition_size_mb), 1) # Repartition for even distribution df_repartitioned = df.repartition(200, "partition_key") # Coalesce to reduce partitions (no shuffle) df_coalesced = df.coalesce(100) # Partition pruning with predicate pushdown df = (spark.read.parquet("s3://bucket/data/") .filter(F.col("date") == "2024-01-01")) # Spark pushes this down # Write with partitioning for future queries (df.write .partitionBy("year", "month", "day") .mode("overwrite") .parquet("s3://bucket/partitioned_output/")) ``` ### Pattern 2: Join Optimization ```python from pyspark.sql import functions as F from pyspark.sql.types import * # 1. Broadcast Join - Small table joins # Best when: One side < 10MB (configurable) small_df = spark.read.parquet("s3://bucket/small_table/") # < 10MB large_df = spark.read.parquet("s3://bucket/large_table/") # TBs # Explicit broadcast hint result = large_df.join( F.broadcast(small_df), on="key", how="left" ) # 2. Sort-Merge Join - Default for large tables # Requires shuffle, but handles any size result = large_df1.join(large_df2, on="key", how="inner") # 3. Bucket Join - Pre-sorted, no shuffle at join time # Write bucketed tables (df.write .bucketBy(200, "customer_id") .sortBy("customer_id") .mode("overwrite") .saveAsTable("bucketed_orders")) # Join bucketed tables (no shuffle!) orders = spark.table("bucketed_orders") customers = spark.table("bucketed_customers") # Same bucket count result = orders.join(customers, on="customer_id") # 4. Skew Join Handling # Enable AQE skew join optimization spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true") spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionFactor", "5") spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes", "256MB") # Manual salting for severe skew def salt_join(df_skewed, df_other, key_col, num_salts=10): """Add salt to distribute skewed keys""" # Add salt to skewed side df_salted = df_skewed.withColumn( "salt", (F.rand() * num_salts).cast("int") ).withColumn( "salted_key", F.concat(F.col(key_col), F.lit("_"), F.col("salt")) ) # Explode other side with all salts df_exploded = df_other.crossJoin( spark.range(num_salts).withColumnRenamed("id", "salt") ).withColumn( "salted_key", F.concat(F.col(key_col), F.lit("_"), F.col("salt")) ) # Join on salted key return df_salted.join(df_exploded, on="salted_key", how="inner") ``` ### Pattern 3: Caching and Persistence ```python from pyspark import StorageLevel # Cache when reusing DataFrame multiple times df = spark.read.parquet("s3://bucket/data/") df_filtered = df.filter(F.col("status") == "active") # Cache in memory (MEMORY_AND_DISK is default) df_filtered.cache() # Or with specific storage level df_filtered.persist(StorageLevel.MEMORY_AND_DISK_SER) # Force materialization df_filtered.count() # Use in multiple actions agg1 = df_filtered.groupBy("category").count() agg2 = df_filtered.groupBy("region").sum("amount") # Unpersist when done df_filtered.unpersist() # Storage levels explained: # MEMORY_ONLY - Fast, but may not fit # MEMORY_AND_DISK - Spills to disk if needed (recommended) # MEMORY_ONLY_SER - Serialized, less memory, more CPU # DISK_ONLY - When memory is tight # OFF_HEAP - Tungsten off-heap memory # Checkpoint for complex lineage spark.sparkContext.setCheckpointDir("s3://bucket/checkpoints/") df_complex = (df .join(other_df, "key") .groupBy("category") .agg(F.sum("amount"))) df_complex.checkpoint() # Breaks lineage, materializes ``` ### Pattern 4: Memory Tuning ```python # Executor memory configuration # spark-submit --executor-memory 8g --executor-cores 4 # Memory breakdown (8GB executor): # - spark.memory.fraction = 0.6 (60% = 4.8GB for execution + storage) # - spark.memory.storageFraction = 0.5 (50% of 4.8GB = 2.4GB for cache) # - Remaining 2.4GB for execution (shuffles, joins, sorts) # - 40% = 3.2GB for user data structures and internal metadata spark = (SparkSession.builder .config("spark.executor.memory", "8g") .config("spark.executor.memoryOverhead", "2g") # For non-JVM memory .config("spark.memory.fraction", "0.6") .config("spark.memory.storageFraction", "0.5") .config("spark.sql.shuffle.partitions", "200") # For memory-intensive operations .config("spark.sql.autoBroadcastJoinThreshold", "50MB") # Prevent OOM on large shuffles .config("spark.sql.files.maxPartitionBytes", "128MB") .getOrCreate()) # Monitor memory usage def print_memory_usage(spark): """Print current memory usage""" sc = spark.sparkContext for executor in sc._jsc.sc().getExecutorMemoryStatus().keySet().toArray(): mem_status = sc._jsc.sc().getExecutorMemoryStatus().get(executor) total = mem_status._1() / (1024**3) free = mem_status._2() / (1024**3) print(f"{executor}: {total:.2f}GB total, {free:.2f}GB free") ``` ### Pattern 5: Shuffle Optimization ```python # Reduce shuffle data size spark.conf.set("spark.sql.shuffle.partitions", "auto") # With AQE spark.conf.set("spark.shuffle.compress", "true") spark.conf.set("spark.shuffle.spill.compress", "true") # Pre-aggregate before shuffle df_optimized = (df # Local aggregation first (combiner) .groupBy("key", "partition_col") .agg(F.sum("value").alias("partial_sum")) # Then global aggregation .groupBy("key") .agg(F.sum("partial_sum").alias("total"))) # Avoid shuffle with map-side operations # BAD: Shuffle for each distinct distinct_count = df.select("category").distinct().count() # GOOD: Approximate distinct (no shuffle) approx_count = df.select(F.approx_count_distinct("category")).collect()[0][0] # Use coalesce instead of repartition when reducing partitions df_reduced = df.coalesce(10) # No shuffle # Optimize shuffle with compression spark.conf.set("spark.io.compression.codec", "lz4") # Fast compression ``` ### Pattern 6: Data Format Optimization ```python # Parquet optimizations (df.write .option("compression", "snappy") # Fast compression .option("parquet.block.size", 128 * 1024 * 1024) # 128MB row groups .parquet("s3://bucket/output/")) # Column pruning - only read needed columns df = (spark.read.parquet("s3://bucket/data/") .select("id", "amount", "date")) # Spark only reads these columns # Predicate pushdown - filter at storage level df = (spark.read.parquet("s3://bucket/partitioned/year=2024/") .filter(F.col("status") == "active")) # Pushed to Parquet reader # Delta Lake optimizations (df.write .format("delta") .option("optimizeWrite", "true") # Bin-packing .option("autoCompact", "true") # Compact small files .mode("overwrite") .save("s3://bucket/delta_table/")) # Z-ordering for multi-dimensional queries spark.sql(""" OPTIMIZE delta.`s3://bucket/delta_table/` ZORDER BY (customer_id, date) """) ``` ### Pattern 7: Monitoring and Debugging ```python # Enable detailed metrics spark.conf.set("spark.sql.codegen.wholeStage", "true") spark.conf.set("spark.sql.execution.arrow.pyspark.enabled", "true") # Explain query plan df.explain(mode="extended") # Modes: simple, extended, codegen, cost, formatted # Get physical plan statistics df.explain(mode="cost") # Monitor task metrics def analyze_stage_metrics(spark): """Analyze recent stage metrics""" status_tracker = spark.sparkContext.statusTracker() for stage_id in status_tracker.getActiveStageIds(): stage_info = status_tracker.getStageInfo(stage_id) print(f"Stage {stage_id}:") print(f" Tasks: {stage_info.numTasks}") print(f" Completed: {stage_info.numCompletedTasks}") print(f" Failed: {stage_info.numFailedTasks}") # Identify data skew def check_partition_skew(df): """Check for partition skew""" partition_counts = (df .withColumn("partition_id", F.spark_partition_id()) .groupBy("partition_id") .count() .orderBy(F.desc("count"))) partition_counts.show(20) stats = partition_counts.select( F.min("count").alias("min"), F.max("count").alias("max"), F.avg("count").alias("avg"), F.stddev("count").alias("stddev") ).collect()[0] skew_ratio = stats["max"] / stats["avg"] print(f"Skew ratio: {skew_ratio:.2f}x (>2x indicates skew)") ``` ## Configuration Cheat Sheet ```python # Production configuration template spark_configs = { # Adaptive Query Execution (AQE) "spark.sql.adaptive.enabled": "true", "spark.sql.adaptive.coalescePartitions.enabled": "true", "spark.sql.adaptive.skewJoin.enabled": "true", # Memory "spark.executor.memory": "8g", "spark.executor.memoryOverhead": "2g", "spark.memory.fraction": "0.6", "spark.memory.storageFraction": "0.5", # Parallelism "spark.sql.shuffle.partitions": "200", "spark.default.parallelism": "200", # Serialization "spark.serializer": "org.apache.spark.serializer.KryoSerializer", "spark.sql.execution.arrow.pyspark.enabled": "true", # Compression "spark.io.compression.codec": "lz4", "spark.shuffle.compress": "true", # Broadcast "spark.sql.autoBroadcastJoinThreshold": "50MB", # File handling "spark.sql.files.maxPartitionBytes": "128MB", "spark.sql.files.openCostInBytes": "4MB", } ``` ## Best Practices ### Do's - **Enable AQE** - Adaptive query execution handles many issues - **Use Parquet/Delta** - Columnar formats with compression - **Broadcast small tables** - Avoid shuffle for small joins - **Monitor Spark UI** - Check for skew, spills, GC - **Right-size partitions** - 128MB - 256MB per partition ### Don'ts - **Don't collect large data** - Keep data distributed - **Don't use UDFs unnecessarily** - Use built-in functions - **Don't over-cache** - Memory is limited - **Don't ignore data skew** - It dominates job time - **Don't use `.count()` for existence** - Use `.take(1)` or `.isEmpty()` ## Resources - [Spark Performance Tuning](https://spark.apache.org/docs/latest/sql-performance-tuning.html) - [Spark Configuration](https://spark.apache.org/docs/latest/configuration.html) - [Databricks Optimization Guide](https://docs.databricks.com/en/optimizations/index.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

projection-patterns

Build read models and projections from event streams. Use when

coding
⭐1
# Projection Patterns Comprehensive guide to building projections and read models for event-sourced systems. ## When to Use This Skill - Building CQRS read models - Creating materialized views from events - Optimizing query performance - Implementing real-time dashboards - Building search indexes from events - Aggregating data across streams ## Core Concepts ### 1. Projection Architecture ``` β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ Event Store │────►│ Projector │────►│ Read Model β”‚ β”‚ β”‚ β”‚ β”‚ β”‚ (Database) β”‚ β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ β”‚ β”‚ Events β”‚ β”‚ β”‚ β”‚ Handler β”‚ β”‚ β”‚ β”‚ Tables β”‚ β”‚ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β”‚ β”‚ Logic β”‚ β”‚ β”‚ β”‚ Views β”‚ β”‚ β”‚ β”‚ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β”‚ β”‚ Cache β”‚ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ ``` ### 2. Projection Types | Type | Description | Use Case | | -------------- | --------------------------- | ---------------------- | | **Live** | Real-time from subscription | Current state queries | | **Catchup** | Process historical events | Rebuilding read models | | **Persistent** | Stores checkpoint | Resume after restart | | **Inline** | Same transaction as write | Strong consistency | ## Templates ### Template 1: Basic Projector ```python from abc import ABC, abstractmethod from dataclasses import dataclass from typing import Dict, Any, Callable, List import asyncpg @dataclass class Event: stream_id: str event_type: str data: dict version: int global_position: int class Projection(ABC): """Base class for projections.""" @property @abstractmethod def name(self) -> str: """Unique projection name for checkpointing.""" pass @abstractmethod def handles(self) -> List[str]: """List of event types this projection handles.""" pass @abstractmethod async def apply(self, event: Event) -> None: """Apply event to the read model.""" pass class Projector: """Runs projections from event store.""" def __init__(self, event_store, checkpoint_store): self.event_store = event_store self.checkpoint_store = checkpoint_store self.projections: List[Projection] = [] def register(self, projection: Projection): self.projections.append(projection) async def run(self, batch_size: int = 100): """Run all projections continuously.""" while True: for projection in self.projections: await self._run_projection(projection, batch_size) await asyncio.sleep(0.1) async def _run_projection(self, projection: Projection, batch_size: int): checkpoint = await self.checkpoint_store.get(projection.name) position = checkpoint or 0 events = await self.event_store.read_all(position, batch_size) for event in events: if event.event_type in projection.handles(): await projection.apply(event) await self.checkpoint_store.save( projection.name, event.global_position ) async def rebuild(self, projection: Projection): """Rebuild a projection from scratch.""" await self.checkpoint_store.delete(projection.name) # Optionally clear read model tables await self._run_projection(projection, batch_size=1000) ``` ### Template 2: Order Summary Projection ```python class OrderSummaryProjection(Projection): """Projects order events to a summary read model.""" def __init__(self, db_pool: asyncpg.Pool): self.pool = db_pool @property def name(self) -> str: return "order_summary" def handles(self) -> List[str]: return [ "OrderCreated", "OrderItemAdded", "OrderItemRemoved", "OrderShipped", "OrderCompleted", "OrderCancelled" ] async def apply(self, event: Event) -> None: handlers = { "OrderCreated": self._handle_created, "OrderItemAdded": self._handle_item_added, "OrderItemRemoved": self._handle_item_removed, "OrderShipped": self._handle_shipped, "OrderCompleted": self._handle_completed, "OrderCancelled": self._handle_cancelled, } handler = handlers.get(event.event_type) if handler: await handler(event) async def _handle_created(self, event: Event): async with self.pool.acquire() as conn: await conn.execute( """ INSERT INTO order_summaries (order_id, customer_id, status, total_amount, item_count, created_at) VALUES ($1, $2, $3, $4, $5, $6) """, event.data['order_id'], event.data['customer_id'], 'pending', 0, 0, event.data['created_at'] ) async def _handle_item_added(self, event: Event): async with self.pool.acquire() as conn: await conn.execute( """ UPDATE order_summaries SET total_amount = total_amount + $2, item_count = item_count + 1, updated_at = NOW() WHERE order_id = $1 """, event.data['order_id'], event.data['price'] * event.data['quantity'] ) async def _handle_item_removed(self, event: Event): async with self.pool.acquire() as conn: await conn.execute( """ UPDATE order_summaries SET total_amount = total_amount - $2, item_count = item_count - 1, updated_at = NOW() WHERE order_id = $1 """, event.data['order_id'], event.data['price'] * event.data['quantity'] ) async def _handle_shipped(self, event: Event): async with self.pool.acquire() as conn: await conn.execute( """ UPDATE order_summaries SET status = 'shipped', shipped_at = $2, updated_at = NOW() WHERE order_id = $1 """, event.data['order_id'], event.data['shipped_at'] ) async def _handle_completed(self, event: Event): async with self.pool.acquire() as conn: await conn.execute( """ UPDATE order_summaries SET status = 'completed', completed_at = $2, updated_at = NOW() WHERE order_id = $1 """, event.data['order_id'], event.data['completed_at'] ) async def _handle_cancelled(self, event: Event): async with self.pool.acquire() as conn: await conn.execute( """ UPDATE order_summaries SET status = 'cancelled', cancelled_at = $2, cancellation_reason = $3, updated_at = NOW() WHERE order_id = $1 """, event.data['order_id'], event.data['cancelled_at'], event.data.get('reason') ) ``` ### Template 3: Elasticsearch Search Projection ```python from elasticsearch import AsyncElasticsearch class ProductSearchProjection(Projection): """Projects product events to Elasticsearch for full-text search.""" def __init__(self, es_client: AsyncElasticsearch): self.es = es_client self.index = "products" @property def name(self) -> str: return "product_search" def handles(self) -> List[str]: return [ "ProductCreated", "ProductUpdated", "ProductPriceChanged", "ProductDeleted" ] async def apply(self, event: Event) -> None: if event.event_type == "ProductCreated": await self.es.index( index=self.index, id=event.data['product_id'], document={ 'name': event.data['name'], 'description': event.data['description'], 'category': event.data['category'], 'price': event.data['price'], 'tags': event.data.get('tags', []), 'created_at': event.data['created_at'] } ) elif event.event_type == "ProductUpdated": await self.es.update( index=self.index, id=event.data['product_id'], doc={ 'name': event.data['name'], 'description': event.data['description'], 'category': event.data['category'], 'tags': event.data.get('tags', []), 'updated_at': event.data['updated_at'] } ) elif event.event_type == "ProductPriceChanged": await self.es.update( index=self.index, id=event.data['product_id'], doc={ 'price': event.data['new_price'], 'price_updated_at': event.data['changed_at'] } ) elif event.event_type == "ProductDeleted": await self.es.delete( index=self.index, id=event.data['product_id'] ) ``` ### Template 4: Aggregating Projection ```python class DailySalesProjection(Projection): """Aggregates sales data by day for reporting.""" def __init__(self, db_pool: asyncpg.Pool): self.pool = db_pool @property def name(self) -> str: return "daily_sales" def handles(self) -> List[str]: return ["OrderCompleted", "OrderRefunded"] async def apply(self, event: Event) -> None: if event.event_type == "OrderCompleted": await self._increment_sales(event) elif event.event_type == "OrderRefunded": await self._decrement_sales(event) async def _increment_sales(self, event: Event): date = event.data['completed_at'][:10] # YYYY-MM-DD async with self.pool.acquire() as conn: await conn.execute( """ INSERT INTO daily_sales (date, total_orders, total_revenue, total_items) VALUES ($1, 1, $2, $3) ON CONFLICT (date) DO UPDATE SET total_orders = daily_sales.total_orders + 1, total_revenue = daily_sales.total_revenue + $2, total_items = daily_sales.total_items + $3, updated_at = NOW() """, date, event.data['total_amount'], event.data['item_count'] ) async def _decrement_sales(self, event: Event): date = event.data['original_completed_at'][:10] async with self.pool.acquire() as conn: await conn.execute( """ UPDATE daily_sales SET total_orders = total_orders - 1, total_revenue = total_revenue - $2, total_refunds = total_refunds + $2, updated_at = NOW() WHERE date = $1 """, date, event.data['refund_amount'] ) ``` ### Template 5: Multi-Table Projection ```python class CustomerActivityProjection(Projection): """Projects customer activity across multiple tables.""" def __init__(self, db_pool: asyncpg.Pool): self.pool = db_pool @property def name(self) -> str: return "customer_activity" def handles(self) -> List[str]: return [ "CustomerCreated", "OrderCompleted", "ReviewSubmitted", "CustomerTierChanged" ] async def apply(self, event: Event) -> None: async with self.pool.acquire() as conn: async with conn.transaction(): if event.event_type == "CustomerCreated": # Insert into customers table await conn.execute( """ INSERT INTO customers (customer_id, email, name, tier, created_at) VALUES ($1, $2, $3, 'bronze', $4) """, event.data['customer_id'], event.data['email'], event.data['name'], event.data['created_at'] ) # Initialize activity summary await conn.execute( """ INSERT INTO customer_activity_summary (customer_id, total_orders, total_spent, total_reviews) VALUES ($1, 0, 0, 0) """, event.data['customer_id'] ) elif event.event_type == "OrderCompleted": # Update activity summary await conn.execute( """ UPDATE customer_activity_summary SET total_orders = total_orders + 1, total_spent = total_spent + $2, last_order_at = $3 WHERE customer_id = $1 """, event.data['customer_id'], event.data['total_amount'], event.data['completed_at'] ) # Insert into order history await conn.execute( """ INSERT INTO customer_order_history (customer_id, order_id, amount, completed_at) VALUES ($1, $2, $3, $4) """, event.data['customer_id'], event.data['order_id'], event.data['total_amount'], event.data['completed_at'] ) elif event.event_type == "ReviewSubmitted": await conn.execute( """ UPDATE customer_activity_summary SET total_reviews = total_reviews + 1, last_review_at = $2 WHERE customer_id = $1 """, event.data['customer_id'], event.data['submitted_at'] ) elif event.event_type == "CustomerTierChanged": await conn.execute( """ UPDATE customers SET tier = $2, updated_at = NOW() WHERE customer_id = $1 """, event.data['customer_id'], event.data['new_tier'] ) ``` ## Best Practices ### Do's - **Make projections idempotent** - Safe to replay - **Use transactions** - For multi-table updates - **Store checkpoints** - Resume after failures - **Monitor lag** - Alert on projection delays - **Plan for rebuilds** - Design for reconstruction ### Don'ts - **Don't couple projections** - Each is independent - **Don't skip error handling** - Log and alert on failures - **Don't ignore ordering** - Events must be processed in order - **Don't over-normalize** - Denormalize for query patterns ## Resources - [CQRS Pattern](https://docs.microsoft.com/en-us/azure/architecture/patterns/cqrs) - [Projection Building Blocks](https://zimarev.com/blog/event-sourcing/projections/)
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

api-design-principles

Master REST and GraphQL API design principles to build intuitive,

coding
⭐1
# API Design Principles Master REST and GraphQL API design principles to build intuitive, scalable, and maintainable APIs that delight developers and stand the test of time. ## When to Use This Skill - Designing new REST or GraphQL APIs - Refactoring existing APIs for better usability - Establishing API design standards for your team - Reviewing API specifications before implementation - Migrating between API paradigms (REST to GraphQL, etc.) - Creating developer-friendly API documentation - Optimizing APIs for specific use cases (mobile, third-party integrations) ## Core Concepts ### 1. RESTful Design Principles **Resource-Oriented Architecture** - Resources are nouns (users, orders, products), not verbs - Use HTTP methods for actions (GET, POST, PUT, PATCH, DELETE) - URLs represent resource hierarchies - Consistent naming conventions **HTTP Methods Semantics:** - `GET`: Retrieve resources (idempotent, safe) - `POST`: Create new resources - `PUT`: Replace entire resource (idempotent) - `PATCH`: Partial resource updates - `DELETE`: Remove resources (idempotent) ### 2. GraphQL Design Principles **Schema-First Development** - Types define your domain model - Queries for reading data - Mutations for modifying data - Subscriptions for real-time updates **Query Structure:** - Clients request exactly what they need - Single endpoint, multiple operations - Strongly typed schema - Introspection built-in ### 3. API Versioning Strategies **URL Versioning:** ``` /api/v1/users /api/v2/users ``` **Header Versioning:** ``` Accept: application/vnd.api+json; version=1 ``` **Query Parameter Versioning:** ``` /api/users?version=1 ``` ## REST API Design Patterns ### Pattern 1: Resource Collection Design ```python # Good: Resource-oriented endpoints GET /api/users # List users (with pagination) POST /api/users # Create user GET /api/users/{id} # Get specific user PUT /api/users/{id} # Replace user PATCH /api/users/{id} # Update user fields DELETE /api/users/{id} # Delete user # Nested resources GET /api/users/{id}/orders # Get user's orders POST /api/users/{id}/orders # Create order for user # Bad: Action-oriented endpoints (avoid) POST /api/createUser POST /api/getUserById POST /api/deleteUser ``` ### Pattern 2: Pagination and Filtering ```python from typing import List, Optional from pydantic import BaseModel, Field class PaginationParams(BaseModel): page: int = Field(1, ge=1, description="Page number") page_size: int = Field(20, ge=1, le=100, description="Items per page") class FilterParams(BaseModel): status: Optional[str] = None created_after: Optional[str] = None search: Optional[str] = None class PaginatedResponse(BaseModel): items: List[dict] total: int page: int page_size: int pages: int @property def has_next(self) -> bool: return self.page < self.pages @property def has_prev(self) -> bool: return self.page > 1 # FastAPI endpoint example from fastapi import FastAPI, Query, Depends app = FastAPI() @app.get("/api/users", response_model=PaginatedResponse) async def list_users( page: int = Query(1, ge=1), page_size: int = Query(20, ge=1, le=100), status: Optional[str] = Query(None), search: Optional[str] = Query(None) ): # Apply filters query = build_query(status=status, search=search) # Count total total = await count_users(query) # Fetch page offset = (page - 1) * page_size users = await fetch_users(query, limit=page_size, offset=offset) return PaginatedResponse( items=users, total=total, page=page, page_size=page_size, pages=(total + page_size - 1) // page_size ) ``` ### Pattern 3: Error Handling and Status Codes ```python from fastapi import HTTPException, status from pydantic import BaseModel class ErrorResponse(BaseModel): error: str message: str details: Optional[dict] = None timestamp: str path: str class ValidationErrorDetail(BaseModel): field: str message: str value: Any # Consistent error responses STATUS_CODES = { "success": 200, "created": 201, "no_content": 204, "bad_request": 400, "unauthorized": 401, "forbidden": 403, "not_found": 404, "conflict": 409, "unprocessable": 422, "internal_error": 500 } def raise_not_found(resource: str, id: str): raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail={ "error": "NotFound", "message": f"{resource} not found", "details": {"id": id} } ) def raise_validation_error(errors: List[ValidationErrorDetail]): raise HTTPException( status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, detail={ "error": "ValidationError", "message": "Request validation failed", "details": {"errors": [e.dict() for e in errors]} } ) # Example usage @app.get("/api/users/{user_id}") async def get_user(user_id: str): user = await fetch_user(user_id) if not user: raise_not_found("User", user_id) return user ``` ### Pattern 4: HATEOAS (Hypermedia as the Engine of Application State) ```python class UserResponse(BaseModel): id: str name: str email: str _links: dict @classmethod def from_user(cls, user: User, base_url: str): return cls( id=user.id, name=user.name, email=user.email, _links={ "self": {"href": f"{base_url}/api/users/{user.id}"}, "orders": {"href": f"{base_url}/api/users/{user.id}/orders"}, "update": { "href": f"{base_url}/api/users/{user.id}", "method": "PATCH" }, "delete": { "href": f"{base_url}/api/users/{user.id}", "method": "DELETE" } } ) ``` ## GraphQL Design Patterns ### Pattern 1: Schema Design ```graphql # schema.graphql # Clear type definitions type User { id: ID! email: String! name: String! createdAt: DateTime! # Relationships orders(first: Int = 20, after: String, status: OrderStatus): OrderConnection! profile: UserProfile } type Order { id: ID! status: OrderStatus! total: Money! items: [OrderItem!]! createdAt: DateTime! # Back-reference user: User! } # Pagination pattern (Relay-style) type OrderConnection { edges: [OrderEdge!]! pageInfo: PageInfo! totalCount: Int! } type OrderEdge { node: Order! cursor: String! } type PageInfo { hasNextPage: Boolean! hasPreviousPage: Boolean! startCursor: String endCursor: String } # Enums for type safety enum OrderStatus { PENDING CONFIRMED SHIPPED DELIVERED CANCELLED } # Custom scalars scalar DateTime scalar Money # Query root type Query { user(id: ID!): User users(first: Int = 20, after: String, search: String): UserConnection! order(id: ID!): Order } # Mutation root type Mutation { createUser(input: CreateUserInput!): CreateUserPayload! updateUser(input: UpdateUserInput!): UpdateUserPayload! deleteUser(id: ID!): DeleteUserPayload! createOrder(input: CreateOrderInput!): CreateOrderPayload! } # Input types for mutations input CreateUserInput { email: String! name: String! password: String! } # Payload types for mutations type CreateUserPayload { user: User errors: [Error!] } type Error { field: String message: String! } ``` ### Pattern 2: Resolver Design ```python from typing import Optional, List from ariadne import QueryType, MutationType, ObjectType from dataclasses import dataclass query = QueryType() mutation = MutationType() user_type = ObjectType("User") @query.field("user") async def resolve_user(obj, info, id: str) -> Optional[dict]: """Resolve single user by ID.""" return await fetch_user_by_id(id) @query.field("users") async def resolve_users( obj, info, first: int = 20, after: Optional[str] = None, search: Optional[str] = None ) -> dict: """Resolve paginated user list.""" # Decode cursor offset = decode_cursor(after) if after else 0 # Fetch users users = await fetch_users( limit=first + 1, # Fetch one extra to check hasNextPage offset=offset, search=search ) # Pagination has_next = len(users) > first if has_next: users = users[:first] edges = [ { "node": user, "cursor": encode_cursor(offset + i) } for i, user in enumerate(users) ] return { "edges": edges, "pageInfo": { "hasNextPage": has_next, "hasPreviousPage": offset > 0, "startCursor": edges[0]["cursor"] if edges else None, "endCursor": edges[-1]["cursor"] if edges else None }, "totalCount": await count_users(search=search) } @user_type.field("orders") async def resolve_user_orders(user: dict, info, first: int = 20) -> dict: """Resolve user's orders (N+1 prevention with DataLoader).""" # Use DataLoader to batch requests loader = info.context["loaders"]["orders_by_user"] orders = await loader.load(user["id"]) return paginate_orders(orders, first) @mutation.field("createUser") async def resolve_create_user(obj, info, input: dict) -> dict: """Create new user.""" try: # Validate input validate_user_input(input) # Create user user = await create_user( email=input["email"], name=input["name"], password=hash_password(input["password"]) ) return { "user": user, "errors": [] } except ValidationError as e: return { "user": None, "errors": [{"field": e.field, "message": e.message}] } ``` ### Pattern 3: DataLoader (N+1 Problem Prevention) ```python from aiodataloader import DataLoader from typing import List, Optional class UserLoader(DataLoader): """Batch load users by ID.""" async def batch_load_fn(self, user_ids: List[str]) -> List[Optional[dict]]: """Load multiple users in single query.""" users = await fetch_users_by_ids(user_ids) # Map results back to input order user_map = {user["id"]: user for user in users} return [user_map.get(user_id) for user_id in user_ids] class OrdersByUserLoader(DataLoader): """Batch load orders by user ID.""" async def batch_load_fn(self, user_ids: List[str]) -> List[List[dict]]: """Load orders for multiple users in single query.""" orders = await fetch_orders_by_user_ids(user_ids) # Group orders by user_id orders_by_user = {} for order in orders: user_id = order["user_id"] if user_id not in orders_by_user: orders_by_user[user_id] = [] orders_by_user[user_id].append(order) # Return in input order return [orders_by_user.get(user_id, []) for user_id in user_ids] # Context setup def create_context(): return { "loaders": { "user": UserLoader(), "orders_by_user": OrdersByUserLoader() } } ``` ## Best Practices ### REST APIs 1. **Consistent Naming**: Use plural nouns for collections (`/users`, not `/user`) 2. **Stateless**: Each request contains all necessary information 3. **Use HTTP Status Codes Correctly**: 2xx success, 4xx client errors, 5xx server errors 4. **Version Your API**: Plan for breaking changes from day one 5. **Pagination**: Always paginate large collections 6. **Rate Limiting**: Protect your API with rate limits 7. **Documentation**: Use OpenAPI/Swagger for interactive docs ### GraphQL APIs 1. **Schema First**: Design schema before writing resolvers 2. **Avoid N+1**: Use DataLoaders for efficient data fetching 3. **Input Validation**: Validate at schema and resolver levels 4. **Error Handling**: Return structured errors in mutation payloads 5. **Pagination**: Use cursor-based pagination (Relay spec) 6. **Deprecation**: Use `@deprecated` directive for gradual migration 7. **Monitoring**: Track query complexity and execution time ## Common Pitfalls - **Over-fetching/Under-fetching (REST)**: Fixed in GraphQL but requires DataLoaders - **Breaking Changes**: Version APIs or use deprecation strategies - **Inconsistent Error Formats**: Standardize error responses - **Missing Rate Limits**: APIs without limits are vulnerable to abuse - **Poor Documentation**: Undocumented APIs frustrate developers - **Ignoring HTTP Semantics**: POST for idempotent operations breaks expectations - **Tight Coupling**: API structure shouldn't mirror database schema ## Resources - **references/rest-best-practices.md**: Comprehensive REST API design guide - **references/graphql-schema-design.md**: GraphQL schema patterns and anti-patterns - **references/api-versioning-strategies.md**: Versioning approaches and migration paths - **assets/rest-api-template.py**: FastAPI REST API template - **assets/graphql-schema-template.graphql**: Complete GraphQL schema example - **assets/api-design-checklist.md**: Pre-implementation review checklist - **scripts/openapi-generator.py**: Generate OpenAPI specs from code
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

workflow-patterns

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

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

gdpr-data-handling

Implement GDPR-compliant data handling with consent management,

data
⭐1
# GDPR Data Handling Practical implementation guide for GDPR-compliant data processing, consent management, and privacy controls. ## When to Use This Skill - Building systems that process EU personal data - Implementing consent management - Handling data subject requests (DSRs) - Conducting GDPR compliance reviews - Designing privacy-first architectures - Creating data processing agreements ## Core Concepts ### 1. Personal Data Categories | Category | Examples | Protection Level | | ---------------------- | --------------------------- | ------------------ | | **Basic** | Name, email, phone | Standard | | **Sensitive (Art. 9)** | Health, religion, ethnicity | Explicit consent | | **Criminal (Art. 10)** | Convictions, offenses | Official authority | | **Children's** | Under 16 data | Parental consent | ### 2. Legal Bases for Processing ``` Article 6 - Lawful Bases: β”œβ”€β”€ Consent: Freely given, specific, informed β”œβ”€β”€ Contract: Necessary for contract performance β”œβ”€β”€ Legal Obligation: Required by law β”œβ”€β”€ Vital Interests: Protecting someone's life β”œβ”€β”€ Public Interest: Official functions └── Legitimate Interest: Balanced against rights ``` ### 3. Data Subject Rights ``` Right to Access (Art. 15) ─┐ Right to Rectification (Art. 16) β”‚ Right to Erasure (Art. 17) β”‚ Must respond Right to Restrict (Art. 18) β”‚ within 1 month Right to Portability (Art. 20) β”‚ Right to Object (Art. 21) β”€β”˜ ``` ## Implementation Patterns ### Pattern 1: Consent Management ```javascript // Consent data model const consentSchema = { userId: String, consents: [ { purpose: String, // 'marketing', 'analytics', etc. granted: Boolean, timestamp: Date, source: String, // 'web_form', 'api', etc. version: String, // Privacy policy version ipAddress: String, // For proof userAgent: String, // For proof }, ], auditLog: [ { action: String, // 'granted', 'withdrawn', 'updated' purpose: String, timestamp: Date, source: String, }, ], }; // Consent service class ConsentManager { async recordConsent(userId, purpose, granted, metadata) { const consent = { purpose, granted, timestamp: new Date(), source: metadata.source, version: await this.getCurrentPolicyVersion(), ipAddress: metadata.ipAddress, userAgent: metadata.userAgent, }; // Store consent await this.db.consents.updateOne( { userId }, { $push: { consents: consent, auditLog: { action: granted ? "granted" : "withdrawn", purpose, timestamp: consent.timestamp, source: metadata.source, }, }, }, { upsert: true }, ); // Emit event for downstream systems await this.eventBus.emit("consent.changed", { userId, purpose, granted, timestamp: consent.timestamp, }); } async hasConsent(userId, purpose) { const record = await this.db.consents.findOne({ userId }); if (!record) return false; const latestConsent = record.consents .filter((c) => c.purpose === purpose) .sort((a, b) => b.timestamp - a.timestamp)[0]; return latestConsent?.granted === true; } async getConsentHistory(userId) { const record = await this.db.consents.findOne({ userId }); return record?.auditLog || []; } } ``` ```html <!-- GDPR-compliant consent UI --> <div class="consent-banner" role="dialog" aria-labelledby="consent-title"> <h2 id="consent-title">Cookie Preferences</h2> <p> We use cookies to improve your experience. Select your preferences below. </p> <form id="consent-form"> <!-- Necessary - always on, no consent needed --> <div class="consent-category"> <input type="checkbox" id="necessary" checked disabled /> <label for="necessary"> <strong>Necessary</strong> <span>Required for the website to function. Cannot be disabled.</span> </label> </div> <!-- Analytics - requires consent --> <div class="consent-category"> <input type="checkbox" id="analytics" name="analytics" /> <label for="analytics"> <strong>Analytics</strong> <span>Help us understand how you use our site.</span> </label> </div> <!-- Marketing - requires consent --> <div class="consent-category"> <input type="checkbox" id="marketing" name="marketing" /> <label for="marketing"> <strong>Marketing</strong> <span>Personalized ads based on your interests.</span> </label> </div> <div class="consent-actions"> <button type="button" id="accept-all">Accept All</button> <button type="button" id="reject-all">Reject All</button> <button type="submit">Save Preferences</button> </div> <p class="consent-links"> <a href="/privacy-policy">Privacy Policy</a> | <a href="/cookie-policy">Cookie Policy</a> </p> </form> </div> ``` ### Pattern 2: Data Subject Access Request (DSAR) ```python from datetime import datetime, timedelta from typing import Dict, List, Optional import json class DSARHandler: """Handle Data Subject Access Requests.""" RESPONSE_DEADLINE_DAYS = 30 EXTENSION_ALLOWED_DAYS = 60 # For complex requests def __init__(self, data_sources: List['DataSource']): self.data_sources = data_sources async def submit_request( self, request_type: str, # 'access', 'erasure', 'rectification', 'portability' user_id: str, verified: bool, details: Optional[Dict] = None ) -> str: """Submit a new DSAR.""" request = { 'id': self.generate_request_id(), 'type': request_type, 'user_id': user_id, 'status': 'pending_verification' if not verified else 'processing', 'submitted_at': datetime.utcnow(), 'deadline': datetime.utcnow() + timedelta(days=self.RESPONSE_DEADLINE_DAYS), 'details': details or {}, 'audit_log': [{ 'action': 'submitted', 'timestamp': datetime.utcnow(), 'details': 'Request received' }] } await self.db.dsar_requests.insert_one(request) await self.notify_dpo(request) return request['id'] async def process_access_request(self, request_id: str) -> Dict: """Process a data access request.""" request = await self.get_request(request_id) if request['type'] != 'access': raise ValueError("Not an access request") # Collect data from all sources user_data = {} for source in self.data_sources: try: data = await source.get_user_data(request['user_id']) user_data[source.name] = data except Exception as e: user_data[source.name] = {'error': str(e)} # Format response response = { 'request_id': request_id, 'generated_at': datetime.utcnow().isoformat(), 'data_categories': list(user_data.keys()), 'data': user_data, 'retention_info': await self.get_retention_info(), 'processing_purposes': await self.get_processing_purposes(), 'third_party_recipients': await self.get_recipients() } # Update request status await self.update_request(request_id, 'completed', response) return response async def process_erasure_request(self, request_id: str) -> Dict: """Process a right to erasure request.""" request = await self.get_request(request_id) if request['type'] != 'erasure': raise ValueError("Not an erasure request") results = {} exceptions = [] for source in self.data_sources: try: # Check for legal exceptions can_delete, reason = await source.can_delete(request['user_id']) if can_delete: await source.delete_user_data(request['user_id']) results[source.name] = 'deleted' else: exceptions.append({ 'source': source.name, 'reason': reason # e.g., 'legal retention requirement' }) results[source.name] = f'retained: {reason}' except Exception as e: results[source.name] = f'error: {str(e)}' response = { 'request_id': request_id, 'completed_at': datetime.utcnow().isoformat(), 'results': results, 'exceptions': exceptions } await self.update_request(request_id, 'completed', response) return response async def process_portability_request(self, request_id: str) -> bytes: """Generate portable data export.""" request = await self.get_request(request_id) user_data = await self.process_access_request(request_id) # Convert to machine-readable format (JSON) portable_data = { 'export_date': datetime.utcnow().isoformat(), 'format_version': '1.0', 'data': user_data['data'] } return json.dumps(portable_data, indent=2, default=str).encode() ``` ### Pattern 3: Data Retention ```python from datetime import datetime, timedelta from enum import Enum class RetentionBasis(Enum): CONSENT = "consent" CONTRACT = "contract" LEGAL_OBLIGATION = "legal_obligation" LEGITIMATE_INTEREST = "legitimate_interest" class DataRetentionPolicy: """Define and enforce data retention policies.""" POLICIES = { 'user_account': { 'retention_period_days': 365 * 3, # 3 years after last activity 'basis': RetentionBasis.CONTRACT, 'trigger': 'last_activity_date', 'archive_before_delete': True }, 'transaction_records': { 'retention_period_days': 365 * 7, # 7 years for tax 'basis': RetentionBasis.LEGAL_OBLIGATION, 'trigger': 'transaction_date', 'archive_before_delete': True, 'legal_reference': 'Tax regulations require 7 year retention' }, 'marketing_consent': { 'retention_period_days': 365 * 2, # 2 years 'basis': RetentionBasis.CONSENT, 'trigger': 'consent_date', 'archive_before_delete': False }, 'support_tickets': { 'retention_period_days': 365 * 2, 'basis': RetentionBasis.LEGITIMATE_INTEREST, 'trigger': 'ticket_closed_date', 'archive_before_delete': True }, 'analytics_data': { 'retention_period_days': 365, # 1 year 'basis': RetentionBasis.CONSENT, 'trigger': 'collection_date', 'archive_before_delete': False, 'anonymize_instead': True } } async def apply_retention_policies(self): """Run retention policy enforcement.""" for data_type, policy in self.POLICIES.items(): cutoff_date = datetime.utcnow() - timedelta( days=policy['retention_period_days'] ) if policy.get('anonymize_instead'): await self.anonymize_old_data(data_type, cutoff_date) else: if policy.get('archive_before_delete'): await self.archive_data(data_type, cutoff_date) await self.delete_old_data(data_type, cutoff_date) await self.log_retention_action(data_type, cutoff_date) async def anonymize_old_data(self, data_type: str, before_date: datetime): """Anonymize data instead of deleting.""" # Example: Replace identifying fields with hashes if data_type == 'analytics_data': await self.db.analytics.update_many( {'collection_date': {'$lt': before_date}}, {'$set': { 'user_id': None, 'ip_address': None, 'device_id': None, 'anonymized': True, 'anonymized_date': datetime.utcnow() }} ) ``` ### Pattern 4: Privacy by Design ```python class PrivacyFirstDataModel: """Example of privacy-by-design data model.""" # Separate PII from behavioral data user_profile_schema = { 'user_id': str, # UUID, not sequential 'email_hash': str, # Hashed for lookups 'created_at': datetime, # Minimal data collection 'preferences': { 'language': str, 'timezone': str } } # Encrypted at rest user_pii_schema = { 'user_id': str, 'email': str, # Encrypted 'name': str, # Encrypted 'phone': str, # Encrypted (optional) 'address': dict, # Encrypted (optional) 'encryption_key_id': str } # Pseudonymized behavioral data analytics_schema = { 'session_id': str, # Not linked to user_id 'pseudonym_id': str, # Rotating pseudonym 'events': list, 'device_category': str, # Generalized, not specific 'country': str, # Not city-level } class DataMinimization: """Implement data minimization principles.""" @staticmethod def collect_only_needed(form_data: dict, purpose: str) -> dict: """Filter form data to only fields needed for purpose.""" REQUIRED_FIELDS = { 'account_creation': ['email', 'password'], 'newsletter': ['email'], 'purchase': ['email', 'name', 'address', 'payment'], 'support': ['email', 'message'] } allowed = REQUIRED_FIELDS.get(purpose, []) return {k: v for k, v in form_data.items() if k in allowed} @staticmethod def generalize_location(ip_address: str) -> str: """Generalize IP to country level only.""" import geoip2.database reader = geoip2.database.Reader('GeoLite2-Country.mmdb') try: response = reader.country(ip_address) return response.country.iso_code except: return 'UNKNOWN' ``` ### Pattern 5: Breach Notification ```python from datetime import datetime from enum import Enum class BreachSeverity(Enum): LOW = "low" MEDIUM = "medium" HIGH = "high" CRITICAL = "critical" class BreachNotificationHandler: """Handle GDPR breach notification requirements.""" AUTHORITY_NOTIFICATION_HOURS = 72 AFFECTED_NOTIFICATION_REQUIRED_SEVERITY = BreachSeverity.HIGH async def report_breach( self, description: str, data_types: List[str], affected_count: int, severity: BreachSeverity ) -> dict: """Report and handle a data breach.""" breach = { 'id': self.generate_breach_id(), 'reported_at': datetime.utcnow(), 'description': description, 'data_types_affected': data_types, 'affected_individuals_count': affected_count, 'severity': severity.value, 'status': 'investigating', 'timeline': [{ 'event': 'breach_reported', 'timestamp': datetime.utcnow(), 'details': description }] } await self.db.breaches.insert_one(breach) # Immediate notifications await self.notify_dpo(breach) await self.notify_security_team(breach) # Authority notification required within 72 hours if self.requires_authority_notification(severity, data_types): breach['authority_notification_deadline'] = ( datetime.utcnow() + timedelta(hours=self.AUTHORITY_NOTIFICATION_HOURS) ) await self.schedule_authority_notification(breach) # Affected individuals notification if severity.value in [BreachSeverity.HIGH.value, BreachSeverity.CRITICAL.value]: await self.schedule_individual_notifications(breach) return breach def requires_authority_notification( self, severity: BreachSeverity, data_types: List[str] ) -> bool: """Determine if supervisory authority must be notified.""" # Always notify for sensitive data sensitive_types = ['health', 'financial', 'credentials', 'biometric'] if any(t in sensitive_types for t in data_types): return True # Notify for medium+ severity return severity in [BreachSeverity.MEDIUM, BreachSeverity.HIGH, BreachSeverity.CRITICAL] async def generate_authority_report(self, breach_id: str) -> dict: """Generate report for supervisory authority.""" breach = await self.get_breach(breach_id) return { 'organization': { 'name': self.config.org_name, 'contact': self.config.dpo_contact, 'registration': self.config.registration_number }, 'breach': { 'nature': breach['description'], 'categories_affected': breach['data_types_affected'], 'approximate_number_affected': breach['affected_individuals_count'], 'likely_consequences': self.assess_consequences(breach), 'measures_taken': await self.get_remediation_measures(breach_id), 'measures_proposed': await self.get_proposed_measures(breach_id) }, 'timeline': breach['timeline'], 'submitted_at': datetime.utcnow().isoformat() } ``` ## Compliance Checklist ```markdown ## GDPR Implementation Checklist ### Legal Basis - [ ] Documented legal basis for each processing activity - [ ] Consent mechanisms meet GDPR requirements - [ ] Legitimate interest assessments completed ### Transparency - [ ] Privacy policy is clear and accessible - [ ] Processing purposes clearly stated - [ ] Data retention periods documented ### Data Subject Rights - [ ] Access request process implemented - [ ] Erasure request process implemented - [ ] Portability export available - [ ] Rectification process available - [ ] Response within 30-day deadline ### Security - [ ] Encryption at rest implemented - [ ] Encryption in transit (TLS) - [ ] Access controls in place - [ ] Audit logging enabled ### Breach Response - [ ] Breach detection mechanisms - [ ] 72-hour notification process - [ ] Breach documentation system ### Documentation - [ ] Records of processing activities (Art. 30) - [ ] Data protection impact assessments - [ ] Data processing agreements with vendors ``` ## Best Practices ### Do's - **Minimize data collection** - Only collect what's needed - **Document everything** - Processing activities, legal bases - **Encrypt PII** - At rest and in transit - **Implement access controls** - Need-to-know basis - **Regular audits** - Verify compliance continuously ### Don'ts - **Don't pre-check consent boxes** - Must be opt-in - **Don't bundle consent** - Separate purposes separately - **Don't retain indefinitely** - Defi
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

track-management

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

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

cqrs-implementation

Implement Command Query Responsibility Segregation for scalable

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