use std::sync::Arc;
use aion_core::{
AssistantSessionEvent, AssistantSessionId, AssistantSessionState, AssistantTurnContext,
};
use aion_store::assistant::{
AssistantSessionListing, AssistantSessionRecord, AssistantSessionStore,
AssistantTranscriptEvent,
};
use aion_store::{InMemoryStore, StoreError};
use async_trait::async_trait;
use chrono::{DateTime, Utc};
use crate::config::{
AssistantAccountConfig, AssistantConfig, AssistantHarnessConfig, ResolvedAssistantConfig,
};
use aion_integration_acp::catalogue::CatalogueHarness;
use super::launch::AssistantEndpoints;
use super::registry::AssistantSessions;
pub(crate) const OPERATOR: &str = "operator";
pub(crate) const HARNESS: &str = "claude-code";
pub(crate) const ACCOUNT: &str = "work";
pub(crate) const ACCOUNT_SOURCE: &str = "AION_FIXTURE_CLAUDE_DIR";
pub(crate) fn config() -> ResolvedAssistantConfig {
AssistantConfig {
harnesses: vec![AssistantHarnessConfig {
name: Some(HARNESS.to_owned()),
accounts: vec![AssistantAccountConfig {
name: Some(ACCOUNT.to_owned()),
env: [("CLAUDE_CONFIG_DIR".to_owned(), ACCOUNT_SOURCE.to_owned())]
.into_iter()
.collect(),
}],
}],
}
.resolved()
}
pub(crate) fn endpoints() -> AssistantEndpoints {
AssistantEndpoints {
base: "http://127.0.0.1:9999".to_owned(),
aion_mcp_enabled: false,
}
}
pub(crate) fn registry() -> (AssistantSessions, Arc<InMemoryStore>) {
let store = Arc::new(InMemoryStore::default());
let sessions = AssistantSessions::new(
Arc::clone(&store) as Arc<dyn AssistantSessionStore>,
config(),
Some(endpoints()),
CATALOGUE,
);
(sessions, store)
}
pub(crate) const CATALOGUE: &[CatalogueHarness] = &[CatalogueHarness {
id: HARNESS,
display_name: "Claude Code (fixture)",
program: "python3",
args: &[STUB_AGENT],
install_hint: "the fixture harness needs `python3` on PATH",
}];
const STUB_AGENT: &str = concat!(
env!("CARGO_MANIFEST_DIR"),
"/src/assistant/sessions/fixtures/stub_agent.py"
);
pub(crate) fn registry_over(store: Arc<dyn AssistantSessionStore>) -> AssistantSessions {
AssistantSessions::new(store, config(), Some(endpoints()), CATALOGUE)
}
pub(crate) async fn open_session(
sessions: &AssistantSessions,
) -> Result<aion_core::AssistantSessionSummary, String> {
let session_id = AssistantSessionId::new_v4();
let now = Utc::now();
let record = AssistantSessionRecord {
session_id,
subject: OPERATOR.to_owned(),
harness: HARNESS.to_owned(),
account: None,
title: None,
created_at: now,
updated_at: now,
turns: 0,
mcp_token_digest: None,
commands: Vec::new(),
config_options: Vec::new(),
};
sessions
.store()
.put_assistant_session(record.clone())
.await
.map_err(|error| error.to_string())?;
sessions
.settle(
session_id,
AssistantSessionState::Dormant,
super::lifecycle::CREATED_AWAITING_FIRST_TURN,
)
.await
.map_err(|error| error.to_string())?;
Ok(record.summary(
AssistantSessionState::Dormant,
Some(super::lifecycle::CREATED_AWAITING_FIRST_TURN.to_owned()),
))
}
pub(crate) async fn opened(
sessions: &AssistantSessions,
session_id: AssistantSessionId,
load_session: bool,
) -> Result<(), String> {
sessions
.recorder(session_id)
.record(AssistantSessionEvent::SessionOpened {
acp_session_ref: format!("acp-{session_id}"),
load_session,
at: Utc::now(),
resumed: false,
})
.await
.map(drop)
.map_err(|error| error.to_string())
}
pub(crate) async fn went_dormant(
sessions: &AssistantSessions,
session_id: AssistantSessionId,
) -> Result<(), String> {
sessions
.settle(
session_id,
AssistantSessionState::Dormant,
super::lifecycle::PROCESS_EXITED_RESUMABLE,
)
.await
.map_err(|error| error.to_string())
}
pub(crate) async fn shared_context(
sessions: &AssistantSessions,
session_id: AssistantSessionId,
url: &str,
) -> Result<(), String> {
sessions
.push_context(
OPERATOR,
session_id,
AssistantTurnContext {
url: Some(url.to_owned()),
concepts: Vec::new(),
document: None,
},
)
.await
.map_err(|error| error.to_string())
}
pub(crate) struct RefusingTranscriptStore {
inner: InMemoryStore,
}
impl RefusingTranscriptStore {
pub(crate) const REFUSAL: &'static str =
"this transcript store refuses every append (test instrument)";
pub(crate) fn new() -> Self {
Self {
inner: InMemoryStore::default(),
}
}
}
#[async_trait]
impl AssistantSessionStore for RefusingTranscriptStore {
async fn put_assistant_session(
&self,
record: AssistantSessionRecord,
) -> Result<(), StoreError> {
self.inner.put_assistant_session(record).await
}
async fn get_assistant_session(
&self,
session_id: &AssistantSessionId,
) -> Result<Option<AssistantSessionRecord>, StoreError> {
self.inner.get_assistant_session(session_id).await
}
async fn list_assistant_sessions(&self) -> Result<AssistantSessionListing, StoreError> {
self.inner.list_assistant_sessions().await
}
async fn append_assistant_transcript_event(
&self,
_session_id: &AssistantSessionId,
_recorded_at: DateTime<Utc>,
_payload: aion_core::Payload,
) -> Result<u64, StoreError> {
Err(StoreError::Backend(Self::REFUSAL.to_owned()))
}
async fn assistant_transcript_head(
&self,
session_id: &AssistantSessionId,
) -> Result<u64, StoreError> {
self.inner.assistant_transcript_head(session_id).await
}
async fn assistant_transcript(
&self,
session_id: &AssistantSessionId,
after: Option<u64>,
) -> Result<Vec<AssistantTranscriptEvent>, StoreError> {
self.inner.assistant_transcript(session_id, after).await
}
async fn put_assistant_default_harness(
&self,
subject: &str,
harness: &str,
) -> Result<(), StoreError> {
self.inner
.put_assistant_default_harness(subject, harness)
.await
}
async fn assistant_default_harness(&self, subject: &str) -> Result<Option<String>, StoreError> {
self.inner.assistant_default_harness(subject).await
}
}