Skip to content
Merged
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
6 changes: 6 additions & 0 deletions modules/frontend/src/conversation_host.rs
Original file line number Diff line number Diff line change
Expand Up @@ -311,6 +311,12 @@ impl ConversationHost {
self.controller.snapshot().cloned()
}

/// Identities of the registered observation-derived facts.
#[must_use]
pub fn derived_fact_ids(&self) -> Vec<crate::conversation_scene::SceneId> {
self.controller.derived_fact_ids()
}

/// Returns the one child surface entity.
#[must_use]
pub fn surface(&self) -> &Entity<ConversationSurface> {
Expand Down
32 changes: 24 additions & 8 deletions modules/frontend/src/conversation_observation_projection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -43,15 +43,22 @@ use artisan_domain::{
ConversationLifecycle, ConversationSnapshot, RunId, TerminalActivityState, ToolAction, TurnId,
};

use crate::conversation_scene::{SCENE_ID_MAX_BYTES, SCENE_MAX_TEXT_BYTES, SceneId};
use crate::conversation_scene::{
SCENE_ID_MAX_BYTES, SCENE_MAX_ITEMS, SCENE_MAX_TEXT_BYTES, SceneId,
};
use crate::conversation_state_machine::{SceneFact, SceneFactKind};
use crate::engine_observation_state::{EngineObservationState, TimelineRow};

/// Maximum activity facts projected in one call.
/// Number of activity facts one projection keeps: the scene room the
/// snapshot's durable items leave.
///
/// Bounded so a pathological retained backlog cannot exceed the aggregate
/// fact registry in one replay; the caller replays again for the remainder.
pub const MAX_PROJECTED_FACTS: usize = 256;
/// Activity slides like the rest of the loaded conversation: the newest
/// facts always fit, and the oldest leave the window first. A long thread
/// never stalls its live turn behind history it can no longer show.
#[must_use]
pub fn activity_window(snapshot: &ConversationSnapshot) -> usize {
SCENE_MAX_ITEMS.saturating_sub(snapshot.items().len())
}

/// Maximum UTF-8 bytes retained in one cumulative reasoning body.
pub const MAX_CUMULATIVE_REASONING_BYTES: usize = SCENE_MAX_TEXT_BYTES;
Expand All @@ -76,6 +83,9 @@ pub struct ActivityProjection {
/// Attributed rows skipped as foreign (thread mismatch is already
/// filtered by the state; run mismatches against settled snapshot runs).
pub rejected: usize,
/// Oldest projectable rows that slid out of the [`activity_window`].
/// Facts registered for them earlier are no longer part of the window.
pub evicted: usize,
}

/// Builds the stable run-scoped scene id for one provider row.
Expand Down Expand Up @@ -137,6 +147,8 @@ pub fn truncate_bounded(text: &str, maximum: usize) -> String {
/// - Ordinals start above the durable watermark (`max durable ordinal + 1`)
/// in `(committed_at, delivery_sequence)` order, so scene order is stable
/// without clock sampling and never collides with durable items.
/// - Only the newest [`activity_window`] facts are kept; older rows count as
/// `evicted`.
/// - Bodies are truncated to scene bounds; only the public reasoning summary
/// is retained.
#[must_use]
Expand All @@ -153,6 +165,7 @@ pub fn project_activities(
facts: Vec::new(),
pending: 0,
rejected: 0,
evicted: 0,
};
}
let known_turns: BTreeSet<&TurnId> =
Expand Down Expand Up @@ -405,10 +418,12 @@ pub fn project_activities(
(left.committed_at_ms, left.delivery_sequence)
.cmp(&(right.committed_at_ms, right.delivery_sequence))
});
candidates.truncate(MAX_PROJECTED_FACTS);
let evicted = candidates.len().saturating_sub(activity_window(snapshot));

let mut facts = Vec::with_capacity(candidates.len());
for (index, candidate) in candidates.into_iter().enumerate() {
// Ordinals count every projectable row, evicted ones included, so a fact
// keeps its place while the window slides past older rows.
let mut facts = Vec::with_capacity(candidates.len() - evicted);
for (index, candidate) in candidates.into_iter().enumerate().skip(evicted) {
let ordinal = durable_watermark
.saturating_add(1)
.saturating_add(index as u64);
Expand Down Expand Up @@ -442,6 +457,7 @@ pub fn project_activities(
facts,
pending,
rejected,
evicted,
}
}

Expand Down
11 changes: 11 additions & 0 deletions modules/frontend/src/conversation_state_machine/impl_dispatch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -273,6 +273,17 @@ impl ConversationStateController {
}))
}

/// Returns the identities of registered facts derived from retained
/// engine observations, in identity order.
#[must_use]
pub fn derived_fact_ids(&self) -> Vec<SceneId> {
self.facts
.values()
.filter(|fact| fact.derived())
.map(|fact| fact.id.clone())
.collect()
}

/// Atomically inserts or updates one bounded non-durable fact.
///
/// An identical fact is a no-op without effects; a changed fact for the
Expand Down
21 changes: 20 additions & 1 deletion modules/frontend/src/native_application/impl_service_events.rs
Original file line number Diff line number Diff line change
Expand Up @@ -828,10 +828,29 @@ impl NativeApplication {
};
let projection =
crate::conversation_observation_projection::project_activities(state, &snapshot);
if projection.facts.is_empty() {
// Facts that slid out of the window leave first, so the scene room
// they held goes to the newest activity.
let window: std::collections::BTreeSet<&crate::conversation_scene::SceneId> =
projection.facts.iter().map(|fact| &fact.id).collect();
let slid_out: Vec<_> = host
.read(cx)
.derived_fact_ids()
.into_iter()
.filter(|id| !window.contains(id))
.collect();
if projection.facts.is_empty() && slid_out.is_empty() {
return;
}
let mut invalidated = false;
for id in slid_out {
let remove = crate::conversation_state_machine::SceneFactCommand::Remove { id };
match host.update(cx, |host, host_cx| {
host.dispatch(ConversationStateEvent::Fact(remove), host_cx)
}) {
Ok(()) => invalidated = true,
Err(_) => break,
}
}
for fact in projection.facts {
let upsert = crate::conversation_state_machine::SceneFactCommand::Upsert(fact);
match host.update(cx, |host, host_cx| {
Expand Down
95 changes: 94 additions & 1 deletion tests/frontend/conversation_observation_projection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,8 @@ use artisan_domain::{
use artisan_frontend::conversation_delivery_machine::{
ConversationDeliveryEffect, ConversationDeliveryEvent,
};
use artisan_frontend::conversation_observation_projection::project_activities;
use artisan_frontend::conversation_observation_projection::{activity_window, project_activities};
use artisan_frontend::conversation_scene::SCENE_MAX_ITEMS;
use artisan_frontend::conversation_scene::{
SceneId, TurnBlock, TurnNarration as SceneTurnNarration,
};
Expand Down Expand Up @@ -1342,3 +1343,95 @@ fn work_keeps_its_first_position_between_assistant_messages_after_completion_and
Some(ConversationLifecycle::Failed)
);
}

#[test]
fn activity_slides_to_the_newest_window_and_frees_room_for_live_work() {
let snapshot = snapshot(
vec![make_turn(TURN_A, 0, ConversationLifecycle::Active)],
vec![make_user("user_a", TURN_A, 1)],
);
let window = activity_window(&snapshot);
assert_eq!(
window,
SCENE_MAX_ITEMS - 1,
"the durable item keeps its room"
);
let mut state = EngineObservationState::new(thread_id());
let apply_tool = |state: &mut EngineObservationState, index: usize| {
let cursor = u64::try_from(index).expect("small index") + 1;
let outcome = state.apply(
cursor,
&attributed_event(
tool_observation(
&format!("obs-tool-{index}"),
cursor,
&format!("tool-{index}"),
ToolAction::Completed,
),
RUN_A,
TURN_A,
1_000 + i64::try_from(index).expect("small index"),
cursor,
),
);
assert!(matches!(outcome, ApplyOutcome::Applied { .. }));
};
for index in 0..window + 3 {
apply_tool(&mut state, index);
}

// The newest rows fill the window; the three oldest slid out.
let projection = project_activities(&state, &snapshot);
assert_eq!(projection.facts.len(), window);
assert_eq!(projection.evicted, 3);
let first = projection
.facts
.first()
.expect("window")
.id
.as_str()
.to_owned();
let last = projection
.facts
.last()
.expect("window")
.id
.as_str()
.to_owned();
assert!(
first.ends_with("-tool-3"),
"oldest kept row is the fourth: {first}"
);
assert!(
last.ends_with(&format!("-tool-{}", window + 2)),
"newest row is kept: {last}"
);

let mut controller = controller_with_snapshot(snapshot.clone());
upsert_all(&mut controller, projection.facts);
assert_eq!(controller.derived_fact_ids().len(), window);

// Live work keeps arriving: the fact that slid out leaves, and the new
// row takes its room instead of stalling at the scene bound.
apply_tool(&mut state, window + 3);
let projection = project_activities(&state, &snapshot);
let kept: std::collections::BTreeSet<_> = projection
.facts
.iter()
.map(|fact| fact.id.clone())
.collect();
for id in controller.derived_fact_ids() {
if !kept.contains(&id) {
controller.remove_fact(id).expect("slid-out fact leaves");
}
}
upsert_all(&mut controller, projection.facts);
let ids = controller.derived_fact_ids();
assert_eq!(ids.len(), window);
assert!(
ids.iter()
.any(|id| id.as_str().ends_with(&format!("-tool-{}", window + 3))),
"the newest row is shown"
);
assert!(!ids.iter().any(|id| id.as_str().ends_with("-tool-3")));
}
Loading