Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
26 commits
Select commit Hold shift + click to select a range
3369aad
fix(sessions): admit newest Codex rollouts before older days
ScriptedAlchemy Sep 28, 2026
4c31655
fix(sessions)!: return deferred import while catch-up continues
ScriptedAlchemy Sep 28, 2026
283b48e
fix(sessions): publish newest Codex day before older days
ScriptedAlchemy Sep 28, 2026
bc7d42b
fix(sessions): publish a yielded Codex window before backoff
ScriptedAlchemy Sep 28, 2026
e9f7b00
fix(sessions): hold projection when a yield reports no stats
ScriptedAlchemy Sep 28, 2026
fcdf3d1
fix(sessions): keep projection ahead of a queued history wake
ScriptedAlchemy Sep 28, 2026
5a54c69
Merge remote-tracking branch 'origin/master' into fleet/land-cloud-pr…
ScriptedAlchemy Sep 28, 2026
dbc415c
fix(sessions): report scheduled imports as deferred, yield per day
ScriptedAlchemy Sep 28, 2026
8f1fd32
fix(sessions): yield only at a day that opens out of project
ScriptedAlchemy Sep 28, 2026
38bd500
Merge remote-tracking branch 'origin/master' into fleet/land-cloud-pr…
ScriptedAlchemy Sep 28, 2026
3bc348e
test(sessions): assert the scheduled import receipt, then search
ScriptedAlchemy Sep 28, 2026
f328496
test(work): wait for imported transcripts before journey reads
ScriptedAlchemy Sep 28, 2026
53dacc3
style(cli): format session sync imports
ScriptedAlchemy Sep 28, 2026
d09b675
fix(ci): install Windows target on pinned toolchain
ScriptedAlchemy Sep 29, 2026
cc6cc21
Merge remote-tracking branch 'origin/master' into cursor/codex-catchu…
ScriptedAlchemy Sep 29, 2026
715ea70
Merge remote-tracking branch 'origin/master' into cursor/codex-catchu…
ScriptedAlchemy Sep 29, 2026
2c06c0b
Merge remote-tracking branch 'origin/master' into cursor/codex-catchu…
ScriptedAlchemy Sep 29, 2026
462471a
fix(sessions): keep a live Codex tail from ending every catch-up pass
ScriptedAlchemy Sep 29, 2026
4179085
test(session-runtime): loosen import hand-off bound for loaded runners
ScriptedAlchemy Sep 29, 2026
a846b16
fix(sessions): avoid complete zero during import catch-up
ScriptedAlchemy Sep 29, 2026
6853963
Merge remote-tracking branch 'origin/master' into fleet/land-cloud-pr…
ScriptedAlchemy Sep 29, 2026
39d99e1
test(sessions): pin catch-up empty answers and hand-off without clocks
ScriptedAlchemy Sep 29, 2026
e6f5037
test(sessions): share one fixed serving-status fixture
ScriptedAlchemy Sep 29, 2026
adf071c
test(cli): wait for code-index readiness before graph reads
ScriptedAlchemy Sep 29, 2026
a576d44
Merge remote-tracking branch 'origin/master' into fleet/land-cloud-pr…
ScriptedAlchemy Sep 29, 2026
7dd58c5
Merge remote-tracking branch 'origin/master' into fleet/land-cloud-pr…
ScriptedAlchemy Sep 29, 2026
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
63 changes: 62 additions & 1 deletion crates/tracedecay-cli/src/sessions_cmd/session_sync.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
use std::path::Path;
use tracedecay_contracts::retrieval::{AdminCliSessionSyncV1, AdminCliSurfaceRequestV1};
use tracedecay_contracts::session_sync::SessionSyncSourceCoverageV1;
use tracedecay_contracts::session_sync::{SessionSyncCoverageV1, SessionSyncSourceCoverageV1};
use tracedecay_contracts::{IdempotencyKey, OperationTermination, RequestId};
use tracedecay_runtime_core::config::ProfileRoot;

Expand Down Expand Up @@ -108,6 +108,19 @@ pub(super) fn session_sync_poll_state(
),
}
})?;
if session_import_deferred_progress(
label,
termination,
&coverage,
&failure_codes,
remaining_work,
) {
println!(
"{label} scheduled ({}); historical catch-up has remaining work {remaining_work}",
operation_id.as_str()
);
return Ok(SessionSyncPollState::Completed);
}
if termination != OperationTermination::Completed || remaining_work > 0 {
let termination = termination_label(termination);
let detail = if failure_codes.is_empty() {
Expand Down Expand Up @@ -144,6 +157,25 @@ fn termination_label(termination: OperationTermination) -> String {
}
}

fn session_import_deferred_progress(
label: &str,
termination: OperationTermination,
coverage: &[SessionSyncSourceCoverageV1],
failure_codes: &[String],
remaining_work: u64,
) -> bool {
label == "session import"
&& termination == OperationTermination::Partial
&& failure_codes.is_empty()
&& remaining_work > 0
&& coverage.iter().all(|entry| {
matches!(
entry.coverage,
SessionSyncCoverageV1::Complete | SessionSyncCoverageV1::Partial { .. }
)
})
}

fn session_sync_remaining_work(coverage: &[SessionSyncSourceCoverageV1]) -> Option<u64> {
if coverage.is_empty() {
return None;
Expand Down Expand Up @@ -379,6 +411,35 @@ mod tests {
}
}

#[test]
fn session_import_accepts_deferred_catch_up_without_treating_it_as_failure() {
let outcome = complete(
OperationTermination::Partial,
vec![SessionSyncCoverageV1::Partial { deferred_units: 1 }],
&[],
);

assert!(matches!(
session_sync_poll_state("session import", outcome).unwrap(),
SessionSyncPollState::Completed
));
}

#[test]
fn session_git_sync_still_rejects_unfinished_coverage() {
let error = session_sync_poll_state(
"session git sync",
complete(
OperationTermination::Partial,
vec![SessionSyncCoverageV1::Partial { deferred_units: 1 }],
&[],
),
)
.expect_err("git sync still requires the bounded pass to finish");

assert!(error.to_string().contains("remaining work"));
}

#[test]
fn session_sync_reports_remaining_coverage_even_if_daemon_mislabels_completion() {
let error = session_sync_poll_state(
Expand Down
4 changes: 2 additions & 2 deletions crates/tracedecay-cli/tests/core_cli_suite/tracedecay_test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -130,8 +130,8 @@ fn project_path_flags_resolve_a_relative_path_from_inside_the_project() {
let stdout = String::from_utf8_lossy(&import.stdout);
assert!(
import.status.success()
&& stdout.starts_with("session import completed (")
&& stdout.ends_with(")\n"),
&& stdout.starts_with("session import scheduled (")
&& stdout.ends_with("); historical catch-up has remaining work 2\n"),
"sessions import --project-path . must import into the CLI's project\nstdout:\n{stdout}\nstderr:\n{}",
String::from_utf8_lossy(&import.stderr)
);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,12 +18,13 @@ use tracedecay_session_memory::context::{
CapabilityDigest, ConfigurationDigest, PolicyDigest, RequestBudgets, ResolvedSessionIdentity,
};
use tracedecay_session_memory::session::{
SessionRequestBinding, SessionRetrievalConfiguration, SessionTemporalQuery,
TaskSessionRetrievalOutcomeV1,
SessionDataFreshness, SessionRequestBinding, SessionRetrievalConfiguration,
SessionTemporalQuery, TaskSessionRetrievalOutcomeV1,
};
use tracedecay_session_temporal_store::execution::TaskSessionRankSelectorV1;
use tracedecay_sessions::serving::{
RefreshWorkerMissing, SessionProjectionServingStatus, SessionProjectionServingStatusPort,
RefreshWorkerMissing, SessionProjectionServingState, SessionProjectionServingStatus,
SessionProjectionServingStatusPort,
};
use tracedecay_store::StoreShardScopeV1;

Expand Down Expand Up @@ -265,6 +266,10 @@ impl SessionApplicationRetrievalPortV1 for DaemonSessionRetrievalService {
Ok(binding) => binding,
Err(outcome) => return *outcome,
};
let history_pending = matches!(
self.refresh_status.serving_status().state,
SessionProjectionServingState::Stale { .. }
);
let outcome = self
.execute_temporal_query_with_context(
context,
Expand All @@ -273,7 +278,28 @@ impl SessionApplicationRetrievalPortV1 for DaemonSessionRetrievalService {
"grant.application.session-retrieval",
)
.await;
self.public_outcome(outcome).await
match self.public_outcome(outcome).await {
SessionRetrievalServiceOutcome::CompleteZero {
temporal,
freshness,
} if history_pending
|| matches!(
self.refresh_status.serving_status().state,
SessionProjectionServingState::Stale { .. }
) =>
{
SessionRetrievalServiceOutcome::Stale {
temporal,
freshness: match freshness {
SessionDataFreshness::Fresh => {
SessionDataFreshness::Stored { generation_lag: 0 }
}
other => other,
},
}
}
outcome => outcome,
}
})
}

Expand Down
91 changes: 76 additions & 15 deletions crates/tracedecay-session-runtime/src/session_retrieval/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ use tracedecay_domain::{
};
use tracedecay_lcm::contracts::{LcmDataFreshness, LcmRetrievalOutcome};
use tracedecay_session_temporal_store::SessionTemporalAccess;
use tracedecay_sessions::serving::{SessionProjectionServingState, SessionProjectionStaleReason};
use tracedecay_store::{
AnchoredObservationWrite, ObservationProjectionStore, ObservationStore, ObservationWrite,
SessionRecord, SessionTemporalSnapshotRequestV1, build_observation_resolution_authorization_v1,
Expand Down Expand Up @@ -1402,20 +1403,6 @@ async fn describe_without_a_refresh_worker_does_not_pretend_history_is_convergin
}
}

struct CurrentRefreshServing;

impl tracedecay_sessions::serving::SessionProjectionServingStatusPort for CurrentRefreshServing {
fn serving_status(&self) -> tracedecay_sessions::serving::SessionProjectionServingStatus {
tracedecay_sessions::serving::SessionProjectionServingStatus {
state: tracedecay_sessions::serving::SessionProjectionServingState::Current,
last_progress_at_unix_micros: None,
backlog: 0,
blocker: None,
retry_class: None,
}
}
}

/// Profile catch-up mounts the real refresh worker's serving-status port. When
/// that port reports current, `RequireFresh` must not be refused as
/// `RefreshWorkerMissing`.
Expand Down Expand Up @@ -1444,7 +1431,9 @@ async fn require_fresh_with_a_current_refresh_worker_is_not_refused_as_worker_mi
let service = DaemonSessionRetrievalService::new_admitted_profile(
harness.registered.clone(),
root.identity().clone(),
Some(std::sync::Arc::new(CurrentRefreshServing)),
Some(std::sync::Arc::new(FixedRefreshServing(
SessionProjectionServingState::Current,
))),
)
.expect("registered retrieval service");
let context = admitted_lookup_context(scope);
Expand Down Expand Up @@ -1482,6 +1471,78 @@ async fn require_fresh_with_a_current_refresh_worker_is_not_refused_as_worker_mi
}
}

struct FixedRefreshServing(SessionProjectionServingState);

impl tracedecay_sessions::serving::SessionProjectionServingStatusPort for FixedRefreshServing {
fn serving_status(&self) -> tracedecay_sessions::serving::SessionProjectionServingStatus {
tracedecay_sessions::serving::SessionProjectionServingStatus {
state: self.0.clone(),
last_progress_at_unix_micros: None,
backlog: 0,
blocker: None,
retry_class: None,
}
}
}

/// An empty answer is complete only once historical catch-up is current.
/// While catch-up is converging the same empty store must answer stale, or a
/// search right after a scheduled import claims nothing matches (#2512).
#[tokio::test]
async fn empty_answer_is_stale_until_historical_catch_up_is_current() {
let harness = tracedecay_global_db::tests::harness::RegisteredGlobalDbHarness::open(
"session-retrieval-empty-during-catch-up",
)
.await;
let root = registered_profile_retrieval_root(&harness.registered);
let scope = root
.identity()
.session_request_scope()
.expect("profile session scope");
let context = admitted_lookup_context(scope);
let answer = |state: SessionProjectionServingState| {
let service = DaemonSessionRetrievalService::new_admitted_profile(
harness.registered.clone(),
root.identity().clone(),
Some(std::sync::Arc::new(FixedRefreshServing(state))),
)
.expect("registered retrieval service");
let query = SessionTemporalQuery::new(
SessionId::new("session.empty.catch-up").expect("session identity"),
None,
"",
None,
TemporalModeV1::Current,
tracedecay_domain::RetrievalGrainV1::Occurrence,
1,
DiversityLimits::unbounded(),
ContextBudget {
max_bytes: APPLICATION_RETRIEVAL_MAX_BYTES,
max_tokens: APPLICATION_RETRIEVAL_MAX_BYTES / 4,
estimator_version: "words-v1".to_owned(),
},
)
.expect("temporal query")
.with_execution_limits(admitted_execution_limits(1));
let context = &context;
async move { service.retrieve_admitted(context, query).await }
};

let current = answer(SessionProjectionServingState::Current).await;
assert!(
matches!(current, SessionRetrievalServiceOutcome::CompleteZero { .. }),
"a current empty store answers complete zero: {current:?}"
);
let converging = answer(SessionProjectionServingState::Stale {
reason: SessionProjectionStaleReason::HistoricalConvergence,
})
.await;
assert!(
matches!(converging, SessionRetrievalServiceOutcome::Stale { .. }),
"an empty store during catch-up must not claim completeness: {converging:?}"
);
}

/// yielding a continuation while records remain.
#[tokio::test]
async fn small_lookup_reads_a_session_larger_than_the_response_budget() {
Expand Down
Loading
Loading