use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::Arc;
use serde::{Deserialize, Serialize};
use crate::agentic::memory_coordinator::MemoryCoordinator;
use crate::agentic::token_ledger::TokenLedger;
use crate::types::{MemoryFact, ModelId, SessionId, SessionRecord};
use foundation_compact::SystemTime;
use foundation_db::traits::DocumentStore;
use foundation_db::StorageResult;
use super::memory_store::MemoryStore;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MemoryAction {
None,
GenerateObservation,
GenerateReflection,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
pub enum MemoryParseStrategy {
#[default]
StructuredText,
TwoStepJson,
}
#[derive(Debug, Clone)]
pub struct MemoryConfig {
pub observation_trigger_tokens: u64,
pub reflection_trigger_tokens: u64,
pub memory_model: Option<ModelId>,
pub parse_strategy: MemoryParseStrategy,
}
impl Default for MemoryConfig {
fn default() -> Self {
Self {
observation_trigger_tokens: 30_000,
reflection_trigger_tokens: 40_000,
memory_model: None,
parse_strategy: MemoryParseStrategy::default(),
}
}
}
pub struct MemoryHierarchy<M, D> {
inner: Arc<MemoryInner<M, D>>,
}
impl<M, D> Clone for MemoryHierarchy<M, D> {
fn clone(&self) -> Self {
Self {
inner: self.inner.clone(),
}
}
}
struct MemoryInner<M, D> {
session_id: SessionId,
coordinator: MemoryCoordinator<M, D>,
ledger: TokenLedger,
config: MemoryConfig,
observation_tokens: AtomicU64,
is_generating: AtomicBool,
}
impl<M: MemoryStore, D: DocumentStore> MemoryHierarchy<M, D> {
#[must_use]
pub fn new(
session_id: SessionId,
coordinator: MemoryCoordinator<M, D>,
ledger: TokenLedger,
config: MemoryConfig,
) -> Self {
Self {
inner: Arc::new(MemoryInner {
session_id,
coordinator,
ledger,
config,
observation_tokens: AtomicU64::new(0),
is_generating: AtomicBool::new(false),
}),
}
}
#[must_use]
pub fn coordinator(&self) -> &MemoryCoordinator<M, D> {
&self.inner.coordinator
}
#[must_use]
pub fn check_triggers(&self) -> MemoryAction {
if self.inner.is_generating.load(Ordering::Relaxed) {
return MemoryAction::None;
}
if self.inner.ledger.rolling() >= self.inner.config.observation_trigger_tokens {
return MemoryAction::GenerateObservation;
}
if self.inner.observation_tokens.load(Ordering::Relaxed)
>= self.inner.config.reflection_trigger_tokens
{
return MemoryAction::GenerateReflection;
}
MemoryAction::None
}
#[must_use]
pub fn observation_memory_tokens(&self) -> u64 {
self.inner.observation_tokens.load(Ordering::Relaxed)
}
pub fn set_observation_memory_tokens(&self, tokens: u64) {
self.inner.observation_tokens.store(tokens, Ordering::Relaxed);
}
pub async fn persist_observation(
&self,
observations: Vec<crate::types::ObservationEntry>,
token_count: u64,
) -> StorageResult<()> {
let record = SessionRecord::Observation {
id: foundation_compact::ids::new_scru128(),
observations,
token_count,
timestamp: SystemTime::now(),
};
self.inner
.coordinator
.record_async(&self.inner.session_id, &record)
.await?;
self.inner
.observation_tokens
.fetch_add(token_count, Ordering::Relaxed);
self.inner.ledger.reset_rolling();
Ok(())
}
pub async fn persist_reflection(
&self,
reflections: Vec<crate::types::ReflectionEntry>,
observation_token_count_before: u64,
) -> StorageResult<()> {
let reflection_token_count = reflections
.iter()
.map(|r| r.summary.len() as u64 / 4)
.sum::<u64>();
let record = SessionRecord::Reflection {
id: foundation_compact::ids::new_scru128(),
reflections,
generated_at: SystemTime::now(),
observation_token_count_before,
reflection_token_count_after: reflection_token_count,
};
self.inner
.coordinator
.record_async(&self.inner.session_id, &record)
.await?;
self.inner.observation_tokens.store(0, Ordering::Relaxed);
Ok(())
}
pub async fn update_working_memory(
&self,
facts: Vec<MemoryFact>,
prev_version: u64,
) -> StorageResult<()> {
let record = SessionRecord::WorkingMemory {
id: foundation_compact::ids::new_scru128(),
facts,
version: prev_version + 1,
timestamp: SystemTime::now(),
};
self.inner
.coordinator
.record_async(&self.inner.session_id, &record)
.await
}
#[must_use]
pub fn begin_generation(&self) -> bool {
self.inner
.is_generating
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Relaxed)
.is_ok()
}
pub fn end_generation(&self) {
self.inner.is_generating.store(false, Ordering::Release);
}
#[must_use]
pub fn is_generating(&self) -> bool {
self.inner.is_generating.load(Ordering::Acquire)
}
#[must_use]
pub fn memory_model(&self) -> Option<&ModelId> {
self.inner.config.memory_model.as_ref()
}
#[must_use]
pub fn config(&self) -> &MemoryConfig {
&self.inner.config
}
#[must_use]
pub fn session_id(&self) -> &SessionId {
&self.inner.session_id
}
}