///|
///| Deterministic projection of an already-produced MoonClaw two-stage result.
///|
///| This package never invokes a model/provider and does not own an agent loop,
///| session, memory, or retry scheduler. MoonClaw owns those runtime concerns and
///|
/// supplies the immutable diagnosis/decision values projected and validated here.
pub(all) enum TwoStageOutcome {
Completed
AwaitingStage1Retry
AwaitingStage2Retry
Cancelled
Failed
} derive(Debug, Eq, ToJson, FromJson)
///|
pub(all) enum StageFailureKind {
SnapshotInvalid
GatewayUnavailable
ContextBudgetExceeded
Stage1Invalid
Stage2Invalid
OperatorCancelled
} derive(Debug, Eq, ToJson, FromJson)
///|
pub(all) struct TwoStageRoutineInput {
request : @routine.RoutineRequest
gateway : @model.ModelGateway
diagnosis : @domain.Stage1Diagnosis
decision : @domain.Stage2Decision
retry_policy : @job.RetryPolicy
stage1_attempt : Int
stage2_attempt : Int
cancel_after_phase : @routine.RoutinePhase?
} derive(Debug, Eq, ToJson, FromJson)
///|
pub(all) struct StageGate {
phase : @routine.RoutinePhase
validation : @domain.ValidationReport
retry_decision : @retry.RetryDecision
} derive(Debug, Eq, ToJson, FromJson)
///|
pub(all) struct TwoStageRoutineResult {
request : @routine.RoutineRequest
outcome : TwoStageOutcome
failure_kind : StageFailureKind?
record : @domain.AnalysisRecord
events : Array[@routine.RoutineEvent]
stage_gates : Array[StageGate]
stage1_prompt : @prompt.PromptPacket?
stage2_prompt : @prompt.PromptPacket?
context_budget : @model.ContextBudgetReport?
stage1_model_request : @model.ModelCallRequest?
stage2_model_request : @model.ModelCallRequest?
artifact_paths : Array[String]
partial_record_available : Bool
summary : String
} derive(Debug, Eq, ToJson, FromJson)
///|
pub fn two_stage_routine_input(
request : @routine.RoutineRequest,
gateway : @model.ModelGateway,
diagnosis : @domain.Stage1Diagnosis,
decision : @domain.Stage2Decision,
retry_policy? : @job.RetryPolicy = @job.default_retry_policy(),
stage1_attempt? : Int = 0,
stage2_attempt? : Int = 0,
cancel_after_phase? : @routine.RoutinePhase,
) -> TwoStageRoutineInput {
{
request,
gateway,
diagnosis,
decision,
retry_policy,
stage1_attempt,
stage2_attempt,
cancel_after_phase,
}
}
///|
fn routine_status_failed(message : String) -> @domain.RoutineStatus {
Failed(message)
}
///|
fn with_failure_status(
record : @domain.AnalysisRecord,
message : String,
) -> @domain.AnalysisRecord {
{ ..record, status: routine_status_failed(message) }
}
///|
fn event(
events : Array[@routine.RoutineEvent],
run_id : @domain.RunId,
phase : @routine.RoutinePhase,
kind : @routine.RoutineEventKind,
message : String,
) -> Unit {
events.push(
@routine.routine_event(run_id, events.length(), phase, kind, message),
)
}
///|
fn cancellation_requested(
input : TwoStageRoutineInput,
phase : @routine.RoutinePhase,
) -> Bool {
match input.cancel_after_phase {
Some(cancel_phase) => cancel_phase == phase
None => false
}
}
///|
fn snapshot_features_available(events : Array[@routine.RoutineEvent]) -> Bool {
events.any(fn(event) {
event.phase is SnapshotMarketData && event.kind is PhaseCompleted
})
}
///|
fn terminal_artifact_paths(
run_id : @domain.RunId,
events : Array[@routine.RoutineEvent],
stage1_prompt : @prompt.PromptPacket?,
stage2_prompt : @prompt.PromptPacket?,
context_budget : @model.ContextBudgetReport?,
stage1_model_request : @model.ModelCallRequest?,
stage2_model_request : @model.ModelCallRequest?,
) -> Array[String] {
let paths = [
"raw/market-data/\{run_id}.json",
"raw/run-events/\{run_id}.json",
"records/analyses/\{run_id}.json",
]
if snapshot_features_available(events) {
paths.push("records/structure-levels/\{run_id}.json")
}
match context_budget {
Some(_) => paths.push("records/context-budgets/\{run_id}.json")
None => ()
}
match stage1_prompt {
Some(_) => paths.push("raw/prompts/\{run_id}/stage1.json")
None => ()
}
match stage2_prompt {
Some(_) => paths.push("raw/prompts/\{run_id}/stage2.json")
None => ()
}
match stage1_model_request {
Some(_) => paths.push("raw/model-replies/\{run_id}/stage1-request.json")
None => ()
}
match stage2_model_request {
Some(_) => paths.push("raw/model-replies/\{run_id}/stage2-request.json")
None => ()
}
paths
}
///|
fn stage_gate(
phase : @routine.RoutinePhase,
report : @domain.ValidationReport,
retry_decision : @retry.RetryDecision,
) -> StageGate {
{ phase, validation: report, retry_decision }
}
///|
fn terminal_result(
input : TwoStageRoutineInput,
outcome : TwoStageOutcome,
failure_kind : StageFailureKind?,
record : @domain.AnalysisRecord,
events : Array[@routine.RoutineEvent],
stage_gates : Array[StageGate],
stage1_prompt : @prompt.PromptPacket?,
stage2_prompt : @prompt.PromptPacket?,
stage1_model_request : @model.ModelCallRequest?,
stage2_model_request : @model.ModelCallRequest?,
summary : String,
context_budget? : @model.ContextBudgetReport,
) -> TwoStageRoutineResult {
let paths = terminal_artifact_paths(
input.request.run_id,
events,
stage1_prompt,
stage2_prompt,
context_budget,
stage1_model_request,
stage2_model_request,
)
{
request: input.request,
outcome,
failure_kind,
record,
events,
stage_gates,
stage1_prompt,
stage2_prompt,
context_budget,
stage1_model_request,
stage2_model_request,
artifact_paths: paths,
partial_record_available: paths.length() > 0,
summary,
}
}
///|
fn cancelled_result(
input : TwoStageRoutineInput,
phase : @routine.RoutinePhase,
record : @domain.AnalysisRecord,
events : Array[@routine.RoutineEvent],
stage_gates : Array[StageGate],
stage1_prompt : @prompt.PromptPacket?,
stage2_prompt : @prompt.PromptPacket?,
stage1_model_request : @model.ModelCallRequest?,
stage2_model_request : @model.ModelCallRequest?,
context_budget? : @model.ContextBudgetReport,
) -> TwoStageRoutineResult {
event(
events,
input.request.run_id,
phase,
RoutineCancelled,
"operator cancelled after \{phase.label()}",
)
match context_budget {
Some(budget) =>
terminal_result(
input,
Cancelled,
Some(OperatorCancelled),
with_failure_status(record, "cancelled after \{phase.label()}"),
events,
stage_gates,
stage1_prompt,
stage2_prompt,
stage1_model_request,
stage2_model_request,
"two-stage routine cancelled after \{phase.label()}",
context_budget=budget,
)
None =>
terminal_result(
input,
Cancelled,
Some(OperatorCancelled),
with_failure_status(record, "cancelled after \{phase.label()}"),
events,
stage_gates,
stage1_prompt,
stage2_prompt,
stage1_model_request,
stage2_model_request,
"two-stage routine cancelled after \{phase.label()}",
)
}
}
///|
fn should_cancel_after(
input : TwoStageRoutineInput,
phase : @routine.RoutinePhase,
record : @domain.AnalysisRecord,
events : Array[@routine.RoutineEvent],
stage_gates : Array[StageGate],
stage1_prompt : @prompt.PromptPacket?,
stage2_prompt : @prompt.PromptPacket?,
stage1_model_request : @model.ModelCallRequest?,
stage2_model_request : @model.ModelCallRequest?,
context_budget? : @model.ContextBudgetReport,
) -> TwoStageRoutineResult? {
if cancellation_requested(input, phase) {
match context_budget {
Some(budget) =>
Some(
cancelled_result(
input,
phase,
record,
events,
stage_gates,
stage1_prompt,
stage2_prompt,
stage1_model_request,
stage2_model_request,
context_budget=budget,
),
)
None =>
Some(
cancelled_result(
input, phase, record, events, stage_gates, stage1_prompt, stage2_prompt,
stage1_model_request, stage2_model_request,
),
)
}
} else {
None
}
}
///|
pub fn run_two_stage_routine(
input : TwoStageRoutineInput,
) -> TwoStageRoutineResult {
let run_id = input.request.run_id
let events : Array[@routine.RoutineEvent] = []
let gates : Array[StageGate] = []
let mut record = @domain.empty_analysis_record(run_id, input.request.snapshot)
event(events, run_id, CreateRunContext, PhaseStarted, "created run context")
match
should_cancel_after(
input,
CreateRunContext,
record,
events,
gates,
None,
None,
None,
None,
) {
Some(result) => return result
None => ()
}
event(events, run_id, CreateRunContext, PhaseCompleted, "run context ready")
event(
events,
run_id,
SnapshotMarketData,
ToolPlanned,
"validating closed-bar market snapshot",
)
let snapshot_validation = @tools.validate_snapshot(input.request.snapshot)
if !snapshot_validation.ok {
record = record.with_validation(snapshot_validation)
gates.push(
stage_gate(
SnapshotMarketData,
snapshot_validation,
@retry.prepare_retry_decision(
@retry.retry_decision_input(
Stage1,
snapshot_validation,
policy=input.retry_policy,
),
),
),
)
event(
events,
run_id,
SnapshotMarketData,
RoutineFailed,
"snapshot validation failed",
)
return terminal_result(
input,
Failed,
Some(SnapshotInvalid),
with_failure_status(record, "snapshot validation failed"),
events,
gates,
None,
None,
None,
None,
"two-stage routine failed before model calls: snapshot invalid",
)
}
match
should_cancel_after(
input,
SnapshotMarketData,
record,
events,
gates,
None,
None,
None,
None,
) {
Some(result) => return result
None => ()
}
let features = @tools.compute_market_features(input.request.snapshot)
let hints = @tools.classify_market_structure(features)
event(
events,
run_id,
SnapshotMarketData,
PhaseCompleted,
"market snapshot and deterministic features ready",
)
if input.gateway.health is Unavailable {
event(
events,
run_id,
Stage1Diagnosis,
RoutineFailed,
"model gateway unavailable",
)
return terminal_result(
input,
Failed,
Some(GatewayUnavailable),
with_failure_status(record, "model gateway unavailable"),
events,
gates,
None,
None,
None,
None,
"two-stage routine failed before stage1: model gateway unavailable",
)
}
let stage1_prompt = @prompt.assemble_stage1_prompt(
@prompt.stage1_prompt_input(input.request, features, hints),
)
let stage1_context_budget = @model.prepare_context_budget_plan(
@model.context_budget_input(run_id, input.gateway, stage1_prompt),
)
let mut context_budget = @model.context_budget_report(run_id, [
stage1_context_budget,
])
let stage1_context_validation = @model.context_budget_validation(
stage1_context_budget,
)
if !stage1_context_validation.ok {
record = record.with_validation(stage1_context_validation)
gates.push(
stage_gate(
Stage1Diagnosis,
stage1_context_validation,
@retry.prepare_retry_decision(
@retry.retry_decision_input(
Stage1,
stage1_context_validation,
policy=input.retry_policy,
),
),
),
)
event(
events,
run_id,
Stage1Diagnosis,
RoutineFailed,
"stage1 prompt exceeds model context budget",
)
return terminal_result(
input,
Failed,
Some(ContextBudgetExceeded),
with_failure_status(record, "stage1 prompt exceeds model context budget"),
events,
gates,
Some(stage1_prompt),
None,
None,
None,
"two-stage routine failed before stage1 model call: context budget exceeded",
context_budget~,
)
}
let stage1_request = @model.model_call_request(
run_id,
input.gateway,
stage1_prompt,
)
event(
events,
run_id,
Stage1Diagnosis,
ModelOutputCaptured,
"captured stage1 output envelope",
)
let routed = @tools.diagnosis_with_routing(input.diagnosis)
record = record.with_stage1(routed)
let stage1_validation = @tools.validate_stage1(routed)
record = record.with_validation(stage1_validation)
let stage1_retry = @retry.prepare_retry_decision(
@retry.retry_decision_input(
Stage1,
stage1_validation,
attempt=input.stage1_attempt,
policy=input.retry_policy,
),
)
gates.push(stage_gate(Stage1Validation, stage1_validation, stage1_retry))
event(
events,
run_id,
Stage1Validation,
ValidationCaptured,
stage1_retry.message,
)
if !stage1_validation.ok {
if stage1_retry.should_retry {
return terminal_result(
input,
AwaitingStage1Retry,
Some(Stage1Invalid),
with_failure_status(record, stage1_retry.message),
events,
gates,
Some(stage1_prompt),
None,
Some(stage1_request),
None,
"two-stage routine paused for bounded stage1 retry",
context_budget~,
)
}
return terminal_result(
input,
Failed,
Some(Stage1Invalid),
with_failure_status(record, stage1_retry.message),
events,
gates,
Some(stage1_prompt),
None,
Some(stage1_request),
None,
"two-stage routine failed: stage1 validation rejected output",
context_budget~,
)
}
match
should_cancel_after(
input,
Stage1Validation,
record,
events,
gates,
Some(stage1_prompt),
None,
Some(stage1_request),
None,
context_budget~,
) {
Some(result) => return result
None => ()
}
event(
events,
run_id,
StrategyRouting,
ToolPlanned,
"routed \{routed.strategy_files_needed.length()} strategy file(s)",
)
let experience_memory = @memory.select_experience_memories(routed)
event(events, run_id, Stage2Decision, ToolPlanned, experience_memory.summary)
let stage2_prompt = @prompt.assemble_stage2_prompt(
@prompt.stage2_prompt_input(
input.request,
routed,
routed.strategy_files_needed,
experience_memories=experience_memory.entries,
),
)
let stage2_context_budget = @model.prepare_context_budget_plan(
@model.context_budget_input(run_id, input.gateway, stage2_prompt),
)
context_budget = @model.context_budget_report(run_id, [
stage1_context_budget, stage2_context_budget,
])
let stage2_context_validation = @model.context_budget_validation(
stage2_context_budget,
)
if !stage2_context_validation.ok {
record = record.with_validation(stage2_context_validation)
gates.push(
stage_gate(
Stage2Decision,
stage2_context_validation,
@retry.prepare_retry_decision(
@retry.retry_decision_input(
Stage2,
stage2_context_validation,
policy=input.retry_policy,
),
),
),
)
event(
events,
run_id,
Stage2Decision,
RoutineFailed,
"stage2 prompt exceeds model context budget",
)
return terminal_result(
input,
Failed,
Some(ContextBudgetExceeded),
with_failure_status(record, "stage2 prompt exceeds model context budget"),
events,
gates,
Some(stage1_prompt),
Some(stage2_prompt),
Some(stage1_request),
None,
"two-stage routine failed before stage2 model call: context budget exceeded",
context_budget~,
)
}
let stage2_request = @model.model_call_request(
run_id,
input.gateway,
stage2_prompt,
)
event(
events,
run_id,
Stage2Decision,
ModelOutputCaptured,
"captured stage2 output envelope",
)
record = record.with_stage2(input.decision)
let stage2_shape_validation = @tools.validate_stage2(input.decision)
let decision_policy = @policy.evaluate_decision_policy(
@policy.decision_policy_input(run_id, routed, input.decision),
)
let stage2_validation = if stage2_shape_validation.ok {
@policy.decision_policy_validation(decision_policy)
} else {
stage2_shape_validation
}
record = record.with_validation(stage2_validation)
let stage2_retry = @retry.prepare_retry_decision(
@retry.retry_decision_input(
Stage2,
stage2_validation,
attempt=input.stage2_attempt,
policy=input.retry_policy,
),
)
gates.push(stage_gate(Stage2Validation, stage2_validation, stage2_retry))
event(
events,
run_id,
Stage2Validation,
ValidationCaptured,
stage2_retry.message,
)
if !stage2_validation.ok {
if stage2_retry.should_retry {
return terminal_result(
input,
AwaitingStage2Retry,
Some(Stage2Invalid),
with_failure_status(record, stage2_retry.message),
events,
gates,
Some(stage1_prompt),
Some(stage2_prompt),
Some(stage1_request),
Some(stage2_request),
"two-stage routine paused for bounded stage2 retry",
context_budget~,
)
}
return terminal_result(
input,
Failed,
Some(Stage2Invalid),
with_failure_status(record, stage2_retry.message),
events,
gates,
Some(stage1_prompt),
Some(stage2_prompt),
Some(stage1_request),
Some(stage2_request),
"two-stage routine failed: stage2 validation rejected output",
context_budget~,
)
}
match
should_cancel_after(
input,
Stage2Validation,
record,
events,
gates,
Some(stage1_prompt),
Some(stage2_prompt),
Some(stage1_request),
Some(stage2_request),
context_budget~,
) {
Some(result) => return result
None => ()
}
record = record.persisted()
event(
events,
run_id,
PersistAnalysis,
RecordPersisted,
"analysis record and evidence paths prepared",
)
event(
events,
run_id,
FollowUpReady,
PhaseCompleted,
"follow-up context ready",
)
terminal_result(
input,
Completed,
None,
record,
events,
gates,
Some(stage1_prompt),
Some(stage2_prompt),
Some(stage1_request),
Some(stage2_request),
"two-stage routine completed with \{events.length()} event(s) and \{gates.length()} validation gate(s)",
context_budget~,
)
}