File size: 5,003 Bytes
52a9af3 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 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 102 103 104 105 106 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 135 136 137 138 139 140 141 142 143 144 145 146 147 148 | //! Handles persistent thread-settings updates and serializes their persistence
//! with checkpoints written directly to storage.
use super::session::Session;
use super::session::SessionSettingsUpdate;
use super::step_settings::StepSettingsUpdate;
use crate::config::ConstraintResult;
use codex_history::RolloutItem;
use codex_protocol::protocol::CodexErrorInfo;
use codex_protocol::protocol::ErrorEvent;
use codex_protocol::protocol::Event;
use codex_protocol::protocol::EventMsg;
use codex_protocol::protocol::ThreadSettingsAppliedEvent;
use codex_protocol::protocol::ThreadSettingsOverrides;
use codex_protocol::protocol::ThreadSettingsSnapshot;
use codex_thread_store::ThreadStoreResult;
use std::sync::Arc;
use tokio::sync::SemaphorePermit;
impl Session {
/// Captures and flushes current settings under the shared persistence permit.
pub(crate) async fn checkpoint_thread_settings(&self) -> ThreadStoreResult<()> {
let _settings_guard = acquire_persistence_lock(self).await;
if let Some(live_thread) = self.live_thread() {
live_thread
.append_items(&[RolloutItem::EventMsg(applied_event(self).await)])
.await?;
live_thread.flush().await?;
}
Ok(())
}
}
/// Applies standalone thread settings and reports invalid overrides through the
/// normal event stream.
pub(super) async fn update(
session: &Arc<Session>,
submission_id: String,
overrides: ThreadSettingsOverrides,
) {
let updates = prepare_update(overrides);
if let Err(error) = apply_update(session, submission_id.clone(), updates).await {
session
.send_event_raw(Event {
id: submission_id,
msg: EventMsg::Error(ErrorEvent {
misalignment: None,
message: format!("invalid thread settings override: {error}"),
codex_error_info: Some(CodexErrorInfo::BadRequest),
}),
})
.await;
} else {
// Standalone settings changes supersede a pending automatic continuation.
session.state.lock().await.last_started_turn_id = None;
}
}
/// Converts protocol overrides into the internal settings update shape.
pub(super) fn prepare_update(overrides: ThreadSettingsOverrides) -> SessionSettingsUpdate {
let ThreadSettingsOverrides {
environments,
runtime_workspace_roots,
profile_workspace_roots,
approval_policy,
approvals_reviewer,
sandbox_policy,
permission_profile,
active_permission_profile,
windows_sandbox_level,
model,
effort,
summary,
service_tier,
collaboration_mode,
personality,
disabled_plugin_ids,
} = overrides;
SessionSettingsUpdate {
step_settings: StepSettingsUpdate {
model,
effort,
collaboration_mode,
reasoning_summary: summary,
service_tier,
personality,
approval_policy,
approvals_reviewer,
},
environments,
runtime_workspace_roots,
profile_workspace_roots,
sandbox_policy,
permission_profile,
active_permission_profile,
windows_sandbox_level,
disabled_plugin_ids,
..Default::default()
}
}
/// Acquires the shared permit before capturing or changing persistent settings.
pub(super) async fn acquire_persistence_lock(session: &Session) -> SemaphorePermit<'_> {
session
.thread_settings_persistence
.acquire()
.await
.unwrap_or_else(|_| unreachable!("thread settings persistence semaphore is never closed"))
}
/// Applies persistent settings and emits the resulting thread-owned snapshot.
pub(super) async fn apply_update(
session: &Session,
submission_id: String,
updates: SessionSettingsUpdate,
) -> ConstraintResult<()> {
let _settings_guard = acquire_persistence_lock(session).await;
let commit = session.update_settings(updates).await?;
emit_applied(session, submission_id, commit.snapshot).await;
Ok(())
}
/// Emits the snapshot published by one successful settings update.
pub(super) async fn emit_applied(
session: &Session,
submission_id: String,
snapshot: ThreadSettingsSnapshot,
) {
let msg = EventMsg::ThreadSettingsApplied(ThreadSettingsAppliedEvent {
thread_id: Some(session.thread_id()),
thread_settings: snapshot,
});
session
.send_event_raw_without_materializing_rollout(Event {
id: submission_id,
msg,
})
.await;
}
/// Builds a current thread-owned snapshot for storage checkpoints.
pub(super) async fn applied_event(session: &Session) -> EventMsg {
EventMsg::ThreadSettingsApplied(ThreadSettingsAppliedEvent {
thread_id: Some(session.thread_id()),
thread_settings: session.thread_settings_snapshot().await,
})
}
|