- 1646
- 1647
pub fn latest_reading(&self) -> Option<crate::turns::ReadingKey> { - 1648
self.chain_to_root() - 1649
.into_iter() - 1650
.rev() - 1651
.find_map(|entry| match &entry.payload { - 1652
EntryPayload::Intent(record) => Some(crate::turns::ReadingKey::from_record(record)), - 1653
_ => None, - 1654
}) - 1655
} - 1656
- 1657
/// The session-derived tail sections for one request. `plan` is the - 1658
/// working-set plan the request will be projected with: the - 1659
/// conversation thread lists only directives that plan leaves out of - 1660
/// the projection (packeted, or behind a reset), so a directive the - 1661
/// model already sees — verbatim in a `Full` turn, as a card's - 1662
/// `asked:` line, or in the open turn — is never restated (one source - 1663
/// per fact, docs/design/68-context-engine.md §6). With `None` every - 1664
/// closed turn is projected at `Full`, so the thread is empty unless a - 1665
/// reset hid something. - 1666
pub fn tail_sections(&self, plan: Option<&WorkingSetPlan>) -> TailSections { - 1667
TailSections { - 1668
intent: self.tail_intent(), - 1669
work_contract: self.tail_work_contract(), - 1670
thread: self.tail_conversation_thread(plan), - 1671
workspace: self.tail_workspace_delta(), - 1672
} - 1673
} - 1674
- 1675
/// The `data.section = "workspace_delta"` activity this turn recorded. - 1676
/// - 1677
/// The value used by the `workspace_delta` activity's `data` key. - 1678
pub const WORKSPACE_DELTA_SECTION: &'static str = "workspace_delta"; - 1679
- 1680
/// The workspace delta recorded for the current turn, tagged — only - 1681
/// when the turn's reading asked for it (`working` / `full`), and only - 1682
/// the activity written after this turn's intent entry, so a previous - 1683
/// turn's delta is never shown as current. - 1684
fn tail_workspace_delta(&self) -> Option<String> { - 1685
let chain = self.chain_to_root(); - 1686
let intent_position = chain - 1687
.iter() - 1688
.rposition(|entry| matches!(entry.payload, EntryPayload::Intent(_)))?; - 1689
let wants = match &chain[intent_position].payload { - 1690
EntryPayload::Intent(record) => { - 1691
crate::turns::ReadingKey::from_record(record).wants_workspace_delta() - 1692
} - 1693
_ => false, - 1694
}; - 1695
if !wants { - 1696
return None; - 1697
} - 1698
let delta = chain[intent_position + 1..] - 1699
.iter() - 1700
.rev() - 1701
.find_map(|entry| match &entry.payload { - 1702
EntryPayload::Activity(activity) - 1703
if activity.data.get("section").map(String::as_str) - 1704
== Some(Self::WORKSPACE_DELTA_SECTION) => - 1705
{ - 1706
activity.detail.clone() - 1707
} - 1708
_ => None, - 1709
})?; - 1710
let trimmed = delta.trim(); - 1711
if trimmed.is_empty() { - 1712
return None; - 1713
} - 1714
// Bounded where it is written (`checkpoints::delta_summary`'s byte - 1715
// budget and excerpt cap), never cut here: a byte truncation panics - 1716
// inside a multi-byte character, and invariant 36 forbids a blind cut. - 1717
Some(format!("<workspace_delta>\n{trimmed}\n</workspace_delta>")) - 1718
} - 1719
- 1720
/// The latest intent note, tagged. Only the newest note applies — it - 1721
/// otherwise wastes context and lets a stale instruction argue with the - 1722
/// current one. - 1723
fn tail_intent(&self) -> Option<String> { - 1724
let note = self - 1725
.chain_to_root() - 1726
.into_iter() - 1727
.rev() - 1728
.find_map(|entry| match &entry.payload { - 1729
EntryPayload::Intent(record) => Some(record.model_visible.clone()), - 1730
_ => None, - 1731
}) - 1732
.flatten()?; - 1733
Some(format!("<intent>\n{note}\n</intent>")) - 1734
} - 1735
- 1736
/// The active work contract's state, tagged, or `None` once it has - 1737
/// settled (completed, failed, cancelled, or unverified). - 1738
fn tail_work_contract(&self) -> Option<String> { - 1739
let work = self.work_projection().ok().flatten()?; - 1740
if matches!( - 1741
work.status, - 1742
crate::types::WorkContractStatus::Completed - 1743
| crate::types::WorkContractStatus::Failed - 1744
| crate::types::WorkContractStatus::Cancelled - 1745
| crate::types::WorkContractStatus::Unverified - 1746
) { - 1747
return None; - 1748
} - 1749
let mut context = format!( - 1750
"<work_contract id=\"{}\" revision=\"{}\">\nObjective: {}\nStatus: {:?}\nItems:\n", - 1751
work.contract.contract_id, work.contract.revision, work.contract.objective, work.status, - 1752
); - 1753
for item in &work.contract.items { - 1754
if let Some(state) = work.items.get(&item.item_id) { - 1755
context.push_str(&format!( - 1756
"- {}: {:?} (owner: {:?})\n", - 1757
item.item_id, state.status, item.owner - 1758
)); - 1759
} - 1760
} - 1761
context.push_str( - 1762
"Rules: use this state for progress; do not claim completion before verification.\n</work_contract>", - 1763
); - 1764
Some(context) - 1765
} - 1766
- 1767
/// The directive texts the projection under `plan` already carries: - 1768
/// the open turn's, every `Full` or `Card` turn's, and — with no plan — - 1769
/// every turn's at or after the reset boundary. - 1770
fn covered_directives( - 1771
&self, - 1772
plan: Option<&WorkingSetPlan>, - 1773
) -> std::collections::HashSet<String> { - 1774
let fidelity_of: HashMap<&str, Fidelity> = plan - 1775
.map(|p| p.per_turn.iter().map(|(id, f)| (id.as_str(), *f)).collect()) - 1776
.unwrap_or_default(); - 1777
TurnIndex::from_log(self) - 1778
.turns - 1779
.iter() - 1780
.filter(|turn| !turn.behind_reset) - 1781
.filter(|turn| { - 1782
!turn.closed - 1783
|| plan.is_none() - 1784
|| matches!( - 1785
fidelity_of.get(turn.id.as_str()), - 1786
Some(Fidelity::Full | Fidelity::Card) - 1787
) - 1788
}) - 1789
.map(|turn| turn.directive.text_content().trim().to_string()) - 1790
.collect() - 1791
} - 1792
- 1793
/// The multi-turn directive timeline, tagged, restricted to directives - 1794
/// the projection under `plan` does not already carry — one source per - 1795
/// fact (docs/design/68-context-engine.md §6): a directive the model - 1796
/// sees in a `Full` turn, on a card line, or in the open turn needs no - 1797
/// restating here; only packeted or reset-hidden ones do. `None` when - 1798
/// there is no active goal spanning multiple turns, a managed work - 1799
/// contract already covers progress, or nothing is left out. - 1800
fn tail_conversation_thread(&self, plan: Option<&WorkingSetPlan>) -> Option<String> { - 1801
let goal = self.goal_state()?; - 1802
if !(goal.revision > 1 || !goal.additions.is_empty()) - 1803
|| goal.control != vak_intent::GoalControlState::Active - 1804
|| goal.objective.trim().is_empty() - 1805
|| self.work_projection().ok().flatten().is_some() - 1806
{ - 1807
return None; - 1808
} - 1809
- 1810
let mut user_directives: Vec<(usize, String)> = Vec::new(); - 1811
let mut turn_counter = 1; - 1812
for entry in self.chain_to_root() { - 1813
if let EntryPayload::Message(record) = &entry.payload - 1814
&& record.message.role == vak_llm::Role::User - 1815
&& record.control_kind().is_none() - 1816
{ - 1817
let text = record.message.text_content(); - 1818
let trimmed = text.trim(); - 1819
if !trimmed.is_empty() && !trimmed.starts_with('<') { - 1820
user_directives.push((turn_counter, trimmed.to_string())); - 1821
turn_counter += 1; - 1822
} - 1823
} - 1824
} - 1825
if user_directives.len() <= 1 { - 1826
return None; - 1827
} - 1828
- 1829
let verbatim: std::collections::HashSet<String> = self.covered_directives(plan); - 1830
let start_idx = user_directives.len().saturating_sub(8); - 1831
let filtered: Vec<&(usize, String)> = user_directives[start_idx..] - 1832
.iter() - 1833
.filter(|(_, req)| !verbatim.contains(req.as_str())) - 1834
.collect(); - 1835
if filtered.is_empty() { - 1836
return None; - 1837
} - 1838
- 1839
// Directive drift (docs/design/68-context-engine.md §7): a listed - 1840
// directive whose reading shares no domain with the CURRENT - 1841
// directive's is marked paused rather than dropped — a - 1842
// re-weighting, not a cut. Turns without a card yet (never closed) - 1843
// have no reading to compare and are never marked. - 1844
let current_domains = self.latest_reading().map(|r| r.domains); - 1845
let index = TurnIndex::from_log(self); - 1846
- 1847
// Only a goal a person stated is presented as the objective. The - 1848
// first message of a conversation is not one by default: rendering - 1849
// "hi" as the primary objective of every later turn misdirects the - 1850
// model more than it orients it. - 1851
let mut thread = format!("<conversation_thread revision=\"{}\">\n", goal.revision); - 1852
if goal.explicit { - 1853
thread.push_str(&format!("Primary objective: {}\n", goal.objective.trim())); - 1854
} - 1855
thread.push_str("User request timeline across turns:\n"); - 1856
for (num, req) in filtered { - 1857
let preview = if req.len() > 200 { - 1858
let head = req - 1859
.char_indices() - 1860
.map(|(idx, _)| idx) - 1861
.nth(200) - 1862
.unwrap_or_else(|| 200.min(req.len())); - 1863
format!("{}...", &req[..head]) - 1864
} else { - 1865
req.clone() - 1866
}; - 1867
let paused = current_domains.as_ref().is_some_and(|current| { - 1868
!current.is_empty() - 1869
&& index.turn_by_number(*num).is_some_and(|turn| { - 1870
turn.card.as_ref().is_some_and(|card| { - 1871
!card.reading.domains.is_empty() - 1872
&& card.reading.domains.iter().all(|d| !current.contains(d)) - 1873
}) - 1874
}) - 1875
}); - 1876
let suffix = if paused { " (earlier, now paused)" } else { "" }; - 1877
thread.push_str(&format!("- Turn {num}: {preview}{suffix}\n")); - 1878
} - 1879
thread.push_str( - 1880
"Rules for multi-turn execution:\n\ - 1881
- Follow the user's intent across conversational drifts without complaint or friction.\n\ - 1882
- Conversational drift across turns is expected: follow along smoothly and adapt immediately.\n\ - 1883
- Resolve references (\"the data\", \"do that\", \"it\", \"something\") against the request timeline above.\n\ - 1884
- If genuinely confused, ask a brief clarification, but NEVER use asking clarification as an exception-handling escape hatch to avoid taking action or using available tools.\n\ - 1885
</conversation_thread>", - 1886
); - 1887
Some(thread) - 1888
} - 1889
- 1890
/// Usage summed over the ACTIVE chain only: usage recorded on - 1891
/// abandoned branches (superseded by branching/compaction) must not - 1892
/// inflate the count. - 1893
pub fn total_usage(&self) -> vak_llm::Usage { - 1894
let mut total = vak_llm::Usage::default(); - 1895
for e in self.chain_to_root() { - 1896
if let EntryPayload::Message(MessageRecord { - 1897
meta: Some(MessageMeta { usage: Some(u), .. }), - 1898
.. - 1899
}) = &e.payload - 1900
{ - 1901
total.input_tokens += u.input_tokens; - 1902
total.output_tokens += u.output_tokens; - 1903
} - 1904
} - 1905
total - 1906
} - 1907
- 1908
/// Proactive retrieval: find the most relevant older entries from the - 1909
/// projected history, returning their entry IDs and messages. - 1910
/// - 1911
/// This is a projection-only operation — it returns entries that are - 1912
/// already in the projection. No ledger mutation, no new entries. - 1913
/// - 1914
/// Scoring uses the same BM25+entity-bonus approach as cross-session - 1915
/// search, but scored against the current query over the in-memory - 1916
/// projection. The result is capped at `cap` entries and excludes the - 1917
/// `exclude_tail` most recent entries. - 1918
pub fn retrieve_relevant_entries( - 1919
&self, - 1920
query: &str, - 1921
cap: usize, - 1922
exclude_tail: usize, - 1923
) -> Vec<(String, vak_llm::Message)> { - 1924
if query.trim().is_empty() || cap == 0 { - 1925
return Vec::new(); - 1926
} - 1927
let tagged = self.derive_keyed_tagged(); - 1928
if tagged.is_empty() { - 1929
return Vec::new(); - 1930
} - 1931
- 1932
let terms: Vec<String> = crate::search::tokenize_impl(query); - 1933
let phrase = crate::search::normalize_impl(query); - 1934
if terms.is_empty() || phrase.is_empty() { - 1935
return Vec::new(); - 1936
} - 1937
- 1938
let candidates: Vec<_> = if tagged.len() > exclude_tail { - 1939
tagged[..tagged.len() - exclude_tail] - 1940
.iter() - 1941
.filter(|(_, _, is_summary, is_control)| !*is_summary && !*is_control) - 1942
.collect() - 1943
} else { - 1944
Vec::new() - 1945
}; - 1946
- 1947
if candidates.is_empty() { - 1948
return Vec::new(); - 1949
} - 1950
- 1951
let mut scored: Vec<(f32, &String, &vak_llm::Message)> = Vec::new(); - 1952
for (id, msg, _, _) in &candidates { - 1953
let text = msg.text_content(); - 1954
let entities = crate::search::extract_entities(&text); - 1955
let normalized = crate::search::normalize_impl(&text); - 1956
let score = crate::search::score_normalized(&normalized, &terms, &phrase, &entities); - 1957
if score > 0.0 { - 1958
scored.push((score, id, msg)); - 1959
} - 1960
} - 1961
- 1962
// Sort by score descending, take the top `cap`. - 1963
scored.sort_by(|a, b| b.0.total_cmp(&a.0)); - 1964
scored - 1965
.into_iter() - 1966
.take(cap) - 1967
.map(|(_, id, m)| (id.clone(), m.clone())) - 1968
.collect() - 1969
} - 1970
} - 1971
- 1972
pub struct SessionPath; - 1973
- 1974
impl SessionPath { - 1975
pub fn sessions_dir(home: &Path, cwd: &Path) -> PathBuf { - 1976
home.join("sessions").join(hash_cwd(cwd)) - 1977
} - 1978
- 1979
pub fn new_session_file(home: &Path, cwd: &Path, session_id: &str) -> PathBuf { - 1980
Self::sessions_dir(home, cwd).join(format!("{session_id}.jsonl")) - 1981
} - 1982
} - 1983
- 1984
/// FNV-1a: a fixed hash whose output does not change across Rust releases, - 1985
/// unlike DefaultHasher (SipHash with randomly-seeded-but-toolchain-chosen - 1986
/// parameters). Session directories must remain reachable after toolchain - 1987
/// upgrades. - 1988
fn hash_cwd(cwd: &Path) -> String { - 1989
let mut h: u64 = 0xcbf2_9ce4_8422_2325; - 1990
for b in cwd.to_string_lossy().as_bytes() { - 1991
h ^= u64::from(*b); - 1992
h = h.wrapping_mul(0x0000_0100_0000_01b3); - 1993
} - 1994
format!("{h:016x}") - 1995
} - 1996
- 1997
#[cfg(all(test, unix))] - 1998
mod tests { - 1999
#![allow(clippy::unwrap_used, clippy::expect_used)] - 2000
use super::*; - 2001
- 2002
/// A child being spawned holds a duplicate of every open descriptor - 2003
/// until it execs; the duplicate here stands in for one. - 2004
#[test] - 2005
fn the_lock_lasts_exactly_as_long_as_the_handle() { - 2006
let dir = tempfile::tempdir().unwrap(); - 2007
let path = dir.path().join("session.jsonl"); - 2008
std::fs::write(&path, "").unwrap(); - 2009
- 2010
let log = SessionLog::open(path.clone()).unwrap(); - 2011
assert!(matches!( - 2012
SessionLog::open(path.clone()), - 2013
Err(SessionError::Locked(_)) - 2014
)); - 2015
- 2016
let duplicate = log.file.file.try_clone().unwrap(); - 2017
drop(log); - 2018
SessionLog::open(path).expect("the dropped handle's lock stayed with a duplicate"); - 2019
drop(duplicate); - 2020
} - 2021
} - 2022
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.