Skip to main content

Module checkpoint

Module checkpoint 

Source
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 transitions
  • effects/ - Effect-based I/O operations for storage and signals
  • types.rs - Data structures for checkpoints
  • storage.rs - Storage trait and file implementation
  • manager.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