Expand description
Checkpoint management for MapReduce workflows
This module provides checkpoint structures and management for all MapReduce phases.
§Architecture
The checkpoint system follows a “pure core, imperative shell” pattern:
pure/- Pure functions for validation, preparation, triggers, state transitionseffects/- Effect-based I/O operations for storage and signalstypes.rs- Data structures for checkpointsstorage.rs- Storage trait and file implementationmanager.rs- High-level checkpoint manager
§Incremental Checkpointing (Spec 162)
Checkpoints are created:
- After every N agent completions (configurable)
- At time intervals (configurable)
- On signal (SIGINT/SIGTERM) for graceful shutdown
- At phase transitions (setup -> map -> reduce)
All checkpoints are stored in ~/.prodigy/state/{repo}/mapreduce/jobs/{job_id}/
Re-exports§
pub use types::AgentInfo;pub use types::AgentState;pub use types::CheckpointConfig;pub use types::CheckpointId;pub use types::CheckpointInfo;pub use types::CheckpointMetadata;pub use types::CheckpointReason;pub use types::CompletedWorkItem;pub use types::DlqItem;pub use types::ErrorState;pub use types::ExecutionState;pub use types::FailedWorkItem;pub use types::MapPhaseResults;pub use types::MapReduceCheckpoint;pub use types::PhaseResult;pub use types::PhaseType;pub use types::ResourceAllocation;pub use types::ResourceState;pub use types::ResumeState;pub use types::ResumeStrategy;pub use types::RetentionPolicy;pub use types::VariableState;pub use types::WorkItem;pub use types::WorkItemBatch;pub use types::WorkItemProgress;pub use types::WorkItemState;pub use storage::CheckpointStorage;pub use storage::CompressionAlgorithm;pub use storage::FileCheckpointStorage;pub use manager::CheckpointManager;pub use reduce::ReducePhaseCheckpoint;pub use reduce::StepResult;pub use pure::calculate_integrity_hash;pub use pure::prepare_checkpoint;pub use pure::reset_in_progress_items;pub use pure::should_checkpoint;pub use pure::transition_work_item;pub use pure::CheckpointTriggerConfig;pub use pure::CheckpointValidationError;pub use pure::WorkItemEvent;pub use pure::WorkItemStatus;pub use effects::load_checkpoint_effect;pub use effects::save_checkpoint_effect;pub use effects::save_checkpoint_on_shutdown;pub use effects::shutdown_signal;pub use effects::CheckpointOnShutdown;pub use effects::CheckpointStorageEnv;pub use effects::CheckpointStorageError;pub use effects::ShutdownSignal;pub use environment::get_checkpoint_job_id;pub use environment::get_checkpoint_storage;pub use environment::get_checkpoint_storage_path;pub use environment::get_items_since_checkpoint;pub use environment::get_trigger_config;pub use environment::is_checkpointing_enabled;pub use environment::with_checkpointing_disabled;pub use environment::with_trigger_config;pub use environment::CheckpointEnv;pub use environment::CheckpointError;pub use environment::MockCheckpointEnvBuilder;pub use incremental::CheckpointStats;pub use incremental::IncrementalCheckpointController;
Modules§
- effects
- Checkpoint effects for I/O operations
- environment
- Checkpoint environment for Reader pattern effects
- incremental
- Incremental checkpointing for MapReduce workflows
- manager
- Enhanced MapReduce checkpoint management
- pure
- Pure functions for checkpoint management
- reduce
- Reduce phase checkpoint structures and methods for MapReduce workflows
- storage
- Checkpoint storage implementations and compression
- types
- Checkpoint data structures and state types