Skip to content
Open
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
70 changes: 70 additions & 0 deletions crates/tracedecay-session-runtime/tests/session_store_read_cost.rs
Original file line number Diff line number Diff line change
Expand Up @@ -35,11 +35,13 @@ use tracedecay_domain::{
use tracedecay_global_db::observation::retention::ObservationRetentionConfig;
use tracedecay_global_db::tests::harness::{HostAdmissionTestRuntimeV1, writer_telemetry};
use tracedecay_global_db::{RegisteredGlobalDb, RegisteredGlobalDbLeaseV1};
use tracedecay_host_admission::session_ingest_authority::GlobalDbSessionIngestAuthority;
use tracedecay_host_admission::{HostAdmissionAuthorities, HostAdmissionFacade};
use tracedecay_lcm::LcmRetentionConfig;
use tracedecay_maintenance::retention::registered_store::run_registered_store_retention;
use tracedecay_privacy::{ObservationRecordParseErrorV1, parse_normalized_observation_record_v1};
use tracedecay_runtime_core::background_cpu::ProcessBackgroundCpuV1;
use tracedecay_runtime_core::config::ProfileRoot;
use tracedecay_session_runtime::session_sync::test_harness::{
SessionTemporalRefreshWakeState, run_session_temporal_refresh_pass,
};
Expand All @@ -51,6 +53,9 @@ use tracedecay_sessions::observation::{
CaptureObservationOutcome, CaptureObservationRequest, ObservationCancellation,
};
use tracedecay_sessions::repository_provenance::RepositoryProvenanceAdmissionContext;
use tracedecay_sessions::runtime::{
TranscriptIngestOutcome, ingest_project_sources_for_provider, with_transcript_source_profile,
};
use tracedecay_store::WAL_SOFT_LIMIT_BYTES;

const PROVIDER: &str = "codex";
Expand Down Expand Up @@ -975,6 +980,71 @@ async fn streamed_message_commits_once_per_durability_boundary() {
);
}

/// One project history pass over every host provider, as the session
/// temporal refresh runs it, reading transcripts under `home`.
async fn project_history_pass(
fixture: &DrainFixture,
database: &RegisteredGlobalDbLeaseV1,
home: &Path,
) -> TranscriptIngestOutcome {
let authority = GlobalDbSessionIngestAuthority::new(database.clone()).with_background_cpu(
Arc::new(ProcessBackgroundCpuV1::new(NonZeroUsize::new(4).unwrap())),
);
let shard = &database.binding().shard_id;
with_transcript_source_profile(
ProfileRoot::under_home(home),
ingest_project_sources_for_provider(
&shard.brain_id,
&shard.profile_id,
&authority,
&fixture.project,
Some(fixture.project_id.clone()),
None,
true,
),
)
.await
}

/// The daemon reruns the history pass every idle minute whether or not a
/// host wrote anything. A pass that finds nothing new must leave the store
/// as it was: every commit it makes is WAL every reader and checkpoint then
/// rereads, for as long as the daemon idles.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn idle_history_pass_commits_nothing() {
let _measured = MEASURED.lock().await;
let fixture = DrainFixture::open().await;
let database = fixture
.runtime
.registered_database_lease(HostAdmissionScope::Project)
.unwrap();
let store = session_store_path(&fixture);
let home = fixture._tmp.path().join("home");
std::fs::create_dir_all(&home).unwrap();

let start = wal_mark(&store);
let first = project_history_pass(&fixture, &database, &home).await;
let settled = wal_mark(&store);
let idle = project_history_pass(&fixture, &database, &home).await;
let idled = wal_mark(&store);

assert!(
first.failures.is_empty() && idle.failures.is_empty(),
"history passes must not fail: first={:?} idle={:?}",
first.failures,
idle.failures,
);
assert!(
wal_commits(&store, start, settled).is_some_and(|commits| commits > 0),
"the first pass records each provider's coverage"
);
assert_eq!(
wal_commits(&store, settled, idled),
Some(0),
"a history pass that finds nothing new must commit nothing"
);
}

/// Levels of every B-tree in the fixture's project session store: the pages
/// one point lookup in that tree reads on a cold reader.
fn session_store_tree_depths(fixture: &DrainFixture) -> BTreeMap<String, u32> {
Expand Down
38 changes: 18 additions & 20 deletions crates/tracedecay-sessions/src/runtime/hosts/kimi.rs
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,7 @@ use crate::runtime::snapshot_observation::{
use crate::runtime::source::{
FileDiscoveryLimit, HostProviderCoverage, TranscriptDiscoveryBounds, TranscriptIngestError,
TranscriptIngestResult, bound_path_list, canonical_framed_sha256, jsonl_file_identity,
persist_host_provider_coverage, run_blocking_transcript_section,
persist_host_provider_coverage, revise_host_record, run_blocking_transcript_section,
};

mod discovery;
Expand Down Expand Up @@ -837,29 +837,27 @@ pub async fn capture_kimi_observations(
&& !cancellation.is_cancelled()
{
let next_frontier = if discovery.reached_end {
Some(ParseOffset {
byte_offset: 0,
mtime: discovery_frontier.mtime.saturating_add(1),
file_id: 0,
})
Some((0, 0))
} else {
last_discovered_entry.map(|entry| ParseOffset {
byte_offset: entry.sequence,
mtime: discovery_frontier.mtime.saturating_add(1),
file_id: entry.sequence,
})
last_discovered_entry.map(|entry| (entry.sequence, entry.sequence))
};
if let Some(next_frontier) = next_frontier
if let Some((byte_offset, file_id)) = next_frontier
&& !cancellation.is_cancelled()
{
facade
.advance_parse_offset(&scope, KIMI_DISCOVERY_FRONTIER_KEY, next_frontier)
.await
.map_err(|outcome| {
crate::runtime::snapshot_observation::host_admission_error(
PROVIDER, outcome,
)
})?;
revise_host_record(
facade,
&scope,
KIMI_DISCOVERY_FRONTIER_KEY,
discovery_frontier,
byte_offset,
file_id,
)
.await
.map_err(|outcome| {
crate::runtime::snapshot_observation::host_admission_error(
PROVIDER, outcome,
)
})?;
}
}
let deferred_units = outcome
Expand Down
31 changes: 15 additions & 16 deletions crates/tracedecay-sessions/src/runtime/hosts/pi.rs
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,8 @@ use crate::runtime::snapshot_observation::MAX_SNAPSHOT_METADATA_BYTES;
use crate::runtime::source::{
FileDiscoveryLimit, FileDiscoveryReport, HostProviderCoverage, TranscriptDiscoveryBounds,
TranscriptIngestError, TranscriptIngestResult, bound_path_list, canonical_framed_sha256,
jsonl_file_identity, persist_host_provider_coverage, run_blocking_transcript_section,
jsonl_file_identity, persist_host_provider_coverage, revise_host_record,
run_blocking_transcript_section,
};

/// Environment override Pi reads for its agent directory.
Expand Down Expand Up @@ -657,23 +658,21 @@ pub async fn capture_pi_observations(
&& !cancellation.is_cancelled()
{
let next_frontier = if discovery.reached_end {
Some(ParseOffset {
byte_offset: 0,
mtime: discovery_frontier.mtime.saturating_add(1),
file_id: 0,
})
Some((0, 0))
} else {
last_discovered_entry.map(|entry| ParseOffset {
byte_offset: entry.sequence,
mtime: discovery_frontier.mtime.saturating_add(1),
file_id: entry.sequence,
})
last_discovered_entry.map(|entry| (entry.sequence, entry.sequence))
};
if let Some(next_frontier) = next_frontier {
facade
.advance_parse_offset(&scope, PI_DISCOVERY_FRONTIER_KEY, next_frontier)
.await
.map_err(admission_error)?;
if let Some((byte_offset, file_id)) = next_frontier {
revise_host_record(
facade,
&scope,
PI_DISCOVERY_FRONTIER_KEY,
discovery_frontier,
byte_offset,
file_id,
)
.await
.map_err(admission_error)?;
}
}
persist_coverage(facade, &scope, &outcome).await?;
Expand Down
19 changes: 9 additions & 10 deletions crates/tracedecay-sessions/src/runtime/ingest/project_provider.rs
Original file line number Diff line number Diff line change
Expand Up @@ -392,16 +392,15 @@ impl<'a> ProjectProviderRun<'a> {
} else {
HostProviderCoverage::Complete
};
if stored_coverage != Some(coverage)
&& let Err(error) = persist_host_provider_coverage(
self.facade,
self.scope,
"codex",
coverage,
u64::from(coverage != HostProviderCoverage::Complete),
None,
)
.await
if let Err(error) = persist_host_provider_coverage(
self.facade,
self.scope,
"codex",
coverage,
u64::from(coverage != HostProviderCoverage::Complete),
None,
)
.await
{
outcome.add_failure(warn_transcript_catch_up_failure(
"codex",
Expand Down
19 changes: 9 additions & 10 deletions crates/tracedecay-sessions/src/runtime/ingest/user_provider.rs
Original file line number Diff line number Diff line change
Expand Up @@ -210,16 +210,15 @@ impl<S: TranscriptIngestStore> UserProviderUnit<'_, S> {
} else {
HostProviderCoverage::Complete
};
if stored_coverage != Some(coverage)
&& let Err(coverage_error) = persist_host_provider_coverage(
self.facade,
&ObservationScopeV1::Profile,
"codex",
coverage,
u64::from(coverage != HostProviderCoverage::Complete),
None,
)
.await
if let Err(coverage_error) = persist_host_provider_coverage(
self.facade,
&ObservationScopeV1::Profile,
"codex",
coverage,
u64::from(coverage != HostProviderCoverage::Complete),
None,
)
.await
{
run.add_failure(warn_transcript_catch_up_failure(
"codex",
Expand Down
41 changes: 34 additions & 7 deletions crates/tracedecay-sessions/src/runtime/source.rs
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,7 @@ use tracedecay_store::{ParseOffset, TranscriptStoreError};
pub use tracedecay_domain::canonical_text::canonical_framed_sha256;

pub use super::host_coverage::HostCoverageReason;
use crate::admission::HostAdmission;
use crate::admission::{HostAdmission, HostAdmissionOutcome};
pub use crate::runtime::shared::{NewRows, StoredCursor, TranscriptIngestStats};
use tracedecay_framing::{WireReadOutcome, read_bounded_to_string};

Expand Down Expand Up @@ -168,20 +168,47 @@ pub(super) async fn persist_host_provider_coverage(
crate::runtime::snapshot_observation::host_admission_error(provider, outcome)
})?
.unwrap_or_default();
revise_host_record(
admission,
scope,
&key,
current,
deferred_units,
coverage.file_id(reason),
)
.await
.map_err(|outcome| {
crate::runtime::snapshot_observation::host_admission_error(provider, outcome)
})
}

/// Writes a host bookkeeping record whose readers consult only its
/// `byte_offset` and `file_id`. Its `mtime` is a revision that lets a changed
/// record move `byte_offset` backwards past the monotonic cursor guard. An
/// unchanged record stays as stored, so a sweep that finds nothing new
/// commits nothing.
pub(super) async fn revise_host_record(
admission: &(impl HostAdmission + ?Sized),
scope: &ObservationScopeV1,
key: &str,
current: ParseOffset,
byte_offset: u64,
file_id: u64,
) -> Result<(), HostAdmissionOutcome> {
if (current.byte_offset, current.file_id) == (byte_offset, file_id) {
return Ok(());
}
admission
.advance_parse_offset(
scope,
&key,
key,
ParseOffset {
byte_offset: deferred_units,
byte_offset,
mtime: current.mtime.saturating_add(1).max(1),
file_id: coverage.file_id(reason),
file_id,
},
)
.await
.map_err(|outcome| {
crate::runtime::snapshot_observation::host_admission_error(provider, outcome)
})
}

#[derive(Debug, Error)]
Expand Down
Loading