Class: Smith::Workflow
- Inherits:
-
Object
- Object
- Smith::Workflow
- Includes:
- ArtifactIntegration, BudgetIntegration, DSL, DataVolumePolicy, DeadlineEnforcement, DefinitionIdentity, Durability, EventIntegration, Execution, GuardrailIntegration, MessageAdmissionBoundary, Persistence, SplitStepPersistence
- Defined in:
- lib/smith/workflow.rb,
lib/smith/workflow/dsl.rb,
lib/smith/workflow/claim.rb,
lib/smith/workflow/graph.rb,
lib/smith/workflow/router.rb,
lib/smith/workflow/parallel.rb,
lib/smith/workflow/pipeline.rb,
lib/smith/workflow/execution.rb,
lib/smith/workflow/graph_dsl.rb,
lib/smith/workflow/branch_env.rb,
lib/smith/workflow/durability.rb,
lib/smith/workflow/identifier.rb,
lib/smith/workflow/run_result.rb,
lib/smith/workflow/transition.rb,
lib/smith/workflow/persistence.rb,
lib/smith/workflow/usage_entry.rb,
lib/smith/workflow/agent_result.rb,
lib/smith/workflow/graph/report.rb,
lib/smith/workflow/step_context.rb,
lib/smith/workflow/graph/metrics.rb,
lib/smith/workflow/graph/targets.rb,
lib/smith/workflow/message_batch.rb,
lib/smith/workflow/prepared_step.rb,
lib/smith/workflow/process_local.rb,
lib/smith/workflow/composite/plan.rb,
lib/smith/workflow/failure_record.rb,
lib/smith/workflow/composite/enums.rb,
lib/smith/workflow/composite/error.rb,
lib/smith/workflow/composite/input.rb,
lib/smith/workflow/execution_frame.rb,
lib/smith/workflow/graph/reference.rb,
lib/smith/workflow/graph/validator.rb,
lib/smith/workflow/retry_execution.rb,
lib/smith/workflow/step_completion.rb,
lib/smith/workflow/string_snapshot.rb,
lib/smith/workflow/composite/branch.rb,
lib/smith/workflow/fanout_execution.rb,
lib/smith/workflow/graph/diagnostic.rb,
lib/smith/workflow/nested_execution.rb,
lib/smith/workflow/worker_execution.rb,
lib/smith/workflow/composite/effects.rb,
lib/smith/workflow/composite/payload.rb,
lib/smith/workflow/composite/planner.rb,
lib/smith/workflow/composite/reducer.rb,
lib/smith/workflow/event_integration.rb,
lib/smith/workflow/message_admission.rb,
lib/smith/workflow/budget_integration.rb,
lib/smith/workflow/composite/contract.rb,
lib/smith/workflow/data_volume_policy.rb,
lib/smith/workflow/deterministic_step.rb,
lib/smith/workflow/graph/reachability.rb,
lib/smith/workflow/optimization_state.rb,
lib/smith/workflow/parallel_execution.rb,
lib/smith/workflow/composite/reduction.rb,
lib/smith/workflow/definition_identity.rb,
lib/smith/workflow/evaluator_optimizer.rb,
lib/smith/workflow/failure_record_text.rb,
lib/smith/workflow/orchestration_state.rb,
lib/smith/workflow/orchestrator_worker.rb,
lib/smith/workflow/artifact_integration.rb,
lib/smith/workflow/deadline_enforcement.rb,
lib/smith/workflow/composite/budget_math.rb,
lib/smith/workflow/composite/outcome_set.rb,
lib/smith/workflow/composite/preparation.rb,
lib/smith/workflow/failure_reconstructor.rb,
lib/smith/workflow/graph/diagnostic_path.rb,
lib/smith/workflow/graph/fanout_contract.rb,
lib/smith/workflow/guardrail_integration.rb,
lib/smith/workflow/parallel/cancellation.rb,
lib/smith/workflow/composite/value_budget.rb,
lib/smith/workflow/failure_record_restore.rb,
lib/smith/workflow/graph/contract_helpers.rb,
lib/smith/workflow/guarded_step_execution.rb,
lib/smith/workflow/parallel_agent_binding.rb,
lib/smith/workflow/prepared_step_dispatch.rb,
lib/smith/workflow/prepared_step_recovery.rb,
lib/smith/workflow/split_step_persistence.rb,
lib/smith/workflow/deterministic_execution.rb,
lib/smith/workflow/failure_detail_snapshot.rb,
lib/smith/workflow/graph/runtime_readiness.rb,
lib/smith/workflow/graph/state_diagnostics.rb,
lib/smith/workflow/parallel/root_execution.rb,
lib/smith/workflow/thread_context_snapshot.rb,
lib/smith/workflow/composite/branch_failure.rb,
lib/smith/workflow/composite/branch_outcome.rb,
lib/smith/workflow/composite/error_evidence.rb,
lib/smith/workflow/composite/payload_digest.rb,
lib/smith/workflow/composite/plan_integrity.rb,
lib/smith/workflow/failure_record_validator.rb,
lib/smith/workflow/message_value_normalizer.rb,
lib/smith/workflow/transition_actionability.rb,
lib/smith/workflow/composite/branch_contract.rb,
lib/smith/workflow/execution_result_snapshot.rb,
lib/smith/workflow/graph/transition_contract.rb,
lib/smith/workflow/graph/transition_snapshot.rb,
lib/smith/workflow/parallel/nested_execution.rb,
lib/smith/workflow/prepared_branch_execution.rb,
lib/smith/workflow/composite/branch_execution.rb,
lib/smith/workflow/composite/budget_allocator.rb,
lib/smith/workflow/composite/effects_baseline.rb,
lib/smith/workflow/graph/execution_successors.rb,
lib/smith/workflow/message_admission_boundary.rb,
lib/smith/workflow/parallel/execution_context.rb,
lib/smith/workflow/composite/effects_preflight.rb,
lib/smith/workflow/graph/identifier_projection.rb,
lib/smith/workflow/graph/optimization_contract.rb,
lib/smith/workflow/composite/execution_contract.rb,
lib/smith/workflow/execution_binding_resolution.rb,
lib/smith/workflow/graph/orchestration_contract.rb,
lib/smith/workflow/graph/transition_diagnostics.rb,
lib/smith/workflow/parallel/cancellation_signal.rb,
lib/smith/workflow/composite/effects_application.rb,
lib/smith/workflow/composite/outcome_accumulator.rb,
lib/smith/workflow/graph/retry_policy_diagnostic.rb,
lib/smith/workflow/prepared_step_execution_scope.rb,
lib/smith/workflow/composite/encoded_value_budget.rb,
lib/smith/workflow/graph/reachability_diagnostics.rb,
lib/smith/workflow/graph/runtime_readiness_report.rb,
lib/smith/workflow/prepared_step_execution_result.rb,
lib/smith/workflow/graph/runtime_readiness_metrics.rb,
lib/smith/workflow/split_step_persistence/boundary.rb,
lib/smith/workflow/split_step_persistence/payloads.rb,
lib/smith/workflow/split_step_persistence/recovery.rb,
lib/smith/workflow/composite/branch_budget_contract.rb,
lib/smith/workflow/composite/fanout_branch_contract.rb,
lib/smith/workflow/split_step_persistence/execution.rb,
lib/smith/workflow/definition_identity/class_methods.rb,
lib/smith/workflow/graph/runtime_binding_diagnostics.rb,
lib/smith/workflow/graph/runtime_readiness_traversal.rb,
lib/smith/workflow/split_step_persistence/checkpoint.rb,
lib/smith/workflow/graph/nested_readiness_diagnostics.rb,
lib/smith/workflow/split_step_persistence/inheritance.rb,
lib/smith/workflow/split_step_persistence/preparation.rb,
lib/smith/workflow/graph/transition_contract_attributes.rb,
lib/smith/workflow/prepared_step_execution_authorization.rb,
lib/smith/workflow/split_step_persistence/boundary_reset.rb,
lib/smith/workflow/split_step_persistence/dispatch_claim.rb,
lib/smith/workflow/split_step_persistence/state_snapshot.rb,
lib/smith/workflow/graph/runtime_readiness_report_builder.rb,
lib/smith/workflow/split_step_persistence/checkpoint_state.rb,
lib/smith/workflow/composite_branch_execution_authorization.rb,
lib/smith/workflow/graph/runtime_binding_diagnostic_builder.rb,
lib/smith/workflow/graph/transition_contract_configurations.rb,
lib/smith/workflow/split_step_persistence/dispatch_boundary.rb,
lib/smith/workflow/split_step_persistence/preparation_claim.rb,
lib/smith/workflow/split_step_persistence/recovery_boundary.rb,
lib/smith/workflow/split_step_persistence/subclass_boundary.rb,
lib/smith/workflow/split_step_persistence/composite_execution.rb,
lib/smith/workflow/split_step_persistence/definition_boundary.rb,
lib/smith/workflow/split_step_persistence/execution_lifecycle.rb,
lib/smith/workflow/split_step_persistence/preparation_payload.rb,
lib/smith/workflow/split_step_persistence/transition_contract.rb,
lib/smith/workflow/graph/transition_optimization_configuration.rb,
lib/smith/workflow/split_step_persistence/preparation_recovery.rb,
lib/smith/workflow/split_step_persistence/restart_safe_adapter.rb,
lib/smith/workflow/split_step_persistence/transaction_identity.rb,
lib/smith/workflow/split_step_persistence/composite_preparation.rb,
lib/smith/workflow/split_step_persistence/dispatch_confirmation.rb,
lib/smith/workflow/split_step_persistence/dispatch_verification.rb,
lib/smith/workflow/split_step_persistence/execution_verification.rb,
lib/smith/workflow/split_step_persistence/recovery_class_methods.rb,
lib/smith/workflow/split_step_persistence/execution_authorization.rb,
lib/smith/workflow/split_step_persistence/replace_exact_signature.rb,
lib/smith/workflow/split_step_persistence/canonical_payload_digest.rb,
lib/smith/workflow/split_step_persistence/composite_branch_effects.rb,
lib/smith/workflow/split_step_persistence/composite_branch_outcome.rb,
lib/smith/workflow/split_step_persistence/execution_result_capture.rb,
lib/smith/workflow/split_step_persistence/composite_branch_execution.rb,
lib/smith/workflow/split_step_persistence/execution_binding_snapshot.rb,
lib/smith/workflow/split_step_persistence/execution_binding_collector.rb,
lib/smith/workflow/split_step_persistence/execution_workflow_snapshot.rb,
lib/smith/workflow/split_step_persistence/transition_contract_freezer.rb,
lib/smith/workflow/split_step_persistence/composite_reduction_execution.rb,
lib/smith/workflow/split_step_persistence/transition_contract_signature.rb,
lib/smith/workflow/split_step_persistence/composite_branch_authorization.rb,
lib/smith/workflow/split_step_persistence/execution_authorization_issuance.rb,
lib/smith/workflow/split_step_persistence/transition_contract_structured_values.rb
Defined Under Namespace
Modules: ArtifactIntegration, BudgetIntegration, Claim, Composite, DSL, DataVolumePolicy, DeadlineEnforcement, DefinitionIdentity, DeterministicExecution, Durability, EvaluatorOptimizer, EventIntegration, Execution, ExecutionBindingResolution, FanoutExecution, GuardedStepExecution, GuardrailIntegration, MessageAdmissionBoundary, NestedExecution, OrchestratorWorker, ParallelExecution, Persistence, RetryExecution, SplitStepPersistence, StepCompletion, StepContext, TransitionActionability Classes: AgentResult, BranchEnv, CompositeBranchExecutionAuthorization, DeterministicStep, ExecutionFrame, ExecutionResultSnapshot, FailureDetailSnapshot, FailureReconstructor, FailureRecord, FailureRecordRestore, FailureRecordText, FailureRecordValidator, Graph, Identifier, MessageAdmission, MessageBatch, MessageValueNormalizer, OptimizationState, OrchestrationState, Parallel, ParallelAgentBinding, Pipeline, PreparedStep, PreparedStepDispatch, PreparedStepExecutionAuthorization, PreparedStepExecutionResult, PreparedStepExecutionScope, PreparedStepRecovery, Router, RunResult, StringSnapshot, Transition, UsageEntry, WorkerExecution
Constant Summary collapse
- DEFAULT_MAX_TRANSITIONS =
100
Constants included from DataVolumePolicy
DataVolumePolicy::LIGHTWEIGHT_SCALARS
Constants included from BudgetIntegration
BudgetIntegration::AGENT_DIM_MAP, BudgetIntegration::BUDGET_DIMENSIONS, BudgetIntegration::COST_DIMENSIONS, BudgetIntegration::TOKEN_DIMENSIONS
Constants included from SplitStepPersistence
SplitStepPersistence::NO_SPLIT_TRANSITION
Instance Attribute Summary collapse
-
#last_prepared_input ⇒ Object
readonly
Returns the value of attribute last_prepared_input.
-
#ledger ⇒ Object
readonly
Returns the value of attribute ledger.
-
#state ⇒ Object
readonly
Returns the value of attribute state.
Class Method Summary collapse
Instance Method Summary collapse
- #advance! ⇒ Object
- #done? ⇒ Boolean
- #failed? ⇒ Boolean
-
#initialize(context: {}, ledger: nil, created_at: nil) ⇒ Workflow
constructor
A new instance of Workflow.
- #pending_transition_name ⇒ Object
- #persisted_keys ⇒ Object
- #record_persisted_key!(key) ⇒ Object
- #run! ⇒ Object
- #session_messages ⇒ Object
- #terminal? ⇒ Boolean
-
#usage_entries ⇒ Object
Public, read-only view of the per-provider-call usage ledger, so hosts can diff usage across a step boundary without paying a full to_state serialization.
Methods included from MessageAdmissionBoundary
Methods included from Durability
#advance_persisted!, #clear_persisted!, #clear_step_in_progress!, included, #mark_step_in_progress!, #persist!, #run_persisted!
Methods included from Persistence
Methods included from DefinitionIdentity
Methods included from DSL
Methods included from SplitStepPersistence
Methods included from SplitStepPersistence::Checkpoint
#complete_persisted_step!, #persist!
Methods included from SplitStepPersistence::CompositeExecution
#execute_prepared_composite_branch!, #prepare_composite_step!, #reduce_prepared_composite_step!
Methods included from SplitStepPersistence::ExecutionAuthorizationIssuance
#authorize_prepared_step_execution!
Methods included from SplitStepPersistence::ExecutionAuthorization
#release_prepared_step_execution!
Methods included from SplitStepPersistence::DispatchConfirmation
#confirm_prepared_step_dispatch!
Methods included from SplitStepPersistence::DispatchClaim
#claim_prepared_step_dispatch!
Methods included from SplitStepPersistence::Preparation
#confirm_prepared_step!, #prepare_persisted_step!, #prepared_persisted_step
Methods included from SplitStepPersistence::StateSnapshot
Methods included from SplitStepPersistence::Boundary
#advance_persisted!, #clear_persisted!, #run_persisted!
Constructor Details
#initialize(context: {}, ledger: nil, created_at: nil) ⇒ Workflow
Returns a new instance of Workflow.
41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 |
# File 'lib/smith/workflow.rb', line 41 def initialize(context: {}, ledger: nil, created_at: nil) @state = self.class.initial_state @context = context @step_count = 0 @next_transition_name = nil @ledger = ledger || build_ledger @created_at = created_at || Time.now.utc.iso8601 @updated_at = @created_at @total_cost = 0.0 @total_tokens = 0 @outcome = nil # Eager init for usage tracking. Both `@usage_mutex` (lazy # init at the call site would race across parallel fan-out # branches) and the durable per-call/output/failure fields # must be present before any agent recording fires. # `restore_state` mirrors these inits because `from_state` uses # `allocate` and bypasses `initialize` — see persistence.rb. @usage_entries = [] @usage_mutex = Mutex.new @last_output = nil # Durable attribution of the most recent serial `execute :agent` step: # { model:, provider: } of the model that actually served it (post # fallback resolution), exposed to the following deterministic step as # `last_agent_model` / `last_agent_provider`. `@pending_agent_execution` # is the transient per-step carrier from `execute_serial_step` to # `complete_step`; it is never persisted. @last_agent_execution = nil @pending_agent_execution = nil @last_failed_step = nil # Optimistic-locking version. Incremented on each persist!; restored # from the persisted payload. Adapters that support store_versioned # raise Smith::PersistenceVersionConflict when expected_version # doesn't match the stored payload's version (i.e., a concurrent # write occurred between this process's restore and persist). @persistence_version = 0 # Digest of the seed_messages produced at construction time. # Compared on restore against the live builder's output when # seed_validation is :warn or :strict; nil when no seed builder # ran or its output was empty. @seed_digest = nil @seed_message_count = 0 # Idempotency marker stamped between persist-before-advance and # persist-after-advance under idempotency_mode :strict; restored # workflows with the marker set raise # Smith::StepInProgressOnRestore. Lax mode leaves it false. @step_in_progress = false # Process-local phase for the split-step persistence API. This is not # serialized: strict restore rejects the durable step_in_progress marker # before an uncertain step can be resumed. @split_step_phase = nil @split_step_transition_name = nil @split_step_mutex = Mutex.new # Set of context keys recorded via deterministic step write_context # writes. Used by persist :auto Context mode to compute the # persisted-context slice. Seeded from the Context class's # also: declaration so explicit input keys round-trip. @persisted_keys = ::Set.new(initial_persist_auto_seed) @persisted_keys_mutex = Mutex.new initialize_tool_result_state end |
Instance Attribute Details
#last_prepared_input ⇒ Object (readonly)
Returns the value of attribute last_prepared_input.
35 36 37 |
# File 'lib/smith/workflow.rb', line 35 def last_prepared_input @last_prepared_input end |
#ledger ⇒ Object (readonly)
Returns the value of attribute ledger.
35 36 37 |
# File 'lib/smith/workflow.rb', line 35 def ledger @ledger end |
#state ⇒ Object (readonly)
Returns the value of attribute state.
35 36 37 |
# File 'lib/smith/workflow.rb', line 35 def state @state end |
Class Method Details
.graph ⇒ Object
5 6 7 8 9 10 11 12 |
# File 'lib/smith/workflow/graph_dsl.rb', line 5 def self.graph Graph.new( workflow_class: self, initial_state: initial_state, states: @states || [], transitions: @transitions || {} ) end |
.runtime_readiness ⇒ Object
18 19 20 |
# File 'lib/smith/workflow/graph_dsl.rb', line 18 def self.runtime_readiness graph.runtime_readiness end |
.validate_graph ⇒ Object
14 15 16 |
# File 'lib/smith/workflow/graph_dsl.rb', line 14 def self.validate_graph graph.validate end |
Instance Method Details
#advance! ⇒ Object
107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 |
# File 'lib/smith/workflow.rb', line 107 def advance! advance_claim = SplitStepPersistence.instance_method(:claim_split_step_advance!).bind_call(self) ensure_transition_budget! @step_work_started = false transition = SplitStepPersistence.instance_method(:resolve_split_step_advance_transition).bind_call(self) transition = resolve_transition if transition.equal?(SplitStepPersistence::NO_SPLIT_TRANSITION) return if transition.nil? @step_work_started = true step_result = if advance_claim == :split_step Execution.instance_method(:execute_step).bind_call(self, transition) else execute_step(transition) end @step_count += 1 @updated_at = Time.now.utc.iso8601 record_step_snapshot(step_result) step_result rescue UnresolvedTransitionError => e step_result = GuardrailIntegration .instance_method(:handle_unresolved_transition_failure) .bind_call(self, e) record_step_snapshot(step_result) step_result ensure SplitStepPersistence.instance_method(:release_split_step_advance!).bind_call(self, advance_claim) if advance_claim end |
#done? ⇒ Boolean
154 155 156 |
# File 'lib/smith/workflow.rb', line 154 def done? state_named?(:done) end |
#failed? ⇒ Boolean
158 159 160 |
# File 'lib/smith/workflow.rb', line 158 def failed? state_named?(:failed) end |
#pending_transition_name ⇒ Object
146 147 148 |
# File 'lib/smith/workflow.rb', line 146 def pending_transition_name @next_transition_name || self.class.first_transition_from(@state)&.name end |
#persisted_keys ⇒ Object
103 104 105 |
# File 'lib/smith/workflow.rb', line 103 def persisted_keys @persisted_keys.dup.freeze end |
#record_persisted_key!(key) ⇒ Object
162 163 164 165 166 |
# File 'lib/smith/workflow.rb', line 162 def record_persisted_key!(key) @persisted_keys_mutex.synchronize do @persisted_keys << key.to_sym end end |
#run! ⇒ Object
136 137 138 139 140 141 142 143 144 |
# File 'lib/smith/workflow.rb', line 136 def run! SplitStepPersistence.instance_method(:ensure_split_step_execution_allowed!).bind_call(self) steps = [] until terminal? step = advance! steps << step if step end build_run_result(steps) end |
#session_messages ⇒ Object
37 38 39 |
# File 'lib/smith/workflow.rb', line 37 def end |
#terminal? ⇒ Boolean
150 151 152 |
# File 'lib/smith/workflow.rb', line 150 def terminal? !self.class.transition_from?(@state) && @next_transition_name.nil? end |
#usage_entries ⇒ Object
Public, read-only view of the per-provider-call usage ledger, so hosts can diff usage across a step boundary without paying a full to_state serialization. The entries themselves are frozen; the returned array is a frozen copy taken under the recording mutex. Tolerates allocated but unrestored instances (no mutex yet), which have no entries.
173 174 175 176 177 178 |
# File 'lib/smith/workflow.rb', line 173 def usage_entries mutex = @usage_mutex return [].freeze unless mutex mutex.synchronize { @usage_entries.dup.freeze } end |