Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
154 changes: 98 additions & 56 deletions apps/signalboxd/src/context_guard.rs
Original file line number Diff line number Diff line change
Expand Up @@ -128,6 +128,11 @@ pub struct ReportedUsageCompaction {
compaction_model: Arc<dyn ContextCompactionModel>,
}

struct ReportedUsageCompactionCandidate {
preview: PreparedActivationPreview,
turn: TurnId,
}

impl fmt::Debug for ReportedUsageCompaction {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
Expand Down Expand Up @@ -167,63 +172,10 @@ impl ReportedUsageCompaction {
&self,
session: SessionId,
) -> Result<(), ReportedUsageCompactionError> {
let Some(preview) = self
.activation
.preview(session, activation_identities())
.await
.map_err(ReportedUsageCompactionError::Activation)?
else {
let Some(candidate) = self.compaction_candidate(session).await? else {
return Ok(());
};
let turn = preview.prepared().turn().turn();
let prospective = self
.model_calls
.preview_activation_operation(
preview.prepared(),
ModelCallId::from_uuid(uuid::Uuid::now_v7()),
)
.await
.map_err(|source| ReportedUsageCompactionError::Model { turn, source })?;
let Some(prospective) = prospective else {
return Ok(());
};
let operation = prospective
.render(self.tools.definitions())
.map_err(|_| ReportedUsageCompactionError::Render(turn))?;
let target = operation.request().call().target();
let selected = self
.runtime_models
.resolve(target)
.ok_or(ReportedUsageCompactionError::ContextWindowUnavailable(turn))?;
let definition = self
.runtime_models
.effective_definition(
selected,
operation.request().model_settings().effective().fast_mode(),
)
.ok_or(ReportedUsageCompactionError::ContextWindowUnavailable(turn))?;
let Some(reported) = self
.model_calls
.latest_reported_usage(
session,
target,
operation.request().call().frontier().snapshot(),
)
.await
.map_err(|source| ReportedUsageCompactionError::Model { turn, source })?
else {
return Ok(());
};
if !reported_usage_requires_compaction(
reported.usage(),
reported.input_includes_cache_tokens(),
reported.output_is_retained(),
reported.projected_unreported_content_bytes(),
u64::from(definition.max_output_tokens()),
u64::from(definition.context_window_tokens()),
) {
return Ok(());
}
let ReportedUsageCompactionCandidate { preview, turn } = candidate;
let applied = match compact_automatically(
&self.model_calls,
&self.model_configuration,
Expand Down Expand Up @@ -290,7 +242,97 @@ impl ReportedUsageCompaction {
context_compaction_id = %applied.compaction.into_uuid(),
"provider-reported usage exhausted reserved context headroom; queued turn compacted before activation"
);
Ok(())
let Some(remaining) = self.compaction_candidate(session).await? else {
return Ok(());
};
let remaining_turn = remaining.turn;
match close_failed_compaction_turn(&self.activation, &self.model_calls, remaining.preview)
.await
.map_err(
|source| ReportedUsageCompactionError::CompactionFailureClosure {
turn: remaining_turn,
source,
},
)? {
CommitCompactionFailurePreviewOutcome::Failed(_) => {
tracing::warn!(
cause_code = "reported_usage_context_still_exceeded",
session_id = %session.as_uuid(),
turn_id = %remaining_turn.as_uuid(),
"automatic compaction did not restore reserved context headroom; the queued turn was closed before provider dispatch"
);
Err(ReportedUsageCompactionError::Compaction {
turn: remaining_turn,
failure_class: OperatorFailureClass::CallerOrHubBug,
cause_code: "reported_usage_context_still_exceeded",
})
}
CommitCompactionFailurePreviewOutcome::Stale => Ok(()),
}
}

async fn compaction_candidate(
&self,
session: SessionId,
) -> Result<Option<ReportedUsageCompactionCandidate>, ReportedUsageCompactionError> {
let Some(preview) = self
.activation
.preview(session, activation_identities())
.await
.map_err(ReportedUsageCompactionError::Activation)?
else {
return Ok(None);
};
let turn = preview.prepared().turn().turn();
let prospective = self
.model_calls
.preview_activation_operation(
preview.prepared(),
ModelCallId::from_uuid(uuid::Uuid::now_v7()),
)
.await
.map_err(|source| ReportedUsageCompactionError::Model { turn, source })?;
let Some(prospective) = prospective else {
return Ok(None);
};
let operation = prospective
.render(self.tools.definitions())
.map_err(|_| ReportedUsageCompactionError::Render(turn))?;
let target = operation.request().call().target();
let selected = self
.runtime_models
.resolve(target)
.ok_or(ReportedUsageCompactionError::ContextWindowUnavailable(turn))?;
let definition = self
.runtime_models
.effective_definition(
selected,
operation.request().model_settings().effective().fast_mode(),
)
.ok_or(ReportedUsageCompactionError::ContextWindowUnavailable(turn))?;
let Some(reported) = self
.model_calls
.latest_reported_usage(
session,
target,
operation.request().call().frontier().snapshot(),
)
.await
.map_err(|source| ReportedUsageCompactionError::Model { turn, source })?
else {
return Ok(None);
};
if !reported_usage_requires_compaction(
reported.usage(),
reported.input_includes_cache_tokens(),
reported.output_is_retained(),
reported.projected_unreported_content_bytes(),
u64::from(definition.max_output_tokens()),
u64::from(definition.context_window_tokens()),
) {
return Ok(None);
}
Ok(Some(ReportedUsageCompactionCandidate { preview, turn }))
}
}

Expand Down
66 changes: 56 additions & 10 deletions apps/signalboxd/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -989,11 +989,13 @@ where
source,
});
}
if let Some(compaction) = reported_usage_compaction {
compaction
.compact_if_needed(session)
.await
.map_err(ActivatedTurnPassError::ReportedUsageCompaction)?;
if let Some(compaction) = reported_usage_compaction
&& let Err(error) = compaction.compact_if_needed(session).await
{
if let Some(recovery) = &occupancy_recovery {
recovery.clear_turn(session);
}
return Err(reported_usage_compaction_failure(&execution, error));
}
let outcome = match activation.await {
Ok(outcome) => outcome,
Expand Down Expand Up @@ -1276,6 +1278,17 @@ where
}
}

fn reported_usage_compaction_failure<Execution, ActivationError>(
execution: &Execution,
error: ReportedUsageCompactionError,
) -> ActivatedTurnPassError<ActivationError, Execution::Error>
where
Execution: ActivatedTurnExecution,
{
report_ambiguous_commit(execution, &error);
ActivatedTurnPassError::ReportedUsageCompaction(error)
}

/// Whether one classified failure left a durable commit outcome the running
/// process cannot decide.
///
Expand Down Expand Up @@ -2467,7 +2480,10 @@ mod tests {
AcceptedInputTurnActivationIdentities, ActivatedTurn, ContextFrontierId,
SemanticTranscriptEntryId, SessionId, TurnAttemptId, TurnId,
};
use signalbox_persistence::turn_liveness::TurnLivenessPersistenceBounds;
use signalbox_persistence::{
start_eligible_turn::{CommitActivationPreviewError, StartEligibleTurnRepositoryError},
turn_liveness::TurnLivenessPersistenceBounds,
};
use tokio::sync::watch;
use uuid::Uuid;

Expand All @@ -2476,12 +2492,12 @@ mod tests {
ActivatedTurnPassError, ApprovalJudgeModelError, ExpiredPassRecoveryPolicy,
FailedApprovalJudgeDisposition, FatalExecutionGuardState, FatalExecutionOccupancyExpiry,
FatalExecutionSignal, FatalExecutionSupervisor, JudgeRequestFields,
MAX_QUOTED_CONTEXT_BYTES, SchedulerPassOccupancyRecovery, SessionAuthorityContext,
TokenUsage, TurnLivenessRepositoryError, TurnPassExecutionStage,
MAX_QUOTED_CONTEXT_BYTES, ReportedUsageCompactionError, SchedulerPassOccupancyRecovery,
SessionAuthorityContext, TokenUsage, TurnLivenessRepositoryError, TurnPassExecutionStage,
activation_session_matches, expired_pass_recovery_retry_delay,
matches_exact_slot_held_turn, reconcile_retained_once, render_dispatch_authority,
render_judge_request_payload, render_session_authority_context, supervise_execution,
supervise_execution_for_session,
render_judge_request_payload, render_session_authority_context,
reported_usage_compaction_failure, supervise_execution, supervise_execution_for_session,
};

fn example_expired_pass_policy() -> ExpiredPassRecoveryPolicy {
Expand Down Expand Up @@ -2586,6 +2602,16 @@ mod tests {
}
}

#[track_caller]
fn assert_reported_usage_compaction_error(
error: ActivatedTurnPassError<ExecutionFailure, ExecutionFailure>,
) {
match error {
ActivatedTurnPassError::ReportedUsageCompaction(_) => {}
other => panic!("expected reported-usage compaction failure, got {other:?}"),
}
}

#[derive(Clone, Copy, Debug, Eq, PartialEq)]
struct CommitAmbiguousActivationFailure;

Expand Down Expand Up @@ -3044,6 +3070,26 @@ mod tests {
assert!(signal.is_triggered());
}

#[test]
fn inv034_ambiguous_reported_usage_failure_closure_raises_the_fatal_recovery_signal() {
let (execution, signal) = FatalExecutionSupervisor::new(NoopExecution);
let source =
CommitActivationPreviewError::Activation(StartEligibleTurnRepositoryError::Database {
source: sqlx::Error::PoolClosed,
commit_ambiguous: true,
});
let error = ReportedUsageCompactionError::CompactionFailureClosure {
turn: TurnId::from_uuid(Uuid::from_u128(11)),
source,
};

let reported: ActivatedTurnPassError<ExecutionFailure, ExecutionFailure> =
reported_usage_compaction_failure(&execution, error);

assert_reported_usage_compaction_error(reported);
assert!(signal.is_triggered());
}

#[test]
fn activation_session_mismatch_raises_the_fatal_signal() {
let (execution, signal) = FatalExecutionSupervisor::new(NoopExecution);
Expand Down
Loading
Loading