Skip to main content
EVOKORE// BROWSE
>

./browse/prompts

17 NODES
🤖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

PR Merge Boundary Validation Runbook

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

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

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

Improvement Cycle Review Wizard

Run a live improvement-cycle workflow that links review artifacts, approvals, tool requests, and operator checkpoints.

productivity
⭐1
# Improvement Cycle Review Wizard Imported from curated first-party documentation sources. ## What this covers Use this workflow when an operator needs a guided improvement cycle with live status, explicit approvals, and linked artifacts. ## Use this when - Operational review sessions with live progress - Approval-heavy improvement cycles - Coordinating review artifacts with execution steps ## Expected outcomes - Live workflow execution remains visible to operators - Artifacts, approvals, and tool requests stay connected - Improvement loops become easier to audit and rerun ## Source synthesis - AGENT33/docs/phase25-26-live-review-walkthrough.md (https://github.com/mattmre/AGENT33/blob/main/docs/phase25-26-live-review-walkthrough.md) - AGENT33/docs/operator-improvement-cycle-and-jupyter.md (https://github.com/mattmre/AGENT33/blob/main/docs/operator-improvement-cycle-and-jupyter.md) ## Dedupe notes Merges AGENT33 live review walkthrough details with the shorter operator-focused improvement-cycle guide. ## Source excerpts ### AGENT33/docs/phase25-26-live-review-walkthrough.md This guide documents the operator flow introduced by the Phase 25 live workflow transport and the Phase 26/27 improvement-cycle review wizard stack. ## Scope - Live workflow execution with run-scoped graph refresh - WebSocket-first status streaming with authenticated SSE fallback - Improvement-cycle preset creation - Explanation artifact generation for `plan_review` and `diff_review` - Linked review creation, risk assessment, L1/L2 signoff, and final approval - Pending tool-approval triage inside the same workflow domain > Note > This walkthrough reflects the review stack built on `codex/session58-phase26-wizard` and validated in `codex/session60-phase22-docs-validation`. If the related PRs are still open, `main` may not yet expose every surface described here. ## Operator Flow ```mermaid sequenceDiagram participant UI as "Frontend Control Plane" participant WF as "Workflow APIs" participant VIZ as "Visualization APIs" participant EXP as "Explanation APIs" participant REV as "Review APIs" participant HITL as "Tool Approval APIs" UI->>WF: "POST /v1/workflows/{name}/execute (single mode, caller run_id)" UI->>VIZ: "GET /v1/visualizations/workflows/{workflow_id}/graph?run_id=..." UI->>WF: "WS /v1/workflows/{run_id}/ws" alt "WebSocket unav ... ### AGENT33/docs/operator-improvement-cycle-and-jupyter.md This guide covers the merged operator surfaces for: - the Phase 26 improvement-cycle review wizard - the Phase 27 canonical workflow presets - the Phase 38 Docker-backed Jupyter kernel workflow Use it when you want the shortest current path from UI entry point to a real workflow run. ## 1. Improvement-Cycle Wizard The wizard is mounted inside the frontend control plane under the `Workflows` domain. ### Entry path 1. Open the frontend at `http://localhost:3000` 2. Authenticate with a bearer token or API key 3. Open `Advanced Settings` 4. Select the `Workflows` domain 5. Use the `Improvement Cycle Wizard` panel at the top of the page ### What the wizard does The wizard stitches together the backend surfaces that previously had to be called manually: - plan review / diff review generation - review creation and risk assessment - L1 and L2 review submission - tool approval request review and decision capture Reference implementation: - frontend: `frontend/src/features/improvement-cycle/ImprovementCycleWizard.tsx` - tests: `frontend/src/features/improvement-cycle/ImprovementCycleWizard.test.tsx` For the detailed review flow, keep using: - [`phase25-26-live-review-walkthrough.md`](phase25-26-live-review-walkthrough.md) ## 2. Canonical Workflow Presets The `Workflows` doma ...
👍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
📝text•6 months ago

Guardrailed Code Review Pipeline

Run a staged code review workflow with explicit risk assessment, reviewer handoffs, and merge gates.

coding
⭐1
# Guardrailed Code Review Pipeline Imported from curated first-party documentation sources. ## What this covers Use this workflow when a change needs structured review evidence, clear approval stages, and a rollback-aware path to merge. ## Use this when - Multi-step L1 and L2 review signoff - High-risk changes that need audit trails - Standardizing reviewer decisions before merge ## Expected outcomes - Review states and approval checkpoints stay explicit - Risk assessment and evidence are captured before merge - The workflow remains small enough to reuse across repositories ## Source synthesis - AGENT33/docs/functionality-and-workflows.md (https://github.com/mattmre/AGENT33/blob/main/docs/functionality-and-workflows.md) - AGENT33/docs/walkthroughs.md (https://github.com/mattmre/AGENT33/blob/main/docs/walkthroughs.md) ## Dedupe notes Combines AGENT33 lifecycle mapping and operator walkthrough material into one review-oriented site entry. ## Source excerpts ### AGENT33/docs/functionality-and-workflows.md ### 4.1 Review Lifecycle States: - `draft -> ready -> l1-review -> l1-approved -> (optional l2-review -> l2-approved) -> approved -> merged` Main APIs: - `/v1/reviews/{id}/assess` - `/v1/reviews/{id}/assign-l1` - `/v1/reviews/{id}/l1` - `/v1/reviews/{id}/assign-l2` - `/v1/reviews/{id}/l2` - `/v1/reviews/{id}/approve` - `/v1/reviews/{id}/merge` ### AGENT33/docs/walkthroughs.md ## 4. Review Lifecycle (Two-Layer Signoff) Create review: ```bash curl -X POST http://localhost:8000/v1/reviews/ \ -H "Authorization: Bearer $TOKEN" \ -H "Content-Type: application/json" \ -d '{"task_id":"TASK-101","branch":"feat/docs-refresh","pr_number":12}' ``` Assess risk: ```bash curl -X POST http://localhost:8000/v1/reviews/<review_id>/assess \ -H "Authorization: Bearer $TOKEN" \ -H "Content-Type: application/json" \ -d '{"triggers":["api-public","security"]}' ``` Move to ready and assign L1: ```bash curl -X POST http://localhost:8000/v1/reviews/<review_id>/ready -H "Authorization: Bearer $TOKEN" curl -X POST http://localhost:8000/v1/reviews/<review_id>/assign-l1 -H "Authorization: Bearer $TOKEN" ``` Submit L1 decision: ```bash curl -X POST http://localhost:8000/v1/reviews/<review_id>/l1 \ -H "Authorization: Bearer $TOKEN" \ -H "Content-Type: application/json" \ -d '{"decision":"approved","issues":[],"comments":"L1 pass"}' ``` If L2 required, continue: ```bash curl -X POST http://localhost:8000/v1/reviews/<review_id>/assign-l2 -H "Authorization: Bearer $TOKEN" curl -X POST http://localhost:8000/v1/reviews/<review_id>/l2 \ -H "Authorization: Bearer $TOKEN" \ -H "Content-Type: application/json" \ -d '{"decision":"approved","issues":[],"commen ...
👍0
👁️0
docs
🤖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

database-migration

Execute database migrations across ORMs and platforms with

data
⭐1
# Database Migration Master database schema and data migrations across ORMs (Sequelize, TypeORM, Prisma), including rollback strategies and zero-downtime deployments. ## When to Use This Skill - Migrating between different ORMs - Performing schema transformations - Moving data between databases - Implementing rollback procedures - Zero-downtime deployments - Database version upgrades - Data model refactoring ## ORM Migrations ### Sequelize Migrations ```javascript // migrations/20231201-create-users.js module.exports = { up: async (queryInterface, Sequelize) => { await queryInterface.createTable("users", { id: { type: Sequelize.INTEGER, primaryKey: true, autoIncrement: true, }, email: { type: Sequelize.STRING, unique: true, allowNull: false, }, createdAt: Sequelize.DATE, updatedAt: Sequelize.DATE, }); }, down: async (queryInterface, Sequelize) => { await queryInterface.dropTable("users"); }, }; // Run: npx sequelize-cli db:migrate // Rollback: npx sequelize-cli db:migrate:undo ``` ### TypeORM Migrations ```typescript // migrations/1701234567-CreateUsers.ts import { MigrationInterface, QueryRunner, Table } from "typeorm"; export class CreateUsers1701234567 implements MigrationInterface { public async up(queryRunner: QueryRunner): Promise<void> { await queryRunner.createTable( new Table({ name: "users", columns: [ { name: "id", type: "int", isPrimary: true, isGenerated: true, generationStrategy: "increment", }, { name: "email", type: "varchar", isUnique: true, }, { name: "created_at", type: "timestamp", default: "CURRENT_TIMESTAMP", }, ], }), ); } public async down(queryRunner: QueryRunner): Promise<void> { await queryRunner.dropTable("users"); } } // Run: npm run typeorm migration:run // Rollback: npm run typeorm migration:revert ``` ### Prisma Migrations ```prisma // schema.prisma model User { id Int @id @default(autoincrement()) email String @unique createdAt DateTime @default(now()) } // Generate migration: npx prisma migrate dev --name create_users // Apply: npx prisma migrate deploy ``` ## Schema Transformations ### Adding Columns with Defaults ```javascript // Safe migration: add column with default module.exports = { up: async (queryInterface, Sequelize) => { await queryInterface.addColumn("users", "status", { type: Sequelize.STRING, defaultValue: "active", allowNull: false, }); }, down: async (queryInterface) => { await queryInterface.removeColumn("users", "status"); }, }; ``` ### Renaming Columns (Zero Downtime) ```javascript // Step 1: Add new column module.exports = { up: async (queryInterface, Sequelize) => { await queryInterface.addColumn("users", "full_name", { type: Sequelize.STRING, }); // Copy data from old column await queryInterface.sequelize.query("UPDATE users SET full_name = name"); }, down: async (queryInterface) => { await queryInterface.removeColumn("users", "full_name"); }, }; // Step 2: Update application to use new column // Step 3: Remove old column module.exports = { up: async (queryInterface) => { await queryInterface.removeColumn("users", "name"); }, down: async (queryInterface, Sequelize) => { await queryInterface.addColumn("users", "name", { type: Sequelize.STRING, }); }, }; ``` ### Changing Column Types ```javascript module.exports = { up: async (queryInterface, Sequelize) => { // For large tables, use multi-step approach // 1. Add new column await queryInterface.addColumn("users", "age_new", { type: Sequelize.INTEGER, }); // 2. Copy and transform data await queryInterface.sequelize.query(` UPDATE users SET age_new = CAST(age AS INTEGER) WHERE age IS NOT NULL `); // 3. Drop old column await queryInterface.removeColumn("users", "age"); // 4. Rename new column await queryInterface.renameColumn("users", "age_new", "age"); }, down: async (queryInterface, Sequelize) => { await queryInterface.changeColumn("users", "age", { type: Sequelize.STRING, }); }, }; ``` ## Data Transformations ### Complex Data Migration ```javascript module.exports = { up: async (queryInterface, Sequelize) => { // Get all records const [users] = await queryInterface.sequelize.query( "SELECT id, address_string FROM users", ); // Transform each record for (const user of users) { const addressParts = user.address_string.split(","); await queryInterface.sequelize.query( `UPDATE users SET street = :street, city = :city, state = :state WHERE id = :id`, { replacements: { id: user.id, street: addressParts[0]?.trim(), city: addressParts[1]?.trim(), state: addressParts[2]?.trim(), }, }, ); } // Drop old column await queryInterface.removeColumn("users", "address_string"); }, down: async (queryInterface, Sequelize) => { // Reconstruct original column await queryInterface.addColumn("users", "address_string", { type: Sequelize.STRING, }); await queryInterface.sequelize.query(` UPDATE users SET address_string = CONCAT(street, ', ', city, ', ', state) `); await queryInterface.removeColumn("users", "street"); await queryInterface.removeColumn("users", "city"); await queryInterface.removeColumn("users", "state"); }, }; ``` ## Rollback Strategies ### Transaction-Based Migrations ```javascript module.exports = { up: async (queryInterface, Sequelize) => { const transaction = await queryInterface.sequelize.transaction(); try { await queryInterface.addColumn( "users", "verified", { type: Sequelize.BOOLEAN, defaultValue: false }, { transaction }, ); await queryInterface.sequelize.query( "UPDATE users SET verified = true WHERE email_verified_at IS NOT NULL", { transaction }, ); await transaction.commit(); } catch (error) { await transaction.rollback(); throw error; } }, down: async (queryInterface) => { await queryInterface.removeColumn("users", "verified"); }, }; ``` ### Checkpoint-Based Rollback ```javascript module.exports = { up: async (queryInterface, Sequelize) => { // Create backup table await queryInterface.sequelize.query( "CREATE TABLE users_backup AS SELECT * FROM users", ); try { // Perform migration await queryInterface.addColumn("users", "new_field", { type: Sequelize.STRING, }); // Verify migration const [result] = await queryInterface.sequelize.query( "SELECT COUNT(*) as count FROM users WHERE new_field IS NULL", ); if (result[0].count > 0) { throw new Error("Migration verification failed"); } // Drop backup await queryInterface.dropTable("users_backup"); } catch (error) { // Restore from backup await queryInterface.sequelize.query("DROP TABLE users"); await queryInterface.sequelize.query( "CREATE TABLE users AS SELECT * FROM users_backup", ); await queryInterface.dropTable("users_backup"); throw error; } }, }; ``` ## Zero-Downtime Migrations ### Blue-Green Deployment Strategy ```javascript // Phase 1: Make changes backward compatible module.exports = { up: async (queryInterface, Sequelize) => { // Add new column (both old and new code can work) await queryInterface.addColumn("users", "email_new", { type: Sequelize.STRING, }); }, }; // Phase 2: Deploy code that writes to both columns // Phase 3: Backfill data module.exports = { up: async (queryInterface) => { await queryInterface.sequelize.query(` UPDATE users SET email_new = email WHERE email_new IS NULL `); }, }; // Phase 4: Deploy code that reads from new column // Phase 5: Remove old column module.exports = { up: async (queryInterface) => { await queryInterface.removeColumn("users", "email"); }, }; ``` ## Cross-Database Migrations ### PostgreSQL to MySQL ```javascript // Handle differences module.exports = { up: async (queryInterface, Sequelize) => { const dialectName = queryInterface.sequelize.getDialect(); if (dialectName === "mysql") { await queryInterface.createTable("users", { id: { type: Sequelize.INTEGER, primaryKey: true, autoIncrement: true, }, data: { type: Sequelize.JSON, // MySQL JSON type }, }); } else if (dialectName === "postgres") { await queryInterface.createTable("users", { id: { type: Sequelize.INTEGER, primaryKey: true, autoIncrement: true, }, data: { type: Sequelize.JSONB, // PostgreSQL JSONB type }, }); } }, }; ``` ## Resources - **references/orm-switching.md**: ORM migration guides - **references/schema-migration.md**: Schema transformation patterns - **references/data-transformation.md**: Data migration scripts - **references/rollback-strategies.md**: Rollback procedures - **assets/schema-migration-template.sql**: SQL migration templates - **assets/data-migration-script.py**: Data migration utilities - **scripts/test-migration.sh**: Migration testing script ## Best Practices 1. **Always Provide Rollback**: Every up() needs a down() 2. **Test Migrations**: Test on staging first 3. **Use Transactions**: Atomic migrations when possible 4. **Backup First**: Always backup before migration 5. **Small Changes**: Break into small, incremental steps 6. **Monitor**: Watch for errors during deployment 7. **Document**: Explain why and how 8. **Idempotent**: Migrations should be rerunnable ## Common Pitfalls - Not testing rollback procedures - Making breaking changes without downtime strategy - Forgetting to handle NULL values - Not considering index performance - Ignoring foreign key constraints - Migrating too much data at once
👍0
👁️0
🤖 Auto-discovered
🤖system prompt•7 months ago

cqrs-implementation

Implement Command Query Responsibility Segregation for scalable

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

langchain-architecture

Design LLM applications using LangChain 1.x and LangGraph for

coding
⭐1
# LangChain & LangGraph Architecture Master modern LangChain 1.x and LangGraph for building sophisticated LLM applications with agents, state management, memory, and tool integration. ## When to Use This Skill - Building autonomous AI agents with tool access - Implementing complex multi-step LLM workflows - Managing conversation memory and state - Integrating LLMs with external data sources and APIs - Creating modular, reusable LLM application components - Implementing document processing pipelines - Building production-grade LLM applications ## Package Structure (LangChain 1.x) ``` langchain (1.2.x) # High-level orchestration langchain-core (1.2.x) # Core abstractions (messages, prompts, tools) langchain-community # Third-party integrations langgraph # Agent orchestration and state management langchain-openai # OpenAI integrations langchain-anthropic # Anthropic/Claude integrations langchain-voyageai # Voyage AI embeddings langchain-pinecone # Pinecone vector store ``` ## Core Concepts ### 1. LangGraph Agents LangGraph is the standard for building agents in 2026. It provides: **Key Features:** - **StateGraph**: Explicit state management with typed state - **Durable Execution**: Agents persist through failures - **Human-in-the-Loop**: Inspect and modify state at any point - **Memory**: Short-term and long-term memory across sessions - **Checkpointing**: Save and resume agent state **Agent Patterns:** - **ReAct**: Reasoning + Acting with `create_react_agent` - **Plan-and-Execute**: Separate planning and execution nodes - **Multi-Agent**: Supervisor routing between specialized agents - **Tool-Calling**: Structured tool invocation with Pydantic schemas ### 2. State Management LangGraph uses TypedDict for explicit state: ```python from typing import Annotated, TypedDict from langgraph.graph import MessagesState # Simple message-based state class AgentState(MessagesState): """Extends MessagesState with custom fields.""" context: Annotated[list, "retrieved documents"] # Custom state for complex agents class CustomState(TypedDict): messages: Annotated[list, "conversation history"] context: Annotated[dict, "retrieved context"] current_step: str results: list ``` ### 3. Memory Systems Modern memory implementations: - **ConversationBufferMemory**: Stores all messages (short conversations) - **ConversationSummaryMemory**: Summarizes older messages (long conversations) - **ConversationTokenBufferMemory**: Token-based windowing - **VectorStoreRetrieverMemory**: Semantic similarity retrieval - **LangGraph Checkpointers**: Persistent state across sessions ### 4. Document Processing Loading, transforming, and storing documents: **Components:** - **Document Loaders**: Load from various sources - **Text Splitters**: Chunk documents intelligently - **Vector Stores**: Store and retrieve embeddings - **Retrievers**: Fetch relevant documents ### 5. Callbacks & Tracing LangSmith is the standard for observability: - Request/response logging - Token usage tracking - Latency monitoring - Error tracking - Trace visualization ## Quick Start ### Modern ReAct Agent with LangGraph ```python from langgraph.prebuilt import create_react_agent from langgraph.checkpoint.memory import MemorySaver from langchain_anthropic import ChatAnthropic from langchain_core.tools import tool import ast import operator # Initialize LLM (Claude Sonnet 4.6 recommended) llm = ChatAnthropic(model="claude-sonnet-4-6", temperature=0) # Define tools with Pydantic schemas @tool def search_database(query: str) -> str: """Search internal database for information.""" # Your database search logic return f"Results for: {query}" @tool def calculate(expression: str) -> str: """Safely evaluate a mathematical expression. Supports: +, -, *, /, **, %, parentheses Example: '(2 + 3) * 4' returns '20' """ # Safe math evaluation using ast allowed_operators = { ast.Add: operator.add, ast.Sub: operator.sub, ast.Mult: operator.mul, ast.Div: operator.truediv, ast.Pow: operator.pow, ast.Mod: operator.mod, ast.USub: operator.neg, } def _eval(node): if isinstance(node, ast.Constant): return node.value elif isinstance(node, ast.BinOp): left = _eval(node.left) right = _eval(node.right) return allowed_operators[type(node.op)](left, right) elif isinstance(node, ast.UnaryOp): operand = _eval(node.operand) return allowed_operators[type(node.op)](operand) else: raise ValueError(f"Unsupported operation: {type(node)}") try: tree = ast.parse(expression, mode='eval') return str(_eval(tree.body)) except Exception as e: return f"Error: {e}" tools = [search_database, calculate] # Create checkpointer for memory persistence checkpointer = MemorySaver() # Create ReAct agent agent = create_react_agent( llm, tools, checkpointer=checkpointer ) # Run agent with thread ID for memory config = {"configurable": {"thread_id": "user-123"}} result = await agent.ainvoke( {"messages": [("user", "Search for Python tutorials and calculate 25 * 4")]}, config=config ) ``` ## Architecture Patterns ### Pattern 1: RAG with LangGraph ```python from langgraph.graph import StateGraph, START, END from langchain_anthropic import ChatAnthropic from langchain_voyageai import VoyageAIEmbeddings from langchain_pinecone import PineconeVectorStore from langchain_core.documents import Document from langchain_core.prompts import ChatPromptTemplate from typing import TypedDict, Annotated class RAGState(TypedDict): question: str context: Annotated[list[Document], "retrieved documents"] answer: str # Initialize components llm = ChatAnthropic(model="claude-sonnet-4-6") embeddings = VoyageAIEmbeddings(model="voyage-3-large") vectorstore = PineconeVectorStore(index_name="docs", embedding=embeddings) retriever = vectorstore.as_retriever(search_kwargs={"k": 4}) # Define nodes async def retrieve(state: RAGState) -> RAGState: """Retrieve relevant documents.""" docs = await retriever.ainvoke(state["question"]) return {"context": docs} async def generate(state: RAGState) -> RAGState: """Generate answer from context.""" prompt = ChatPromptTemplate.from_template( """Answer based on the context below. If you cannot answer, say so. Context: {context} Question: {question} Answer:""" ) context_text = "\n\n".join(doc.page_content for doc in state["context"]) response = await llm.ainvoke( prompt.format(context=context_text, question=state["question"]) ) return {"answer": response.content} # Build graph builder = StateGraph(RAGState) builder.add_node("retrieve", retrieve) builder.add_node("generate", generate) builder.add_edge(START, "retrieve") builder.add_edge("retrieve", "generate") builder.add_edge("generate", END) rag_chain = builder.compile() # Use the chain result = await rag_chain.ainvoke({"question": "What is the main topic?"}) ``` ### Pattern 2: Custom Agent with Structured Tools ```python from langchain_core.tools import StructuredTool from pydantic import BaseModel, Field class SearchInput(BaseModel): """Input for database search.""" query: str = Field(description="Search query") filters: dict = Field(default={}, description="Optional filters") class EmailInput(BaseModel): """Input for sending email.""" recipient: str = Field(description="Email recipient") subject: str = Field(description="Email subject") content: str = Field(description="Email body") async def search_database(query: str, filters: dict = {}) -> str: """Search internal database for information.""" # Your database search logic return f"Results for '{query}' with filters {filters}" async def send_email(recipient: str, subject: str, content: str) -> str: """Send an email to specified recipient.""" # Email sending logic return f"Email sent to {recipient}" tools = [ StructuredTool.from_function( coroutine=search_database, name="search_database", description="Search internal database", args_schema=SearchInput ), StructuredTool.from_function( coroutine=send_email, name="send_email", description="Send an email", args_schema=EmailInput ) ] agent = create_react_agent(llm, tools) ``` ### Pattern 3: Multi-Step Workflow with StateGraph ```python from langgraph.graph import StateGraph, START, END from typing import TypedDict, Literal class WorkflowState(TypedDict): text: str entities: list analysis: str summary: str current_step: str async def extract_entities(state: WorkflowState) -> WorkflowState: """Extract key entities from text.""" prompt = f"Extract key entities from: {state['text']}\n\nReturn as JSON list." response = await llm.ainvoke(prompt) return {"entities": response.content, "current_step": "analyze"} async def analyze_entities(state: WorkflowState) -> WorkflowState: """Analyze extracted entities.""" prompt = f"Analyze these entities: {state['entities']}\n\nProvide insights." response = await llm.ainvoke(prompt) return {"analysis": response.content, "current_step": "summarize"} async def generate_summary(state: WorkflowState) -> WorkflowState: """Generate final summary.""" prompt = f"""Summarize: Entities: {state['entities']} Analysis: {state['analysis']} Provide a concise summary.""" response = await llm.ainvoke(prompt) return {"summary": response.content, "current_step": "complete"} def route_step(state: WorkflowState) -> Literal["analyze", "summarize", "end"]: """Route to next step based on current state.""" step = state.get("current_step", "extract") if step == "analyze": return "analyze" elif step == "summarize": return "summarize" return "end" # Build workflow builder = StateGraph(WorkflowState) builder.add_node("extract", extract_entities) builder.add_node("analyze", analyze_entities) builder.add_node("summarize", generate_summary) builder.add_edge(START, "extract") builder.add_conditional_edges("extract", route_step, { "analyze": "analyze", "summarize": "summarize", "end": END }) builder.add_conditional_edges("analyze", route_step, { "summarize": "summarize", "end": END }) builder.add_edge("summarize", END) workflow = builder.compile() ``` ### Pattern 4: Multi-Agent Orchestration ```python from langgraph.graph import StateGraph, START, END from langgraph.prebuilt import create_react_agent from langchain_core.messages import HumanMessage from typing import Literal class MultiAgentState(TypedDict): messages: list next_agent: str # Create specialized agents researcher = create_react_agent(llm, research_tools) writer = create_react_agent(llm, writing_tools) reviewer = create_react_agent(llm, review_tools) async def supervisor(state: MultiAgentState) -> MultiAgentState: """Route to appropriate agent based on task.""" prompt = f"""Based on the conversation, which agent should handle this? Options: - researcher: For finding information - writer: For creating content - reviewer: For reviewing and editing - FINISH: Task is complete Messages: {state['messages']} Respond with just the agent name.""" response = await llm.ainvoke(prompt) return {"next_agent": response.content.strip().lower()} def route_to_agent(state: MultiAgentState) -> Literal["researcher", "writer", "reviewer", "end"]: """Route based on supervisor decision.""" next_agent = state.get("next_agent", "").lower() if next_agent == "finish": return "end" return next_agent if next_agent in ["researcher", "writer", "reviewer"] else "end" # Build multi-agent graph builder = StateGraph(MultiAgentState) builder.add_node("supervisor", supervisor) builder.add_node("researcher", researcher) builder.add_node("writer", writer) builder.add_node("reviewer", reviewer) builder.add_edge(START, "supervisor") builder.add_conditional_edges("supervisor", route_to_agent, { "researcher": "researcher", "writer": "writer", "reviewer": "reviewer", "end": END }) # Each agent returns to supervisor for agent in ["researcher", "writer", "reviewer"]: builder.add_edge(agent, "supervisor") multi_agent = builder.compile() ``` ## Memory Management ### Token-Based Memory with LangGraph ```python from langgraph.checkpoint.memory import MemorySaver from langgraph.prebuilt import create_react_agent # In-memory checkpointer (development) checkpointer = MemorySaver() # Create agent with persistent memory agent = create_react_agent(llm, tools, checkpointer=checkpointer) # Each thread_id maintains separate conversation config = {"configurable": {"thread_id": "session-abc123"}} # Messages persist across invocations with same thread_id result1 = await agent.ainvoke({"messages": [("user", "My name is Alice")]}, config) result2 = await agent.ainvoke({"messages": [("user", "What's my name?")]}, config) # Agent remembers: "Your name is Alice" ``` ### Production Memory with PostgreSQL ```python from langgraph.checkpoint.postgres import PostgresSaver # Production checkpointer checkpointer = PostgresSaver.from_conn_string( "postgresql://user:pass@localhost/langgraph" ) agent = create_react_agent(llm, tools, checkpointer=checkpointer) ``` ### Vector Store Memory for Long-Term Context ```python from langchain_community.vectorstores import Chroma from langchain_voyageai import VoyageAIEmbeddings embeddings = VoyageAIEmbeddings(model="voyage-3-large") memory_store = Chroma( collection_name="conversation_memory", embedding_function=embeddings, persist_directory="./memory_db" ) async def retrieve_relevant_memory(query: str, k: int = 5) -> list: """Retrieve relevant past conversations.""" docs = await memory_store.asimilarity_search(query, k=k) return [doc.page_content for doc in docs] async def store_memory(content: str, metadata: dict = {}): """Store conversation in long-term memory.""" await memory_store.aadd_texts([content], metadatas=[metadata]) ``` ## Callback System & LangSmith ### LangSmith Tracing ```python import os from langchain_anthropic import ChatAnthropic # Enable LangSmith tracing os.environ["LANGCHAIN_TRACING_V2"] = "true" os.environ["LANGCHAIN_API_KEY"] = "your-api-key" os.environ["LANGCHAIN_PROJECT"] = "my-project" # All LangChain/LangGraph operations are automatically traced llm = ChatAnthropic(model="claude-sonnet-4-6") ``` ### Custom Callback Handler ```python from langchain_core.callbacks import BaseCallbackHandler from typing import Any, Dict, List class CustomCallbackHandler(BaseCallbackHandler): def on_llm_start( self, serialized: Dict[str, Any], prompts: List[str], **kwargs ) -> None: print(f"LLM started with {len(prompts)} prompts") def on_llm_end(self, response, **kwargs) -> None: print(f"LLM completed: {len(response.generations)} generations") def on_llm_error(self, error: Exception, **kwargs) -> None: print(f"LLM error: {error}") def on_tool_start( self, serialized: Dict[str, Any], input_str: str, **kwargs ) -> None: print(f"Tool started: {serialized.get('name')}") def on_tool_end(self, output: str, **kwargs) -> None: print(f"Tool completed: {output[:100]}...") # Use callbacks result = await agent.ainvoke( {"messages": [("user", "query")]}, config={"callbacks": [CustomCallbackHandler()]} ) ``` ## Streaming Responses ```python from langchain_anthropic import ChatAnthropic llm = ChatAnthropic(model="claude-sonnet-4-6", streaming=True) # Stream tokens async for chunk in llm.astream("Tell me a story"): print(chunk.content, end="", flush=True) # Stream agent events async for event in agent.astream_events( {"messages": [("user", "Search and summarize")]}, version="v2" ): if event["event"] == "on_chat_model_stream": print(event["data"]["chunk"].content, end="") elif event["event"] == "on_tool_start": print(f"\n[Using tool: {event['name']}]") ``` ## Testing Strategies ```python import pytest from unittest.mock import AsyncMock, patch @pytest.mark.asyncio async def test_agent_tool_selection(): """Test agent selects correct tool.""" with patch.object(llm, 'ainvoke') as mock_llm: mock_llm.return_value = AsyncMock(content="Using search_database") result = await agent.ainvoke({ "messages": [("user", "search for documents")] }) # Verify tool was called assert "search_database" in str(result) @pytest.mark.asyncio async def test_memory_persistence(): """Test memory persists across invocations.""" config = {"configurable": {"thread_id": "test-thread"}} # First message await agent.ainvoke( {"messages": [("user", "Remember: the code is 12345")]}, config ) # Second message should remember result = await agent.ainvoke( {"messages": [("user", "What was the code?")]}, config ) assert "12345" in result["messages"][-1].content ``` ## Performance Optimization ### 1. Caching with Redis ```python from langchain_community.cache import RedisCache from langchain_core.globals import set_llm_cache import redis redis_client = redis.Redis.from_url("redis://localhost:6379") set_llm_cache(RedisCache(redis_client)) ``` ### 2. Async Batch Processing ```python import asyncio from langchain_core.documents import Document async def process_documents(documents: list[Document]) -> list: """Process documents in parallel.""" tasks = [process_single(doc) for doc in documents] return await asyncio.gather(*tasks) async def process_single(doc: Document) -> dict: """Process a single document.""" chunks = text_splitter.split_documents([doc]) embeddings = await embeddings_model.aembed_documents( [c.page_content for c in chunks] ) return {"doc_id": doc.metadata.get("id"), "embeddings": embeddings} ``` ### 3. Connection Pooling ```python from langchain_pinecone import PineconeVectorStore from pinecone import Pinecone # Reuse Pinecone client pc = Pinecone(api_key=os.environ["PINECONE_API_KEY"]) index = pc.Index("my-index") # Create vector store with existing index vectorstore = PineconeVectorStore(index=index, embedding=embeddings) ``` ## Resources - [LangChain Documentation](https://python.langchain.com/docs/) - [LangGraph Documentation](https://langchain-ai.github.io/langgraph/) - [LangSmith Platform](https://smith.langchain.com/) - [LangChain GitHub](https://github.com/langchain-ai/langchain) - [LangGraph GitHub](https://github.com/langchain-ai/langgraph) ## Common Pitfalls 1. **Using Deprecated APIs**: Use LangGraph for agents, not `initialize_agent` 2. **Memory Overflow**: Use checkpointers with TTL for long-running agents 3. **Poor Tool Descriptions**: Clear descriptions help LLM select correct tools 4. **Context Window Exceeded**: Use summarization or sliding window memory 5. **No Error Handling**: Wrap too
👍0
👁️0
🤖 Auto-discovered
🤖system prompt•7 months ago

memory-forensics

Master memory forensics techniques including memory acquisition,

security
⭐1
# Memory Forensics Comprehensive techniques for acquiring, analyzing, and extracting artifacts from memory dumps for incident response and malware analysis. ## Memory Acquisition ### Live Acquisition Tools #### Windows ```powershell # WinPmem (Recommended) winpmem_mini_x64.exe memory.raw # DumpIt DumpIt.exe # Belkasoft RAM Capturer # GUI-based, outputs raw format # Magnet RAM Capture # GUI-based, outputs raw format ``` #### Linux ```bash # LiME (Linux Memory Extractor) sudo insmod lime.ko "path=/tmp/memory.lime format=lime" # /dev/mem (limited, requires permissions) sudo dd if=/dev/mem of=memory.raw bs=1M # /proc/kcore (ELF format) sudo cp /proc/kcore memory.elf ``` #### macOS ```bash # osxpmem sudo ./osxpmem -o memory.raw # MacQuisition (commercial) ``` ### Virtual Machine Memory ```bash # VMware: .vmem file is raw memory cp vm.vmem memory.raw # VirtualBox: Use debug console vboxmanage debugvm "VMName" dumpvmcore --filename memory.elf # QEMU virsh dump <domain> memory.raw --memory-only # Hyper-V # Checkpoint contains memory state ``` ## Volatility 3 Framework ### Installation and Setup ```bash # Install Volatility 3 pip install volatility3 # Install symbol tables (Windows) # Download from https://downloads.volatilityfoundation.org/volatility3/symbols/ # Basic usage vol -f memory.raw <plugin> # With symbol path vol -f memory.raw -s /path/to/symbols windows.pslist ``` ### Essential Plugins #### Process Analysis ```bash # List processes vol -f memory.raw windows.pslist # Process tree (parent-child relationships) vol -f memory.raw windows.pstree # Hidden process detection vol -f memory.raw windows.psscan # Process memory dumps vol -f memory.raw windows.memmap --pid <PID> --dump # Process environment variables vol -f memory.raw windows.envars --pid <PID> # Command line arguments vol -f memory.raw windows.cmdline ``` #### Network Analysis ```bash # Network connections vol -f memory.raw windows.netscan # Network connection state vol -f memory.raw windows.netstat ``` #### DLL and Module Analysis ```bash # Loaded DLLs per process vol -f memory.raw windows.dlllist --pid <PID> # Find hidden/injected DLLs vol -f memory.raw windows.ldrmodules # Kernel modules vol -f memory.raw windows.modules # Module dumps vol -f memory.raw windows.moddump --pid <PID> ``` #### Memory Injection Detection ```bash # Detect code injection vol -f memory.raw windows.malfind # VAD (Virtual Address Descriptor) analysis vol -f memory.raw windows.vadinfo --pid <PID> # Dump suspicious memory regions vol -f memory.raw windows.vadyarascan --yara-rules rules.yar ``` #### Registry Analysis ```bash # List registry hives vol -f memory.raw windows.registry.hivelist # Print registry key vol -f memory.raw windows.registry.printkey --key "Software\Microsoft\Windows\CurrentVersion\Run" # Dump registry hive vol -f memory.raw windows.registry.hivescan --dump ``` #### File System Artifacts ```bash # Scan for file objects vol -f memory.raw windows.filescan # Dump files from memory vol -f memory.raw windows.dumpfiles --pid <PID> # MFT analysis vol -f memory.raw windows.mftscan ``` ### Linux Analysis ```bash # Process listing vol -f memory.raw linux.pslist # Process tree vol -f memory.raw linux.pstree # Bash history vol -f memory.raw linux.bash # Network connections vol -f memory.raw linux.sockstat # Loaded kernel modules vol -f memory.raw linux.lsmod # Mount points vol -f memory.raw linux.mount # Environment variables vol -f memory.raw linux.envars ``` ### macOS Analysis ```bash # Process listing vol -f memory.raw mac.pslist # Process tree vol -f memory.raw mac.pstree # Network connections vol -f memory.raw mac.netstat # Kernel extensions vol -f memory.raw mac.lsmod ``` ## Analysis Workflows ### Malware Analysis Workflow ```bash # 1. Initial process survey vol -f memory.raw windows.pstree > processes.txt vol -f memory.raw windows.pslist > pslist.txt # 2. Network connections vol -f memory.raw windows.netscan > network.txt # 3. Detect injection vol -f memory.raw windows.malfind > malfind.txt # 4. Analyze suspicious processes vol -f memory.raw windows.dlllist --pid <PID> vol -f memory.raw windows.handles --pid <PID> # 5. Dump suspicious executables vol -f memory.raw windows.pslist --pid <PID> --dump # 6. Extract strings from dumps strings -a pid.<PID>.exe > strings.txt # 7. YARA scanning vol -f memory.raw windows.yarascan --yara-rules malware.yar ``` ### Incident Response Workflow ```bash # 1. Timeline of events vol -f memory.raw windows.timeliner > timeline.csv # 2. User activity vol -f memory.raw windows.cmdline vol -f memory.raw windows.consoles # 3. Persistence mechanisms vol -f memory.raw windows.registry.printkey \ --key "Software\Microsoft\Windows\CurrentVersion\Run" # 4. Services vol -f memory.raw windows.svcscan # 5. Scheduled tasks vol -f memory.raw windows.scheduled_tasks # 6. Recent files vol -f memory.raw windows.filescan | grep -i "recent" ``` ## Data Structures ### Windows Process Structures ```c // EPROCESS (Executive Process) typedef struct _EPROCESS { KPROCESS Pcb; // Kernel process block EX_PUSH_LOCK ProcessLock; LARGE_INTEGER CreateTime; LARGE_INTEGER ExitTime; // ... LIST_ENTRY ActiveProcessLinks; // Doubly-linked list ULONG_PTR UniqueProcessId; // PID // ... PEB* Peb; // Process Environment Block // ... } EPROCESS; // PEB (Process Environment Block) typedef struct _PEB { BOOLEAN InheritedAddressSpace; BOOLEAN ReadImageFileExecOptions; BOOLEAN BeingDebugged; // Anti-debug check // ... PVOID ImageBaseAddress; // Base address of executable PPEB_LDR_DATA Ldr; // Loader data (DLL list) PRTL_USER_PROCESS_PARAMETERS ProcessParameters; // ... } PEB; ``` ### VAD (Virtual Address Descriptor) ```c typedef struct _MMVAD { MMVAD_SHORT Core; union { ULONG LongFlags; MMVAD_FLAGS VadFlags; } u; // ... PVOID FirstPrototypePte; PVOID LastContiguousPte; // ... PFILE_OBJECT FileObject; } MMVAD; // Memory protection flags #define PAGE_EXECUTE 0x10 #define PAGE_EXECUTE_READ 0x20 #define PAGE_EXECUTE_READWRITE 0x40 #define PAGE_EXECUTE_WRITECOPY 0x80 ``` ## Detection Patterns ### Process Injection Indicators ```python # Malfind indicators # - PAGE_EXECUTE_READWRITE protection (suspicious) # - MZ header in non-image VAD region # - Shellcode patterns at allocation start # Common injection techniques # 1. Classic DLL Injection # - VirtualAllocEx + WriteProcessMemory + CreateRemoteThread # 2. Process Hollowing # - CreateProcess (SUSPENDED) + NtUnmapViewOfSection + WriteProcessMemory # 3. APC Injection # - QueueUserAPC targeting alertable threads # 4. Thread Execution Hijacking # - SuspendThread + SetThreadContext + ResumeThread ``` ### Rootkit Detection ```bash # Compare process lists vol -f memory.raw windows.pslist > pslist.txt vol -f memory.raw windows.psscan > psscan.txt diff pslist.txt psscan.txt # Hidden processes # Check for DKOM (Direct Kernel Object Manipulation) vol -f memory.raw windows.callbacks # Detect hooked functions vol -f memory.raw windows.ssdt # System Service Descriptor Table # Driver analysis vol -f memory.raw windows.driverscan vol -f memory.raw windows.driverirp ``` ### Credential Extraction ```bash # Dump hashes (requires hivelist first) vol -f memory.raw windows.hashdump # LSA secrets vol -f memory.raw windows.lsadump # Cached domain credentials vol -f memory.raw windows.cachedump # Mimikatz-style extraction # Requires specific plugins/tools ``` ## YARA Integration ### Writing Memory YARA Rules ```yara rule Suspicious_Injection { meta: description = "Detects common injection shellcode" strings: // Common shellcode patterns $mz = { 4D 5A } $shellcode1 = { 55 8B EC 83 EC } // Function prologue $api_hash = { 68 ?? ?? ?? ?? 68 ?? ?? ?? ?? E8 } // Push hash, call condition: $mz at 0 or any of ($shellcode*) } rule Cobalt_Strike_Beacon { meta: description = "Detects Cobalt Strike beacon in memory" strings: $config = { 00 01 00 01 00 02 } $sleep = "sleeptime" $beacon = "%s (admin)" wide condition: 2 of them } ``` ### Scanning Memory ```bash # Scan all process memory vol -f memory.raw windows.yarascan --yara-rules rules.yar # Scan specific process vol -f memory.raw windows.yarascan --yara-rules rules.yar --pid 1234 # Scan kernel memory vol -f memory.raw windows.yarascan --yara-rules rules.yar --kernel ``` ## String Analysis ### Extracting Strings ```bash # Basic string extraction strings -a memory.raw > all_strings.txt # Unicode strings strings -el memory.raw >> all_strings.txt # Targeted extraction from process dump vol -f memory.raw windows.memmap --pid 1234 --dump strings -a pid.1234.dmp > process_strings.txt # Pattern matching grep -E "(https?://|[0-9]{1,3}\.[0-9]{1,3}\.[0-9]{1,3}\.[0-9]{1,3})" all_strings.txt ``` ### FLOSS for Obfuscated Strings ```bash # FLOSS extracts obfuscated strings floss malware.exe > floss_output.txt # From memory dump floss pid.1234.dmp ``` ## Best Practices ### Acquisition Best Practices 1. **Minimize footprint**: Use lightweight acquisition tools 2. **Document everything**: Record time, tool, and hash of capture 3. **Verify integrity**: Hash memory dump immediately after capture 4. **Chain of custody**: Maintain proper forensic handling ### Analysis Best Practices 1. **Start broad**: Get overview before deep diving 2. **Cross-reference**: Use multiple plugins for same data 3. **Timeline correlation**: Correlate memory findings with disk/network 4. **Document findings**: Keep detailed notes and screenshots 5. **Validate results**: Verify findings through multiple methods ### Common Pitfalls - **Stale data**: Memory is volatile, analyze promptly - **Incomplete dumps**: Verify dump size matches expected RAM - **Symbol issues**: Ensure correct symbol files for OS version - **Smear**: Memory may change during acquisition - **Encryption**: Some data may be encrypted in memory
👍0
👁️0
🤖 Auto-discovered
🤖system prompt•7 months ago

godot-gdscript-patterns

Master Godot 4 GDScript patterns including signals, scenes, state

coding
⭐1
# Godot GDScript Patterns Production patterns for Godot 4.x game development with GDScript, covering architecture, signals, scenes, and optimization. ## When to Use This Skill - Building games with Godot 4 - Implementing game systems in GDScript - Designing scene architecture - Managing game state - Optimizing GDScript performance - Learning Godot best practices ## Core Concepts ### 1. Godot Architecture ``` Node: Base building block ├── Scene: Reusable node tree (saved as .tscn) ├── Resource: Data container (saved as .tres) ├── Signal: Event communication └── Group: Node categorization ``` ### 2. GDScript Basics ```gdscript class_name Player extends CharacterBody2D # Signals signal health_changed(new_health: int) signal died # Exports (Inspector-editable) @export var speed: float = 200.0 @export var max_health: int = 100 @export_range(0, 1) var damage_reduction: float = 0.0 @export_group("Combat") @export var attack_damage: int = 10 @export var attack_cooldown: float = 0.5 # Onready (initialized when ready) @onready var sprite: Sprite2D = $Sprite2D @onready var animation: AnimationPlayer = $AnimationPlayer @onready var hitbox: Area2D = $Hitbox # Private variables (convention: underscore prefix) var _health: int var _can_attack: bool = true func _ready() -> void: _health = max_health func _physics_process(delta: float) -> void: var direction := Input.get_vector("left", "right", "up", "down") velocity = direction * speed move_and_slide() func take_damage(amount: int) -> void: var actual_damage := int(amount * (1.0 - damage_reduction)) _health = max(_health - actual_damage, 0) health_changed.emit(_health) if _health <= 0: died.emit() ``` ## Patterns ### Pattern 1: State Machine ```gdscript # state_machine.gd class_name StateMachine extends Node signal state_changed(from_state: StringName, to_state: StringName) @export var initial_state: State var current_state: State var states: Dictionary = {} func _ready() -> void: # Register all State children for child in get_children(): if child is State: states[child.name] = child child.state_machine = self child.process_mode = Node.PROCESS_MODE_DISABLED # Start initial state if initial_state: current_state = initial_state current_state.process_mode = Node.PROCESS_MODE_INHERIT current_state.enter() func _process(delta: float) -> void: if current_state: current_state.update(delta) func _physics_process(delta: float) -> void: if current_state: current_state.physics_update(delta) func _unhandled_input(event: InputEvent) -> void: if current_state: current_state.handle_input(event) func transition_to(state_name: StringName, msg: Dictionary = {}) -> void: if not states.has(state_name): push_error("State '%s' not found" % state_name) return var previous_state := current_state previous_state.exit() previous_state.process_mode = Node.PROCESS_MODE_DISABLED current_state = states[state_name] current_state.process_mode = Node.PROCESS_MODE_INHERIT current_state.enter(msg) state_changed.emit(previous_state.name, current_state.name) ``` ```gdscript # state.gd class_name State extends Node var state_machine: StateMachine func enter(_msg: Dictionary = {}) -> void: pass func exit() -> void: pass func update(_delta: float) -> void: pass func physics_update(_delta: float) -> void: pass func handle_input(_event: InputEvent) -> void: pass ``` ```gdscript # player_idle.gd class_name PlayerIdle extends State @export var player: Player func enter(_msg: Dictionary = {}) -> void: player.animation.play("idle") func physics_update(_delta: float) -> void: var direction := Input.get_vector("left", "right", "up", "down") if direction != Vector2.ZERO: state_machine.transition_to("Move") func handle_input(event: InputEvent) -> void: if event.is_action_pressed("attack"): state_machine.transition_to("Attack") elif event.is_action_pressed("jump"): state_machine.transition_to("Jump") ``` ### Pattern 2: Autoload Singletons ```gdscript # game_manager.gd (Add to Project Settings > Autoload) extends Node signal game_started signal game_paused(is_paused: bool) signal game_over(won: bool) signal score_changed(new_score: int) enum GameState { MENU, PLAYING, PAUSED, GAME_OVER } var state: GameState = GameState.MENU var score: int = 0: set(value): score = value score_changed.emit(score) var high_score: int = 0 func _ready() -> void: process_mode = Node.PROCESS_MODE_ALWAYS _load_high_score() func _input(event: InputEvent) -> void: if event.is_action_pressed("pause") and state == GameState.PLAYING: toggle_pause() func start_game() -> void: score = 0 state = GameState.PLAYING game_started.emit() func toggle_pause() -> void: var is_paused := state != GameState.PAUSED if is_paused: state = GameState.PAUSED get_tree().paused = true else: state = GameState.PLAYING get_tree().paused = false game_paused.emit(is_paused) func end_game(won: bool) -> void: state = GameState.GAME_OVER if score > high_score: high_score = score _save_high_score() game_over.emit(won) func add_score(points: int) -> void: score += points func _load_high_score() -> void: if FileAccess.file_exists("user://high_score.save"): var file := FileAccess.open("user://high_score.save", FileAccess.READ) high_score = file.get_32() func _save_high_score() -> void: var file := FileAccess.open("user://high_score.save", FileAccess.WRITE) file.store_32(high_score) ``` ```gdscript # event_bus.gd (Global signal bus) extends Node # Player events signal player_spawned(player: Node2D) signal player_died(player: Node2D) signal player_health_changed(health: int, max_health: int) # Enemy events signal enemy_spawned(enemy: Node2D) signal enemy_died(enemy: Node2D, position: Vector2) # Item events signal item_collected(item_type: StringName, value: int) signal powerup_activated(powerup_type: StringName) # Level events signal level_started(level_number: int) signal level_completed(level_number: int, time: float) signal checkpoint_reached(checkpoint_id: int) ``` ### Pattern 3: Resource-based Data ```gdscript # weapon_data.gd class_name WeaponData extends Resource @export var name: StringName @export var damage: int @export var attack_speed: float @export var range: float @export_multiline var description: String @export var icon: Texture2D @export var projectile_scene: PackedScene @export var sound_attack: AudioStream ``` ```gdscript # character_stats.gd class_name CharacterStats extends Resource signal stat_changed(stat_name: StringName, new_value: float) @export var max_health: float = 100.0 @export var attack: float = 10.0 @export var defense: float = 5.0 @export var speed: float = 200.0 # Runtime values (not saved) var _current_health: float func _init() -> void: _current_health = max_health func get_current_health() -> float: return _current_health func take_damage(amount: float) -> float: var actual_damage := maxf(amount - defense, 1.0) _current_health = maxf(_current_health - actual_damage, 0.0) stat_changed.emit("health", _current_health) return actual_damage func heal(amount: float) -> void: _current_health = minf(_current_health + amount, max_health) stat_changed.emit("health", _current_health) func duplicate_for_runtime() -> CharacterStats: var copy := duplicate() as CharacterStats copy._current_health = copy.max_health return copy ``` ```gdscript # Using resources class_name Character extends CharacterBody2D @export var base_stats: CharacterStats @export var weapon: WeaponData var stats: CharacterStats func _ready() -> void: # Create runtime copy to avoid modifying the resource stats = base_stats.duplicate_for_runtime() stats.stat_changed.connect(_on_stat_changed) func attack() -> void: if weapon: print("Attacking with %s for %d damage" % [weapon.name, weapon.damage]) func _on_stat_changed(stat_name: StringName, value: float) -> void: if stat_name == "health" and value <= 0: die() ``` ### Pattern 4: Object Pooling ```gdscript # object_pool.gd class_name ObjectPool extends Node @export var pooled_scene: PackedScene @export var initial_size: int = 10 @export var can_grow: bool = true var _available: Array[Node] = [] var _in_use: Array[Node] = [] func _ready() -> void: _initialize_pool() func _initialize_pool() -> void: for i in initial_size: _create_instance() func _create_instance() -> Node: var instance := pooled_scene.instantiate() instance.process_mode = Node.PROCESS_MODE_DISABLED instance.visible = false add_child(instance) _available.append(instance) # Connect return signal if exists if instance.has_signal("returned_to_pool"): instance.returned_to_pool.connect(_return_to_pool.bind(instance)) return instance func get_instance() -> Node: var instance: Node if _available.is_empty(): if can_grow: instance = _create_instance() _available.erase(instance) else: push_warning("Pool exhausted and cannot grow") return null else: instance = _available.pop_back() instance.process_mode = Node.PROCESS_MODE_INHERIT instance.visible = true _in_use.append(instance) if instance.has_method("on_spawn"): instance.on_spawn() return instance func _return_to_pool(instance: Node) -> void: if not instance in _in_use: return _in_use.erase(instance) if instance.has_method("on_despawn"): instance.on_despawn() instance.process_mode = Node.PROCESS_MODE_DISABLED instance.visible = false _available.append(instance) func return_all() -> void: for instance in _in_use.duplicate(): _return_to_pool(instance) ``` ```gdscript # pooled_bullet.gd class_name PooledBullet extends Area2D signal returned_to_pool @export var speed: float = 500.0 @export var lifetime: float = 5.0 var direction: Vector2 var _timer: float func on_spawn() -> void: _timer = lifetime func on_despawn() -> void: direction = Vector2.ZERO func initialize(pos: Vector2, dir: Vector2) -> void: global_position = pos direction = dir.normalized() rotation = direction.angle() func _physics_process(delta: float) -> void: position += direction * speed * delta _timer -= delta if _timer <= 0: returned_to_pool.emit() func _on_body_entered(body: Node2D) -> void: if body.has_method("take_damage"): body.take_damage(10) returned_to_pool.emit() ``` ### Pattern 5: Component System ```gdscript # health_component.gd class_name HealthComponent extends Node signal health_changed(current: int, maximum: int) signal damaged(amount: int, source: Node) signal healed(amount: int) signal died @export var max_health: int = 100 @export var invincibility_time: float = 0.0 var current_health: int: set(value): var old := current_health current_health = clampi(value, 0, max_health) if current_health != old: health_changed.emit(current_health, max_health) var _invincible: bool = false func _ready() -> void: current_health = max_health func take_damage(amount: int, source: Node = null) -> int: if _invincible or current_health <= 0: return 0 var actual := mini(amount, current_health) current_health -= actual damaged.emit(actual, source) if current_health <= 0: died.emit() elif invincibility_time > 0: _start_invincibility() return actual func heal(amount: int) -> int: var actual := mini(amount, max_health - current_health) current_health += actual if actual > 0: healed.emit(actual) return actual func _start_invincibility() -> void: _invincible = true await get_tree().create_timer(invincibility_time).timeout _invincible = false ``` ```gdscript # hitbox_component.gd class_name HitboxComponent extends Area2D signal hit(hurtbox: HurtboxComponent) @export var damage: int = 10 @export var knockback_force: float = 200.0 var owner_node: Node func _ready() -> void: owner_node = get_parent() area_entered.connect(_on_area_entered) func _on_area_entered(area: Area2D) -> void: if area is HurtboxComponent: var hurtbox := area as HurtboxComponent if hurtbox.owner_node != owner_node: hit.emit(hurtbox) hurtbox.receive_hit(self) ``` ```gdscript # hurtbox_component.gd class_name HurtboxComponent extends Area2D signal hurt(hitbox: HitboxComponent) @export var health_component: HealthComponent var owner_node: Node func _ready() -> void: owner_node = get_parent() func receive_hit(hitbox: HitboxComponent) -> void: hurt.emit(hitbox) if health_component: health_component.take_damage(hitbox.damage, hitbox.owner_node) ``` ### Pattern 6: Scene Management ```gdscript # scene_manager.gd (Autoload) extends Node signal scene_loading_started(scene_path: String) signal scene_loading_progress(progress: float) signal scene_loaded(scene: Node) signal transition_started signal transition_finished @export var transition_scene: PackedScene @export var loading_scene: PackedScene var _current_scene: Node var _transition: CanvasLayer var _loader: ResourceLoader func _ready() -> void: _current_scene = get_tree().current_scene if transition_scene: _transition = transition_scene.instantiate() add_child(_transition) _transition.visible = false func change_scene(scene_path: String, with_transition: bool = true) -> void: if with_transition: await _play_transition_out() _load_scene(scene_path) func change_scene_packed(scene: PackedScene, with_transition: bool = true) -> void: if with_transition: await _play_transition_out() _swap_scene(scene.instantiate()) func _load_scene(path: String) -> void: scene_loading_started.emit(path) # Check if already loaded if ResourceLoader.has_cached(path): var scene := load(path) as PackedScene _swap_scene(scene.instantiate()) return # Async loading ResourceLoader.load_threaded_request(path) while true: var progress := [] var status := ResourceLoader.load_threaded_get_status(path, progress) match status: ResourceLoader.THREAD_LOAD_IN_PROGRESS: scene_loading_progress.emit(progress[0]) await get_tree().process_frame ResourceLoader.THREAD_LOAD_LOADED: var scene := ResourceLoader.load_threaded_get(path) as PackedScene _swap_scene(scene.instantiate()) return _: push_error("Failed to load scene: %s" % path) return func _swap_scene(new_scene: Node) -> void: if _current_scene: _current_scene.queue_free() _current_scene = new_scene get_tree().root.add_child(_current_scene) get_tree().current_scene = _current_scene scene_loaded.emit(_current_scene) await _play_transition_in() func _play_transition_out() -> void: if not _transition: return transition_started.emit() _transition.visible = true if _transition.has_method("transition_out"): await _transition.transition_out() else: await get_tree().create_timer(0.3).timeout func _play_transition_in() -> void: if not _transition: transition_finished.emit() return if _transition.has_method("transition_in"): await _transition.transition_in() else: await get_tree().create_timer(0.3).timeout _transition.visible = false transition_finished.emit() ``` ### Pattern 7: Save System ```gdscript # save_manager.gd (Autoload) extends Node const SAVE_PATH := "user://savegame.save" const ENCRYPTION_KEY := "your_secret_key_here" signal save_completed signal load_completed signal save_error(message: String) func save_game(data: Dictionary) -> void: var file := FileAccess.open_encrypted_with_pass( SAVE_PATH, FileAccess.WRITE, ENCRYPTION_KEY ) if file == null: save_error.emit("Could not open save file") return var json := JSON.stringify(data) file.store_string(json) file.close() save_completed.emit() func load_game() -> Dictionary: if not FileAccess.file_exists(SAVE_PATH): return {} var file := FileAccess.open_encrypted_with_pass( SAVE_PATH, FileAccess.READ, ENCRYPTION_KEY ) if file == null: save_error.emit("Could not open save file") return {} var json := file.get_as_text() file.close() var parsed := JSON.parse_string(json) if parsed == null: save_error.emit("Could not parse save data") return {} load_completed.emit() return parsed func delete_save() -> void: if FileAccess.file_exists(SAVE_PATH): DirAccess.remove_absolute(SAVE_PATH) func has_save() -> bool: return FileAccess.file_exists(SAVE_PATH) ``` ```gdscript # saveable.gd (Attach to saveable nodes) class_name Saveable extends Node @export var save_id: String func _ready() -> void: if save_id.is_empty(): save_id = str(get_path()) func get_save_data() -> Dictionary: var parent := get_parent() var data := {"id": save_id} if parent is Node2D: data["position"] = {"x": parent.position.x, "y": parent.position.y} if parent.has_method("get_custom_save_data"): data.merge(parent.get_custom_save_data()) return data func load_save_data(data: Dictionary) -> void: var parent := get_parent() if data.has("position") and parent is Node2D: parent.position = Vector2(data.position.x, data.position.y) if parent.has_method("load_custom_save_data"): parent.load_custom_save_data(data) ``` ## Performance Tips ```gdscript # 1. Cache node references @onready var sprite := $Sprite2D # Good # $Sprite2D in _process() # Bad - repeated lookup # 2. Use object pooling for frequent spawning # See Pattern 4 # 3. Avoid allocations in hot paths var _reusable_array: Array = [] func _process(_delta: float) -> void: _reusable_array.clear() # Reuse instead of creating new # 4. Use static typing func calculate(value: float) -> float: # Good return value * 2.0 # 5. Disable processing when not needed func _on_off_screen() -> void: set_process(false) set_physics_process(false) ``` ## Best Practices ### Do's - **Use signals for decoupling** - Avoid direct references - **Type everything** - Static typing catches errors - **Use resources for data** - Separate data from logic - **Pool frequently spawned objects** - Avoid GC hitches - **Use Autoloads sparingly** - Only for truly global systems ### Don'ts - **Don't use `get_node()` in loops** - Cache ref
👍0
👁️0
🤖 Auto-discovered
🤖system prompt•7 months ago

event-store-design

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

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

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

track-management

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

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

workflow-patterns

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

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

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