use super::*;
pub(super) type SharedRx = Arc<Mutex<mpsc::Receiver<AgentEvent>>>;
pub(super) type SharedManifestRx =
Arc<Mutex<tokio::sync::broadcast::Receiver<LocalWorkspaceManifestSnapshot>>>;
pub(super) type SharedActiveSession = Arc<std::sync::Mutex<Arc<AgentSession>>>;
pub(super) type StreamJoin = tokio::task::JoinHandle<()>;
pub(super) type HostToolAbort = tokio::task::AbortHandle;
#[derive(Clone, Copy, PartialEq)]
pub(super) enum State {
Idle,
Streaming,
Awaiting,
Rebuilding,
}
#[derive(Clone, Copy, Debug)]
pub(super) enum ViewportAnchor {
Bottom,
Transcript(TranscriptAnchor),
Absolute(usize),
}
#[derive(Clone)]
#[allow(clippy::enum_variant_names)]
pub(super) enum Action {
ScrollUp,
ScrollDown,
ScrollTop,
ScrollBottom,
}
pub(super) static UPGRADE_ON_EXIT: std::sync::atomic::AtomicBool =
std::sync::atomic::AtomicBool::new(false);
pub(super) static LATEST: std::sync::Mutex<Option<String>> = std::sync::Mutex::new(None);
#[derive(Clone, Debug, PartialEq, Eq)]
pub(super) struct AutoReviewKey {
pub(super) session_id: String,
pub(super) revision: u64,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub(super) struct AutoReviewTicket {
pub(super) id: u64,
pub(super) key: AutoReviewKey,
}
#[derive(Debug)]
pub(super) struct AutoReviewTracker {
pub(super) revision: u64,
pub(super) reviewed: Option<AutoReviewKey>,
pub(super) inflight: Option<AutoReviewTicket>,
pub(super) next_ticket_id: u64,
}
impl AutoReviewTracker {
pub(super) fn new(revision: u64) -> Self {
Self {
revision,
reviewed: None,
inflight: None,
next_ticket_id: 0,
}
}
pub(super) fn on_user_turn(&mut self) {
self.revision = self.revision.wrapping_add(1);
}
pub(super) fn current_key(&self, session_id: &str) -> AutoReviewKey {
AutoReviewKey {
session_id: session_id.to_string(),
revision: self.revision,
}
}
pub(super) fn current_is_reviewed(&self, session_id: &str) -> bool {
self.reviewed
.as_ref()
.is_some_and(|key| key.session_id == session_id && key.revision == self.revision)
}
pub(super) fn begin(
&mut self,
session_id: &str,
has_user_turn: bool,
) -> Option<AutoReviewTicket> {
let key = self.current_key(session_id);
if self.reviewed.as_ref() == Some(&key) {
return None;
}
self.reviewed = Some(key.clone());
if !has_user_turn {
return None;
}
self.next_ticket_id = self.next_ticket_id.wrapping_add(1);
let ticket = AutoReviewTicket {
id: self.next_ticket_id,
key,
};
self.inflight = Some(ticket.clone());
Some(ticket)
}
pub(super) fn accept(&mut self, ticket: &AutoReviewTicket, session_id: &str) -> bool {
if self.inflight.as_ref() != Some(ticket) {
return false;
}
self.inflight = None;
ticket.key.session_id == session_id
&& ticket.key.revision == self.revision
&& self.reviewed.as_ref() == Some(&ticket.key)
}
}
pub(super) fn auto_review_history_has_user_turn(history: &[Message]) -> bool {
history
.iter()
.any(|message| message.role == "user" && !message.text().trim().is_empty())
}
pub(super) enum SessionRebuildAction {
Model {
model: String,
source: ModelSelectionSource,
llm_override: Option<LlmOverride>,
context_limit: u32,
},
Effort {
selected: usize,
codex_effort: Option<CodexEffortStatus>,
},
GoalStart {
generation: u64,
previous_effort: usize,
previous_goal: Option<String>,
previous_goal_since: Option<Instant>,
},
GoalRestore,
Compact {
summary: String,
session_id: String,
},
Fork {
session_id: String,
},
Clear {
session_id: String,
},
Reload {
skill_count: usize,
},
Refresh {
failure_context: Option<&'static str>,
},
}
pub(super) struct SessionRebuildProfile {
pub(super) session_id: String,
pub(super) model: Option<String>,
pub(super) effort: usize,
pub(super) context_limit: u32,
pub(super) llm_override: Option<LlmOverride>,
pub(super) compact_summary: Option<String>,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(super) enum SessionRebuildMode {
ResumeExisting,
CreateFresh,
}
pub(super) enum Msg {
Term(Event),
Agent {
source: SharedRx,
event: Box<AgentEvent>,
},
Submit(String),
StreamStarted {
token: u64,
session: Arc<AgentSession>,
rx: SharedRx,
join: StreamJoin,
},
StreamEnded(SharedRx),
StreamJoinSettled {
token: u64,
synthesis: Option<(String, String)>,
},
DiscardedStreamSettled,
QuitReady,
StreamError {
token: u64,
error: String,
},
WorkspaceManifest(Box<LocalWorkspaceManifestSnapshot>),
WorkspaceManifestStopped,
SpinnerTick,
StreamCommitTick,
BannerTick,
UltracodeTick {
epoch: u64,
},
ModalConfirm {
tool_id: String,
approved: bool,
approve_all_pending: bool,
},
BackgroundSubagentFinished {
session_id: String,
generation: u64,
task_id: String,
agent: String,
output: String,
outcome: SubagentOutcome,
finished_ms: u64,
},
BackgroundSubagentWatchStopped {
session_id: String,
generation: u64,
task_id: String,
},
SubagentSnapshots {
session_id: String,
generation: u64,
request_id: u64,
snapshots: Vec<RestoredSubagentSnapshot>,
},
DeepResearchSubagentsSettled {
session_id: String,
generation: u64,
exit: DeepResearchSettlementExit,
settlements: Vec<DeepResearchSubagentSettlement>,
},
DeepResearchJournalFinalized {
run_id: String,
exit: DeepResearchSettlementExit,
result: Result<ResearchRunProjection, String>,
},
DeepResearchJournalEventRecorded {
run_id: String,
result: Result<ResearchRunProjection, String>,
},
Resume,
Interrupted {
goal_cancelled: bool,
status_entry: TranscriptEntryId,
},
ShellOutput(String),
ResearchDiagnostic(Result<String, String>),
DeepResearchWorkflowCompleted {
query: String,
os_runtime: bool,
args: serde_json::Value,
result: Result<ToolCallResult, String>,
convergence: ConvergenceDecision,
accepted_evidence: Vec<AcceptedEvidence>,
},
DeepResearchSynthesisTimedOut {
token: u64,
},
DeepResearchSynthesisTimedOutAfterCancel {
token: u64,
status: String,
streamed_text: String,
report_completed: bool,
},
UpdatePlan(Option<String>),
UpdateRepair {
status_entry: TranscriptEntryId,
result: Result<Vec<String>, String>,
},
OsLogin {
status_entry: TranscriptEntryId,
result: Result<String, String>,
},
SshKeySynced(crate::a3s_os::SshKeyOutcome),
OsRefreshed(Result<crate::a3s_os::StoredOsSession, String>),
OsGatewayModels {
login_at_ms: u64,
result: Result<Vec<crate::a3s_os::GatewayModel>, String>,
},
AccountModels {
provider: crate::account_providers::AccountProvider,
result: Result<Vec<String>, String>,
},
GoalContinue {
generation: u64,
prompt: String,
},
GoalCleared,
CodexModels(Result<Vec<crate::account_providers::codex::CodexModel>, String>),
SessionRebuilt {
request_id: u64,
action: SessionRebuildAction,
result: Box<panels::model::SessionRebuildResult>,
},
Forked {
request_id: u64,
result: Result<String, String>,
},
MemoryLoaded(MemPanelData),
MemoryForgotten(Result<(String, MemPanelData), String>),
AssetListLoaded(Result<panels::asset_resources::AssetListFetch, String>),
RuntimeActivityLoaded(Result<panels::asset_resources::RuntimeActivityFetch, String>),
KbAdded(String),
CtxResults {
status_entry: TranscriptEntryId,
result: Result<String, String>,
},
CtxWindow {
status_entry: TranscriptEntryId,
result: Result<(String, String), String>,
},
CtxSaved(Result<String, String>),
SleepSaved(Result<usize, String>),
FlowOsCompleted {
status_entry: TranscriptEntryId,
result: Result<panels::flow::FlowOsResult, String>,
},
AgentOsCompleted {
status_entry: TranscriptEntryId,
result: Result<panels::agent::AgentOsResult, String>,
},
McpOsCompleted {
status_entry: TranscriptEntryId,
result: Result<panels::mcp::McpOsResult, String>,
},
SkillOsCompleted {
status_entry: TranscriptEntryId,
result: Result<panels::skill::SkillOsResult, String>,
},
OkfOsCompleted {
status_entry: TranscriptEntryId,
result: Result<panels::okf::OkfOsResult, String>,
},
AssetCloned {
status_entry: TranscriptEntryId,
result: Result<asset_clone::AssetCloneResult, String>,
},
CtxMemorySource(Result<(String, String), String>),
AutoReview {
ticket: AutoReviewTicket,
text: String,
},
Compacted(Result<Option<String>, String>),
UpdateCheck(Option<String>),
}
pub(super) struct RestoredSubagentSnapshot {
pub(super) snapshot: a3s_code_core::SubagentTaskSnapshot,
pub(super) parent_result_expected: bool,
}
pub(super) struct DeepResearchSubagentSettlement {
pub(super) task_id: String,
pub(super) agent: String,
pub(super) output: String,
pub(super) outcome: SubagentOutcome,
pub(super) finished_ms: u64,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(super) enum DeepResearchSettlementExit {
ReportReady,
Interrupted,
}
impl DeepResearchSettlementExit {
pub(super) fn opens_report(self) -> bool {
matches!(self, Self::ReportReady)
}
}
impl From<Event> for Msg {
fn from(event: Event) -> Self {
Msg::Term(event)
}
}